From 26f2119e87d9fd01b07f6762cb8a8b1e529e9cb3 Mon Sep 17 00:00:00 2001 From: Shivam Sharma <91240327+shivamhwp@users.noreply.github.com> Date: Sun, 23 Aug 2026 06:37:21 +0530 Subject: [PATCH 1/3] fix(server): a draft can retry its first send after a failed bootstrap A new thread's first send creates the thread and starts the turn in one bootstrap. When the bootstrap fails partway (worktree prep, setup script, a dropped connection), the server rolls back with thread.delete. That is a soft delete, the draft keeps its client-minted thread id, and the retry's thread.create hit requireThreadAbsent, which treated the tombstone as a live thread: "Thread already exists and cannot be created twice", on every retry, until the draft was abandoned. requireThreadAbsent now only blocks on a live row. Each per-thread projector drops its rows for the old incarnation when it applies thread.created, so a re-created id starts with an empty timeline and per-projector replay stays deterministic. A replayed thread.deleted that a later thread.created supersedes no longer removes attachment files, since those already belong to the new incarnation. Both thread.create paths in ws.ts wait for the deletion reactor to drain first, so the old incarnation's session stop and terminal close always finish before the new thread can own those resources. Co-Authored-By: Claude Fable 5 --- .../Layers/OrchestrationEngine.test.ts | 3 + .../Layers/ProjectionPipeline.test.ts | 252 +++++++++++++++++- .../Layers/ProjectionPipeline.ts | 52 +++- .../orchestration/commandInvariants.test.ts | 30 +++ .../src/orchestration/commandInvariants.ts | 6 +- .../Layers/OrchestrationEventStore.ts | 35 +++ .../Layers/ProjectionPendingApprovals.ts | 17 ++ .../Services/OrchestrationEventStore.ts | 13 + .../Services/ProjectionPendingApprovals.ts | 7 + apps/server/src/server.test.ts | 125 ++++++++- apps/server/src/ws.ts | 14 +- 11 files changed, 543 insertions(+), 11 deletions(-) diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 382c253fe60b..ba23d56b5e07 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -113,6 +113,7 @@ describe("OrchestrationEngine", () => { detail: "historical replay should not be used during bootstrap", }), ), + hasEventAfter: () => Effect.succeed(false), }; const projectionSnapshot = { @@ -812,6 +813,7 @@ describe("OrchestrationEngine", () => { readAll() { return Stream.fromIterable(events); }, + hasEventAfter: () => Effect.succeed(false), }; const ServerConfigLayer = ServerConfig.layerTest(process.cwd(), { @@ -1048,6 +1050,7 @@ describe("OrchestrationEngine", () => { readAll() { return Stream.fromIterable(events); }, + hasEventAfter: () => Effect.succeed(false), }; let shouldFailProjection = true; diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index cd95293aa635..732f88e3e07a 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -9,6 +9,7 @@ import { TurnId, ProviderInstanceId, } from "@t3tools/contracts"; +import * as Option from "effect/Option"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; @@ -31,6 +32,7 @@ import { OrchestrationProjectionPipelineLive, } from "./ProjectionPipeline.ts"; import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts"; +import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts"; import * as ThreadPlanProgress from "../ThreadPlanProgress.ts"; import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts"; @@ -1242,6 +1244,111 @@ it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-projection-atta }, ); +it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-projection-attachments-replay-")))( + "OrchestrationProjectionPipeline", + (it) => { + it.effect("replaying a superseded thread.deleted keeps the re-created thread's files", () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const { attachmentsDir } = yield* ServerConfig; + const now = "2026-01-01T00:00:00.000Z"; + const projectId = ProjectId.make("project-replay"); + const retriedThreadId = ThreadId.make("thread-replay-retried"); + const goneThreadId = ThreadId.make("thread-replay-gone"); + const retriedAttachmentPath = path.join( + attachmentsDir, + "thread-replay-retried-00000000-0000-4000-8000-000000000001.png", + ); + const goneAttachmentPath = path.join( + attachmentsDir, + "thread-replay-gone-00000000-0000-4000-8000-000000000002.png", + ); + const threadCreated = (threadId: ThreadId, suffix: string) => + eventStore.append({ + type: "thread.created", + eventId: EventId.make(`evt-replay-create-${suffix}`), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: now, + commandId: CommandId.make(`cmd-replay-create-${suffix}`), + causationEventId: null, + correlationId: CorrelationId.make(`cmd-replay-create-${suffix}`), + metadata: {}, + payload: { + threadId, + projectId, + title: `Thread ${suffix}`, + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + }, + }); + const threadDeleted = (threadId: ThreadId, suffix: string) => + eventStore.append({ + type: "thread.deleted", + eventId: EventId.make(`evt-replay-delete-${suffix}`), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: now, + commandId: CommandId.make(`cmd-replay-delete-${suffix}`), + causationEventId: null, + correlationId: CorrelationId.make(`cmd-replay-delete-${suffix}`), + metadata: {}, + payload: { threadId, deletedAt: now }, + }); + + yield* eventStore.append({ + type: "project.created", + eventId: EventId.make("evt-replay-project"), + aggregateKind: "project", + aggregateId: projectId, + occurredAt: now, + commandId: CommandId.make("cmd-replay-project"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-replay-project"), + metadata: {}, + payload: { + projectId, + title: "Replay", + workspaceRoot: "/tmp/project-replay", + defaultModelSelection: null, + scripts: [], + createdAt: now, + updatedAt: now, + }, + }); + // A failed first send: create, roll back, then the draft retries the id. + yield* threadCreated(retriedThreadId, "retried-1"); + yield* threadDeleted(retriedThreadId, "retried"); + yield* threadCreated(retriedThreadId, "retried-2"); + // A thread that was deleted for good. + yield* threadCreated(goneThreadId, "gone"); + yield* threadDeleted(goneThreadId, "gone"); + + // Files on disk are not event-sourced: by the time anything replays, + // the retried thread's attachments already belong to its second life. + yield* fileSystem.makeDirectory(attachmentsDir, { recursive: true }); + yield* fileSystem.writeFileString(retriedAttachmentPath, "second incarnation"); + yield* fileSystem.writeFileString(goneAttachmentPath, "gone"); + + yield* projectionPipeline.bootstrap; + + assert.isTrue(yield* exists(retriedAttachmentPath)); + assert.isFalse(yield* exists(goneAttachmentPath)); + }), + ); + }, +); + it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { it.effect("resumes from projector last_applied_sequence without replaying older events", () => Effect.gen(function* () { @@ -2744,7 +2851,7 @@ it.effect("restores pending turn-start metadata across projection pipeline resta const engineLayer = it.layer( OrchestrationEngineLive.pipe( - Layer.provide(OrchestrationProjectionSnapshotQueryLive), + Layer.provideMerge(OrchestrationProjectionSnapshotQueryLive), Layer.provide(ThreadBackgroundLiveness.layer), Layer.provide(ThreadPlanProgress.layer), Layer.provide(OrchestrationProjectionPipelineLive), @@ -2861,4 +2968,147 @@ engineLayer("OrchestrationProjectionPipeline via engine dispatch", (it) => { ]); }), ); + + it.effect("re-creating a deleted thread id starts from an empty projection", () => + Effect.gen(function* () { + const engine = yield* OrchestrationEngineService; + const snapshotQuery = yield* ProjectionSnapshotQuery; + const sql = yield* SqlClient.SqlClient; + const createdAt = "2026-01-01T00:00:00.000Z"; + const projectId = ProjectId.make("project-retry"); + const threadId = ThreadId.make("thread-retry"); + const modelSelection = { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }; + const createThread = (commandId: string, title: string) => + engine.dispatch({ + type: "thread.create", + commandId: CommandId.make(commandId), + threadId, + projectId, + title, + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt, + }); + const countRowsForThread = (table: string) => + sql<{ readonly count: number }>` + SELECT COUNT(*) AS count FROM ${sql(table)} WHERE thread_id = ${threadId} + `.pipe(Effect.map((rows) => rows[0]?.count ?? 0)); + const perThreadTables = [ + "projection_thread_messages", + "projection_thread_activities", + "projection_thread_sessions", + "projection_turns", + "projection_thread_proposed_plans", + "projection_pending_approvals", + ]; + + yield* engine.dispatch({ + type: "project.create", + commandId: CommandId.make("cmd-retry-project"), + projectId, + title: "Retry Project", + workspaceRoot: "/tmp/project-retry", + defaultModelSelection: modelSelection, + createdAt, + }); + + // First attempt: the thread gets a turn, a message, an activity, and a + // running session before its bootstrap fails and the server rolls back. + yield* createThread("cmd-retry-create-1", "First attempt"); + yield* engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-retry-turn-1"), + threadId, + message: { + messageId: MessageId.make("message-retry-1"), + role: "user", + text: "first attempt", + attachments: [], + }, + runtimeMode: "full-access", + interactionMode: "default", + createdAt, + }); + yield* engine.dispatch({ + type: "thread.activity.append", + commandId: CommandId.make("cmd-retry-activity-1"), + threadId, + activity: { + id: EventId.make("activity-retry-1"), + tone: "info", + kind: "approval.requested", + summary: "approval requested", + payload: { requestId: "request-retry-1" }, + turnId: null, + createdAt, + }, + createdAt, + }); + yield* engine.dispatch({ + type: "thread.proposed-plan.upsert", + commandId: CommandId.make("cmd-retry-plan-1"), + threadId, + proposedPlan: { + id: "plan-retry-1", + turnId: null, + planMarkdown: "# Plan", + implementedAt: null, + implementationThreadId: null, + createdAt, + updatedAt: createdAt, + }, + createdAt, + }); + yield* engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-retry-session-1"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: TurnId.make("turn-retry-1"), + lastError: null, + updatedAt: createdAt, + }, + createdAt, + }); + for (const table of perThreadTables) { + assert.isAbove(yield* countRowsForThread(table), 0, `${table} should be populated`); + } + const populatedShell = Option.getOrThrow(yield* snapshotQuery.getThreadShellById(threadId)); + assert.isTrue(populatedShell.hasPendingApprovals); + assert.isTrue(populatedShell.hasActionableProposedPlan); + + yield* engine.dispatch({ + type: "thread.delete", + commandId: CommandId.make("cmd-retry-delete"), + threadId, + }); + assert.isTrue(Option.isNone(yield* snapshotQuery.getThreadShellById(threadId))); + + // Retry from the same draft reuses the thread id. + yield* createThread("cmd-retry-create-2", "Second attempt"); + + const shell = Option.getOrThrow(yield* snapshotQuery.getThreadShellById(threadId)); + assert.strictEqual(shell.title, "Second attempt"); + assert.isFalse(shell.hasPendingApprovals); + assert.isFalse(shell.hasActionableProposedPlan); + for (const table of perThreadTables) { + assert.strictEqual(yield* countRowsForThread(table), 0, `${table} should be empty`); + } + const detail = Option.getOrThrow(yield* snapshotQuery.getThreadDetailById(threadId)); + assert.deepEqual(detail.messages, []); + assert.deepEqual(detail.activities, []); + assert.isNull(detail.latestTurn); + assert.isNull(detail.session); + }), + ); }); diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index 8eb8cdb561b8..19a521d4da43 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -862,7 +862,18 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti } case "thread.deleted": { - attachmentSideEffects.deletedThreadIds.add(event.payload.threadId); + // A draft retry can re-create this id later in the log. During + // replay the attachment files on disk already belong to that later + // incarnation, so only an unsuperseded deletion removes them. + const recreatedLater = yield* eventStore.hasEventAfter({ + aggregateKind: "thread", + aggregateId: event.payload.threadId, + type: "thread.created", + sequenceExclusive: event.sequence, + }); + if (!recreatedLater) { + attachmentSideEffects.deletedThreadIds.add(event.payload.threadId); + } const existingRow = yield* projectionThreadRepository.getById({ threadId: event.payload.threadId, }); @@ -978,6 +989,15 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti "applyThreadMessagesProjection", )(function* (event, attachmentSideEffects) { switch (event.type) { + // A draft retry re-creates a soft-deleted thread id. Every projector + // drops its own rows for the old incarnation here so replay from any + // per-projector cursor rebuilds the new thread without stale history. + case "thread.created": + yield* projectionThreadMessageRepository.deleteByThreadId({ + threadId: event.payload.threadId, + }); + return; + case "thread.message-sent": { const existingMessage = yield* projectionThreadMessageRepository.getByMessageId({ messageId: event.payload.messageId, @@ -1057,6 +1077,12 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti "applyThreadProposedPlansProjection", )(function* (event, _attachmentSideEffects) { switch (event.type) { + case "thread.created": + yield* projectionThreadProposedPlanRepository.deleteByThreadId({ + threadId: event.payload.threadId, + }); + return; + case "thread.proposed-plan-upserted": yield* projectionThreadProposedPlanRepository.upsert({ planId: event.payload.proposedPlan.id, @@ -1108,6 +1134,12 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti "applyThreadActivitiesProjection", )(function* (event, _attachmentSideEffects) { switch (event.type) { + case "thread.created": + yield* projectionThreadActivityRepository.deleteByThreadId({ + threadId: event.payload.threadId, + }); + return; + case "thread.activity-appended": yield* projectionThreadActivityRepository.upsert({ activityId: event.payload.activity.id, @@ -1159,6 +1191,12 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti const applyThreadSessionsProjection: ProjectorDefinition["apply"] = Effect.fn( "applyThreadSessionsProjection", )(function* (event, _attachmentSideEffects) { + if (event.type === "thread.created") { + yield* projectionThreadSessionRepository.deleteByThreadId({ + threadId: event.payload.threadId, + }); + return; + } if (event.type !== "thread.session-set") { return; } @@ -1178,6 +1216,12 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti "applyThreadTurnsProjection", )(function* (event, _attachmentSideEffects) { switch (event.type) { + case "thread.created": + yield* projectionTurnRepository.deleteByThreadId({ + threadId: event.payload.threadId, + }); + return; + case "thread.turn-start-requested": { yield* projectionTurnRepository.replacePendingTurnStart({ threadId: event.payload.threadId, @@ -1515,6 +1559,12 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti "applyPendingApprovalsProjection", )(function* (event, _attachmentSideEffects) { switch (event.type) { + case "thread.created": + yield* projectionPendingApprovalRepository.deleteByThreadId({ + threadId: event.payload.threadId, + }); + return; + case "thread.activity-appended": { const requestId = extractActivityRequestId(event.payload.activity.payload) ?? diff --git a/apps/server/src/orchestration/commandInvariants.test.ts b/apps/server/src/orchestration/commandInvariants.test.ts index 52aac1f0c105..9aaeba943423 100644 --- a/apps/server/src/orchestration/commandInvariants.test.ts +++ b/apps/server/src/orchestration/commandInvariants.test.ts @@ -199,4 +199,34 @@ describe("commandInvariants", () => { ), ).rejects.toThrow("already exists"); }); + + it("lets a draft retry re-create a thread id after its first attempt was deleted", async () => { + const threadId = ThreadId.make("thread-1"); + const firstAttempt = readModel.threads.find((thread) => thread.id === threadId)!; + const afterRollback: OrchestrationReadModel = { + ...readModel, + threads: readModel.threads.map((thread) => + thread.id === threadId ? { ...thread, deletedAt: now, updatedAt: now } : thread, + ), + }; + const retry: OrchestrationCommand = { + type: "thread.create", + commandId: CommandId.make("cmd-retry"), + threadId, + projectId: firstAttempt.projectId, + title: firstAttempt.title, + modelSelection: firstAttempt.modelSelection, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + branch: null, + worktreePath: null, + createdAt: now, + }; + + await expect( + Effect.runPromise( + requireThreadAbsent({ readModel: afterRollback, command: retry, threadId }), + ), + ).resolves.toBeUndefined(); + }); }); diff --git a/apps/server/src/orchestration/commandInvariants.ts b/apps/server/src/orchestration/commandInvariants.ts index b59ded77f4f4..8f1e3e89851c 100644 --- a/apps/server/src/orchestration/commandInvariants.ts +++ b/apps/server/src/orchestration/commandInvariants.ts @@ -156,7 +156,11 @@ export function requireThreadAbsent(input: { readonly command: OrchestrationCommand; readonly threadId: ThreadId; }): Effect.Effect { - if (!findThreadById(input.readModel, input.threadId)) { + // Thread deletion is a soft delete and a draft keeps its client-minted id + // across retries, so only a live row blocks creation. Projectors reset the + // thread's rows when the id is created again. + const existing = findThreadById(input.readModel, input.threadId); + if (existing === undefined || existing.deletedAt !== null) { return Effect.void; } return Effect.fail( diff --git a/apps/server/src/persistence/Layers/OrchestrationEventStore.ts b/apps/server/src/persistence/Layers/OrchestrationEventStore.ts index 18d0e9aa578b..edd8620c3b3f 100644 --- a/apps/server/src/persistence/Layers/OrchestrationEventStore.ts +++ b/apps/server/src/persistence/Layers/OrchestrationEventStore.ts @@ -15,6 +15,7 @@ import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as SqlSchema from "effect/unstable/sql/SqlSchema"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; @@ -60,6 +61,13 @@ const OrchestrationEventPersistedRowSchema = Schema.Struct({ metadata: EventMetadataFromJsonString, }); +const HasEventAfterRequestSchema = Schema.Struct({ + aggregateKind: Schema.String, + aggregateId: Schema.String, + type: Schema.String, + sequenceExclusive: NonNegativeInt, +}); + const ReadFromSequenceRequestSchema = Schema.Struct({ sequenceExclusive: NonNegativeInt, limit: Schema.Number, @@ -260,10 +268,37 @@ const makeEventStore = Effect.gen(function* () { return readPage(sequenceExclusive, normalizedLimit); }; + const findEventAfter = SqlSchema.findOneOption({ + Request: HasEventAfterRequestSchema, + Result: Schema.Struct({ sequence: Schema.Number }), + execute: (request) => + sql` + SELECT sequence + FROM orchestration_events + WHERE aggregate_kind = ${request.aggregateKind} + AND stream_id = ${request.aggregateId} + AND event_type = ${request.type} + AND sequence > ${request.sequenceExclusive} + LIMIT 1 + `, + }); + + const hasEventAfter: OrchestrationEventStoreShape["hasEventAfter"] = (input) => + findEventAfter(input).pipe( + Effect.map(Option.isSome), + Effect.mapError( + toPersistenceSqlOrDecodeError( + "OrchestrationEventStore.hasEventAfter:query", + "OrchestrationEventStore.hasEventAfter:decodeRow", + ), + ), + ); + return { append, readFromSequence, readAll: () => readFromSequence(0, Number.MAX_SAFE_INTEGER), + hasEventAfter, } satisfies OrchestrationEventStoreShape; }); diff --git a/apps/server/src/persistence/Layers/ProjectionPendingApprovals.ts b/apps/server/src/persistence/Layers/ProjectionPendingApprovals.ts index 253f6e13b977..3b159a9e1715 100644 --- a/apps/server/src/persistence/Layers/ProjectionPendingApprovals.ts +++ b/apps/server/src/persistence/Layers/ProjectionPendingApprovals.ts @@ -95,6 +95,15 @@ const makeProjectionPendingApprovalRepository = Effect.gen(function* () { `, }); + const deleteProjectionPendingApprovalRowsByThread = SqlSchema.void({ + Request: ListProjectionPendingApprovalsInput, + execute: ({ threadId }) => + sql` + DELETE FROM projection_pending_approvals + WHERE thread_id = ${threadId} + `, + }); + const upsert: ProjectionPendingApprovalRepositoryShape["upsert"] = (row) => upsertProjectionPendingApprovalRow(row).pipe( Effect.mapError(toPersistenceSqlError("ProjectionPendingApprovalRepository.upsert:query")), @@ -123,11 +132,19 @@ const makeProjectionPendingApprovalRepository = Effect.gen(function* () { ), ); + const deleteByThreadId: ProjectionPendingApprovalRepositoryShape["deleteByThreadId"] = (input) => + deleteProjectionPendingApprovalRowsByThread(input).pipe( + Effect.mapError( + toPersistenceSqlError("ProjectionPendingApprovalRepository.deleteByThreadId:query"), + ), + ); + return { upsert, listByThreadId, getByRequestId, deleteByRequestId, + deleteByThreadId, } satisfies ProjectionPendingApprovalRepositoryShape; }); diff --git a/apps/server/src/persistence/Services/OrchestrationEventStore.ts b/apps/server/src/persistence/Services/OrchestrationEventStore.ts index 8b465e7713e1..488210ab74a5 100644 --- a/apps/server/src/persistence/Services/OrchestrationEventStore.ts +++ b/apps/server/src/persistence/Services/OrchestrationEventStore.ts @@ -52,6 +52,19 @@ export interface OrchestrationEventStoreShape { * @returns Stream containing all stored events. */ readonly readAll: () => Stream.Stream; + + /** + * Check whether an aggregate has an event of the given type after a sequence. + * + * Used during replay to tell whether a later event supersedes the one being + * applied, without streaming the rest of the log. + */ + readonly hasEventAfter: (input: { + readonly aggregateKind: OrchestrationEvent["aggregateKind"]; + readonly aggregateId: string; + readonly type: OrchestrationEvent["type"]; + readonly sequenceExclusive: number; + }) => Effect.Effect; } /** diff --git a/apps/server/src/persistence/Services/ProjectionPendingApprovals.ts b/apps/server/src/persistence/Services/ProjectionPendingApprovals.ts index 967e6da9d3af..40b0d1ae03b6 100644 --- a/apps/server/src/persistence/Services/ProjectionPendingApprovals.ts +++ b/apps/server/src/persistence/Services/ProjectionPendingApprovals.ts @@ -82,6 +82,13 @@ export interface ProjectionPendingApprovalRepositoryShape { readonly deleteByRequestId: ( input: DeleteProjectionPendingApprovalInput, ) => Effect.Effect; + + /** + * Delete every pending approval row for a thread. + */ + readonly deleteByThreadId: ( + input: ListProjectionPendingApprovalsInput, + ) => Effect.Effect; } /** diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index a9a2c3fa10d6..8e141515fced 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -115,6 +115,7 @@ import * as RemoteOpenTargets from "./environment/RemoteOpenTargets.ts"; import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts"; import { OrchestrationListenerCallbackError } from "./orchestration/Errors.ts"; import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSnapshotQuery.ts"; +import { ThreadDeletionReactor } from "./orchestration/Services/ThreadDeletionReactor.ts"; import { SqlitePersistenceMemory } from "./persistence/Layers/Sqlite.ts"; import { PersistenceSqlError } from "./persistence/Errors.ts"; import * as ProviderRegistry from "./provider/Services/ProviderRegistry.ts"; @@ -410,6 +411,7 @@ const buildAppUnderTest = (options?: { >; terminalManager?: Partial; orchestrationEngine?: Partial; + threadDeletionReactor?: Partial; analyticsService?: Partial; projectionSnapshotQuery?: Partial; checkpointDiffQuery?: Partial; @@ -786,13 +788,20 @@ const buildAppUnderTest = (options?: { ), ), Layer.provide( - Layer.mock(OrchestrationEngine.OrchestrationEngineService)({ - readEvents: () => Stream.empty, - dispatch: () => Effect.succeed({ sequence: 0 }), - streamDomainEvents: Stream.empty, - latestSequence: Effect.succeed(0), - ...options?.layers?.orchestrationEngine, - }), + Layer.mergeAll( + Layer.mock(OrchestrationEngine.OrchestrationEngineService)({ + readEvents: () => Stream.empty, + dispatch: () => Effect.succeed({ sequence: 0 }), + streamDomainEvents: Stream.empty, + latestSequence: Effect.succeed(0), + ...options?.layers?.orchestrationEngine, + }), + Layer.mock(ThreadDeletionReactor)({ + start: () => Effect.void, + drain: Effect.void, + ...options?.layers?.threadDeletionReactor, + }), + ), ), Layer.provide( Layer.mock(ProjectionSnapshotQuery.ProjectionSnapshotQuery)({ @@ -8204,6 +8213,108 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }).pipe(Effect.provide(NodeHttpServer.layerTest)), ); + it.effect("waits for deletion cleanup before re-creating a thread id", () => + Effect.gen(function* () { + // A draft retry reuses the thread id its failed bootstrap deleted. The + // deletion reactor stops sessions and closes terminals by that id, so + // both thread.create paths must let its queue drain first. + const trace: Array = []; + const drainRequested = yield* Deferred.make(); + const cleanupDone = yield* Deferred.make(); + yield* buildAppUnderTest({ + layers: { + threadDeletionReactor: { + drain: Effect.gen(function* () { + trace.push("drain"); + yield* Deferred.succeed(drainRequested, undefined); + yield* Deferred.await(cleanupDone); + }), + }, + orchestrationEngine: { + dispatch: (command) => + Effect.sync(() => { + trace.push(command.type); + return { sequence: trace.length }; + }), + readEvents: () => Stream.empty, + }, + }, + }); + + const createdAt = "2026-01-01T00:00:00.000Z"; + const threadId = ThreadId.make("thread-retry-after-delete"); + const wsUrl = yield* getWsServerUrl("/ws"); + + yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + Effect.gen(function* () { + const directCreate = yield* Effect.forkChild( + client[ORCHESTRATION_WS_METHODS.dispatchCommand]({ + type: "thread.create", + commandId: CommandId.make("cmd-retry-create"), + threadId, + projectId: defaultProjectId, + title: "Retry", + modelSelection: defaultModelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt, + }), + ); + yield* Deferred.await(drainRequested); + assert.deepEqual(trace, ["drain"]); + yield* Deferred.succeed(cleanupDone, undefined); + yield* Fiber.join(directCreate); + }), + ), + ); + assert.deepEqual(trace, ["drain", "thread.create"]); + + // Cleanup is already released; the bootstrap path must still drain first. + trace.length = 0; + yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + Effect.gen(function* () { + const bootstrapCreate = yield* Effect.forkChild( + client[ORCHESTRATION_WS_METHODS.dispatchCommand]({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-retry-bootstrap"), + threadId, + message: { + messageId: MessageId.make("msg-retry-bootstrap"), + role: "user", + text: "hello", + attachments: [], + }, + modelSelection: defaultModelSelection, + runtimeMode: "full-access", + interactionMode: "default", + bootstrap: { + createThread: { + projectId: defaultProjectId, + title: "Retry", + modelSelection: defaultModelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt, + }, + runSetupScript: false, + }, + createdAt, + }), + ); + yield* Fiber.join(bootstrapCreate); + }), + ), + ); + assert.deepEqual(trace, ["drain", "thread.create", "thread.turn.start"]); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + it.effect("does not report a deleted bootstrap thread when cleanup fails", () => Effect.gen(function* () { const dispatchedCommands: Array = []; diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 226c82cdb1ac..6380edf9aa71 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -81,6 +81,7 @@ import { } from "./orchestration/Normalizer.ts"; import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts"; import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSnapshotQuery.ts"; +import { ThreadDeletionReactor } from "./orchestration/Services/ThreadDeletionReactor.ts"; import { observeRpcEffect as instrumentRpcEffect, observeRpcStream as instrumentRpcStream, @@ -431,6 +432,12 @@ const makeWsRpcLayer = ( const crypto = yield* Crypto.Crypto; const projectionSnapshotQuery = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; const orchestrationEngine = yield* OrchestrationEngine.OrchestrationEngineService; + const threadDeletionReactor = yield* ThreadDeletionReactor; + // A draft retry re-creates the thread id its failed bootstrap deleted. + // Session stop and terminal close are keyed by that id and run in the + // deletion reactor's queue, so let the old incarnation's cleanup finish + // before the new one can own those resources. + const awaitThreadDeletionCleanup = threadDeletionReactor.drain; const analytics = yield* AnalyticsService.AnalyticsService; // Every command dispatched on this connection carries the connecting // client's origin, including server-generated bootstrap sub-commands: @@ -998,6 +1005,7 @@ const makeWsRpcLayer = ( const bootstrapProgram = Effect.gen(function* () { if (bootstrap?.createThread) { + yield* awaitThreadDeletionCleanup; yield* dispatchFromClient({ type: "thread.create", commandId: yield* serverCommandId("bootstrap-thread-create"), @@ -1097,7 +1105,11 @@ const makeWsRpcLayer = ( const dispatchEffect = normalizedCommand.type === "thread.turn.start" && normalizedCommand.bootstrap ? dispatchBootstrapTurnStart(normalizedCommand) - : dispatchFromClient(normalizedCommand).pipe( + : (normalizedCommand.type === "thread.create" + ? awaitThreadDeletionCleanup + : Effect.void + ).pipe( + Effect.andThen(dispatchFromClient(normalizedCommand)), Effect.mapError((cause) => toDispatchCommandError(cause, "Failed to dispatch orchestration command"), ), From 1c06a88319fba982c910da6869b17661445be562 Mon Sep 17 00:00:00 2001 From: Shivam Sharma <91240327+shivamhwp@users.noreply.github.com> Date: Wed, 26 Aug 2026 02:55:50 +0530 Subject: [PATCH 2/3] fix(server): deletion drain covers events still in flight to the subscriber The drain that gates thread.create only waited for the reactor's queue to empty, but thread.deleted reaches that queue through an asynchronous subscriber. A retry arriving between publish and enqueue could re-create the id and then have the old cleanup stop the new session. The reactor now tracks the highest event sequence its subscriber has handed on, and drain first waits for that to reach the engine's latestSequence before draining the worker. Co-Authored-By: Claude Fable 5 --- .../Layers/ThreadDeletionReactor.test.ts | 105 +++++++++++++++++- .../Layers/ThreadDeletionReactor.ts | 37 ++++-- .../Services/ThreadDeletionReactor.ts | 6 +- 3 files changed, 137 insertions(+), 11 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts index 34b1b995a3ad..11ed1dc62124 100644 --- a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts @@ -1,10 +1,35 @@ -import { ThreadId } from "@t3tools/contracts"; +import { + CommandId, + CorrelationId, + EventId, + type OrchestrationEvent, + ThreadId, +} from "@t3tools/contracts"; +import { it as effectIt } from "@effect/vitest"; import * as Cause from "effect/Cause"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Ref from "effect/Ref"; +import * as Stream from "effect/Stream"; import { describe, expect, it } from "vite-plus/test"; -import { logCleanupCauseUnlessInterrupted } from "./ThreadDeletionReactor.ts"; +import { + ProviderService, + type ProviderServiceShape, +} from "../../provider/Services/ProviderService.ts"; +import * as TerminalManager from "../../terminal/Manager.ts"; +import { + OrchestrationEngineService, + type OrchestrationEngineShape, +} from "../Services/OrchestrationEngine.ts"; +import { ThreadDeletionReactor } from "../Services/ThreadDeletionReactor.ts"; +import { + logCleanupCauseUnlessInterrupted, + ThreadDeletionReactorLive, +} from "./ThreadDeletionReactor.ts"; describe("logCleanupCauseUnlessInterrupted", () => { const threadId = ThreadId.make("thread-deletion-reactor-test"); @@ -36,3 +61,79 @@ describe("logCleanupCauseUnlessInterrupted", () => { } }); }); + +describe("ThreadDeletionReactor drain", () => { + const now = "2026-01-01T00:00:00.000Z"; + const threadId = ThreadId.make("thread-deletion-reactor-drain"); + const deletedEvent = (sequence: number): OrchestrationEvent => ({ + sequence, + eventId: EventId.make(`evt-deleted-${sequence}`), + aggregateKind: "thread", + aggregateId: threadId, + type: "thread.deleted", + occurredAt: now, + commandId: CommandId.make(`cmd-deleted-${sequence}`), + causationEventId: null, + correlationId: CorrelationId.make(`cmd-deleted-${sequence}`), + metadata: {}, + payload: { threadId, deletedAt: now }, + }); + + effectIt.effect("waits for a published deletion the subscriber has not consumed yet", () => + Effect.gen(function* () { + const stops: Array = []; + const firstCleanupDone = yield* Deferred.make(); + // The engine has already committed and published sequence 2, but the + // subscriber has not received it yet: the stream releases it on demand. + const releaseSecondEvent = yield* Deferred.make(); + const latestSequence = yield* Ref.make(0); + const engine = { + latestSequence: Ref.get(latestSequence), + streamDomainEvents: Stream.concat( + Stream.make(deletedEvent(1)), + Stream.fromEffect(Deferred.await(releaseSecondEvent)).pipe( + Stream.map(() => deletedEvent(2)), + ), + ), + } as unknown as OrchestrationEngineShape; + const providerService = { + stopSession: () => + Effect.gen(function* () { + stops.push(stops.length + 1); + if (stops.length === 1) { + yield* Deferred.succeed(firstCleanupDone, undefined); + } + }), + } as unknown as ProviderServiceShape; + const terminalManager = { + close: () => Effect.void, + } as unknown as TerminalManager.TerminalManager["Service"]; + const layer = ThreadDeletionReactorLive.pipe( + Layer.provide(Layer.succeed(ProviderService, providerService)), + Layer.provide(Layer.succeed(TerminalManager.TerminalManager, terminalManager)), + Layer.provide(Layer.succeed(OrchestrationEngineService, engine)), + ); + + yield* Effect.scoped( + Effect.gen(function* () { + const reactor = yield* ThreadDeletionReactor; + yield* reactor.start(); + yield* Deferred.await(firstCleanupDone); + + // Sequence 1 is fully cleaned and the worker queue is idle. Sequence + // 2 is committed and published but still in flight to the subscriber. + yield* Ref.set(latestSequence, 2); + const drained = yield* Effect.forkChild(reactor.drain); + yield* Effect.yieldNow; + yield* Effect.yieldNow; + expect(stops).toEqual([1]); + expect(drained.pollUnsafe()).toBeUndefined(); + + yield* Deferred.succeed(releaseSecondEvent, undefined); + yield* Fiber.join(drained); + expect(stops).toEqual([1, 2]); + }), + ).pipe(Effect.provide(layer)); + }), + ); +}); diff --git a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts index a026f5ad81bd..28ed3adf094e 100644 --- a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts +++ b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts @@ -4,6 +4,7 @@ import * as Cause from "effect/Cause"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Stream from "effect/Stream"; +import * as SubscriptionRef from "effect/SubscriptionRef"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; import * as TerminalManager from "../../terminal/Manager.ts"; @@ -80,20 +81,42 @@ const make = Effect.gen(function* () { const worker = yield* makeDrainableWorker(processThreadDeletedSafely); + // Highest event sequence the subscriber has handed to the worker. Events + // are published after they commit, so a caller that reads latestSequence + // and waits for this to catch up knows every deletion up to that point is + // at least enqueued; the worker drain then covers the in-flight ones. + const seenSequence = yield* SubscriptionRef.make(0); + const noteSeen = (sequence: number) => + SubscriptionRef.update(seenSequence, (seen) => Math.max(seen, sequence)); + const start: ThreadDeletionReactorShape["start"] = Effect.fn("start")(function* () { yield* forkParked( - Stream.runForEach(orchestrationEngine.streamDomainEvents, (event) => { - if (event.type !== "thread.deleted") { - return Effect.void; - } - return worker.enqueue(event); - }), + Stream.runForEach( + orchestrationEngine.streamDomainEvents.pipe( + // Events that landed before the subscription are not replayed, so + // start the watermark at the current head instead of zero. + Stream.onStart(orchestrationEngine.latestSequence.pipe(Effect.flatMap(noteSeen))), + ), + (event) => + (event.type === "thread.deleted" ? worker.enqueue(event) : Effect.void).pipe( + Effect.andThen(noteSeen(event.sequence)), + ), + ), + ); + }); + + const drain: ThreadDeletionReactorShape["drain"] = Effect.gen(function* () { + const target = yield* orchestrationEngine.latestSequence; + yield* SubscriptionRef.changes(seenSequence).pipe( + Stream.filter((seen) => seen >= target), + Stream.runHead, ); + yield* worker.drain; }); return { start, - drain: worker.drain, + drain, } satisfies ThreadDeletionReactorShape; }); diff --git a/apps/server/src/orchestration/Services/ThreadDeletionReactor.ts b/apps/server/src/orchestration/Services/ThreadDeletionReactor.ts index 7c6718965a63..c47415c019c3 100644 --- a/apps/server/src/orchestration/Services/ThreadDeletionReactor.ts +++ b/apps/server/src/orchestration/Services/ThreadDeletionReactor.ts @@ -23,8 +23,10 @@ export interface ThreadDeletionReactorShape { readonly start: () => Effect.Effect; /** - * Resolves when the internal processing queue is empty and idle. - * Intended for test use to replace timing-sensitive sleeps. + * Resolves once every thread.deleted published up to now has been handed + * to the worker and the worker's queue is empty and idle. Callers that + * re-create a deleted thread id wait on this so the old incarnation's + * cleanup cannot run against the new one. */ readonly drain: Effect.Effect; } From e6b5f34931d7b444de9a104665942e6e06c0a448 Mon Sep 17 00:00:00 2001 From: Shivam Sharma <91240327+shivamhwp@users.noreply.github.com> Date: Sat, 29 Aug 2026 02:32:30 +0530 Subject: [PATCH 3/3] fix(server): fence deletion cleanup with thread creation --- .../OrchestrationEngineHarness.integration.ts | 2 +- .../Layers/OrchestrationReactor.test.ts | 2 +- .../Layers/ThreadDeletionReactor.test.ts | 2 +- .../Layers/ThreadDeletionReactor.ts | 15 ++++++----- .../Services/ThreadDeletionReactor.ts | 10 +++---- apps/server/src/server.test.ts | 27 ++++++++++--------- apps/server/src/ws.ts | 27 ++++++++++--------- 7 files changed, 46 insertions(+), 39 deletions(-) diff --git a/apps/server/integration/OrchestrationEngineHarness.integration.ts b/apps/server/integration/OrchestrationEngineHarness.integration.ts index f332b080ceea..6ade6025bcbc 100644 --- a/apps/server/integration/OrchestrationEngineHarness.integration.ts +++ b/apps/server/integration/OrchestrationEngineHarness.integration.ts @@ -373,7 +373,7 @@ export const makeOrchestrationIntegrationHarness = ( Layer.provideMerge( Layer.succeed(ThreadDeletionReactor, { start: () => Effect.void, - drain: Effect.void, + drainThrough: () => Effect.void, }), ), Layer.provideMerge( diff --git a/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts b/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts index 300d1526bb9a..b05ce3b1e235 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts @@ -61,7 +61,7 @@ describe("OrchestrationReactor", () => { started.push("thread-deletion-reactor"); return Effect.void; }, - drain: Effect.void, + drainThrough: () => Effect.void, }), ), Layer.provideMerge( diff --git a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts index 11ed1dc62124..f83f1dd1b9fa 100644 --- a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts @@ -123,7 +123,7 @@ describe("ThreadDeletionReactor drain", () => { // Sequence 1 is fully cleaned and the worker queue is idle. Sequence // 2 is committed and published but still in flight to the subscriber. yield* Ref.set(latestSequence, 2); - const drained = yield* Effect.forkChild(reactor.drain); + const drained = yield* Effect.forkChild(reactor.drainThrough(2)); yield* Effect.yieldNow; yield* Effect.yieldNow; expect(stops).toEqual([1]); diff --git a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts index 28ed3adf094e..14a92a5eaef5 100644 --- a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts +++ b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts @@ -81,10 +81,10 @@ const make = Effect.gen(function* () { const worker = yield* makeDrainableWorker(processThreadDeletedSafely); - // Highest event sequence the subscriber has handed to the worker. Events - // are published after they commit, so a caller that reads latestSequence - // and waits for this to catch up knows every deletion up to that point is - // at least enqueued; the worker drain then covers the in-flight ones. + // Highest event sequence the subscriber has handed to the worker. Waiting + // through a successful thread.created sequence covers every deletion that + // was ahead of that create in the engine queue; the worker drain then covers + // the in-flight cleanup. const seenSequence = yield* SubscriptionRef.make(0); const noteSeen = (sequence: number) => SubscriptionRef.update(seenSequence, (seen) => Math.max(seen, sequence)); @@ -105,8 +105,9 @@ const make = Effect.gen(function* () { ); }); - const drain: ThreadDeletionReactorShape["drain"] = Effect.gen(function* () { - const target = yield* orchestrationEngine.latestSequence; + const drainThrough: ThreadDeletionReactorShape["drainThrough"] = Effect.fn( + "ThreadDeletionReactor.drainThrough", + )(function* (target) { yield* SubscriptionRef.changes(seenSequence).pipe( Stream.filter((seen) => seen >= target), Stream.runHead, @@ -116,7 +117,7 @@ const make = Effect.gen(function* () { return { start, - drain, + drainThrough, } satisfies ThreadDeletionReactorShape; }); diff --git a/apps/server/src/orchestration/Services/ThreadDeletionReactor.ts b/apps/server/src/orchestration/Services/ThreadDeletionReactor.ts index c47415c019c3..cdbb70919a8e 100644 --- a/apps/server/src/orchestration/Services/ThreadDeletionReactor.ts +++ b/apps/server/src/orchestration/Services/ThreadDeletionReactor.ts @@ -23,12 +23,12 @@ export interface ThreadDeletionReactorShape { readonly start: () => Effect.Effect; /** - * Resolves once every thread.deleted published up to now has been handed - * to the worker and the worker's queue is empty and idle. Callers that - * re-create a deleted thread id wait on this so the old incarnation's - * cleanup cannot run against the new one. + * Resolves once every thread.deleted at or before the supplied event + * sequence has been handed to the worker and the worker is empty and idle. + * A successful thread.create sequence is the fence callers use before the + * new incarnation can own runtime resources. */ - readonly drain: Effect.Effect; + readonly drainThrough: (sequence: number) => Effect.Effect; } /** diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 22064922a211..f4a75e4bc305 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -798,7 +798,7 @@ const buildAppUnderTest = (options?: { }), Layer.mock(ThreadDeletionReactor)({ start: () => Effect.void, - drain: Effect.void, + drainThrough: () => Effect.void, ...options?.layers?.threadDeletionReactor, }), ), @@ -8449,22 +8449,24 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }).pipe(Effect.provide(NodeHttpServer.layerTest)), ); - it.effect("waits for deletion cleanup before re-creating a thread id", () => + it.effect("drains deletion cleanup through the re-created thread event", () => Effect.gen(function* () { // A draft retry reuses the thread id its failed bootstrap deleted. The // deletion reactor stops sessions and closes terminals by that id, so - // both thread.create paths must let its queue drain first. + // both thread.create paths use the created event as a fence, then drain + // cleanup before handing the new incarnation to resource-owning work. const trace: Array = []; const drainRequested = yield* Deferred.make(); const cleanupDone = yield* Deferred.make(); yield* buildAppUnderTest({ layers: { threadDeletionReactor: { - drain: Effect.gen(function* () { - trace.push("drain"); - yield* Deferred.succeed(drainRequested, undefined); - yield* Deferred.await(cleanupDone); - }), + drainThrough: (sequence) => + Effect.gen(function* () { + trace.push(`drain:${sequence}`); + yield* Deferred.succeed(drainRequested, undefined); + yield* Deferred.await(cleanupDone); + }), }, orchestrationEngine: { dispatch: (command) => @@ -8500,15 +8502,16 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }), ); yield* Deferred.await(drainRequested); - assert.deepEqual(trace, ["drain"]); + assert.deepEqual(trace, ["thread.create", "drain:1"]); yield* Deferred.succeed(cleanupDone, undefined); yield* Fiber.join(directCreate); }), ), ); - assert.deepEqual(trace, ["drain", "thread.create"]); + assert.deepEqual(trace, ["thread.create", "drain:1"]); - // Cleanup is already released; the bootstrap path must still drain first. + // Cleanup is already released; the bootstrap path must still drain + // between creating the thread and starting its turn. trace.length = 0; yield* Effect.scoped( withWsRpcClient(wsUrl, (client) => @@ -8547,7 +8550,7 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }), ), ); - assert.deepEqual(trace, ["drain", "thread.create", "thread.turn.start"]); + assert.deepEqual(trace, ["thread.create", "drain:1", "thread.turn.start"]); }).pipe(Effect.provide(NodeHttpServer.layerTest)), ); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index a9720af9f934..00b2cc589cf1 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -462,11 +462,6 @@ const makeWsRpcLayer = ( const projectionSnapshotQuery = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; const orchestrationEngine = yield* OrchestrationEngine.OrchestrationEngineService; const threadDeletionReactor = yield* ThreadDeletionReactor; - // A draft retry re-creates the thread id its failed bootstrap deleted. - // Session stop and terminal close are keyed by that id and run in the - // deletion reactor's queue, so let the old incarnation's cleanup finish - // before the new one can own those resources. - const awaitThreadDeletionCleanup = threadDeletionReactor.drain; const analytics = yield* AnalyticsService.AnalyticsService; // Every command dispatched on this connection carries the connecting // client's origin, including server-generated bootstrap sub-commands: @@ -1033,8 +1028,7 @@ const makeWsRpcLayer = ( const bootstrapProgram = Effect.gen(function* () { if (bootstrap?.createThread) { - yield* awaitThreadDeletionCleanup; - yield* dispatchFromClient({ + const created = yield* dispatchFromClient({ type: "thread.create", commandId: yield* serverCommandId("bootstrap-thread-create"), threadId: command.threadId, @@ -1047,6 +1041,11 @@ const makeWsRpcLayer = ( worktreePath: bootstrap.createThread.worktreePath, createdAt: bootstrap.createThread.createdAt, }); + // The successful create is a fence in the engine command queue: + // every delete for the prior incarnation committed before it. + // Drain through that event before setup or turn start can own + // terminals and provider sessions under the reused thread id. + yield* threadDeletionReactor.drainThrough(created.sequence); createdThread = true; } @@ -1133,11 +1132,15 @@ const makeWsRpcLayer = ( const dispatchEffect = normalizedCommand.type === "thread.turn.start" && normalizedCommand.bootstrap ? dispatchBootstrapTurnStart(normalizedCommand) - : (normalizedCommand.type === "thread.create" - ? awaitThreadDeletionCleanup - : Effect.void - ).pipe( - Effect.andThen(dispatchFromClient(normalizedCommand)), + : dispatchFromClient(normalizedCommand).pipe( + Effect.tap(({ sequence }) => + // Returning from thread.create is the handoff point at which + // clients may start resources for the new incarnation. Use + // its event sequence as the exact deletion-cleanup fence. + normalizedCommand.type === "thread.create" + ? threadDeletionReactor.drainThrough(sequence) + : Effect.void, + ), Effect.mapError((cause) => toDispatchCommandError(cause, "Failed to dispatch orchestration command"), ),