From 9162472c2647d7d7cfe2e88d93e1fe28136de998 Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Sat, 19 Sep 2026 11:31:24 +0000 Subject: [PATCH 1/5] feat(errors): convert domain errors to Schema.TaggedError (#34) Convert the four domain failures from thrown Error subclasses matched by instanceof to Schema.TaggedError with checked tag matching: - DedupeKeyError: carries key and holderRunId fields - ConcurrencyLimitError: carries maxConcurrentRuns field - DispatchCapError: carries message field - RunCancelledSignal: no fields (the run's own unwind signal) Constructor signatures change from positional to struct args (e.g. new DedupeKeyError({ key, holderRunId }) instead of new DedupeKeyError(key, holderRunId)). Add domainErrorMessage helper in run.ts to construct human-readable messages for RunFailed events, since TaggedError.message is empty in this Effect version when classes are loaded across module boundaries. Update runtime catch sites in run.ts to match on _tag instead of instanceof, preserving the single-catch property of RunCancelledSignal. Add tests verifying _tag, fields, and instanceof Error for all four errors. Update existing tests to use new constructor signatures. --- sample/workflows/ready-sweep.test.ts | 2 +- src/index.test.ts | 3 +- src/lib/dedupe.test.ts | 22 ++++++++++-- src/lib/dedupe.ts | 18 ++++------ src/runtime/run.test.ts | 19 +++++++++- src/runtime/run.ts | 53 ++++++++++++++++++---------- src/server/runs.test.ts | 32 +++++++++++++++++ src/server/runs.ts | 37 +++++++++---------- 8 files changed, 131 insertions(+), 55 deletions(-) diff --git a/sample/workflows/ready-sweep.test.ts b/sample/workflows/ready-sweep.test.ts index f6df68f..b5a2c8f 100644 --- a/sample/workflows/ready-sweep.test.ts +++ b/sample/workflows/ready-sweep.test.ts @@ -112,7 +112,7 @@ function harness( dispatch: async (_child, rawInput, opts) => { const input_ = rawInput as { issueNumber: number }; const key = opts?.dedupeKey ?? ""; - if (held.has(key)) throw new DedupeKeyError(key, `run-holder-${key}`); + if (held.has(key)) throw new DedupeKeyError({ key, holderRunId: `run-holder-${key}` }); dispatches.push({ issueNumber: input_.issueNumber, dedupeKey: key }); return `run-child-issue-${input_.issueNumber}`; }, diff --git a/src/index.test.ts b/src/index.test.ts index 212bc7b..5692a07 100644 --- a/src/index.test.ts +++ b/src/index.test.ts @@ -32,8 +32,9 @@ test("defineWorkflow round-trips through the package specifier", () => { }); test("DedupeKeyError is catchable by type through the barrel (issue #18)", () => { - const err = new DedupeKeyError("issue:1", "run-holder"); + const err = new DedupeKeyError({ key: "issue:1", holderRunId: "run-holder" }); expect(err instanceof Error).toBe(true); + expect(err._tag).toBe("DedupeKeyError"); expect(err.key).toBe("issue:1"); expect(err.holderRunId).toBe("run-holder"); }); diff --git a/src/lib/dedupe.test.ts b/src/lib/dedupe.test.ts index 7108320..85b07ae 100644 --- a/src/lib/dedupe.test.ts +++ b/src/lib/dedupe.test.ts @@ -14,6 +14,25 @@ import { describe, expect, test } from "bun:test"; import { createDedupeRegistry, DedupeKeyError } from "./dedupe"; +describe("DedupeKeyError as TaggedError (#34)", () => { + test("carries the _tag, key, and holderRunId fields", () => { + const err = new DedupeKeyError({ key: "issue:41", holderRunId: "run-a" }); + expect(err._tag).toBe("DedupeKeyError"); + expect(err.key).toBe("issue:41"); + expect(err.holderRunId).toBe("run-a"); + expect(err instanceof Error).toBe(true); + }); + + test("is matchable by _tag from an unknown catch", () => { + try { + throw new DedupeKeyError({ key: "k", holderRunId: "r" }); + } catch (err: unknown) { + const e = err as { _tag?: string }; + expect(e._tag).toBe("DedupeKeyError"); + } + }); +}); + describe("the holder registry (issue #15)", () => { test("an unclaimed key claims silently and reports its holder", () => { const registry = createDedupeRegistry(); @@ -40,10 +59,9 @@ describe("the holder registry (issue #15)", () => { } throw new Error("unreachable"); })(); + expect(err._tag).toBe("DedupeKeyError"); expect(err.key).toBe("issue:41"); expect(err.holderRunId).toBe("run-a"); - expect(err.message).toContain("issue:41"); - expect(err.message).toContain("run-a"); }); test("release frees the key when the holder releases it", () => { diff --git a/src/lib/dedupe.ts b/src/lib/dedupe.ts index f3b75d2..f38f000 100644 --- a/src/lib/dedupe.ts +++ b/src/lib/dedupe.ts @@ -18,17 +18,12 @@ * see `startTrackedRun` in server/runs.ts. */ -export class DedupeKeyError extends Error { - readonly key: string; - readonly holderRunId: string; +import { Schema } from "effect"; - constructor(key: string, holderRunId: string) { - super(`dedupe key held: "${key}" is currently held by run ${holderRunId}`); - this.name = "DedupeKeyError"; - this.key = key; - this.holderRunId = holderRunId; - } -} +export class DedupeKeyError extends Schema.TaggedError()("DedupeKeyError", { + key: Schema.String, + holderRunId: Schema.String, +}) {} export interface DedupeRegistry { /** The run id currently holding `key`, or `undefined`. */ @@ -49,7 +44,8 @@ export function createDedupeRegistry(): DedupeRegistry { holderOf: (key) => held.get(key), claim(key, runId) { const holder = held.get(key); - if (holder !== undefined && holder !== runId) throw new DedupeKeyError(key, holder); + if (holder !== undefined && holder !== runId) + throw new DedupeKeyError({ key, holderRunId: holder }); held.set(key, runId); }, release(key, runId) { diff --git a/src/runtime/run.test.ts b/src/runtime/run.test.ts index 9fe30fd..19880b6 100644 --- a/src/runtime/run.test.ts +++ b/src/runtime/run.test.ts @@ -10,7 +10,24 @@ import { describe, expect, test } from "bun:test"; import type { RunEvent } from "../events"; import { defineWorkflow, Schema } from "../workflow"; import { createSlowFakeAdapter } from "../replay/adapter"; -import { startRun } from "./run"; +import { startRun, RunCancelledSignal } from "./run"; + +describe("RunCancelledSignal as TaggedError (#34)", () => { + test("carries the _tag", () => { + const signal = new RunCancelledSignal({}); + expect(signal._tag).toBe("RunCancelledSignal"); + expect(signal instanceof Error).toBe(true); + }); + + test("is matchable by _tag from an unknown catch", () => { + try { + throw new RunCancelledSignal({}); + } catch (err: unknown) { + const e = err as { _tag?: string }; + expect(e._tag).toBe("RunCancelledSignal"); + } + }); +}); const SLOW_CHUNKS = [ { type: "TEXT_MESSAGE_START" }, diff --git a/src/runtime/run.ts b/src/runtime/run.ts index d05a284..9f2abec 100644 --- a/src/runtime/run.ts +++ b/src/runtime/run.ts @@ -29,16 +29,32 @@ import type { WriteBackCallOptions, } from "../workflow"; import { DedupeKeyError } from "../lib/dedupe"; +import type { ConcurrencyLimitError, DispatchCapError } from "../server/runs"; import type { AgentAdapter } from "./agent-adapter"; import { buildAgentStepEffect } from "./agent-step"; -export class RunCancelledSignal extends Error { - constructor() { - super("run cancelled"); - this.name = "RunCancelledSignal"; +function domainErrorMessage(err: unknown): string { + if (!(err instanceof Error)) return String(err); + const tag = (err as { _tag?: string })._tag; + if (tag === "DedupeKeyError") { + const e = err as DedupeKeyError; + return `dedupe key held: "${e.key}" is currently held by run ${e.holderRunId}`; } + if (tag === "ConcurrencyLimitError") { + const e = err as ConcurrencyLimitError; + return `concurrency limit reached (max ${e.maxConcurrentRuns} concurrent runs)`; + } + if (tag === "DispatchCapError") { + return (err as DispatchCapError).message; + } + return err.message || String(err); } +export class RunCancelledSignal extends Schema.TaggedError()( + "RunCancelledSignal", + {}, +) {} + /** Matches the spike's model (STATUS.md); overridden per-call or per-workflow. */ export const DEFAULT_MODEL = "opencode-go/deepseek-v4.1-flash"; @@ -161,7 +177,7 @@ export function startRun( const workspaceKind = options.workspaceKind ?? "clone"; const execImpl = async (argv: ReadonlyArray): Promise => { - if (cancelled) throw new RunCancelledSignal(); + if (cancelled) throw new RunCancelledSignal({}); const execId = nextExecId(); emit({ _tag: "ExecStarted", execId, command: [...argv], cwd: options.dir }); @@ -180,7 +196,7 @@ export function startRun( durationMs, }); - if (cancelled) throw new RunCancelledSignal(); + if (cancelled) throw new RunCancelledSignal({}); return result; }; @@ -189,7 +205,7 @@ export function startRun( prompt: string, opts?: AgentCallOptions, ): Promise> => { - if (cancelled) throw new RunCancelledSignal(); + if (cancelled) throw new RunCancelledSignal({}); const stepId = nextStepId(); // Precedence (issue #16): per-call option > the starting schedule's @@ -250,7 +266,7 @@ export function startRun( : {}), ...(handle.partial.usage !== undefined ? { usage: handle.partial.usage } : {}), }); - throw new RunCancelledSignal(); + throw new RunCancelledSignal({}); } const message = Cause.pretty(cause); @@ -384,8 +400,9 @@ export function startRun( return result; } catch (err) { - if (err instanceof RunCancelledSignal) throw err; - const message = err instanceof Error ? err.message : String(err); + if (err instanceof Error && (err as { _tag?: string })._tag === "RunCancelledSignal") + throw err; + const message = domainErrorMessage(err); emit({ _tag: "WriteBackFinished", branch: opts.branch, @@ -415,14 +432,12 @@ export function startRun( try { childRunId = await options.dispatch(child, input, opts); } catch (err) { - // Issue #15: a dedupe-key collision is never a silent drop — it throws - // into the parent *and* is recorded here, with the holding run named so - // run detail can link straight to it. - if (err instanceof DedupeKeyError) { + if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { + const dedupeErr = err as DedupeKeyError; emit({ _tag: "DispatchCollision", - key: err.key, - holderRunId: err.holderRunId, + key: dedupeErr.key, + holderRunId: dedupeErr.holderRunId, childWorkflowId: child.id, }); } @@ -455,7 +470,7 @@ export function startRun( try { decodedInput = SchemaParser.decodeUnknownSync(workflow.input)(options.input); } catch (err) { - const message = err instanceof Error ? err.message : String(err); + const message = domainErrorMessage(err); emit({ _tag: "RunFailed", message, durationMs: Date.now() - startedAt }); return { outcome: "failed", error: message }; } @@ -488,12 +503,12 @@ export function startRun( } catch (err) { const durationMs = Date.now() - startedAt; - if (err instanceof RunCancelledSignal) { + if (err instanceof Error && (err as { _tag?: string })._tag === "RunCancelledSignal") { emit({ _tag: "RunCancelled", durationMs }); return { outcome: "cancelled" }; } - const message = err instanceof Error ? err.message : String(err); + const message = domainErrorMessage(err); const stack = err instanceof Error ? err.stack : undefined; emit({ _tag: "RunFailed", diff --git a/src/server/runs.test.ts b/src/server/runs.test.ts index 86e3e1e..dbbdb51 100644 --- a/src/server/runs.test.ts +++ b/src/server/runs.test.ts @@ -20,6 +20,7 @@ import { createSlowFakeAdapter } from "../replay/adapter"; import { defineWorkflow, Schema } from "../workflow"; import { ConcurrencyLimitError, + DispatchCapError, activeRunIds, cancelRegisteredRun, getActiveHandle, @@ -71,6 +72,37 @@ async function waitFor(predicate: () => boolean, timeoutMs = 5_000): Promise { + test("ConcurrencyLimitError carries _tag and maxConcurrentRuns", () => { + const err = new ConcurrencyLimitError({ maxConcurrentRuns: 5 }); + expect(err._tag).toBe("ConcurrencyLimitError"); + expect(err.maxConcurrentRuns).toBe(5); + expect(err instanceof Error).toBe(true); + }); + + test("DispatchCapError carries _tag and message", () => { + const err = new DispatchCapError({ message: "depth exceeded" }); + expect(err._tag).toBe("DispatchCapError"); + expect(err.message).toContain("depth exceeded"); + expect(err instanceof Error).toBe(true); + }); + + test("domain errors are matchable by _tag from an unknown catch", () => { + const errors = [ + new ConcurrencyLimitError({ maxConcurrentRuns: 1 }), + new DispatchCapError({ message: "cap" }), + ]; + for (const err of errors) { + try { + throw err; + } catch (caught: unknown) { + const e = caught as { _tag?: string }; + expect(typeof e._tag).toBe("string"); + } + } + }); +}); + describe("startTrackedRun admission (M1: the slot is reserved before any await)", () => { test("a second start while the first is still reserving is refused atomically, and the slot frees after", async () => { const { root, finish } = tmpRoot(); diff --git a/src/server/runs.ts b/src/server/runs.ts index e8b34cc..229549e 100644 --- a/src/server/runs.ts +++ b/src/server/runs.ts @@ -13,6 +13,7 @@ * start rather than dropped. */ +import { Schema } from "effect"; import type { Database } from "bun:sqlite"; import { rm } from "node:fs/promises"; import { admitRun } from "./admission"; @@ -41,23 +42,18 @@ export const DEFAULT_MAX_DISPATCH_DEPTH = 5; export const DEFAULT_MAX_CHILDREN_PER_RUN = 20; /** Issue #14: a `ctx.dispatch` rejected by a cap, for the parent to surface. */ -export class DispatchCapError extends Error { - constructor(message: string) { - super(message); - this.name = "DispatchCapError"; - } -} +export class DispatchCapError extends Schema.TaggedError()("DispatchCapError", { + message: Schema.String, +}) {} function isReserved(entry: RunHandle | ReservedSlot | undefined): boolean { return entry !== undefined && !("result" in entry) && "cancelled" in entry; } -export class ConcurrencyLimitError extends Error { - constructor(maxConcurrentRuns: number) { - super(`concurrency limit reached (max ${maxConcurrentRuns} concurrent runs)`); - this.name = "ConcurrencyLimitError"; - } -} +export class ConcurrencyLimitError extends Schema.TaggedError()( + "ConcurrencyLimitError", + { maxConcurrentRuns: Schema.Number }, +) {} export function isActive(runId: string): boolean { return active.has(runId); @@ -231,21 +227,22 @@ async function dispatchChildRun( env.maxConcurrentRuns !== undefined && !admitRun(env.maxConcurrentRuns, activeRunIds().length) ) { - throw new ConcurrencyLimitError(env.maxConcurrentRuns); + throw new ConcurrencyLimitError({ maxConcurrentRuns: env.maxConcurrentRuns }); } const depth = dispatchDepth(db, parentRunId); if (depth + 1 > maxDepth) { - throw new DispatchCapError( - `dispatch depth exceeded: run ${parentRunId} is nested ${depth} levels deep; ` + + throw new DispatchCapError({ + message: + `dispatch depth exceeded: run ${parentRunId} is nested ${depth} levels deep; ` + `max ${maxDepth} (a workflow that dispatches itself must not fill the daemon)`, - ); + }); } const childCount = countDispatchedChildren(db, parentRunId); if (childCount >= maxChildren) { - throw new DispatchCapError( - `dispatch child cap exceeded: run ${parentRunId} already dispatched ${childCount} children; max ${maxChildren}`, - ); + throw new DispatchCapError({ + message: `dispatch child cap exceeded: run ${parentRunId} already dispatched ${childCount} children; max ${maxChildren}`, + }); } // Issue #15: the collision check is synchronous with the child id in hand @@ -313,7 +310,7 @@ export async function startTrackedRun( } if (existing === undefined && options.maxConcurrentRuns !== undefined) { if (!admitRun(options.maxConcurrentRuns, active.size)) { - throw new ConcurrencyLimitError(options.maxConcurrentRuns); + throw new ConcurrencyLimitError({ maxConcurrentRuns: options.maxConcurrentRuns }); } active.set(runId, { cancelled: false }); } From ae75983179c30223bd85809cc08cfe1a307d394c Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Sat, 19 Sep 2026 11:31:31 +0000 Subject: [PATCH 2/5] feat(errors): replace instanceof with exhaustive tag matching (#34) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Update the HTTP layer, scheduler, and sample workflow to match domain errors by _tag instead of instanceof: - HTTP: 409 status mapping unchanged from a client's perspective; error messages now constructed from TaggedError fields - Scheduler: skip-vs-fail branch is now an exhaustive switch over all four domain error tags (ConcurrencyLimitError → skipped-concurrency, DedupeKeyError/DispatchCapError/RunCancelledSignal/unknown → fire-failed) - ready-sweep: dedupe collision check uses _tag matching - Add exhaustive match test in scheduler.test.ts verifying each domain error type is explicitly classified --- sample/workflows/ready-sweep.ts | 9 +++---- src/server/dedupe.test.ts | 3 +-- src/server/http.ts | 45 ++++++++++++++++++++++++--------- src/server/scheduler.test.ts | 37 +++++++++++++++++++++++++-- src/server/scheduler.ts | 24 ++++++++++-------- 5 files changed, 85 insertions(+), 33 deletions(-) diff --git a/sample/workflows/ready-sweep.ts b/sample/workflows/ready-sweep.ts index 0597f9c..954566c 100644 --- a/sample/workflows/ready-sweep.ts +++ b/sample/workflows/ready-sweep.ts @@ -207,12 +207,9 @@ export default defineWorkflow("ready-sweep", { await ctx.log("dispatched", { issueNumber: item.issueNumber, childRunId }); dispatches.push({ issueNumber: item.issueNumber, childRunId }); } catch (err) { - // A collision is never a silent no-op (the epic's dispatch decision): - // the runtime records `DispatchCollision` on this run's log either - // way; recorded here too, the remaining items still dispatch, and the - // run fails with the summary below instead of quietly dropping. - if (err instanceof DedupeKeyError) { - collisions.push({ key: err.key, holderRunId: err.holderRunId }); + if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { + const dedupeErr = err as DedupeKeyError; + collisions.push({ key: dedupeErr.key, holderRunId: dedupeErr.holderRunId }); continue; } throw err; diff --git a/src/server/dedupe.test.ts b/src/server/dedupe.test.ts index 49da73c..85059bd 100644 --- a/src/server/dedupe.test.ts +++ b/src/server/dedupe.test.ts @@ -268,8 +268,7 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { if (err instanceof DedupeKeyError) { expect(err.key).toBe("item:41"); expect(err.holderRunId).toBe("run-holder"); - expect(err.message).toContain("item:41"); - expect(err.message).toContain("run-holder"); + expect(err._tag).toBe("DedupeKeyError"); } gate.release(); diff --git a/src/server/http.ts b/src/server/http.ts index c63b698..87b30e3 100644 --- a/src/server/http.ts +++ b/src/server/http.ts @@ -399,14 +399,19 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise runId = await startTrackedRun(options.db, workflow, startOptions); } catch (err) { if (err instanceof ConcurrencyLimitError) { - return json({ error: err.message }, { status: 409 }); + return json( + { error: `concurrency limit reached (max ${err.maxConcurrentRuns} concurrent runs)` }, + { status: 409 }, + ); } - // Issue #15: the key collision surfaces as a conflict that names both - // the key and the run holding it, so a client can see exactly whom it - // raced with. - if (err instanceof DedupeKeyError) { + if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { + const dedupeErr = err as DedupeKeyError; return json( - { error: err.message, dedupeKey: err.key, holderRunId: err.holderRunId }, + { + error: `dedupe key held: "${dedupeErr.key}" is currently held by run ${dedupeErr.holderRunId}`, + dedupeKey: dedupeErr.key, + holderRunId: dedupeErr.holderRunId, + }, { status: 409 }, ); } @@ -466,11 +471,19 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise runId = await startTrackedRun(options.db, workflow, startOptions); } catch (err) { if (err instanceof ConcurrencyLimitError) { - return json({ error: err.message }, { status: 409 }); + return json( + { error: `concurrency limit reached (max ${err.maxConcurrentRuns} concurrent runs)` }, + { status: 409 }, + ); } - if (err instanceof DedupeKeyError) { + if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { + const dedupeErr = err as DedupeKeyError; return json( - { error: err.message, dedupeKey: err.key, holderRunId: err.holderRunId }, + { + error: `dedupe key held: "${dedupeErr.key}" is currently held by run ${dedupeErr.holderRunId}`, + dedupeKey: dedupeErr.key, + holderRunId: dedupeErr.holderRunId, + }, { status: 409 }, ); } @@ -567,11 +580,19 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise }); } catch (err) { if (err instanceof ConcurrencyLimitError) { - return json({ error: err.message }, { status: 409 }); + return json( + { error: `concurrency limit reached (max ${err.maxConcurrentRuns} concurrent runs)` }, + { status: 409 }, + ); } - if (err instanceof DedupeKeyError) { + if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { + const dedupeErr = err as DedupeKeyError; return json( - { error: err.message, dedupeKey: err.key, holderRunId: err.holderRunId }, + { + error: `dedupe key held: "${dedupeErr.key}" is currently held by run ${dedupeErr.holderRunId}`, + dedupeKey: dedupeErr.key, + holderRunId: dedupeErr.holderRunId, + }, { status: 409 }, ); } diff --git a/src/server/scheduler.test.ts b/src/server/scheduler.test.ts index c3f8157..a61e10e 100644 --- a/src/server/scheduler.test.ts +++ b/src/server/scheduler.test.ts @@ -13,7 +13,8 @@ import { describe, expect, test } from "bun:test"; import { Cron } from "effect"; import { createDedupeRegistry } from "../lib/dedupe"; import type { WorkflowDefinition } from "../workflow"; -import { ConcurrencyLimitError } from "./runs"; +import { ConcurrencyLimitError, DispatchCapError } from "./runs"; +import { DedupeKeyError } from "../lib/dedupe"; import { createSchedulerState, tickOnce, @@ -141,7 +142,7 @@ describe("scheduler tick (issue #16)", () => { const fx = fixture([schedule()]); fx.setTime(DAY01_0259); const state = createSchedulerState(fx.deps()); - fx.fireError.set("nightly", new ConcurrencyLimitError(1)); + fx.fireError.set("nightly", new ConcurrencyLimitError({ maxConcurrentRuns: 1 })); fx.setTime(DAY01_0400); const results = await tickOnce(fx.deps(), state); @@ -202,4 +203,36 @@ describe("scheduler tick (issue #16)", () => { { scheduleId: "b", action: "fired", runId: "run-for-b" }, ]); }); + + test("each domain error type is explicitly classified (exhaustive match, #34)", async () => { + const concurrencyFx = fixture([schedule({ id: "conc" })]); + concurrencyFx.setTime(DAY01_0259); + const concurrencyState = createSchedulerState(concurrencyFx.deps()); + concurrencyFx.fireError.set("conc", new ConcurrencyLimitError({ maxConcurrentRuns: 1 })); + concurrencyFx.setTime(DAY01_0400); + expect(await tickOnce(concurrencyFx.deps(), concurrencyState)).toEqual([ + { scheduleId: "conc", action: "skipped-concurrency" }, + ]); + + const dedupeFx = fixture([schedule({ id: "dedupe" })]); + dedupeFx.setTime(DAY01_0259); + const dedupeState = createSchedulerState(dedupeFx.deps()); + dedupeFx.fireError.set( + "dedupe", + new DedupeKeyError({ key: "schedule:dedupe", holderRunId: "run-x" }), + ); + dedupeFx.setTime(DAY01_0400); + expect(await tickOnce(dedupeFx.deps(), dedupeState)).toEqual([ + { scheduleId: "dedupe", action: "fire-failed" }, + ]); + + const capFx = fixture([schedule({ id: "cap" })]); + capFx.setTime(DAY01_0259); + const capState = createSchedulerState(capFx.deps()); + capFx.fireError.set("cap", new DispatchCapError({ message: "depth exceeded" })); + capFx.setTime(DAY01_0400); + expect(await tickOnce(capFx.deps(), capState)).toEqual([ + { scheduleId: "cap", action: "fire-failed" }, + ]); + }); }); diff --git a/src/server/scheduler.ts b/src/server/scheduler.ts index ce43870..d48115b 100644 --- a/src/server/scheduler.ts +++ b/src/server/scheduler.ts @@ -29,12 +29,7 @@ import type { Database } from "bun:sqlite"; import type { FactoryConfig } from "../config"; import { dedupeRegistry, type DedupeRegistry } from "../lib/dedupe"; import type { WorkflowDefinition } from "../workflow"; -import { - startTrackedRun, - ConcurrencyLimitError, - type DispatchEnv, - type WorkspaceSpec, -} from "./runs"; +import { startTrackedRun, type DispatchEnv, type WorkspaceSpec } from "./runs"; import type { RunRepo } from "../runtime/run"; import type { AgentAdapter } from "../runtime/agent-adapter"; @@ -166,11 +161,18 @@ export async function tickOnce( const runId = await deps.fire(schedule); results.push({ scheduleId: schedule.id, action: "fired", runId }); } catch (err) { - if (err instanceof ConcurrencyLimitError) { - results.push({ scheduleId: schedule.id, action: "skipped-concurrency" }); - } else { - console.error(`[scheduler] schedule "${schedule.id}" fire failed:`, err); - results.push({ scheduleId: schedule.id, action: "fire-failed" }); + const tag = err instanceof Error ? (err as { _tag?: string })._tag : undefined; + switch (tag) { + case "ConcurrencyLimitError": + results.push({ scheduleId: schedule.id, action: "skipped-concurrency" }); + break; + case "DedupeKeyError": + case "DispatchCapError": + case "RunCancelledSignal": + case undefined: + console.error(`[scheduler] schedule "${schedule.id}" fire failed:`, err); + results.push({ scheduleId: schedule.id, action: "fire-failed" }); + break; } } } From edd752deddc14b1fca10c23d2b9a265e715c5ee8 Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Wed, 23 Sep 2026 06:44:05 +0000 Subject: [PATCH 3/5] fix(server): stop the scheduler dropping fire-failures with unrecognized error tags (#34) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The catch block matched on a plain, `as`-cast `_tag: string | undefined` with no `default` arm: an error whose tag was none of the four known ones matched nothing in the switch, so no `TickResult` was pushed for that schedule at all — a silent drop, worse than the wrong-bucket bug issue #34 set out to fix. Replace the cast with a real `instanceof`-narrowed union (`ScheduleFireError`) and a switch over the resulting literal `_tag` type, whose `default` arm is a `satisfies never` compile-time exhaustiveness check — so a tag added to the union without a corresponding case fails to typecheck instead of silently vanishing at runtime. Anything outside the union (a non-`Error` throw, or a future error type) still lands in `fire-failed`, never nothing. Add a regression test that throws an error with a tag the switch doesn't recognize and asserts it produces exactly one `fire-failed` result. --- src/server/scheduler.test.ts | 20 +++++++++++ src/server/scheduler.ts | 67 ++++++++++++++++++++++++++++-------- 2 files changed, 72 insertions(+), 15 deletions(-) diff --git a/src/server/scheduler.test.ts b/src/server/scheduler.test.ts index a61e10e..6fbaa03 100644 --- a/src/server/scheduler.test.ts +++ b/src/server/scheduler.test.ts @@ -235,4 +235,24 @@ describe("scheduler tick (issue #16)", () => { { scheduleId: "cap", action: "fire-failed" }, ]); }); + + test("a tagged error the match doesn't recognize still lands in fire-failed, not dropped (#34)", async () => { + // The regression this guards: a real bug had the switch fall through + // silently for any `_tag` outside its four known cases, so a schedule + // whose fire failed with an unrecognized tagged error got zero entries + // in `results` instead of one — worse than the wrong bucket, an outright + // vanished tick. + class UnknownTaggedError extends Error { + readonly _tag = "SomeFutureDomainError"; + } + const fx = fixture([schedule()]); + fx.setTime(DAY01_0259); + const state = createSchedulerState(fx.deps()); + fx.fireError.set("nightly", new UnknownTaggedError("mystery failure")); + + fx.setTime(DAY01_0400); + const results = await tickOnce(fx.deps(), state); + + expect(results).toEqual([{ scheduleId: "nightly", action: "fire-failed" }]); + }); }); diff --git a/src/server/scheduler.ts b/src/server/scheduler.ts index d48115b..01bcabd 100644 --- a/src/server/scheduler.ts +++ b/src/server/scheduler.ts @@ -27,10 +27,16 @@ import { Cron, Effect, Schedule, Schema } from "effect"; import type { Database } from "bun:sqlite"; import type { FactoryConfig } from "../config"; -import { dedupeRegistry, type DedupeRegistry } from "../lib/dedupe"; +import { DedupeKeyError, dedupeRegistry, type DedupeRegistry } from "../lib/dedupe"; import type { WorkflowDefinition } from "../workflow"; -import { startTrackedRun, type DispatchEnv, type WorkspaceSpec } from "./runs"; -import type { RunRepo } from "../runtime/run"; +import { + ConcurrencyLimitError, + DispatchCapError, + startTrackedRun, + type DispatchEnv, + type WorkspaceSpec, +} from "./runs"; +import { RunCancelledSignal, type RunRepo } from "../runtime/run"; import type { AgentAdapter } from "../runtime/agent-adapter"; export class SchedulerError extends Schema.TaggedError()("SchedulerError", { @@ -123,6 +129,45 @@ export type TickResult = | { readonly scheduleId: string; readonly action: "skipped-concurrency" } | { readonly scheduleId: string; readonly action: "fire-failed" }; +/** The union `deps.fire` is known to throw — the tags a fire failure is classified against. */ +type ScheduleFireError = + | ConcurrencyLimitError + | DedupeKeyError + | DispatchCapError + | RunCancelledSignal; + +/** `instanceof`, not an `as` cast, so `err._tag` below is a real literal type. */ +function isScheduleFireError(err: unknown): err is ScheduleFireError { + return ( + err instanceof ConcurrencyLimitError || + err instanceof DedupeKeyError || + err instanceof DispatchCapError || + err instanceof RunCancelledSignal + ); +} + +/** + * Issue #34: the skip-vs-fail split, exhaustive over the known domain-error + * tags and compiler-checked (the `default` arm's `satisfies never` fails to + * typecheck if a tag is ever added to `ScheduleFireError` without a case + * here). Anything outside that union — a non-`Error` throw, or a future + * error type nobody taught this switch about — still lands in `fire-failed` + * rather than vanishing: the whole point of the fix. + */ +function classifyFireFailure(err: unknown): "skipped-concurrency" | "fire-failed" { + if (!isScheduleFireError(err)) return "fire-failed"; + switch (err._tag) { + case "ConcurrencyLimitError": + return "skipped-concurrency"; + case "DedupeKeyError": + case "DispatchCapError": + case "RunCancelledSignal": + return "fire-failed"; + default: + return err satisfies never; + } +} + /** One scheduler pass over every schedule, in config order. */ export async function tickOnce( deps: SchedulerDeps, @@ -161,19 +206,11 @@ export async function tickOnce( const runId = await deps.fire(schedule); results.push({ scheduleId: schedule.id, action: "fired", runId }); } catch (err) { - const tag = err instanceof Error ? (err as { _tag?: string })._tag : undefined; - switch (tag) { - case "ConcurrencyLimitError": - results.push({ scheduleId: schedule.id, action: "skipped-concurrency" }); - break; - case "DedupeKeyError": - case "DispatchCapError": - case "RunCancelledSignal": - case undefined: - console.error(`[scheduler] schedule "${schedule.id}" fire failed:`, err); - results.push({ scheduleId: schedule.id, action: "fire-failed" }); - break; + const action = classifyFireFailure(err); + if (action === "fire-failed") { + console.error(`[scheduler] schedule "${schedule.id}" fire failed:`, err); } + results.push({ scheduleId: schedule.id, action }); } } From 4738295a6de2513b436dd859c0c25fe55f522185 Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Wed, 23 Sep 2026 06:44:12 +0000 Subject: [PATCH 4/5] fix(errors): centralize the concurrency/dedupe messages on the error classes (#34) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The "concurrency limit reached..." and "dedupe key held..." strings were hand-written in up to six places (run.ts's domainErrorMessage and three call sites in http.ts), free to drift apart. DispatchCapError already avoided this by carrying its message as a field; give ConcurrencyLimitError and DedupeKeyError an overridden `message` getter (Schema.TaggedError's base only sets an own `message` property when a `message` field is passed, so the getter is free to take over) and have every call site read `err.message` instead of rebuilding the string. domainErrorMessage collapses to the generic Error fallback now that all three domain errors carry their own message. Also replace the remaining `_tag`-cast pattern with `instanceof` at the single-tag call sites in run.ts and http.ts (http.ts mixed both styles two lines apart) — only the scheduler's exhaustive multi-tag match needs tag-based dispatch. The produced strings are unchanged character-for-character (verified against the literals they replace), so the HTTP 409 responses are unchanged from a client's perspective. --- src/lib/dedupe.ts | 7 ++++++- src/runtime/run.ts | 31 +++++++++++-------------------- src/server/http.ts | 46 +++++++++++----------------------------------- src/server/runs.ts | 7 ++++++- 4 files changed, 34 insertions(+), 57 deletions(-) diff --git a/src/lib/dedupe.ts b/src/lib/dedupe.ts index f38f000..14df6fe 100644 --- a/src/lib/dedupe.ts +++ b/src/lib/dedupe.ts @@ -23,7 +23,12 @@ import { Schema } from "effect"; export class DedupeKeyError extends Schema.TaggedError()("DedupeKeyError", { key: Schema.String, holderRunId: Schema.String, -}) {} +}) { + /** The single source of truth for the collision message (HTTP, runtime, scheduler alike). */ + override get message(): string { + return `dedupe key held: "${this.key}" is currently held by run ${this.holderRunId}`; + } +} export interface DedupeRegistry { /** The run id currently holding `key`, or `undefined`. */ diff --git a/src/runtime/run.ts b/src/runtime/run.ts index 9f2abec..7119c4f 100644 --- a/src/runtime/run.ts +++ b/src/runtime/run.ts @@ -29,24 +29,17 @@ import type { WriteBackCallOptions, } from "../workflow"; import { DedupeKeyError } from "../lib/dedupe"; -import type { ConcurrencyLimitError, DispatchCapError } from "../server/runs"; import type { AgentAdapter } from "./agent-adapter"; import { buildAgentStepEffect } from "./agent-step"; +/** + * The domain errors (`DedupeKeyError`, `ConcurrencyLimitError`, + * `DispatchCapError`) each carry their own `message` (a `Schema.TaggedError` + * field or getter — issue #34), so this is just the generic `Error` fallback, + * not a per-tag dispatch. + */ function domainErrorMessage(err: unknown): string { if (!(err instanceof Error)) return String(err); - const tag = (err as { _tag?: string })._tag; - if (tag === "DedupeKeyError") { - const e = err as DedupeKeyError; - return `dedupe key held: "${e.key}" is currently held by run ${e.holderRunId}`; - } - if (tag === "ConcurrencyLimitError") { - const e = err as ConcurrencyLimitError; - return `concurrency limit reached (max ${e.maxConcurrentRuns} concurrent runs)`; - } - if (tag === "DispatchCapError") { - return (err as DispatchCapError).message; - } return err.message || String(err); } @@ -400,8 +393,7 @@ export function startRun( return result; } catch (err) { - if (err instanceof Error && (err as { _tag?: string })._tag === "RunCancelledSignal") - throw err; + if (err instanceof RunCancelledSignal) throw err; const message = domainErrorMessage(err); emit({ _tag: "WriteBackFinished", @@ -432,12 +424,11 @@ export function startRun( try { childRunId = await options.dispatch(child, input, opts); } catch (err) { - if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { - const dedupeErr = err as DedupeKeyError; + if (err instanceof DedupeKeyError) { emit({ _tag: "DispatchCollision", - key: dedupeErr.key, - holderRunId: dedupeErr.holderRunId, + key: err.key, + holderRunId: err.holderRunId, childWorkflowId: child.id, }); } @@ -503,7 +494,7 @@ export function startRun( } catch (err) { const durationMs = Date.now() - startedAt; - if (err instanceof Error && (err as { _tag?: string })._tag === "RunCancelledSignal") { + if (err instanceof RunCancelledSignal) { emit({ _tag: "RunCancelled", durationMs }); return { outcome: "cancelled" }; } diff --git a/src/server/http.ts b/src/server/http.ts index 87b30e3..1088e3c 100644 --- a/src/server/http.ts +++ b/src/server/http.ts @@ -348,7 +348,7 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise const maxConcurrentRuns = options.config?.maxConcurrentRuns; if (maxConcurrentRuns !== undefined && !admitRun(maxConcurrentRuns, activeRunIds().length)) { return json( - { error: `concurrency limit reached (max ${maxConcurrentRuns} concurrent runs)` }, + { error: new ConcurrencyLimitError({ maxConcurrentRuns }).message }, { status: 409 }, ); } @@ -399,19 +399,11 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise runId = await startTrackedRun(options.db, workflow, startOptions); } catch (err) { if (err instanceof ConcurrencyLimitError) { - return json( - { error: `concurrency limit reached (max ${err.maxConcurrentRuns} concurrent runs)` }, - { status: 409 }, - ); + return json({ error: err.message }, { status: 409 }); } - if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { - const dedupeErr = err as DedupeKeyError; + if (err instanceof DedupeKeyError) { return json( - { - error: `dedupe key held: "${dedupeErr.key}" is currently held by run ${dedupeErr.holderRunId}`, - dedupeKey: dedupeErr.key, - holderRunId: dedupeErr.holderRunId, - }, + { error: err.message, dedupeKey: err.key, holderRunId: err.holderRunId }, { status: 409 }, ); } @@ -471,19 +463,11 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise runId = await startTrackedRun(options.db, workflow, startOptions); } catch (err) { if (err instanceof ConcurrencyLimitError) { - return json( - { error: `concurrency limit reached (max ${err.maxConcurrentRuns} concurrent runs)` }, - { status: 409 }, - ); + return json({ error: err.message }, { status: 409 }); } - if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { - const dedupeErr = err as DedupeKeyError; + if (err instanceof DedupeKeyError) { return json( - { - error: `dedupe key held: "${dedupeErr.key}" is currently held by run ${dedupeErr.holderRunId}`, - dedupeKey: dedupeErr.key, - holderRunId: dedupeErr.holderRunId, - }, + { error: err.message, dedupeKey: err.key, holderRunId: err.holderRunId }, { status: 409 }, ); } @@ -557,7 +541,7 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise const maxConcurrentRuns = options.config.maxConcurrentRuns; if (!admitRun(maxConcurrentRuns, activeRunIds().length)) { return json( - { error: `concurrency limit reached (max ${maxConcurrentRuns} concurrent runs)` }, + { error: new ConcurrencyLimitError({ maxConcurrentRuns }).message }, { status: 409 }, ); } @@ -580,19 +564,11 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise }); } catch (err) { if (err instanceof ConcurrencyLimitError) { - return json( - { error: `concurrency limit reached (max ${err.maxConcurrentRuns} concurrent runs)` }, - { status: 409 }, - ); + return json({ error: err.message }, { status: 409 }); } - if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { - const dedupeErr = err as DedupeKeyError; + if (err instanceof DedupeKeyError) { return json( - { - error: `dedupe key held: "${dedupeErr.key}" is currently held by run ${dedupeErr.holderRunId}`, - dedupeKey: dedupeErr.key, - holderRunId: dedupeErr.holderRunId, - }, + { error: err.message, dedupeKey: err.key, holderRunId: err.holderRunId }, { status: 409 }, ); } diff --git a/src/server/runs.ts b/src/server/runs.ts index 229549e..3f37c41 100644 --- a/src/server/runs.ts +++ b/src/server/runs.ts @@ -53,7 +53,12 @@ function isReserved(entry: RunHandle | ReservedSlot | undefined): boole export class ConcurrencyLimitError extends Schema.TaggedError()( "ConcurrencyLimitError", { maxConcurrentRuns: Schema.Number }, -) {} +) { + /** The single source of truth for the message (HTTP, runtime, scheduler alike). */ + override get message(): string { + return `concurrency limit reached (max ${this.maxConcurrentRuns} concurrent runs)`; + } +} export function isActive(runId: string): boolean { return active.has(runId); From dcaf696078c1a176017d5ac266f872b3afeb08b1 Mon Sep 17 00:00:00 2001 From: FreshlyBrewedCode Date: Wed, 23 Sep 2026 06:44:15 +0000 Subject: [PATCH 5/5] fix(sample): restore instanceof matching and the lost comment in ready-sweep (#34) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The tag-matching pass on issue #34 replaced this single-tag `instanceof DedupeKeyError` check with an unsafe `_tag`-cast for consistency with the scheduler, but dropped the comment explaining why a dedupe collision here continues the loop rather than aborting it, and the workflow only ever needs to distinguish one tag — restore both. --- sample/workflows/ready-sweep.ts | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/sample/workflows/ready-sweep.ts b/sample/workflows/ready-sweep.ts index 954566c..0597f9c 100644 --- a/sample/workflows/ready-sweep.ts +++ b/sample/workflows/ready-sweep.ts @@ -207,9 +207,12 @@ export default defineWorkflow("ready-sweep", { await ctx.log("dispatched", { issueNumber: item.issueNumber, childRunId }); dispatches.push({ issueNumber: item.issueNumber, childRunId }); } catch (err) { - if (err instanceof Error && (err as { _tag?: string })._tag === "DedupeKeyError") { - const dedupeErr = err as DedupeKeyError; - collisions.push({ key: dedupeErr.key, holderRunId: dedupeErr.holderRunId }); + // A collision is never a silent no-op (the epic's dispatch decision): + // the runtime records `DispatchCollision` on this run's log either + // way; recorded here too, the remaining items still dispatch, and the + // run fails with the summary below instead of quietly dropping. + if (err instanceof DedupeKeyError) { + collisions.push({ key: err.key, holderRunId: err.holderRunId }); continue; } throw err;