diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 03043a9f3..89a33de08 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -11,6 +11,21 @@ concurrency: cancel-in-progress: true jobs: + diff-check: + name: Diff check (PR) + if: github.event_name == 'pull_request' + runs-on: macos-latest + timeout-minutes: 5 + + steps: + - name: Checkout + uses: actions/checkout@v4 + with: + fetch-depth: 0 + + - name: Git diff check + run: git diff --check "${{ github.event.pull_request.base.sha }}...HEAD" + smoke: name: Smoke (${{ matrix.os }}) runs-on: ${{ matrix.os }} @@ -64,6 +79,45 @@ jobs: - name: Doctor run: node dist/cli.js doctor + chat-swarm-continuation: + name: Chat Swarm continuation ${{ matrix.name }} (macOS) + runs-on: macos-latest + timeout-minutes: 10 + strategy: + fail-fast: false + matrix: + include: + - name: coordinator + file: src/chat-swarm-continuation-coordinator.test.ts + - name: durable-store + file: src/chat-swarm-continuation-store.test.ts + - name: restart-fence + file: src/chat-swarm-continuation-restart-fence.test.ts + - name: retired-carrier + file: src/chat-swarm-continuation-stale-carrier.test.ts + - name: terminal-replay + file: src/chat-swarm-continuation-terminal-replay.test.ts + - name: corruption-fence + file: src/chat-swarm-continuation-corruption-fence.test.ts + - name: aggregate-domain + file: src/chat-swarm-continuation-domain.test.ts + + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Setup Node + uses: actions/setup-node@v4 + with: + node-version: 22 + cache: npm + + - name: Install dependencies + run: npm ci + + - name: Run focused continuation test + run: npx tsx ${{ matrix.file }} + local-agent-sessions: name: Local agent sessions (macOS) runs-on: macos-latest diff --git a/package.json b/package.json index 61a87ddb5..0e7643b8c 100644 --- a/package.json +++ b/package.json @@ -32,7 +32,7 @@ "postinstall": "node scripts/fix-node-pty-permissions.mjs", "test:node-pty-permissions": "node --import tsx --test src/node-pty-postinstall.test.ts", "start": "node dist/cli.js serve", - "test": "npm run test:node-pty-permissions && npm run test:carrier-binding && tsx src/coordination-reader-loader.test.ts && tsx src/control-plane-ownership.test.ts && tsx src/control-plane-handoff.test.ts && tsx src/control-plane-continuation.test.ts && tsx src/deployment-convergence.test.ts && tsx src/current-completion-matrix.test.ts && tsx src/control-plane-consumer.test.ts && tsx src/local-agent-cline-catalog.test.ts && tsx src/chat-swarm-contract.test.ts && tsx src/chat-swarm-store.test.ts && tsx src/chat-swarm-runtime-owner.test.ts && tsx src/local-agent-opencode-mcp-catalog.test.ts && tsx src/chat-swarm-coordinator.test.ts && tsx src/chat-swarm-peer-runtime.test.ts && tsx src/chat-swarm-tools.test.ts && tsx src/chat-swarm-lifecycle.test.ts && tsx src/chat-swarm-carrier.test.ts && tsx src/chat-swarm-peer-admission.test.ts && tsx src/chat-swarm-task-ledger.test.ts && tsx src/git-worktrees.test.ts && tsx src/execution-protocol.test.ts && tsx src/durable-operations-ci.test.ts && tsx src/git-candidate.test.ts && tsx src/config.test.ts && tsx src/onboarding.test.ts && tsx src/cli-workspace.test.ts && tsx src/request-meta.test.ts && tsx src/incoming-artifacts.test.ts && tsx src/artifact-download.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions-ci.test.ts && tsx src/codex-goal-sessions-ci.test.ts && tsx src/mcp-sessions.test.ts && tsx src/cutover-state.test.ts && tsx src/cutover-state-recovery-guard.test.ts && tsx src/cutover-build-ready.test.ts && tsx src/cutover-orchestration.test.ts && tsx src/mcp-cutover.test.ts && tsx src/cutover-restart.test.ts && tsx src/cutover-http.test.ts && tsx src/cutover-recovery.test.ts && tsx src/cutover-binding-repair.test.ts && tsx src/capability-manifest.test.ts && tsx src/server-shutdown.test.ts && tsx src/codex-runtime.test.ts && tsx src/local-agent-config.test.ts && tsx src/local-agent-catalog.test.ts && tsx src/local-agent-presentation.test.ts && tsx src/local-agent-runtime.test.ts && tsx src/local-agent-daemon-lifecycle.test.ts && tsx src/local-agent-daemon-protocol.test.ts && tsx src/local-agent-daemon.test.ts && tsx src/local-agent-codex.test.ts && tsx src/local-agent-opencode.test.ts && tsx src/local-agent-opencode-catalog.test.ts && tsx src/local-agent-acp.test.ts && tsx src/local-agent-grok.test.ts && tsx src/local-agent-pi-sandbox.test.ts && tsx src/local-agent-pi.test.ts && tsx src/local-agent-claude.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-profile-source.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-toolchains.test.ts && tsx src/local-agent-idle-policy.test.ts && tsx src/local-agent-capacity.test.ts && tsx src/local-agent-execution-contract.test.ts && tsx src/local-agent-continuation.test.ts && tsx src/provider-scratch.test.ts && tsx src/git-integration-ci.test.ts && tsx src/repository-intelligence.test.ts && tsx src/local-agent-store.test.ts && tsx src/local-agent-manager.test.ts && tsx src/roots.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/workspace-conversation.test.ts && tsx src/conversation-isolation.test.ts && tsx src/review-checkpoints.test.ts && tsx src/server-ci.test.ts && tsx src/oauth-store.test.ts && tsx src/cli-ci.test.ts && tsx src/oauth-json-and-child-reap.test.ts && tsx src/local-agent-errors.test.ts && tsx src/provider-output-redaction.test.ts && tsx src/host-operation-policy.test.ts && tsx src/host-operations.test.ts && tsx src/host-operations-http.test.ts", + "test": "npm run test:node-pty-permissions && npm run test:carrier-binding && tsx src/coordination-reader-loader.test.ts && tsx src/control-plane-ownership.test.ts && tsx src/control-plane-handoff.test.ts && tsx src/control-plane-continuation.test.ts && tsx src/deployment-convergence.test.ts && tsx src/current-completion-matrix.test.ts && tsx src/control-plane-consumer.test.ts && tsx src/local-agent-cline-catalog.test.ts && tsx src/chat-swarm-contract.test.ts && tsx src/chat-swarm-store.test.ts && tsx src/chat-swarm-runtime-owner.test.ts && tsx src/local-agent-opencode-mcp-catalog.test.ts && tsx src/chat-swarm-coordinator.test.ts && tsx src/chat-swarm-peer-runtime.test.ts && tsx src/chat-swarm-tools.test.ts && tsx src/chat-swarm-lifecycle.test.ts && tsx src/chat-swarm-carrier.test.ts && tsx src/chat-swarm-continuation-domain.test.ts && tsx src/chat-swarm-continuation-store.test.ts && tsx src/chat-swarm-peer-admission.test.ts && tsx src/chat-swarm-task-ledger.test.ts && tsx src/git-worktrees.test.ts && tsx src/execution-protocol.test.ts && tsx src/durable-operations-ci.test.ts && tsx src/git-candidate.test.ts && tsx src/config.test.ts && tsx src/onboarding.test.ts && tsx src/cli-workspace.test.ts && tsx src/request-meta.test.ts && tsx src/incoming-artifacts.test.ts && tsx src/artifact-download.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions-ci.test.ts && tsx src/codex-goal-sessions-ci.test.ts && tsx src/mcp-sessions.test.ts && tsx src/cutover-state.test.ts && tsx src/cutover-state-recovery-guard.test.ts && tsx src/cutover-build-ready.test.ts && tsx src/cutover-orchestration.test.ts && tsx src/mcp-cutover.test.ts && tsx src/cutover-restart.test.ts && tsx src/cutover-http.test.ts && tsx src/cutover-recovery.test.ts && tsx src/cutover-binding-repair.test.ts && tsx src/capability-manifest.test.ts && tsx src/server-shutdown.test.ts && tsx src/codex-runtime.test.ts && tsx src/local-agent-config.test.ts && tsx src/local-agent-catalog.test.ts && tsx src/local-agent-presentation.test.ts && tsx src/local-agent-runtime.test.ts && tsx src/local-agent-daemon-lifecycle.test.ts && tsx src/local-agent-daemon-protocol.test.ts && tsx src/local-agent-daemon.test.ts && tsx src/local-agent-codex.test.ts && tsx src/local-agent-opencode.test.ts && tsx src/local-agent-opencode-catalog.test.ts && tsx src/local-agent-acp.test.ts && tsx src/local-agent-grok.test.ts && tsx src/local-agent-pi-sandbox.test.ts && tsx src/local-agent-pi.test.ts && tsx src/local-agent-claude.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-profile-source.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-toolchains.test.ts && tsx src/local-agent-idle-policy.test.ts && tsx src/local-agent-capacity.test.ts && tsx src/local-agent-execution-contract.test.ts && tsx src/local-agent-continuation.test.ts && tsx src/provider-scratch.test.ts && tsx src/git-integration-ci.test.ts && tsx src/repository-intelligence.test.ts && tsx src/local-agent-store.test.ts && tsx src/local-agent-manager.test.ts && tsx src/roots.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/workspace-conversation.test.ts && tsx src/conversation-isolation.test.ts && tsx src/review-checkpoints.test.ts && tsx src/server-ci.test.ts && tsx src/oauth-store.test.ts && tsx src/cli-ci.test.ts && tsx src/oauth-json-and-child-reap.test.ts && tsx src/local-agent-errors.test.ts && tsx src/provider-output-redaction.test.ts && tsx src/host-operation-policy.test.ts && tsx src/host-operations.test.ts && tsx src/host-operations-http.test.ts", "typecheck": "tsc -p tsconfig.json --noEmit", "test:carrier-binding": "tsx src/carrier-binding.test.ts && tsx src/carrier-binding-http.test.ts", "test:local-agent-sessions": "node --import tsx --test src/local-agent-sessions.test.ts" @@ -85,4 +85,4 @@ "optionalDependencies": { "node-pty": "1.1.0" } -} +} \ No newline at end of file diff --git a/src/chat-swarm-continuation-contract.ts b/src/chat-swarm-continuation-contract.ts new file mode 100644 index 000000000..fd46d172c --- /dev/null +++ b/src/chat-swarm-continuation-contract.ts @@ -0,0 +1,43 @@ +export const CHAT_SWARM_CONTINUATION_STATES = [ + "PENDING", + "APPROVED", + "EXPIRED", + "SUPERSEDED", + "RECONCILE_REQUIRED", +] as const; + +export type ChatSwarmContinuationState = (typeof CHAT_SWARM_CONTINUATION_STATES)[number]; + +export interface ChatSwarmContinuationRequest { + id: string; + swarmId: string; + workerId: string; + attemptKey: string; + requestHash: string; + sourceEpoch: number; + targetEpoch: number; + sourceCarrierFingerprint: string; + targetCarrierFingerprint: string; + checkpointHash: string; + ttlSeconds: number; + version: number; + status: ChatSwarmContinuationState; + requestedAt: string; + expiresAt: string; + approvedAt?: string; +} + +export interface CreateContinuationRequestInput { + swarmId: string; + workerId: string; + attemptKey: string; + sourceEpoch: number; + ttlSeconds?: number; +} + +export interface ApproveContinuationRequestInput { + swarmId: string; + requestId: string; + expectedRequestVersion: number; + expectedSwarmVersion: number; +} diff --git a/src/chat-swarm-continuation-coordinator.test.ts b/src/chat-swarm-continuation-coordinator.test.ts new file mode 100644 index 000000000..6ddcf4795 --- /dev/null +++ b/src/chat-swarm-continuation-coordinator.test.ts @@ -0,0 +1,147 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { ChatSwarmError } from "./chat-swarm-contract.js"; +import { ChatSwarmContinuationCoordinator } from "./chat-swarm-continuation-coordinator.js"; +import { ChatSwarmContinuationStore } from "./chat-swarm-continuation-store.js"; +import { ChatSwarmCoordinator } from "./chat-swarm-coordinator.js"; +import { ChatSwarmStore } from "./chat-swarm-store.js"; + +function fixture() { + const root = mkdtempSync(join(tmpdir(), "swarm-continuation-coordinator-")); + const swarmStore = new ChatSwarmStore(root); + const swarmCoordinator = new ChatSwarmCoordinator(swarmStore); + const continuationStore = new ChatSwarmContinuationStore(root); + const continuation = new ChatSwarmContinuationCoordinator(continuationStore, swarmCoordinator); + const ownerMeta = { "openai/session": "continuation-owner" }; + const attackerMeta = { "openai/session": "continuation-attacker" }; + const sourceMeta = { "openai/conversation_id": "continuation-source" }; + const targetMeta = { "openai/conversation_id": "continuation-target" }; + const otherTargetMeta = { "openai/conversation_id": "continuation-other-target" }; + const swarm = swarmCoordinator.createSwarm(ownerMeta, { workerLimit: 3, metadata: { test: true } }); + const worker = swarmCoordinator.joinWorker(sourceMeta, swarm.id, { + label: "Worker-01", + runtimeKind: "mcp_peer", + }); + swarmCoordinator.checkpoint( + sourceMeta, + worker.id, + 0, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { summary: "safe checkpoint" }, + ); + return { + root, + swarmStore, + continuationStore, + continuation, + ownerMeta, + attackerMeta, + sourceMeta, + targetMeta, + otherTargetMeta, + swarm, + worker, + }; +} + +function cleanup(f: ReturnType) { + f.continuationStore.close(); + f.swarmStore.close(); + rmSync(f.root, { recursive: true, force: true }); +} + +test("target request identity comes from authenticated metadata and exact target can read it", () => { + const f = fixture(); + try { + const created = f.continuation.request(f.targetMeta, { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey: "target-bound-1", + sourceEpoch: 0, + }); + assert.equal(created.created, true); + assert.equal(f.continuation.targetStatus(f.targetMeta, f.swarm.id)?.id, created.request.id); + assert.equal(f.continuation.targetStatus(f.otherTargetMeta, f.swarm.id), undefined); + + assert.throws( + () => f.continuation.request(f.sourceMeta, { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey: "source-cannot-replace-itself", + sourceEpoch: 0, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_INPUT", + ); + } finally { + cleanup(f); + } +}); + +test("only exact owner may approve, including idempotent approved readback", () => { + const f = fixture(); + try { + const created = f.continuation.request(f.targetMeta, { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey: "owner-bound-1", + sourceEpoch: 0, + }); + + assert.throws( + () => f.continuation.approve(f.attackerMeta, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "OWNERSHIP_CONFLICT", + ); + + const approved = f.continuation.approve(f.ownerMeta, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }); + assert.equal(approved.status, "APPROVED"); + + assert.throws( + () => f.continuation.approve(f.attackerMeta, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "OWNERSHIP_CONFLICT", + ); + } finally { + cleanup(f); + } +}); + +test("same owner cannot reconcile a continuation through a different swarm id", () => { + const f = fixture(); + try { + const created = f.continuation.request(f.targetMeta, { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey: "cross-swarm-reconcile", + sourceEpoch: 0, + }); + const otherSwarm = f.continuation.swarmCoordinator.createSwarm(f.ownerMeta, { + workerLimit: 1, + metadata: { test: "other-swarm" }, + }); + + assert.throws( + () => f.continuation.reconcileNoEffect(f.ownerMeta, otherSwarm.id, created.request.id), + (error: unknown) => error instanceof ChatSwarmError && error.code === "OWNERSHIP_CONFLICT", + ); + assert.equal(f.continuationStore.getRequest(created.request.id)?.status, "PENDING"); + } finally { + cleanup(f); + } +}); diff --git a/src/chat-swarm-continuation-coordinator.ts b/src/chat-swarm-continuation-coordinator.ts new file mode 100644 index 000000000..f1cad52cc --- /dev/null +++ b/src/chat-swarm-continuation-coordinator.ts @@ -0,0 +1,56 @@ +import { ChatSwarmError } from "./chat-swarm-contract.js"; +import type { + ApproveContinuationRequestInput, + ChatSwarmContinuationRequest, + CreateContinuationRequestInput, +} from "./chat-swarm-continuation-contract.js"; +import { ChatSwarmContinuationStore } from "./chat-swarm-continuation-store.js"; +import { ChatSwarmCoordinator } from "./chat-swarm-coordinator.js"; +import { resolveChatSwarmIdentity } from "./request-meta.js"; + +export class ChatSwarmContinuationCoordinator { + constructor( + readonly store: ChatSwarmContinuationStore, + readonly swarmCoordinator: ChatSwarmCoordinator, + ) {} + + request( + meta: unknown, + input: CreateContinuationRequestInput, + ): { request: ChatSwarmContinuationRequest; created: boolean } { + const identity = resolveChatSwarmIdentity(meta); + return this.store.createRequest(identity.fingerprint, input); + } + + targetStatus(meta: unknown, swarmId: string): ChatSwarmContinuationRequest | undefined { + const identity = resolveChatSwarmIdentity(meta); + return this.store.getLatestForTarget(swarmId, identity.fingerprint); + } + + approve(meta: unknown, input: ApproveContinuationRequestInput): ChatSwarmContinuationRequest { + this.swarmCoordinator.assertOwnerForLifecycle(meta, input.swarmId); + const identity = resolveChatSwarmIdentity(meta); + return this.store.approveRequest(identity.fingerprint, input); + } + + reconcileNoEffect( + meta: unknown, + swarmId: string, + requestId: string, + ): ChatSwarmContinuationRequest { + this.swarmCoordinator.assertOwnerForLifecycle(meta, swarmId); + const request = this.store.getRequest(requestId); + if (!request || request.swarmId !== swarmId) { + throw new ChatSwarmError( + "OWNERSHIP_CONFLICT", + "continuation request does not belong to asserted swarm", + ); + } + const identity = resolveChatSwarmIdentity(meta); + return this.store.reconcileUnknownNoEffect(identity.fingerprint, requestId); + } + + recoverAfterRestart(): number { + return this.store.recoverAfterRestart(); + } +} diff --git a/src/chat-swarm-continuation-corruption-fence.test.ts b/src/chat-swarm-continuation-corruption-fence.test.ts new file mode 100644 index 000000000..7af1f1e57 --- /dev/null +++ b/src/chat-swarm-continuation-corruption-fence.test.ts @@ -0,0 +1,206 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { ChatSwarmError } from "./chat-swarm-contract.js"; +import { ChatSwarmContinuationStore } from "./chat-swarm-continuation-store.js"; +import { ChatSwarmCoordinator } from "./chat-swarm-coordinator.js"; +import { ChatSwarmStore } from "./chat-swarm-store.js"; +import { openDatabase } from "./db/client.js"; +import { resolveChatSwarmIdentity } from "./request-meta.js"; + +function fixture() { + const root = mkdtempSync(join(tmpdir(), "swarm-continuation-corruption-fence-")); + const swarmStore = new ChatSwarmStore(root); + const continuation = new ChatSwarmContinuationStore(root); + const coordinator = new ChatSwarmCoordinator(swarmStore); + const ownerMeta = { "openai/session": "continuation-owner" }; + const sourceMeta = { "openai/conversation_id": "continuation-source" }; + const targetA = resolveChatSwarmIdentity({ "openai/conversation_id": "continuation-target-a" }).fingerprint; + const targetB = resolveChatSwarmIdentity({ "openai/conversation_id": "continuation-target-b" }).fingerprint; + const ownerFingerprint = resolveChatSwarmIdentity(ownerMeta).fingerprint; + const swarm = coordinator.createSwarm(ownerMeta, { workerLimit: 2, metadata: { test: true } }); + const worker = coordinator.joinWorker(sourceMeta, swarm.id, { + label: "Worker-01", + runtimeKind: "mcp_peer", + }); + coordinator.checkpoint( + sourceMeta, + worker.id, + 0, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { summary: "safe checkpoint" }, + ); + return { + root, + swarmStore, + continuation, + coordinator, + ownerMeta, + sourceMeta, + targetA, + targetB, + ownerFingerprint, + swarm, + worker, + }; +} + +function cleanup(f: ReturnType): void { + f.continuation.close(); + f.swarmStore.close(); + rmSync(f.root, { recursive: true, force: true }); +} + +function createRequest(f: ReturnType, attemptKey: string, target = f.targetA) { + return f.continuation.createRequest(target, { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey, + sourceEpoch: 0, + }); +} + +test("corrupt pending continuation identity fences all new continuation work", () => { + const f = fixture(); + try { + const first = createRequest(f, "original-attempt"); + + const database = openDatabase(f.root); + try { + const row = database.sqlite.prepare( + "select request_json from durable_operations where operation_id=?", + ).get(first.request.id) as { request_json: string }; + const request = JSON.parse(row.request_json) as Record; + request.attemptKey = "tampered-attempt"; + database.sqlite.prepare( + "update durable_operations set request_json=? where operation_id=?", + ).run(JSON.stringify(request), first.request.id); + } finally { + database.close(); + } + + assert.throws( + () => createRequest(f, "new-attempt-must-not-bypass-corruption", f.targetB), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + assert.throws( + () => f.continuation.getLatestForTarget(f.swarm.id, f.targetB), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + assert.equal(f.swarmStore.getWorker(f.worker.id)?.continuationEpoch, 0); + } finally { + cleanup(f); + } +}); + +test("shifted absolute expiry with unchanged TTL fails closed", () => { + const f = fixture(); + try { + const created = createRequest(f, "shifted-expiry"); + const database = openDatabase(f.root); + try { + const row = database.sqlite.prepare( + "select request_json from durable_operations where operation_id=?", + ).get(created.request.id) as { request_json: string }; + const request = JSON.parse(row.request_json) as Record; + const requestedAt = Date.parse(String(request.requestedAt)); + const expiresAt = Date.parse(String(request.expiresAt)); + request.requestedAt = new Date(requestedAt + 60_000).toISOString(); + request.expiresAt = new Date(expiresAt + 60_000).toISOString(); + database.sqlite.prepare( + "update durable_operations set request_json=? where operation_id=?", + ).run(JSON.stringify(request), created.request.id); + } finally { + database.close(); + } + + assert.throws( + () => f.continuation.getRequest(created.request.id), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + } finally { + cleanup(f); + } +}); + +test("succeeded continuation without approval receipt fails closed", () => { + const f = fixture(); + try { + const created = createRequest(f, "approved-with-receipt"); + const approved = f.continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }); + assert.equal(approved.status, "APPROVED"); + + const database = openDatabase(f.root); + try { + database.sqlite.prepare( + "update durable_operations set receipt_json=null where operation_id=?", + ).run(created.request.id); + } finally { + database.close(); + } + + assert.throws( + () => f.continuation.getRequest(created.request.id), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + } finally { + cleanup(f); + } +}); + +test("pending continuation with forged approval receipt fails closed", () => { + const f = fixture(); + try { + const created = createRequest(f, "pending-forged-receipt"); + const database = openDatabase(f.root); + try { + database.sqlite.prepare( + "update durable_operations set receipt_json=? where operation_id=?", + ).run(JSON.stringify({ schema: "forged" }), created.request.id); + } finally { + database.close(); + } + + assert.throws( + () => f.continuation.getRequest(created.request.id), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + assert.throws( + () => createRequest(f, "pending-forged-receipt-new", f.targetB), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + } finally { + cleanup(f); + } +}); + +test("outcome-unknown continuation with wrong reconciliation code fails closed", () => { + const f = fixture(); + try { + const created = createRequest(f, "unknown-wrong-error-code"); + assert.equal(f.continuation.recoverAfterRestart(), 1); + + const database = openDatabase(f.root); + try { + database.sqlite.prepare( + "update durable_operations set error_code='WRONG_CODE' where operation_id=?", + ).run(created.request.id); + } finally { + database.close(); + } + + assert.throws( + () => f.continuation.getRequest(created.request.id), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + } finally { + cleanup(f); + } +}); diff --git a/src/chat-swarm-continuation-domain.test.ts b/src/chat-swarm-continuation-domain.test.ts new file mode 100644 index 000000000..0631b1ed8 --- /dev/null +++ b/src/chat-swarm-continuation-domain.test.ts @@ -0,0 +1,123 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import "./chat-swarm-continuation-coordinator.test.js"; +import "./chat-swarm-continuation-corruption-fence.test.js"; +import "./chat-swarm-continuation-restart-fence.test.js"; +import "./chat-swarm-continuation-stale-carrier.test.js"; +import "./chat-swarm-continuation-terminal-replay.test.js"; +import { ChatSwarmError } from "./chat-swarm-contract.js"; +import { + assertContinuationCommitAllowed, + prepareContinuationMaterial, + type WorkerContinuationSnapshot, +} from "./chat-swarm-continuation-domain.js"; +import type { ChatSwarmContinuationRequest } from "./chat-swarm-continuation-contract.js"; + +const source = "a".repeat(64); +const target = "b".repeat(64); + +function worker(overrides: Partial = {}): WorkerContinuationSnapshot { + return { + swarmId: "swarm-1", + workerId: "worker-1", + lifecycleState: "AVAILABLE", + continuationEpoch: 4, + carrierConversationFingerprint: source, + checkpoint: { lastTaskId: "task-9", summary: "safe boundary" }, + ...overrides, + }; +} + +function createInput(attemptKey = "cont-1", sourceEpoch = 4) { + return { + swarmId: "swarm-1", + workerId: "worker-1", + attemptKey, + sourceEpoch, + }; +} + +function request(overrides: Partial = {}): ChatSwarmContinuationRequest { + const material = prepareContinuationMaterial(worker(), createInput(), target); + return { + id: "contreq-1", + swarmId: "swarm-1", + workerId: "worker-1", + attemptKey: "cont-1", + requestHash: "d".repeat(64), + sourceEpoch: material.sourceEpoch, + targetEpoch: material.targetEpoch, + sourceCarrierFingerprint: material.sourceCarrierFingerprint, + targetCarrierFingerprint: material.targetCarrierFingerprint, + checkpointHash: material.checkpointHash, + ttlSeconds: 15 * 60, + version: 1, + status: "PENDING", + requestedAt: "2026-09-13T00:00:00.000Z", + expiresAt: "2026-09-13T00:15:00.000Z", + ...overrides, + }; +} + +test("prepare binds exact epoch, authenticated target carrier and checkpoint", () => { + const material = prepareContinuationMaterial(worker(), createInput(), target); + assert.equal(material.sourceEpoch, 4); + assert.equal(material.targetEpoch, 5); + assert.equal(material.sourceCarrierFingerprint, source); + assert.equal(material.targetCarrierFingerprint, target); + assert.match(material.checkpointHash, /^[0-9a-f]{64}$/); + + const otherTarget = "c".repeat(64); + const other = prepareContinuationMaterial(worker(), createInput(), otherTarget); + assert.notEqual(material.targetCarrierFingerprint, other.targetCarrierFingerprint); +}); + +test("prepare rejects unsafe active worker and stale epoch", () => { + assert.throws( + () => prepareContinuationMaterial( + worker({ lifecycleState: "BUSY", currentTaskId: "task-live" }), + createInput("cont-2"), + target, + ), + (error: unknown) => error instanceof ChatSwarmError && error.code === "RECONCILIATION_REQUIRED", + ); + + assert.throws( + () => prepareContinuationMaterial(worker(), createInput("cont-3", 3), target), + (error: unknown) => error instanceof ChatSwarmError && error.code === "OWNERSHIP_CONFLICT", + ); +}); + +test("prepare rejects reuse of the current carrier as target", () => { + assert.throws( + () => prepareContinuationMaterial(worker(), createInput("cont-same"), source), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_INPUT", + ); +}); + +test("commit guard rejects changed source carrier and checkpoint", () => { + const prepared = request(); + assert.doesNotThrow(() => assertContinuationCommitAllowed(worker(), prepared)); + + assert.throws( + () => assertContinuationCommitAllowed(worker({ carrierConversationFingerprint: "c".repeat(64) }), prepared), + (error: unknown) => error instanceof ChatSwarmError && error.code === "CAS_DRIFT", + ); + + assert.throws( + () => assertContinuationCommitAllowed(worker({ checkpoint: { summary: "changed" } }), prepared), + (error: unknown) => error instanceof ChatSwarmError && error.code === "CAS_DRIFT", + ); +}); + +test("commit guard rejects malformed target epoch and non-pending request", () => { + assert.throws( + () => assertContinuationCommitAllowed(worker(), request({ targetEpoch: 7 })), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + + assert.throws( + () => assertContinuationCommitAllowed(worker(), request({ status: "APPROVED" })), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); +}); diff --git a/src/chat-swarm-continuation-domain.ts b/src/chat-swarm-continuation-domain.ts new file mode 100644 index 000000000..96c0ed3b5 --- /dev/null +++ b/src/chat-swarm-continuation-domain.ts @@ -0,0 +1,113 @@ +import { + canonicalize, + ChatSwarmError, + hashContent, +} from "./chat-swarm-contract.js"; +import type { + ChatSwarmContinuationRequest, + CreateContinuationRequestInput, +} from "./chat-swarm-continuation-contract.js"; + +const SHA256 = /^[0-9a-f]{64}$/; + +export interface WorkerContinuationSnapshot { + swarmId: string; + workerId: string; + lifecycleState: "AVAILABLE" | "BUSY" | "DISABLED" | "RECONCILE_REQUIRED"; + currentTaskId?: string; + continuationEpoch: number; + carrierConversationFingerprint?: string; + checkpoint?: Record; +} + +export interface PreparedContinuationMaterial { + sourceEpoch: number; + targetEpoch: number; + sourceCarrierFingerprint: string; + targetCarrierFingerprint: string; + checkpointHash: string; +} + +export function prepareContinuationMaterial( + worker: WorkerContinuationSnapshot, + input: CreateContinuationRequestInput, + authenticatedTargetCarrierFingerprint: string, +): PreparedContinuationMaterial { + if (worker.swarmId !== input.swarmId || worker.workerId !== input.workerId) { + throw new ChatSwarmError("OWNERSHIP_CONFLICT", "worker does not match requested swarm identity"); + } + if (worker.lifecycleState !== "AVAILABLE" || worker.currentTaskId) { + throw new ChatSwarmError("RECONCILIATION_REQUIRED", "worker is not at a safe continuation boundary"); + } + if (worker.continuationEpoch !== input.sourceEpoch) { + throw new ChatSwarmError("OWNERSHIP_CONFLICT", "worker continuation epoch changed"); + } + const sourceCarrierFingerprint = requireFingerprint( + worker.carrierConversationFingerprint, + "source carrier fingerprint", + ); + const targetCarrierFingerprint = requireFingerprint( + authenticatedTargetCarrierFingerprint, + "target carrier fingerprint", + ); + if (sourceCarrierFingerprint === targetCarrierFingerprint) { + throw new ChatSwarmError("INVALID_INPUT", "target carrier must differ from source carrier"); + } + if (!worker.checkpoint || Array.isArray(worker.checkpoint) || typeof worker.checkpoint !== "object") { + throw new ChatSwarmError("RECONCILIATION_REQUIRED", "worker has no bounded continuation checkpoint"); + } + + const checkpointJson = JSON.stringify(canonicalize(worker.checkpoint)); + const checkpointHash = hashContent(checkpointJson); + const targetEpoch = input.sourceEpoch + 1; + + return { + sourceEpoch: input.sourceEpoch, + targetEpoch, + sourceCarrierFingerprint, + targetCarrierFingerprint, + checkpointHash, + }; +} + +export function assertContinuationCommitAllowed( + worker: WorkerContinuationSnapshot, + request: ChatSwarmContinuationRequest, +): void { + if (request.status !== "PENDING") { + throw new ChatSwarmError("INVALID_STATE", "continuation request is not pending"); + } + if (worker.swarmId !== request.swarmId || worker.workerId !== request.workerId) { + throw new ChatSwarmError("OWNERSHIP_CONFLICT", "continuation request targets another worker"); + } + if (worker.lifecycleState !== "AVAILABLE" || worker.currentTaskId) { + throw new ChatSwarmError("RECONCILIATION_REQUIRED", "worker is not at a safe continuation boundary"); + } + if (worker.continuationEpoch !== request.sourceEpoch) { + throw new ChatSwarmError("CAS_DRIFT", "worker continuation epoch changed before transfer"); + } + if (request.targetEpoch !== request.sourceEpoch + 1) { + throw new ChatSwarmError("INVALID_STATE", "continuation epoch binding is malformed"); + } + const sourceCarrierFingerprint = requireFingerprint( + worker.carrierConversationFingerprint, + "source carrier fingerprint", + ); + if (sourceCarrierFingerprint !== request.sourceCarrierFingerprint) { + throw new ChatSwarmError("CAS_DRIFT", "worker source carrier changed before transfer"); + } + if (!worker.checkpoint || Array.isArray(worker.checkpoint) || typeof worker.checkpoint !== "object") { + throw new ChatSwarmError("RECONCILIATION_REQUIRED", "worker checkpoint is unavailable"); + } + const checkpointHash = hashContent(JSON.stringify(canonicalize(worker.checkpoint))); + if (checkpointHash !== request.checkpointHash) { + throw new ChatSwarmError("CAS_DRIFT", "worker checkpoint changed before transfer"); + } +} + +function requireFingerprint(value: string | undefined, label: string): string { + if (!value || !SHA256.test(value)) { + throw new ChatSwarmError("INVALID_INPUT", `${label} must be a SHA-256 fingerprint`); + } + return value; +} diff --git a/src/chat-swarm-continuation-restart-fence.test.ts b/src/chat-swarm-continuation-restart-fence.test.ts new file mode 100644 index 000000000..0b472d844 --- /dev/null +++ b/src/chat-swarm-continuation-restart-fence.test.ts @@ -0,0 +1,75 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { ChatSwarmError } from "./chat-swarm-contract.js"; +import { ChatSwarmContinuationStore } from "./chat-swarm-continuation-store.js"; +import { ChatSwarmCoordinator } from "./chat-swarm-coordinator.js"; +import { ChatSwarmStore } from "./chat-swarm-store.js"; +import { resolveChatSwarmIdentity } from "./request-meta.js"; + +test("unresolved continuation fences any second replacement attempt", () => { + const root = mkdtempSync(join(tmpdir(), "swarm-continuation-restart-fence-")); + const swarmStore = new ChatSwarmStore(root); + const continuation = new ChatSwarmContinuationStore(root); + const coordinator = new ChatSwarmCoordinator(swarmStore); + const ownerMeta = { "openai/session": "continuation-owner" }; + const sourceMeta = { "openai/conversation_id": "continuation-source" }; + const targetA = resolveChatSwarmIdentity({ "openai/conversation_id": "continuation-target-a" }).fingerprint; + const targetB = resolveChatSwarmIdentity({ "openai/conversation_id": "continuation-target-b" }).fingerprint; + const ownerFingerprint = resolveChatSwarmIdentity(ownerMeta).fingerprint; + + try { + const swarm = coordinator.createSwarm(ownerMeta, { workerLimit: 2, metadata: { test: true } }); + const worker = coordinator.joinWorker(sourceMeta, swarm.id, { + label: "Worker-01", + runtimeKind: "mcp_peer", + }); + coordinator.checkpoint( + sourceMeta, + worker.id, + 0, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { summary: "safe checkpoint" }, + ); + + const first = continuation.createRequest(targetA, { + swarmId: swarm.id, + workerId: worker.id, + attemptKey: "first-replacement", + sourceEpoch: 0, + }); + assert.equal(continuation.recoverAfterRestart(), 1); + assert.equal(continuation.getRequest(first.request.id)?.status, "RECONCILE_REQUIRED"); + + assert.throws( + () => continuation.createRequest(targetB, { + swarmId: swarm.id, + workerId: worker.id, + attemptKey: "second-replacement", + sourceEpoch: 0, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "RECONCILIATION_REQUIRED", + ); + assert.equal(continuation.getLatestForTarget(swarm.id, targetB), undefined); + assert.equal(swarmStore.getWorker(worker.id)?.continuationEpoch, 0); + assert.notEqual(swarmStore.getWorker(worker.id)?.carrierConversationFingerprint, targetB); + + const reconciled = continuation.reconcileUnknownNoEffect(ownerFingerprint, first.request.id); + assert.equal(reconciled.status, "PENDING"); + + const second = continuation.createRequest(targetB, { + swarmId: swarm.id, + workerId: worker.id, + attemptKey: "second-replacement", + sourceEpoch: 0, + }); + assert.equal(second.created, true); + assert.equal(second.request.status, "PENDING"); + } finally { + continuation.close(); + swarmStore.close(); + rmSync(root, { recursive: true, force: true }); + } +}); diff --git a/src/chat-swarm-continuation-stale-carrier.test.ts b/src/chat-swarm-continuation-stale-carrier.test.ts new file mode 100644 index 000000000..d60f87e8c --- /dev/null +++ b/src/chat-swarm-continuation-stale-carrier.test.ts @@ -0,0 +1,86 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { ChatSwarmError } from "./chat-swarm-contract.js"; +import { ChatSwarmContinuationStore } from "./chat-swarm-continuation-store.js"; +import { ChatSwarmCoordinator } from "./chat-swarm-coordinator.js"; +import { ChatSwarmStore } from "./chat-swarm-store.js"; +import { resolveChatSwarmIdentity } from "./request-meta.js"; + +test("successful continuation permanently fences the retired source carrier", () => { + const root = mkdtempSync(join(tmpdir(), "swarm-continuation-stale-carrier-")); + const swarmStore = new ChatSwarmStore(root); + const continuation = new ChatSwarmContinuationStore(root); + const coordinator = new ChatSwarmCoordinator(swarmStore); + const ownerMeta = { "openai/session": "continuation-owner" }; + const carrierAMeta = { "openai/conversation_id": "continuation-carrier-a" }; + const carrierBMeta = { "openai/conversation_id": "continuation-carrier-b" }; + const carrierA = resolveChatSwarmIdentity(carrierAMeta).fingerprint; + const carrierB = resolveChatSwarmIdentity(carrierBMeta).fingerprint; + const carrierC = resolveChatSwarmIdentity({ "openai/conversation_id": "continuation-carrier-c" }).fingerprint; + const ownerFingerprint = resolveChatSwarmIdentity(ownerMeta).fingerprint; + + try { + const swarm = coordinator.createSwarm(ownerMeta, { workerLimit: 2, metadata: { test: true } }); + const worker = coordinator.joinWorker(carrierAMeta, swarm.id, { + label: "Worker-01", + runtimeKind: "mcp_peer", + }); + coordinator.checkpoint( + carrierAMeta, + worker.id, + 0, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { summary: "epoch zero checkpoint" }, + ); + + const first = continuation.createRequest(carrierB, { + swarmId: swarm.id, + workerId: worker.id, + attemptKey: "a-to-b", + sourceEpoch: 0, + }); + continuation.approveRequest(ownerFingerprint, { + swarmId: swarm.id, + requestId: first.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }); + assert.equal(swarmStore.getWorker(worker.id)?.carrierConversationFingerprint, carrierB); + assert.equal(swarmStore.getWorker(worker.id)?.continuationEpoch, 1); + + coordinator.checkpoint( + carrierBMeta, + worker.id, + 1, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { summary: "epoch one checkpoint" }, + ); + + assert.throws( + () => continuation.createRequest(carrierA, { + swarmId: swarm.id, + workerId: worker.id, + attemptKey: "b-back-to-retired-a", + sourceEpoch: 1, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "OWNERSHIP_CONFLICT", + ); + assert.equal(continuation.getLatestForTarget(swarm.id, carrierA), undefined); + + const fresh = continuation.createRequest(carrierC, { + swarmId: swarm.id, + workerId: worker.id, + attemptKey: "b-to-fresh-c", + sourceEpoch: 1, + }); + assert.equal(fresh.created, true); + assert.equal(fresh.request.targetCarrierFingerprint, carrierC); + } finally { + continuation.close(); + swarmStore.close(); + rmSync(root, { recursive: true, force: true }); + } +}); diff --git a/src/chat-swarm-continuation-store.test.ts b/src/chat-swarm-continuation-store.test.ts new file mode 100644 index 000000000..19e3767d1 --- /dev/null +++ b/src/chat-swarm-continuation-store.test.ts @@ -0,0 +1,335 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { ChatSwarmError } from "./chat-swarm-contract.js"; +import { ChatSwarmCoordinator } from "./chat-swarm-coordinator.js"; +import { ChatSwarmStore } from "./chat-swarm-store.js"; +import { ChatSwarmContinuationStore } from "./chat-swarm-continuation-store.js"; +import { openDatabase } from "./db/client.js"; +import { resolveChatSwarmIdentity } from "./request-meta.js"; + +function fixture() { + const root = mkdtempSync(join(tmpdir(), "swarm-continuation-store-")); + const swarmStore = new ChatSwarmStore(root); + const coordinator = new ChatSwarmCoordinator(swarmStore); + const ownerMeta = { "openai/session": "continuation-owner" }; + const sourceMeta = { "openai/conversation_id": "continuation-source" }; + const targetMeta = { "openai/conversation_id": "continuation-target" }; + const ownerFingerprint = resolveChatSwarmIdentity(ownerMeta).fingerprint; + const targetFingerprint = resolveChatSwarmIdentity(targetMeta).fingerprint; + const swarm = coordinator.createSwarm(ownerMeta, { workerLimit: 3, metadata: { test: true } }); + const worker = coordinator.joinWorker(sourceMeta, swarm.id, { + label: "Worker-01", + runtimeKind: "mcp_peer", + }); + coordinator.checkpoint( + sourceMeta, + worker.id, + 0, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { lastTaskId: "task-safe", summary: "safe checkpoint" }, + ); + return { + root, + swarmStore, + coordinator, + ownerMeta, + sourceMeta, + targetMeta, + ownerFingerprint, + targetFingerprint, + swarm, + worker, + }; +} + +function cleanup(f: ReturnType, continuation?: ChatSwarmContinuationStore) { + continuation?.close(); + f.swarmStore.close(); + rmSync(f.root, { recursive: true, force: true }); +} + +function createInput(f: ReturnType, attemptKey = "continuation-attempt-1") { + return { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey, + sourceEpoch: 0, + }; +} + +test("durable continuation request replays exactly and changed material conflicts on the same attempt", () => { + const f = fixture(); + const continuation = new ChatSwarmContinuationStore(f.root); + try { + const first = continuation.createRequest(f.targetFingerprint, createInput(f)); + assert.equal(first.created, true); + assert.equal(first.request.status, "PENDING"); + assert.equal(first.request.version, 1); + assert.equal(first.request.sourceEpoch, 0); + assert.equal(first.request.targetEpoch, 1); + assert.equal(first.request.ttlSeconds, 15 * 60); + + const replay = continuation.createRequest(f.targetFingerprint, createInput(f)); + assert.equal(replay.created, false); + assert.equal(replay.request.id, first.request.id); + assert.equal(replay.request.requestHash, first.request.requestHash); + + assert.throws( + () => continuation.createRequest(f.targetFingerprint, { + ...createInput(f), + ttlSeconds: 1, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "REPLAY_CONFLICT", + ); + + const otherTarget = resolveChatSwarmIdentity({ + "openai/conversation_id": "continuation-target-other", + }).fingerprint; + assert.throws( + () => continuation.createRequest(otherTarget, createInput(f)), + (error: unknown) => error instanceof ChatSwarmError && error.code === "REPLAY_CONFLICT", + ); + + const second = continuation.createRequest( + otherTarget, + createInput(f, "continuation-attempt-2"), + ); + assert.equal(second.created, true); + assert.notEqual(second.request.id, first.request.id); + assert.equal(second.request.sourceEpoch, 0); + assert.equal(second.request.targetEpoch, 1); + } finally { + cleanup(f, continuation); + } +}); + +test("parallel prepared targets terminate losers as SUPERSEDED at winner commit", () => { + const f = fixture(); + const continuation = new ChatSwarmContinuationStore(f.root); + try { + const targetA = f.targetFingerprint; + const targetB = resolveChatSwarmIdentity({ + "openai/conversation_id": "continuation-target-b", + }).fingerprint; + const requestA = continuation.createRequest(targetA, createInput(f, "parallel-a")).request; + const requestB = continuation.createRequest(targetB, createInput(f, "parallel-b")).request; + + const winner = continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: requestA.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }); + assert.equal(winner.status, "APPROVED"); + assert.equal(f.swarmStore.getWorker(f.worker.id)!.carrierConversationFingerprint, targetA); + + const loser = continuation.getRequest(requestB.id)!; + assert.equal(loser.status, "SUPERSEDED"); + assert.equal(loser.version, 2); + assert.throws( + () => continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: requestB.id, + expectedRequestVersion: 2, + expectedSwarmVersion: 2, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "OWNERSHIP_CONFLICT", + ); + + assert.equal(continuation.recoverAfterRestart(), 0); + assert.equal(continuation.getRequest(requestB.id)?.status, "SUPERSEDED"); + assert.equal(f.swarmStore.getWorker(f.worker.id)!.continuationEpoch, 1); + assert.equal(f.swarmStore.getWorker(f.worker.id)!.carrierConversationFingerprint, targetA); + } finally { + cleanup(f, continuation); + } +}); + +test("owner approval atomically transfers epoch and fences the old carrier", () => { + const f = fixture(); + const continuation = new ChatSwarmContinuationStore(f.root); + try { + const created = continuation.createRequest(f.targetFingerprint, createInput(f)); + const approved = continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }); + assert.equal(approved.status, "APPROVED"); + assert.equal(approved.version, 2); + assert.equal(approved.targetEpoch, 1); + assert.ok(approved.approvedAt); + + const rebound = f.swarmStore.getWorker(f.worker.id)!; + assert.equal(rebound.continuationEpoch, 1); + assert.equal(rebound.carrierConversationFingerprint, f.targetFingerprint); + assert.equal(f.coordinator.peerStatus(f.targetMeta, f.swarm.id).state, "BOUND"); + + assert.throws( + () => f.coordinator.nextTask(f.sourceMeta, f.worker.id), + (error: unknown) => error instanceof ChatSwarmError && error.code === "OWNERSHIP_CONFLICT", + ); + + const replay = continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }); + assert.equal(replay.status, "APPROVED"); + assert.equal(replay.version, 2); + assert.equal(f.swarmStore.getSwarm(f.swarm.id)!.revision, 2); + } finally { + cleanup(f, continuation); + } +}); + +test("checkpoint drift and active work fail closed before epoch transfer", () => { + const f = fixture(); + const continuation = new ChatSwarmContinuationStore(f.root); + try { + const created = continuation.createRequest(f.targetFingerprint, createInput(f)); + f.coordinator.checkpoint( + f.sourceMeta, + f.worker.id, + 0, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { lastTaskId: "task-safe", summary: "changed checkpoint" }, + ); + assert.throws( + () => continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "CAS_DRIFT", + ); + } finally { + cleanup(f, continuation); + } + + const busy = fixture(); + const continuationBusy = new ChatSwarmContinuationStore(busy.root); + try { + busy.coordinator.dispatch(busy.ownerMeta, { + swarmId: busy.swarm.id, + taskKey: "busy-task", + prompt: "bounded read task", + preferredWorkerId: busy.worker.id, + }); + assert.throws( + () => continuationBusy.createRequest(busy.targetFingerprint, createInput(busy)), + (error: unknown) => error instanceof ChatSwarmError && error.code === "RECONCILIATION_REQUIRED", + ); + } finally { + cleanup(busy, continuationBusy); + } +}); + +test("restart marks pending continuation unknown and exact no-effect reconciliation restores same request", () => { + const f = fixture(); + const continuation = new ChatSwarmContinuationStore(f.root); + try { + const created = continuation.createRequest(f.targetFingerprint, createInput(f)); + assert.equal(continuation.recoverAfterRestart(), 1); + const unknown = continuation.getRequest(created.request.id)!; + assert.equal(unknown.status, "RECONCILE_REQUIRED"); + assert.equal(unknown.version, 2); + + assert.throws( + () => continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 2, + expectedSwarmVersion: 1, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "RECONCILIATION_REQUIRED", + ); + + const reconciled = continuation.reconcileUnknownNoEffect(f.ownerFingerprint, created.request.id); + assert.equal(reconciled.status, "PENDING"); + assert.equal(reconciled.version, 3); + assert.equal(reconciled.id, created.request.id); + + const approved = continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 3, + expectedSwarmVersion: 1, + }); + assert.equal(approved.status, "APPROVED"); + assert.equal(approved.version, 4); + } finally { + cleanup(f, continuation); + } +}); + +test("expired continuation cannot transfer and target readback remains bounded to exact target", () => { + const f = fixture(); + let now = new Date("2026-09-13T05:00:00.000Z"); + const continuation = new ChatSwarmContinuationStore(f.root, () => now); + try { + const created = continuation.createRequest(f.targetFingerprint, { + ...createInput(f), + ttlSeconds: 1, + }); + const targetRead = continuation.getLatestForTarget(f.swarm.id, f.targetFingerprint); + assert.equal(targetRead?.id, created.request.id); + + const otherFingerprint = resolveChatSwarmIdentity({ + "openai/conversation_id": "not-the-target", + }).fingerprint; + assert.equal(continuation.getLatestForTarget(f.swarm.id, otherFingerprint), undefined); + + now = new Date("2026-09-13T05:00:02.000Z"); + assert.equal(continuation.getRequest(created.request.id)?.status, "EXPIRED"); + assert.throws( + () => continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "EXPIRED", + ); + const expired = continuation.getRequest(created.request.id)!; + assert.equal(expired.status, "EXPIRED"); + assert.equal(expired.version, 2); + assert.equal(f.swarmStore.getWorker(f.worker.id)!.continuationEpoch, 0); + } finally { + cleanup(f, continuation); + } +}); + +test("persisted continuation attempt identity tamper fails closed", () => { + const f = fixture(); + const continuation = new ChatSwarmContinuationStore(f.root); + try { + const created = continuation.createRequest(f.targetFingerprint, createInput(f, "tamper-attempt")); + const database = openDatabase(f.root); + try { + const row = database.sqlite.prepare( + "select request_json from durable_operations where operation_id=?", + ).get(created.request.id) as { request_json: string }; + const request = JSON.parse(row.request_json) as Record; + request.attemptKey = "tampered-attempt"; + database.sqlite.prepare( + "update durable_operations set request_json=? where operation_id=?", + ).run(JSON.stringify(request), created.request.id); + } finally { + database.close(); + } + + assert.throws( + () => continuation.getRequest(created.request.id), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); + } finally { + cleanup(f, continuation); + } +}); diff --git a/src/chat-swarm-continuation-store.ts b/src/chat-swarm-continuation-store.ts new file mode 100644 index 000000000..95709cf10 --- /dev/null +++ b/src/chat-swarm-continuation-store.ts @@ -0,0 +1,779 @@ +import { createHash } from "node:crypto"; +import { resolve } from "node:path"; +import { openDatabase, type DatabaseHandle } from "./db/client.js"; +import { + assertBounded, + ChatSwarmError, + MAX_ID_BYTES, + newId, +} from "./chat-swarm-contract.js"; +import type { + ApproveContinuationRequestInput, + ChatSwarmContinuationRequest, + CreateContinuationRequestInput, +} from "./chat-swarm-continuation-contract.js"; +import { + assertContinuationCommitAllowed, + prepareContinuationMaterial, + type WorkerContinuationSnapshot, +} from "./chat-swarm-continuation-domain.js"; + +const KIND = "chat_swarm_continuation"; +const REQUEST_SCHEMA = "devspace.chat_swarm_continuation.v1"; +const RECEIPT_SCHEMA = "devspace.chat_swarm_continuation_receipt.v1"; +const SHA256 = /^[0-9a-f]{64}$/; +const DEFAULT_TTL_SECONDS = 15 * 60; +const MAX_TTL_SECONDS = 60 * 60; +const MAX_PENDING_PER_WORKER = 10; +const CONTINUATION_RESTART_MESSAGE = + "DevSpace restarted while continuation was pending; reconcile exact binding before approval."; + +type Row = Record; + +interface DurableRow { + operation_id: string; + attempt_key: string; + request_hash: string; + kind: string; + authority_mode: string; + scope_root: string; + status: string; + request_json: string; + receipt_json: string | null; + error_code: string | null; + error_message: string | null; + created_at: string; + updated_at: string; +} + +interface PersistedRequest { + schema: typeof REQUEST_SCHEMA; + version: number; + swarmId: string; + workerId: string; + attemptKey: string; + sourceEpoch: number; + targetEpoch: number; + sourceCarrierFingerprint: string; + targetCarrierFingerprint: string; + checkpointHash: string; + ttlSeconds: number; + requestedAt: string; + expiresAt: string; +} + +interface PersistedReceipt { + schema: typeof RECEIPT_SCHEMA; + approvedAt: string; + targetEpoch: number; + targetCarrierFingerprint: string; + checkpointHash: string; +} + +type ApprovalOutcome = + | { ok: true; request: ChatSwarmContinuationRequest } + | { + ok: false; + code: "EXPIRED" | "RECONCILIATION_REQUIRED" | "OWNERSHIP_CONFLICT"; + message: string; + }; + +export class ChatSwarmContinuationStore { + private readonly database: DatabaseHandle; + private readonly scopeRoot: string; + + constructor( + stateDir: string, + private readonly clock: () => Date = () => new Date(), + ) { + this.database = openDatabase(stateDir); + this.scopeRoot = resolve(stateDir); + } + + close(): void { + this.database.close(); + } + + createRequest( + authenticatedTargetCarrierFingerprint: string, + input: CreateContinuationRequestInput, + ): { request: ChatSwarmContinuationRequest; created: boolean } { + assertFingerprint(authenticatedTargetCarrierFingerprint, "authenticated target carrier fingerprint"); + assertBounded(input.swarmId, MAX_ID_BYTES, "swarmId"); + assertBounded(input.workerId, MAX_ID_BYTES, "workerId"); + assertBounded(input.attemptKey, MAX_ID_BYTES, "attemptKey"); + const ttlSeconds = input.ttlSeconds ?? DEFAULT_TTL_SECONDS; + if (!Number.isInteger(ttlSeconds) || ttlSeconds < 1 || ttlSeconds > MAX_TTL_SECONDS) { + throw new ChatSwarmError("INVALID_INPUT", `ttlSeconds must be between 1 and ${MAX_TTL_SECONDS}`); + } + + const tx = this.database.sqlite.transaction(() => { + const durableAttemptKey = stableAttemptKey(input.swarmId, input.workerId, input.attemptKey); + const existing = this.getDurableByAttempt(durableAttemptKey); + if (existing) { + if (existing.kind !== KIND) { + throw new ChatSwarmError("REPLAY_CONFLICT", "continuation attemptKey is bound to different material"); + } + const persisted = readPersistedRequest(existing); + const replayHash = continuationMaterialHash({ + ...persisted, + swarmId: input.swarmId, + workerId: input.workerId, + attemptKey: input.attemptKey, + sourceEpoch: input.sourceEpoch, + targetEpoch: input.sourceEpoch + 1, + targetCarrierFingerprint: authenticatedTargetCarrierFingerprint, + ttlSeconds, + }); + if (existing.request_hash !== replayHash) { + throw new ChatSwarmError("REPLAY_CONFLICT", "continuation attemptKey is bound to different material"); + } + const replay = this.toRequest(existing); + if (replay.status === "EXPIRED" && (existing.status === "started" || existing.status === "outcome_unknown")) { + this.persistExpired(replay.id, this.nowIso()); + return { request: this.getRequest(replay.id)!, created: false }; + } + return { request: replay, created: false }; + } + + this.requireActiveSwarm(input.swarmId); + const worker = this.workerSnapshot(input.workerId); + const material = prepareContinuationMaterial( + worker, + input, + authenticatedTargetCarrierFingerprint, + ); + this.assertTargetCarrierAvailable(authenticatedTargetCarrierFingerprint, input.workerId); + this.assertTargetCarrierNotRetired(authenticatedTargetCarrierFingerprint, input.workerId); + + const now = this.nowIso(); + this.persistExpiredPendingForWorker(input.workerId, now); + this.assertNoUnresolvedContinuation(input.workerId, now); + const pendingCount = this.listPendingDurableForWorker(input.workerId).length; + if (pendingCount >= MAX_PENDING_PER_WORKER) { + throw new ChatSwarmError("CAPACITY_FULL", "too many pending continuation requests for worker"); + } + + const requestedAt = now; + const expiresAt = new Date(Date.parse(requestedAt) + ttlSeconds * 1000).toISOString(); + const request: PersistedRequest = { + schema: REQUEST_SCHEMA, + version: 1, + swarmId: input.swarmId, + workerId: input.workerId, + attemptKey: input.attemptKey, + sourceEpoch: material.sourceEpoch, + targetEpoch: material.targetEpoch, + sourceCarrierFingerprint: material.sourceCarrierFingerprint, + targetCarrierFingerprint: material.targetCarrierFingerprint, + checkpointHash: material.checkpointHash, + ttlSeconds, + requestedAt, + expiresAt, + }; + const requestHash = continuationMaterialHash(request); + const operationId = newId("continuation"); + this.database.sqlite.prepare(` + insert into durable_operations ( + operation_id,attempt_key,request_hash,kind,authority_mode,scope_root, + workspace_id,status,retry_safe,request_json,receipt_json,error_code, + error_message,created_at,updated_at + ) values (?,?,?,?,?,?,null,'started','false',?,null,null,null,?,?) + `).run( + operationId, + durableAttemptKey, + requestHash, + KIND, + "OWNER_DIRECT", + this.scopeRoot, + JSON.stringify(request), + requestedAt, + requestedAt, + ); + return { request: this.getRequest(operationId)!, created: true }; + }); + + return tx.immediate(); + } + + approveRequest( + ownerIdentityFingerprint: string, + input: ApproveContinuationRequestInput, + ): ChatSwarmContinuationRequest { + assertFingerprint(ownerIdentityFingerprint, "owner identity fingerprint"); + assertBounded(input.swarmId, MAX_ID_BYTES, "swarmId"); + assertBounded(input.requestId, MAX_ID_BYTES, "requestId"); + if (!Number.isInteger(input.expectedRequestVersion) || input.expectedRequestVersion < 1) { + throw new ChatSwarmError("INVALID_INPUT", "expectedRequestVersion must be a positive integer"); + } + if (!Number.isInteger(input.expectedSwarmVersion) || input.expectedSwarmVersion < 1) { + throw new ChatSwarmError("INVALID_INPUT", "expectedSwarmVersion must be a positive integer"); + } + + const tx = this.database.sqlite.transaction((): ApprovalOutcome => { + const durable = this.requireDurable(input.requestId); + const request = this.toRequest(durable); + if (request.swarmId !== input.swarmId) { + throw new ChatSwarmError("OWNERSHIP_CONFLICT", "continuation request belongs to another swarm"); + } + + const swarm = this.requireActiveSwarm(request.swarmId); + if (String(swarm.owner_identity_fingerprint) !== ownerIdentityFingerprint) { + throw new ChatSwarmError("OWNERSHIP_CONFLICT", "controller identity does not own swarm"); + } + + if (request.status === "APPROVED") { + this.assertCommittedBinding(request); + return { ok: true, request }; + } + if (request.status === "EXPIRED") { + this.persistExpired(input.requestId, this.nowIso()); + return { ok: false, code: "EXPIRED", message: "continuation request has expired" }; + } + if (request.status === "SUPERSEDED") { + return { + ok: false, + code: "OWNERSHIP_CONFLICT", + message: "another continuation target already advanced this worker epoch", + }; + } + if (request.status === "RECONCILE_REQUIRED") { + return { + ok: false, + code: "RECONCILIATION_REQUIRED", + message: "continuation outcome requires explicit reconciliation", + }; + } + if (request.version !== input.expectedRequestVersion) { + throw new ChatSwarmError( + "VERSION_CONFLICT", + `expected request version ${input.expectedRequestVersion} but found ${request.version}`, + ); + } + + const now = this.nowIso(); + if (Date.parse(request.expiresAt) <= Date.parse(now)) { + this.persistExpired(request.id, now); + return { ok: false, code: "EXPIRED", message: "continuation request has expired" }; + } + + const swarmRevision = Number(swarm.revision ?? 1); + if (swarmRevision !== input.expectedSwarmVersion) { + throw new ChatSwarmError( + "CAS_DRIFT", + `swarm revision mismatch: expected ${input.expectedSwarmVersion} but found ${swarmRevision}`, + ); + } + + const worker = this.workerSnapshot(request.workerId); + assertContinuationCommitAllowed(worker, request); + this.assertTargetCarrierAvailable(request.targetCarrierFingerprint, request.workerId); + this.assertTargetCarrierNotRetired(request.targetCarrierFingerprint, request.workerId); + + const workerUpdate = this.database.sqlite.prepare(` + update chat_swarm_workers + set carrier_conversation_fingerprint=?,continuation_epoch=?,lease_json=null,updated_at=? + where id=? and swarm_id=? and continuation_epoch=? + and carrier_conversation_fingerprint=? and lifecycle_state='AVAILABLE' + and current_task_id is null + `).run( + request.targetCarrierFingerprint, + request.targetEpoch, + now, + request.workerId, + request.swarmId, + request.sourceEpoch, + request.sourceCarrierFingerprint, + ); + if (workerUpdate.changes !== 1) { + throw new ChatSwarmError("CAS_DRIFT", "worker binding changed during continuation transfer"); + } + + const receipt: PersistedReceipt = { + schema: RECEIPT_SCHEMA, + approvedAt: now, + targetEpoch: request.targetEpoch, + targetCarrierFingerprint: request.targetCarrierFingerprint, + checkpointHash: request.checkpointHash, + }; + const nextRequest = persistedFromRequest(request, request.version + 1); + const operationUpdate = this.database.sqlite.prepare(` + update durable_operations + set status='succeeded',retry_safe='false',request_json=?,receipt_json=?,error_code=null, + error_message=null,updated_at=? + where operation_id=? and kind=? and status='started' + `).run( + JSON.stringify(nextRequest), + JSON.stringify(receipt), + now, + request.id, + KIND, + ); + if (operationUpdate.changes !== 1) { + throw new ChatSwarmError("VERSION_CONFLICT", "continuation request changed during transfer"); + } + + const swarmUpdate = this.database.sqlite.prepare(` + update chat_swarms set revision=revision+1,updated_at=? + where id=? and revision=? and status='ACTIVE' + `).run(now, request.swarmId, swarmRevision); + if (swarmUpdate.changes !== 1) { + throw new ChatSwarmError("CAS_DRIFT", "swarm revision changed during continuation transfer"); + } + + this.supersedeCompetingRequests(request, now); + return { ok: true, request: this.getRequest(request.id)! }; + }); + + const outcome = tx.immediate(); + if (!outcome.ok) throw new ChatSwarmError(outcome.code, outcome.message); + return outcome.request; + } + + getRequest(requestId: string): ChatSwarmContinuationRequest | undefined { + const row = this.database.sqlite.prepare( + "select * from durable_operations where operation_id=? and kind=? limit 1", + ).get(requestId, KIND) as DurableRow | undefined; + return row ? this.toRequest(row) : undefined; + } + + getLatestForTarget( + swarmId: string, + authenticatedTargetCarrierFingerprint: string, + ): ChatSwarmContinuationRequest | undefined { + assertFingerprint(authenticatedTargetCarrierFingerprint, "authenticated target carrier fingerprint"); + const rows = this.database.sqlite.prepare( + "select * from durable_operations where kind=? and scope_root=? order by created_at desc,operation_id desc", + ).all(KIND, this.scopeRoot) as DurableRow[]; + for (const row of rows) { + const request = this.toRequest(row); + if ( + request.swarmId === swarmId && + request.targetCarrierFingerprint === authenticatedTargetCarrierFingerprint + ) return request; + } + return undefined; + } + + recoverAfterRestart(): number { + const now = this.nowIso(); + const tx = this.database.sqlite.transaction(() => { + const rows = this.database.sqlite.prepare(` + select * from durable_operations + where kind=? and scope_root=? and status in ('started','outcome_unknown') + `).all(KIND, this.scopeRoot) as DurableRow[]; + let changed = 0; + for (const row of rows) { + const persisted = readPersistedRequest(row); + if (Date.parse(persisted.expiresAt) <= Date.parse(now)) { + this.persistExpired(row.operation_id, now); + changed += 1; + continue; + } + if (row.status === "outcome_unknown" && row.error_message === CONTINUATION_RESTART_MESSAGE) { + continue; + } + const request = this.toRequest(row); + const nextRequest = persistedFromRequest(request, request.version + 1); + const result = this.database.sqlite.prepare(` + update durable_operations + set status='outcome_unknown',retry_safe='false',request_json=?, + error_code='RECONCILIATION_REQUIRED',error_message=?,updated_at=? + where operation_id=? and kind=? and status in ('started','outcome_unknown') + `).run( + JSON.stringify(nextRequest), + CONTINUATION_RESTART_MESSAGE, + now, + row.operation_id, + KIND, + ); + changed += Number(result.changes); + } + return changed; + }); + return tx.immediate(); + } + + reconcileUnknownNoEffect( + ownerIdentityFingerprint: string, + requestId: string, + ): ChatSwarmContinuationRequest { + assertFingerprint(ownerIdentityFingerprint, "owner identity fingerprint"); + const tx = this.database.sqlite.transaction((): ApprovalOutcome => { + const durable = this.requireDurable(requestId); + const request = this.toRequest(durable); + if (request.status !== "RECONCILE_REQUIRED") { + throw new ChatSwarmError("INVALID_STATE", "continuation request does not require reconciliation"); + } + const swarm = this.requireActiveSwarm(request.swarmId); + if (String(swarm.owner_identity_fingerprint) !== ownerIdentityFingerprint) { + throw new ChatSwarmError("OWNERSHIP_CONFLICT", "controller identity does not own swarm"); + } + const now = this.nowIso(); + if (Date.parse(request.expiresAt) <= Date.parse(now)) { + this.persistExpired(request.id, now); + return { ok: false, code: "EXPIRED", message: "continuation request expired during reconciliation" }; + } + const worker = this.workerSnapshot(request.workerId); + const pendingView: ChatSwarmContinuationRequest = { ...request, status: "PENDING" }; + assertContinuationCommitAllowed(worker, pendingView); + const nextRequest = persistedFromRequest(request, request.version + 1); + const result = this.database.sqlite.prepare(` + update durable_operations + set status='started',request_json=?,error_code=null,error_message=null,updated_at=? + where operation_id=? and kind=? and status='outcome_unknown' + `).run(JSON.stringify(nextRequest), now, request.id, KIND); + if (result.changes !== 1) { + throw new ChatSwarmError("CAS_DRIFT", "continuation reconciliation changed concurrently"); + } + return { ok: true, request: this.getRequest(request.id)! }; + }); + + const outcome = tx.immediate(); + if (!outcome.ok) throw new ChatSwarmError(outcome.code, outcome.message); + return outcome.request; + } + + private requireActiveSwarm(swarmId: string): Row { + const row = this.database.sqlite.prepare( + "select * from chat_swarms where id=? and status='ACTIVE'", + ).get(swarmId) as Row | undefined; + if (!row) throw new ChatSwarmError("NOT_FOUND", "swarm not found or not active"); + return row; + } + + private workerSnapshot(workerId: string): WorkerContinuationSnapshot { + const row = this.database.sqlite.prepare( + "select * from chat_swarm_workers where id=?", + ).get(workerId) as Row | undefined; + if (!row) throw new ChatSwarmError("NOT_FOUND", "worker not found"); + let checkpoint: Record | undefined; + if (row.checkpoint_json != null) { + try { + const parsed = JSON.parse(String(row.checkpoint_json)); + if (!parsed || Array.isArray(parsed) || typeof parsed !== "object") throw new Error("bad checkpoint"); + checkpoint = parsed as Record; + } catch { + throw new ChatSwarmError("INVALID_STATE", "corrupt persisted worker checkpoint"); + } + } + return { + swarmId: String(row.swarm_id), + workerId: String(row.id), + lifecycleState: String(row.lifecycle_state) as WorkerContinuationSnapshot["lifecycleState"], + currentTaskId: row.current_task_id == null ? undefined : String(row.current_task_id), + continuationEpoch: Number(row.continuation_epoch), + carrierConversationFingerprint: row.carrier_conversation_fingerprint == null + ? undefined + : String(row.carrier_conversation_fingerprint), + checkpoint, + }; + } + + private assertTargetCarrierAvailable(targetCarrierFingerprint: string, workerId: string): void { + const row = this.database.sqlite.prepare(` + select id from chat_swarm_workers + where carrier_conversation_fingerprint=? and lifecycle_state <> 'DISABLED' + limit 1 + `).get(targetCarrierFingerprint) as Row | undefined; + if (row && String(row.id) !== workerId) { + throw new ChatSwarmError("OWNERSHIP_CONFLICT", "target carrier is already bound to another worker"); + } + } + + private assertTargetCarrierNotRetired(targetCarrierFingerprint: string, workerId: string): void { + const rows = this.database.sqlite.prepare( + "select * from durable_operations where kind=? and scope_root=? and status='succeeded'", + ).all(KIND, this.scopeRoot) as DurableRow[]; + for (const row of rows) { + const request = this.toRequest(row); + if ( + request.workerId === workerId && + request.sourceCarrierFingerprint === targetCarrierFingerprint + ) { + throw new ChatSwarmError( + "OWNERSHIP_CONFLICT", + "target carrier is retired by an accepted continuation and cannot be rebound", + ); + } + } + } + + private getDurableByAttempt(durableAttemptKey: string): DurableRow | undefined { + return this.database.sqlite.prepare( + "select * from durable_operations where scope_root=? and attempt_key=? limit 1", + ).get(this.scopeRoot, durableAttemptKey) as DurableRow | undefined; + } + + private listPendingDurableForWorker(workerId: string): DurableRow[] { + const rows = this.database.sqlite.prepare( + "select * from durable_operations where kind=? and scope_root=? and status='started'", + ).all(KIND, this.scopeRoot) as DurableRow[]; + return rows.filter((row) => this.toRequest(row).workerId === workerId); + } + + private assertNoUnresolvedContinuation(workerId: string, now: string): void { + const rows = this.database.sqlite.prepare( + "select * from durable_operations where kind=? and scope_root=? and status='outcome_unknown'", + ).all(KIND, this.scopeRoot) as DurableRow[]; + for (const row of rows) { + const request = this.toRequest(row); + if (request.workerId !== workerId) continue; + if (Date.parse(request.expiresAt) <= Date.parse(now)) { + this.persistExpired(request.id, now); + continue; + } + throw new ChatSwarmError( + "RECONCILIATION_REQUIRED", + "worker has an unresolved continuation outcome; reconcile it before creating another replacement", + ); + } + } + + private persistExpiredPendingForWorker(workerId: string, now: string): void { + for (const row of this.listPendingDurableForWorker(workerId)) { + const request = this.toRequest(row); + if (Date.parse(request.expiresAt) <= Date.parse(now)) this.persistExpired(row.operation_id, now); + } + } + + private supersedeCompetingRequests( + winner: ChatSwarmContinuationRequest, + now: string, + ): void { + for (const row of this.listPendingDurableForWorker(winner.workerId)) { + if (row.operation_id === winner.id) continue; + const other = this.toRequest(row); + if ( + other.sourceEpoch !== winner.sourceEpoch || + other.sourceCarrierFingerprint !== winner.sourceCarrierFingerprint + ) continue; + if (Date.parse(other.expiresAt) <= Date.parse(now)) { + this.persistExpired(other.id, now); + continue; + } + const nextRequest = persistedFromRequest(other, other.version + 1); + this.database.sqlite.prepare(` + update durable_operations + set status='failed',retry_safe='false',request_json=?,error_code='SUPERSEDED', + error_message='Another continuation target atomically advanced this worker epoch.',updated_at=? + where operation_id=? and kind=? and status='started' + `).run(JSON.stringify(nextRequest), now, other.id, KIND); + } + } + + private requireDurable(requestId: string): DurableRow { + const row = this.database.sqlite.prepare( + "select * from durable_operations where operation_id=? and kind=? limit 1", + ).get(requestId, KIND) as DurableRow | undefined; + if (!row) throw new ChatSwarmError("REQUEST_NOT_FOUND", "continuation request not found"); + return row; + } + + private persistExpired(requestId: string, now: string): void { + const durable = this.requireDurable(requestId); + if (durable.status !== "started" && durable.status !== "outcome_unknown") return; + const request = this.toRequest(durable); + const nextRequest = persistedFromRequest(request, request.version + 1); + this.database.sqlite.prepare(` + update durable_operations + set status='failed',retry_safe='false',request_json=?,error_code='EXPIRED', + error_message='Continuation request expired before transfer.',updated_at=? + where operation_id=? and kind=? and status in ('started','outcome_unknown') + `).run(JSON.stringify(nextRequest), now, requestId, KIND); + } + + private assertCommittedBinding(request: ChatSwarmContinuationRequest): void { + const worker = this.workerSnapshot(request.workerId); + if ( + worker.continuationEpoch !== request.targetEpoch || + worker.carrierConversationFingerprint !== request.targetCarrierFingerprint + ) { + throw new ChatSwarmError("RECONCILIATION_REQUIRED", "approved continuation does not match current worker binding"); + } + } + + private toRequest(row: DurableRow): ChatSwarmContinuationRequest { + if (row.kind !== KIND || row.authority_mode !== "OWNER_DIRECT") { + throw new ChatSwarmError("INVALID_STATE", "durable continuation operation authority is malformed"); + } + const persisted = readPersistedRequest(row); + if (row.scope_root !== this.scopeRoot) { + throw new ChatSwarmError("INVALID_STATE", "durable continuation scope root mismatch"); + } + if (row.created_at !== persisted.requestedAt) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation creation timestamp mismatch"); + } + if (stableAttemptKey(persisted.swarmId, persisted.workerId, persisted.attemptKey) !== row.attempt_key) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation attempt identity mismatch"); + } + const requestHash = continuationMaterialHash(persisted); + if (requestHash !== row.request_hash) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation request hash mismatch"); + } + + let status: ChatSwarmContinuationRequest["status"]; + if (row.status === "started") { + if (row.receipt_json !== null || row.error_code !== null || row.error_message !== null) { + throw new ChatSwarmError("INVALID_STATE", "pending continuation terminal metadata is malformed"); + } + status = Date.parse(persisted.expiresAt) <= this.clock().getTime() ? "EXPIRED" : "PENDING"; + } else if (row.status === "succeeded") { + if (row.receipt_json === null || row.error_code !== null || row.error_message !== null) { + throw new ChatSwarmError("INVALID_STATE", "approved continuation terminal metadata is malformed"); + } + status = "APPROVED"; + } else if (row.status === "failed" && (row.error_code === "EXPIRED" || row.error_code === "SUPERSEDED")) { + if (row.receipt_json !== null || !row.error_message) { + throw new ChatSwarmError("INVALID_STATE", "failed continuation terminal metadata is malformed"); + } + status = row.error_code === "EXPIRED" ? "EXPIRED" : "SUPERSEDED"; + } else if (row.status === "outcome_unknown") { + if ( + row.receipt_json !== null || + row.error_code !== "RECONCILIATION_REQUIRED" || + !row.error_message + ) { + throw new ChatSwarmError("INVALID_STATE", "unresolved continuation terminal metadata is malformed"); + } + status = "RECONCILE_REQUIRED"; + } else { + throw new ChatSwarmError("INVALID_STATE", `unsupported durable continuation status '${row.status}'`); + } + + let approvedAt: string | undefined; + if (status === "APPROVED") { + try { + const receipt = JSON.parse(row.receipt_json!) as PersistedReceipt; + if (receipt.schema !== RECEIPT_SCHEMA) throw new Error("receipt schema mismatch"); + if ( + receipt.targetEpoch !== persisted.targetEpoch || + receipt.targetCarrierFingerprint !== persisted.targetCarrierFingerprint || + receipt.checkpointHash !== persisted.checkpointHash || + !Number.isFinite(Date.parse(receipt.approvedAt)) + ) throw new Error("receipt binding mismatch"); + approvedAt = receipt.approvedAt; + } catch { + throw new ChatSwarmError("INVALID_STATE", "corrupt persisted continuation receipt"); + } + } + + return { + id: row.operation_id, + swarmId: persisted.swarmId, + workerId: persisted.workerId, + attemptKey: persisted.attemptKey, + requestHash: row.request_hash, + sourceEpoch: persisted.sourceEpoch, + targetEpoch: persisted.targetEpoch, + sourceCarrierFingerprint: persisted.sourceCarrierFingerprint, + targetCarrierFingerprint: persisted.targetCarrierFingerprint, + checkpointHash: persisted.checkpointHash, + ttlSeconds: persisted.ttlSeconds, + version: persisted.version, + status, + requestedAt: persisted.requestedAt, + expiresAt: persisted.expiresAt, + approvedAt, + }; + } + + private nowIso(): string { + return this.clock().toISOString(); + } +} + +function stableAttemptKey(swarmId: string, workerId: string, attemptKey: string): string { + return `continuation:${createHash("sha256") + .update(JSON.stringify({ swarmId, workerId, attemptKey })) + .digest("hex")}`; +} + +function continuationMaterialHash(request: PersistedRequest): string { + return createHash("sha256").update(JSON.stringify({ + attemptKey: request.attemptKey, + checkpointHash: request.checkpointHash, + expiresAt: request.expiresAt, + requestedAt: request.requestedAt, + sourceCarrierFingerprint: request.sourceCarrierFingerprint, + sourceEpoch: request.sourceEpoch, + swarmId: request.swarmId, + targetCarrierFingerprint: request.targetCarrierFingerprint, + targetEpoch: request.targetEpoch, + ttlSeconds: request.ttlSeconds, + workerId: request.workerId, + })).digest("hex"); +} + +function persistedFromRequest( + request: ChatSwarmContinuationRequest, + version: number, +): PersistedRequest { + return { + schema: REQUEST_SCHEMA, + version, + swarmId: request.swarmId, + workerId: request.workerId, + attemptKey: request.attemptKey, + sourceEpoch: request.sourceEpoch, + targetEpoch: request.targetEpoch, + sourceCarrierFingerprint: request.sourceCarrierFingerprint, + targetCarrierFingerprint: request.targetCarrierFingerprint, + checkpointHash: request.checkpointHash, + ttlSeconds: request.ttlSeconds, + requestedAt: request.requestedAt, + expiresAt: request.expiresAt, + }; +} + +function readPersistedRequest(row: DurableRow): PersistedRequest { + let request: PersistedRequest; + try { + request = JSON.parse(row.request_json) as PersistedRequest; + } catch { + throw new ChatSwarmError("INVALID_STATE", "corrupt persisted continuation request"); + } + validatePersistedRequest(request); + return request; +} + +function validatePersistedRequest(request: PersistedRequest): void { + if (!request || request.schema !== REQUEST_SCHEMA) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation request schema mismatch"); + } + if (!Number.isSafeInteger(request.version) || request.version < 1) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation request version is malformed"); + } + if (!Number.isInteger(request.ttlSeconds) || request.ttlSeconds < 1 || request.ttlSeconds > MAX_TTL_SECONDS) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation ttl is malformed"); + } + for (const [label, value] of [ + ["source carrier fingerprint", request.sourceCarrierFingerprint], + ["target carrier fingerprint", request.targetCarrierFingerprint], + ["checkpoint hash", request.checkpointHash], + ] as const) assertFingerprint(value, label); + if ( + !Number.isSafeInteger(request.sourceEpoch) || + request.sourceEpoch < 0 || + request.targetEpoch !== request.sourceEpoch + 1 + ) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation epoch binding is malformed"); + } + if (!request.swarmId || !request.workerId || !request.attemptKey) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation identity is incomplete"); + } + const requestedMs = Date.parse(request.requestedAt); + const expiresMs = Date.parse(request.expiresAt); + if (!Number.isFinite(requestedMs) || !Number.isFinite(expiresMs)) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation timestamps are malformed"); + } + if (expiresMs - requestedMs !== request.ttlSeconds * 1000) { + throw new ChatSwarmError("INVALID_STATE", "persisted continuation expiry does not match ttl"); + } +} + +function assertFingerprint(value: string, label: string): void { + if (!SHA256.test(value)) { + throw new ChatSwarmError("INVALID_INPUT", `${label} must be a SHA-256 fingerprint`); + } +} diff --git a/src/chat-swarm-continuation-terminal-replay.test.ts b/src/chat-swarm-continuation-terminal-replay.test.ts new file mode 100644 index 000000000..de1a17810 --- /dev/null +++ b/src/chat-swarm-continuation-terminal-replay.test.ts @@ -0,0 +1,134 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { ChatSwarmError } from "./chat-swarm-contract.js"; +import { ChatSwarmContinuationStore } from "./chat-swarm-continuation-store.js"; +import { ChatSwarmCoordinator } from "./chat-swarm-coordinator.js"; +import { ChatSwarmStore } from "./chat-swarm-store.js"; +import { resolveChatSwarmIdentity } from "./request-meta.js"; + +function fixture(clock?: () => Date) { + const root = mkdtempSync(join(tmpdir(), "swarm-continuation-terminal-replay-")); + const swarmStore = new ChatSwarmStore(root); + const continuation = new ChatSwarmContinuationStore(root, clock); + const coordinator = new ChatSwarmCoordinator(swarmStore); + const ownerMeta = { "openai/session": "continuation-owner" }; + const sourceMeta = { "openai/conversation_id": "continuation-source" }; + const targetA = resolveChatSwarmIdentity({ "openai/conversation_id": "continuation-target-a" }).fingerprint; + const targetB = resolveChatSwarmIdentity({ "openai/conversation_id": "continuation-target-b" }).fingerprint; + const ownerFingerprint = resolveChatSwarmIdentity(ownerMeta).fingerprint; + const swarm = coordinator.createSwarm(ownerMeta, { workerLimit: 2, metadata: { test: true } }); + const worker = coordinator.joinWorker(sourceMeta, swarm.id, { + label: "Worker-01", + runtimeKind: "mcp_peer", + }); + coordinator.checkpoint( + sourceMeta, + worker.id, + 0, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { summary: "safe checkpoint" }, + ); + return { root, swarmStore, continuation, coordinator, ownerFingerprint, swarm, worker, targetA, targetB }; +} + +function cleanup(f: ReturnType) { + f.continuation.close(); + f.swarmStore.close(); + rmSync(f.root, { recursive: true, force: true }); +} + +test("approved continuation exact replay returns the same durable result after epoch advances", () => { + const f = fixture(); + try { + const input = { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey: "approved-replay", + sourceEpoch: 0, + }; + const created = f.continuation.createRequest(f.targetA, input); + const approved = f.continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: created.request.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }); + assert.equal(approved.status, "APPROVED"); + + const replay = f.continuation.createRequest(f.targetA, input); + assert.equal(replay.created, false); + assert.equal(replay.request.id, created.request.id); + assert.equal(replay.request.status, "APPROVED"); + assert.equal(replay.request.version, approved.version); + + assert.throws( + () => f.continuation.createRequest(f.targetA, { ...input, ttlSeconds: 1 }), + (error: unknown) => error instanceof ChatSwarmError && error.code === "REPLAY_CONFLICT", + ); + } finally { + cleanup(f); + } +}); + +test("superseded continuation exact replay returns the same terminal loser", () => { + const f = fixture(); + try { + const inputA = { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey: "winner-replay", + sourceEpoch: 0, + }; + const inputB = { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey: "loser-replay", + sourceEpoch: 0, + }; + const winner = f.continuation.createRequest(f.targetA, inputA).request; + const loser = f.continuation.createRequest(f.targetB, inputB).request; + f.continuation.approveRequest(f.ownerFingerprint, { + swarmId: f.swarm.id, + requestId: winner.id, + expectedRequestVersion: 1, + expectedSwarmVersion: 1, + }); + assert.equal(f.continuation.getRequest(loser.id)?.status, "SUPERSEDED"); + + const replay = f.continuation.createRequest(f.targetB, inputB); + assert.equal(replay.created, false); + assert.equal(replay.request.id, loser.id); + assert.equal(replay.request.status, "SUPERSEDED"); + } finally { + cleanup(f); + } +}); + +test("expired exact replay durably persists terminal expiry", () => { + let now = new Date("2026-09-13T06:00:00.000Z"); + const f = fixture(() => now); + try { + const input = { + swarmId: f.swarm.id, + workerId: f.worker.id, + attemptKey: "expired-replay", + sourceEpoch: 0, + ttlSeconds: 1, + }; + const created = f.continuation.createRequest(f.targetA, input); + now = new Date("2026-09-13T06:00:02.000Z"); + + const replay = f.continuation.createRequest(f.targetA, input); + assert.equal(replay.created, false); + assert.equal(replay.request.id, created.request.id); + assert.equal(replay.request.status, "EXPIRED"); + assert.equal(replay.request.version, 2); + assert.equal(f.continuation.getRequest(created.request.id)?.status, "EXPIRED"); + assert.equal(f.continuation.recoverAfterRestart(), 0); + } finally { + cleanup(f); + } +}); diff --git a/src/chat-swarm-lifecycle.test.ts b/src/chat-swarm-lifecycle.test.ts index de43fba06..f2d28a0bb 100644 --- a/src/chat-swarm-lifecycle.test.ts +++ b/src/chat-swarm-lifecycle.test.ts @@ -6,6 +6,7 @@ import test from "node:test"; import { ChatSwarmError } from "./chat-swarm-contract.js"; import { ChatSwarmLifecycle, type ChatSwarmLifecycleMode } from "./chat-swarm-lifecycle.js"; import type { ChatSwarmCarrierAdapter } from "./chat-swarm-carrier.js"; +import { DurableOperationStore } from "./durable-operations.js"; function fixture() { const root = mkdtempSync(join(tmpdir(), "devspace-chat-swarm-lifecycle-")); @@ -42,16 +43,92 @@ test("shares one coordinator and keeps startup recovery explicit", () => { } }); +test("startup recovery normalizes continuation after generic durable-operation fencing", () => { + const f = fixture(); + try { + const owner = { "openai/session": "continuation-owner" }; + const source = { "openai/conversation_id": "continuation-source" }; + const target = { "openai/conversation_id": "continuation-target" }; + const swarm = f.lifecycle.coordinator!.createSwarm(owner, { workerLimit: 2 }); + const worker = f.lifecycle.coordinator!.joinWorker(source, swarm.id, { + label: "peer", + runtimeKind: "mcp_peer", + }); + f.lifecycle.coordinator!.checkpoint( + source, + worker.id, + 0, + new Date(Date.now() + 60 * 60 * 1000).toISOString(), + { summary: "safe checkpoint" }, + ); + const request = f.lifecycle.continuationCoordinator!.request(target, { + swarmId: swarm.id, + workerId: worker.id, + attemptKey: "restart-fence", + sourceEpoch: 0, + }).request; + assert.equal(request.status, "PENDING"); + assert.equal(request.version, 1); + + const genericDurableStore = new DurableOperationStore(f.root); + try { + assert.equal(genericDurableStore.markInterruptedUnknown(), 1); + } finally { + genericDurableStore.close(); + } + const genericFenced = f.lifecycle.continuationStore!.getRequest(request.id)!; + assert.equal(genericFenced.status, "RECONCILE_REQUIRED"); + assert.equal(genericFenced.version, 1); + + const second = new ChatSwarmLifecycle({ stateDir: f.root }); + try { + assert.equal(second.recoverAfterStartup(), 0); + const recovered = f.lifecycle.continuationStore!.getRequest(request.id)!; + assert.equal(recovered.status, "RECONCILE_REQUIRED"); + assert.equal(recovered.version, 2); + assert.equal(f.lifecycle.store!.getWorker(worker.id)!.continuationEpoch, 0); + assert.equal( + f.lifecycle.store!.getWorker(worker.id)!.carrierConversationFingerprint, + worker.carrierConversationFingerprint, + ); + assert.equal(second.continuationCoordinator!.recoverAfterRestart(), 0); + assert.equal(f.lifecycle.continuationStore!.getRequest(request.id)!.version, 2); + } finally { + second.close(); + } + } finally { + cleanup(f); + } +}); + test("drain permits current-task next and uncertainty-safe completion, but rejects new claims", () => { const f = fixture(); try { f.setMode("drain"); - for (const action of ["create", "join", "dispatch", "close"] as const) { - assert.throws(() => f.lifecycle.admit(action), (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE"); + for (const action of [ + "create", + "join", + "dispatch", + "close", + "continuation_request", + "continuation_approve", + ] as const) { + assert.throws( + () => f.lifecycle.admit(action), + (error: unknown) => error instanceof ChatSwarmError && error.code === "INVALID_STATE", + ); } assert.doesNotThrow(() => f.lifecycle.admit("next", { existingTask: true })); assert.throws(() => f.lifecycle.admit("next", { existingTask: false }), ChatSwarmError); - for (const action of ["submit", "status", "collect", "cancel", "reconcile"] as const) { + for (const action of [ + "submit", + "status", + "collect", + "cancel", + "reconcile", + "continuation_status", + "continuation_reconcile", + ] as const) { assert.doesNotThrow(() => f.lifecycle.admit(action)); } } finally { @@ -63,10 +140,26 @@ test("reconcile-only permits inspection and explicit recovery actions only", () const f = fixture(); try { f.setMode("reconcile-only"); - for (const action of ["status", "collect", "cancel", "reconcile", "submit"] as const) { + for (const action of [ + "status", + "collect", + "cancel", + "reconcile", + "submit", + "continuation_status", + "continuation_reconcile", + ] as const) { assert.doesNotThrow(() => f.lifecycle.admit(action)); } - for (const action of ["create", "join", "dispatch", "next", "close"] as const) { + for (const action of [ + "create", + "join", + "dispatch", + "next", + "close", + "continuation_request", + "continuation_approve", + ] as const) { assert.throws(() => f.lifecycle.admit(action), ChatSwarmError); } } finally { @@ -80,6 +173,8 @@ test("disabled lifecycle has no store and closes safely", () => { try { assert.equal(lifecycle.store, undefined); assert.equal(lifecycle.coordinator, undefined); + assert.equal(lifecycle.continuationStore, undefined); + assert.equal(lifecycle.continuationCoordinator, undefined); assert.throws(() => lifecycle.recoverAfterStartup(), ChatSwarmError); assert.throws(() => lifecycle.admit("status"), ChatSwarmError); } finally { @@ -89,8 +184,66 @@ test("disabled lifecycle has no store and closes safely", () => { }); test("carrier effects require injected adapter, owner, normal mode, and explicit startup recovery", async () => { - const root = mkdtempSync(join(tmpdir(), "devspace-chat-swarm-carrier-lifecycle-")); let mode: ChatSwarmLifecycleMode = "normal"; - const adapter: ChatSwarmCarrierAdapter = { kind: "fake", capabilities: () => ({ boundedWait: "UNSUPPORTED", eventWake: "UNSUPPORTED", resultReadback: "UNSUPPORTED", durableReplay: "SUPPORTED" }), ensureExisting: async (input) => ({ disposition: "READY", operationId: input.operationId, swarmId: input.swarmId, workerId: input.workerId, expectedEpoch: input.expectedEpoch, carrierKind: input.carrierKind, carrierFingerprint: input.carrierFingerprint, remoteMayContinue: false }), wake: async (input) => ({ disposition: "UNSUPPORTED", operationId: input.operationId, swarmId: input.swarmId, workerId: input.workerId, expectedEpoch: input.expectedEpoch, carrierKind: input.carrierKind, carrierFingerprint: input.carrierFingerprint, remoteMayContinue: false }) }; - const lifecycle = new ChatSwarmLifecycle({ stateDir: root, mode: () => mode, carrierAdapter: adapter }); const owner = { "openai/session": "owner" }; - try { const swarm = lifecycle.coordinator!.createSwarm(owner, { workerLimit: 1 }); lifecycle.store!.createWorker({ swarmId: swarm.id, label: "peer", runtimeKind: "mcp_peer", carrierConversationFingerprint: "a".repeat(64) }); const status = lifecycle.carrierStatus(owner, swarm.id); assert.equal(status.workers.length, 1); mode = "drain"; await assert.rejects(lifecycle.ensureCarriers(owner, { swarmId: swarm.id, capacity: 1, adapterConfigHash: "b".repeat(64) }), /carrier effects are unavailable/); mode = "normal"; assert.equal(lifecycle.recoverAfterStartup(), 0); } finally { lifecycle.close(); rmSync(root, { recursive: true, force: true }); } + const root = mkdtempSync(join(tmpdir(), "devspace-chat-swarm-carrier-lifecycle-")); + let mode: ChatSwarmLifecycleMode = "normal"; + const adapter: ChatSwarmCarrierAdapter = { + kind: "fake", + capabilities: () => ({ + boundedWait: "UNSUPPORTED", + eventWake: "UNSUPPORTED", + resultReadback: "UNSUPPORTED", + durableReplay: "SUPPORTED", + }), + ensureExisting: async (input) => ({ + disposition: "READY", + operationId: input.operationId, + swarmId: input.swarmId, + workerId: input.workerId, + expectedEpoch: input.expectedEpoch, + carrierKind: input.carrierKind, + carrierFingerprint: input.carrierFingerprint, + remoteMayContinue: false, + }), + wake: async (input) => ({ + disposition: "UNSUPPORTED", + operationId: input.operationId, + swarmId: input.swarmId, + workerId: input.workerId, + expectedEpoch: input.expectedEpoch, + carrierKind: input.carrierKind, + carrierFingerprint: input.carrierFingerprint, + remoteMayContinue: false, + }), + }; + const lifecycle = new ChatSwarmLifecycle({ + stateDir: root, + mode: () => mode, + carrierAdapter: adapter, + }); + const owner = { "openai/session": "owner" }; + try { + const swarm = lifecycle.coordinator!.createSwarm(owner, { workerLimit: 1 }); + lifecycle.store!.createWorker({ + swarmId: swarm.id, + label: "peer", + runtimeKind: "mcp_peer", + carrierConversationFingerprint: "a".repeat(64), + }); + const status = lifecycle.carrierStatus(owner, swarm.id); + assert.equal(status.workers.length, 1); + mode = "drain"; + await assert.rejects( + lifecycle.ensureCarriers(owner, { + swarmId: swarm.id, + capacity: 1, + adapterConfigHash: "b".repeat(64), + }), + /carrier effects are unavailable/, + ); + mode = "normal"; + assert.equal(lifecycle.recoverAfterStartup(), 0); + } finally { + lifecycle.close(); + rmSync(root, { recursive: true, force: true }); + } }); diff --git a/src/chat-swarm-lifecycle.ts b/src/chat-swarm-lifecycle.ts index c59df7335..06f21a65d 100644 --- a/src/chat-swarm-lifecycle.ts +++ b/src/chat-swarm-lifecycle.ts @@ -1,7 +1,14 @@ import { ChatSwarmError } from "./chat-swarm-contract.js"; import { ChatSwarmCoordinator } from "./chat-swarm-coordinator.js"; import { ChatSwarmStore } from "./chat-swarm-store.js"; -import { ChatSwarmCarrierManager, type ChatSwarmCarrierAdapter, type CarrierResult, type CarrierStatusResult } from "./chat-swarm-carrier.js"; +import { + ChatSwarmCarrierManager, + type ChatSwarmCarrierAdapter, + type CarrierResult, + type CarrierStatusResult, +} from "./chat-swarm-carrier.js"; +import { ChatSwarmContinuationCoordinator } from "./chat-swarm-continuation-coordinator.js"; +import { ChatSwarmContinuationStore } from "./chat-swarm-continuation-store.js"; export type ChatSwarmLifecycleMode = "normal" | "drain" | "reconcile-only"; export type ChatSwarmLifecycleAction = @@ -19,7 +26,11 @@ export type ChatSwarmLifecycleAction = | "inspect" | "tasks" | "join_request" - | "approve_join"; + | "approve_join" + | "continuation_request" + | "continuation_status" + | "continuation_approve" + | "continuation_reconcile"; export interface ChatSwarmLifecycleOptions { stateDir: string; @@ -33,14 +44,16 @@ export interface ChatSwarmAdmissionContext { } /** - * Owns one process-level Swarm store/coordinator pair. Construction opens the - * durable store but deliberately does not perform restart recovery; the server - * must call recoverAfterRestart once it has established its startup authority. + * Owns one process-level Swarm state domain. Construction opens the durable + * stores but deliberately does not perform restart recovery; the server must + * call recoverAfterStartup once it has established its startup authority. */ export class ChatSwarmLifecycle { readonly enabled: boolean; readonly store?: ChatSwarmStore; readonly coordinator?: ChatSwarmCoordinator; + readonly continuationStore?: ChatSwarmContinuationStore; + readonly continuationCoordinator?: ChatSwarmContinuationCoordinator; private readonly modeProvider: () => ChatSwarmLifecycleMode; private closed = false; private readonly carrierManager?: ChatSwarmCarrierManager; @@ -50,29 +63,91 @@ export class ChatSwarmLifecycle { this.modeProvider = options.mode ?? (() => "normal"); if (this.enabled) { const store = new ChatSwarmStore(options.stateDir); + const coordinator = new ChatSwarmCoordinator(store); + const continuationStore = new ChatSwarmContinuationStore(options.stateDir); this.store = store; - this.coordinator = new ChatSwarmCoordinator(store); - if (options.carrierAdapter) this.carrierManager = new ChatSwarmCarrierManager(store, this.coordinator, options.carrierAdapter); + this.coordinator = coordinator; + this.continuationStore = continuationStore; + this.continuationCoordinator = new ChatSwarmContinuationCoordinator( + continuationStore, + coordinator, + ); + if (options.carrierAdapter) { + this.carrierManager = new ChatSwarmCarrierManager( + store, + coordinator, + options.carrierAdapter, + ); + } } } recoverAfterStartup(): number { this.requireEnabled(); this.store!.fenceCarrierOperations(); - return this.store!.recoverAfterRestart(); + const taskRecoveryCount = this.store!.recoverAfterRestart(); + this.continuationCoordinator!.recoverAfterRestart(); + return taskRecoveryCount; } - async ensureCarriers(meta: unknown, input: { swarmId: string; capacity: number; adapterConfigHash: string }): Promise { this.requireCarrierEffects(); return this.carrierManager!.ensure(meta, input.swarmId, input.capacity, input.adapterConfigHash); } - carrierStatus(meta: unknown, swarmId: string): CarrierStatusResult { this.requireEnabled(); if (!this.carrierManager) throw new ChatSwarmError("INVALID_STATE", "carrier adapter is unavailable"); return this.carrierManager.status(meta, swarmId); } - async carrierWake(meta: unknown, input: { swarmId: string; workerId: string; expectedEpoch: number; taskId?: string; adapterConfigHash: string }): Promise { this.requireCarrierEffects(); return this.carrierManager!.wake(meta, input); } + async ensureCarriers( + meta: unknown, + input: { swarmId: string; capacity: number; adapterConfigHash: string }, + ): Promise { + this.requireCarrierEffects(); + return this.carrierManager!.ensure( + meta, + input.swarmId, + input.capacity, + input.adapterConfigHash, + ); + } + + carrierStatus(meta: unknown, swarmId: string): CarrierStatusResult { + this.requireEnabled(); + if (!this.carrierManager) { + throw new ChatSwarmError("INVALID_STATE", "carrier adapter is unavailable"); + } + return this.carrierManager.status(meta, swarmId); + } + + async carrierWake( + meta: unknown, + input: { + swarmId: string; + workerId: string; + expectedEpoch: number; + taskId?: string; + adapterConfigHash: string; + }, + ): Promise { + this.requireCarrierEffects(); + return this.carrierManager!.wake(meta, input); + } - admit(action: ChatSwarmLifecycleAction, context: ChatSwarmAdmissionContext = {}): void { + admit( + action: ChatSwarmLifecycleAction, + context: ChatSwarmAdmissionContext = {}, + ): void { this.requireEnabled(); const mode = this.modeProvider(); if (mode === "normal") return; if (action === "next" && context.existingTask === true) return; - if (["status", "collect", "cancel", "reconcile", "submit", "peer_status", "inspect", "tasks"].includes(action)) return; + if ( + [ + "status", + "collect", + "cancel", + "reconcile", + "submit", + "peer_status", + "inspect", + "tasks", + "continuation_status", + "continuation_reconcile", + ].includes(action) + ) return; throw new ChatSwarmError( "INVALID_STATE", @@ -83,14 +158,35 @@ export class ChatSwarmLifecycle { close(): void { if (this.closed) return; this.closed = true; + this.continuationStore?.close(); this.store?.close(); } private requireEnabled(): void { - if (!this.enabled || !this.store || !this.coordinator) { + if ( + !this.enabled || + !this.store || + !this.coordinator || + !this.continuationStore || + !this.continuationCoordinator + ) { throw new ChatSwarmError("INVALID_STATE", "chat swarm feature is disabled"); } - if (this.closed) throw new ChatSwarmError("INVALID_STATE", "chat swarm lifecycle is closed"); + if (this.closed) { + throw new ChatSwarmError("INVALID_STATE", "chat swarm lifecycle is closed"); + } + } + + private requireCarrierEffects(): void { + this.requireEnabled(); + if (this.modeProvider() !== "normal") { + throw new ChatSwarmError( + "INVALID_STATE", + "carrier effects are unavailable while lifecycle is not normal", + ); + } + if (!this.carrierManager) { + throw new ChatSwarmError("INVALID_STATE", "carrier adapter is unavailable"); + } } - private requireCarrierEffects(): void { this.requireEnabled(); if (this.modeProvider() !== "normal") throw new ChatSwarmError("INVALID_STATE", "carrier effects are unavailable while lifecycle is not normal"); if (!this.carrierManager) throw new ChatSwarmError("INVALID_STATE", "carrier adapter is unavailable"); } }