Files
windmill/backend/windmill-worker/src/result_processor.rs
Yoaquim Cintrón 7bf6ac2b69 fix: enrich OTEL spans with job_kind, trigger_kind, trigger, created_by, and script_hash (#8718)
Add five new attributes to the `job` and `job_postprocessing` tracing spans
so that OTEL-consuming backends (Sentry, Honeycomb, Datadog, etc.) can
filter and group telemetry by how a job was triggered and what type it is.

New span attributes:
- `job_kind`     — Script, Flow, AppScript, AIAgent, Preview, etc.
- `created_by`   — the user or system identity that queued the job
- `trigger_kind` — schedule, webhook, kafka, http, sqs, etc.
- `trigger`      — the schedule/trigger path (when applicable)
- `runnable_id`  — the id of the runnable that ran

Also adds `JobKind::as_str()` for a consistent lowercase string
representation, following the same pattern as `ScriptLang::as_str()`.

Existing attributes (job_id, workspace_id, script_path, language, tag,
flow_step_id, parent_job, root_job) are unchanged.

Note: the EE `full_job` span in `otel_ee.rs` and the log records emitted
by `job_logger_ee.rs` would also benefit from these attributes. This PR
covers only the public-repo spans; a follow-up EE change would propagate
the same fields to logs and the full_job span.
2026-04-05 14:11:22 +00:00

1211 lines
42 KiB
Rust

#[cfg(feature = "otel")]
use opentelemetry::trace::FutureExt;
use serde::{Deserialize, Serialize};
use sqlx::types::Json;
use std::{
collections::HashMap,
sync::{
atomic::{AtomicBool, AtomicU16, Ordering},
Arc,
},
};
use tracing::{field, Instrument};
#[cfg(not(feature = "otel"))]
use windmill_common::otel_oss::FutureExt;
use uuid::Uuid;
/// Set by the result processor when a WAC child completion makes suspend reach 0,
/// signaling the worker main loop to check for suspended jobs immediately.
pub static WAC_SUSPEND_READY: AtomicBool = AtomicBool::new(false);
use windmill_common::{
add_time,
error::{self, Error},
flow_status::FlowJobDuration,
jobs::JobKind,
utils::WarnAfterExt,
worker::{error_to_value, to_raw_value, Connection, WORKER_GROUP},
worker_group_job_stats::{accumulate_job_stats, flush_stats_to_db, JobStatsMap},
KillpillSender, DB,
};
#[cfg(feature = "benchmark")]
use windmill_common::bench::{BenchmarkInfo, BenchmarkIter};
use windmill_queue::{
append_logs, get_mini_completed_job, CanceledBy, FlowRunners, JobCompleted, MiniCompletedJob,
MiniPulledJob, ValidableJson, WrappedError, INIT_SCRIPT_TAG,
};
use serde_json::{json, value::RawValue, Value};
use tokio::{sync::Notify, task::JoinHandle};
use windmill_queue::{add_completed_job, add_completed_job_error};
use crate::{
bash_executor::ANSI_ESCAPE_RE,
common::{read_result, save_in_cache},
otel_oss::add_root_flow_job_to_otlp,
worker_flow::update_flow_status_after_job_completion,
JobCompletedReceiver, JobCompletedSender, SameWorkerSender, SendResult, SendResultPayload,
UpdateFlow, SAME_WORKER_REQUIREMENTS,
};
use windmill_common::client::AuthedClient;
#[derive(Debug, Deserialize)]
struct ErrorMessage {
message: String,
name: String,
}
async fn process_jc(
jc: JobCompleted,
worker_name: &str,
base_internal_url: &str,
db: &DB,
worker_dir: &str,
same_worker_tx: Option<&SameWorkerSender>,
job_completed_sender: &JobCompletedSender,
stats_map: &JobStatsMap,
killpill_rx: &tokio::sync::broadcast::Receiver<()>,
#[cfg(feature = "benchmark")] bench: &mut BenchmarkIter,
#[cfg(feature = "benchmark")] bench_infos: &mut BenchmarkInfo,
) {
let success: bool = jc.success;
let span = if success {
tracing::span!(
tracing::Level::INFO,
"job_postprocessing",
job_id = %jc.job.id, root_job = field::Empty, workspace_id = %jc.job.workspace_id, worker = %worker_name,tag = %jc.job.tag,
// hostname = %hostname,
language = field::Empty,
script_path = field::Empty,
flow_step_id = field::Empty,
parent_job = field::Empty,
job_kind = %jc.job.kind.as_str(),
created_by = %jc.job.created_by,
trigger_kind = field::Empty,
trigger = field::Empty,
script_hash = field::Empty,
otel.name = field::Empty,
success = %success,
labels = field::Empty,
)
} else {
tracing::span!(
tracing::Level::INFO,
"job_postprocessing",
job_id = %jc.job.id, root_job = field::Empty, workspace_id = %jc.job.workspace_id, worker = %worker_name,tag = %jc.job.tag,
// hostname = %hostname,
language = field::Empty,
script_path = field::Empty,
flow_step_id = field::Empty,
parent_job = field::Empty,
job_kind = %jc.job.kind.as_str(),
created_by = %jc.job.created_by,
trigger_kind = field::Empty,
trigger = field::Empty,
script_hash = field::Empty,
otel.name = field::Empty,
success = %success,
error.message = field::Empty,
error.name = field::Empty,
labels = field::Empty,
)
};
let rj = if let Some(root_job) = jc.job.flow_innermost_root_job {
root_job
} else {
jc.job.id
};
if let Some(labels) = jc.result.wm_labels() {
if !labels.is_empty() {
span.record("labels", labels.join(","));
}
}
windmill_common::otel_oss::set_span_parent(&span, &rj);
if let Some(lg) = jc.job.script_lang.as_ref() {
span.record("language", lg.as_str());
}
if let Some(step_id) = jc.job.flow_step_id.as_ref() {
span.record(
"otel.name",
format!("job_postprocessing {}", step_id).as_str(),
);
span.record("flow_step_id", step_id.as_str());
} else {
span.record("otel.name", "job postprocessing");
}
if let Some(parent_job) = jc.job.parent_job.as_ref() {
span.record("parent_job", parent_job.to_string().as_str());
}
if let Some(script_path) = jc.job.runnable_path.as_ref() {
span.record("script_path", script_path.as_str());
}
if let Some(root_job) = jc.job.flow_innermost_root_job.as_ref() {
span.record("root_job", root_job.to_string().as_str());
}
if let Some(trigger_kind) = jc.job.trigger_kind.as_ref() {
span.record("trigger_kind", trigger_kind.to_string().as_str());
}
if let Some(trigger) = jc.job.trigger.as_ref() {
span.record("trigger", trigger.as_str());
}
if let Some(script_hash) = jc.job.runnable_id.as_ref() {
span.record("script_hash", script_hash.to_string().as_str());
}
if !success {
if let Ok(result_error) = serde_json::from_str::<ErrorMessage>(jc.result.get()) {
span.record("error.message", result_error.message.as_str());
span.record("error.name", result_error.name.as_str());
}
}
// Extract stats info before moving jc
let duration_ms = jc.duration.clone();
let script_lang = jc.job.script_lang.clone();
let workspace_id = jc.job.workspace_id.clone();
let root_job = handle_receive_completed_job(
jc,
&base_internal_url,
&db,
worker_dir,
same_worker_tx,
&worker_name,
job_completed_sender.clone(),
killpill_rx,
#[cfg(feature = "benchmark")]
bench,
)
.instrument(span)
.warn_after_seconds(10)
.await;
if let Some(root_job) = root_job {
add_root_flow_job_to_otlp(&root_job, success);
#[cfg(feature = "benchmark")]
if bench_infos.count_top_level(root_job.id) {
bench_infos
.shared_iters
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
}
// Accumulate job stats if duration is available
if let Some(duration_ms) = duration_ms {
accumulate_job_stats(
stats_map,
&*WORKER_GROUP,
script_lang,
&workspace_id,
duration_ms,
)
.await;
}
}
enum JobCompletedRx {
JobCompleted(SendResult),
Killpill,
WakeUp,
}
pub fn start_background_processor(
job_completed_rx: JobCompletedReceiver,
job_completed_sender: JobCompletedSender,
same_worker_queue_size: Arc<AtomicU16>,
job_completed_processor_is_done: Arc<AtomicBool>,
wake_up_notify: Arc<Notify>,
last_processing_duration: Arc<AtomicU16>,
base_internal_url: String,
db: DB,
worker_dir: String,
same_worker_tx: SameWorkerSender,
worker_name: String,
killpill_tx: KillpillSender,
is_dedicated_worker: bool,
stats_map: JobStatsMap,
) -> JoinHandle<()> {
tokio::spawn(async move {
let mut has_been_killed = false;
let JobCompletedReceiver { bounded_rx, mut killpill_rx, unbounded_rx } = job_completed_rx;
#[cfg(feature = "benchmark")]
let mut infos = BenchmarkInfo::new(windmill_common::bench::shared_bench_iters());
// Start periodic stats flush task
let db_clone = db.clone();
let stats_map_clone = stats_map.clone();
let mut killpill_rx_clone = killpill_rx.resubscribe();
let flush_handle = tokio::spawn(async move {
let mut interval = tokio::time::interval(tokio::time::Duration::from_secs(900)); // Flush every 15 min
loop {
tokio::select! {
_ = interval.tick() => {
if let Err(e) = flush_stats_to_db(&db_clone, &stats_map_clone).await {
tracing::error!("Failed to flush worker group job stats: {}", e);
}
}
_ = killpill_rx_clone.recv() => {
tracing::info!("bg processor received killpill signal, flushing remaining stats");
break;
}
}
}
});
//if we have been killed, we want to drain the queue of jobs
while let Some(sr) = {
if has_been_killed {
tracing::info!("bg processor is killed, draining. same_worker_queue_size: {}, unbounded_rx: {}, bounded_rx: {}", same_worker_queue_size.load(Ordering::SeqCst), unbounded_rx.len(), bounded_rx.len())
}
if has_been_killed && same_worker_queue_size.load(Ordering::SeqCst) == 0 {
unbounded_rx
.try_recv()
.ok()
.map(JobCompletedRx::JobCompleted)
.or_else(|| bounded_rx.try_recv().ok().map(JobCompletedRx::JobCompleted))
} else {
tokio::select! {
biased;
result = unbounded_rx.recv_async() => {
result.ok().map(JobCompletedRx::JobCompleted)
}
result = bounded_rx.recv_async() => {
result.ok().map(JobCompletedRx::JobCompleted)
},
_ = wake_up_notify.notified() => {
tracing::info!("bg processor received wake up signal, checking if same worker queue is empty");
Some(JobCompletedRx::WakeUp)
},
_ = killpill_rx.recv() => {
tracing::info!("bg processor received killpill signal, queuing killpill job");
Some(JobCompletedRx::Killpill)
}
}
}
} {
#[cfg(feature = "benchmark")]
let mut bench = BenchmarkIter::new();
match sr {
JobCompletedRx::JobCompleted(SendResult {
result: SendResultPayload::JobCompleted(jc),
time,
}) => {
let is_init_script_and_failure =
!jc.success && jc.job.tag.as_str() == INIT_SCRIPT_TAG;
let is_dependency_job = matches!(
jc.job.kind,
JobKind::Dependencies | JobKind::FlowDependencies
);
#[cfg(feature = "benchmark")]
let bench_job_id = jc.job.id;
#[cfg(feature = "benchmark")]
let is_top_level_job = jc.job.parent_job.is_none();
process_jc(
jc,
&worker_name,
&base_internal_url,
&db,
&worker_dir,
Some(&same_worker_tx),
&job_completed_sender,
&stats_map,
&killpill_rx,
#[cfg(feature = "benchmark")]
&mut bench,
#[cfg(feature = "benchmark")]
&mut infos,
)
.warn_after_seconds(10)
.await;
if is_init_script_and_failure {
tracing::error!("init script errored, exiting");
killpill_tx.send();
break;
}
if is_dependency_job && is_dedicated_worker {
tracing::error!("Dedicated worker executed a dependency job, a new script has been deployed. Exiting expecting to be restarted.");
sqlx::query!(
"UPDATE config SET config = config WHERE name = $1",
format!("worker__{}", *WORKER_GROUP)
)
.execute(&db)
.await
.expect("update config to trigger restart of all dedicated workers at that config");
killpill_tx.send();
}
add_time!(bench, "job completed processed");
#[cfg(feature = "benchmark")]
{
if infos.add_iter(bench, bench_job_id, is_top_level_job) {
infos.shared_iters.fetch_add(1, Ordering::Relaxed);
}
}
last_processing_duration
.store(time.elapsed().as_secs() as u16, Ordering::SeqCst);
}
JobCompletedRx::JobCompleted(SendResult {
result:
SendResultPayload::UpdateFlow(UpdateFlow {
flow,
w_id,
success,
result,
worker_dir,
stop_early_override,
token,
}),
time,
}) => {
// let r;
tracing::info!(parent_flow = %flow, "updating flow status after job completion");
if let Err(e) = update_flow_status_after_job_completion(
&db,
&AuthedClient::new(
base_internal_url.to_string(),
w_id.clone(),
token.clone(),
None,
),
flow,
&Uuid::nil(),
&w_id,
success,
None,
Arc::new(result),
None,
true,
&same_worker_tx,
&worker_dir,
stop_early_override,
&worker_name,
job_completed_sender.clone(),
None,
&killpill_rx,
#[cfg(feature = "benchmark")]
&mut bench,
)
.await
{
tracing::error!("Error updating flow status after job completion for {flow} on {worker_name}: {e:#}");
}
#[cfg(feature = "benchmark")]
{
if infos.add_iter(bench, flow, true) {
infos.shared_iters.fetch_add(1, Ordering::Relaxed);
}
}
last_processing_duration
.store(time.elapsed().as_secs() as u16, Ordering::SeqCst);
}
JobCompletedRx::Killpill => {
tracing::info!("killpill job received, processing only same worker jobs");
has_been_killed = true;
}
JobCompletedRx::WakeUp => {}
}
}
// Flush any remaining stats before shutting down
tracing::info!("flushing remaining stats before shutting down");
let flush_result =
tokio::time::timeout(std::time::Duration::from_secs(10), flush_handle).await;
match flush_result {
Ok(Ok(())) => tracing::info!("Stats flushed successfully"),
Ok(Err(join_err)) => tracing::error!("Stats flush task failed: {}", join_err),
Err(_) => tracing::error!("Stats flush timed out after 10 seconds"),
}
job_completed_processor_is_done.store(true, Ordering::SeqCst);
tracing::info!("finished processing all completed jobs");
#[cfg(feature = "benchmark")]
{
infos
.write_to_file("profiling_result_processor.json")
.expect("write to file profiling");
}
})
}
async fn send_job_completed(job_completed_tx: JobCompletedSender, jc: JobCompleted) {
if let Err(e) = job_completed_tx
.send_job(jc, true)
.with_context(windmill_common::otel_oss::otel_ctx())
.await
{
tracing::error!("send job completed failed, triggering worker shutdown: {e:#}");
job_completed_tx.send_worker_killpill();
}
}
pub async fn process_result(
job: MiniCompletedJob,
result: error::Result<Arc<Box<RawValue>>>,
job_dir: &str,
job_completed_tx: JobCompletedSender,
mem_peak: i32,
canceled_by: Option<CanceledBy>,
cached_res_path: Option<String>,
token: &str,
result_columns: Option<Vec<String>>,
preprocessed_args: Option<HashMap<String, Box<RawValue>>>,
conn: &Connection,
duration: Option<i64>,
has_stream: bool,
flow_runners: Option<Arc<FlowRunners>>,
) -> error::Result<bool> {
match result {
Ok(result) => {
send_job_completed(
job_completed_tx,
JobCompleted {
job,
preprocessed_args,
result,
result_columns,
mem_peak,
canceled_by,
success: true,
cached_res_path,
token: token.to_string(),
duration,
has_stream: Some(has_stream),
from_cache: None,
flow_runners,
done_tx: None,
},
)
.with_context(windmill_common::otel_oss::otel_ctx())
.await;
Ok(true)
}
Err(e) => {
let error_value = match e {
Error::ExitStatus(program, i) => {
let res = read_result(job_dir, None).await.ok();
if res.as_ref().is_some_and(|x| !x.get().is_empty()) {
res.unwrap()
} else {
match conn {
Connection::Sql(db) => {
let last_10_log_lines = sqlx::query_scalar!(
"SELECT right(logs, 600) FROM job_logs WHERE job_id = $1 AND workspace_id = $2 ORDER BY created_at DESC LIMIT 1",
&job.id,
&job.workspace_id
).fetch_one(db).await.ok().flatten().unwrap_or("".to_string());
let log_lines = last_10_log_lines
.split("CODE EXECUTION ---")
.last()
.unwrap_or(&last_10_log_lines);
extract_error_value(
&program,
log_lines,
i,
job.flow_step_id.clone(),
)
}
Connection::Http(_) => {
to_raw_value(&"See logs for more details".to_string())
}
}
}
}
Error::ExecutionRawError(e) => to_raw_value(&e),
err @ _ => to_raw_value(&SerializedError {
message: format!("execution error:\n{err:#}",),
name: "ExecutionErr".to_string(),
step_id: job.flow_step_id.clone(),
exit_code: None,
}),
};
send_job_completed(
job_completed_tx,
JobCompleted {
job,
result: Arc::new(to_raw_value(&error_value)),
result_columns: None,
preprocessed_args: None,
mem_peak,
canceled_by,
success: false,
cached_res_path,
token: token.to_string(),
duration,
has_stream: Some(has_stream),
from_cache: None,
flow_runners,
done_tx: None,
},
)
.with_context(windmill_common::otel_oss::otel_ctx())
.await;
Ok(false)
}
}
}
pub async fn handle_receive_completed_job(
jc: JobCompleted,
base_internal_url: &str,
db: &DB,
worker_dir: &str,
same_worker_tx: Option<&SameWorkerSender>,
worker_name: &str,
job_completed_tx: JobCompletedSender,
killpill_rx: &tokio::sync::broadcast::Receiver<()>,
#[cfg(feature = "benchmark")] bench: &mut BenchmarkIter,
) -> Option<Arc<MiniPulledJob>> {
let token = jc.token.clone();
let workspace = jc.job.workspace_id.clone();
let client = AuthedClient::new(base_internal_url.to_string(), workspace, token, None);
let job = jc.job.clone();
let mem_peak = jc.mem_peak.clone();
let canceled_by = jc.canceled_by.clone();
let processed_completed_job = process_completed_job(
jc,
&client,
db,
&worker_dir,
same_worker_tx.clone(),
worker_name,
job_completed_tx.clone(),
killpill_rx,
#[cfg(feature = "benchmark")]
bench,
)
.warn_after_seconds(10)
.await;
match processed_completed_job {
Err(err) => {
handle_job_error(
db,
&client,
&job,
mem_peak,
canceled_by,
err,
false,
same_worker_tx.clone(),
&worker_dir,
worker_name,
job_completed_tx,
killpill_rx,
#[cfg(feature = "benchmark")]
bench,
)
.await;
None
}
Ok(r) => r,
}
}
pub async fn process_completed_job(
JobCompleted {
job,
result,
mem_peak,
success,
cached_res_path,
canceled_by,
duration,
result_columns,
preprocessed_args,
from_cache,
flow_runners,
done_tx,
..
}: JobCompleted,
client: &AuthedClient,
db: &DB,
worker_dir: &str,
same_worker_tx: Option<&SameWorkerSender>,
worker_name: &str,
job_completed_tx: JobCompletedSender,
killpill_rx: &tokio::sync::broadcast::Receiver<()>,
#[cfg(feature = "benchmark")] bench: &mut BenchmarkIter,
) -> error::Result<Option<Arc<MiniPulledJob>>> {
if success {
// println!("bef completed job{:?}", SystemTime::now());
if let Some(cached_path) = cached_res_path {
save_in_cache(db, client, &job, cached_path, result.clone()).await;
}
let is_flow_step = job.is_flow_step();
let parent_job = job.parent_job.clone();
let job_id = job.id.clone();
let workspace_id = job.workspace_id.clone();
let started_at = job.started_at.clone();
if job.flow_step_id.as_deref() == Some("preprocessor") {
// Do this before inserting to `v2_job_completed` for backwards compatibility
// when we set `flow_status->_metadata->preprocessed_args` to true.
sqlx::query!(
r#"UPDATE v2_job SET
args = '{"reason":"PREPROCESSOR_ARGS_ARE_DISCARDED"}'::jsonb,
preprocessed = TRUE
WHERE id = $1 AND preprocessed = FALSE"#,
job.id
)
.execute(db)
.await
.map_err(|e| {
Error::InternalErr(format!(
"error while deleting args of preprocessing step: {e:#}"
))
})?;
} else if let Some(preprocessed_args) = preprocessed_args {
// Update script args to preprocessed args
sqlx::query!(
"UPDATE v2_job SET args = $1, preprocessed = TRUE WHERE id = $2",
Json(preprocessed_args) as Json<HashMap<String, Box<RawValue>>>,
job.id
)
.execute(db)
.await?;
}
add_time!(bench, "pre add_completed_job");
let (_, duration, wac_job_ids) = add_completed_job(
db,
&job,
true,
false,
Json(&result),
result_columns,
mem_peak.to_owned(),
canceled_by.clone(),
false,
duration,
from_cache.unwrap_or(false),
)
.await?;
drop(job);
add_time!(bench, "add_completed_job END");
if is_flow_step {
if let Some(parent_job) = parent_job {
// tracing::info!(parent_flow = %parent_job, subflow = %job_id, "updating flow status (2)");
let r = update_flow_status_after_job_completion(
db,
client,
parent_job,
&job_id,
&workspace_id,
true,
canceled_by,
result,
started_at.map(|x| FlowJobDuration { started_at: x, duration_ms: duration }),
false,
&same_worker_tx.expect(SAME_WORKER_REQUIREMENTS).to_owned(),
&worker_dir,
None,
worker_name,
job_completed_tx,
flow_runners,
killpill_rx,
#[cfg(feature = "benchmark")]
bench,
)
.warn_after_seconds(10)
.await?;
add_time!(bench, "updated flow status END");
if let Some(done_tx) = done_tx {
done_tx
.send(())
.expect("done receiver should still be alive");
}
return Ok(r);
}
} else if let Some(parent_job) = parent_job {
// wac_job_ids is piggybacked from the duration write in
// add_completed_job — no extra query needed.
if let Some(job_ids) = wac_job_ids {
if let Ok(Some(_)) = handle_wac_child_completion(
db,
&job_id,
parent_job,
&workspace_id,
result,
true,
job_ids,
)
.await
{
if let Some(done_tx) = done_tx {
done_tx
.send(())
.expect("done receiver should still be alive");
}
return Ok(None);
}
}
}
} else {
let result = add_completed_job_error(
db,
&job,
mem_peak.to_owned(),
canceled_by.clone(),
serde_json::from_str(result.get()).unwrap_or_else(
|_| json!({ "message": format!("Non serializable error: {}", result.get()) }),
),
worker_name,
false,
None,
)
.await?;
if job.is_flow_step() {
if let Some(parent_job) = job.parent_job {
tracing::error!(parent_flow = %parent_job, subflow = %job.id, "process completed job error, updating flow status");
let r = update_flow_status_after_job_completion(
db,
client,
parent_job,
&job.id,
&job.workspace_id,
false,
canceled_by,
Arc::new(serde_json::value::to_raw_value(&result).unwrap()),
duration.and_then(|d| {
job.started_at.map(|started_at| FlowJobDuration {
started_at: started_at,
duration_ms: d,
})
}),
false,
&same_worker_tx.expect(SAME_WORKER_REQUIREMENTS).to_owned(),
&worker_dir,
None,
worker_name,
job_completed_tx,
flow_runners,
killpill_rx,
#[cfg(feature = "benchmark")]
bench,
)
.warn_after_seconds(10)
.await?;
if let Some(done_tx) = done_tx {
done_tx
.send(())
.expect("done receiver should still be alive");
}
return Ok(r);
}
} else if let Some(parent_job) = job.parent_job {
// WAC child failed — query job_ids from parent (errors are rare,
// so the extra read is acceptable here).
let job_ids_json: Option<Option<Value>> = sqlx::query_scalar(
"SELECT workflow_as_code_status->'_checkpoint'->'pending_steps'->'job_ids' \
FROM v2_job_status WHERE id = $1",
)
.bind(&parent_job)
.fetch_optional(db)
.await?;
if let Some(Some(job_ids)) = job_ids_json {
let err_result = Arc::new(serde_json::value::to_raw_value(&result).unwrap());
if let Ok(Some(_)) = handle_wac_child_completion(
db,
&job.id,
parent_job,
&job.workspace_id,
err_result,
false,
job_ids,
)
.await
{
if let Some(done_tx) = done_tx {
done_tx
.send(())
.expect("done receiver should still be alive");
}
return Ok(None);
}
}
}
}
return Ok(None);
}
/// Handle a WAC v2 child job completion.
/// Returns Ok(Some(())) if the parent was a WAC job and was handled,
/// Ok(None) if the parent is not a WAC job (caller should fall through).
///
/// CONCURRENCY: Multiple parallel children may complete simultaneously on
/// different workers. We use atomic SQL operations throughout:
/// - `completed_steps` is merged via `jsonb_set(... || jsonb_build_object(...))`
/// — PostgreSQL serialises concurrent UPDATEs on the same row, so each
/// worker sees the previous worker's writes.
/// - The suspend counter (set to N at dispatch time) is decremented atomically
/// with `RETURNING` to determine the "all done" condition.
pub(crate) async fn handle_wac_child_completion(
db: &DB,
child_job_id: &Uuid,
parent_job_id: Uuid,
workspace_id: &str,
result: Arc<Box<RawValue>>,
success: bool,
job_ids_value: Value,
) -> error::Result<Option<()>> {
let job_ids = match job_ids_value {
Value::Object(m) => m,
_ => return Ok(None), // Not a WAC parent or no pending steps
};
let child_id_str = child_job_id.to_string();
let step_key = job_ids.iter().find_map(|(key, val)| {
if val.as_str() == Some(&child_id_str) {
Some(key.clone())
} else {
None
}
});
let step_key = match step_key {
Some(k) => k,
None => {
if !success {
// No step key and failed — can't store error, fail parent immediately
tracing::error!(
parent_job = %parent_job_id,
child_job = %child_job_id,
"WAC v2 child job failed but no step key found, failing parent"
);
sqlx::query!(
"UPDATE v2_job_queue SET suspend = 0, suspend_until = NULL WHERE id = $1",
parent_job_id,
)
.execute(db)
.await?;
let parent_mini = get_mini_completed_job(&parent_job_id, workspace_id, db).await?;
if let Some(parent_mini) = parent_mini {
let child_err: Value =
serde_json::from_str(result.get()).unwrap_or(Value::Null);
let err_value = json!({
"message": format!("WAC child job {} failed (no step key)", child_job_id),
"error": child_err,
});
let _ = windmill_queue::add_completed_job_error(
db,
&parent_mini,
0,
None,
err_value,
"wac_child_handler",
false,
None,
)
.await;
}
return Ok(Some(()));
}
tracing::warn!(
parent_job = %parent_job_id,
child_job = %child_job_id,
"WAC v2 child completed but no matching step key found in checkpoint, decrementing suspend to avoid parent hang"
);
// Still decrement suspend so the parent doesn't hang indefinitely
let _ = sqlx::query_scalar!(
"UPDATE v2_job_queue \
SET suspend = GREATEST(suspend - 1, 0) \
WHERE id = $1 \
RETURNING suspend",
parent_job_id,
)
.fetch_optional(db)
.await?;
return Ok(Some(()));
}
};
// Build result — wrap errors with _error marker so workflow try/catch can handle them
let result_value: Value = if success {
serde_json::from_str(result.get()).unwrap_or(Value::Null)
} else {
let child_err: Value = serde_json::from_str(result.get()).unwrap_or(Value::Null);
tracing::info!(
parent_job = %parent_job_id,
child_job = %child_job_id,
step_key = %step_key,
"WAC v2 child job failed, storing error for workflow try/catch"
);
json!({
"__wmill_error": true,
"message": format!("WAC task '{}' failed (child job {})", step_key, child_job_id),
"child_job_id": child_job_id.to_string(),
"step_key": step_key,
"result": child_err,
})
};
tracing::info!(
parent_job = %parent_job_id,
child_job = %child_job_id,
step_key = %step_key,
success = success,
"WAC v2 child job completed"
);
// Use a transaction to ensure completed_steps merge + suspend decrement
// are atomic. Without this, a crash between the two could strand the parent.
let result_json = serde_json::to_value(&result_value)
.map_err(|e| error::Error::InternalErr(format!("Failed to serialize step result: {e}")))?;
let mut tx = db.begin().await?;
// Merge the completed step into the checkpoint.
// Uses `|| jsonb_build_object(key, value)` so concurrent children on
// different workers don't overwrite each other — PostgreSQL serialises
// concurrent UPDATEs on the same row and each sees the previous write.
sqlx::query(
"UPDATE v2_job_status SET workflow_as_code_status = jsonb_set(
workflow_as_code_status,
'{_checkpoint,completed_steps}',
COALESCE(workflow_as_code_status->'_checkpoint'->'completed_steps', '{}'::jsonb)
|| jsonb_build_object($2::text, $3::jsonb)
) WHERE id = $1",
)
.bind(&parent_job_id)
.bind(&step_key)
.bind(&result_json)
.execute(&mut *tx)
.await
.map_err(|e| error::Error::InternalErr(format!("Failed to add WAC completed step: {e}")))?;
// Decrement the suspend counter. The counter was set to N (number of
// children) at dispatch time. When it reaches 0 all children are done.
// Keep suspend_until non-null so the suspended pull query
// (`WHERE suspend_until IS NOT NULL AND suspend <= 0`) picks up the parent.
let new_suspend: Option<i32> = sqlx::query_scalar!(
"UPDATE v2_job_queue \
SET suspend = GREATEST(suspend - 1, 0) \
WHERE id = $1 \
RETURNING suspend",
parent_job_id,
)
.fetch_optional(&mut *tx)
.await?;
let all_done = new_suspend == Some(0);
if all_done {
// Clear pending_steps from checkpoint since all children are complete.
// This is cosmetic — the next replay will overwrite it anyway — but
// keeps the checkpoint clean for frontend display.
let _ = sqlx::query(
"UPDATE v2_job_status SET workflow_as_code_status = \
workflow_as_code_status #- '{_checkpoint,pending_steps}' \
WHERE id = $1",
)
.bind(&parent_job_id)
.execute(&mut *tx)
.await;
}
tx.commit().await?;
if all_done {
tracing::info!(
parent_job = %parent_job_id,
"WAC v2 all child jobs completed, unsuspending parent"
);
WAC_SUSPEND_READY.store(true, Ordering::Relaxed);
}
Ok(Some(()))
}
pub async fn handle_non_flow_job_error(
db: &DB,
job: &MiniCompletedJob,
mem_peak: i32,
canceled_by: Option<CanceledBy>,
err_string: String,
err_json: Value,
worker_name: &str,
) -> Result<WrappedError, Error> {
append_logs(
&job.id,
&job.workspace_id,
format!("Unexpected error during job execution:\n{err_string}"),
&db.into(),
)
.await;
add_completed_job_error(
db,
job,
mem_peak,
canceled_by,
err_json,
worker_name,
false,
None,
)
.await
}
#[tracing::instrument(name = "job_error", level = "info", skip_all, fields(job_id = %job.id))]
pub async fn handle_job_error(
db: &DB,
client: &AuthedClient,
job: &MiniCompletedJob,
mem_peak: i32,
canceled_by: Option<CanceledBy>,
err: Error,
unrecoverable: bool,
same_worker_tx: Option<&SameWorkerSender>,
worker_dir: &str,
worker_name: &str,
job_completed_tx: JobCompletedSender,
killpill_rx: &tokio::sync::broadcast::Receiver<()>,
#[cfg(feature = "benchmark")] bench: &mut BenchmarkIter,
) {
let err_string = format!("{}: {}", err.name(), err.to_string());
let err_json = error_to_value(&err);
let update_job_future = || async {
handle_non_flow_job_error(
db,
job,
mem_peak,
canceled_by.clone(),
err_string,
err_json.clone(),
worker_name,
)
.warn_after_seconds(10)
.await
};
let update_job_future = if job.is_flow_step() || job.is_flow() {
let (flow, job_status_to_update) = if let Some(parent_job_id) = job.parent_job {
if let Err(e) = update_job_future().await {
tracing::error!(
"error updating job future for job {} for handle_job_error: {e:#}",
job.id
);
}
(parent_job_id, job.id)
} else {
(job.id, Uuid::nil())
};
let wrapped_error = WrappedError { error: err_json.clone() };
tracing::error!(parent_flow = %flow, subflow = %job_status_to_update, "handle job error, updating flow status: {err_json:?}");
let updated_flow = update_flow_status_after_job_completion(
db,
client,
flow,
&job_status_to_update,
&job.workspace_id,
false,
canceled_by.clone(),
Arc::new(serde_json::value::to_raw_value(&wrapped_error).unwrap()),
None,
unrecoverable,
&same_worker_tx.expect(SAME_WORKER_REQUIREMENTS).clone(),
worker_dir,
None,
worker_name,
job_completed_tx.clone(),
None,
killpill_rx,
#[cfg(feature = "benchmark")]
bench,
)
.await;
if let Err(err) = updated_flow {
if let Some(parent_job_id) = job.parent_job {
if let Ok(Some(parent_job)) =
get_mini_completed_job(&parent_job_id, &job.workspace_id, db)
.warn_after_seconds(10)
.await
{
let e = json!({"message": err.to_string(), "name": "InternalErr"});
append_logs(
&parent_job.id,
&job.workspace_id,
format!("Unexpected error during flow job error handling:\n{err}"),
&db.into(),
)
.await;
let _ = add_completed_job_error(
db,
&parent_job,
mem_peak,
canceled_by,
e,
worker_name,
false,
None,
)
.warn_after_seconds(10)
.await;
}
}
}
None
} else {
Some(update_job_future)
};
if let Some(f) = update_job_future {
let _ = f().await;
}
}
#[derive(Debug, Serialize)]
pub struct SerializedError {
pub message: String,
pub name: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub step_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub exit_code: Option<i32>,
}
pub fn extract_error_value(
program: &str,
log_lines: &str,
i: i32,
step_id: Option<String>,
) -> Box<RawValue> {
return to_raw_value(&SerializedError {
message: format!(
"exit code for \"{program}\": {i}, last log lines:\n{}",
ANSI_ESCAPE_RE.replace_all(log_lines.trim(), "").to_string()
),
name: "ExecutionErr".to_string(),
step_id,
exit_code: Some(i),
});
}