Files
windmill/backend/windmill-worker/src/result_processor.rs
T
Ruben Fiszel 12f92e3ab7 [ee] feat(backend): native script retry without one-step-flow wrapping (#9688)
* feat(backend): native script retry without one-step-flow wrapping

Schedules and data pipelines that retry a single script previously wrapped
it in a one-step flow (JobKind::SingleStepFlow), creating extra job rows, a
v2_job_status row, and UI projection complexity. This adds native retry on a
plain JobKind::Script job.

- RetrySettings: flatten Retry into a deduped retry_settings table, carried
  via the existing runnable_settings_handle (lazy, off the hot path).
- push() materializes a bare-script-with-retry SingleStepFlow into a native
  Script job (gated on min-version + no handlers/retry_if).
- add_completed_job re-pushes the next attempt on failure with backoff,
  tracking the attempt counter in v2_job_queue.extras and the chain via
  parent_job; schedule completion handlers fire only on the terminal attempt.
- frontend: ScriptRetryChain shows the attempt chain on the run page.

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

* feat(backend): native retry_if eval + per-occurrence schedule handlers

Extends native script retry to the two cases that previously stayed on the
one-step-flow path:

- retry_if: evaluated natively on the failure path via a feature-gated
  windmill-jseval dep (quickjs) over the failure result + flow_input; push
  materializes such policies natively only when quickjs is available.
- on_failure_times / on_recovery: apply_schedule_handlers now resolves each
  past scheduled occurrence's terminal status across its native-retry chain
  (root OR any parent_job=root child succeeded) and excludes the current
  occurrence, so the counting is per-occurrence rather than per-attempt.

All scheduled-script retries now go native (schedule.rs gate removed).

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

* refactor(backend): always materialize retry_if natively; unsupported without quickjs

retry_if is evaluated by the worker (which always has quickjs), not the
pusher, so gating materialization on the pusher's feature was wrong. The
flow path was never a real fallback either — the flow runtime needs quickjs
to evaluate retry_if too. retry_if now always goes native; on a worker
without quickjs it is unsupported and fails closed (no retry).

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

* feat(backend): un-park asset-cascade (pipeline) retry

Native retry resolves the blocker that parked pipeline retry: a retried
subscriber is now a Script job (not a one-step flow / flow step), so it
stays eligible for asset dispatch and can trigger its own downstream on
recovery.

- scripts.rs: persist // retry <count> [<delay>] to script_trigger on asset
  edges (was dropped with a TODO warning).
- asset_dispatch.rs: is_eligible_kind keys off flow_step_id, not parent_job,
  so native-retry attempts dispatch on success while flow steps stay excluded.
- tests: retry-bearing subscriber now dispatches as a native Script carrying
  the policy in runnable_settings_handle; native-retry attempt is eligible.

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

* fix(backend): cap native retry interval, lazy result serialization, idempotent retry push

Hardening from a self-review of the native retry path:
- Cap the backoff at MAX_RETRY_INTERVAL to match the flow-runtime path
  (evaluate_retry); the exponential formula could otherwise schedule up to
  ~18h vs the flow path's 6h.
- Serialize the failure result lazily (only when a retry_if policy needs it),
  so the common failure no longer pays the serialization on the failure path.
- Push each retry with a deterministic id per (root, attempt). If a worker
  dies between enqueueing the retry and finalizing the current attempt, the
  reaper re-handles the attempt and lands here again — push rejects the
  duplicate id, so the retry is enqueued exactly once (no double-retry).

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

* fix(backend): defer schedule handlers idempotently on retry-push replay (review P1)

Address local-review findings:
- P1: retry_pending was derived from the retry push *result*, so on a worker
  crash + reaper replay the duplicate-id push returned Err → retry_pending
  flipped to false → apply_schedule_handlers fired for the non-terminal
  attempt (and the terminal attempt later fired them again). Pre-check whether
  the deterministic retry id already exists and report it as pending without
  re-pushing, so the handler-deferral invariant is crash-idempotent too.
- P2: refresh the stale 'wrap the script in a one-step flow' comment in the
  asset-cascade retry push — it now materializes a native Script.
- Add RetrySettings <-> Retry round-trip unit tests (clamping edges).

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

* test(backend): native retry chain + per-occurrence status sqlx tests

Close the two integration-test gaps flagged in local review:
- chains_attempts_and_is_idempotent: drives maybe_enqueue_native_script_retry
  through attempt0 -> retry1 -> retry2 -> exhausted (counter, backoff, max-attempts)
  and asserts crash-replay idempotency (the P1 fix: a replayed completion reports
  pending without double-enqueueing).
- per_occurrence_status_counts_recovered_as_success: pins the exact per-occurrence
  terminal-status query from jobs_ee::apply_schedule_handlers — a retried-but-
  recovered occurrence counts as success, retries (parent_job set) are excluded
  from occurrence counting, and the current occurrence is excluded.
- canceled_job_does_not_retry: cancellation wins over a pending retry.

Runtime sqlx API (no .sqlx cache entry needed).

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

* fix(frontend): exclude schedule handlers from the retry-attempt chain

The retry chain listed all script children of the root by parent_job, but
schedule completion handlers (on_failure/on_recovery/on_success) are also
script children — when the occurrence has no retries, the handler's parent is
the root itself, so a successful, never-retried job rendered a bogus
'Retries (1)' badge pointing at the handler. Filter children to re-runs of the
same script (matching script_hash); real retries keep the root's hash, handlers
run a different script.

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

* feat(frontend): surface schedule handlers on the run page

Extend the run-page chain component with schedule completion handlers:
- A 'Handlers' row on a scheduled job links to the on_failure/on_recovery/
  on_success runs that fired for that occurrence (found as children of the
  terminal attempt, identified by their synthetic created_by).
- A handler's own run page now shows a 'Failure/Recovery/Success handler'
  label with a link back to the run it handled and its schedule. on_recovery
  and on_success share created_by, disambiguated by the recovery-only
  error_started_at arg.

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

* fix(backend): restore folder_default_permissioned_as sqlx caches dropped by prepare

An earlier `cargo sqlx prepare` on this branch ran before #8801's
folder_default_permissioned_as test merged in, so it pruned the 3 query caches
that test needs; cargo_test then failed under SQLX_OFFLINE. Restore them from main.

* fix(backend): only cascade assets from native retry attempts, not handlers (review P1)

is_eligible_kind keyed dispatch on flow_step_id alone, so every parented Script
child became asset-eligible — including schedule/error/recovery handlers (Script
jobs with parent_job set and no flow_step_id). A handler that declares assets
would then trigger a cascade the old parent_job IS NULL guard prevented. Gate
parented jobs on being a genuine retry attempt: a re-run of the SAME runnable as
its chain parent (handlers run a different script). Runtime query, no sqlx cache.

* fix(backend): cache the private-gated retry_setting asset-dispatch test query

The same prepare-without-private that dropped the folder_default caches also
pruned the cache for the retry_setting_dispatches_subscriber_as_native_script
test query (asset_trigger_dispatch.rs:721). Regenerated with --features private.

* fix(backend): exclude handler children from per-occurrence recovery (review)

A scheduled occurrence's on_failure/on_success handler runs as a successful
child (parent_job = occurrence), and the per-occurrence success EXISTS counted
ANY successful child — so a failed occurrence whose error handler succeeded was
marked 'recovered', breaking on_recovery (test_script/flow_schedule_handlers in
the merge) and on_failure_times counting. EE query now scopes the EXISTS to
same-runnable children (only native retry attempts); regenerate sqlx cache + bump
ee-repo-ref. native_retry_test gains a handler-child regression case.

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

* test(backend): scheduled-script retry is a native Script, not SingleStepFlow

test_push_script_with_retry / test_try_schedule_with_retry (from main) asserted
the old SingleStepFlow wrapping for scheduled-script retry; this PR makes it a
native Script. Update both to assert kind='script' and that the retry policy is
carried via runnable_settings_handle.

* fix(backend): preserve dedicated_worker on native retry + saturate count casts (cubic)

Address cubic CI review:
- P1: the SingleStepFlow->native Script materialization dropped dedicated_worker,
  so a dedicated-worker scheduled script lost its dedicated pool on retry. Resolve
  it from the script row in push so the materialized Script keeps the dedicated tag.
- P2: saturate the u32->i32 retry-attempt narrowings (RetrySettings::from) and the
  u32->i16 // retry count narrowing (scripts.rs) instead of wrapping.

* fix(backend): use a retry-specific signal, not runnable equality (codex review)

Address Codex CI review:
- P1: is_native_retry_attempt treated any same-runnable parented Script child as
  a retry. WAC v2 inline children have that exact shape, so an inline child of an
  asset producer would cascade. Use a retry-specific signal instead: the job
  carries a retry_settings policy (always re-inserted by maybe_enqueue) and has no
  flow_innermost_root_job. Apply the same flow_innermost guard to the EE
  per-occurrence EXISTS (WAC inline children must not count as a recovery).
- P1: the deterministic retry-id pre-check raced with push; a concurrent duplicate
  now resolves as 'retry pending' (re-check on the duplicate-id error) instead of
  flipping retry_pending to false and firing handlers early.
- Tests: native_retry + asset_trigger_dispatch gain WAC-inline-child cases.

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

* refactor(backend): explicit native_retry_attempt marker, drop heuristics

Replace the per-site "is this a retry?" inference (parent_job + runnable match +
flow_innermost / retry_settings) with one explicit marker: a sparse
native_retry_attempt(job_id, attempt) table, written in maybe_enqueue. The marker
also carries the attempt counter (previously in v2_job_queue.extras), so it's the
single source of truth.

- asset_dispatch: is_native_retry_attempt is now one indexed EXISTS on the marker.
- EE per-occurrence query: joins the marker instead of guessing by runnable/flow_innermost.
- maybe_enqueue: reads/writes the marker (persistent) instead of queue extras.
- Lifecycle: swept with the job in retention (log_cleanup), no FK to keep bulk delete cheap.
- Eliminates handler / WAC-inline-child misclassification by construction.

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

* fix(backend): sweep native_retry_attempt markers in the periodic retention path too (codex)

The marker has no FK and relies on retention cleanup; log_cleanup.rs swept it but
the periodic monitor.rs path deleted v2_job rows without it, orphaning markers.
Add the same WHERE job_id = ANY(...) sweep there.

* fix(backend): widen native_retry_attempt.attempt to integer (cubic)

The smallint column was cast to/from u32 and could wrap a retry chain longer than
i16::MAX into premature exhaustion. Use integer, matching the retry policy's i32
attempt count, so no narrowing occurs on the maybe_enqueue read/write path.

* feat(frontend): mark retries via is_retry on listJobs; drop SAVEPOINT

- Expose an is_retry flag on jobs (UnifiedJob/CompletedJob/QueuedJob + openapi),
  computed from the native_retry_attempt marker. The run-page chain now filters
  retry attempts by is_retry instead of the script_hash heuristic, so WAC v2
  inline children (same script, parent_job) no longer render as retries (codex).
- Revert the marker-cleanup SAVEPOINT (an unused pattern in this codebase): keep
  the plain catch-and-continue matching the other side-table deletes; the table is
  created by a startup migration so it always exists when cleanup runs.

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

* fix(backend): mark is_retry sqlx(default) so non-list job queries can omit it

The single-job GET query maps directly to CompletedJob/QueuedJob via FromRow but
does not select is_retry, which errored with "no column found". Only the list
endpoint populates the marker; #[sqlx(default)] lets every other query omit the
column and default to None.

* feat(backend): select is_retry in single-job GET too for consistency

The list endpoint already exposes the marker; populate it on the single-job GET
(both completed and queued variants) as well so a run loaded directly reflects
its retry status. #[sqlx(default)] stays as a safety net for any other query.

* feat(backend): reap orphaned native_retry_attempt markers via periodic sweep

The marker has no FK to v2_job (to keep the hot bulk retention delete cheap), so
direct job deletions (workspace/job delete, schedule clearing) would leave marker
rows orphaned. Rather than add explicit cleanup to every v2_job delete site (which
must then be remembered for every future path), reap orphans in the periodic
delete_expired_items pass: DELETE FROM native_retry_attempt WHERE NOT EXISTS (the
job). The table is sparse so the anti-join drives off it and probes v2_job by PK —
cheap. Retention still sweeps markers inline (keeps the table small so this stays
cheap); a transient orphan is harmless (nothing reads is_retry for a gone job).

* fix(frontend): include flow handlers in retry chain handler row (codex)

Schedule on_failure/on_recovery/on_success handlers can be flow paths (flow/...),
whose handler job is a flow, not a script. The chain fetched children with
jobKinds:'script', hiding flow handlers. Drop the kind filter — retry attempts
are still selected by is_retry and handlers by created_by, so both kinds surface.

* fix(backend): carry concurrency/debouncing settings into native retries

maybe_enqueue re-pushed the next attempt with ConcurrencySettings/DebouncingSettings
::default(), dropping the script/pipeline concurrency settings the failed job carried
in its runnable_settings_handle. A retry of a concurrency-limited script then inserted
no concurrency_key and ran unbounded. Resolve both from the same handle (cached) and
pass them in the payload, which push forwards to the materialized retry. Adds a
regression test asserting the retry's handle resolves to the concurrency settings.

* fix(backend): carry concurrency/debounce into scheduled-retry root + document retry-helper auth (codex)

P1a (schedule.rs): the scheduled-retry materialization fetched the script's
concurrency/debounce settings but passed ConcurrencySettings/DebouncingSettings
::default() into the SingleStepFlow payload, so the root attempt's handle held only
the retry policy and the whole chain ran unbounded. Pass the fetched settings.
Regression test asserts the root handle resolves to retry + concurrency.

P1b (jobs.rs): document maybe_enqueue_native_script_retry's authorization contract
— it is pub only for the integration test; the sole production caller is the worker
completion path passing a DB-derived, already-authorized MiniCompletedJob.

* docs(backend): attach native-retry auth contract to the function itself (codex)

The doc block was merged with eval_retry_if's doc and bound to that function,
leaving maybe_enqueue_native_script_retry undocumented. Split them: eval_retry_if
keeps its own doc; the native-retry + authorization contract now sits directly
above maybe_enqueue_native_script_retry.

* docs(backend): regenerate served openapi-deref with is_retry + fix stale comments (codex)

- Regenerate openapi-deref.{yaml,json} (served from lib.rs): they were stale since
  1.734.0 and lacked is_retry on QueuedJob/CompletedJob, so clients reading the
  served spec couldn't see the field. Now current at 1.739.0.
- schedule.rs: a retry_if gate is evaluated at failure time and fails closed without
  quickjs (no retry); it does not fall back to a flow path.
- windmill-types jobs.rs: is_retry is selected by both the list and single-job GET
  endpoints (not list-only).

* docs(backend): fix remaining stale retry_if/quickjs comments (codex)

The retry_if block and the push materialization comments claimed push keeps
retry_if on a flow path / the worker always has quickjs. The code always
materializes native retry and the no-quickjs eval_retry_if path fails closed —
correct the comments to that constraint.

* docs(backend): fix stale quickjs-fallback + schedule-handler-restriction comments (codex)

- Cargo.toml quickjs feature: without quickjs a retry_if gate cannot be evaluated
  and the job does not retry (no one-step-flow fallback).
- jobs.rs handler-defer comment: apply_schedule_handlers resolves per-occurrence
  failure/recovery status across the retry chain, so the old 'restricted to
  schedules whose handlers don't need per-occurrence counting' claim is dropped.

---------

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-25 13:22:26 +00:00

1344 lines
48 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, asset_dispatch, get_mini_completed_job, is_pre_shaped_wm_failure_result,
CanceledBy, FlowRunners, JobCompleted, MiniCompletedJob, MiniPulledJob, ValidableJson,
WrappedError, INIT_SCRIPT_TAG, MANUAL_FAILURE_ERROR_NAME,
};
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,
}
#[derive(Debug, Deserialize)]
struct NestedErrorMessage {
error: ErrorMessage,
}
/// Extract `{ name, message }` from a result. Accepts both the standard
/// top-level shape (regular runtime errors) and the nested `{ error: { name,
/// message }, ... }` shape produced by the wm_failure injection.
///
/// For wm_failure-injected results, we prefer the nested error: a successful
/// run may legitimately contain top-level `name`/`message` fields (user data
/// named `name`/`message`), and we want OTel to record the ManualFailure
/// rather than the user's sibling fields.
fn extract_error_message(raw: &str) -> Option<ErrorMessage> {
let nested = serde_json::from_str::<NestedErrorMessage>(raw)
.ok()
.map(|n| n.error);
if matches!(&nested, Some(em) if em.name == MANUAL_FAILURE_ERROR_NAME) {
return nested;
}
if let Ok(em) = serde_json::from_str::<ErrorMessage>(raw) {
return Some(em);
}
nested
}
/// Returns the post-processing `success` value (after any `wm_failure`
/// override). Callers use this to make worker-loop decisions that depend on
/// whether the job ultimately succeeded — e.g. the init-script killpill.
async fn process_jc(
mut 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,
) -> bool {
// Parse `wm_labels` and `wm_failure` together (single `from_str`)
// so we don't deserialize the whole result twice on every job.
let metadata = jc.result.result_metadata();
// If the script returned a `wm_failure: <string>` field in its
// result, tag the run as a failure. Inject an `error: { name, message }`
// at the top level so error handlers / UI / OTel see the standard error
// shape, while preserving sibling fields (`windmill_status_code`,
// `windmill_content_type`, `windmill_headers`, the user's data) at the
// top level so sync webhook responses still honor them.
if jc.success {
if let Some(failure_msg) = metadata.wm_failure.as_ref() {
if let Ok(Value::Object(mut map)) = serde_json::from_str::<Value>(jc.result.get()) {
map.insert(
"error".to_string(),
json!({ "name": MANUAL_FAILURE_ERROR_NAME, "message": failure_msg }),
);
if let Ok(raw) = serde_json::value::to_raw_value(&Value::Object(map)) {
jc.result = Arc::new(raw);
}
}
jc.success = false;
}
}
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,
otel.status_code = "ERROR",
otel.status_message = 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) = metadata.wm_labels.as_ref() {
if !labels.is_empty() {
span.record("labels", labels.join(","));
}
}
// The secondary `job_postprocessing` span stays on the UUID-derived context
// (MiniCompletedJob carries no args, so the inbound traceparent isn't
// available here); the primary job span is relocated in `create_span_with_name`.
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 Some(result_error) = extract_error_message(jc.result.get()) {
span.record("error.message", result_error.message.as_str());
span.record("error.name", result_error.name.as_str());
span.record(
"otel.status_message",
crate::worker::truncate_description(&result_error.message).as_str(),
);
} else {
span.record("otel.status_message", "Job failed");
}
}
// 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;
}
success
}
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 = 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 returns the post-override success value so a
// job that flipped to failure via `wm_failure` still
// triggers the init-script killpill.
let final_success = 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 && !final_success {
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<crate::worker::JobOutcome> {
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(crate::worker::JobOutcome::Completed)
}
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,
}),
};
// Use the structured error message that was just extracted (the
// user-facing script error) rather than the generic Error string.
// Pull `.message` out of the JSON object if present, otherwise
// accept a bare string (e.g. agent-worker "See logs for more
// details"), and only fall back to "Job failed" when the value
// carries no readable description.
let description = serde_json::from_str::<serde_json::Value>(error_value.get())
.ok()
.and_then(|value| {
value
.get("message")
.and_then(|m| m.as_str())
.or_else(|| value.as_str())
.map(|m| crate::worker::truncate_description(m))
})
.unwrap_or_else(|| "Job failed".to_string());
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(crate::worker::JobOutcome::Failed { description })
}
}
}
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 {
// The job was already completed by another worker (e.g. this worker was
// declared a zombie and the job restarted+finished elsewhere, then this
// worker caught up). The job genuinely succeeded; routing this through
// `handle_job_error` would propagate a spurious "AlreadyCompleted"
// failure up the parent flow. Drop it instead, mirroring the
// `JobOutcome::AlreadyCompleted` guard on the execution path.
Err(err @ Error::AlreadyCompleted(_)) => {
tracing::info!(
job_id = %job.id,
"job already completed by another worker, skipping result processing: {err:#}"
);
None
}
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, but preserve a
// resolved pipeline `partition` (injected before the body ran
// by resolve_partition_for_job). Run identity is immutable —
// the preprocessor must not change or drop it, or the asset
// cascade would read no partition for this producer.
windmill_common::partition::merge_args_preserving_partition(
db,
job.id,
preprocessed_args,
)
.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?;
// Asset-trigger fan-out: best-effort, never propagates errors.
// Internal eligibility checks gate to top-level Script/Preview runs;
// see windmill_queue::asset_dispatch.
asset_dispatch::dispatch_asset_triggers(db, &job).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 {
// The result already carries our injected
// `error: { name: "ManualFailure", ... }` marker when process_jc
// retagged a successful run as a failure — store it as-is to preserve
// sibling fields like `windmill_status_code`. We check for the
// injected marker specifically (not just the presence of a
// `wm_failure` field) so a real runtime failure whose raw
// result happens to contain a `wm_failure` field still goes
// through the standard `WrappedError { error: ... }` wrap path.
let downstream_result: Arc<Box<RawValue>> = if is_pre_shaped_wm_failure_result(result.get())
{
windmill_queue::add_completed_job_pre_shaped_failure(
db,
&job,
mem_peak.to_owned(),
canceled_by.clone(),
Json(&*result),
worker_name,
false,
None,
)
.await?;
result.clone()
} else {
let wrapped = 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?;
Arc::new(serde_json::value::to_raw_value(&wrapped).unwrap())
};
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,
downstream_result,
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 {
if let Ok(Some(_)) = handle_wac_child_completion(
db,
&job.id,
parent_job,
&job.workspace_id,
downstream_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.clone(),
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),
});
}