diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index ff5fb76ee0..1013baa2e0 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -1068,6 +1068,18 @@ impl TokioPgConnection { } } +/// Without these, a server that vanishes without closing the socket (a failover, +/// a dropped route) leaves a query waiting on a read for the OS default of two +/// hours. The server's kernel answers the probes, so a slow query is unaffected. +fn set_pg_keepalive(config: &mut tokio_postgres::Config) { + config + .keepalives(true) + .keepalives_idle(std::time::Duration::from_secs(60)) + .keepalives_interval(std::time::Duration::from_secs(10)) + .keepalives_retries(6) + .tcp_user_timeout(std::time::Duration::from_secs(120)); +} + impl PgDatabase { /// The role the connection logs in as, whichever way it authenticates. pub fn login_name(&self) -> &str { @@ -1109,6 +1121,12 @@ impl PgDatabase { self.options.as_deref().filter(|o| !o.is_empty()) } + fn uri_config(&self) -> Result { + let mut config: tokio_postgres::Config = self.to_uri().parse().map_err(to_anyhow)?; + set_pg_keepalive(&mut config); + Ok(config) + } + pub async fn connect( &self, main_db: Option<&DB>, @@ -1255,10 +1273,8 @@ impl PgDatabase { let (client, connection) = tokio::time::timeout( std::time::Duration::from_secs(20), - tokio_postgres::connect( - &self.to_uri(), - MakeTlsConnector::new(connector.build().map_err(to_anyhow)?), - ), + self.uri_config()? + .connect(MakeTlsConnector::new(connector.build().map_err(to_anyhow)?)), ) .await .map_err(to_anyhow)? @@ -1269,7 +1285,7 @@ impl PgDatabase { tracing::info!("Creating new connection"); let (client, connection) = tokio::time::timeout( std::time::Duration::from_secs(20), - tokio_postgres::connect(&self.to_uri(), NoTls), + self.uri_config()?.connect(NoTls), ) .await .map_err(to_anyhow)? @@ -1386,6 +1402,7 @@ impl PgDatabase { if let Some(options) = self.non_empty_options() { config.options(options); } + set_pg_keepalive(&mut config); let (client, connection) = tokio::time::timeout( std::time::Duration::from_secs(20), diff --git a/backend/windmill-worker/src/pg_executor.rs b/backend/windmill-worker/src/pg_executor.rs index 668e44aad2..2de56ddf08 100644 --- a/backend/windmill-worker/src/pg_executor.rs +++ b/backend/windmill-worker/src/pg_executor.rs @@ -59,6 +59,96 @@ lazy_static! { pub static ref CACHE_HITS: AtomicU64 = AtomicU64::new(0); } +/// How far a job's statements have got, read by its stall warning. +struct PgProgress { + start: std::time::Instant, + statement: AtomicUsize, + rows: AtomicU64, + last_progress_ms: AtomicU64, +} + +impl PgProgress { + fn new() -> Self { + Self { + start: std::time::Instant::now(), + statement: AtomicUsize::new(0), + rows: AtomicU64::new(0), + last_progress_ms: AtomicU64::new(0), + } + } + + fn elapsed_ms(&self) -> u64 { + self.start.elapsed().as_millis() as u64 + } + + fn start_statement(&self, index: usize) { + self.statement.store(index, Ordering::Relaxed); + self.last_progress_ms + .store(self.elapsed_ms(), Ordering::Relaxed); + } + + fn row(&self) { + self.rows.fetch_add(1, Ordering::Relaxed); + self.last_progress_ms + .store(self.elapsed_ms(), Ordering::Relaxed); + } + + fn rows(&self) -> u64 { + self.rows.load(Ordering::Relaxed) + } +} + +const PG_STALL_WARNING_AFTER: Duration = Duration::from_secs(5 * 60); +/// Steps faster than this are not logged: a log line is a write to the main DB, +/// and a script split into many statements would pay one per statement. +const PG_SLOW_STEP: Duration = Duration::from_secs(1); + +fn connection_kind(fresh_connection: bool) -> &'static str { + if fresh_connection { + "new connection" + } else { + "worker's cached connection" + } +} + +/// Warns once per stall in the job log; never completes. A statement still +/// computing its first row also trips it, which is why it only warns: the job +/// timeout stays the one thing that stops a job. +async fn warn_on_stalled_statement( + progress: &PgProgress, + statement_count: usize, + fresh_connection: bool, + job_id: Uuid, + workspace_id: &str, + conn: &Connection, +) -> std::convert::Infallible { + let mut warned_for = None; + loop { + tokio::time::sleep(Duration::from_secs(30)).await; + let last = progress.last_progress_ms.load(Ordering::Relaxed); + let stalled = Duration::from_millis(progress.elapsed_ms().saturating_sub(last)); + if stalled < PG_STALL_WARNING_AFTER || warned_for == Some(last) { + continue; + } + warned_for = Some(last); + windmill_queue::append_logs( + &job_id, + workspace_id, + format!( + "No new row for {} min on statement {}/{statement_count} ({} rows so far, {}). \ + A query still running on the database is normal; if the database shows no \ + active query for this connection, the job is stuck in the worker.\n", + stalled.as_secs() / 60, + progress.statement.load(Ordering::Relaxed) + 1, + progress.rows(), + connection_kind(fresh_connection), + ), + conn, + ) + .await; + } +} + /// One reusable connection per worker process, checked out by a job for the /// duration of its query and checked back in afterwards. /// @@ -641,6 +731,7 @@ fn do_postgresql_inner<'a>( workspace_id: &'a str, log_conn: &'a Connection, raw_output: bool, + progress: &'a PgProgress, ) -> error::Result>>>> { let mut query_params = vec![]; let mut param_types: Vec = vec![]; @@ -791,10 +882,13 @@ fn do_postgresql_inner<'a>( if skip_collect { futures::pin_mut!(rows); - while rows.try_next().await.map_err(to_anyhow)?.is_some() {} + while rows.try_next().await.map_err(to_anyhow)?.is_some() { + progress.row(); + } } else if let Some(ref s3) = s3 { let format_state_ref = &format_state; let rows_stream = rows.map_err(to_anyhow).map(move |row_result| { + progress.row(); row_result.and_then(|row| { postgres_row_to_json_value_with_state(row, format_state_ref).map_err(to_anyhow) }) @@ -832,6 +926,7 @@ fn do_postgresql_inner<'a>( let mut column_names: Option> = None; while let Some(row) = rows.try_next().await.map_err(to_anyhow)? { + progress.row(); if column_names.is_none() { column_names = Some( row.columns() @@ -1027,6 +1122,8 @@ pub async fn do_postgresql( auth_mode.cache_key_segment() ); + let connect_started = std::time::Instant::now(); + let mut cached_connection_failed_reset = false; let mut lease = if *CLOUD_HOSTED { PgConnectionLease::uncached() } else { @@ -1091,6 +1188,7 @@ pub async fn do_postgresql( CACHE_HITS.fetch_add(1, std::sync::atomic::Ordering::Relaxed); } else { tracing::info!("Cached connection is stale, creating new one"); + cached_connection_failed_reset = true; lease.discard(); } } @@ -1105,6 +1203,27 @@ pub async fn do_postgresql( }); } + let fresh_connection = lease.slot == PgLeaseSlot::Uncached; + let log_progress = !run_inline && !annotations.prepare; + let connect_time = connect_started.elapsed(); + if log_progress && connect_time >= PG_SLOW_STEP { + windmill_queue::append_logs( + &job.id, + &job.workspace_id, + format!( + "Getting a database connection took {} ms ({})\n", + connect_time.as_millis(), + if cached_connection_failed_reset { + "the cached connection failed its reset, then a new connection" + } else { + connection_kind(fresh_connection) + } + ), + conn, + ) + .await; + } + let (mut sig, _) = parse_pgsql_sig_with_typed_schema(&query) .map_err(|x| Error::ExecutionErr(x.to_string()))?; @@ -1139,6 +1258,9 @@ pub async fn do_postgresql( let size = AtomicUsize::new(0); let size_ref = &size; + let progress = PgProgress::new(); + let progress_ref = &progress; + let statement_count = queries.len(); let result_f = async move { let mut results = vec![]; // Session reset (DISCARD ALL) is now handled eagerly when validating @@ -1176,6 +1298,9 @@ pub async fn do_postgresql( let skip_collect = collection_strategy.collect_last_statement_only(queries.len()) && i < queries.len() - 1; let is_last = i == queries.len() - 1; + progress_ref.start_statement(i); + let rows_before = progress_ref.rows(); + let statement_started = std::time::Instant::now(); let result = do_postgresql_inner( query.to_string(), ¶m_idx_to_arg_and_value, @@ -1197,8 +1322,24 @@ pub async fn do_postgresql( &job.workspace_id, conn, annotations.raw_output && is_last && !skip_collect, + progress_ref, )? .await?; + let statement_time = statement_started.elapsed(); + if log_progress && statement_time >= PG_SLOW_STEP { + windmill_queue::append_logs( + &job.id, + &job.workspace_id, + format!( + "Statement {}/{statement_count}: {} rows in {} ms\n", + i + 1, + progress_ref.rows() - rows_before, + statement_time.as_millis() + ), + conn, + ) + .await; + } results.push(result); } @@ -1211,6 +1352,16 @@ pub async fn do_postgresql( collection_strategy.collect(results) } }; + let result_f = async { + if log_progress { + tokio::select! { + result = result_f => result, + never = warn_on_stalled_statement(progress_ref, statement_count, fresh_connection, job.id, &job.workspace_id, conn) => match never {}, + } + } else { + result_f.await + } + }; let result = if run_inline { result_f.await