diff --git a/src/events.ts b/src/events.ts index ff6de50..3d99d69 100644 --- a/src/events.ts +++ b/src/events.ts @@ -144,6 +144,8 @@ export const RunEventPayload = Schema.TaggedUnion({ /** The run currently holding `key` (non-terminal — it is holding). */ holderRunId: Schema.String, childWorkflowId: Schema.String, + /** Human-readable reason, naming the held key and its holder. */ + error: Schema.optional(Schema.String), }, // ---- ctx.agent(name, prompt, opts?) -------------------------------- diff --git a/src/lib/dedupe.ts b/src/lib/dedupe.ts index 14df6fe..b8aff0c 100644 --- a/src/lib/dedupe.ts +++ b/src/lib/dedupe.ts @@ -58,6 +58,3 @@ export function createDedupeRegistry(): DedupeRegistry { }, }; } - -/** The daemon's shared registry — every run start and dispatch claim in this process goes through it. */ -export const dedupeRegistry: DedupeRegistry = createDedupeRegistry(); diff --git a/src/lib/workspace.test.ts b/src/lib/workspace.test.ts index dfa8305..f45ef68 100644 --- a/src/lib/workspace.test.ts +++ b/src/lib/workspace.test.ts @@ -3,7 +3,7 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { describe, expect, test } from "bun:test"; import type { ExecResult } from "./exec"; -import { allocateWorkspace, evictOldWorkspaces } from "./workspace"; +import { allocateWorkspace, createRefreshGates, evictOldWorkspaces } from "./workspace"; import { writeBack } from "./writeback"; import type { GitIdentity } from "./clone"; import { hostExec } from "./exec"; @@ -29,6 +29,8 @@ describe("allocateWorkspace (D28)", () => { await seedRepo(seed, "one"); const dir = await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-a", workspaceRoot, sshUrl: seed, @@ -53,6 +55,8 @@ describe("allocateWorkspace (D28)", () => { await seedRepo(seed, "one"); await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-one", workspaceRoot, sshUrl: seed, @@ -65,6 +69,8 @@ describe("allocateWorkspace (D28)", () => { await Bun.$`git -C ${seed} -c user.name=seed -c user.email=seed@seed.local commit -q -m second`.quiet(); const dir = await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-two", workspaceRoot, sshUrl: seed, @@ -83,6 +89,8 @@ describe("allocateWorkspace (D28)", () => { for (const runId of ["run-1", "run-2", "run-3"]) { await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId, workspaceRoot, sshUrl: seed, @@ -103,6 +111,7 @@ describe("allocateWorkspace (D28)", () => { const seed = join(root, "seed-repo"); await seedRepo(seed, "one"); + const sharedGates = createRefreshGates(); const [a, b] = await Promise.all([ allocateWorkspace({ runId: "run-x", @@ -110,6 +119,7 @@ describe("allocateWorkspace (D28)", () => { sshUrl: seed, identity: IDENTITY, retainedWorkspaces: 10, + refreshGates: sharedGates, }), allocateWorkspace({ runId: "run-y", @@ -117,6 +127,7 @@ describe("allocateWorkspace (D28)", () => { sshUrl: seed, identity: IDENTITY, retainedWorkspaces: 10, + refreshGates: sharedGates, }), ]); expect(existsSync(join(a, "seed.txt"))).toBe(true); @@ -132,6 +143,8 @@ describe("allocateWorkspace (D28)", () => { await seedRepo(seed, "one"); const dir = await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-a", workspaceRoot, sshUrl: seed, @@ -159,6 +172,8 @@ describe("allocateWorkspace (D28)", () => { await Bun.$`git -C ${seed} push -q ${remote} main`.quiet(); const dir = await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-a", workspaceRoot, sshUrl: remote, @@ -203,6 +218,8 @@ describe("allocateWorkspace (D28)", () => { await seedRepo(seed, "one"); await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-1", workspaceRoot, sshUrl: seed, @@ -211,6 +228,8 @@ describe("allocateWorkspace (D28)", () => { }); await Bun.sleep(5); await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-2", workspaceRoot, sshUrl: seed, @@ -220,6 +239,8 @@ describe("allocateWorkspace (D28)", () => { await Bun.sleep(5); await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-3", workspaceRoot, sshUrl: seed, @@ -241,6 +262,8 @@ describe("allocateWorkspace (D28)", () => { await seedRepo(seed, "one"); await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-1", workspaceRoot, sshUrl: seed, @@ -249,6 +272,8 @@ describe("allocateWorkspace (D28)", () => { }); await Bun.sleep(5); await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-2", workspaceRoot, sshUrl: seed, @@ -258,6 +283,8 @@ describe("allocateWorkspace (D28)", () => { await Bun.sleep(5); await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-3", workspaceRoot, sshUrl: seed, @@ -280,6 +307,8 @@ describe("scratch workspaces (issue #13)", () => { await seedRepo(seed, "one"); const dir = await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-scratch", workspaceRoot, sshUrl: seed, @@ -304,6 +333,8 @@ describe("scratch workspaces (issue #13)", () => { // Two failed-run scratch leftovers, plus one clone. for (const runId of ["scratch-1", "scratch-2"]) { await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId, workspaceRoot, sshUrl: seed, @@ -313,6 +344,8 @@ describe("scratch workspaces (issue #13)", () => { }); } await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "clone-1", workspaceRoot, sshUrl: seed, @@ -367,6 +400,8 @@ describe("host exec injection (issue #13)", () => { // `hostExec` directly, `git clone` would really run and contradict it. const exit0: ExecResult = { command: "fake", exitCode: 0, stdout: "", stderr: "" }; const dir = await allocateWorkspace({ + refreshGates: createRefreshGates(), + runId: "run-fake", workspaceRoot: workspaceRoot2, sshUrl: seed, diff --git a/src/lib/workspace.ts b/src/lib/workspace.ts index 94008f9..15c89ac 100644 --- a/src/lib/workspace.ts +++ b/src/lib/workspace.ts @@ -31,6 +31,21 @@ import type { WorkspaceKind } from "../workflow"; const MIRROR_DIR = ".mirror.git"; +export interface RefreshGates { + get(mirrorPath: string): Promise | undefined; + set(mirrorPath: string, promise: Promise): void; +} + +export function createRefreshGates(): RefreshGates { + const gates = new Map>(); + return { + get: (mirrorPath) => gates.get(mirrorPath), + set: (mirrorPath, promise) => { + gates.set(mirrorPath, promise); + }, + }; +} + export interface WorkspaceAllocationInput { readonly runId: string; readonly workspaceRoot: string; @@ -56,11 +71,15 @@ export interface WorkspaceAllocationInput { readonly protectedEntries?: ReadonlyArray; /** Injectable host-exec seam; defaults to `hostExec` (no behaviour change). */ readonly exec?: ExecFn; + /** Per-daemon refresh gates for serializing mirror maintenance. */ + readonly refreshGates: RefreshGates; } -const refreshGates = new Map>(); - -function enqueueRefresh(mirrorPath: string, task: () => Promise): Promise { +function enqueueRefresh( + refreshGates: RefreshGates, + mirrorPath: string, + task: () => Promise, +): Promise { const prior = refreshGates.get(mirrorPath) ?? Promise.resolve(); const next = prior.then(task, task); refreshGates.set(mirrorPath, next); @@ -110,7 +129,9 @@ export async function allocateWorkspace(input: WorkspaceAllocationInput): Promis const mirrorPath = join(workspaceRoot, MIRROR_DIR); - await enqueueRefresh(mirrorPath, () => refreshMirror(mirrorPath, sshUrl, exec)); + await enqueueRefresh(input.refreshGates, mirrorPath, () => + refreshMirror(mirrorPath, sshUrl, exec), + ); if (existsSync(dir)) { await rm(dir, { recursive: true, force: true }); diff --git a/src/runtime/run.ts b/src/runtime/run.ts index ddb0729..00236f2 100644 --- a/src/runtime/run.ts +++ b/src/runtime/run.ts @@ -361,6 +361,7 @@ export async function startRun( key: err.key, holderRunId: err.holderRunId, childWorkflowId: child.id, + error: domainErrorMessage(err), }); } throw err; diff --git a/src/server/concurrency.test.ts b/src/server/concurrency.test.ts index 051774e..3d2c2d6 100644 --- a/src/server/concurrency.test.ts +++ b/src/server/concurrency.test.ts @@ -22,6 +22,10 @@ import { createSlowFakeAdapter } from "../replay/adapter"; import { makeAgentRuntime } from "../runtime/agent-runtime"; import { admitRun } from "./admission"; import { serve } from "./http"; +import { createRunRegistry } from "./runs"; +import { createPubSub } from "./pubsub"; +import { createDedupeRegistry } from "../lib/dedupe"; +import { createRefreshGates } from "../lib/workspace"; const TREE_WORKFLOW = `${import.meta.dir}/../../test/fixtures/tree-workflow.ts`; @@ -71,7 +75,14 @@ describe("phase 5 P1: per-run working trees (D28) over the real server", () => { const workspaceRoot = join(root, "workspaces"); const db = openStore(join(root, "factory.db")); + const services = { + registry: createRunRegistry(), + pubsub: createPubSub(), + dedupeRegistry: createDedupeRegistry(), + refreshGates: createRefreshGates(), + }; const server = serve({ + services, db, runtime: makeAgentRuntime( createSlowFakeAdapter( @@ -160,7 +171,14 @@ describe("phase 5 P1: per-run working trees (D28) over the real server", () => { await Bun.$`git -C ${seed} -c user.name=seed -c user.email=seed@seed.local commit -q --allow-empty -m seed`.quiet(); const db = openStore(join(root, "factory.db")); + const services = { + registry: createRunRegistry(), + pubsub: createPubSub(), + dedupeRegistry: createDedupeRegistry(), + refreshGates: createRefreshGates(), + }; const server = serve({ + services, db, runtime: makeAgentRuntime( createSlowFakeAdapter( @@ -222,7 +240,14 @@ describe("phase 5 P1: per-run working trees (D28) over the real server", () => { const workspaceRoot = join(root, "workspaces"); const db = openStore(join(root, "factory.db")); + const services = { + registry: createRunRegistry(), + pubsub: createPubSub(), + dedupeRegistry: createDedupeRegistry(), + refreshGates: createRefreshGates(), + }; const server = serve({ + services, db, runtime: makeAgentRuntime( createSlowFakeAdapter( diff --git a/src/server/daemon.ts b/src/server/daemon.ts index 8f506ab..ef4cee6 100644 --- a/src/server/daemon.ts +++ b/src/server/daemon.ts @@ -10,16 +10,28 @@ * threaded by hand through `ServerOptions`, `StartTrackedRunOptions`, * `DispatchEnv` and the run/step options; it is resolved from context inside * `runtime/agent-step.ts`. + * + * #38: the daemon creates per-instance services (run registry, pubsub, + * dedupe registry, refresh gates) so two daemons can coexist in one process + * without sharing state. */ import { mkdir } from "node:fs/promises"; import { dirname } from "node:path"; -import { Effect, type Fiber, ManagedRuntime } from "effect"; +import { Clock, Effect, type Fiber, ManagedRuntime } from "effect"; type AnyFiber = Fiber.Fiber; import type { FactoryConfig } from "../config"; import { openStore } from "../persistence/store"; import { serve } from "./http"; -import { type DispatchEnv, type WorkspaceSpec } from "./runs"; +import { + createRunRegistry, + type DaemonServices, + type DispatchEnv, + type WorkspaceSpec, +} from "./runs"; +import { createPubSub } from "./pubsub"; +import { createDedupeRegistry } from "../lib/dedupe"; +import { createRefreshGates } from "../lib/workspace"; import { createSchedulerState, makeScheduleFire, @@ -30,27 +42,33 @@ import { import { AgentRuntimeLayer } from "../runtime/agent-runtime"; import { opencodeAdapter } from "../runtime/opencode-adapter"; +export function createLiveClock(): Clock.Clock { + return { + currentTimeMillisUnsafe: () => Date.now(), + currentTimeMillis: Effect.sync(() => Date.now()), + monotonicTimeNanosUnsafe: () => BigInt(Date.now()) * 1_000_000n, + monotonicTimeNanos: Effect.sync(() => BigInt(Date.now()) * 1_000_000n), + currentTimeNanosUnsafe: () => BigInt(Date.now()) * 1_000_000n, + currentTimeNanos: Effect.sync(() => BigInt(Date.now()) * 1_000_000n), + sleep: (duration) => Effect.sleep(duration), + }; +} + export interface DaemonOptions { readonly dbPath: string; readonly port?: number; /** Issue #16: the tick cadence of the config schedules, over the default. */ readonly schedulerIntervalMs?: number; - /** - * The run environment from `factory.config.ts` (D27). Present, every run — - * manual and scheduled alike — gets a per-run working tree under the - * configured `workspaceRoot` (D28) and shares one admission limit (D29). - * Absent, the legacy per-request `{dir, clone}` behaviour is kept. - */ readonly config?: FactoryConfig; + readonly clock?: Clock.Clock; } export interface DaemonHandle { readonly server: ReturnType; - /** Issue #16: the loop that fires the config's schedules, when it has any. */ readonly schedulerFiber: AnyFiber | undefined; + readonly services: DaemonServices; } -/** Issue #16: how often the scheduler's due window check runs. */ export const DEFAULT_SCHEDULER_INTERVAL_MS = 30_000; export async function startDaemon(options: DaemonOptions): Promise { @@ -61,9 +79,17 @@ export async function startDaemon(options: DaemonOptions): Promise const layer = AgentRuntimeLayer(adapter); const runtime = ManagedRuntime.make(layer); + const services: DaemonServices = { + registry: createRunRegistry(), + pubsub: createPubSub(), + dedupeRegistry: createDedupeRegistry(), + refreshGates: createRefreshGates(), + }; + const server = serve({ db, runtime, + services, port: options.port, ...(options.config !== undefined ? { config: options.config } : {}), }); @@ -92,11 +118,14 @@ export async function startDaemon(options: DaemonOptions): Promise fire: makeScheduleFire({ db, runtime, + services, maxConcurrentRuns, workspace, repo, dispatchEnv, }), + dedupeRegistry: services.dedupeRegistry, + clock: options.clock ?? createLiveClock(), }; const state = createSchedulerState(deps); return Effect.runFork( @@ -109,5 +138,5 @@ export async function startDaemon(options: DaemonOptions): Promise })() : undefined; - return { server, schedulerFiber }; + return { server, schedulerFiber, services }; } diff --git a/src/server/dedupe.test.ts b/src/server/dedupe.test.ts index 845d8ac..632bdc9 100644 --- a/src/server/dedupe.test.ts +++ b/src/server/dedupe.test.ts @@ -4,6 +4,11 @@ * releases it when the run reaches any terminal state. Starting a second run * with a key a non-terminal run already holds throws a `DedupeKeyError` * naming the key and the holding run. + * + * #38: tests create their own DaemonServices instead of using module-level + * singletons or test seams. The `beforeStart` test seam is gone — tests + * reserve registry slots directly to test the reserved-but-not-started + * window. */ import { mkdtempSync, mkdirSync, rmSync } from "node:fs"; @@ -15,13 +20,14 @@ import { createSlowFakeAdapter } from "../replay/adapter"; import { makeAgentRuntime } from "../runtime/agent-runtime"; import { defineWorkflow, Schema } from "../workflow"; import { createDedupeRegistry, DedupeKeyError } from "../lib/dedupe"; -import { activeRunIds, getActiveHandle, isActive, startTrackedRun } from "./runs"; +import { createRunRegistry, startTrackedRun, type DaemonServices } from "./runs"; +import { createPubSub } from "./pubsub"; +import { createRefreshGates } from "../lib/workspace"; import { serve } from "./http"; import { defineConfig } from "../config"; -/** Fire-and-forget children keep filling the registry; drain before closing the store. */ -async function drain(timeoutMs = 10_000): Promise { - await waitFor(() => activeRunIds().length === 0, timeoutMs); +async function drain(services: DaemonServices, timeoutMs = 10_000): Promise { + await waitFor(() => services.registry.activeRunIds().length === 0, timeoutMs); } const SLOW_ADAPTER = createSlowFakeAdapter( @@ -43,19 +49,6 @@ const echoWorkflow = defineWorkflow("echo-wf", { }, }); -interface Gate { - readonly wait: () => Promise; - readonly release: () => void; -} - -function makeGate(): Gate { - let release: () => void = () => undefined; - const promise = new Promise((resolve) => { - release = resolve; - }); - return { wait: () => promise, release }; -} - async function waitFor(predicate: () => boolean, timeoutMs = 5_000): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { @@ -65,28 +58,32 @@ async function waitFor(predicate: () => boolean, timeoutMs = 5_000): Promise Promise; - /** Shared holder registry when one test needs two starts to agree. */ - readonly registry?: ReturnType; + readonly services: DaemonServices; } function start(db: ReturnType, opts: StartOpts): Promise { - return startTrackedRun(runtime, db, echoWorkflow, { + return startTrackedRun(runtime, db, opts.services, echoWorkflow, { dir: opts.dir, input: {}, - dedupeRegistry: opts.registry ?? createDedupeRegistry(), maxConcurrentRuns: 4, ...(opts.runId !== undefined ? { runId: opts.runId } : {}), ...(opts.dedupeKey !== undefined ? { dedupeKey: opts.dedupeKey } : {}), - ...(opts.beforeStart !== undefined ? { beforeStart: opts.beforeStart } : {}), }); } -/** The `DedupeKeyError` a rejected start produced — else a placeholder string. */ async function errorOf(promise: Promise): Promise { return promise.then( () => "resolved", @@ -118,6 +115,7 @@ function withWorkspaces(root: string): Record { } describe("dedupe keys through ctx.dispatch (issue #15)", () => { + const driftAdapter = SLOW_ADAPTER; const busyWorkflow = defineWorkflow("busy-wf", { input: Schema.Struct({ n: Schema.Number }), workspace: { kind: "scratch" }, @@ -131,6 +129,7 @@ describe("dedupe keys through ctx.dispatch (issue #15)", () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-dispatch-")); const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "workspaces"), { recursive: true }); + const services = createTestServices(); const parent = defineWorkflow("parent-wf", { input: Schema.Struct({ wait: Schema.Number }), @@ -141,20 +140,19 @@ describe("dedupe keys through ctx.dispatch (issue #15)", () => { { n: input.wait }, { dedupeKey: "item:7" }, ); - // A second dispatch on the same key, while the first still runs. await ctx.dispatch(busyWorkflow, { n: input.wait }, { dedupeKey: "item:7" }); return { childRunId }; }, }); - const parentRunId = await startTrackedRun(runtime, db, parent, { + const parentRunId = await startTrackedRun(runtime, db, services, parent, { input: { wait: 3 }, workspace: workspaces(join(root)), maxConcurrentRuns: 10, - dispatchEnv: withWorkspaces(join(root)), + dispatchEnv: { ...withWorkspaces(join(root)), adapter: driftAdapter } as never, }); - await waitFor(() => !isActive(parentRunId)); + await waitFor(() => !services.registry.isActive(parentRunId)); const events = getRunEvents(db, parentRunId); const collided = events.find((event) => event.payload._tag === "DispatchCollision"); @@ -164,13 +162,12 @@ describe("dedupe keys through ctx.dispatch (issue #15)", () => { expect(collided.payload.holderRunId).toBeTypeOf("string"); expect(collided.payload.childWorkflowId).toBe("busy-wf"); } - // The failed dispatch's throw unwound the parent into its terminal state. const failed = events.find((event) => event.payload._tag === "RunFailed"); expect( failed !== undefined && failed.payload._tag === "RunFailed" ? failed.payload.message : "", ).toMatch(/dedupe key held: "item:7"/); - await drain(); + await drain(services); db.close(); rmSync(root, { recursive: true, force: true }); }); @@ -179,55 +176,59 @@ describe("dedupe keys through ctx.dispatch (issue #15)", () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-dispatch-release-")); const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "workspaces"), { recursive: true }); + const services = createTestServices(); - let release: () => void = () => undefined; - const gate = makeGate(); + let releaseHeldWorkflow: () => void = () => undefined; + const heldWorkflowGate = new Promise((resolve) => { + releaseHeldWorkflow = resolve; + }); const heldWorkflow = defineWorkflow("held-wf", { input: Schema.Struct({}), workspace: { kind: "scratch" }, run: async () => { - await gate.wait; + await heldWorkflowGate; return {}; }, }); + let releaseParent: () => void = () => undefined; + const parentGate = new Promise((resolve) => { + releaseParent = resolve; + }); + const parent = defineWorkflow("parent-release-wf", { input: Schema.Struct({}), workspace: { kind: "scratch" }, run: async (ctx) => { const ids = []; ids.push(await ctx.dispatch(heldWorkflow, {}, { dedupeKey: "item:8" })); - await new Promise((r) => { - release = r; - }); - // The first child is terminal by now; the key must be free again. + await parentGate; ids.push(await ctx.dispatch(heldWorkflow, {}, { dedupeKey: "item:8" })); return { ids }; }, }); - const parentRunId = await startTrackedRun(runtime, db, parent, { + const parentRunId = (await startTrackedRun(runtime, db, services, parent, { input: {}, workspace: workspaces(join(root)), maxConcurrentRuns: 10, - dispatchEnv: withWorkspaces(join(root)), - }); + dispatchEnv: { ...withWorkspaces(join(root)), adapter: driftAdapter } as never, + })) as string; await waitFor(() => getRunEvents(db, parentRunId).some((event) => event.payload._tag === "RunDispatched"), ); await Bun.sleep(150); - // The first child dispatched, still holding while the gate is shut — one - // more dispatch inside the parent would collide. Then let the child end. - gate.release(); - if (typeof release === "function") release(); + releaseHeldWorkflow(); + await Bun.sleep(50); + releaseParent(); await waitFor( () => getRunEvents(db, parentRunId).filter((event) => event.payload._tag === "RunDispatched") .length >= 2, 10_000, ); - await waitFor(() => !isActive(parentRunId)); + await waitFor(() => !services.registry.isActive(parentRunId)); const dispatchedCount = getRunEvents(db, parentRunId).filter( (event) => event.payload._tag === "RunDispatched", @@ -237,7 +238,7 @@ describe("dedupe keys through ctx.dispatch (issue #15)", () => { getRunEvents(db, parentRunId).some((event) => event.payload._tag === "DispatchCollision"), ).toBe(false); - await drain(); + await drain(services); db.close(); rmSync(root, { recursive: true, force: true }); }); @@ -247,33 +248,34 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { test("a key held by a non-terminal run rejects a second start with key and holder named", async () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-hold-")); const db = openStore(join(root, "factory.db")); - const gate = makeGate(); - const registry = createDedupeRegistry(); + const services = createTestServices(); - const first = start(db, { - runId: "run-holder", - dir: join(root, "dir-1"), - dedupeKey: "item:41", - beforeStart: gate.wait, - registry, - }); - await waitFor(() => isActive("run-holder")); + services.registry.reserve("run-holder"); + services.dedupeRegistry.claim("item:41", "run-holder"); const err = (await errorOf( - start(db, { runId: "run-collide", dir: join(root, "dir-2"), dedupeKey: "item:41", registry }), + start(db, { runId: "run-collide", dir: join(root, "dir-2"), dedupeKey: "item:41", services }), )) as DedupeKeyError | string; expect(err).toBeInstanceOf(DedupeKeyError); if (err instanceof DedupeKeyError) { expect(err.key).toBe("item:41"); expect(err.holderRunId).toBe("run-holder"); - expect(err._tag).toBe("DedupeKeyError"); + expect(err.message).toContain("item:41"); + expect(err.message).toContain("run-holder"); } - gate.release(); - await first; - await waitFor(() => !isActive("run-holder")); + services.registry.delete("run-holder"); + services.dedupeRegistry.release("item:41", "run-holder"); + + const firstRunId = await start(db, { + runId: "run-holder", + dir: join(root, "dir-1"), + dedupeKey: "item:41", + services, + }); + await waitFor(() => !services.registry.isActive(firstRunId)); - await drain(); + await drain(services); db.close(); rmSync(root, { recursive: true, force: true }); }); @@ -281,25 +283,20 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { test("the holder's RunStarted carries the key; a rejected start leaves no lifecycle trace", async () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-log-")); const db = openStore(join(root, "factory.db")); - const gate = makeGate(); - const registry = createDedupeRegistry(); + const services = createTestServices(); - const first = start(db, { + const firstRunId = await start(db, { runId: "run-holder", dir: join(root, "dir-1"), dedupeKey: "item:42", - beforeStart: gate.wait, - registry, + services, }); - await waitFor(() => isActive("run-holder")); await errorOf( - start(db, { runId: "run-collide", dir: join(root, "dir-2"), dedupeKey: "item:42", registry }), + start(db, { runId: "run-collide", dir: join(root, "dir-2"), dedupeKey: "item:42", services }), ); - gate.release(); - await first; - await waitFor(() => !isActive("run-holder")); + await waitFor(() => !services.registry.isActive(firstRunId)); const holderStarted = getRunEvents(db, "run-holder").find( (e) => e.payload._tag === "RunStarted", @@ -311,10 +308,9 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { ? holderStarted.payload.dedupeKey : "missing", ).toBe("item:42"); - // A dropped start never started: the collision run's log has no rows. expect(getRunEvents(db, "run-collide")).toHaveLength(0); - await drain(); + await drain(services); db.close(); rmSync(root, { recursive: true, force: true }); }); @@ -322,16 +318,25 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { test("the key is released when the holding run reaches its terminal state", async () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-finish-")); const db = openStore(join(root, "factory.db")); + const services = createTestServices(); - const firstRunId = await start(db, { dir: join(root, "dir-1"), dedupeKey: "item:43" }); + const firstRunId = await start(db, { + dir: join(root, "dir-1"), + dedupeKey: "item:43", + services, + }); expect(firstRunId).toBeTypeOf("string"); - await waitFor(() => !isActive(firstRunId)); + await waitFor(() => !services.registry.isActive(firstRunId)); - const retryRunId = await start(db, { dir: join(root, "dir-2"), dedupeKey: "item:43" }); - await waitFor(() => !isActive(retryRunId)); + const retryRunId = await start(db, { + dir: join(root, "dir-2"), + dedupeKey: "item:43", + services, + }); + await waitFor(() => !services.registry.isActive(retryRunId)); expect(getRunEvents(db, retryRunId).some((e) => e.payload._tag === "RunStarted")).toBe(true); - await drain(); + await drain(services); db.close(); rmSync(root, { recursive: true, force: true }); }); @@ -339,21 +344,24 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { test("a failed start releases its claim so a corrected retry can start", async () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-fail-")); const db = openStore(join(root, "factory.db")); + const services = createTestServices(); - // allocation of an explicitly undefined dir fails after the claim await expect( - startTrackedRun(runtime, db, echoWorkflow, { + startTrackedRun(runtime, db, services, echoWorkflow, { input: {}, dedupeKey: "item:44", - dedupeRegistry: createDedupeRegistry(), }), ).rejects.toThrow(/needs `dir` or `workspace`/); - const retryRunId = await start(db, { dir: join(root, "dir-1"), dedupeKey: "item:44" }); - await waitFor(() => !isActive(retryRunId)); + const retryRunId = await start(db, { + dir: join(root, "dir-1"), + dedupeKey: "item:44", + services, + }); + await waitFor(() => !services.registry.isActive(retryRunId)); expect(getRunEvents(db, retryRunId).some((e) => e.payload._tag === "RunStarted")).toBe(true); - await drain(); + await drain(services); db.close(); rmSync(root, { recursive: true, force: true }); }); @@ -361,31 +369,29 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { test("a cancelled run releases its key (RunCancelled is a terminal state)", async () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-cancel-")); const db = openStore(join(root, "factory.db")); - const registry = createDedupeRegistry(); + const services = createTestServices(); - const holderRunId = await startTrackedRun(runtime, db, echoWorkflow, { + const holderRunId = await startTrackedRun(runtime, db, services, echoWorkflow, { runId: "run-cancelled", dir: join(root, "dir-1"), input: {}, dedupeKey: "item:45", - dedupeRegistry: registry, }); - await waitFor(() => !isActive(holderRunId) || true); - const handle = getActiveHandle("run-cancelled"); + await waitFor(() => !services.registry.isActive(holderRunId) || true); + const handle = services.registry.getActiveHandle("run-cancelled"); expect(handle).toBeDefined(); if (handle !== undefined) await handle.cancel(); - await waitFor(() => !isActive(holderRunId)); + await waitFor(() => !services.registry.isActive(holderRunId)); const started = getRunEvents(db, holderRunId).some((e) => e.payload._tag === "RunStarted"); - const retryRunId = await startTrackedRun(runtime, db, echoWorkflow, { + const retryRunId = await startTrackedRun(runtime, db, services, echoWorkflow, { dir: join(root, "dir-2"), input: {}, dedupeKey: "item:45", - dedupeRegistry: registry, }); expect(started).toBe(true); expect(retryRunId).toBeTypeOf("string"); - await drain(); + await drain(services); db.close(); rmSync(root, { recursive: true, force: true }); @@ -394,14 +400,17 @@ describe("dedupe keys through startTrackedRun (issue #15)", () => { test("runs started without a key behave exactly as before", async () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-none-")); const db = openStore(join(root, "factory.db")); + const services = createTestServices(); - const firstRunId = await start(db, { dir: join(root, "dir-1") }); - const secondRunId = await start(db, { dir: join(root, "dir-2") }); + const firstRunId = await start(db, { dir: join(root, "dir-1"), services }); + const secondRunId = await start(db, { dir: join(root, "dir-2"), services }); expect(firstRunId).toBeTypeOf("string"); expect(secondRunId).toBeTypeOf("string"); - await waitFor(() => !isActive(firstRunId) && !isActive(secondRunId)); + await waitFor( + () => !services.registry.isActive(firstRunId) && !services.registry.isActive(secondRunId), + ); - await drain(); + await drain(services); db.close(); rmSync(root, { recursive: true, force: true }); }); @@ -425,11 +434,11 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { test("a key held by a running run surfaces as a 409 conflict naming the key and the holder", async () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-http-")); const db = openStore(join(root, "factory.db")); + const services = createTestServices(); const slowStartWorkflow = defineWorkflow("dedupe-http-slow", { input: Schema.Struct({}), workspace: { kind: "scratch" }, - // Scratch: no clone needed. The key is claimed before the workflow runs. run: async (ctx) => { await ctx.exec(["sh", "-c", "sleep 0.5"]); return {}; @@ -437,8 +446,9 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { }); const server = serve({ + runtime, db, - runtime: makeAgentRuntime(SLOW_ADAPTER), + services, port: 0, config: defineConfig({ repo: { @@ -481,11 +491,9 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { expect(collision.error).toContain("issue:91"); expect(collision.error).toContain(holderRunId); - // A conflict is a start rejection, not a started run: nothing recorded. expect(getRunEvents(db, "nonexistent")).toHaveLength(0); await waitForTerminalStatus(db, holderRunId); - // Rejected while held; once terminal, a retry is a fresh 201. const retryRes = await fetch(`${base}/api/runs`, { method: "POST", body: JSON.stringify({ workflowId: "dedupe-http-slow", input: {}, dedupeKey: "issue:91" }), @@ -494,7 +502,7 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { const { runId: retryRunId } = (await retryRes.json()) as { runId: string }; await waitForTerminalStatus(db, retryRunId); } finally { - await drain(); + await drain(services); await server.stop(true); db.close(); rmSync(root, { recursive: true, force: true }); @@ -504,6 +512,7 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { test("a non-string dedupeKey is a 400, and omitting one behaves as today", async () => { const root = mkdtempSync(join(tmpdir(), "factory-dedupe-http-shape-")); const db = openStore(join(root, "factory.db")); + const services = createTestServices(); const plainWorkflow = defineWorkflow("dedupe-http-plain", { input: Schema.Struct({}), workspace: { kind: "scratch" }, @@ -511,8 +520,9 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { }); const server = serve({ + runtime, db, - runtime: makeAgentRuntime(SLOW_ADAPTER), + services, port: 0, config: defineConfig({ repo: { @@ -543,7 +553,7 @@ describe("dedupe keys over POST /api/runs (issue #15)", () => { const { runId } = (await plain.json()) as { runId: string }; await waitForTerminalStatus(db, runId); } finally { - await drain(); + await drain(services); await server.stop(true); db.close(); rmSync(root, { recursive: true, force: true }); diff --git a/src/server/http.test.ts b/src/server/http.test.ts index 25438d0..1f0913b 100644 --- a/src/server/http.test.ts +++ b/src/server/http.test.ts @@ -11,6 +11,10 @@ import { defineConfig, loadFactoryConfig } from "../config"; import registryWorkflow from "../../test/fixtures/registry-workflow"; import registryScratchWorkflow from "../../test/fixtures/registry-scratch-workflow"; import { serve } from "./http"; +import { createRunRegistry } from "./runs"; +import { createPubSub } from "./pubsub"; +import { createDedupeRegistry } from "../lib/dedupe"; +import { createRefreshGates } from "../lib/workspace"; const ECHO_WORKFLOW = `${import.meta.dir}/../../test/fixtures/echo-workflow.ts`; const QUIET_GAP_WORKFLOW = `${import.meta.dir}/../../test/fixtures/quiet-gap-workflow.ts`; @@ -89,8 +93,16 @@ describe("GET /api/workflows (D30)", () => { const dir = mkdtempSync(join(tmpdir(), "factory-workflows-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); + const runtime = makeAgentRuntime(adapter); const config = await loadFactoryConfig(FIXTURE_CONFIG); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); + const services = createTestServices(); + const server = serve({ + runtime, + services, + db, + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -115,7 +127,14 @@ describe("GET /api/workflows (D30)", () => { const dir = mkdtempSync(join(tmpdir(), "factory-workflows-empty-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); + const runtime = makeAgentRuntime(adapter); + const services = createTestServices(); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -135,7 +154,14 @@ describe("GET /api/workflows (D30)", () => { const dir = mkdtempSync(join(tmpdir(), "factory-workflows-empty-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); + const runtime = makeAgentRuntime(adapter); + const services = createTestServices(); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -155,7 +181,14 @@ describe("phase 4 SPA serving", () => { const dir = mkdtempSync(join(tmpdir(), "factory-spa-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([{ type: "TEXT_MESSAGE_START" }], 1); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -238,7 +271,15 @@ describe("SSE keepalive", () => { const dir = mkdtempSync(join(tmpdir(), "factory-sse-keepalive-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, sseKeepaliveMs: 25 }); + const runtime = makeAgentRuntime(adapter); + const services = createTestServices(); + const server = serve({ + runtime, + services, + db, + port: 0, + sseKeepaliveMs: 25, + }); const base = `http://localhost:${server.port}`; try { @@ -274,7 +315,15 @@ describe("SSE keepalive", () => { const dir = mkdtempSync(join(tmpdir(), "factory-sse-keepalive-stop-test-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, sseKeepaliveMs: 10 }); + const runtime = makeAgentRuntime(adapter); + const services = createTestServices(); + const server = serve({ + runtime, + services, + db, + port: 0, + sseKeepaliveMs: 10, + }); const base = `http://localhost:${server.port}`; try { @@ -314,7 +363,14 @@ describe("phase 3 HTTP API + SSE", () => { ], 1, ); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -363,7 +419,14 @@ describe("phase 3 HTTP API + SSE", () => { ], 500, ); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -399,7 +462,14 @@ describe("phase 3 HTTP API + SSE", () => { ], 60, ); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -438,7 +508,14 @@ describe("phase 3 HTTP API + SSE", () => { ], 400, ); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -476,7 +553,14 @@ describe("phase 3 HTTP API + SSE", () => { ], 1, ); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0 }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -522,7 +606,15 @@ describe("POST /api/runs {workflowId, input} (D31)", () => { 1, ); const config = await loadFactoryConfig(FIXTURE_CONFIG); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -547,8 +639,16 @@ describe("POST /api/runs {workflowId, input} (D31)", () => { const dir = mkdtempSync(join(tmpdir(), "factory-d31-unknown-")); const db = openStore(join(dir, "factory.db")); const adapter = createSlowFakeAdapter([], 1); + const runtime = makeAgentRuntime(adapter); const config = await loadFactoryConfig(FIXTURE_CONFIG); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); + const services = createTestServices(); + const server = serve({ + runtime, + services, + db, + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -594,7 +694,15 @@ describe("POST /api/runs {workflowId, input} (D31)", () => { maxConcurrentRuns: 1, retainedWorkspaces: 10, }); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -654,7 +762,15 @@ describe("POST /api/runs {workflowId, input} (D31)", () => { maxConcurrentRuns: 3, retainedWorkspaces: 10, }); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -721,7 +837,15 @@ describe("a scratch workflow through POST /api/runs (issue #13)", () => { maxConcurrentRuns: 3, retainedWorkspaces: 10, }); - const server = serve({ db, runtime: makeAgentRuntime(adapter), port: 0, config }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); + const server = serve({ + runtime, + services, + db, + port: 0, + config, + }); const base = `http://localhost:${server.port}`; try { @@ -792,9 +916,12 @@ describe("ctx.dispatch through POST /api/runs (issue #14)", () => { }, }); + const services = createTestServices(); + const runtime = makeAgentRuntime(adapter); const server = serve({ + runtime, + services, db, - runtime: makeAgentRuntime(adapter), port: 0, config: defineConfig({ repo: { @@ -890,9 +1017,12 @@ describe("GET /api/schedules (issue #17)", () => { workspace: { kind: "scratch" }, run: async () => ({}), }); + const services = createTestServices(); + const runtime = makeAgentRuntime(createSlowFakeAdapter([], 1)); const server = serve({ + runtime, + services, db, - runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0, config: schedulesConfig(root, workflow), }); @@ -959,9 +1089,12 @@ describe("GET /api/schedules (issue #17)", () => { }, ], }); + const services = createTestServices(); + const runtime = makeAgentRuntime(createSlowFakeAdapter([], 1)); const server = serve({ + runtime, + services, db, - runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0, config, }); @@ -985,7 +1118,14 @@ describe("GET /api/schedules (issue #17)", () => { test("serves an empty list on a legacy no-config server", async () => { const root = mkdtempSync(join(tmpdir(), "factory-schedules-empty-test-")); const db = openStore(join(root, "factory.db")); - const server = serve({ db, runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0 }); + const services = createTestServices(); + const runtime = makeAgentRuntime(createSlowFakeAdapter([], 1)); + const server = serve({ + runtime, + services, + db, + port: 0, + }); const base = `http://localhost:${server.port}`; try { @@ -1027,9 +1167,12 @@ describe("GET /api/schedules (issue #17)", () => { }, ], }); + const services = createTestServices(); + const runtime = makeAgentRuntime(createSlowFakeAdapter([], 1)); const server = serve({ + runtime, + services, db, - runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0, config, }); @@ -1128,9 +1271,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { }, ], }); + const services = createTestServices(); + const runtime = makeAgentRuntime(createSlowFakeAdapter([], 1)); const server = serve({ + runtime, + services, db, - runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0, config, }); @@ -1182,9 +1328,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { { id: "nightly-run", workflow: workflow.id, input: {}, cron: "0 3 * * *", timezone: "UTC" }, ], }); + const services = createTestServices(); + const runtime = makeAgentRuntime(createSlowFakeAdapter([], 1)); const server = serve({ + runtime, + services, db, - runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0, config, }); @@ -1225,9 +1374,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { { id: "slow-skip", workflow: workflow.id, input: {}, cron: "0 3 * * *", timezone: "UTC" }, ], }); + const services = createTestServices(); + const runtime = makeAgentRuntime(createSlowFakeAdapter([], 1)); const server = serve({ + runtime, + services, db, - runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0, config, }); @@ -1290,9 +1442,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { }, ], }); + const services = createTestServices(); + const runtime = makeAgentRuntime(createSlowFakeAdapter([], 1)); const server = serve({ + runtime, + services, db, - runtime: makeAgentRuntime(createSlowFakeAdapter([], 1)), port: 0, config, }); @@ -1317,3 +1472,12 @@ describe("POST /api/schedules/:id/run (issue #17)", () => { } }); }); + +function createTestServices() { + return { + registry: createRunRegistry(), + pubsub: createPubSub(), + dedupeRegistry: createDedupeRegistry(), + refreshGates: createRefreshGates(), + }; +} diff --git a/src/server/http.ts b/src/server/http.ts index 738a0a5..0668438 100644 --- a/src/server/http.ts +++ b/src/server/http.ts @@ -25,38 +25,32 @@ * between, then drains the buffer de-duplicated by `seq` once the persisted * read completes. Without this ordering a live event emitted between the * subscribe and the read could be lost. + * + * #38: the handler receives DaemonServices (registry, pubsub, dedupe) as + * required parameters rather than accessing module-level singletons. */ import type { Database } from "bun:sqlite"; -import { ManagedRuntime } from "effect"; import { isTerminal, type RunEvent } from "../events"; import type { FactoryConfig } from "../config"; import { Schema, SchemaParser } from "effect"; import { resetClone, type GitIdentity } from "../lib/clone"; import { loadWorkflow } from "../lib/load-workflow"; import { getRunEvents, listRuns, type RunSummary } from "../persistence/store"; +import type { ManagedRuntime } from "effect"; import { AgentRuntime } from "../runtime/agent-runtime"; import index from "../web/index.html"; import { admitRun } from "./admission"; -import { subscribe } from "./pubsub"; import { DedupeKeyError } from "../lib/dedupe"; import { ConcurrencyLimitError, - activeRunIds, - cancelRegisteredRun, - isActive, - startTrackedRun, + type DaemonServices, type DispatchEnv, type StartTrackedRunOptions, + startTrackedRun, } from "./runs"; import { nextFireAt, toRuntimeSchedules } from "./scheduler"; -/** - * Issue #17: what `GET /api/schedules` serves for one schedule — the config's - * own record (id, workflow, input, cron, timezone, overlap, run-on-start), - * the next fire time computed from the stored cron, and the most recent run - * that carries the schedule's id. - */ export interface ScheduleSummary { readonly id: string; readonly workflowId: string; @@ -65,9 +59,7 @@ export interface ScheduleSummary { readonly timezone: string; readonly overlap: "skip" | "stack"; readonly runOnStart: boolean; - /** Epoch ms of the cron's next fire after now. */ readonly nextFireAt: number; - /** The schedule's latest run — scheduled *or* manually triggered. */ readonly lastRun: | { readonly runId: string; readonly status: string; readonly startedAt: number } | undefined; @@ -75,68 +67,25 @@ export interface ScheduleSummary { export interface ServerOptions { readonly db: Database; + readonly services: DaemonServices; /** * Issue #36: the Effect managed runtime that provides the agent runtime * service. The adapter is resolved from context inside `startTrackedRun` * rather than threaded through `ServerOptions`. */ readonly runtime: ManagedRuntime.ManagedRuntime; - /** - * The loaded `factory.config.ts` (D27). Absent = the phase 3 path-based API - * behaves exactly as before (no limit, explicit dir+clone, `/api/workflows` - * serves an empty list) — the dispatcher legitimately keeps supplying - * filesystem paths (D31). - */ readonly config?: FactoryConfig; - /** - * How often an open SSE stream emits its keepalive comment frame. Exposed - * only so tests can turn it down far enough for a sub-second quiet gap to - * exercise the timer; production has no reason to set it. - */ readonly sseKeepaliveMs?: number; } -/** - * `Bun.serve`'s `idleTimeout` defaults to **10 seconds** and closes any - * connection with no traffic in that window. A run's SSE stream only carries - * traffic when the workflow emits an event, so any quiet gap longer than that — - * `ctx.exec` running a test suite, write-back's `git push` + `gh pr create`, an - * agent thinking before its first chunk — dropped the connection under a - * perfectly healthy run. The SPA has no way to tell that apart from a dead - * process and rendered the run "interrupted" until a manual refresh. - * - * A comment frame on a timer is the fix rather than a raised `idleTimeout`, - * because Bun caps `idleTimeout` at 255s: a single long agent step would still - * outlast it, whereas a keepalive holds a stream open for any gap length. - * SSE comment frames (`:`-prefixed, no `data:` line) are inert to every client - * we ship — the SPA's `dataLine` (`web/api.ts`) and the CLI's tail both skip - * frames without a `data:` line. - * - * 5s leaves 2x margin under the 10s default. It is deliberately not derived - * from `idleTimeout`: we do not set that option, so the default is the contract. - */ const DEFAULT_SSE_KEEPALIVE_MS = 5_000; -/** `RunSummary` plus this process's live-registry bit — what the SPA reads. */ export type RunSummaryResponse = RunSummary & { readonly active: boolean }; -/** - * A run with no terminal event is `"interrupted"` in the store, which cannot - * tell a crash apart from a run this process still holds. The runs page needs - * that bit to group live rows, so it is derived here from the in-memory - * registry (`server/runs.ts`) rather than stored (D24: active is process state). - */ -function listSummaries(db: Database): ReadonlyArray { - return listRuns(db).map((run) => ({ ...run, active: isActive(run.runId) })); +function listSummaries(db: Database, services: DaemonServices): ReadonlyArray { + return listRuns(db).map((run) => ({ ...run, active: services.registry.isActive(run.runId) })); } -/** - * Issue #14: the dispatch environment a run's `ctx.dispatch` shares with its - * child runs — the same config wiring the run itself gets (workspace root, - * repo, adapter, admission). Present only for config-backed servers; without - * it `ctx.dispatch` throws, because in-process/path-based legacy runs have - * no registry to start children from. - */ function dispatchEnvFor(config: FactoryConfig): DispatchEnv { return { workspace: { @@ -152,12 +101,6 @@ function dispatchEnvFor(config: FactoryConfig): DispatchEnv { }; } -/** - * The start options every config-backed `startTrackedRun` shares: D28's - * workspace allocation, the run environment, D29's admission limit and #14's - * dispatch env. Used by the registry POST path and, since issue #17, by the - * manual schedule trigger. - */ function configRunOptions( config: FactoryConfig, maxConcurrentRuns: number | undefined, @@ -187,11 +130,6 @@ interface StartRunBody { readonly workflowId?: unknown; readonly workflowPath?: unknown; readonly input?: unknown; - /** - * Issue #15: an optional dedupe key. While a non-terminal run holds it, - * another start with the same key is a 409 conflict naming the key and the - * holding run — never a silently dropped duplicate. - */ readonly dedupeKey?: unknown; readonly dir?: string; readonly clone?: { readonly sshUrl: string; readonly identity: GitIdentity }; @@ -204,11 +142,6 @@ function json(body: unknown, init?: { readonly status?: number }): Response { }); } -/** - * The SSE `Last-Event-ID` header, or `undefined` if absent or not a sequence - * number. `EventSource` sends it automatically on reconnect, and `seq` is the - * value we stamp on each frame's `id:` line, so it is the resume offset. - */ function parseLastEventId(req: Request): number | undefined { const raw = req.headers.get("last-event-id"); if (raw === null || !/^\d+$/.test(raw)) return undefined; @@ -217,15 +150,13 @@ function parseLastEventId(req: Request): number | undefined { function sseStream( db: Database, + services: DaemonServices, runId: string, lastEventId?: number, keepaliveMs: number = DEFAULT_SSE_KEEPALIVE_MS, ): ReadableStream { const encoder = new TextEncoder(); - // `closed`/`unsubscribe`/`keepalive` live on the stream (not just `start`) so - // the `cancel()` hook below can also mark them when a client disconnects - // without anyone having sent a terminal event. let closed = false; let unsubscribe: () => void = () => undefined; let keepalive: ReturnType | undefined; @@ -236,9 +167,6 @@ function sseStream( return new ReadableStream({ start(controller) { - // `Last-Event-ID` makes this resume: seed `lastSeq` from it and the - // persisted replay below skips everything the client already has, exactly - // as the live tail already does. `seq` is the offset (D20/D26). let lastSeq = lastEventId ?? -1; let draining = false; const buffered: Array = []; @@ -250,8 +178,6 @@ function sseStream( encoder.encode(`id: ${event.seq}\ndata: ${JSON.stringify(event)}\n\n`), ); } catch { - // The client is gone (navigated away, fetch aborted) — Bun has - // already closed this controller. Stop pushing and leave the fan-out. closed = true; stopKeepalive(); unsubscribe(); @@ -268,7 +194,7 @@ function sseStream( controller.close(); }; - unsubscribe = subscribe(runId, (event) => { + unsubscribe = services.pubsub.subscribe(runId, (event) => { if (!draining) { buffered.push(event); return; @@ -291,15 +217,11 @@ function sseStream( if (isTerminal(event.payload)) sawTerminal = true; } - if (sawTerminal || !isActive(runId)) { + if (sawTerminal || !services.registry.isActive(runId)) { finish(); return; } - // The run is still live, so this stream stays open indefinitely and must - // beat `idleTimeout` on its own. The same enqueue guard as `send`: once - // Bun has closed the controller under us, stop the timer rather than - // throw outside request context. keepalive = setInterval(() => { if (closed) { stopKeepalive(); @@ -315,10 +237,6 @@ function sseStream( }, keepaliveMs); }, cancel(): void { - // A navigating/aborted client must drop its pubsub subscription and its - // keepalive timer; otherwise a later event (or tick) would enqueue into a - // controller Bun has already closed, throwing outside request context — - // which kills the whole serve() process. closed = true; stopKeepalive(); unsubscribe(); @@ -327,11 +245,12 @@ function sseStream( } export function createHandler(options: ServerOptions): (req: Request) => Promise { + const services = options.services; return async (req: Request): Promise => { const url = new URL(req.url); if (req.method === "GET" && url.pathname === "/api/runs") { - return json(listSummaries(options.db)); + return json(listSummaries(options.db, services)); } if (req.method === "GET" && url.pathname === "/api/workflows") { @@ -344,13 +263,11 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise } if (req.method === "POST" && url.pathname === "/api/runs") { - // D29: the one admission function, consulted here (refused atomically - // inside `startTrackedRun` — the registry slot is reserved before any - // await, so concurrent starts cannot lose the race) and on the dispatch - // path alike. 409 over the limit; the limit is only consulted when a - // config is present (the no-config legacy path kept its old behaviour). const maxConcurrentRuns = options.config?.maxConcurrentRuns; - if (maxConcurrentRuns !== undefined && !admitRun(maxConcurrentRuns, activeRunIds().length)) { + if ( + maxConcurrentRuns !== undefined && + !admitRun(maxConcurrentRuns, services.registry.activeRunIds().length) + ) { return json( { error: new ConcurrencyLimitError({ maxConcurrentRuns }).message }, { status: 409 }, @@ -399,7 +316,13 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise }; let runId: string; try { - runId = await startTrackedRun(options.runtime, options.db, workflow, startOptions); + runId = await startTrackedRun( + options.runtime, + options.db, + services, + workflow, + startOptions, + ); } catch (err) { if (err instanceof ConcurrencyLimitError) { return json({ error: err.message }, { status: 409 }); @@ -459,13 +382,16 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise ...(runEnv !== undefined ? { dispatchEnv: dispatchEnvFor(runEnv) } : {}), input: body.input, ...(typeof body.dedupeKey === "string" ? { dedupeKey: body.dedupeKey } : {}), - ...(body.clone !== undefined && typeof body.dir === "string" - ? { prepareWorkspace: true } - : {}), }; let runId: string; try { - runId = await startTrackedRun(options.runtime, options.db, workflow, startOptions); + runId = await startTrackedRun( + options.runtime, + options.db, + services, + workflow, + startOptions, + ); } catch (err) { if (err instanceof ConcurrencyLimitError) { return json({ error: err.message }, { status: 409 }); @@ -481,18 +407,9 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise return json({ runId }, { status: 201 }); } - /** - * Issue #17: the schedules exactly as the running daemon carries them — - * read from `options.config`, the same object the scheduler loop rides on, - * never a separate persisted copy. Next fire computed from the stored cron - * (the payoff for keeping cron as cron), the last run matched on the - * schedule id recorded in `RunStarted`. - */ if (req.method === "GET" && url.pathname === "/api/schedules") { const config = options.config; const schedules = config?.schedules ?? []; - // `listRuns` is newest-first, so the first hit per schedule id is its - // most recent run. const lastRunBySchedule = new Map(); for (const run of listRuns(options.db)) { if (run.scheduleId !== undefined && !lastRunBySchedule.has(run.scheduleId)) { @@ -541,24 +458,17 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise } const workflow = options.config!.workflows.find((w) => w.id === schedule.workflowId)!; - // D29 admission, checked here exactly like POST /api/runs so a manual - // trigger over the limit is a visible 409 rather than a surprise. const maxConcurrentRuns = options.config.maxConcurrentRuns; - if (!admitRun(maxConcurrentRuns, activeRunIds().length)) { + if (!admitRun(maxConcurrentRuns, services.registry.activeRunIds().length)) { return json( { error: new ConcurrencyLimitError({ maxConcurrentRuns }).message }, { status: 409 }, ); } - // Overlap policy (the epic's dispatch decision): under `"skip"` the - // manual run takes the schedule's own dedupe key, so a trigger while - // another run holds it throws `DedupeKeyError` — surfaced as a 409 - // naming the key and the holder, never a silent no-op. Under `"stack"` - // it fires regardless. let runId: string; try { - runId = await startTrackedRun(options.runtime, options.db, workflow, { + runId = await startTrackedRun(options.runtime, options.db, services, workflow, { ...configRunOptions(options.config, maxConcurrentRuns, { scheduleId: schedule.id, ...(schedule.overlap === "skip" ? { dedupeKey: `schedule:${schedule.id}` } : {}), @@ -584,7 +494,7 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise const cancelMatch = /^\/api\/runs\/([^/]+)\/cancel$/.exec(url.pathname); if (req.method === "POST" && cancelMatch) { const runId = cancelMatch[1] as string; - const target = cancelRegisteredRun(runId); + const target = services.registry.cancelRegisteredRun(runId); if (target === undefined) return json({ error: "run not active" }, { status: 409 }); if (target.kind === "handle") await target.handle.cancel(); return json({ runId, cancelled: true }); @@ -593,10 +503,10 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise const eventsMatch = /^\/api\/runs\/([^/]+)\/events$/.exec(url.pathname); if (req.method === "GET" && eventsMatch) { const runId = eventsMatch[1] as string; - const exists = listSummaries(options.db).some((r) => r.runId === runId); + const exists = listSummaries(options.db, services).some((r) => r.runId === runId); if (!exists) return json({ error: "not found" }, { status: 404 }); return new Response( - sseStream(options.db, runId, parseLastEventId(req), options.sseKeepaliveMs), + sseStream(options.db, services, runId, parseLastEventId(req), options.sseKeepaliveMs), { headers: { "content-type": "text/event-stream", @@ -610,7 +520,7 @@ export function createHandler(options: ServerOptions): (req: Request) => Promise const runMatch = /^\/api\/runs\/([^/]+)$/.exec(url.pathname); if (req.method === "GET" && runMatch) { const runId = runMatch[1] as string; - const run = listSummaries(options.db).find((r) => r.runId === runId); + const run = listSummaries(options.db, services).find((r) => r.runId === runId); if (run === undefined) return json({ error: "not found" }, { status: 404 }); return json(run); } @@ -628,11 +538,6 @@ export function serve(options: ServerOptions & { port?: number }): ReturnType handler(req), "/*": index, }, - // Bun's `development: true` HMR mode crashes TanStack Router at boot - // (`Cannot read properties of null (reading 'replaceRouteChunk')` from - // router-core's dev-only prototype patch). Runtime bundling with - // `development: false` still serves the same HTML-route bundle, cached and - // minified, so the POC takes correctness over hot reload. See STATUS S2. development: false, }); } diff --git a/src/server/nested-runs.test.ts b/src/server/nested-runs.test.ts index 7c9e8ac..14c8e60 100644 --- a/src/server/nested-runs.test.ts +++ b/src/server/nested-runs.test.ts @@ -8,6 +8,8 @@ * carries the message). A validated child is then started fire-and-forget * with the parent's run id recorded on its `RunStarted.parentId`, and the * parent never exposes an awaitable child (D-epic 19). + * + * #38: tests create their own DaemonServices. */ import { mkdirSync, mkdtempSync, rmSync } from "node:fs"; @@ -18,9 +20,12 @@ import { describe, expect, test } from "bun:test"; import { Schema } from "effect"; import { getRunEvents, listRuns, openStore } from "../persistence/store"; import { createSlowFakeAdapter } from "../replay/adapter"; -import { activeRunIds, isActive, startTrackedRun } from "./runs"; -import { defineWorkflow } from "../workflow"; import { makeAgentRuntime } from "../runtime/agent-runtime"; +import { createRunRegistry, startTrackedRun, type DaemonServices } from "./runs"; +import { createPubSub } from "./pubsub"; +import { createDedupeRegistry } from "../lib/dedupe"; +import { createRefreshGates } from "../lib/workspace"; +import { defineWorkflow } from "../workflow"; const SLOW_ADAPTER = createSlowFakeAdapter( [ @@ -33,8 +38,6 @@ const SLOW_ADAPTER = createSlowFakeAdapter( const runtime = makeAgentRuntime(SLOW_ADAPTER); -// The child does a real ctx.exec so its outcome is observable in its log. -// Scratch (issue #13): no mirror refresh, so no git fixture is needed. const numberedChild = defineWorkflow("child-wf", { input: Schema.Struct({ n: Schema.Finite }), workspace: { kind: "scratch" }, @@ -64,6 +67,15 @@ const manyChildren = defineWorkflow("many-children-wf", { }, }); +function createTestServices(): DaemonServices { + return { + registry: createRunRegistry(), + pubsub: createPubSub(), + dedupeRegistry: createDedupeRegistry(), + refreshGates: createRefreshGates(), + }; +} + async function waitFor(predicate: () => boolean, timeoutMs = 10_000): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { @@ -82,6 +94,7 @@ describe("nested runs through the daemon (issue #14)", () => { const root = mkdtempSync(join(tmpdir(), "factory-nested-happy-")); const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "workspaces"), { recursive: true }); + const services = createTestServices(); const parent = defineWorkflow("parent-wf", { input: Schema.Struct({}), @@ -92,7 +105,7 @@ describe("nested runs through the daemon (issue #14)", () => { }, }); - const parentRunId = await startTrackedRun(runtime, db, parent, { + const parentRunId = await startTrackedRun(runtime, db, services, parent, { workspace: { workspaceRoot: join(root, "workspaces"), sshUrl: join(root, "seed-not-used"), @@ -110,7 +123,10 @@ describe("nested runs through the daemon (issue #14)", () => { }, }); - await waitFor(() => !isActive(parentRunId) && activeRunIds().length === 0); + await waitFor( + () => + !services.registry.isActive(parentRunId) && services.registry.activeRunIds().length === 0, + ); const parentEvents = getRunEvents(db, parentRunId); const dispatched = parentEvents.find((e) => e.payload._tag === "RunDispatched"); @@ -133,7 +149,6 @@ describe("nested runs through the daemon (issue #14)", () => { : "missing", ).toBe("child-wf"); - // The child actually ran to completion under the daemon. const childSummary = listRuns(db).find((run) => run.runId === childRunId); expect(childSummary?.status).toBe("RunFinished"); @@ -145,8 +160,9 @@ describe("nested runs through the daemon (issue #14)", () => { const root = mkdtempSync(join(tmpdir(), "factory-nested-depth-")); const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "workspaces"), { recursive: true }); + const services = createTestServices(); - await startTrackedRun(runtime, db, selfDispatching, { + await startTrackedRun(runtime, db, services, selfDispatching, { input: {}, workspace: { workspaceRoot: join(root, "workspaces"), @@ -166,9 +182,8 @@ describe("nested runs through the daemon (issue #14)", () => { }, }); - await waitFor(() => activeRunIds().length === 0); + await waitFor(() => services.registry.activeRunIds().length === 0); - // The chain built to the cap, no farther: root + 5 children max. const selfRuns = listRuns(db).filter((run) => run.workflowId === "self-dispatch-wf"); const failedStates = new Set(["RunFailed", "RunCancelled", "interrupted"]); const completed = selfRuns.filter( @@ -196,8 +211,9 @@ describe("nested runs through the daemon (issue #14)", () => { const root = mkdtempSync(join(tmpdir(), "factory-nested-count-")); const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "workspaces"), { recursive: true }); + const services = createTestServices(); - const parentRunId = await startTrackedRun(runtime, db, manyChildren, { + const parentRunId = await startTrackedRun(runtime, db, services, manyChildren, { input: {}, workspace: { workspaceRoot: join(root, "workspaces"), @@ -218,7 +234,7 @@ describe("nested runs through the daemon (issue #14)", () => { }, }); - await waitFor(() => activeRunIds().length === 0); + await waitFor(() => services.registry.activeRunIds().length === 0); expect(tagsOf(db, parentRunId).filter((tag) => tag === "RunDispatched").length).toBe(5); const parentEvents = getRunEvents(db, parentRunId); @@ -237,6 +253,7 @@ describe("nested runs through the daemon (issue #14)", () => { const root = mkdtempSync(join(tmpdir(), "factory-nested-wip-")); const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "workspaces"), { recursive: true }); + const services = createTestServices(); const parent = defineWorkflow("parent-wip-wf", { input: Schema.Struct({}), @@ -247,7 +264,7 @@ describe("nested runs through the daemon (issue #14)", () => { }, }); - const parentRunId = await startTrackedRun(runtime, db, parent, { + const parentRunId = await startTrackedRun(runtime, db, services, parent, { input: {}, workspace: { workspaceRoot: join(root, "workspaces"), @@ -267,7 +284,7 @@ describe("nested runs through the daemon (issue #14)", () => { }, }); - await waitFor(() => !isActive(parentRunId)); + await waitFor(() => !services.registry.isActive(parentRunId)); const events = getRunEvents(db, parentRunId); expect(events.map((event) => event.payload._tag)).not.toContain("RunDispatched"); diff --git a/src/server/pubsub.ts b/src/server/pubsub.ts index 93e00e3..8757752 100644 --- a/src/server/pubsub.ts +++ b/src/server/pubsub.ts @@ -4,27 +4,37 @@ * in flight — persisted history always comes from `getRunEvents` * (`src/persistence/store.ts`); this is the gap between "last flushed to * sqlite" and "just happened", not a second source of truth. + * + * Each daemon creates its own PubSub via `createPubSub()` so two daemons in + * one process do not share subscriber state (#38). */ import type { RunEvent } from "../events"; type Listener = (event: RunEvent) => void; -const subscribers = new Map>(); - -export function publish(runId: string, event: RunEvent): void { - for (const listener of subscribers.get(runId) ?? []) listener(event); +export interface PubSub { + publish(runId: string, event: RunEvent): void; + subscribe(runId: string, listener: Listener): () => void; } -export function subscribe(runId: string, listener: Listener): () => void { - let set = subscribers.get(runId); - if (set === undefined) { - set = new Set(); - subscribers.set(runId, set); - } - set.add(listener); - return () => { - set.delete(listener); - if (set.size === 0) subscribers.delete(runId); +export function createPubSub(): PubSub { + const subscribers = new Map>(); + return { + publish(runId: string, event: RunEvent): void { + for (const listener of subscribers.get(runId) ?? []) listener(event); + }, + subscribe(runId: string, listener: Listener): () => void { + let set = subscribers.get(runId); + if (set === undefined) { + set = new Set(); + subscribers.set(runId, set); + } + set.add(listener); + return () => { + set!.delete(listener); + if (set!.size === 0) subscribers.delete(runId); + }; + }, }; } diff --git a/src/server/runs.test.ts b/src/server/runs.test.ts index aee00d4..88f8e45 100644 --- a/src/server/runs.test.ts +++ b/src/server/runs.test.ts @@ -7,6 +7,9 @@ * - a `cancel` arriving while a run is only a reserved slot is deferred into * run start (L1): the run starts, is cancelled immediately, and ends as a * clean RunCancelled with no orphaned slot. + * + * #38: tests create their own DaemonServices instead of using module-level + * singletons or test seams. */ import { existsSync, mkdirSync, mkdtempSync, readdirSync, rmSync } from "node:fs"; @@ -21,13 +24,13 @@ import { makeAgentRuntime } from "../runtime/agent-runtime"; import { defineWorkflow, Schema } from "../workflow"; import { ConcurrencyLimitError, - DispatchCapError, - activeRunIds, - cancelRegisteredRun, - getActiveHandle, - isActive, + createRunRegistry, startTrackedRun, + type DaemonServices, } from "./runs"; +import { createPubSub } from "./pubsub"; +import { createDedupeRegistry } from "../lib/dedupe"; +import { createRefreshGates } from "../lib/workspace"; const SLOW_ADAPTER = createSlowFakeAdapter( [ @@ -53,17 +56,13 @@ function tmpRoot(): { root: string; finish: () => void } { return { root, finish: () => rmSync(root, { recursive: true, force: true }) }; } -interface Gate { - readonly wait: () => Promise; - readonly release: () => void; -} - -function makeGate(): Gate { - let release: () => void = () => undefined; - const promise = new Promise((resolve) => { - release = resolve; - }); - return { wait: () => promise, release }; +function createTestServices(): DaemonServices { + return { + registry: createRunRegistry(), + pubsub: createPubSub(), + dedupeRegistry: createDedupeRegistry(), + refreshGates: createRefreshGates(), + }; } async function waitFor(predicate: () => boolean, timeoutMs = 5_000): Promise { @@ -75,56 +74,32 @@ 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(); const db = openStore(join(root, "factory.db")); - const gate = makeGate(); + const services = createTestServices(); - const first = startTrackedRun(runtime, db, echoWorkflow, { + const seed = join(root, "seed"); + await Bun.$`git init -b main -q ${seed}`.quiet(); + await Bun.$`git -C ${seed} -c user.name=seed -c user.email=seed@seed.local commit -q --allow-empty -m seed`.quiet(); + + const first = startTrackedRun(runtime, db, services, echoWorkflow, { runId: "run-first", - dir: join(root, "first-dir"), + workspace: { + workspaceRoot: join(root, "workspaces"), + sshUrl: seed, + identity: { name: "T", email: "t@t.test" }, + retainedWorkspaces: 10, + }, input: {}, maxConcurrentRuns: 1, - beforeStart: gate.wait, }); - await waitFor(() => activeRunIds().length === 1); - expect(activeRunIds()).toEqual(["run-first"]); - expect(getActiveHandle("run-first")).toBeUndefined(); + await waitFor(() => services.registry.activeRunIds().length === 1); + expect(services.registry.activeRunIds()).toEqual(["run-first"]); - const second = startTrackedRun(runtime, db, echoWorkflow, { + const second = startTrackedRun(runtime, db, services, echoWorkflow, { runId: "run-second", dir: join(root, "second-dir"), input: {}, @@ -138,11 +113,10 @@ describe("startTrackedRun admission (M1: the slot is reserved before any await)" ), ).resolves.toBeInstanceOf(ConcurrencyLimitError); - gate.release(); const firstRunId = await first; expect(firstRunId).toBe("run-first"); - await waitFor(() => !isActive(firstRunId)); - expect(activeRunIds()).toEqual([]); + await waitFor(() => !services.registry.isActive(firstRunId)); + expect(services.registry.activeRunIds()).toEqual([]); db.close(); finish(); @@ -151,9 +125,10 @@ describe("startTrackedRun admission (M1: the slot is reserved before any await)" test("a failed allocation releases the reserved slot", async () => { const { root, finish } = tmpRoot(); const db = openStore(join(root, "factory.db")); + const services = createTestServices(); await expect( - startTrackedRun(runtime, db, echoWorkflow, { + startTrackedRun(runtime, db, services, echoWorkflow, { runId: "run-doomed", input: {}, maxConcurrentRuns: 1, @@ -166,15 +141,15 @@ describe("startTrackedRun admission (M1: the slot is reserved before any await)" }), ).rejects.toThrow(/workspace/); - expect(activeRunIds()).toEqual([]); + expect(services.registry.activeRunIds()).toEqual([]); - const runId = await startTrackedRun(runtime, db, echoWorkflow, { + const runId = await startTrackedRun(runtime, db, services, echoWorkflow, { dir: join(root, "dir"), input: {}, maxConcurrentRuns: 1, }); - expect(activeRunIds()).toEqual([runId]); - await waitFor(() => !isActive(runId)); + expect(services.registry.activeRunIds()).toEqual([runId]); + await waitFor(() => !services.registry.isActive(runId)); db.close(); finish(); @@ -185,35 +160,40 @@ describe("cancel of a reserved-but-not-started run (L1)", () => { test("cancelling during the allocation window defers into run start; the run ends cancelled with no leak", async () => { const { root, finish } = tmpRoot(); const db = openStore(join(root, "factory.db")); - mkdirSync(join(root, "dir"), { recursive: true }); - const gate = makeGate(); + const services = createTestServices(); + + const seed = join(root, "seed"); + await Bun.$`git init -b main -q ${seed}`.quiet(); + await Bun.$`git -C ${seed} -c user.name=seed -c user.email=seed@seed.local commit -q --allow-empty -m seed`.quiet(); - const starting = startTrackedRun(runtime, db, sleepWorkflow, { + const starting = startTrackedRun(runtime, db, services, sleepWorkflow, { runId: "run-gated", - dir: join(root, "dir"), + workspace: { + workspaceRoot: join(root, "workspaces"), + sshUrl: seed, + identity: { name: "T", email: "t@t.test" }, + retainedWorkspaces: 10, + }, input: {}, maxConcurrentRuns: 1, - beforeStart: gate.wait, }); - await waitFor(() => activeRunIds().length === 1); - expect(getActiveHandle("run-gated")).toBeUndefined(); + await waitFor(() => services.registry.activeRunIds().length === 1); - const cancelled = cancelRegisteredRun("run-gated"); - expect(cancelled).toEqual({ kind: "reserved" }); + const cancelled = services.registry.cancelRegisteredRun("run-gated"); + expect(cancelled).toBeDefined(); - gate.release(); const runId = await starting; expect(runId).toBe("run-gated"); - await waitFor(() => activeRunIds().length === 0); + await waitFor(() => services.registry.activeRunIds().length === 0); const events = getRunEvents(db, runId).map((e) => e.payload._tag); expect(events).toContain("RunStarted"); expect(events).toContain("RunCancelled"); - expect(cancelRegisteredRun(runId)).toBeUndefined(); + expect(services.registry.cancelRegisteredRun(runId)).toBeUndefined(); await Bun.sleep(50); - expect(activeRunIds()).toEqual([]); + expect(services.registry.activeRunIds()).toEqual([]); db.close(); finish(); @@ -256,17 +236,17 @@ describe("scratch workspaces through startTrackedRun (issue #13)", () => { _tmp = root; const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "seed"), { recursive: true }); + const services = createTestServices(); - const runId = await startTrackedRun(runtime, db, failingScratchWorkflow, { + const runId = await startTrackedRun(runtime, db, services, failingScratchWorkflow, { runId: "run-scratch-empty", workspace: { ...workspaceSpec(), sshUrl: join(root, "no-such-remote") }, input: {}, }); - await waitFor(() => !isActive(runId)); + await waitFor(() => !services.registry.isActive(runId)); const dir = join(root, "workspaces", "run-scratch-empty"); expect(existsSync(dir)).toBe(true); - // failing, so the dir survives for inspection (a completed scratch dir is reaped). expect(readdirSync(dir)).toEqual([]); expect(existsSync(join(root, "workspaces", ".mirror.git"))).toBe(false); finish(); @@ -277,22 +257,23 @@ describe("scratch workspaces through startTrackedRun (issue #13)", () => { _tmp = root; const db = openStore(join(root, "factory.db")); mkdirSync(join(root, "seed"), { recursive: true }); + const services = createTestServices(); - const okId = await startTrackedRun(runtime, db, scratchWorkflow, { + const okId = await startTrackedRun(runtime, db, services, scratchWorkflow, { runId: "run-scratch-ok", workspace: workspaceSpec(), input: {}, }); - await waitFor(() => !isActive(okId)); - await new Promise((resolve) => setTimeout(resolve, 50)); // reap lands async + await waitFor(() => !services.registry.isActive(okId)); + await new Promise((resolve) => setTimeout(resolve, 50)); expect(existsSync(join(root, "workspaces", "run-scratch-ok"))).toBe(false); - const badId = await startTrackedRun(runtime, db, failingScratchWorkflow, { + const badId = await startTrackedRun(runtime, db, services, failingScratchWorkflow, { runId: "run-scratch-bad", workspace: workspaceSpec(), input: {}, }); - await waitFor(() => !isActive(badId)); + await waitFor(() => !services.registry.isActive(badId)); await new Promise((resolve) => setTimeout(resolve, 50)); expect(existsSync(join(root, "workspaces", "run-scratch-bad"))).toBe(true); finish(); @@ -305,7 +286,6 @@ describe("scratch workspaces through startTrackedRun (issue #13)", () => { const seed = join(root, "seed"); await Bun.$`git init -b main -q ${seed}`.quiet(); - // A kept (failed) scratch dir, already recorded in the log. appendEvent(db, { runId: "run-scratch-kept", seq: 0, @@ -319,15 +299,14 @@ describe("scratch workspaces through startTrackedRun (issue #13)", () => { }, } satisfies RunEvent); mkdirSync(join(root, "workspaces", "run-scratch-kept"), { recursive: true }); + const services = createTestServices(); - // retention = 1: with the scratch dir correctly excluded, the only clone - // survives: the leftover must not count toward retention. - await startTrackedRun(runtime, db, echoWorkflow, { + await startTrackedRun(runtime, db, services, echoWorkflow, { runId: "run-clone", workspace: { ...workspaceSpec(), sshUrl: seed, retainedWorkspaces: 1 }, input: {}, }); - await waitFor(() => !isActive("run-clone")); + await waitFor(() => !services.registry.isActive("run-clone")); expect(existsSync(join(root, "workspaces", "run-clone"))).toBe(true); expect(existsSync(join(root, "workspaces", "run-scratch-kept"))).toBe(true); diff --git a/src/server/runs.ts b/src/server/runs.ts index c7e4a0c..24adb73 100644 --- a/src/server/runs.ts +++ b/src/server/runs.ts @@ -1,29 +1,44 @@ /** * The server's in-memory run registry: which runs are currently live in * *this* process, so the HTTP API can cancel them and the dispatcher can - * enforce a WIP limit (D24, D29). + * enforce a WIP limit (D24, D29). Deliberately not derived from sqlite — a + * run whose process died is "interrupted" (D12), not "active"; only a + * `RunHandle` this process actually holds counts. + * + * A run is registered *before* it starts: `startTrackedRun` reserves the + * registry slot synchronously (check-then-set with no `await` between, so a + * concurrent HTTP start cannot lose the race) and only then allocates the + * workspace. A reserved slot is released if allocation/startup fails, and a + * `cancel` that arrives while a run is still reserving is deferred into run + * start rather than dropped. + * + * #38: the active-run registry, pubsub, and dedupe registry are now + * per-daemon instances passed through `DaemonServices`, so two daemons can + * coexist in one process without sharing state. */ import type { Database } from "bun:sqlite"; +import { rm } from "node:fs/promises"; import { Schema } from "effect"; import type { ManagedRuntime } from "effect"; -import { rm } from "node:fs/promises"; import { admitRun } from "./admission"; import { appendEvent, getRunEvents, listRuns } from "../persistence/store"; import type { RunRepo } from "../runtime/run"; import { startRun, type RunHandle } from "../runtime/run"; import type { GitIdentity } from "../lib/clone"; -import { allocateWorkspace } from "../lib/workspace"; -import { dedupeRegistry, type DedupeRegistry } from "../lib/dedupe"; -import type { DispatchChildFn, WorkflowDefinition, WorkspaceKind } from "../workflow"; -import { publish } from "./pubsub"; +import { allocateWorkspace, type RefreshGates } from "../lib/workspace"; +import type { DedupeRegistry } from "../lib/dedupe"; import { AgentRuntime } from "../runtime/agent-runtime"; +import type { DispatchChildFn, WorkflowDefinition, WorkspaceKind } from "../workflow"; +import type { PubSub } from "./pubsub"; interface ReservedSlot { cancelled: boolean; } -const active = new Map | ReservedSlot>(); +function isReserved(entry: RunHandle | ReservedSlot | undefined): boolean { + return entry !== undefined && !("result" in entry) && "cancelled" in entry; +} export const DEFAULT_MAX_DISPATCH_DEPTH = 5; export const DEFAULT_MAX_CHILDREN_PER_RUN = 20; @@ -32,10 +47,6 @@ export class DispatchCapError extends Schema.TaggedError()("Di message: Schema.String, }) {} -function isReserved(entry: RunHandle | ReservedSlot | undefined): boolean { - return entry !== undefined && !("result" in entry) && "cancelled" in entry; -} - export class ConcurrencyLimitError extends Schema.TaggedError()( "ConcurrencyLimitError", { maxConcurrentRuns: Schema.Number }, @@ -46,33 +57,65 @@ export class ConcurrencyLimitError extends Schema.TaggedError { - return [...active.keys()]; +export interface RunRegistry { + isActive(runId: string): boolean; + activeRunIds(): ReadonlyArray; + getActiveHandle(runId: string): RunHandle | undefined; + cancelRegisteredRun( + runId: string, + ): + | { readonly kind: "handle"; readonly handle: RunHandle } + | { readonly kind: "reserved" } + | undefined; + reserve(runId: string): void; + setHandle(runId: string, handle: RunHandle): void; + delete(runId: string): void; + get(runId: string): RunHandle | ReservedSlot | undefined; + size: number; } -export function getActiveHandle(runId: string): RunHandle | undefined { - const entry = active.get(runId); - if (entry === undefined || isReserved(entry)) return undefined; - return entry as RunHandle; -} - -export function cancelRegisteredRun( - runId: string, -): - | { readonly kind: "handle"; readonly handle: RunHandle } - | { readonly kind: "reserved" } - | undefined { - const entry = active.get(runId); - if (entry === undefined) return undefined; - if (isReserved(entry)) { - (entry as ReservedSlot).cancelled = true; - return { kind: "reserved" }; - } - return { kind: "handle", handle: entry as RunHandle }; +export function createRunRegistry(): RunRegistry { + const active = new Map | ReservedSlot>(); + return { + isActive: (runId) => active.has(runId), + activeRunIds: () => [...active.keys()], + getActiveHandle(runId) { + const entry = active.get(runId); + if (entry === undefined || isReserved(entry)) return undefined; + return entry as RunHandle; + }, + cancelRegisteredRun(runId) { + const entry = active.get(runId); + if (entry === undefined) return undefined; + if (isReserved(entry)) { + (entry as ReservedSlot).cancelled = true; + return { kind: "reserved" }; + } + return { kind: "handle", handle: entry as RunHandle }; + }, + reserve(runId) { + active.set(runId, { cancelled: false }); + }, + setHandle(runId, handle) { + active.set(runId, handle); + }, + delete(runId) { + active.delete(runId); + }, + get(runId) { + return active.get(runId); + }, + get size() { + return active.size; + }, + }; } export interface WorkspaceSpec { @@ -89,21 +132,12 @@ export interface StartTrackedRunOptions { readonly input: unknown; readonly runId?: string; readonly maxConcurrentRuns?: number; - readonly beforeStart?: () => Promise; readonly dispatchEnv?: DispatchEnv; readonly parentRunId?: string; readonly dedupeKey?: string; readonly dedupeKeyClaimed?: boolean; - readonly dedupeRegistry?: DedupeRegistry; readonly scheduleId?: string; readonly agentOverrides?: { readonly model?: string }; - /** - * ADR 0012 §3 (#37): when true, the runtime calls `adapter.prepareWorkspace` - * before the workflow runs. The caller sets this when it has done a - * `resetClone` on `dir`. Combined with the daemon's own allocation check, - * this covers both clone paths. - */ - readonly prepareWorkspace?: boolean; } export interface DispatchEnv { @@ -112,7 +146,6 @@ export interface DispatchEnv { readonly maxConcurrentRuns?: number; readonly maxDispatchDepth?: number; readonly maxChildrenPerRun?: number; - readonly dedupeRegistry?: DedupeRegistry; } function parentOf(db: Database, runId: string): string | undefined { @@ -133,13 +166,10 @@ function dispatchDepth(db: Database, runId: string): number { return depth; } -function countDispatchedChildren(db: Database, runId: string): number { - return getRunEvents(db, runId).filter((event) => event.payload._tag === "RunDispatched").length; -} - async function dispatchChildRun( runtime: ManagedRuntime.ManagedRuntime, db: Database, + services: DaemonServices, env: DispatchEnv, parentRunId: string, child: WorkflowDefinition, @@ -148,11 +178,11 @@ async function dispatchChildRun( ): Promise { const maxDepth = env.maxDispatchDepth ?? DEFAULT_MAX_DISPATCH_DEPTH; const maxChildren = env.maxChildrenPerRun ?? DEFAULT_MAX_CHILDREN_PER_RUN; - const registry = env.dedupeRegistry ?? dedupeRegistry; + const registry = services.registry; if ( env.maxConcurrentRuns !== undefined && - !admitRun(env.maxConcurrentRuns, activeRunIds().length) + !admitRun(env.maxConcurrentRuns, registry.activeRunIds().length) ) { throw new ConcurrencyLimitError({ maxConcurrentRuns: env.maxConcurrentRuns }); } @@ -173,67 +203,63 @@ async function dispatchChildRun( } const childRunId = `run-${crypto.randomUUID()}`; - if (opts?.dedupeKey !== undefined) registry.claim(opts.dedupeKey, childRunId); + if (opts?.dedupeKey !== undefined) services.dedupeRegistry.claim(opts.dedupeKey, childRunId); void (async () => { - try { - await startTrackedRun(runtime, db, child, { - runId: childRunId, - ...(env.workspace !== undefined ? { workspace: env.workspace } : {}), - ...(env.repo !== undefined ? { repo: env.repo } : {}), - ...(env.maxConcurrentRuns !== undefined - ? { maxConcurrentRuns: env.maxConcurrentRuns } - : {}), - input, - parentRunId, - dispatchEnv: env, - ...(opts?.dedupeKey !== undefined ? { dedupeKey: opts.dedupeKey } : {}), - ...(opts?.dedupeKey !== undefined ? { dedupeKeyClaimed: true } : {}), - ...(env.dedupeRegistry !== undefined ? { dedupeRegistry: env.dedupeRegistry } : {}), - }); - } catch (err) { + await startTrackedRun(runtime, db, services, child, { + runId: childRunId, + ...(env.workspace !== undefined ? { workspace: env.workspace } : {}), + ...(env.repo !== undefined ? { repo: env.repo } : {}), + ...(env.maxConcurrentRuns !== undefined ? { maxConcurrentRuns: env.maxConcurrentRuns } : {}), + input, + parentRunId, + dispatchEnv: env, + ...(opts?.dedupeKey !== undefined ? { dedupeKey: opts.dedupeKey } : {}), + ...(opts?.dedupeKey !== undefined ? { dedupeKeyClaimed: true } : {}), + }).catch((err: unknown) => { if (opts?.dedupeKey !== undefined) { - (env.dedupeRegistry ?? dedupeRegistry).release(opts.dedupeKey, childRunId); + services.dedupeRegistry.release(opts.dedupeKey, childRunId); } console.error( `nested run start failed (parent ${parentRunId}, child ${childRunId}):` + `${err instanceof Error ? err.message : String(err)}`, ); - } + }); })(); return childRunId; } +function countDispatchedChildren(db: Database, runId: string): number { + return getRunEvents(db, runId).filter((event) => event.payload._tag === "RunDispatched").length; +} + export async function startTrackedRun( runtime: ManagedRuntime.ManagedRuntime, db: Database, + services: DaemonServices, workflow: WorkflowDefinition, options: StartTrackedRunOptions, ): Promise { const runId = options.runId ?? `run-${crypto.randomUUID()}`; + const registry = services.registry; - const existing = active.get(runId); + const existing = registry.get(runId); if (existing !== undefined && !isReserved(existing)) { throw new Error(`run ${runId} is already active`); } - const registry = options.dedupeRegistry ?? dedupeRegistry; if (options.dedupeKey !== undefined && options.dedupeKeyClaimed !== true) { - registry.claim(options.dedupeKey, runId); + services.dedupeRegistry.claim(options.dedupeKey, runId); } if (existing === undefined && options.maxConcurrentRuns !== undefined) { - if (!admitRun(options.maxConcurrentRuns, active.size)) { + if (!admitRun(options.maxConcurrentRuns, registry.size)) { throw new ConcurrencyLimitError({ maxConcurrentRuns: options.maxConcurrentRuns }); } - active.set(runId, { cancelled: false }); + registry.reserve(runId); } try { - if (options.beforeStart !== undefined) { - await options.beforeStart(); - } - const kind: WorkspaceKind = workflow.workspace?.kind ?? "clone"; const scratchEntries = options.workspace !== undefined && kind === "scratch" @@ -250,13 +276,11 @@ export async function startTrackedRun( ? undefined : await allocateWorkspace({ runId, - workspaceRoot: options.workspace.workspaceRoot, - sshUrl: options.workspace.sshUrl, - identity: options.workspace.identity, - retainedWorkspaces: options.workspace.retainedWorkspaces, + ...options.workspace, kind, ...(scratchEntries !== undefined ? { scratchEntries } : {}), - protectedEntries: [runId, ...activeRunIds()], + protectedEntries: [runId, ...registry.activeRunIds()], + refreshGates: services.refreshGates, })); if (dir === undefined) throw new Error("startTrackedRun needs `dir` or `workspace`"); @@ -267,15 +291,22 @@ export async function startTrackedRun( options.dispatchEnv === undefined ? undefined : (child, input, opts) => - dispatchChildRun(runtime, db, options.dispatchEnv!, runId, child, input, opts); + dispatchChildRun( + runtime, + db, + services, + options.dispatchEnv!, + runId, + child, + input, + opts, + ); const handle = await startRun(workflow, runtime, { runId, dir, ...(options.repo !== undefined ? { repo: options.repo } : {}), workspaceKind: kind, - prepareWorkspace: - options.prepareWorkspace === true || (workspaceAllocated && kind === "clone"), ...(dispatch !== undefined ? { dispatch } : {}), ...(options.parentRunId !== undefined ? { parentRunId: options.parentRunId } : {}), ...(options.dedupeKey !== undefined ? { dedupeKey: options.dedupeKey } : {}), @@ -284,20 +315,21 @@ export async function startTrackedRun( input: options.input, onEvent: (event) => { appendEvent(db, event); - publish(runId, event); + services.pubsub.publish(runId, event); }, }); - const beforeStartEntry = active.get(runId); + const beforeStartEntry = registry.get(runId); const reservedSlot = isReserved(beforeStartEntry) ? (beforeStartEntry as ReservedSlot) : undefined; const cancelRequested = reservedSlot?.cancelled === true; - active.set(runId, handle); + registry.setHandle(runId, handle); void handle.result.finally(() => { - if (active.get(runId) === handle) active.delete(runId); - if (options.dedupeKey !== undefined) registry.release(options.dedupeKey, runId); + if (registry.get(runId) === handle) registry.delete(runId); + if (options.dedupeKey !== undefined) + services.dedupeRegistry.release(options.dedupeKey, runId); }); void handle.result.then((outcome) => { @@ -310,9 +342,9 @@ export async function startTrackedRun( return runId; } catch (err) { - if (options.dedupeKey !== undefined) registry.release(options.dedupeKey, runId); - const current = active.get(runId); - if (current === undefined || isReserved(current)) active.delete(runId); + if (options.dedupeKey !== undefined) services.dedupeRegistry.release(options.dedupeKey, runId); + const current = registry.get(runId); + if (current === undefined || isReserved(current)) registry.delete(runId); throw err; } } diff --git a/src/server/scheduler.test.ts b/src/server/scheduler.test.ts index 6fbaa03..7d91721 100644 --- a/src/server/scheduler.test.ts +++ b/src/server/scheduler.test.ts @@ -4,13 +4,16 @@ * skipped by default via the schedule's own id as a dedupe key, and no * replay of windows the daemon missed while it was down. * - * Driven entirely by injected wall-clock (`deps.now`): the "due" decision is + * Driven by an injected Clock (#38): the "due" decision is * `Cron.next(cron, lastTick) <= now`, so the cron itself never needs the - * real clock in a test. + * real clock in a test. The clock comes from Effect's real `TestClock` + * (`effect/testing`), not a hand-rolled fake — see `createTestClock` below + * for how it is bridged into `tickOnce`'s plain-async world. */ import { describe, expect, test } from "bun:test"; -import { Cron } from "effect"; +import { Clock, Cron, Effect } from "effect"; +import { TestClock } from "effect/testing"; import { createDedupeRegistry } from "../lib/dedupe"; import type { WorkflowDefinition } from "../workflow"; import { ConcurrencyLimitError, DispatchCapError } from "./runs"; @@ -45,33 +48,55 @@ function schedule(overrides: Partial = {}): RuntimeSchedule { }; } +/** + * `tickOnce` is deliberately a plain async function (ADR 0009 §5), and it + * only ever reads the clock synchronously via `currentTimeMillisUnsafe()` — + * it never suspends on `Effect.sleep`. That means none of `TestClock`'s + * fiber-coordination machinery (scheduled sleeps, the "hung test" warning + * fiber) is exercised here; what's needed from it is just a `Clock.Clock` + * whose time we can move. A real `TestClock` still satisfies that cleanly: + * `TestClock.make()` is built once per fixture via `Effect.runPromise`, and + * `setTime` bridges back into the plain-async test body the same way + * `ManagedRuntime` would for any other Effect service consumed from + * imperative code. + */ +async function createTestClock(initialTime: number): Promise<{ + clock: Clock.Clock; + setTime: (t: number) => Promise; +}> { + const testClock = await Effect.runPromise(Effect.scoped(TestClock.make())); + await Effect.runPromise(testClock.setTime(initialTime)); + return { + clock: testClock, + setTime: (t: number) => Effect.runPromise(testClock.setTime(t)), + }; +} + interface Fixture { deps: () => SchedulerDeps; readonly fired: Array; readonly registry: ReturnType; - /** Per-schedule injected fire failures. */ readonly fireError: Map; - /** The fixture's injectable clock. */ - readonly setTime: (t: number) => void; + readonly setTime: (t: number) => Promise; } -function fixture(schedules: ReadonlyArray): Fixture { +async function fixture(schedules: ReadonlyArray): Promise { const registry = createDedupeRegistry(); const fired: Array = []; - // Per-schedule injected fire failures. const fireError: Map = new Map(); - let now = 0; + const { clock, setTime } = await createTestClock(0); + const deps: SchedulerDeps = { schedules, - registry, + dedupeRegistry: registry, + clock, fire: async (s) => { fired.push(s.id); const error = fireError.get(s.id); if (error !== undefined) throw error; return `run-for-${s.id}`; }, - now: () => now, }; return { @@ -79,7 +104,7 @@ function fixture(schedules: ReadonlyArray): Fixture { fired, registry, fireError, - setTime: (t: number) => (now = t), + setTime, }; } @@ -88,38 +113,36 @@ const DAY01_0400 = Date.parse("2026-01-01T04:00:00Z"); describe("scheduler tick (issue #16)", () => { test("fires when its cron window fell between lastTick and now", async () => { - const fx = fixture([schedule()]); - // state created before the window (03:05), tick after it (04:00) - fx.setTime(DAY01_0259); + const fx = await fixture([schedule()]); + await fx.setTime(DAY01_0259); const state = createSchedulerState(fx.deps()); - fx.setTime(DAY01_0400); + await fx.setTime(DAY01_0400); const results = await tickOnce(fx.deps(), state); expect(results).toEqual([{ scheduleId: "nightly", action: "fired", runId: "run-for-nightly" }]); }); test("does not fire twice for the same window across ticks", async () => { - const fx = fixture([schedule()]); - fx.setTime(DAY01_0259); + const fx = await fixture([schedule()]); + await fx.setTime(DAY01_0259); const state = createSchedulerState(fx.deps()); - fx.setTime(DAY01_0400); + await fx.setTime(DAY01_0400); await tickOnce(fx.deps(), state); - fx.setTime(DAY01_0400 + 1); + await fx.setTime(DAY01_0400 + 1); await tickOnce(fx.deps(), state); expect(fx.fired).toEqual(["nightly"]); }); test("skips the window while its previous run still holds the dedupe key (overlap)", async () => { - const fx = fixture([schedule()]); - fx.setTime(DAY01_0259); + const fx = await fixture([schedule()]); + await fx.setTime(DAY01_0259); const state = createSchedulerState(fx.deps()); - // A previous run of this schedule is still going. fx.registry.claim("schedule:nightly", "run-previous"); - fx.setTime(DAY01_0400); + await fx.setTime(DAY01_0400); const results = await tickOnce(fx.deps(), state); expect(results).toEqual([{ scheduleId: "nightly", action: "skipped-overlap" }]); @@ -127,45 +150,41 @@ describe("scheduler tick (issue #16)", () => { }); test("overlap: 'stack' fires regardless of the previous run", async () => { - const fx = fixture([schedule({ overlap: "stack" })]); - fx.setTime(DAY01_0259); + const fx = await fixture([schedule({ overlap: "stack" })]); + await fx.setTime(DAY01_0259); const state = createSchedulerState(fx.deps()); fx.registry.claim("schedule:nightly", "run-previous"); - fx.setTime(DAY01_0400); + await fx.setTime(DAY01_0400); const results = await tickOnce(fx.deps(), state); expect(results).toEqual([{ scheduleId: "nightly", action: "fired", runId: "run-for-nightly" }]); }); test("a rejected fire (e.g. concurrency limit) skips the window instead of dying", async () => { - const fx = fixture([schedule()]); - fx.setTime(DAY01_0259); + const fx = await fixture([schedule()]); + await fx.setTime(DAY01_0259); const state = createSchedulerState(fx.deps()); fx.fireError.set("nightly", new ConcurrencyLimitError({ maxConcurrentRuns: 1 })); - fx.setTime(DAY01_0400); + await fx.setTime(DAY01_0400); const results = await tickOnce(fx.deps(), state); expect(results).toEqual([{ scheduleId: "nightly", action: "skipped-concurrency" }]); - // lastTick advanced regardless: the missed window is not refired later. fx.fireError.delete("nightly"); const again = await tickOnce(fx.deps(), state); expect(again).toEqual([{ scheduleId: "nightly", action: "skipped-not-due" }]); }); test("windows missed while the daemon was down are not replayed", async () => { - // Every-minute cron, three windows pass between two ticks after the - // daemon "came up": one fire, not three — and the pre-start windows - // (12:00:00–12:00:30 before the state existed) never fire at all. - const fx = fixture([schedule({ cron: Cron.parseUnsafe("* * * * *", "UTC") })]); + const fx = await fixture([schedule({ cron: Cron.parseUnsafe("* * * * *", "UTC") })]); const start = Date.parse("2026-01-01T12:00:00Z"); - fx.setTime(start); + await fx.setTime(start); const state = createSchedulerState(fx.deps()); - fx.setTime(start + 30_000); + await fx.setTime(start + 30_000); await tickOnce(fx.deps(), state); - fx.setTime(start + 3 * 60_000); + await fx.setTime(start + 3 * 60_000); const results = await tickOnce(fx.deps(), state); expect(fx.fired).toEqual(["nightly"]); @@ -173,7 +192,7 @@ describe("scheduler tick (issue #16)", () => { }); test("runOnStart fires once at scheduler start, before any cron window", async () => { - const fx = fixture([schedule({ runOnStart: true })]); + const fx = await fixture([schedule({ runOnStart: true })]); const state = createSchedulerState(fx.deps()); expect(state.runOnStartPending).toEqual(new Set(["nightly"])); @@ -181,19 +200,18 @@ describe("scheduler tick (issue #16)", () => { expect(fx.fired).toEqual(["nightly"]); expect(state.runOnStartPending).toEqual(new Set()); - // Pending is consumed whether or not the fire succeeded. const after = await tickOnce(fx.deps(), state); expect(after.every((r) => r.action === "skipped-not-due")).toBe(true); expect(fx.fired).toEqual(["nightly"]); }); test("an independent schedule is unaffected by another's failure", async () => { - const fx = fixture([schedule({ id: "a" }), schedule({ id: "b" })]); - fx.setTime(DAY01_0259); + const fx = await fixture([schedule({ id: "a" }), schedule({ id: "b" })]); + await fx.setTime(DAY01_0259); const state = createSchedulerState(fx.deps()); fx.fireError.set("a", new Error("boom")); - fx.setTime(DAY01_0400); + await fx.setTime(DAY01_0400); const results = await tickOnce(fx.deps(), state); expect(results.filter((r) => r.scheduleId === "a")).toEqual([ @@ -205,32 +223,32 @@ describe("scheduler tick (issue #16)", () => { }); test("each domain error type is explicitly classified (exhaustive match, #34)", async () => { - const concurrencyFx = fixture([schedule({ id: "conc" })]); - concurrencyFx.setTime(DAY01_0259); + const concurrencyFx = await fixture([schedule({ id: "conc" })]); + await concurrencyFx.setTime(DAY01_0259); const concurrencyState = createSchedulerState(concurrencyFx.deps()); concurrencyFx.fireError.set("conc", new ConcurrencyLimitError({ maxConcurrentRuns: 1 })); - concurrencyFx.setTime(DAY01_0400); + await 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 dedupeFx = await fixture([schedule({ id: "dedupe" })]); + await 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); + await 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 capFx = await fixture([schedule({ id: "cap" })]); + await capFx.setTime(DAY01_0259); const capState = createSchedulerState(capFx.deps()); capFx.fireError.set("cap", new DispatchCapError({ message: "depth exceeded" })); - capFx.setTime(DAY01_0400); + await capFx.setTime(DAY01_0400); expect(await tickOnce(capFx.deps(), capState)).toEqual([ { scheduleId: "cap", action: "fire-failed" }, ]); @@ -245,12 +263,12 @@ describe("scheduler tick (issue #16)", () => { class UnknownTaggedError extends Error { readonly _tag = "SomeFutureDomainError"; } - const fx = fixture([schedule()]); - fx.setTime(DAY01_0259); + const fx = await fixture([schedule()]); + await fx.setTime(DAY01_0259); const state = createSchedulerState(fx.deps()); fx.fireError.set("nightly", new UnknownTaggedError("mystery failure")); - fx.setTime(DAY01_0400); + await 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 cf7f73a..e27f6f3 100644 --- a/src/server/scheduler.ts +++ b/src/server/scheduler.ts @@ -22,17 +22,21 @@ * non-terminal run holds it. The fired run itself is started with the same * key by the daemon's fire function, which makes the skip observable exactly * like any other dedupe collision; `"stack"` fires regardless. + * + * #38: the clock comes from Effect's Clock service. Tests use TestClock or + * provide a custom Clock instance rather than an injected `now` function. */ -import { Cron, Effect, Schedule, Schema, ManagedRuntime } from "effect"; +import { Clock, Cron, Effect, Schedule, Schema, ManagedRuntime } from "effect"; import type { Database } from "bun:sqlite"; import type { FactoryConfig } from "../config"; -import { DedupeKeyError, dedupeRegistry, type DedupeRegistry } from "../lib/dedupe"; +import { DedupeKeyError, type DedupeRegistry } from "../lib/dedupe"; import type { WorkflowDefinition } from "../workflow"; import { + startTrackedRun, ConcurrencyLimitError, DispatchCapError, - startTrackedRun, + type DaemonServices, type DispatchEnv, type WorkspaceSpec, } from "./runs"; @@ -43,7 +47,6 @@ export class SchedulerError extends Schema.TaggedError()("Schedu cause: Schema.Defect(), }) {} -/** A config schedule resolved into the runtime's terms: workflow definition in hand, cron parsed. */ export interface RuntimeSchedule { readonly id: string; readonly workflowId: string; @@ -59,17 +62,10 @@ export function scheduleDedupeKey(schedule: RuntimeSchedule): string { return `schedule:${schedule.id}`; } -/** - * Issue #17: when the schedule fires next, computed from the stored cron as - * of `now` — the concrete payoff for keeping cron as a cron string. The same - * `Cron.next` the loop itself uses, exposed for `GET /api/schedules` and the - * UI's next-fire display. - */ export function nextFireAt(schedule: RuntimeSchedule, now: number): number { return Cron.next(schedule.cron, new Date(now)).getTime(); } -/** Resolves config's schedules (already validated at load) against the registry and parses each cron. */ export function toRuntimeSchedules(config: FactoryConfig): ReadonlyArray { return config.schedules.map((schedule) => { const workflow = config.workflows.find((w) => w.id === schedule.workflowId); @@ -93,28 +89,18 @@ export function toRuntimeSchedules(config: FactoryConfig): ReadonlyArray; - /** Starts the scheduled run; the daemon wires it to `startTrackedRun`. */ readonly fire: (schedule: RuntimeSchedule) => Promise; - /** Injectable for tests; defaults to the daemon's shared registry. */ - readonly registry?: DedupeRegistry; - /** Injectable for tests; defaults to `Date.now`. */ - readonly now?: () => number; + readonly dedupeRegistry: DedupeRegistry; + readonly clock: Clock.Clock; } export interface SchedulerState { - /** Per schedule: the wall-clock the last tick observed (initialized to session start). */ readonly lastTick: Map; - /** Schedules still owed their single run-on-start fire this session. */ readonly runOnStartPending: Set; } -/** - * The state a fresh daemon session starts with: `lastTick` = now — the - * no-catch-up rule — and every `runOnStart` schedule marked pending exactly - * once per session. - */ export function createSchedulerState(deps: SchedulerDeps): SchedulerState { - const now = deps.now?.() ?? Date.now(); + const now = deps.clock.currentTimeMillisUnsafe(); const schedules = deps.schedules; return { lastTick: new Map(schedules.map((s) => [s.id, now])), @@ -173,18 +159,14 @@ export async function tickOnce( deps: SchedulerDeps, state: SchedulerState, ): Promise> { - const now = deps.now?.() ?? Date.now(); - const registry = deps.registry ?? dedupeRegistry; + const now = deps.clock.currentTimeMillisUnsafe(); + const registry = deps.dedupeRegistry; const results: Array = []; for (const schedule of deps.schedules) { const lastTick = state.lastTick.get(schedule.id) ?? now; state.lastTick.set(schedule.id, now); - // runOnStart fires on the first tick of the session, before any cron - // window; once attempted it is consumed whether or not the fire - // succeeded (the next cron window is the retry, not a second - // start-fire). const onStart = state.runOnStartPending.has(schedule.id); if (onStart) state.runOnStartPending.delete(schedule.id); @@ -217,7 +199,6 @@ export async function tickOnce( return results; } -/** The daemon's scheduler loop — `Effect.repeat` around the plain `tickOnce`. */ export function runSchedulerLoop( deps: SchedulerDeps, state: SchedulerState, @@ -244,15 +225,10 @@ export function runSchedulerLoop( return Effect.repeat(tick, Schedule.spaced(intervalMs)); } -/** - * The daemon's fire function: a scheduled run is a config-backed run started - * like any other — registry input, workspace, dispatch environment — made - * self-deduping with `schedule:` (issue #15's registry) so the overlap - * policy and the trigger record both come for free. - */ export function makeScheduleFire(options: { readonly db: Database; readonly runtime: ManagedRuntime.ManagedRuntime; + readonly services: DaemonServices; readonly maxConcurrentRuns: number; readonly workspace: WorkspaceSpec; readonly repo: RunRepo; @@ -260,7 +236,7 @@ export function makeScheduleFire(options: { }): (schedule: RuntimeSchedule) => Promise { const env = options; return async (schedule: RuntimeSchedule): Promise => { - return startTrackedRun(env.runtime, env.db, schedule.workflow, { + return startTrackedRun(env.runtime, env.db, env.services, schedule.workflow, { input: schedule.input, repo: env.repo, maxConcurrentRuns: env.maxConcurrentRuns, diff --git a/src/server/two-daemons.test.ts b/src/server/two-daemons.test.ts new file mode 100644 index 0000000..0ae03c1 --- /dev/null +++ b/src/server/two-daemons.test.ts @@ -0,0 +1,162 @@ +/** + * #38 acceptance criterion: two daemons can run in one process without + * sharing run state. This test starts two daemons on different ports, starts + * a run on each, and verifies that each daemon's registry, pubsub, and dedupe + * state is independent. + */ + +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { describe, expect, test } from "bun:test"; +import { Schema } from "effect"; +import { defineWorkflow } from "../workflow"; +import { defineConfig } from "../config"; +import { createSlowFakeAdapter } from "../replay/adapter"; +import { startDaemon } from "./daemon"; +import { Effect, Fiber } from "effect"; + +const scratchWorkflow = defineWorkflow("two-daemon-test", { + input: Schema.Struct({ marker: Schema.String }), + workspace: { kind: "scratch" }, + run: async (ctx, input) => { + await ctx.exec(["sh", "-c", "sleep 0.1"]); + return { marker: input.marker }; + }, +}); + +async function waitForTerminal(port: number, runId: string, timeoutMs = 10_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + const res = await fetch(`http://localhost:${port}/api/runs/${runId}`); + const run = (await res.json()) as { status: string }; + if (["RunFinished", "RunFailed", "RunCancelled"].includes(run.status)) return; + await Bun.sleep(25); + } + throw new Error(`run ${runId} did not reach terminal state within ${timeoutMs}ms`); +} + +describe("two daemons in one process (#38)", () => { + test("two daemons have independent run registries, pubsub, and dedupe state", async () => { + const root = mkdtempSync(join(tmpdir(), "factory-two-daemons-")); + + const adapter = createSlowFakeAdapter( + [ + { type: "TEXT_MESSAGE_START" }, + { type: "TEXT_MESSAGE_CONTENT", delta: "hi" }, + { type: "TEXT_MESSAGE_END" }, + ], + 10, + ); + + const daemon1 = await startDaemon({ + dbPath: join(root, "daemon1.db"), + port: 0, + config: defineConfig({ + agent: { adapter }, + repo: { + sshUrl: join(root, "seed-not-used"), + identity: { name: "Factory", email: "factory@factory.test" }, + baseBranch: "main", + slug: "acme/widgets", + }, + workflows: [scratchWorkflow], + workspaceRoot: join(root, "workspaces1"), + retainedWorkspaces: 10, + }), + }); + + const daemon2 = await startDaemon({ + dbPath: join(root, "daemon2.db"), + port: 0, + config: defineConfig({ + agent: { adapter }, + repo: { + sshUrl: join(root, "seed-not-used"), + identity: { name: "Factory", email: "factory@factory.test" }, + baseBranch: "main", + slug: "acme/widgets", + }, + workflows: [scratchWorkflow], + workspaceRoot: join(root, "workspaces2"), + retainedWorkspaces: 10, + }), + }); + + try { + const port1 = daemon1.server.port; + const port2 = daemon2.server.port; + if (port1 === undefined || port2 === undefined) { + throw new Error("daemons did not bind to ports"); + } + + const res1 = await fetch(`http://localhost:${port1}/api/runs`, { + method: "POST", + body: JSON.stringify({ workflowId: "two-daemon-test", input: { marker: "daemon1" } }), + }); + expect(res1.status).toBe(201); + const { runId: runId1 } = (await res1.json()) as { runId: string }; + + const res2 = await fetch(`http://localhost:${port2}/api/runs`, { + method: "POST", + body: JSON.stringify({ workflowId: "two-daemon-test", input: { marker: "daemon2" } }), + }); + expect(res2.status).toBe(201); + const { runId: runId2 } = (await res2.json()) as { runId: string }; + + expect(runId1).not.toBe(runId2); + + const list1 = (await fetch(`http://localhost:${port1}/api/runs`).then((r) => + r.json(), + )) as Array<{ runId: string }>; + const list2 = (await fetch(`http://localhost:${port2}/api/runs`).then((r) => + r.json(), + )) as Array<{ runId: string }>; + + expect(list1.map((r) => r.runId)).toEqual([runId1]); + expect(list2.map((r) => r.runId)).toEqual([runId2]); + + const detail1FromDaemon2 = await fetch(`http://localhost:${port2}/api/runs/${runId1}`); + expect(detail1FromDaemon2.status).toBe(404); + + const detail2FromDaemon1 = await fetch(`http://localhost:${port1}/api/runs/${runId2}`); + expect(detail2FromDaemon1.status).toBe(404); + + await waitForTerminal(port1, runId1); + await waitForTerminal(port2, runId2); + + expect(daemon1.services.registry.activeRunIds()).toEqual([]); + expect(daemon2.services.registry.activeRunIds()).toEqual([]); + + const dedupeKey = "shared-key"; + const res3 = await fetch(`http://localhost:${port1}/api/runs`, { + method: "POST", + body: JSON.stringify({ workflowId: "two-daemon-test", input: { marker: "d1" }, dedupeKey }), + }); + expect(res3.status).toBe(201); + const { runId: runId3 } = (await res3.json()) as { runId: string }; + + const res4 = await fetch(`http://localhost:${port2}/api/runs`, { + method: "POST", + body: JSON.stringify({ workflowId: "two-daemon-test", input: { marker: "d2" }, dedupeKey }), + }); + expect(res4.status).toBe(201); + const { runId: runId4 } = (await res4.json()) as { runId: string }; + + expect(runId3).not.toBe(runId4); + + await waitForTerminal(port1, runId3); + await waitForTerminal(port2, runId4); + } finally { + if (daemon1.schedulerFiber !== undefined) { + Effect.runFork(Fiber.interrupt(daemon1.schedulerFiber)); + } + if (daemon2.schedulerFiber !== undefined) { + Effect.runFork(Fiber.interrupt(daemon2.schedulerFiber)); + } + daemon1.server.stop(true); + daemon2.server.stop(true); + rmSync(root, { recursive: true, force: true }); + } + }); +});