fix: never reuse a pg connection a script left inside a transaction (#11497)

* fix: never reuse a pg connection a script left inside a transaction

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KHjj2iG8kKPCaFuK6bqXvi

* fix: close a pg connection whose job failed instead of caching it

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KHjj2iG8kKPCaFuK6bqXvi

* docs: describe every caller of PgConnectionLease::discard

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KHjj2iG8kKPCaFuK6bqXvi

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-10-03 09:03:44 +02:00
committed by GitHub
co-authored by Claude Opus 5.5
parent f2393e5b24
commit b7f74cd562
4 changed files with 109 additions and 15 deletions
+4 -4
View File
@@ -9497,7 +9497,7 @@ dependencies = [
[[package]]
name = "postgres-native-tls"
version = "0.5.2"
source = "git+https://github.com/windmill-labs/rust-postgres?rev=7f4374f9237237759305942727e136742b1e7f8b#7f4374f9237237759305942727e136742b1e7f8b"
source = "git+https://github.com/windmill-labs/rust-postgres?rev=b3a6443aa1ae90a4f6c4b742103ba80584598f2e#b3a6443aa1ae90a4f6c4b742103ba80584598f2e"
dependencies = [
"native-tls",
"tokio",
@@ -9520,7 +9520,7 @@ dependencies = [
[[package]]
name = "postgres-protocol"
version = "0.6.12"
source = "git+https://github.com/windmill-labs/rust-postgres?rev=7f4374f9237237759305942727e136742b1e7f8b#7f4374f9237237759305942727e136742b1e7f8b"
source = "git+https://github.com/windmill-labs/rust-postgres?rev=b3a6443aa1ae90a4f6c4b742103ba80584598f2e#b3a6443aa1ae90a4f6c4b742103ba80584598f2e"
dependencies = [
"base64 0.22.1",
"byteorder",
@@ -9537,7 +9537,7 @@ dependencies = [
[[package]]
name = "postgres-types"
version = "0.2.11"
source = "git+https://github.com/windmill-labs/rust-postgres?rev=7f4374f9237237759305942727e136742b1e7f8b#7f4374f9237237759305942727e136742b1e7f8b"
source = "git+https://github.com/windmill-labs/rust-postgres?rev=b3a6443aa1ae90a4f6c4b742103ba80584598f2e#b3a6443aa1ae90a4f6c4b742103ba80584598f2e"
dependencies = [
"array-init",
"bit-vec 0.6.3",
@@ -13417,7 +13417,7 @@ dependencies = [
[[package]]
name = "tokio-postgres"
version = "0.7.15"
source = "git+https://github.com/windmill-labs/rust-postgres?rev=7f4374f9237237759305942727e136742b1e7f8b#7f4374f9237237759305942727e136742b1e7f8b"
source = "git+https://github.com/windmill-labs/rust-postgres?rev=b3a6443aa1ae90a4f6c4b742103ba80584598f2e#b3a6443aa1ae90a4f6c4b742103ba80584598f2e"
dependencies = [
"async-trait",
"byteorder",
+8 -7
View File
@@ -254,16 +254,17 @@ tiberius = { git = "https://github.com/prisma/tiberius", rev = "59db57960a14b422
# 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
# (custom enums / domains, citext, postgis) don't deadlock a large result,
# and `Client::transaction_status`, which keeps a connection left inside a
# transaction out of the executor's cache. Each pinned commit carries a `windmill-*` tag so it stays reachable when the
# branch is rebased onto MaterializeInc master.
#
# 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" }
tokio-postgres = { git = "https://github.com/windmill-labs/rust-postgres", rev = "b3a6443aa1ae90a4f6c4b742103ba80584598f2e" }
postgres-types = { git = "https://github.com/windmill-labs/rust-postgres", rev = "b3a6443aa1ae90a4f6c4b742103ba80584598f2e" }
postgres-protocol = { git = "https://github.com/windmill-labs/rust-postgres", rev = "b3a6443aa1ae90a4f6c4b742103ba80584598f2e" }
[dependencies]
anyhow.workspace = true
@@ -616,8 +617,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/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" }
rust-postgres = { package = "tokio-postgres", git = "https://github.com/windmill-labs/rust-postgres", rev = "b3a6443aa1ae90a4f6c4b742103ba80584598f2e"}
rust-postgres-native-tls = { package = "postgres-native-tls", git = "https://github.com/windmill-labs/rust-postgres", features = ["runtime"], rev = "b3a6443aa1ae90a4f6c4b742103ba80584598f2e" }
bit-vec = "=0.6.3"
mappable-rc = "^0"
mysql_async = { version = "*", default-features = false, features = ["minimal", "default", "native-tls-tls", "rust_decimal"]}
+66
View File
@@ -1795,6 +1795,72 @@ async fn test_postgresql_cached_connection_released_for_other_key(
Ok(())
}
/// A transaction a script leaves open must not carry over into the next job
/// that would reuse the connection.
#[sqlx::test(fixtures("base"))]
#[serial(pg_cache)]
async fn test_postgresql_open_transaction_not_reused(db: Pool<Postgres>) -> anyhow::Result<()> {
use windmill_worker::pg_executor::clear_pg_cache;
initialize_tracing().await;
clear_pg_cache().await;
let server = ApiServer::start(db.clone()).await?;
let port = server.addr.port();
let dbname: String = sqlx::query_scalar("SELECT current_database()")
.fetch_one(&db)
.await?;
let run = |content: &str| {
RunJob::from(JobPayload::Code(RawCode {
hash: None,
content: content.to_string(),
path: None,
lock: None,
language: ScriptLang::Postgresql,
cache_ttl: None,
cache_ignore_s3_path: None,
dedicated_worker: None,
concurrency_settings: windmill_common::runnable_settings::ConcurrencySettings::default(
)
.into(),
debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(),
modules: None,
tag: None,
}))
.arg(
"database",
json!({"host": "localhost", "port": 5432, "dbname": dbname, "user": "postgres", "password": "changeme"}),
)
.run_until_complete(&db, false, port)
};
run("SELECT 1 as n;").await.json_result().unwrap();
let opened = run("BEGIN; SELECT 1 as n;").await.json_result().unwrap();
assert_eq!(opened, json!([{"n": 1}]));
run("SELECT 2 as n;").await.json_result().unwrap();
// A connection kept in the cache with the transaction open would sit idle
// in it, holding its locks, and run the next job inside it.
let mut idle_in_transaction = -1;
for _ in 0..20 {
idle_in_transaction = sqlx::query_scalar::<_, i64>(
"SELECT count(*) FROM pg_stat_activity
WHERE datname = current_database() AND state LIKE 'idle in transaction%'",
)
.fetch_one(&db)
.await?;
if idle_in_transaction == 0 {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
assert_eq!(idle_in_transaction, 0);
clear_pg_cache().await;
Ok(())
}
/// Runs multiple PG jobs through a SINGLE worker (like production) to verify
/// that SET ROLE / search_path changes do not leak across jobs.
#[sqlx::test(fixtures("base"))]
+31 -4
View File
@@ -235,8 +235,9 @@ impl PgConnectionLease {
Self { slot, conn, cache_on_release: false }
}
/// Gives up a cached connection that failed its probe; the job then
/// connects on its own.
/// Closes the lease's connection so it is never cached again: a cached one
/// that failed its probe (the job then connects on its own), or one a job
/// failed on or left inside a transaction.
fn discard(&mut self) {
let conn = self.conn.take();
if self.slot == PgLeaseSlot::CheckedOut {
@@ -1380,11 +1381,37 @@ pub async fn do_postgresql(
)
.await
}
.map_err(|e| map_s3object_jsonb_overflow(e, had_s3object_input))?;
.map_err(|e| map_s3object_jsonb_overflow(e, had_s3object_input));
// A failed job can leave the connection inside a transaction, or with its
// query still running after a timeout or cancel, and the status byte can
// lag behind an error. Never reuse it.
let result = match result {
Ok(result) => result,
Err(e) => {
lease.discard();
return Err(e);
}
};
*mem_peak = size.load(Ordering::Relaxed) as i32;
if !*CLOUD_HOSTED && lease.slot == PgLeaseSlot::Uncached {
// A transaction the script left open would otherwise carry over into the
// next job that reuses the connection. Closing the connection rolls it
// back, as it always has without the cache. The status comes with the last
// reply, so a script that ends cleanly pays nothing for this check.
if lease.client().transaction_status() != tokio_postgres::TransactionStatus::Idle {
lease.discard();
if !run_inline {
windmill_queue::append_logs(
&job.id,
&job.workspace_id,
"The script ended inside an open transaction, which was rolled back. \
End it with COMMIT to keep its changes.\n",
conn,
)
.await;
}
} else if !*CLOUD_HOSTED && lease.slot == PgLeaseSlot::Uncached {
lease.cache_on_release = is_most_used_conn(&database_string).await;
}
drop(lease);