Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion sample/workflows/ready-sweep.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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}`;
},
Expand Down
3 changes: 2 additions & 1 deletion src/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
});
Expand Down
22 changes: 20 additions & 2 deletions src/lib/dedupe.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand All @@ -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", () => {
Expand Down
19 changes: 10 additions & 9 deletions src/lib/dedupe.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,15 @@
* 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>()("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}`;
}
}

Expand All @@ -49,7 +49,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) {
Expand Down
19 changes: 18 additions & 1 deletion src/runtime/run.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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" },
Expand Down
36 changes: 21 additions & 15 deletions src/runtime/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,13 +32,22 @@ import { DedupeKeyError } from "../lib/dedupe";
import type { AgentAdapter } from "./agent-adapter";
import { buildAgentStepEffect } from "./agent-step";

export class RunCancelledSignal extends Error {
constructor() {
super("run cancelled");
this.name = "RunCancelledSignal";
}
/**
* 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);
return err.message || String(err);
}

export class RunCancelledSignal extends Schema.TaggedError<RunCancelledSignal>()(
"RunCancelledSignal",
{},
) {}

/** Matches the spike's model (STATUS.md); overridden per-call or per-workflow. */
export const DEFAULT_MODEL = "opencode-go/deepseek-v4.1-flash";

Expand Down Expand Up @@ -161,7 +170,7 @@ export function startRun<I, O>(
const workspaceKind = options.workspaceKind ?? "clone";

const execImpl = async (argv: ReadonlyArray<string>): Promise<ExecResult> => {
if (cancelled) throw new RunCancelledSignal();
if (cancelled) throw new RunCancelledSignal({});

const execId = nextExecId();
emit({ _tag: "ExecStarted", execId, command: [...argv], cwd: options.dir });
Expand All @@ -180,7 +189,7 @@ export function startRun<I, O>(
durationMs,
});

if (cancelled) throw new RunCancelledSignal();
if (cancelled) throw new RunCancelledSignal({});
return result;
};

Expand All @@ -189,7 +198,7 @@ export function startRun<I, O>(
prompt: string,
opts?: AgentCallOptions,
): Promise<AgentResult<O2>> => {
if (cancelled) throw new RunCancelledSignal();
if (cancelled) throw new RunCancelledSignal({});

const stepId = nextStepId();
// Precedence (issue #16): per-call option > the starting schedule's
Expand Down Expand Up @@ -250,7 +259,7 @@ export function startRun<I, O>(
: {}),
...(handle.partial.usage !== undefined ? { usage: handle.partial.usage } : {}),
});
throw new RunCancelledSignal();
throw new RunCancelledSignal({});
}

const message = Cause.pretty(cause);
Expand Down Expand Up @@ -385,7 +394,7 @@ export function startRun<I, O>(
return result;
} catch (err) {
if (err instanceof RunCancelledSignal) throw err;
const message = err instanceof Error ? err.message : String(err);
const message = domainErrorMessage(err);
emit({
_tag: "WriteBackFinished",
branch: opts.branch,
Expand Down Expand Up @@ -415,9 +424,6 @@ export function startRun<I, O>(
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) {
emit({
_tag: "DispatchCollision",
Expand Down Expand Up @@ -455,7 +461,7 @@ export function startRun<I, O>(
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 };
}
Expand Down Expand Up @@ -493,7 +499,7 @@ export function startRun<I, O>(
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",
Expand Down
3 changes: 1 addition & 2 deletions src/server/dedupe.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
7 changes: 2 additions & 5 deletions src/server/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 },
);
}
Expand Down Expand Up @@ -401,9 +401,6 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise
if (err instanceof ConcurrencyLimitError) {
return json({ error: err.message }, { 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) {
return json(
{ error: err.message, dedupeKey: err.key, holderRunId: err.holderRunId },
Expand Down Expand Up @@ -544,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 },
);
}
Expand Down
32 changes: 32 additions & 0 deletions src/server/runs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import { createSlowFakeAdapter } from "../replay/adapter";
import { defineWorkflow, Schema } from "../workflow";
import {
ConcurrencyLimitError,
DispatchCapError,
activeRunIds,
cancelRegisteredRun,
getActiveHandle,
Expand Down Expand Up @@ -71,6 +72,37 @@ async function waitFor(predicate: () => boolean, timeoutMs = 5_000): Promise<voi
throw new Error("condition not met within timeout");
}

describe("domain errors as TaggedError (#34)", () => {
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();
Expand Down
38 changes: 20 additions & 18 deletions src/server/runs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -41,21 +42,21 @@ 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>()("DispatchCapError", {
message: Schema.String,
}) {}

function isReserved(entry: RunHandle<unknown> | 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>()(
"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)`;
}
}

Expand Down Expand Up @@ -231,21 +232,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
Expand Down Expand Up @@ -313,7 +315,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 });
}
Expand Down
Loading
Loading