From 2190d1ea70be18acd1fbf4cc6e5899a92dc7a0ca Mon Sep 17 00:00:00 2001 From: Adam Firestone Date: Wed, 12 Aug 2026 15:44:12 -0500 Subject: [PATCH] fix(server): remove retracted messages after revert Long-thread reverts could backfill a retracted null-turn user message when capped checkpoint retention left a count shortfall. Exclude completed retraction message IDs from both projection retention passes and clean existing completed-retraction orphans in migration 044. Model and harness: GPT-5 via Codex --- .../Layers/ProjectionPipeline.test.ts | 55 ++++++ .../Layers/ProjectionPipeline.ts | 15 +- .../src/orchestration/projector.test.ts | 164 ++++++++++++++++++ apps/server/src/orchestration/projector.ts | 11 ++ apps/server/src/persistence/Migrations.ts | 2 + ...CleanupCompletedRetractionMessages.test.ts | 77 ++++++++ .../044_CleanupCompletedRetractionMessages.ts | 28 +++ 7 files changed, 350 insertions(+), 2 deletions(-) create mode 100644 apps/server/src/persistence/Migrations/044_CleanupCompletedRetractionMessages.test.ts create mode 100644 apps/server/src/persistence/Migrations/044_CleanupCompletedRetractionMessages.ts diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index 6526bb584753..f4b2dc66a02e 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -2633,6 +2633,61 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { ]); }), ); + + it.effect("excludes a retracted message from SQLite fallback retention", () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const threadId = ThreadId.make("thread-retraction-message-exclusion"); + + yield* sql` + INSERT INTO projection_thread_messages ( + message_id, thread_id, turn_id, role, text, is_streaming, created_at, updated_at + ) VALUES + ( + 'message-retracted', ${threadId}, NULL, 'user', 'retracted', 0, + '2026-03-01T10:00:00.000Z', '2026-03-01T10:00:00.000Z' + ), + ( + 'message-legitimate', ${threadId}, NULL, 'user', 'legitimate', 0, + '2026-03-01T10:00:01.000Z', '2026-03-01T10:00:01.000Z' + ) + `; + + const savedEvent = yield* eventStore.append({ + type: "thread.reverted", + eventId: EventId.make("evt-retraction-message-exclusion"), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: "2026-03-01T10:00:02.000Z", + commandId: CommandId.make("cmd-retraction-message-exclusion"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-retraction-message-exclusion"), + metadata: {}, + payload: { + threadId, + turnCount: 1, + retraction: { + requestId: CommandId.make("request-retraction-message-exclusion"), + messageId: MessageId.make("message-retracted"), + turnId: null, + firstUserMessage: false, + completedAt: "2026-03-01T10:00:02.000Z", + }, + }, + }); + yield* projectionPipeline.projectEvent(savedEvent); + + const messageRows = yield* sql<{ readonly messageId: string }>` + SELECT message_id AS "messageId" + FROM projection_thread_messages + WHERE thread_id = ${threadId} + ORDER BY message_id + `; + assert.deepEqual(messageRows, [{ messageId: "message-legitimate" }]); + }), + ); }); it.layer(makeProjectionPipelinePrefixedTestLayer("t3-pending-turn-terminal-test-"))( diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index 6f544418785e..40abc81d1cbd 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -217,6 +217,7 @@ function retainProjectionMessagesAfterRevert( messages: ReadonlyArray, turns: ReadonlyArray, turnCount: number, + excludedMessageIds: ReadonlySet, ): ReadonlyArray { const retainedMessageIds = new Set(); const retainedTurnIds = new Set(); @@ -230,15 +231,18 @@ function retainProjectionMessagesAfterRevert( if (turn.turnId !== null) { retainedTurnIds.add(turn.turnId); } - if (turn.pendingMessageId !== null) { + if (turn.pendingMessageId !== null && !excludedMessageIds.has(turn.pendingMessageId)) { retainedMessageIds.add(turn.pendingMessageId); } - if (turn.assistantMessageId !== null) { + 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; @@ -257,6 +261,7 @@ function retainProjectionMessagesAfterRevert( .filter( (message) => message.role === "user" && + !excludedMessageIds.has(message.messageId) && !retainedMessageIds.has(message.messageId) && (message.turnId === null || retainedTurnIds.has(message.turnId)), ) @@ -280,6 +285,7 @@ function retainProjectionMessagesAfterRevert( .filter( (message) => message.role === "assistant" && + !excludedMessageIds.has(message.messageId) && !retainedMessageIds.has(message.messageId) && (message.turnId === null || retainedTurnIds.has(message.turnId)), ) @@ -1063,10 +1069,15 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti const existingTurns = yield* projectionTurnRepository.listByThreadId({ threadId: event.payload.threadId, }); + const excludedMessageIds = + event.payload.retraction === undefined + ? new Set() + : new Set([event.payload.retraction.messageId]); const keptRows = retainProjectionMessagesAfterRevert( existingRows, existingTurns, event.payload.turnCount, + excludedMessageIds, ); if (keptRows.length === existingRows.length) { return; diff --git a/apps/server/src/orchestration/projector.test.ts b/apps/server/src/orchestration/projector.test.ts index 329c56f69171..e0ca8fd9d241 100644 --- a/apps/server/src/orchestration/projector.test.ts +++ b/apps/server/src/orchestration/projector.test.ts @@ -911,6 +911,170 @@ describe("orchestration projector", () => { ).toEqual([{ id: "assistant-keep", role: "assistant", turnId: "turn-1" }]); }); + it("excludes retracted messages while preserving fallback retention on long threads", async () => { + const createdAt = "2026-03-01T09:00:00.000Z"; + const afterCreate = await Effect.runPromise( + projectEvent( + createEmptyReadModel(createdAt), + makeEvent({ + sequence: 1, + type: "thread.created", + aggregateKind: "thread", + aggregateId: "thread-long-retraction", + occurredAt: createdAt, + commandId: "cmd-create-long-retraction", + payload: { + threadId: "thread-long-retraction", + projectId: "project-1", + title: "long retraction", + modelSelection: { + provider: ProviderDriverKind.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt, + updatedAt: createdAt, + }, + }), + ), + ); + + const checkpointEvents: ReadonlyArray = Array.from( + { length: 501 }, + (_, index) => + makeEvent({ + sequence: index + 2, + type: "thread.turn-diff-completed", + aggregateKind: "thread", + aggregateId: "thread-long-retraction", + occurredAt: "2026-03-01T09:00:01.000Z", + commandId: `cmd-checkpoint-${index}`, + payload: { + threadId: "thread-long-retraction", + turnId: `turn-${index}`, + checkpointTurnCount: index + 1, + checkpointRef: `refs/t3/checkpoints/thread-long-retraction/turn/${index + 1}`, + status: "ready", + files: [], + assistantMessageId: null, + completedAt: "2026-03-01T09:00:01.000Z", + }, + }), + ); + const afterCheckpoints = await checkpointEvents.reduce< + Promise> + >( + (statePromise, event) => + statePromise.then((state) => Effect.runPromise(projectEvent(state, event))), + Promise.resolve(afterCreate), + ); + + const legitimateMessageEvents: ReadonlyArray = Array.from( + { length: 501 }, + (_, index) => + makeEvent({ + sequence: index + 503, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-long-retraction", + occurredAt: "2026-03-01T09:00:02.000Z", + commandId: `cmd-message-${index}`, + payload: { + threadId: "thread-long-retraction", + messageId: `legitimate-user-${String(index).padStart(3, "0")}`, + role: "user", + text: `legitimate ${index}`, + turnId: null, + streaming: false, + createdAt: "2026-03-01T09:00:02.000Z", + updatedAt: "2026-03-01T09:00:02.000Z", + }, + }), + ); + const longThreadState = await legitimateMessageEvents.reduce< + Promise> + >( + (statePromise, event) => + statePromise.then((state) => Effect.runPromise(projectEvent(state, event))), + Promise.resolve(afterCheckpoints), + ); + + const ordinaryRevert = await Effect.runPromise( + projectEvent( + longThreadState, + makeEvent({ + sequence: 1_004, + type: "thread.reverted", + aggregateKind: "thread", + aggregateId: "thread-long-retraction", + occurredAt: "2026-03-01T09:00:03.000Z", + commandId: "cmd-ordinary-revert", + payload: { + threadId: "thread-long-retraction", + turnCount: 501, + }, + }), + ), + ); + expect(ordinaryRevert.threads[0]?.messages).toHaveLength(501); + + const withRetractedMessage = await Effect.runPromise( + projectEvent( + longThreadState, + makeEvent({ + sequence: 1_004, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-long-retraction", + occurredAt: "2026-03-01T09:00:03.000Z", + commandId: "cmd-retracted-message", + payload: { + threadId: "thread-long-retraction", + messageId: "000-retracted-user", + role: "user", + text: "retracted", + turnId: null, + streaming: false, + createdAt: "2026-03-01T09:00:01.500Z", + updatedAt: "2026-03-01T09:00:01.500Z", + }, + }), + ), + ); + const retractionRevert = await Effect.runPromise( + projectEvent( + withRetractedMessage, + makeEvent({ + sequence: 1_005, + type: "thread.reverted", + aggregateKind: "thread", + aggregateId: "thread-long-retraction", + occurredAt: "2026-03-01T09:00:04.000Z", + commandId: "cmd-retraction-revert", + payload: { + threadId: "thread-long-retraction", + turnCount: 501, + retraction: { + requestId: "request-retraction", + messageId: "000-retracted-user", + turnId: "turn-500", + firstUserMessage: false, + completedAt: "2026-03-01T09:00:04.000Z", + }, + }, + }), + ), + ); + + expect(retractionRevert.threads[0]?.messages).toHaveLength(501); + expect(retractionRevert.threads[0]?.messages.map((message) => message.id)).not.toContain( + "000-retracted-user", + ); + expect(retractionRevert.threads[0]?.messages.at(-1)?.id).toBe("legitimate-user-500"); + }); + it("caps message and checkpoint retention for long-lived threads", async () => { const createdAt = "2026-03-01T10:00:00.000Z"; const model = createEmptyReadModel(createdAt); diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index ad8b5ae47630..7f468cf6fbb7 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -119,9 +119,13 @@ function retainThreadMessagesAfterRevert( messages: ReadonlyArray, retainedTurnIds: ReadonlySet, turnCount: number, + 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; @@ -140,6 +144,7 @@ function retainThreadMessagesAfterRevert( .filter( (message) => message.role === "user" && + !excludedMessageIds.has(message.id) && !retainedMessageIds.has(message.id) && (message.turnId === null || retainedTurnIds.has(message.turnId)), ) @@ -162,6 +167,7 @@ function retainThreadMessagesAfterRevert( .filter( (message) => message.role === "assistant" && + !excludedMessageIds.has(message.id) && !retainedMessageIds.has(message.id) && (message.turnId === null || retainedTurnIds.has(message.turnId)), ) @@ -823,10 +829,15 @@ export function projectEvent( .toSorted((left, right) => left.checkpointTurnCount - right.checkpointTurnCount) .slice(-MAX_THREAD_CHECKPOINTS); const retainedTurnIds = new Set(checkpoints.map((checkpoint) => checkpoint.turnId)); + const excludedMessageIds = + payload.retraction === undefined + ? new Set() + : new Set([payload.retraction.messageId]); const messages = retainThreadMessagesAfterRevert( thread.messages, retainedTurnIds, payload.turnCount, + excludedMessageIds, ).slice(-MAX_THREAD_MESSAGES); const proposedPlans = retainThreadProposedPlansAfterRevert( thread.proposedPlans, diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index b185bdb89a09..862fe53b9d1b 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -56,6 +56,7 @@ import Migration0040 from "./Migrations/040_ProjectionProjectFaviconPath.ts"; 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"; /** * Migration loader with all migrations defined inline. @@ -111,6 +112,7 @@ export const migrationEntries = [ [41, "ProjectionTurnRetractions", Migration0041], [42, "ProjectionTurnDispatchOwnership", Migration0042], [43, "ProjectionManagedWorktrees", Migration0043], + [44, "CleanupCompletedRetractionMessages", Migration0044], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/044_CleanupCompletedRetractionMessages.test.ts b/apps/server/src/persistence/Migrations/044_CleanupCompletedRetractionMessages.test.ts new file mode 100644 index 000000000000..5049b7df1194 --- /dev/null +++ b/apps/server/src/persistence/Migrations/044_CleanupCompletedRetractionMessages.test.ts @@ -0,0 +1,77 @@ +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 { runMigrations } from "../Migrations.ts"; +import * as NodeSqliteClient from "../NodeSqliteClient.ts"; + +const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layerMemory())); + +layer("044_CleanupCompletedRetractionMessages", (it) => { + it.effect("removes only messages belonging to completed retractions and is idempotent", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + + yield* runMigrations({ toMigrationInclusive: 43 }); + yield* sql` + INSERT INTO projection_thread_messages ( + message_id, thread_id, turn_id, role, text, is_streaming, created_at, updated_at + ) VALUES + ( + 'message-completed', 'thread-1', NULL, 'user', 'completed', 0, + '2026-01-01T00:00:00.000Z', '2026-01-01T00:00:00.000Z' + ), + ( + 'message-requested', 'thread-1', NULL, 'user', 'requested', 0, + '2026-01-01T00:00:01.000Z', '2026-01-01T00:00:01.000Z' + ), + ( + 'message-unrelated', 'thread-1', 'turn-1', 'assistant', 'unrelated', 0, + '2026-01-01T00:00:02.000Z', '2026-01-01T00:00:02.000Z' + ) + `; + yield* sql` + INSERT INTO projection_turn_retractions ( + request_id, thread_id, message_id, baseline_turn_count, + baseline_checkpoint_ref, target_turn_id, provider_send_claimed, + first_user_message, requested_at, status, completed_at, failed_at + ) VALUES + ( + 'request-completed', 'thread-1', 'message-completed', 1, + 'refs/t3/thread/thread-1/turn/1', 'turn-1', 1, + 0, '2026-01-01T00:00:03.000Z', 'completed', + '2026-01-01T00:00:04.000Z', NULL + ), + ( + 'request-requested', 'thread-1', 'message-requested', 1, + 'refs/t3/thread/thread-1/turn/1', 'turn-1', 0, + 0, '2026-01-01T00:00:05.000Z', 'requested', NULL, NULL + ) + `; + + const firstRun = yield* runMigrations({ toMigrationInclusive: 44 }); + assert.deepEqual(firstRun, [[44, "CleanupCompletedRetractionMessages"]]); + + const rowsAfterFirstRun = yield* sql<{ readonly messageId: string }>` + SELECT message_id AS "messageId" + FROM projection_thread_messages + ORDER BY message_id + `; + assert.deepEqual(rowsAfterFirstRun, [ + { messageId: "message-requested" }, + { messageId: "message-unrelated" }, + ]); + + const secondRun = yield* runMigrations({ toMigrationInclusive: 44 }); + assert.deepEqual(secondRun, []); + + const rowsAfterSecondRun = yield* sql<{ readonly messageId: string }>` + SELECT message_id AS "messageId" + FROM projection_thread_messages + ORDER BY message_id + `; + assert.deepEqual(rowsAfterSecondRun, rowsAfterFirstRun); + }), + ); +}); diff --git a/apps/server/src/persistence/Migrations/044_CleanupCompletedRetractionMessages.ts b/apps/server/src/persistence/Migrations/044_CleanupCompletedRetractionMessages.ts new file mode 100644 index 000000000000..f126bb5c5d19 --- /dev/null +++ b/apps/server/src/persistence/Migrations/044_CleanupCompletedRetractionMessages.ts @@ -0,0 +1,28 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +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' + AND name IN ('projection_thread_messages', 'projection_turn_retractions') + `; + const tableNames = new Set(tables.map((table) => table.name)); + if ( + !tableNames.has("projection_thread_messages") || + !tableNames.has("projection_turn_retractions") + ) { + return; + } + + yield* sql` + DELETE FROM projection_thread_messages + WHERE message_id IN ( + SELECT message_id + FROM projection_turn_retractions + WHERE status = 'completed' + ) + `; +});