From 3e3a41d418d4ee3fe060bd3acf3324f311d89c8a Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Tue, 8 Sep 2026 10:19:26 +0000 Subject: [PATCH] feat: report a WAC task failure the workflow body never awaited (#11017) * feat: warn when a WAC task fails and the body never awaited it Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01Q6FHYvk7z4eGvhZXKZB9JF * fix: report unawaited WAC failures on the failing round and in stream order Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01Q6FHYvk7z4eGvhZXKZB9JF * docs: state the WAC warn-placement invariant where the wrapper enforces it Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01Q6FHYvk7z4eGvhZXKZB9JF --------- Co-authored-by: Claude Opus 5 (1M context) --- backend/windmill-worker/src/bun_executor.rs | 10 +- typescript-client/client.ts | 24 ++++- .../tests/workflow_unawaited_failure.test.ts | 101 ++++++++++++++++++ 3 files changed, 132 insertions(+), 3 deletions(-) create mode 100644 typescript-client/tests/workflow_unawaited_failure.test.ts diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 7c112f5985..a9844e5cfb 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -1899,8 +1899,12 @@ pub async fn handle_bun_job( // Kept comment-free — this string is written out per job. // `_takePendingStepFailure` / `_takePendingSuspend` hand back what the body - // caught and swallowed; honour them instead of reporting a `complete` (see - // `_pendingStepFailure` in client.ts). Optional: npm clients may predate them. + // caught and swallowed; honour them instead of reporting a bare `complete` + // (see client.ts). Optional: npm clients may predate them. + // `_warnUnobservedTaskFailures` reports what the body never awaited, and so + // belongs only on the paths that end the round for good. A round that + // dispatches, sleeps, checkpoints or waits for approval replays later and + // re-registers the same failures from the checkpoint — keep those quiet. let wrapper_content = if is_wac_v2 { format!( r#" @@ -1958,6 +1962,7 @@ async function run() {{ if (trailing.length > 0) {{ return {{ type: "dispatch", mode: trailing.length > 1 ? "parallel" : "sequential", steps: trailing }}; }} + ctx._warnUnobservedTaskFailures?.(); return {{ type: "complete", result: result ?? null }}; }} catch (e) {{ setWorkflowCtx(null); @@ -1977,6 +1982,7 @@ async function run() {{ }} return {{ type: "dispatch", mode: dispatch.mode ?? "sequential", steps: dispatch.steps ?? [] }}; }} + ctx._warnUnobservedTaskFailures?.(); const failed = ctx._takePendingStepFailure?.(); if (failed) {{ throw failed.error; diff --git a/typescript-client/client.ts b/typescript-client/client.ts index 480123f3e4..23ca08786c 100644 --- a/typescript-client/client.ts +++ b/typescript-client/client.ts @@ -1804,6 +1804,11 @@ export class WorkflowCtx { * into a `complete` — the parent would then record the caught branch's value as * a successful step. Boxed: the thrown value may be any falsy value. */ private _pendingStepFailure: { error: unknown } | null = null; + /** Failed tasks whose rejection nothing has consumed, by step key. An unawaited + * task is still dispatched and still fails, but nothing drives the rejecting + * thenable it returned. The first `.then()` on that thenable drops the entry, + * so what remains is only what the body never looked at. */ + private _unobservedTaskFailures = new Map(); /** When set, the task matching this key executes its inner function directly */ _executingKey: string | null; /** Serializes fast-path POSTs across concurrent step() calls within one @@ -1876,7 +1881,8 @@ export class WorkflowCtx { continue; } const err = taskErrorFromMarker(value, `Task '${name}' failed`); - return { then: (_resolve: any, reject?: any) => { if (reject) reject(err); else throw err; } } as PromiseLike; + this._unobservedTaskFailures.set(baseKey, err); + return { then: (_resolve: any, reject?: any) => { this._unobservedTaskFailures.delete(baseKey); if (reject) reject(err); else throw err; } } as PromiseLike; } return { then: (resolve: any) => resolve(value) }; } @@ -2220,6 +2226,22 @@ export class WorkflowCtx { this._pendingStepFailure = null; return f; } + + /** Report the task failures the body never looked at, and forget them. Which + * rounds may call this is the runner's constraint, stated where it is enforced. */ + _warnUnobservedTaskFailures(): void { + // A child round replays the body just to reach one step, so the failures it + // re-registers from the checkpoint are the parent round's to report. + if (this._executingKey !== null) return; + for (const [key, err] of this._unobservedTaskFailures) { + // stdout, like every other `--- WAC:` marker: the two streams are merged + // without preserving order, so a warning on stderr floats away from them. + console.log( + `\n--- WAC: task '${key}' failed but was never awaited, so the workflow result does not reflect it: ${err.message} ---`, + ); + } + this._unobservedTaskFailures.clear(); + } } export async function sleep(seconds: number): Promise { diff --git a/typescript-client/tests/workflow_unawaited_failure.test.ts b/typescript-client/tests/workflow_unawaited_failure.test.ts new file mode 100644 index 0000000000..a7c0e25f03 --- /dev/null +++ b/typescript-client/tests/workflow_unawaited_failure.test.ts @@ -0,0 +1,101 @@ +/** + * A task the body never awaits still runs and can still fail, and the workflow + * result cannot express that. Against the real client, not the inline mirror. + * + * Run with: bun test typescript-client/tests/workflow_unawaited_failure.test.ts + */ +import { expect, test, describe, mock, beforeEach, afterEach } from "bun:test"; + +mock.module("../services.gen", () => ({ + ResourceService: {}, + VariableService: {}, + JobService: {}, + HelpersService: {}, + AppService: {}, + MetricsService: {}, + OidcService: {}, + UserService: {}, + KafkaTriggerService: {}, +})); +mock.module("../core/OpenAPI", () => ({ + OpenAPI: { BASE: "http://localhost:8000/api", TOKEN: "tok" }, +})); + +const { WorkflowCtx, task, setWorkflowCtx } = await import("../client.ts"); + +const notify = task(async function notify() { + throw new Error("boom"); +}); + +const retried = task(async function retried() { + throw new Error("boom"); +}, { retry: { attempts: 1 } }); + +const marker = { __wmill_error: true, message: "boom", error: { name: "Error", message: "boom" } }; + +/** The checkpoint a replay reads after the dispatched task failed. */ +const failed = { completed_steps: { notify: marker } }; + +let reported: string[]; +const realLog = console.log; +beforeEach(() => { + reported = []; + console.log = (m: any) => reported.push(String(m)); +}); +afterEach(() => { + console.log = realLog; + setWorkflowCtx(null); +}); + +describe("unawaited task failure", () => { + test("is reported when the body never looked at it", async () => { + const ctx = new WorkflowCtx(failed); + setWorkflowCtx(ctx); + + notify(); + + ctx._warnUnobservedTaskFailures(); + expect(reported).toHaveLength(1); + expect(reported[0]).toContain("task 'notify' failed but was never awaited"); + expect(reported[0]).toContain("boom"); + }); + + // The body is free to hold the handle and await it further down, so the + // failure has to be judged at the end of the round rather than at the call. + test("is not reported when the body awaits it later", async () => { + const ctx = new WorkflowCtx(failed); + setWorkflowCtx(ctx); + + const handle = notify(); + await expect(Promise.resolve(handle)).rejects.toThrow("boom"); + + ctx._warnUnobservedTaskFailures(); + expect(reported).toEqual([]); + }); + + // Only the attempt handed back to the body counts: the ones a retry moved + // past are not failures the workflow was ever in a position to see. + test("is not reported when a retry recovered from it", async () => { + const ctx = new WorkflowCtx({ + completed_steps: { retried: marker, "retried#retry2": null, "retried#2": 1 }, + }); + setWorkflowCtx(ctx); + + retried(); + + ctx._warnUnobservedTaskFailures(); + expect(reported).toEqual([]); + }); + + // A child round replays the whole body to reach one step, so it re-registers + // every checkpointed failure; reporting them here duplicates them per child. + test("is not reported by a child round", async () => { + const ctx = new WorkflowCtx({ ...failed, _executing_key: "other" }); + setWorkflowCtx(ctx); + + notify(); + + ctx._warnUnobservedTaskFailures(); + expect(reported).toEqual([]); + }); +});