diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 10f5a71687..dcbaaede68 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -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", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 263cb62907..e32920b180 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -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"]} diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 4ec106ecb2..dec5abfb23 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -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) -> 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"))] diff --git a/backend/windmill-worker/src/pg_executor.rs b/backend/windmill-worker/src/pg_executor.rs index 2de56ddf08..3222d7ca3e 100644 --- a/backend/windmill-worker/src/pg_executor.rs +++ b/backend/windmill-worker/src/pg_executor.rs @@ -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);