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) <noreply@anthropic.com>

* fix: carry the parked-rows delivery fix in the postgres fork

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: describe postgres queries before streaming instead of parking rows

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* perf: skip the postgres describe for statements that return no rows

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: look for returning in the whole postgres statement

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: read the leading postgres keyword past nested block comments

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: end leading postgres line comments at CR too

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-10-02 17:21:20 +02:00
committed by GitHub
co-authored by Claude Opus 5.5
parent ae093b7666
commit 14801fdd68
3 changed files with 206 additions and 68 deletions
+71 -34
View File
@@ -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",
]
+16 -31
View File
@@ -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"]}
+119 -3
View File
@@ -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<String> {
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<Type> {
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<Box<serde_json::value::RawValue>> = 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, &param_types).await {
return Err(wrap_param_encoding_error(e, &param_meta, &param_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