Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 55 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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-"))(
Expand Down
15 changes: 13 additions & 2 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -217,6 +217,7 @@ function retainProjectionMessagesAfterRevert(
messages: ReadonlyArray<ProjectionThreadMessage>,
turns: ReadonlyArray<ProjectionTurn>,
turnCount: number,
excludedMessageIds: ReadonlySet<string>,
): ReadonlyArray<ProjectionThreadMessage> {
const retainedMessageIds = new Set<string>();
const retainedTurnIds = new Set<string>();
Expand All @@ -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;
Expand All @@ -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)),
)
Expand All @@ -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)),
)
Expand Down Expand Up @@ -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<string>()
: new Set([event.payload.retraction.messageId]);
const keptRows = retainProjectionMessagesAfterRevert(
existingRows,
existingTurns,
event.payload.turnCount,
excludedMessageIds,
);
if (keptRows.length === existingRows.length) {
return;
Expand Down
164 changes: 164 additions & 0 deletions apps/server/src/orchestration/projector.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<OrchestrationEvent> = 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<ReturnType<typeof createEmptyReadModel>>
>(
(statePromise, event) =>
statePromise.then((state) => Effect.runPromise(projectEvent(state, event))),
Promise.resolve(afterCreate),
);

const legitimateMessageEvents: ReadonlyArray<OrchestrationEvent> = 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<ReturnType<typeof createEmptyReadModel>>
>(
(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);
Expand Down
11 changes: 11 additions & 0 deletions apps/server/src/orchestration/projector.ts
Original file line number Diff line number Diff line change
Expand Up @@ -119,9 +119,13 @@ function retainThreadMessagesAfterRevert(
messages: ReadonlyArray<OrchestrationMessage>,
retainedTurnIds: ReadonlySet<string>,
turnCount: number,
excludedMessageIds: ReadonlySet<string>,
): ReadonlyArray<OrchestrationMessage> {
const retainedMessageIds = new Set<string>();
for (const message of messages) {
if (excludedMessageIds.has(message.id)) {
continue;
}
if (message.role === "system") {
retainedMessageIds.add(message.id);
continue;
Expand All @@ -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)),
)
Expand All @@ -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)),
)
Expand Down Expand Up @@ -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<string>()
: new Set([payload.retraction.messageId]);
const messages = retainThreadMessagesAfterRevert(
thread.messages,
retainedTurnIds,
payload.turnCount,
excludedMessageIds,
).slice(-MAX_THREAD_MESSAGES);
const proposedPlans = retainThreadProposedPlansAfterRevert(
thread.proposedPlans,
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/persistence/Migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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);
Expand Down
Loading
Loading