feat: enrich hanging flow error with worker and service log info (#8800)

* feat: enrich hanging flow error with worker and service log info

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

* fix: address PR review on hanging flow diagnostics

- Widen log_file lookup window to [-90s, +30s] around worker last ping
  so the batch containing the crash is captured (log files are
  minute-aligned; looking forward only was missing the relevant bucket).
- Log a warning on log_file query errors instead of silently swallowing,
  so a misconfigured table is not reported as "no log files found".
- Note that service log download URLs require S3/parquet collection.
- Fix memory display when only worker_memory_total is known.
- Regenerate sqlx offline cache for the new/modified queries.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-04-10 13:19:20 -04:00
committed by GitHub
parent b783bf2d83
commit 59c457a138
4 changed files with 286 additions and 60 deletions

View File

@@ -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"
}

View File

@@ -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<str>\", 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<str>",
"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"
}

View File

@@ -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<str>\", 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<str>",
"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"
}

View File

@@ -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<str>", r.ping AS last_ping, j.same_worker AS "same_worker?"
COALESCE(s.flow_status, s.workflow_as_code_status) AS "flow_status: Box<str>", 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::<Vec<_>>()
.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 {