diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 764300076b..14cb8c54e9 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -453,6 +453,26 @@ async fn get_mem_peak(pid: Option, nsjail: bool) -> i32 { } } +pub fn sizeof_val(v: &serde_json::Value) -> usize { + std::mem::size_of::() + + match v { + serde_json::Value::Null => 0, + serde_json::Value::Bool(_) => 0, + serde_json::Value::Number(_) => 4, // Incorrect if arbitrary_precision is enabled. oh well + serde_json::Value::String(s) => s.capacity(), + serde_json::Value::Array(a) => a.iter().map(sizeof_val).sum(), + serde_json::Value::Object(o) => o + .iter() + .map(|(k, v)| { + std::mem::size_of::() + + k.capacity() + + sizeof_val(v) + + std::mem::size_of::() * 3 + }) + .sum(), + } +} + pub async fn run_future_with_polling_update_job_poller( job_id: Uuid, timeout: Option, diff --git a/backend/windmill-worker/src/pg_executor.rs b/backend/windmill-worker/src/pg_executor.rs index 36e1b1c379..1ecd6ab9d9 100644 --- a/backend/windmill-worker/src/pg_executor.rs +++ b/backend/windmill-worker/src/pg_executor.rs @@ -31,8 +31,8 @@ use windmill_parser::Typ; use windmill_parser_sql::{parse_db_resource, parse_pgsql_sig}; use windmill_queue::CanceledBy; -use crate::common::{build_args_values, run_future_with_polling_update_job_poller}; -use crate::AuthedClientBackgroundTask; +use crate::common::{build_args_values, run_future_with_polling_update_job_poller, sizeof_val}; +use crate::{AuthedClientBackgroundTask, MAX_RESULT_SIZE}; use bytes::{Buf, BytesMut}; use lazy_static::lazy_static; use urlencoding::encode; @@ -230,15 +230,31 @@ pub async fn do_postgresql( .unwrap_or_default(), ); - let result = rows - .into_iter() - .map(|x: Row| postgres_row_to_json_value(x)) - .collect::, _>>()?; + let mut siz = 0; + let mut res: Vec = vec![]; + for row in rows.into_iter() { + let r = postgres_row_to_json_value(row); + if let Ok(v) = r.as_ref() { + let size = sizeof_val(v); + siz += size; + } + if *CLOUD_HOSTED && siz > MAX_RESULT_SIZE { + return Err(anyhow::anyhow!( + "Query result too large for cloud (size > {})", + MAX_RESULT_SIZE + )); + } + if let Ok(v) = r { + res.push(v); + } else { + return Err(to_anyhow(r.err().unwrap())); + } + } - Ok(result) + Ok((res, siz)) }; - let result = run_future_with_polling_update_job_poller( + let (result, size) = run_future_with_polling_update_job_poller( job.id, job.timeout, db, @@ -250,6 +266,8 @@ pub async fn do_postgresql( ) .await?; + *mem_peak = size as i32; + RUNNING.store(false, std::sync::atomic::Ordering::Relaxed); if let Some(handle) = handle {