diff --git a/backend/.sqlx/query-32de3e7e8d1378059766fcb22ddd7da34d98b6b62311fc20dd99722eeb47355a.json b/backend/.sqlx/query-32de3e7e8d1378059766fcb22ddd7da34d98b6b62311fc20dd99722eeb47355a.json new file mode 100644 index 0000000000..5843a897ac --- /dev/null +++ b/backend/.sqlx/query-32de3e7e8d1378059766fcb22ddd7da34d98b6b62311fc20dd99722eeb47355a.json @@ -0,0 +1,41 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT file_path, log_ts, ok_lines, err_lines\n FROM log_file\n WHERE hostname = $1\n AND log_ts BETWEEN $2::timestamp - interval '90 seconds' AND $2::timestamp + interval '30 seconds'\n ORDER BY log_ts DESC\n LIMIT 10", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "file_path", + "type_info": "Varchar" + }, + { + "ordinal": 1, + "name": "log_ts", + "type_info": "Timestamp" + }, + { + "ordinal": 2, + "name": "ok_lines", + "type_info": "Int8" + }, + { + "ordinal": 3, + "name": "err_lines", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Text", + "Timestamp" + ] + }, + "nullable": [ + false, + false, + true, + true + ] + }, + "hash": "32de3e7e8d1378059766fcb22ddd7da34d98b6b62311fc20dd99722eeb47355a" +} diff --git a/backend/.sqlx/query-4d80985dd0794a01a2af18ae7abf4a3ab8ba3d162ed5d04735caea7295da0b20.json b/backend/.sqlx/query-4d80985dd0794a01a2af18ae7abf4a3ab8ba3d162ed5d04735caea7295da0b20.json deleted file mode 100644 index 4847a9a2b1..0000000000 --- a/backend/.sqlx/query-4d80985dd0794a01a2af18ae7abf4a3ab8ba3d162ed5d04735caea7295da0b20.json +++ /dev/null @@ -1,58 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n SELECT\n j.id AS \"id!\", j.workspace_id AS \"workspace_id!\", j.parent_job, j.flow_step_id IS NOT NULL AS \"is_flow_step?\",\n COALESCE(s.flow_status, s.workflow_as_code_status) AS \"flow_status: Box\", r.ping AS last_ping, j.same_worker AS \"same_worker?\"\n FROM v2_job_queue q JOIN v2_job j USING (id) LEFT JOIN v2_job_runtime r USING (id) LEFT JOIN v2_job_status s USING (id)\n WHERE q.running = true AND q.suspend = 0 AND q.suspend_until IS null AND q.scheduled_for <= now()\n AND (j.kind = 'flow' OR j.kind = 'flowpreview' OR j.kind = 'flownode')\n AND r.ping IS NOT NULL AND r.ping < NOW() - ($1 || ' seconds')::interval\n AND q.canceled_by IS NULL\n\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id!", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "workspace_id!", - "type_info": "Varchar" - }, - { - "ordinal": 2, - "name": "parent_job", - "type_info": "Uuid" - }, - { - "ordinal": 3, - "name": "is_flow_step?", - "type_info": "Bool" - }, - { - "ordinal": 4, - "name": "flow_status: Box", - "type_info": "Jsonb" - }, - { - "ordinal": 5, - "name": "last_ping", - "type_info": "Timestamptz" - }, - { - "ordinal": 6, - "name": "same_worker?", - "type_info": "Bool" - } - ], - "parameters": { - "Left": [ - "Text" - ] - }, - "nullable": [ - false, - false, - true, - null, - null, - true, - false - ] - }, - "hash": "4d80985dd0794a01a2af18ae7abf4a3ab8ba3d162ed5d04735caea7295da0b20" -} diff --git a/backend/.sqlx/query-bd000b7ab6438826093ba556af819741b493a91babaeb103852d3912f520719f.json b/backend/.sqlx/query-bd000b7ab6438826093ba556af819741b493a91babaeb103852d3912f520719f.json new file mode 100644 index 0000000000..42b90ef229 --- /dev/null +++ b/backend/.sqlx/query-bd000b7ab6438826093ba556af819741b493a91babaeb103852d3912f520719f.json @@ -0,0 +1,112 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT\n j.id AS \"id!\", j.workspace_id AS \"workspace_id!\", j.parent_job, j.flow_step_id IS NOT NULL AS \"is_flow_step?\",\n COALESCE(s.flow_status, s.workflow_as_code_status) AS \"flow_status: Box\", r.ping AS last_ping, j.same_worker AS \"same_worker?\",\n q.worker AS \"worker?\",\n wp.ping_at AS \"worker_last_ping?\",\n wp.memory_usage AS \"worker_memory_usage?\",\n wp.wm_memory_usage AS \"worker_wm_memory_usage?\",\n wp.memory AS \"worker_memory_total?\",\n wp.worker_group AS \"worker_group?\",\n wp.wm_version AS \"worker_version?\",\n wp.current_job_id AS \"worker_current_job_id?\",\n wp.worker_instance AS \"worker_instance?\"\n FROM v2_job_queue q JOIN v2_job j USING (id) LEFT JOIN v2_job_runtime r USING (id) LEFT JOIN v2_job_status s USING (id)\n LEFT JOIN worker_ping wp ON wp.worker = q.worker\n WHERE q.running = true AND q.suspend = 0 AND q.suspend_until IS null AND q.scheduled_for <= now()\n AND (j.kind = 'flow' OR j.kind = 'flowpreview' OR j.kind = 'flownode')\n AND r.ping IS NOT NULL AND r.ping < NOW() - ($1 || ' seconds')::interval\n AND q.canceled_by IS NULL\n\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "workspace_id!", + "type_info": "Varchar" + }, + { + "ordinal": 2, + "name": "parent_job", + "type_info": "Uuid" + }, + { + "ordinal": 3, + "name": "is_flow_step?", + "type_info": "Bool" + }, + { + "ordinal": 4, + "name": "flow_status: Box", + "type_info": "Jsonb" + }, + { + "ordinal": 5, + "name": "last_ping", + "type_info": "Timestamptz" + }, + { + "ordinal": 6, + "name": "same_worker?", + "type_info": "Bool" + }, + { + "ordinal": 7, + "name": "worker?", + "type_info": "Varchar" + }, + { + "ordinal": 8, + "name": "worker_last_ping?", + "type_info": "Timestamptz" + }, + { + "ordinal": 9, + "name": "worker_memory_usage?", + "type_info": "Int8" + }, + { + "ordinal": 10, + "name": "worker_wm_memory_usage?", + "type_info": "Int8" + }, + { + "ordinal": 11, + "name": "worker_memory_total?", + "type_info": "Int8" + }, + { + "ordinal": 12, + "name": "worker_group?", + "type_info": "Varchar" + }, + { + "ordinal": 13, + "name": "worker_version?", + "type_info": "Varchar" + }, + { + "ordinal": 14, + "name": "worker_current_job_id?", + "type_info": "Uuid" + }, + { + "ordinal": 15, + "name": "worker_instance?", + "type_info": "Varchar" + } + ], + "parameters": { + "Left": [ + "Text" + ] + }, + "nullable": [ + false, + false, + true, + null, + null, + true, + false, + true, + false, + true, + true, + true, + false, + false, + true, + false + ] + }, + "hash": "bd000b7ab6438826093ba556af819741b493a91babaeb103852d3912f520719f" +} diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 5044d26a80..af2a66031c 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -3309,8 +3309,18 @@ async fn handle_zombie_flows(db: &DB) -> error::Result<()> { r#" SELECT j.id AS "id!", j.workspace_id AS "workspace_id!", j.parent_job, j.flow_step_id IS NOT NULL AS "is_flow_step?", - COALESCE(s.flow_status, s.workflow_as_code_status) AS "flow_status: Box", r.ping AS last_ping, j.same_worker AS "same_worker?" + COALESCE(s.flow_status, s.workflow_as_code_status) AS "flow_status: Box", r.ping AS last_ping, j.same_worker AS "same_worker?", + q.worker AS "worker?", + wp.ping_at AS "worker_last_ping?", + wp.memory_usage AS "worker_memory_usage?", + wp.wm_memory_usage AS "worker_wm_memory_usage?", + wp.memory AS "worker_memory_total?", + wp.worker_group AS "worker_group?", + wp.wm_version AS "worker_version?", + wp.current_job_id AS "worker_current_job_id?", + wp.worker_instance AS "worker_instance?" FROM v2_job_queue q JOIN v2_job j USING (id) LEFT JOIN v2_job_runtime r USING (id) LEFT JOIN v2_job_status s USING (id) + LEFT JOIN worker_ping wp ON wp.worker = q.worker WHERE q.running = true AND q.suspend = 0 AND q.suspend_until IS null AND q.scheduled_for <= now() AND (j.kind = 'flow' OR j.kind = 'flowpreview' OR j.kind = 'flownode') AND r.ping IS NOT NULL AND r.ping < NOW() - ($1 || ' seconds')::interval @@ -3377,8 +3387,129 @@ async fn handle_zombie_flows(db: &DB) -> error::Result<()> { let now = now_from_db(db).await?; let base_url = BASE_URL.read().await; let workspace_id = flow.workspace_id.clone(); + + let fmt_mb = |b: i64| format!("{:.1} MB", b as f64 / 1024.0 / 1024.0); + let worker_ping_stale = flow + .worker_last_ping + .map(|wp| (now - wp).num_seconds() > 60); + let worker_info = if let Some(worker_name) = flow.worker.as_deref() { + let mut s = format!("\nWorker handling the flow: {worker_name}"); + match (flow.worker_group.as_deref(), flow.worker_version.as_deref()) { + (Some(g), Some(v)) => s.push_str(&format!(" (group: {g}, version: {v})")), + (Some(g), None) => s.push_str(&format!(" (group: {g})")), + (None, Some(v)) => s.push_str(&format!(" (version: {v})")), + (None, None) => {} + } + if let Some(wp) = flow.worker_last_ping { + let age = (now - wp).num_seconds(); + let status = if worker_ping_stale.unwrap_or(false) { + "APPEARS DEAD OR CRASHED — worker ping is stale" + } else { + "worker still pinging — likely deadlocked or blocking on the state transition" + }; + s.push_str(&format!("\nWorker last ping: {wp} ({age}s ago) — {status}")); + if let Some(cjid) = flow.worker_current_job_id { + if cjid != flow.id { + s.push_str(&format!("\nWorker has since moved on to job {cjid}")); + } + } + } else { + s.push_str("\nWorker last ping: unknown (no worker_ping record)"); + } + let total_suffix = flow + .worker_memory_total + .map(|t| format!(" (total available: {})", fmt_mb(t))) + .unwrap_or_default(); + match (flow.worker_memory_usage, flow.worker_wm_memory_usage) { + (Some(host), Some(wm)) => s.push_str(&format!( + "\nWorker memory at last ping: host={}, wm process={}{}", + fmt_mb(host), + fmt_mb(wm), + total_suffix + )), + (Some(host), None) => s.push_str(&format!( + "\nWorker memory at last ping: host={}{}", + fmt_mb(host), + total_suffix + )), + (None, Some(wm)) => s.push_str(&format!( + "\nWorker memory at last ping: wm process={}{}", + fmt_mb(wm), + total_suffix + )), + (None, None) => { + if let Some(total) = flow.worker_memory_total { + s.push_str(&format!("\nWorker total memory: {}", fmt_mb(total))); + } + } + } + s + } else { + "\nWorker handling the flow: unknown (no worker recorded on v2_job_queue)" + .to_string() + }; + + let hint = match worker_ping_stale { + Some(true) => "\nThe worker's own ping is stale — it likely crashed, was killed (OOM?), or lost network connectivity. Check worker logs, k8s/docker events, or host dmesg around the worker's last ping time.", + Some(false) => "\nThe worker is still pinging — it is likely deadlocked or stuck in a blocking call during the state transition. Capture a stack trace from the worker process (e.g. via SIGQUIT or a debugger).", + None => "", + }; + + let service_logs_info = match (flow.worker_instance.as_deref(), flow.worker_last_ping) { + (Some(host), Some(wlp)) => { + let after = (wlp - chrono::Duration::seconds(90)) + .to_rfc3339_opts(chrono::SecondsFormat::Millis, true); + let before = (wlp + chrono::Duration::seconds(30)) + .to_rfc3339_opts(chrono::SecondsFormat::Millis, true); + let log_files_result = sqlx::query!( + "SELECT file_path, log_ts, ok_lines, err_lines + FROM log_file + WHERE hostname = $1 + AND log_ts BETWEEN $2::timestamp - interval '90 seconds' AND $2::timestamp + interval '30 seconds' + ORDER BY log_ts DESC + LIMIT 10", + host, + wlp.naive_utc(), + ) + .fetch_all(db) + .await; + let log_files = match log_files_result { + Ok(rows) => rows, + Err(e) => { + tracing::warn!( + "failed to query log_file for hanging-flow diagnostics (host={host}): {e:#}" + ); + Vec::new() + } + }; + if log_files.is_empty() { + format!( + "\nService logs: no log_file rows found for hostname '{host}' between {after} and {before}. If service log collection is enabled (requires S3/parquet), try /api/service_logs/list_files?after={after}&before={before}." + ) + } else { + let listed = log_files + .iter() + .map(|f| { + format!( + " - {} (log_ts: {}, ok_lines: {}, err_lines: {}) — GET /api/service_logs/get_log_file/{}/{}", + f.file_path, + f.log_ts, + f.ok_lines.unwrap_or(0), + f.err_lines.unwrap_or(0), + host, + f.file_path, + ) + }) + .collect::>() + .join("\n"); + format!("\nService logs for worker instance '{host}' around last ping ({after} to {before}) (download URLs require S3/parquet service log collection):\n{listed}") + } + } + _ => String::new(), + }; + let reason = format!( - "{} was hanging in between 2 steps. Last ping: {last_ping:?} (now: {now})", + "{} was hanging in between 2 steps. Last ping: {last_ping:?} (now: {now}){worker_info}{hint}{service_logs_info}", if flow.is_flow_step.unwrap_or(false) && flow.parent_job.is_some() { format!("Flow was cancelled because subflow {id} ({base_url}/run/{id}?workspace={workspace_id})") } else {