diff --git a/packages/core/src/ledger.test.ts b/packages/core/src/ledger.test.ts index f2c8be0..e4e6ce4 100644 --- a/packages/core/src/ledger.test.ts +++ b/packages/core/src/ledger.test.ts @@ -10,88 +10,19 @@ import { import { tmpdir } from "node:os"; import { join } from "node:path"; import { DatabaseSync } from "node:sqlite"; -import { Worker } from "node:worker_threads"; import { afterEach, describe, expect, it } from "vitest"; +import { runBlockedWorkers } from "./testing/blocked-sqlite-workers.js"; import { SqliteLedger } from "./ledger.js"; import type { ArtifactSubmission, RunRecord } from "./types.js"; const ledgers: SqliteLedger[] = []; -type WorkerResult = - | { kind: "result"; ok: true; artifact: unknown } - | { kind: "result"; ok: false; error: string }; - -function workerMessage(worker: Worker, kind: string): Promise { - return new Promise((resolvePromise, rejectPromise) => { - const onMessage = (message: { kind?: unknown }): void => { - if (message.kind !== kind) return; - cleanup(); - resolvePromise(message as T); - }; - const onError = (error: Error): void => { - cleanup(); - rejectPromise(error); - }; - const onExit = (code: number): void => { - cleanup(); - rejectPromise( - new Error( - `Supporting worker exited with code ${String(code)} before ${kind}`, - ), - ); - }; - const cleanup = (): void => { - worker.off("message", onMessage); - worker.off("error", onError); - worker.off("exit", onExit); - }; - worker.on("message", onMessage); - worker.on("error", onError); - worker.on("exit", onExit); - }); -} - -async function runBlockedWorkers( - databasePath: string, - workerData: readonly Record[], -): Promise { - const workerUrl = new URL( - "../../../test/helpers/supporting-artifact-worker.ts", - import.meta.url, - ); - const workers = workerData.map( - (data) => - new Worker(workerUrl, { - workerData: data, - execArgv: ["--import", "tsx"], - }), - ); - await Promise.all( - workers.map(async (worker) => workerMessage(worker, "ready")), - ); - const blocker = new DatabaseSync(databasePath); - blocker.exec("PRAGMA busy_timeout = 5000; BEGIN IMMEDIATE;"); - let locked = true; - try { - const starting = workers.map(async (worker) => - workerMessage(worker, "starting"), - ); - const results = workers.map(async (worker) => - workerMessage(worker, "result"), - ); - for (const worker of workers) worker.postMessage("go"); - await Promise.all(starting); - await new Promise((resolvePromise) => setTimeout(resolvePromise, 100)); - blocker.exec("COMMIT"); - locked = false; - return await Promise.all(results); - } finally { - if (locked) blocker.exec("ROLLBACK"); - blocker.close(); - } -} +const supportingWorkerUrl = new URL( + "./testing/supporting-artifact-worker.ts", + import.meta.url, +); function createRun(): RunRecord { return { @@ -345,6 +276,7 @@ describe("SQLite ledger and content-addressed artifacts", () => { decisionSummary: "Stored bounded evidence.", }; const results = await runBlockedWorkers( + supportingWorkerUrl, ledger.databasePath, Array.from({ length: 2 }, () => ({ kind: "ledger", @@ -377,20 +309,24 @@ describe("SQLite ledger and content-addressed artifacts", () => { eventType: "evidence_captured", decisionSummary: "Stored bounded evidence.", }; - const results = await runBlockedWorkers(ledger.databasePath, [ - { - kind: "ledger", - stateDirectory: ledger.rootDirectory, - artifact: base, - event, - }, - { - kind: "ledger", - stateDirectory: ledger.rootDirectory, - artifact: { ...base, body: { content: "second bounded result" } }, - event, - }, - ]); + const results = await runBlockedWorkers( + supportingWorkerUrl, + ledger.databasePath, + [ + { + kind: "ledger", + stateDirectory: ledger.rootDirectory, + artifact: base, + event, + }, + { + kind: "ledger", + stateDirectory: ledger.rootDirectory, + artifact: { ...base, body: { content: "second bounded result" } }, + event, + }, + ], + ); expect(results.filter((result) => result.ok)).toHaveLength(1); expect( diff --git a/packages/core/src/ledger.ts b/packages/core/src/ledger.ts index 6c0de2e..cddbbd6 100644 --- a/packages/core/src/ledger.ts +++ b/packages/core/src/ledger.ts @@ -454,6 +454,35 @@ export class SqliteLedger { } } + private reconcileExistingSupportingArtifact( + existing: HydratedArtifact, + artifact: ArtifactSubmission, + ): StoredArtifact { + const expectedSourceRefs = artifact.sourceRefs ?? []; + const expectedRedaction = artifact.redaction ?? "none"; + if ( + existing.type !== artifact.type || + existing.schemaVersion !== artifact.schemaVersion || + existing.producer !== artifact.producer || + existing.sha256 !== sha256Json(artifact.body) || + canonicalJson(existing.sourceRefs) !== + canonicalJson(expectedSourceRefs) || + existing.redaction !== expectedRedaction + ) { + throw new Error( + `Supporting artifact replay conflicts with immutable artifact: ${artifact.id}`, + ); + } + const { body: _body, ...stored } = existing; + return stored; + } + + private isUniqueConstraintError(error: unknown): boolean { + return ( + error instanceof Error && /UNIQUE constraint failed/i.test(error.message) + ); + } + appendSupportingArtifact( artifact: ArtifactSubmission, event: Omit & { phase?: Phase }, @@ -464,29 +493,29 @@ export class SqliteLedger { const run = this.requireRun(artifact.runId); const existing = this.getArtifact(artifact.runId, artifact.id); if (existing) { - const expectedSourceRefs = artifact.sourceRefs ?? []; - const expectedRedaction = artifact.redaction ?? "none"; - if ( - existing.type !== artifact.type || - existing.schemaVersion !== artifact.schemaVersion || - existing.producer !== artifact.producer || - existing.sha256 !== sha256Json(artifact.body) || - canonicalJson(existing.sourceRefs) !== - canonicalJson(expectedSourceRefs) || - existing.redaction !== expectedRedaction - ) { - throw new Error( - `Supporting artifact replay conflicts with immutable artifact: ${artifact.id}`, - ); - } - const { body: _body, ...stored } = existing; + const stored = this.reconcileExistingSupportingArtifact( + existing, + artifact, + ); this.database.exec("COMMIT"); return stored; } this.assertArtifactSlot(artifact.runId, artifact.id); if (quota) this.assertSupportingArtifactQuota(artifact, quota); const stored = this.prepareArtifact(artifact); - this.insertPreparedArtifact(stored); + try { + this.insertPreparedArtifact(stored); + } catch (error) { + if (!this.isUniqueConstraintError(error)) throw error; + const raced = this.getArtifact(artifact.runId, artifact.id); + if (!raced) throw error; + const replay = this.reconcileExistingSupportingArtifact( + raced, + artifact, + ); + this.database.exec("COMMIT"); + return replay; + } this.insertEvent(run, { ...event, phase: event.phase ?? run.phase, diff --git a/packages/core/src/testing/blocked-sqlite-workers.ts b/packages/core/src/testing/blocked-sqlite-workers.ts new file mode 100644 index 0000000..cf6ecc6 --- /dev/null +++ b/packages/core/src/testing/blocked-sqlite-workers.ts @@ -0,0 +1,97 @@ +import { DatabaseSync } from "node:sqlite"; +import { Worker } from "node:worker_threads"; + +export type BlockedWorkerResult = + | { kind: "result"; ok: true; artifact: unknown } + | { kind: "result"; ok: false; error: string }; + +function workerMessage(worker: Worker, kind: string): Promise { + return new Promise((resolvePromise, rejectPromise) => { + const onMessage = (message: { kind?: unknown }): void => { + if (message.kind !== kind) return; + cleanup(); + resolvePromise(message as T); + }; + const onError = (error: Error): void => { + cleanup(); + rejectPromise(error); + }; + const onExit = (code: number): void => { + cleanup(); + rejectPromise( + new Error( + `Supporting worker exited with code ${String(code)} before ${kind}`, + ), + ); + }; + const cleanup = (): void => { + worker.off("message", onMessage); + worker.off("error", onError); + worker.off("exit", onExit); + }; + worker.on("message", onMessage); + worker.on("error", onError); + worker.on("exit", onExit); + }); +} + +function createContentionBarrier(workerCount: number): SharedArrayBuffer { + const buffer = new SharedArrayBuffer(Int32Array.BYTES_PER_ELEMENT); + const counter = new Int32Array(buffer); + counter[0] = workerCount; + return buffer; +} + +async function waitForContentionBarrier( + buffer: SharedArrayBuffer, + timeoutMs = 5000, +): Promise { + const counter = new Int32Array(buffer); + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (Atomics.load(counter, 0) <= 0) return; + Atomics.wait(counter, 0, Atomics.load(counter, 0), 25); + } + throw new Error( + `Timed out waiting for ${String(Atomics.load(counter, 0))} workers to reach ledger contention`, + ); +} + +export async function runBlockedWorkers>( + workerModuleUrl: URL, + databasePath: string, + workerData: readonly T[], + options?: { execArgv?: string[] }, +): Promise { + const contentionBarrier = createContentionBarrier(workerData.length); + const workers = workerData.map( + (data) => + new Worker(workerModuleUrl, { + workerData: { ...data, contentionBarrier }, + execArgv: options?.execArgv ?? ["--import", "tsx"], + }), + ); + await Promise.all( + workers.map(async (worker) => workerMessage(worker, "ready")), + ); + const blocker = new DatabaseSync(databasePath); + blocker.exec("PRAGMA busy_timeout = 5000; BEGIN IMMEDIATE;"); + let locked = true; + try { + const starting = workers.map(async (worker) => + workerMessage(worker, "starting"), + ); + const results = workers.map(async (worker) => + workerMessage(worker, "result"), + ); + for (const worker of workers) worker.postMessage("go"); + await Promise.all(starting); + await waitForContentionBarrier(contentionBarrier); + blocker.exec("COMMIT"); + locked = false; + return await Promise.all(results); + } finally { + if (locked) blocker.exec("ROLLBACK"); + blocker.close(); + } +} diff --git a/test/helpers/supporting-artifact-worker.ts b/packages/core/src/testing/supporting-artifact-worker.ts similarity index 72% rename from test/helpers/supporting-artifact-worker.ts rename to packages/core/src/testing/supporting-artifact-worker.ts index a693a25..6f2cba2 100644 --- a/test/helpers/supporting-artifact-worker.ts +++ b/packages/core/src/testing/supporting-artifact-worker.ts @@ -4,9 +4,9 @@ import { SqliteLedger, type SubmissionEvent, type SupportingArtifactQuota, -} from "../../packages/core/src/ledger.ts"; -import type { ArtifactSubmission } from "../../packages/core/src/types.ts"; -import { TelicService } from "../../packages/mcp/src/service.ts"; +} from "../ledger.js"; +import type { ArtifactSubmission } from "../types.js"; +import { TelicService } from "../../../mcp/src/service.js"; type LedgerWorkerData = { kind: "ledger"; @@ -25,9 +25,22 @@ type ServiceWorkerData = { artifact: ArtifactSubmission; }; +type WorkerData = (LedgerWorkerData | ServiceWorkerData) & { + contentionBarrier?: SharedArrayBuffer; +}; + +function signalContentionReady( + contentionBarrier: SharedArrayBuffer | undefined, +): void { + if (contentionBarrier === undefined) return; + const counter = new Int32Array(contentionBarrier); + const remaining = Atomics.sub(counter, 0, 1) - 1; + if (remaining === 0) Atomics.notify(counter, 0); +} + if (parentPort === null) throw new Error("Supporting worker requires a parent"); -const data = workerData as LedgerWorkerData | ServiceWorkerData; +const data = workerData as WorkerData; const target = data.kind === "ledger" ? new SqliteLedger(data.stateDirectory) @@ -40,6 +53,7 @@ parentPort.postMessage({ kind: "ready" }); parentPort.once("message", (message: unknown) => { if (message !== "go") return; parentPort.postMessage({ kind: "starting" }); + signalContentionReady(data.contentionBarrier); try { const artifact = data.kind === "ledger" diff --git a/packages/core/tsconfig.json b/packages/core/tsconfig.json index 70d64c9..4bd9607 100644 --- a/packages/core/tsconfig.json +++ b/packages/core/tsconfig.json @@ -6,5 +6,6 @@ "outDir": "dist", "tsBuildInfoFile": "dist/.tsbuildinfo" }, - "include": ["src/**/*.ts"] + "include": ["src/**/*.ts"], + "exclude": ["src/testing/supporting-artifact-worker.ts"] } diff --git a/packages/mcp/src/pipeline.e2e.test.ts b/packages/mcp/src/pipeline.e2e.test.ts index 3cf9685..501b238 100644 --- a/packages/mcp/src/pipeline.e2e.test.ts +++ b/packages/mcp/src/pipeline.e2e.test.ts @@ -1,8 +1,6 @@ import { mkdirSync, mkdtempSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { DatabaseSync } from "node:sqlite"; -import { Worker } from "node:worker_threads"; import { afterEach, describe, expect, it } from "vitest"; @@ -10,43 +8,17 @@ import { NO_PERMISSIONS, VALID_ARTIFACT_BODIES, } from "../../protocol/test/test-helpers.js"; +import { + runBlockedWorkers, + type BlockedWorkerResult, +} from "../../core/src/testing/blocked-sqlite-workers.js"; import { TelicService } from "./service.js"; const services: TelicService[] = []; - -type SupportingWorkerResult = - | { kind: "result"; ok: true; artifact: unknown } - | { kind: "result"; ok: false; error: string }; - -function supportingWorkerMessage(worker: Worker, kind: string): Promise { - return new Promise((resolvePromise, rejectPromise) => { - const onMessage = (message: { kind?: unknown }): void => { - if (message.kind !== kind) return; - cleanup(); - resolvePromise(message as T); - }; - const onError = (error: Error): void => { - cleanup(); - rejectPromise(error); - }; - const onExit = (code: number): void => { - cleanup(); - rejectPromise( - new Error( - `Supporting worker exited with code ${String(code)} before ${kind}`, - ), - ); - }; - const cleanup = (): void => { - worker.off("message", onMessage); - worker.off("error", onError); - worker.off("exit", onExit); - }; - worker.on("message", onMessage); - worker.on("error", onError); - worker.on("exit", onExit); - }); -} +const supportingWorkerUrl = new URL( + "../../core/src/testing/supporting-artifact-worker.ts", + import.meta.url, +); async function runBlockedServiceWorkers( harness: Awaited>, @@ -55,53 +27,24 @@ async function runBlockedServiceWorkers( producer: string; body: Record; }[], -): Promise { - const workerUrl = new URL( - "../../../test/helpers/supporting-artifact-worker.ts", - import.meta.url, - ); - const workers = submissions.map( - (submission) => - new Worker(workerUrl, { - workerData: { - kind: "service", - repositoryRoot: harness.service.repositoryRoot, - stateDirectory: harness.service.stateDirectory, - artifact: { - id: submission.body.id, - runId: harness.started.run.runId, - type: submission.type, - schemaVersion: "1.0", - producer: submission.producer, - body: submission.body, - }, - }, - execArgv: ["--import", "tsx"], - }), - ); - await Promise.all( - workers.map(async (worker) => supportingWorkerMessage(worker, "ready")), +): Promise { + return runBlockedWorkers( + supportingWorkerUrl, + harness.service.ledger.databasePath, + submissions.map((submission) => ({ + kind: "service" as const, + repositoryRoot: harness.service.repositoryRoot, + stateDirectory: harness.service.stateDirectory, + artifact: { + id: submission.body.id, + runId: harness.started.run.runId, + type: submission.type, + schemaVersion: "1.0", + producer: submission.producer, + body: submission.body, + }, + })), ); - const blocker = new DatabaseSync(harness.service.ledger.databasePath); - blocker.exec("PRAGMA busy_timeout = 5000; BEGIN IMMEDIATE;"); - let locked = true; - try { - const starting = workers.map(async (worker) => - supportingWorkerMessage(worker, "starting"), - ); - const results = workers.map(async (worker) => - supportingWorkerMessage(worker, "result"), - ); - for (const worker of workers) worker.postMessage("go"); - await Promise.all(starting); - await new Promise((resolvePromise) => setTimeout(resolvePromise, 100)); - blocker.exec("COMMIT"); - locked = false; - return await Promise.all(results); - } finally { - if (locked) blocker.exec("ROLLBACK"); - blocker.close(); - } } afterEach(() => {