From ca2eff7ae78d0785f9cf2c971405507437bc6c0e Mon Sep 17 00:00:00 2001 From: Daniel Gordon Date: Mon, 27 Jul 2026 18:55:49 +0000 Subject: [PATCH] fix(orchestration): refresh the engine read model when a command targets an externally appended aggregate The t3 import CLI runs a second in-process OrchestrationEngine over the same SQLite file while the server is running. Its events and projections land correctly, but the running server's in-memory command read model is hydrated once at startup and only folds events the engine itself commits, so imported threads render in the UI (SQL reads) while every command against them - thread.delete included - is rejected by requireThread with "Thread ... does not exist". The post-failure reconcile cannot recover: it reads events above the in-memory watermark, and the engine's own later appends leapfrog the externally appended sequences. On an aggregate miss for non-create commands, re-read the command read model from the SQL snapshot (the same source as startup hydration), clamped so snapshotSequence/updatedAt never rewind. No sequence-based gate: the engine's own commits advance the projection bookkeeping past external sequences, so such a comparison reads "equal" in exactly the failing case. Regression test simulates the external writer by appending a decider-produced thread.created via the event store + projection pipeline directly, advances the watermark with an unrelated dispatch, and asserts a first-attempt thread.delete succeeds. Co-Authored-By: Claude Fable 5 --- .../Layers/OrchestrationEngine.test.ts | 190 +++++++++++++++++- .../Layers/OrchestrationEngine.ts | 53 ++++- 2 files changed, 238 insertions(+), 5 deletions(-) diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 731002ea783..2a8efe2e972 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -28,6 +28,7 @@ import { type OrchestrationEventStoreShape, } from "../../persistence/Services/OrchestrationEventStore.ts"; import * as RepositoryIdentityResolver from "../../project/RepositoryIdentityResolver.ts"; +import { decideOrchestrationCommand } from "../decider.ts"; import { OrchestrationEngineLive } from "./OrchestrationEngine.ts"; import { OrchestrationProjectionPipelineLive } from "./ProjectionPipeline.ts"; import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts"; @@ -54,8 +55,9 @@ async function createOrchestrationSystem() { Layer.provide(OrchestrationProjectionPipelineLive), ), OrchestrationProjectionSnapshotQueryLive, + OrchestrationProjectionPipelineLive, ).pipe( - Layer.provide(OrchestrationEventStoreLive), + Layer.provideMerge(OrchestrationEventStoreLive), Layer.provide(OrchestrationCommandReceiptRepositoryLive), Layer.provide(RepositoryIdentityResolver.layer), Layer.provide(SqlitePersistenceMemory), @@ -65,8 +67,15 @@ async function createOrchestrationSystem() { const runtime = ManagedRuntime.make(orchestrationLayer); const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery)); + const eventStore = await runtime.runPromise(Effect.service(OrchestrationEventStore)); + const projectionPipeline = await runtime.runPromise( + Effect.service(OrchestrationProjectionPipeline), + ); return { engine, + eventStore, + projectionPipeline, + getCommandReadModel: () => runtime.runPromise(snapshotQuery.getCommandReadModel()), readModel: () => runtime.runPromise(snapshotQuery.getSnapshot()), run: (effect: Effect.Effect) => runtime.runPromise(effect), dispose: () => runtime.dispose(), @@ -420,6 +429,185 @@ describe("OrchestrationEngine", () => { await system.dispose(); }); + it("dispatches commands against threads appended by an external writer", async () => { + const system = await createOrchestrationSystem(); + const { engine, eventStore, projectionPipeline } = system; + const createdAt = now(); + const projectId = asProjectId("project-external-writer"); + const threadAId = ThreadId.make("thread-engine"); + const threadBId = ThreadId.make("thread-external"); + + await system.run( + engine.dispatch({ + type: "project.create", + commandId: CommandId.make("cmd-external-project-create"), + projectId, + title: "External Writer Project", + workspaceRoot: "/tmp/project-external-writer", + defaultModelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + createdAt, + }), + ); + await system.run( + engine.dispatch({ + type: "thread.create", + commandId: CommandId.make("cmd-external-thread-a-create"), + threadId: threadAId, + projectId, + title: "engine thread", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + branch: null, + worktreePath: null, + createdAt, + }), + ); + + const externalWriterReadModel = await system.getCommandReadModel(); + await system.run( + Effect.gen(function* () { + const eventBase = yield* decideOrchestrationCommand({ + command: { + type: "thread.create", + commandId: CommandId.make("cmd-external-thread-b-create"), + threadId: threadBId, + projectId, + title: "external thread", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + branch: null, + worktreePath: null, + createdAt, + }, + readModel: externalWriterReadModel, + }); + const eventBases = Array.isArray(eventBase) ? eventBase : [eventBase]; + for (const event of eventBases) { + const savedEvent = yield* eventStore.append(event); + yield* projectionPipeline.projectEvent(savedEvent); + } + }).pipe(Effect.provide(NodeServices.layer)), + ); + + await system.run( + engine.dispatch({ + type: "thread.meta.update", + commandId: CommandId.make("cmd-external-watermark-advance"), + threadId: threadAId, + title: "engine thread updated", + }), + ); + + await system.run( + engine.dispatch({ + type: "thread.delete", + commandId: CommandId.make("cmd-external-delete"), + threadId: threadBId, + }), + ); + + const events = await system.run( + Stream.runCollect(engine.readEvents(0)).pipe( + Effect.map((chunk): OrchestrationEvent[] => Array.from(chunk)), + ), + ); + expect( + events.some((event) => event.type === "thread.deleted" && event.aggregateId === threadBId), + ).toBe(true); + + await system.dispose(); + }); + + it("dispatches thread creation into a project appended by an external writer", async () => { + const system = await createOrchestrationSystem(); + const { engine, eventStore, projectionPipeline } = system; + const createdAt = now(); + const projectId = asProjectId("project-external-import"); + const threadId = ThreadId.make("thread-external-import"); + + const externalWriterReadModel = await system.getCommandReadModel(); + await system.run( + Effect.gen(function* () { + const eventBase = yield* decideOrchestrationCommand({ + command: { + type: "project.create", + commandId: CommandId.make("cmd-external-project-create"), + projectId, + title: "External Import Project", + workspaceRoot: "/tmp/project-external-import", + defaultModelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + createdAt, + }, + readModel: externalWriterReadModel, + }); + const eventBases = Array.isArray(eventBase) ? eventBase : [eventBase]; + for (const event of eventBases) { + const savedEvent = yield* eventStore.append(event); + yield* projectionPipeline.projectEvent(savedEvent); + } + }).pipe(Effect.provide(NodeServices.layer)), + ); + + await system.run( + engine.dispatch({ + type: "project.create", + commandId: CommandId.make("cmd-watermark-project-create"), + projectId: asProjectId("project-watermark"), + title: "Watermark Project", + workspaceRoot: "/tmp/project-watermark", + defaultModelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + createdAt, + }), + ); + + await system.run( + engine.dispatch({ + type: "thread.create", + commandId: CommandId.make("cmd-thread-create-in-external-project"), + threadId, + projectId, + title: "thread in imported project", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + branch: null, + worktreePath: null, + createdAt, + }), + ); + + const events = await system.run( + Stream.runCollect(engine.readEvents(0)).pipe( + Effect.map((chunk): OrchestrationEvent[] => Array.from(chunk)), + ), + ); + expect( + events.some((event) => event.type === "thread.created" && event.aggregateId === threadId), + ).toBe(true); + + await system.dispose(); + }); + it("streams persisted domain events in order", async () => { const system = await createOrchestrationSystem(); const { engine } = system; diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts index 19184915ac7..1c9e455019b 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts @@ -37,6 +37,7 @@ import { type OrchestrationDispatchError, type OrchestrationProjectorDecodeError, } from "../Errors.ts"; +import { findProjectById, findThreadById } from "../commandInvariants.ts"; import { decideOrchestrationCommand } from "../decider.ts"; import { createEmptyReadModel, projectEvent } from "../projector.ts"; import { OrchestrationProjectionPipeline } from "../Services/ProjectionPipeline.ts"; @@ -56,10 +57,11 @@ interface CommandEnvelope { startedAtMs: number; } -function commandToAggregateRef(command: OrchestrationCommand): { - readonly aggregateKind: "project" | "thread"; - readonly aggregateId: ProjectId | ThreadId; -} { +function commandToAggregateRef( + command: OrchestrationCommand, +): + | { readonly aggregateKind: "project"; readonly aggregateId: ProjectId } + | { readonly aggregateKind: "thread"; readonly aggregateId: ThreadId } { switch (command.type) { case "project.create": case "project.meta.update": @@ -150,6 +152,49 @@ const makeOrchestrationEngine = Effect.gen(function* () { }); } + // External writers (e.g. the `t3 import` CLI) append events and projections + // to the shared database from another process, so an aggregate missing from + // the in-memory model may still exist in SQL; refresh before rejecting. + const mustExistAggregate = + envelope.command.type === "project.create" + ? undefined + : envelope.command.type === "thread.create" + ? { + aggregateKind: "project" as const, + aggregateId: envelope.command.projectId, + } + : aggregateRef; + const aggregateMissingFromMemory = + mustExistAggregate !== undefined && + (mustExistAggregate.aggregateKind === "thread" + ? findThreadById(commandReadModel, mustExistAggregate.aggregateId) === undefined + : findProjectById(commandReadModel, mustExistAggregate.aggregateId) === undefined); + if (aggregateMissingFromMemory) { + const mustExistAggregateId = mustExistAggregate.aggregateId; + yield* Effect.gen(function* () { + const refreshed = yield* projectionSnapshotQuery.getCommandReadModel(); + commandReadModel = + refreshed.snapshotSequence >= commandReadModel.snapshotSequence + ? refreshed + : { + ...refreshed, + snapshotSequence: commandReadModel.snapshotSequence, + updatedAt: commandReadModel.updatedAt, + }; + }).pipe( + Effect.catch(() => + Effect.logWarning( + "failed to refresh orchestration read model before command dispatch", + ).pipe( + Effect.annotateLogs({ + commandId: envelope.command.commandId, + aggregateId: mustExistAggregateId, + }), + ), + ), + ); + } + const eventBase = yield* decideOrchestrationCommand({ command: envelope.command, readModel: commandReadModel,