feat: log pg job progress, warn on stalls, detect dead db connections (#11491)

* feat: log pg job progress, warn on stalls, detect dead db connections

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

* fix: log only slow pg steps to keep job log writes off fast statements

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

* fix: say when a slow connect included a failed cached-connection reset

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-02 22:26:37 +02:00
committed by GitHub
co-authored by Claude Opus 5.5
parent eed7e7d9e6
commit d7227726be
2 changed files with 174 additions and 6 deletions
+22 -5
View File
@@ -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<tokio_postgres::Config, error::Error> {
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),
+152 -1
View File
@@ -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<BoxFuture<'a, error::Result<Vec<Box<RawValue>>>>> {
let mut query_params = vec![];
let mut param_types: Vec<Type> = 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<Vec<String>> = 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(),
&param_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