diff --git a/apps/server/integration/orchestrationEngine.integration.test.ts b/apps/server/integration/orchestrationEngine.integration.test.ts index ccfb9c467421..b922cadd6d4c 100644 --- a/apps/server/integration/orchestrationEngine.integration.test.ts +++ b/apps/server/integration/orchestrationEngine.integration.test.ts @@ -851,6 +851,7 @@ it.live("reverts to an earlier checkpoint and trims checkpoint projections + git [ { role: "user", text: "First edit" }, { role: "assistant", text: "Updated README to v2.\n" }, + { role: "user", text: "Second edit" }, ], ); assert.equal( diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts index aaedda18c71e..b6bbe0b5c2d2 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts @@ -1265,7 +1265,7 @@ describe("CheckpointReactor", () => { ); await waitForEvent(harness.engine, (event) => event.type === "thread.reverted"); - const thread = await waitForThread(harness.readModel, (entry) => entry.messages.length === 0); + const thread = await waitForThread(harness.readModel, (entry) => entry.messages.length === 1); expect(thread.checkpoints).toHaveLength(0); expect(harness.provider.rollbackConversation).toHaveBeenCalledTimes(1); diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index f4b2dc66a02e..ad929a6f55d1 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -1332,6 +1332,283 @@ it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-projection-atta }, ); +it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-bootstrap-attachments-")))( + "OrchestrationProjectionPipeline attachment bootstrap reconciliation", + (it) => { + it.effect("reconciles files once from the final replayed projection", () => + 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-02T00:00:00.000Z"; + const liveThreadId = ThreadId.make("bootstrap-attachment-live"); + const deletedThreadId = ThreadId.make("bootstrap-attachment-deleted"); + const oldAttachmentId = "bootstrap-attachment-live-00000000-0000-4000-8000-000000000001"; + const checkpointlessAttachmentId = + "bootstrap-attachment-live-00000000-0000-4000-8000-000000000006"; + const laterAttachmentId = "bootstrap-attachment-live-00000000-0000-4000-8000-000000000002"; + const orphanAttachmentId = "bootstrap-attachment-live-00000000-0000-4000-8000-000000000003"; + const deletedAttachmentId = + "bootstrap-attachment-deleted-00000000-0000-4000-8000-000000000004"; + + yield* eventStore.append({ + type: "thread.created", + eventId: EventId.make("evt-bootstrap-attachments-1"), + aggregateKind: "thread", + aggregateId: liveThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-1"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-1"), + metadata: {}, + payload: { + threadId: liveThreadId, + projectId: ProjectId.make("project-bootstrap-attachments"), + title: "Live attachments", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.message-sent", + eventId: EventId.make("evt-bootstrap-attachments-checkpointless"), + aggregateKind: "thread", + aggregateId: liveThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-checkpointless"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-checkpointless"), + metadata: {}, + payload: { + threadId: liveThreadId, + messageId: MessageId.make("message-bootstrap-checkpointless"), + role: "assistant", + text: "checkpointless history", + attachments: [ + { + type: "image", + id: checkpointlessAttachmentId, + name: "checkpointless.png", + mimeType: "image/png", + sizeBytes: 4, + }, + ], + turnId: TurnId.make("turn-bootstrap-checkpointless"), + streaming: false, + createdAt: now, + updatedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.turn-diff-completed", + eventId: EventId.make("evt-bootstrap-attachments-2"), + aggregateKind: "thread", + aggregateId: liveThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-2"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-2"), + metadata: {}, + payload: { + threadId: liveThreadId, + turnId: TurnId.make("turn-bootstrap-old"), + checkpointTurnCount: 1, + checkpointRef: CheckpointRef.make( + "refs/t3/checkpoints/bootstrap-attachment-live/turn/1", + ), + status: "ready", + files: [], + assistantMessageId: MessageId.make("message-bootstrap-old"), + completedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.message-sent", + eventId: EventId.make("evt-bootstrap-attachments-3"), + aggregateKind: "thread", + aggregateId: liveThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-3"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-3"), + metadata: {}, + payload: { + threadId: liveThreadId, + messageId: MessageId.make("message-bootstrap-old"), + role: "assistant", + text: "removed by the historical revert", + attachments: [ + { + type: "image", + id: oldAttachmentId, + name: "old.png", + mimeType: "image/png", + sizeBytes: 3, + }, + ], + turnId: TurnId.make("turn-bootstrap-old"), + streaming: false, + createdAt: now, + updatedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.reverted", + eventId: EventId.make("evt-bootstrap-attachments-4"), + aggregateKind: "thread", + aggregateId: liveThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-4"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-4"), + metadata: {}, + payload: { + threadId: liveThreadId, + turnCount: 0, + }, + }); + yield* eventStore.append({ + type: "thread.message-sent", + eventId: EventId.make("evt-bootstrap-attachments-5"), + aggregateKind: "thread", + aggregateId: liveThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-5"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-5"), + metadata: {}, + payload: { + threadId: liveThreadId, + messageId: MessageId.make("message-bootstrap-later"), + role: "user", + text: "sent after the revert", + attachments: [ + { + type: "image", + id: laterAttachmentId, + name: "later.png", + mimeType: "image/png", + sizeBytes: 5, + }, + ], + turnId: null, + streaming: false, + createdAt: now, + updatedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.created", + eventId: EventId.make("evt-bootstrap-attachments-6"), + aggregateKind: "thread", + aggregateId: deletedThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-6"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-6"), + metadata: {}, + payload: { + threadId: deletedThreadId, + projectId: ProjectId.make("project-bootstrap-attachments"), + title: "Deleted attachments", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.message-sent", + eventId: EventId.make("evt-bootstrap-attachments-7"), + aggregateKind: "thread", + aggregateId: deletedThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-7"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-7"), + metadata: {}, + payload: { + threadId: deletedThreadId, + messageId: MessageId.make("message-bootstrap-deleted"), + role: "user", + text: "deleted thread", + attachments: [ + { + type: "image", + id: deletedAttachmentId, + name: "deleted.png", + mimeType: "image/png", + sizeBytes: 7, + }, + ], + turnId: null, + streaming: false, + createdAt: now, + updatedAt: now, + }, + }); + yield* eventStore.append({ + type: "thread.deleted", + eventId: EventId.make("evt-bootstrap-attachments-8"), + aggregateKind: "thread", + aggregateId: deletedThreadId, + occurredAt: now, + commandId: CommandId.make("cmd-bootstrap-attachments-8"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-bootstrap-attachments-8"), + metadata: {}, + payload: { + threadId: deletedThreadId, + deletedAt: now, + }, + }); + + const oldPath = path.join(attachmentsDir, `${oldAttachmentId}.png`); + const checkpointlessPath = path.join(attachmentsDir, `${checkpointlessAttachmentId}.png`); + const laterPath = path.join(attachmentsDir, `${laterAttachmentId}.png`); + const orphanPath = path.join(attachmentsDir, `${orphanAttachmentId}.png`); + const deletedPath = path.join(attachmentsDir, `${deletedAttachmentId}.png`); + const invalidPath = path.join(attachmentsDir, "not-an-attachment.txt"); + const attachmentLikeDirectory = path.join( + attachmentsDir, + "bootstrap-attachment-live-00000000-0000-4000-8000-000000000005.png", + ); + yield* fileSystem.makeDirectory(attachmentsDir, { recursive: true }); + yield* Effect.forEach( + [oldPath, checkpointlessPath, laterPath, orphanPath, deletedPath, invalidPath], + (filePath) => fileSystem.writeFileString(filePath, "test"), + { concurrency: 1 }, + ); + yield* fileSystem.makeDirectory(attachmentLikeDirectory); + + yield* projectionPipeline.bootstrap; + + assert.isFalse(yield* exists(oldPath)); + assert.isTrue(yield* exists(checkpointlessPath)); + assert.isTrue(yield* exists(laterPath)); + assert.isFalse(yield* exists(orphanPath)); + assert.isFalse(yield* exists(deletedPath)); + assert.isTrue(yield* exists(invalidPath)); + assert.isTrue(yield* exists(attachmentLikeDirectory)); + }), + ); + }, +); + it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { it.effect("resumes from projector last_applied_sequence without replaying older events", () => Effect.gen(function* () { @@ -2428,7 +2705,7 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { }), ); - it.effect("does not fallback-retain messages whose turnId is removed by revert", () => + it.effect("removes messages tied to post-baseline checkpoint turns", () => Effect.gen(function* () { const projectionPipeline = yield* OrchestrationProjectionPipeline; const eventStore = yield* OrchestrationEventStore; @@ -2634,7 +2911,7 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { }), ); - it.effect("excludes a retracted message from SQLite fallback retention", () => + it.effect("excludes a completed retraction message from SQLite projection", () => Effect.gen(function* () { const projectionPipeline = yield* OrchestrationProjectionPipeline; const eventStore = yield* OrchestrationEventStore; diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index 40abc81d1cbd..b28f7c9e46d1 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -59,6 +59,7 @@ import { toSafeThreadAttachmentSegment, } from "../../attachmentStore.ts"; import { checkpointRefForThreadTurn } from "../../checkpointing/Utils.ts"; +import { collectRevertedTurnIds } from "../RevertRetention.ts"; export const ORCHESTRATION_PROJECTOR_NAMES = { projects: "projection.projects", @@ -215,131 +216,33 @@ function deriveHasActionableProposedPlan(input: { function retainProjectionMessagesAfterRevert( messages: ReadonlyArray, - turns: ReadonlyArray, - turnCount: number, + revertedTurnIds: ReadonlySet, excludedMessageIds: ReadonlySet, ): ReadonlyArray { - const retainedMessageIds = new Set(); - const retainedTurnIds = new Set(); - const keptTurns = turns.filter( - (turn) => - turn.turnId !== null && - turn.checkpointTurnCount !== null && - turn.checkpointTurnCount <= turnCount, + return messages.filter( + (message) => + !excludedMessageIds.has(message.messageId) && + (message.role === "system" || + message.turnId === null || + !revertedTurnIds.has(message.turnId)), ); - for (const turn of keptTurns) { - if (turn.turnId !== null) { - retainedTurnIds.add(turn.turnId); - } - if (turn.pendingMessageId !== null && !excludedMessageIds.has(turn.pendingMessageId)) { - retainedMessageIds.add(turn.pendingMessageId); - } - if (turn.assistantMessageId !== null && !excludedMessageIds.has(turn.assistantMessageId)) { - retainedMessageIds.add(turn.assistantMessageId); - } - } - - for (const message of messages) { - if (excludedMessageIds.has(message.messageId)) { - continue; - } - if (message.role === "system") { - retainedMessageIds.add(message.messageId); - continue; - } - if (message.turnId !== null && retainedTurnIds.has(message.turnId)) { - retainedMessageIds.add(message.messageId); - } - } - - const retainedUserCount = messages.filter( - (message) => message.role === "user" && retainedMessageIds.has(message.messageId), - ).length; - const missingUserCount = Math.max(0, turnCount - retainedUserCount); - if (missingUserCount > 0) { - const fallbackUserMessages = messages - .filter( - (message) => - message.role === "user" && - !excludedMessageIds.has(message.messageId) && - !retainedMessageIds.has(message.messageId) && - (message.turnId === null || retainedTurnIds.has(message.turnId)), - ) - .toSorted( - (left, right) => - left.createdAt.localeCompare(right.createdAt) || - left.messageId.localeCompare(right.messageId), - ) - .slice(0, missingUserCount); - for (const message of fallbackUserMessages) { - retainedMessageIds.add(message.messageId); - } - } - - const retainedAssistantCount = messages.filter( - (message) => message.role === "assistant" && retainedMessageIds.has(message.messageId), - ).length; - const missingAssistantCount = Math.max(0, turnCount - retainedAssistantCount); - if (missingAssistantCount > 0) { - const fallbackAssistantMessages = messages - .filter( - (message) => - message.role === "assistant" && - !excludedMessageIds.has(message.messageId) && - !retainedMessageIds.has(message.messageId) && - (message.turnId === null || retainedTurnIds.has(message.turnId)), - ) - .toSorted( - (left, right) => - left.createdAt.localeCompare(right.createdAt) || - left.messageId.localeCompare(right.messageId), - ) - .slice(0, missingAssistantCount); - for (const message of fallbackAssistantMessages) { - retainedMessageIds.add(message.messageId); - } - } - - return messages.filter((message) => retainedMessageIds.has(message.messageId)); } function retainProjectionActivitiesAfterRevert( activities: ReadonlyArray, - turns: ReadonlyArray, - turnCount: number, + revertedTurnIds: ReadonlySet, ): ReadonlyArray { - const retainedTurnIds = new Set( - turns - .filter( - (turn) => - turn.turnId !== null && - turn.checkpointTurnCount !== null && - turn.checkpointTurnCount <= turnCount, - ) - .flatMap((turn) => (turn.turnId === null ? [] : [turn.turnId])), - ); return activities.filter( - (activity) => activity.turnId === null || retainedTurnIds.has(activity.turnId), + (activity) => activity.turnId === null || !revertedTurnIds.has(activity.turnId), ); } function retainProjectionProposedPlansAfterRevert( proposedPlans: ReadonlyArray, - turns: ReadonlyArray, - turnCount: number, + revertedTurnIds: ReadonlySet, ): ReadonlyArray { - const retainedTurnIds = new Set( - turns - .filter( - (turn) => - turn.turnId !== null && - turn.checkpointTurnCount !== null && - turn.checkpointTurnCount <= turnCount, - ) - .flatMap((turn) => (turn.turnId === null ? [] : [turn.turnId])), - ); return proposedPlans.filter( - (proposedPlan) => proposedPlan.turnId === null || retainedTurnIds.has(proposedPlan.turnId), + (proposedPlan) => proposedPlan.turnId === null || !revertedTurnIds.has(proposedPlan.turnId), ); } @@ -502,6 +405,25 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti const path = yield* Path.Path; const serverConfig = yield* ServerConfig; + const getRevertedTurnIds = Effect.fn("getRevertedTurnIds")(function* (input: { + readonly threadId: ThreadId; + readonly baselineTurnCount: number; + readonly retractionTurnId: string | null; + }) { + const [turns, thread, session] = yield* Effect.all([ + projectionTurnRepository.listByThreadId({ threadId: input.threadId }), + projectionThreadRepository.getById({ threadId: input.threadId }), + projectionThreadSessionRepository.getByThreadId({ threadId: input.threadId }), + ]); + return collectRevertedTurnIds({ + turns, + baselineTurnCount: input.baselineTurnCount, + retractionTurnId: input.retractionTurnId, + latestTurnId: Option.isSome(thread) ? thread.value.latestTurnId : null, + activeTurnId: Option.isSome(session) ? session.value.activeTurnId : null, + }); + }); + const isCompletedRetractedTurn = Effect.fn("isCompletedRetractedTurn")(function* ( threadId: ThreadId, turnId: string | null, @@ -1066,8 +988,10 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti return; } - const existingTurns = yield* projectionTurnRepository.listByThreadId({ + const revertedTurnIds = yield* getRevertedTurnIds({ threadId: event.payload.threadId, + baselineTurnCount: event.payload.turnCount, + retractionTurnId: event.payload.retraction?.turnId ?? null, }); const excludedMessageIds = event.payload.retraction === undefined @@ -1075,8 +999,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti : new Set([event.payload.retraction.messageId]); const keptRows = retainProjectionMessagesAfterRevert( existingRows, - existingTurns, - event.payload.turnCount, + revertedTurnIds, excludedMessageIds, ); if (keptRows.length === existingRows.length) { @@ -1135,14 +1058,12 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti return; } - const existingTurns = yield* projectionTurnRepository.listByThreadId({ + const revertedTurnIds = yield* getRevertedTurnIds({ threadId: event.payload.threadId, + baselineTurnCount: event.payload.turnCount, + retractionTurnId: event.payload.retraction?.turnId ?? null, }); - const keptRows = retainProjectionProposedPlansAfterRevert( - existingRows, - existingTurns, - event.payload.turnCount, - ); + const keptRows = retainProjectionProposedPlansAfterRevert(existingRows, revertedTurnIds); if (keptRows.length === existingRows.length) { return; } @@ -1188,14 +1109,12 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti if (existingRows.length === 0) { return; } - const existingTurns = yield* projectionTurnRepository.listByThreadId({ + const revertedTurnIds = yield* getRevertedTurnIds({ threadId: event.payload.threadId, + baselineTurnCount: event.payload.turnCount, + retractionTurnId: event.payload.retraction?.turnId ?? null, }); - const keptRows = retainProjectionActivitiesAfterRevert( - existingRows, - existingTurns, - event.payload.turnCount, - ); + const keptRows = retainProjectionActivitiesAfterRevert(existingRows, revertedTurnIds); if (keptRows.length === existingRows.length) { return; } @@ -1799,9 +1718,95 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti }, ]; + const reconcileAttachmentFiles = Effect.fn("reconcileAttachmentFiles")(function* () { + const entries = yield* fileSystem + .readDirectory(serverConfig.attachmentsDir, { recursive: false }) + .pipe(Effect.orElseSucceed(() => [] as Array)); + const entriesByThreadSegment = new Map< + string, + Array<{ readonly relativePath: string; readonly absolutePath: string }> + >(); + + yield* Effect.forEach( + entries, + Effect.fn("collectAttachmentFileForReconciliation")(function* (entry) { + const relativePath = entry.replace(/^[/\\]+/, "").replace(/\\/g, "/"); + if (relativePath.length === 0 || relativePath.includes("/")) { + return; + } + const attachmentId = parseAttachmentIdFromRelativePath(relativePath); + if (!attachmentId) { + return; + } + const threadSegment = parseThreadSegmentFromAttachmentId(attachmentId); + if (!threadSegment) { + return; + } + const absolutePath = path.join(serverConfig.attachmentsDir, relativePath); + const fileInfo = yield* fileSystem + .stat(absolutePath) + .pipe(Effect.orElseSucceed(() => null)); + if (!fileInfo || fileInfo.type !== "File") { + return; + } + const segmentEntries = entriesByThreadSegment.get(threadSegment) ?? []; + segmentEntries.push({ relativePath, absolutePath }); + entriesByThreadSegment.set(threadSegment, segmentEntries); + }), + { concurrency: 1 }, + ); + + const liveThreadRows = yield* sql<{ readonly threadId: string }>` + SELECT thread_id AS "threadId" + FROM projection_threads + WHERE deleted_at IS NULL + `; + const liveThreadIdsBySegment = new Map>(); + for (const row of liveThreadRows) { + const threadSegment = toSafeThreadAttachmentSegment(row.threadId); + if (!threadSegment) { + continue; + } + const threadIds = liveThreadIdsBySegment.get(threadSegment) ?? []; + threadIds.push(ThreadId.make(row.threadId)); + liveThreadIdsBySegment.set(threadSegment, threadIds); + } + + yield* Effect.forEach( + entriesByThreadSegment.entries(), + Effect.fn("reconcileThreadAttachmentFiles")(function* ([threadSegment, segmentEntries]) { + const liveThreadIds = liveThreadIdsBySegment.get(threadSegment) ?? []; + const keptRelativePaths = new Set(); + yield* Effect.forEach( + liveThreadIds, + Effect.fn("collectProjectedThreadAttachmentPaths")(function* (threadId) { + const messages = yield* projectionThreadMessageRepository.listByThreadId({ + threadId, + }); + for (const relativePath of collectThreadAttachmentRelativePaths(threadId, messages)) { + keptRelativePaths.add(relativePath); + } + }), + { concurrency: 1 }, + ); + + yield* Effect.forEach( + segmentEntries, + ({ relativePath, absolutePath }) => + liveThreadIds.length === 0 || !keptRelativePaths.has(relativePath) + ? fileSystem.remove(absolutePath, { force: true }) + : Effect.void, + { concurrency: 1 }, + ); + }), + { concurrency: 1 }, + ); + }); + const runProjectorForEvent = Effect.fn("runProjectorForEvent")(function* ( projector: ProjectorDefinition, event: OrchestrationEvent, + mode: "bootstrap" | "live", ) { const attachmentSideEffects: AttachmentSideEffects = { deletedThreadIds: new Set(), @@ -1820,36 +1825,53 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti ), ); - yield* runAttachmentSideEffects(attachmentSideEffects).pipe( - Effect.catch((cause) => - Effect.logWarning("failed to apply projected attachment side-effects", { - projector: projector.name, - sequence: event.sequence, - eventType: event.type, - cause, - }), - ), - ); + if (mode === "live") { + yield* runAttachmentSideEffects(attachmentSideEffects).pipe( + Effect.catch((cause) => + Effect.logWarning("failed to apply projected attachment side-effects", { + projector: projector.name, + sequence: event.sequence, + eventType: event.type, + cause, + }), + ), + ); + } }); - const bootstrapProjector = (projector: ProjectorDefinition) => - projectionStateRepository - .getByProjector({ + const bootstrapProjectors = Effect.gen(function* () { + const projectorStates = yield* Effect.forEach(projectors, (projector) => + projectionStateRepository.getByProjector({ projector: projector.name, - }) - .pipe( - Effect.flatMap((stateRow) => - Stream.runForEach( - eventStore.readFromSequence( - Option.isSome(stateRow) ? stateRow.value.lastAppliedSequence : 0, - ), - (event) => runProjectorForEvent(projector, event), - ), - ), - ); + }), + ); + const lastAppliedSequenceByProjector = new Map( + projectors.map((projector, index) => { + const state = projectorStates[index]; + return [ + projector.name, + state !== undefined && Option.isSome(state) ? state.value.lastAppliedSequence : 0, + ] as const; + }), + ); + const firstSequence = Math.min(...lastAppliedSequenceByProjector.values()); + + // Match live event-major ordering so a revert sees turn/session evidence + // before later projectors apply that same event's trims. + yield* Stream.runForEach(eventStore.readFromSequence(firstSequence), (event) => + Effect.forEach( + projectors, + (projector) => + event.sequence > (lastAppliedSequenceByProjector.get(projector.name) ?? 0) + ? runProjectorForEvent(projector, event, "bootstrap") + : Effect.void, + { concurrency: 1 }, + ), + ); + }); const projectEvent: OrchestrationProjectionPipelineShape["projectEvent"] = (event) => - Effect.forEach(projectors, (projector) => runProjectorForEvent(projector, event), { + Effect.forEach(projectors, (projector) => runProjectorForEvent(projector, event, "live"), { concurrency: 1, }).pipe( Effect.provideService(FileSystem.FileSystem, fileSystem), @@ -1861,11 +1883,14 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti ), ); - const bootstrap: OrchestrationProjectionPipelineShape["bootstrap"] = Effect.forEach( - projectors, - bootstrapProjector, - { concurrency: 1 }, - ).pipe( + const bootstrap: OrchestrationProjectionPipelineShape["bootstrap"] = bootstrapProjectors.pipe( + Effect.andThen( + reconcileAttachmentFiles().pipe( + Effect.catch((cause) => + Effect.logWarning("failed to reconcile projected attachment files", { cause }), + ), + ), + ), Effect.provideService(FileSystem.FileSystem, fileSystem), Effect.provideService(Path.Path, path), Effect.provideService(ServerConfig, serverConfig), diff --git a/apps/server/src/orchestration/RevertRetention.test.ts b/apps/server/src/orchestration/RevertRetention.test.ts new file mode 100644 index 000000000000..be82c874787b --- /dev/null +++ b/apps/server/src/orchestration/RevertRetention.test.ts @@ -0,0 +1,38 @@ +import { describe, expect, it } from "vite-plus/test"; + +import { collectRevertedTurnIds } from "./RevertRetention.ts"; + +describe("collectRevertedTurnIds", () => { + it("uses post-baseline checkpoints, the retraction target, and uncheckpointed head evidence", () => { + const revertedTurnIds = collectRevertedTurnIds({ + turns: [ + { turnId: "turn-retained", checkpointTurnCount: 2 }, + { turnId: "turn-after-baseline", checkpointTurnCount: 3 }, + { turnId: "turn-checkpointless-older", checkpointTurnCount: null }, + ], + baselineTurnCount: 2, + retractionTurnId: "turn-retraction-target", + latestTurnId: "turn-latest-uncheckpointed", + activeTurnId: "turn-active-uncheckpointed", + }); + + expect([...revertedTurnIds].toSorted()).toEqual([ + "turn-active-uncheckpointed", + "turn-after-baseline", + "turn-latest-uncheckpointed", + "turn-retraction-target", + ]); + }); + + it("does not treat a latest or active turn with a retained checkpoint as reverted", () => { + const revertedTurnIds = collectRevertedTurnIds({ + turns: [{ turnId: "turn-retained", checkpointTurnCount: 2 }], + baselineTurnCount: 2, + retractionTurnId: null, + latestTurnId: "turn-retained", + activeTurnId: "turn-retained", + }); + + expect([...revertedTurnIds]).toEqual([]); + }); +}); diff --git a/apps/server/src/orchestration/RevertRetention.ts b/apps/server/src/orchestration/RevertRetention.ts new file mode 100644 index 000000000000..d3858e1d7121 --- /dev/null +++ b/apps/server/src/orchestration/RevertRetention.ts @@ -0,0 +1,42 @@ +interface RevertTurnEvidence { + readonly turnId: string | null; + readonly checkpointTurnCount: number | null; +} + +/** + * Checkpoints prove which turns cross the baseline. The retraction target is + * authoritative; latest/session ids cover current work that has no checkpoint. + */ +export function collectRevertedTurnIds(input: { + readonly turns: ReadonlyArray; + readonly baselineTurnCount: number; + readonly retractionTurnId: string | null; + readonly latestTurnId: string | null; + readonly activeTurnId: string | null; +}): ReadonlySet { + const revertedTurnIds = new Set(); + const retainedCheckpointTurnIds = new Set(); + + for (const turn of input.turns) { + if (turn.turnId === null || turn.checkpointTurnCount === null) { + continue; + } + if (turn.checkpointTurnCount > input.baselineTurnCount) { + revertedTurnIds.add(turn.turnId); + } else { + retainedCheckpointTurnIds.add(turn.turnId); + } + } + + if (input.retractionTurnId !== null) { + revertedTurnIds.add(input.retractionTurnId); + } + + for (const fallbackTurnId of [input.latestTurnId, input.activeTurnId]) { + if (fallbackTurnId !== null && !retainedCheckpointTurnIds.has(fallbackTurnId)) { + revertedTurnIds.add(fallbackTurnId); + } + } + + return revertedTurnIds; +} diff --git a/apps/server/src/orchestration/projector.test.ts b/apps/server/src/orchestration/projector.test.ts index e0ca8fd9d241..b20cbf4e6109 100644 --- a/apps/server/src/orchestration/projector.test.ts +++ b/apps/server/src/orchestration/projector.test.ts @@ -749,6 +749,7 @@ describe("orchestration projector", () => { [ { role: "user", text: "First edit" }, { role: "assistant", text: "Updated README to v2.\n" }, + { role: "user", text: "Second edit" }, ], ); expect( @@ -758,7 +759,320 @@ describe("orchestration projector", () => { expect(thread?.latestTurn?.turnId).toBe("turn-1"); }); - it("does not fallback-retain messages tied to removed turn IDs", async () => { + it("retains checkpointless history across sequential reverts", async () => { + const createdAt = "2026-02-24T10:00:00.000Z"; + const afterCreate = await Effect.runPromise( + projectEvent( + createEmptyReadModel(createdAt), + makeEvent({ + sequence: 1, + type: "thread.created", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: createdAt, + commandId: "cmd-denylist-create", + payload: { + threadId: "thread-denylist", + projectId: "project-1", + title: "denylist", + modelSelection: { + provider: ProviderDriverKind.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt, + updatedAt: createdAt, + }, + }), + ), + ); + + const beforeFirstRevert: ReadonlyArray = [ + makeEvent({ + sequence: 2, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:01.000Z", + commandId: "cmd-early-user", + payload: { + threadId: "thread-denylist", + messageId: "message-early-user", + role: "user", + text: "checkpointless head", + turnId: "turn-early-checkpointless", + streaming: false, + createdAt: "2026-02-24T10:00:01.000Z", + updatedAt: "2026-02-24T10:00:01.000Z", + }, + }), + makeEvent({ + sequence: 3, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:02.000Z", + commandId: "cmd-early-assistant", + payload: { + threadId: "thread-denylist", + messageId: "message-early-assistant", + role: "assistant", + text: "checkpointless response", + turnId: "turn-early-checkpointless", + streaming: false, + createdAt: "2026-02-24T10:00:02.000Z", + updatedAt: "2026-02-24T10:00:02.000Z", + }, + }), + makeEvent({ + sequence: 4, + type: "thread.activity-appended", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:03.000Z", + commandId: "cmd-early-activity", + payload: { + threadId: "thread-denylist", + activity: { + id: "activity-early", + tone: "tool", + kind: "tool.completed", + summary: "checkpointless activity", + payload: {}, + turnId: "turn-early-checkpointless", + createdAt: "2026-02-24T10:00:03.000Z", + }, + }, + }), + makeEvent({ + sequence: 5, + type: "thread.proposed-plan-upserted", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:04.000Z", + commandId: "cmd-early-plan", + payload: { + threadId: "thread-denylist", + proposedPlan: { + id: "plan-early", + turnId: "turn-early-checkpointless", + planMarkdown: "Keep the head", + implementedAt: null, + implementationThreadId: null, + createdAt: "2026-02-24T10:00:04.000Z", + updatedAt: "2026-02-24T10:00:04.000Z", + }, + }, + }), + makeEvent({ + sequence: 6, + type: "thread.turn-diff-completed", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:05.000Z", + commandId: "cmd-kept-checkpoint", + payload: { + threadId: "thread-denylist", + turnId: "turn-kept", + checkpointTurnCount: 1, + checkpointRef: "refs/t3/checkpoints/thread-denylist/turn/1", + status: "ready", + files: [], + assistantMessageId: "message-kept", + completedAt: "2026-02-24T10:00:05.000Z", + }, + }), + makeEvent({ + sequence: 7, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:06.000Z", + commandId: "cmd-kept-message", + payload: { + threadId: "thread-denylist", + messageId: "message-kept", + role: "assistant", + text: "kept checkpoint", + turnId: "turn-kept", + streaming: false, + createdAt: "2026-02-24T10:00:06.000Z", + updatedAt: "2026-02-24T10:00:06.000Z", + }, + }), + makeEvent({ + sequence: 8, + type: "thread.turn-diff-completed", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:07.000Z", + commandId: "cmd-removed-checkpoint", + payload: { + threadId: "thread-denylist", + turnId: "turn-removed", + checkpointTurnCount: 2, + checkpointRef: "refs/t3/checkpoints/thread-denylist/turn/2", + status: "ready", + files: [], + assistantMessageId: "message-removed", + completedAt: "2026-02-24T10:00:07.000Z", + }, + }), + makeEvent({ + sequence: 9, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:08.000Z", + commandId: "cmd-removed-message", + payload: { + threadId: "thread-denylist", + messageId: "message-removed", + role: "assistant", + text: "remove this turn", + turnId: "turn-removed", + streaming: false, + createdAt: "2026-02-24T10:00:08.000Z", + updatedAt: "2026-02-24T10:00:08.000Z", + }, + }), + makeEvent({ + sequence: 10, + type: "thread.activity-appended", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:09.000Z", + commandId: "cmd-removed-activity", + payload: { + threadId: "thread-denylist", + activity: { + id: "activity-removed", + tone: "tool", + kind: "tool.completed", + summary: "removed activity", + payload: {}, + turnId: "turn-removed", + createdAt: "2026-02-24T10:00:09.000Z", + }, + }, + }), + makeEvent({ + sequence: 11, + type: "thread.proposed-plan-upserted", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:10.000Z", + commandId: "cmd-removed-plan", + payload: { + threadId: "thread-denylist", + proposedPlan: { + id: "plan-removed", + turnId: "turn-removed", + planMarkdown: "Remove this plan", + implementedAt: null, + implementationThreadId: null, + createdAt: "2026-02-24T10:00:10.000Z", + updatedAt: "2026-02-24T10:00:10.000Z", + }, + }, + }), + makeEvent({ + sequence: 12, + type: "thread.reverted", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:11.000Z", + commandId: "cmd-first-revert", + payload: { + threadId: "thread-denylist", + turnCount: 1, + }, + }), + ]; + const afterFirstRevert = await beforeFirstRevert.reduce< + Promise> + >( + (statePromise, event) => + statePromise.then((state) => Effect.runPromise(projectEvent(state, event))), + Promise.resolve(afterCreate), + ); + + expect(afterFirstRevert.threads[0]?.messages.map((message) => message.id)).toEqual([ + "message-early-user", + "message-early-assistant", + "message-kept", + ]); + expect(afterFirstRevert.threads[0]?.activities.map((activity) => activity.id)).toEqual([ + "activity-early", + ]); + expect(afterFirstRevert.threads[0]?.proposedPlans.map((plan) => plan.id)).toEqual([ + "plan-early", + ]); + + const afterRetractedMessage = await Effect.runPromise( + projectEvent( + afterFirstRevert, + makeEvent({ + sequence: 13, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:12.000Z", + commandId: "cmd-retracted-message", + payload: { + threadId: "thread-denylist", + messageId: "message-retracted", + role: "user", + text: "retract me", + turnId: "turn-retracted", + streaming: false, + createdAt: "2026-02-24T10:00:12.000Z", + updatedAt: "2026-02-24T10:00:12.000Z", + }, + }), + ), + ); + const afterSecondRevert = await Effect.runPromise( + projectEvent( + afterRetractedMessage, + makeEvent({ + sequence: 14, + type: "thread.reverted", + aggregateKind: "thread", + aggregateId: "thread-denylist", + occurredAt: "2026-02-24T10:00:13.000Z", + commandId: "cmd-second-revert", + payload: { + threadId: "thread-denylist", + turnCount: 1, + retraction: { + requestId: "request-second-revert", + messageId: "message-retracted", + turnId: "turn-retracted", + firstUserMessage: false, + completedAt: "2026-02-24T10:00:13.000Z", + }, + }, + }), + ), + ); + + expect(afterSecondRevert.threads[0]?.messages.map((message) => message.id)).toEqual([ + "message-early-user", + "message-early-assistant", + "message-kept", + ]); + expect(afterSecondRevert.threads[0]?.activities.map((activity) => activity.id)).toEqual([ + "activity-early", + ]); + expect(afterSecondRevert.threads[0]?.proposedPlans.map((plan) => plan.id)).toEqual([ + "plan-early", + ]); + }); + + it("removes messages tied to post-baseline checkpoint turns", async () => { const createdAt = "2026-02-26T12:00:00.000Z"; const model = createEmptyReadModel(createdAt); @@ -911,7 +1225,7 @@ describe("orchestration projector", () => { ).toEqual([{ id: "assistant-keep", role: "assistant", turnId: "turn-1" }]); }); - it("excludes retracted messages while preserving fallback retention on long threads", async () => { + it("excludes retracted messages while preserving long thread history", async () => { const createdAt = "2026-03-01T09:00:00.000Z"; const afterCreate = await Effect.runPromise( projectEvent( diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index 7f468cf6fbb7..4d88c65236c3 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -16,6 +16,7 @@ import * as Schema from "effect/Schema"; import { checkpointRefForThreadTurn } from "../checkpointing/Utils.ts"; import { toProjectorDecodeError, type OrchestrationProjectorDecodeError } from "./Errors.ts"; +import { collectRevertedTurnIds } from "./RevertRetention.ts"; import { MessageSentPayloadSchema, ProjectCreatedPayload, @@ -117,88 +118,33 @@ function decodeForEvent( function retainThreadMessagesAfterRevert( messages: ReadonlyArray, - retainedTurnIds: ReadonlySet, - turnCount: number, + revertedTurnIds: ReadonlySet, excludedMessageIds: ReadonlySet, ): ReadonlyArray { - const retainedMessageIds = new Set(); - for (const message of messages) { - if (excludedMessageIds.has(message.id)) { - continue; - } - if (message.role === "system") { - retainedMessageIds.add(message.id); - continue; - } - if (message.turnId !== null && retainedTurnIds.has(message.turnId)) { - retainedMessageIds.add(message.id); - } - } - - const retainedUserCount = messages.filter( - (message) => message.role === "user" && retainedMessageIds.has(message.id), - ).length; - const missingUserCount = Math.max(0, turnCount - retainedUserCount); - if (missingUserCount > 0) { - const fallbackUserMessages = messages - .filter( - (message) => - message.role === "user" && - !excludedMessageIds.has(message.id) && - !retainedMessageIds.has(message.id) && - (message.turnId === null || retainedTurnIds.has(message.turnId)), - ) - .toSorted( - (left, right) => - left.createdAt.localeCompare(right.createdAt) || left.id.localeCompare(right.id), - ) - .slice(0, missingUserCount); - for (const message of fallbackUserMessages) { - retainedMessageIds.add(message.id); - } - } - - const retainedAssistantCount = messages.filter( - (message) => message.role === "assistant" && retainedMessageIds.has(message.id), - ).length; - const missingAssistantCount = Math.max(0, turnCount - retainedAssistantCount); - if (missingAssistantCount > 0) { - const fallbackAssistantMessages = messages - .filter( - (message) => - message.role === "assistant" && - !excludedMessageIds.has(message.id) && - !retainedMessageIds.has(message.id) && - (message.turnId === null || retainedTurnIds.has(message.turnId)), - ) - .toSorted( - (left, right) => - left.createdAt.localeCompare(right.createdAt) || left.id.localeCompare(right.id), - ) - .slice(0, missingAssistantCount); - for (const message of fallbackAssistantMessages) { - retainedMessageIds.add(message.id); - } - } - - return messages.filter((message) => retainedMessageIds.has(message.id)); + return messages.filter( + (message) => + !excludedMessageIds.has(message.id) && + (message.role === "system" || + message.turnId === null || + !revertedTurnIds.has(message.turnId)), + ); } function retainThreadActivitiesAfterRevert( activities: ReadonlyArray, - retainedTurnIds: ReadonlySet, + revertedTurnIds: ReadonlySet, ): ReadonlyArray { return activities.filter( - (activity) => activity.turnId === null || retainedTurnIds.has(activity.turnId), + (activity) => activity.turnId === null || !revertedTurnIds.has(activity.turnId), ); } function retainThreadProposedPlansAfterRevert( proposedPlans: ReadonlyArray, - retainedTurnIds: ReadonlySet, + revertedTurnIds: ReadonlySet, ): ReadonlyArray { return proposedPlans.filter( - (proposedPlan) => proposedPlan.turnId === null || retainedTurnIds.has(proposedPlan.turnId), + (proposedPlan) => proposedPlan.turnId === null || !revertedTurnIds.has(proposedPlan.turnId), ); } @@ -828,22 +774,27 @@ export function projectEvent( .filter((entry) => entry.checkpointTurnCount <= payload.turnCount) .toSorted((left, right) => left.checkpointTurnCount - right.checkpointTurnCount) .slice(-MAX_THREAD_CHECKPOINTS); - const retainedTurnIds = new Set(checkpoints.map((checkpoint) => checkpoint.turnId)); + const revertedTurnIds = collectRevertedTurnIds({ + turns: thread.checkpoints, + baselineTurnCount: payload.turnCount, + retractionTurnId: payload.retraction?.turnId ?? null, + latestTurnId: thread.latestTurn?.turnId ?? null, + activeTurnId: thread.session?.activeTurnId ?? null, + }); const excludedMessageIds = payload.retraction === undefined ? new Set() : new Set([payload.retraction.messageId]); const messages = retainThreadMessagesAfterRevert( thread.messages, - retainedTurnIds, - payload.turnCount, + revertedTurnIds, excludedMessageIds, ).slice(-MAX_THREAD_MESSAGES); const proposedPlans = retainThreadProposedPlansAfterRevert( thread.proposedPlans, - retainedTurnIds, + revertedTurnIds, ).slice(-200); - const activities = retainThreadActivitiesAfterRevert(thread.activities, retainedTurnIds); + const activities = retainThreadActivitiesAfterRevert(thread.activities, revertedTurnIds); const latestCheckpoint = checkpoints.at(-1) ?? null; const latestTurn = diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index 862fe53b9d1b..a14e914baf97 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -57,6 +57,7 @@ import Migration0041 from "./Migrations/041_ProjectionTurnRetractions.ts"; import Migration0042 from "./Migrations/042_ProjectionTurnDispatchOwnership.ts"; import Migration0043 from "./Migrations/043_ProjectionManagedWorktrees.ts"; import Migration0044 from "./Migrations/044_CleanupCompletedRetractionMessages.ts"; +import Migration0045 from "./Migrations/045_RebuildProjectionsFromEvents.ts"; /** * Migration loader with all migrations defined inline. @@ -113,6 +114,7 @@ export const migrationEntries = [ [42, "ProjectionTurnDispatchOwnership", Migration0042], [43, "ProjectionManagedWorktrees", Migration0043], [44, "CleanupCompletedRetractionMessages", Migration0044], + [45, "RebuildProjectionsFromEvents", Migration0045], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/045_RebuildProjectionsFromEvents.test.ts b/apps/server/src/persistence/Migrations/045_RebuildProjectionsFromEvents.test.ts new file mode 100644 index 000000000000..7eb840090aa3 --- /dev/null +++ b/apps/server/src/persistence/Migrations/045_RebuildProjectionsFromEvents.test.ts @@ -0,0 +1,99 @@ +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +import { migrationManifest, runMigrations } from "../Migrations.ts"; +import rebuildProjectionsFromEvents, { + projectionTableNames, +} from "./045_RebuildProjectionsFromEvents.ts"; +import * as NodeSqliteClient from "../NodeSqliteClient.ts"; + +const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layerMemory())); + +layer("045_RebuildProjectionsFromEvents", (it) => { + it.effect("clears every event-derived projection and resets cursors idempotently", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + + for (const tableName of projectionTableNames) { + yield* sql.unsafe(`CREATE TABLE ${tableName} (value TEXT NOT NULL)`).unprepared; + yield* sql.unsafe(`INSERT INTO ${tableName} (value) VALUES ('projected')`).unprepared; + } + yield* sql` + CREATE TABLE projection_state ( + projector TEXT PRIMARY KEY, + last_applied_sequence INTEGER NOT NULL, + updated_at TEXT NOT NULL + ) + `; + yield* sql` + INSERT INTO projection_state (projector, last_applied_sequence, updated_at) + VALUES ('projection.threads', 42, '2026-01-01T00:00:00.000Z') + `; + + yield* rebuildProjectionsFromEvents; + yield* rebuildProjectionsFromEvents; + + for (const tableName of projectionTableNames) { + const rows = yield* sql.unsafe<{ readonly count: number }>( + `SELECT COUNT(*) AS count FROM ${tableName}`, + ).unprepared; + assert.equal(rows[0]?.count, 0); + } + const stateRows = yield* sql<{ + readonly projector: string; + readonly lastAppliedSequence: number; + }>` + SELECT + projector, + last_applied_sequence AS "lastAppliedSequence" + FROM projection_state + `; + assert.deepEqual(stateRows, [{ projector: "projection.threads", lastAppliedSequence: 0 }]); + + yield* sql`DROP TABLE projection_thread_proposed_plans`; + yield* rebuildProjectionsFromEvents; + assert.deepEqual(migrationManifest.at(-1), [45, "RebuildProjectionsFromEvents"]); + }), + ); +}); + +it.layer(Layer.fresh(Layer.mergeAll(NodeSqliteClient.layerMemory())))( + "045_RebuildProjectionsFromEvents registration", + (it) => { + it.effect("runs after migration 044 against the full projection schema", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + + yield* runMigrations({ toMigrationInclusive: 44 }); + yield* sql` + INSERT INTO projection_thread_messages ( + message_id, thread_id, turn_id, role, text, is_streaming, created_at, updated_at + ) VALUES ( + 'message-before-rebuild', 'thread-before-rebuild', NULL, 'user', 'stale', 0, + '2026-01-01T00:00:00.000Z', '2026-01-01T00:00:00.000Z' + ) + `; + yield* sql` + INSERT INTO projection_state (projector, last_applied_sequence, updated_at) + VALUES ('projection.thread-messages', 99, '2026-01-01T00:00:00.000Z') + `; + + const migrations = yield* runMigrations({ toMigrationInclusive: 45 }); + assert.deepEqual(migrations, [[45, "RebuildProjectionsFromEvents"]]); + + const messageRows = yield* sql<{ readonly count: number }>` + SELECT COUNT(*) AS count + FROM projection_thread_messages + `; + assert.equal(messageRows[0]?.count, 0); + const stateRows = yield* sql<{ readonly lastAppliedSequence: number }>` + SELECT last_applied_sequence AS "lastAppliedSequence" + FROM projection_state + `; + assert.deepEqual(stateRows, [{ lastAppliedSequence: 0 }]); + }), + ); + }, +); diff --git a/apps/server/src/persistence/Migrations/045_RebuildProjectionsFromEvents.ts b/apps/server/src/persistence/Migrations/045_RebuildProjectionsFromEvents.ts new file mode 100644 index 000000000000..2011ae271d47 --- /dev/null +++ b/apps/server/src/persistence/Migrations/045_RebuildProjectionsFromEvents.ts @@ -0,0 +1,57 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +const projectionTableNames = [ + "projection_projects", + "projection_threads", + "projection_thread_messages", + "projection_thread_activities", + "projection_thread_sessions", + "projection_turns", + "projection_pending_approvals", + "projection_thread_proposed_plans", + "projection_turn_retractions", +] as const; + +export default Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const tables = yield* sql<{ readonly name: string }>` + SELECT name + FROM sqlite_master + WHERE type = 'table' + `; + const tableNames = new Set(tables.map((table) => table.name)); + + if (tableNames.has("projection_projects")) { + yield* sql`DELETE FROM projection_projects`; + } + if (tableNames.has("projection_threads")) { + yield* sql`DELETE FROM projection_threads`; + } + if (tableNames.has("projection_thread_messages")) { + yield* sql`DELETE FROM projection_thread_messages`; + } + if (tableNames.has("projection_thread_activities")) { + yield* sql`DELETE FROM projection_thread_activities`; + } + if (tableNames.has("projection_thread_sessions")) { + yield* sql`DELETE FROM projection_thread_sessions`; + } + if (tableNames.has("projection_turns")) { + yield* sql`DELETE FROM projection_turns`; + } + if (tableNames.has("projection_pending_approvals")) { + yield* sql`DELETE FROM projection_pending_approvals`; + } + if (tableNames.has("projection_thread_proposed_plans")) { + yield* sql`DELETE FROM projection_thread_proposed_plans`; + } + if (tableNames.has("projection_turn_retractions")) { + yield* sql`DELETE FROM projection_turn_retractions`; + } + if (tableNames.has("projection_state")) { + yield* sql`UPDATE projection_state SET last_applied_sequence = 0`; + } +}); + +export { projectionTableNames }; diff --git a/apps/web/src/connection/storage.ts b/apps/web/src/connection/storage.ts index bbe234eed0de..766d56f14dfe 100644 --- a/apps/web/src/connection/storage.ts +++ b/apps/web/src/connection/storage.ts @@ -62,8 +62,9 @@ const StoredShellSnapshotJson = Schema.fromJsonString(StoredShellSnapshot); // orphaned retracted messages out-of-band: the deletion emits no events, so // a warm cache resuming via `afterSequence` would render the ghost messages // forever. Any server-side row surgery needs a bump here to reach clients. +// v6 pairs with server migration 045's full projection rebuild. const StoredThreadSnapshot = Schema.Struct({ - schemaVersion: Schema.Literal(5), + schemaVersion: Schema.Literal(6), environmentId: EnvironmentId, threadId: ThreadId, snapshot: OrchestrationThreadDetailSnapshot, @@ -571,7 +572,7 @@ export const connectionStorageLayer = Layer.effectContext( saveThread: (environmentId, snapshot) => Effect.gen(function* () { const encoded = yield* encodeStoredThreadSnapshot({ - schemaVersion: 5, + schemaVersion: 6, environmentId, threadId: snapshot.thread.id, snapshot,