From 14801fdd68407add824bfa2c679c0e8cb4788610 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Fri, 2 Oct 2026 17:21:20 +0200 Subject: [PATCH] fix: unstick postgres jobs on large results with custom-typed columns (#11475) * fix: unstick postgres jobs on large results with custom-typed columns Co-Authored-By: Claude Opus 5.5 (1M context) * fix: carry the parked-rows delivery fix in the postgres fork Co-Authored-By: Claude Opus 5.5 (1M context) * fix: describe postgres queries before streaming instead of parking rows Co-Authored-By: Claude Opus 5.5 (1M context) * perf: skip the postgres describe for statements that return no rows Co-Authored-By: Claude Opus 5.5 (1M context) * fix: look for returning in the whole postgres statement Co-Authored-By: Claude Opus 5.5 (1M context) * fix: read the leading postgres keyword past nested block comments Co-Authored-By: Claude Opus 5.5 (1M context) * fix: end leading postgres line comments at CR too Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) --- backend/Cargo.lock | 105 ++++++++++++------ backend/Cargo.toml | 47 +++----- backend/windmill-worker/src/pg_executor.rs | 122 ++++++++++++++++++++- 3 files changed, 206 insertions(+), 68 deletions(-) diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 271ec93546..7f660b6b51 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -1254,7 +1254,7 @@ dependencies = [ "bytes", "form_urlencoded", "hex", - "hmac", + "hmac 0.12.1", "http 0.2.12", "http 1.5.0", "percent-encoding", @@ -2489,7 +2489,7 @@ version = "0.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7b77c319abfd5219629c45c34c89ba945ed3c5e49fcde9d16b6c3885f118a730" dependencies = [ - "const-oid", + "const-oid 0.9.6", "der", "spki", "x509-cert", @@ -2568,6 +2568,12 @@ version = "0.9.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8" +[[package]] +name = "const-oid" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6ef517f0926dd24a1582492c791b6a4818a4d94e789a334894aa15b0d12f55c" + [[package]] name = "const-random" version = "0.1.18" @@ -3894,7 +3900,7 @@ dependencies = [ "aws-lc-rs", "base64 0.22.1", "cbc", - "const-oid", + "const-oid 0.9.6", "ctr", "curve25519-dalek", "deno_core", @@ -4306,7 +4312,7 @@ version = "0.7.10" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e7c1832837b905bbfb5101e07cc24c8deddf52f93225eee6ead5f4d63d53ddcb" dependencies = [ - "const-oid", + "const-oid 0.9.6", "der_derive", "flagset", "pem-rfc7468", @@ -4501,7 +4507,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" dependencies = [ "block-buffer 0.10.4", - "const-oid", + "const-oid 0.9.6", "crypto-common 0.1.7", "subtle", ] @@ -4513,6 +4519,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2" dependencies = [ "block-buffer 0.12.1", + "const-oid 0.10.2", "crypto-common 0.2.2", "ctutils", ] @@ -5675,7 +5682,7 @@ dependencies = [ "google-cloud-metadata", "google-cloud-token", "hex", - "hmac", + "hmac 0.12.1", "home", "jsonwebtoken 9.3.1", "path-clean", @@ -6042,7 +6049,7 @@ version = "0.12.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7b5f8eb2ad728638ea2c7d47a21db23b7b58a72ed6a38256b8a1849f15fbbdf7" dependencies = [ - "hmac", + "hmac 0.12.1", ] [[package]] @@ -6054,6 +6061,15 @@ dependencies = [ "digest 0.10.7", ] +[[package]] +name = "hmac" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6303bc9732ae41b04cb554b844a762b4115a61bfaa81e3e83050991eeb56863f" +dependencies = [ + "digest 0.11.3", +] + [[package]] name = "home" version = "0.5.12" @@ -7709,6 +7725,16 @@ dependencies = [ "digest 0.10.7", ] +[[package]] +name = "md-5" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69b6441f590336821bb897fb28fc622898ccceb1d6cea3fde5ea86b090c4de98" +dependencies = [ + "cfg-if", + "digest 0.11.3", +] + [[package]] name = "md5" version = "0.6.1" @@ -8559,7 +8585,7 @@ dependencies = [ "chrono", "dyn-clone", "ed25519-dalek", - "hmac", + "hmac 0.12.1", "http 1.5.0", "itertools 0.10.5", "log", @@ -8916,7 +8942,7 @@ dependencies = [ "der", "des 0.8.1", "hex", - "hmac", + "hmac 0.12.1", "pkcs12", "pkcs5", "rand 0.9.5", @@ -9080,7 +9106,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8ed6a7761f76e3b9f92dfb0a60a6a6477c61024b775147ff0973a02653abaf2" dependencies = [ "digest 0.10.7", - "hmac", + "hmac 0.12.1", ] [[package]] @@ -9366,7 +9392,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "695b3df3d3cc1015f12d70235e35b6b79befc5fa7a9b95b951eab1dd07c9efc2" dependencies = [ "cms", - "const-oid", + "const-oid 0.9.6", "der", "digest 0.10.7", "spki", @@ -9471,7 +9497,7 @@ dependencies = [ [[package]] name = "postgres-native-tls" version = "0.5.2" -source = "git+https://github.com/MaterializeInc/rust-postgres?rev=78c1222577bb091d69bc22b1bc7ad01c14675abe#78c1222577bb091d69bc22b1bc7ad01c14675abe" +source = "git+https://github.com/windmill-labs/rust-postgres?rev=7f4374f9237237759305942727e136742b1e7f8b#7f4374f9237237759305942727e136742b1e7f8b" dependencies = [ "native-tls", "tokio", @@ -9493,25 +9519,25 @@ dependencies = [ [[package]] name = "postgres-protocol" -version = "0.6.9" -source = "git+https://github.com/MaterializeInc/rust-postgres?rev=78c1222577bb091d69bc22b1bc7ad01c14675abe#78c1222577bb091d69bc22b1bc7ad01c14675abe" +version = "0.6.12" +source = "git+https://github.com/windmill-labs/rust-postgres?rev=7f4374f9237237759305942727e136742b1e7f8b#7f4374f9237237759305942727e136742b1e7f8b" dependencies = [ "base64 0.22.1", "byteorder", "bytes", "fallible-iterator", - "hmac", - "md-5 0.10.6", + "hmac 0.13.0", + "md-5 0.11.0", "memchr", - "rand 0.9.5", - "sha2 0.10.9", + "rand 0.10.3", + "sha2 0.11.0", "stringprep", ] [[package]] name = "postgres-types" version = "0.2.11" -source = "git+https://github.com/MaterializeInc/rust-postgres?rev=78c1222577bb091d69bc22b1bc7ad01c14675abe#78c1222577bb091d69bc22b1bc7ad01c14675abe" +source = "git+https://github.com/windmill-labs/rust-postgres?rev=7f4374f9237237759305942727e136742b1e7f8b#7f4374f9237237759305942727e136742b1e7f8b" dependencies = [ "array-init", "bit-vec 0.6.3", @@ -10471,7 +10497,7 @@ version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8dd2a808d456c4a54e300a23e9f5a67e122c3024119acbfd73e3bf664491cb2" dependencies = [ - "hmac", + "hmac 0.12.1", "subtle", ] @@ -10634,7 +10660,7 @@ version = "0.9.10" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b8573f03f5883dcaebdfcf4725caa1ecb9c15b2ef50c43a07b816e06799bb12d" dependencies = [ - "const-oid", + "const-oid 0.9.6", "digest 0.10.7", "num-bigint-dig", "num-integer", @@ -11658,6 +11684,17 @@ dependencies = [ "digest 0.10.7", ] +[[package]] +name = "sha2" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "446ba717509524cb3f22f17ecc096f10f4822d76ab5c0b9822c5f9c284e825f4" +dependencies = [ + "cfg-if", + "cpufeatures 0.3.1", + "digest 0.11.3", +] + [[package]] name = "sharded-slab" version = "0.1.7" @@ -12077,7 +12114,7 @@ dependencies = [ "generic-array", "hex", "hkdf", - "hmac", + "hmac 0.12.1", "itoa", "log", "md-5 0.10.6", @@ -12117,7 +12154,7 @@ dependencies = [ "futures-util", "hex", "hkdf", - "hmac", + "hmac 0.12.1", "home", "itoa", "log", @@ -13380,7 +13417,7 @@ dependencies = [ [[package]] name = "tokio-postgres" version = "0.7.15" -source = "git+https://github.com/MaterializeInc/rust-postgres?rev=78c1222577bb091d69bc22b1bc7ad01c14675abe#78c1222577bb091d69bc22b1bc7ad01c14675abe" +source = "git+https://github.com/windmill-labs/rust-postgres?rev=7f4374f9237237759305942727e136742b1e7f8b#7f4374f9237237759305942727e136742b1e7f8b" dependencies = [ "async-trait", "byteorder", @@ -14808,7 +14845,7 @@ dependencies = [ "gethostname", "git-version", "hex", - "hmac", + "hmac 0.12.1", "jsonwebtoken 10.4.0", "lazy_static", "once_cell", @@ -14956,7 +14993,7 @@ dependencies = [ "futures", "git-version", "hex", - "hmac", + "hmac 0.12.1", "http 1.5.0", "hyper 1.11.1", "indexmap 2.14.2", @@ -15291,7 +15328,7 @@ dependencies = [ "base64 0.22.1", "futures", "hex", - "hmac", + "hmac 0.12.1", "rand 0.9.5", "rdkafka", "reqwest 0.13.5", @@ -15640,7 +15677,7 @@ dependencies = [ "git-version", "globset", "hex", - "hmac", + "hmac 0.12.1", "hyper 1.11.1", "indexmap 2.14.2", "itertools 0.14.0", @@ -15843,7 +15880,7 @@ dependencies = [ "backon", "base64 0.22.1", "chrono", - "hmac", + "hmac 0.12.1", "http 1.5.0", "itertools 0.14.0", "lazy_static", @@ -15875,7 +15912,7 @@ dependencies = [ "base64 0.22.1", "chrono", "hex", - "hmac", + "hmac 0.12.1", "itertools 0.14.0", "lazy_static", "reqwest 0.12.28", @@ -16226,7 +16263,7 @@ dependencies = [ "futures", "futures-core", "hex", - "hmac", + "hmac 0.12.1", "itertools 0.14.0", "lazy_static", "once_cell", @@ -16519,7 +16556,7 @@ dependencies = [ "constant_time_eq 0.3.1", "futures", "hex", - "hmac", + "hmac 0.12.1", "http 1.5.0", "hyper 1.11.1", "itertools 0.14.0", @@ -16750,7 +16787,7 @@ dependencies = [ "gcp_auth", "git-version", "hex", - "hmac", + "hmac 0.12.1", "hudsucker", "hyper-http-proxy", "hyper-rustls 0.27.10", @@ -17438,7 +17475,7 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1301e935010a701ae5f8655edc0ad17c44bad3ac5ce8c39185f75453b720ae94" dependencies = [ - "const-oid", + "const-oid 0.9.6", "der", "spki", ] diff --git a/backend/Cargo.toml b/backend/Cargo.toml index d2bec22817..a17197ee9f 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -249,36 +249,21 @@ sqlx-sqlite = { git = "https://github.com/windmill-labs/sqlx", rev = "6bdaee94fa object_store = { git = "https://github.com/apache/arrow-rs-object-store", rev = "36752c975d4f29e20b57c91f81a10872dcd48ae7" } # Use tiberius main branch for libgssapi 0.8.1 fix (https://github.com/prisma/tiberius/issues/343) tiberius = { git = "https://github.com/prisma/tiberius", rev = "59db57960a14b422fb3a1309aa4aa47880896ff8" } -# Pin tokio-postgres / postgres-types / postgres-protocol to the -# MaterializeInc fork. windmill-trigger-postgres already pulled this -# fork in transitively for the postgres-replication crate -# (CopyBothDuplex, LogicalReplicationStream, TupleData with binary -# tuple support) which upstream rust-postgres has declined to merge -# since 2021 (PR #752 → #778, both still unmerged). +# Pin tokio-postgres / postgres-types / postgres-protocol to +# windmill-labs/rust-postgres, branch `windmill-describe-first`: the +# MaterializeInc fork (postgres-replication, which upstream has not merged) +# plus `Client::describe_typed`, which pg_executor needs before +# query_typed_raw so result columns of types the client has not seen yet +# (custom enums / domains, citext, postgis) don't deadlock a large result. +# Each pinned commit carries a `windmill-*` tag so it stays reachable when the +# branch is rebased onto MaterializeInc master. # -# MI also carries a mitigation for the -# Client::query_typed_raw / Client::prepare deadlock on result columns -# whose Oid the client doesn't know about yet (citext, custom enums / -# domains, postgis): MI's 2025-12-11 PR #33 resized the per-request -# response channel from mpsc::channel(1) → mpsc::channel(1024). -# bounded(1024) is sufficient for any realistic typeinfo deferral -# (need ~2-3 batches) but leaves a theoretical failure mode at -# >~64 MB results with a custom-Oid column. The strict-correct fix is -# mpsc::unbounded(); a follow-up PR to MI is open proposing that. -# -# The [patch.crates-io] entries below force windmill-worker's -# pg_executor (which imports `tokio_postgres::` directly from -# crates.io) onto the same fork as windmill-trigger-postgres, so the -# deadlock mitigation reaches both consumers. -# -# Upstream deadlock PRs (open, not on the critical path now that MI -# is mitigated): -# https://github.com/rust-postgres/rust-postgres/pull/1348 -# https://github.com/rust-postgres/rust-postgres/pull/1349 -# Reproducer: https://github.com/rubenfiszel/tokio-postgres-deadlock-repro -tokio-postgres = { git = "https://github.com/MaterializeInc/rust-postgres", rev = "78c1222577bb091d69bc22b1bc7ad01c14675abe" } -postgres-types = { git = "https://github.com/MaterializeInc/rust-postgres", rev = "78c1222577bb091d69bc22b1bc7ad01c14675abe" } -postgres-protocol = { git = "https://github.com/MaterializeInc/rust-postgres", rev = "78c1222577bb091d69bc22b1bc7ad01c14675abe" } +# The [patch.crates-io] entries force windmill-worker's pg_executor (which +# imports `tokio_postgres::` from crates.io) onto the same fork as +# windmill-trigger-postgres. +tokio-postgres = { git = "https://github.com/windmill-labs/rust-postgres", rev = "7f4374f9237237759305942727e136742b1e7f8b" } +postgres-types = { git = "https://github.com/windmill-labs/rust-postgres", rev = "7f4374f9237237759305942727e136742b1e7f8b" } +postgres-protocol = { git = "https://github.com/windmill-labs/rust-postgres", rev = "7f4374f9237237759305942727e136742b1e7f8b" } [dependencies] anyhow.workspace = true @@ -630,8 +615,8 @@ convert_case = "0.6.0" getrandom = "0.2" tokio-postgres = {version = "^0.7", features = ["array-impls", "with-serde_json-1", "with-chrono-0_4", "with-uuid-1", "with-bit-vec-0_6"]} postgres-protocol = "0.6" -rust-postgres = { package = "tokio-postgres", git = "https://github.com/MaterializeInc/rust-postgres", rev = "78c1222577bb091d69bc22b1bc7ad01c14675abe"} -rust-postgres-native-tls = { package = "postgres-native-tls", git = "https://github.com/MaterializeInc/rust-postgres", features = ["runtime"], rev = "78c1222577bb091d69bc22b1bc7ad01c14675abe" } +rust-postgres = { package = "tokio-postgres", git = "https://github.com/windmill-labs/rust-postgres", rev = "7f4374f9237237759305942727e136742b1e7f8b"} +rust-postgres-native-tls = { package = "postgres-native-tls", git = "https://github.com/windmill-labs/rust-postgres", features = ["runtime"], rev = "7f4374f9237237759305942727e136742b1e7f8b" } bit-vec = "=0.6.3" mappable-rc = "^0" mysql_async = { version = "*", default-features = false, features = ["minimal", "default", "native-tls-tls", "rust_decimal"]} diff --git a/backend/windmill-worker/src/pg_executor.rs b/backend/windmill-worker/src/pg_executor.rs index 22a59f82e6..f2d48446bc 100644 --- a/backend/windmill-worker/src/pg_executor.rs +++ b/backend/windmill-worker/src/pg_executor.rs @@ -361,6 +361,79 @@ fn wrap_param_encoding_error( to_anyhow(err).into() } +/// Whether a statement may return a row stream, and so needs its column types described +/// before it runs. Anything not recognised as returning at most a row counts as a stream: +/// a wrong guess only costs a round trip, while a missed stream can hang the job. +fn can_stream_rows(query: &str) -> bool { + let Some(keyword) = leading_keyword(query) else { + return true; + }; + let rowless = matches!( + keyword.as_str(), + "insert" + | "update" + | "delete" + | "merge" + | "create" + | "alter" + | "drop" + | "truncate" + | "grant" + | "revoke" + | "comment" + | "set" + | "reset" + | "call" + | "do" + | "begin" + | "commit" + | "rollback" + | "lock" + | "vacuum" + | "analyze" + | "refresh" + ); + !rowless || query.to_ascii_lowercase().contains("returning") +} + +/// The first keyword of a statement, past whitespace and comments, following Postgres's +/// scanner: a `--` comment ends at CR or LF, and block comments nest, so +/// `/* /* a */ INSERT */ SELECT` starts with SELECT. `None` when the comments never end. +fn leading_keyword(query: &str) -> Option { + let mut rest = query; + loop { + rest = rest.trim_start(); + if let Some(after) = rest.strip_prefix("--") { + rest = &after[after.find(['\n', '\r'])?..]; + } else if rest.starts_with("/*") { + let mut depth = 0usize; + let mut i = 0; + let bytes = rest.as_bytes(); + loop { + match bytes.get(i..i + 2)? { + b"/*" => depth += 1, + b"*/" => depth -= 1, + _ => { + i += 1; + continue; + } + } + i += 2; + if depth == 0 { + break; + } + } + rest = &rest[i..]; + } else { + break; + } + } + let end = rest + .find(|c: char| !(c.is_ascii_alphanumeric() || c == '_' || c == '$')) + .unwrap_or(rest.len()); + Some(rest[..end].to_ascii_lowercase()) +} + fn otyp_to_pg_type(otyp: &str) -> error::Result { let base = otyp.trim_end_matches("[]"); let is_array = otyp.ends_with("[]"); @@ -509,15 +582,25 @@ fn do_postgresql_inner<'a>( let result_f = async move { let mut res: Vec> = vec![]; - // Always prefer query_typed_raw (unnamed prepared statement). It is sent as - // a single Parse+Bind+Execute+Sync round-trip, so it survives transaction-mode - // connection poolers (PgBouncer/Supabase pooler/RDS Proxy) where named + // Always prefer query_typed_raw (unnamed prepared statement). It and its + // describe are each a self-contained round-trip ending in Sync, so they + // survive transaction-mode connection poolers (PgBouncer/Supabase + // pooler/RDS Proxy) where named // statements ("s0", "s1", ...) can be reported missing because the prepare // and the execute land on different backend connections. Fall back to // prepare + query_raw only when an arg has a type unsupported by // otyp_to_pg_type (e.g. custom enum, geometry, …) — in that case we lose // pooler safety, but the query at least runs against a direct connection. let rows = if all_types_resolved { + // query_typed_raw looks up result column types this connection has not seen + // (enums, domains, extension types) while the rows already stream, and on a + // large result that lookup waits behind them forever. Describing first resolves + // them up front, still without a named statement. + if can_stream_rows(&query) { + if let Err(e) = client.describe_typed(&query, ¶m_types).await { + return Err(wrap_param_encoding_error(e, ¶m_meta, ¶m_types)); + } + } let typed_params = query_params .iter() .zip(param_types.iter()) @@ -2510,6 +2593,39 @@ mod tests { windmill_parser_sql::parse_pg_typ(arg_t) } + #[test] + fn only_rowless_statements_skip_the_describe() { + for q in [ + "SELECT * FROM t", + "-- $1 n (int)\nselect $1", + "WITH x AS (DELETE FROM t RETURNING *) SELECT * FROM x", + "INSERT INTO t VALUES (1) RETURNING id", + "/* c */ UPDATE t SET a = 1\nreturning *", + "UPDATE t SET note = $$a;b$$ RETURNING *", + "/* /* inner */ INSERT */ SELECT * FROM t", + "/* never closed INSERT", + "-- INSERT", + "-- header\rSELECT m, pad AS\nupdate FROM t", + "insert_rows()", + "TABLE t", + "VALUES (1)", + "EXPLAIN SELECT 1", + ] { + assert!(can_stream_rows(q), "{q}"); + } + for q in [ + "INSERT INTO t VALUES (1)", + "-- $1 n (int)\nUPDATE t SET a = $1", + "delete from t", + "CREATE TABLE t (a int)", + "SET search_path TO x", + "CALL p()", + "/* a /* b */ c */\n-- d\nINSERT INTO t VALUES (1)", + ] { + assert!(!can_stream_rows(q), "{q}"); + } + } + #[test] fn convert_val_null_for_every_known_arg_t() { // `Value::Null` for every type the parser may resolve, plus an unknown