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
20 changes: 16 additions & 4 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1709,11 +1709,23 @@ const make = Effect.gen(function* () {
const shouldApplyFallbackCompletionText =
!existingAssistantMessage || existingAssistantMessage.text.length === 0;

const shouldSkipRedundantCompletion =
Option.isNone(activeAssistantMessageId) &&
// Queue-drain race: after turn.completed finalizes a segment under turn A,
// the next provider session can re-emit assistant.complete for the same
// message id with turn B. Re-dispatching rebinds turnId in the projector
// and orphans the final from turn A (Discord loses it under Working).
const alreadyBoundToOtherTurn =
existingAssistantMessage !== undefined &&
existingAssistantMessage.role === "assistant" &&
existingAssistantMessage.turnId !== null &&
turnId !== undefined &&
hasAssistantMessagesForTurn &&
(assistantCompletion.fallbackText?.trim().length ?? 0) === 0;
!sameId(existingAssistantMessage.turnId, turnId);

const shouldSkipRedundantCompletion =
alreadyBoundToOtherTurn ||
(Option.isNone(activeAssistantMessageId) &&
turnId !== undefined &&
hasAssistantMessagesForTurn &&
(assistantCompletion.fallbackText?.trim().length ?? 0) === 0);

if (!shouldSkipRedundantCompletion) {
if (turnId && Option.isNone(activeAssistantMessageId)) {
Expand Down
115 changes: 115 additions & 0 deletions apps/server/src/orchestration/projector.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -503,6 +503,121 @@ describe("orchestration projector", () => {
expect(message?.updatedAt).toBe(completeAt);
});

it("does not rebind an assistant message turnId when a later complete races under the next turn", async () => {
// Production race (t3vm thread 16feaadd…): queue drain re-emits
// assistant.complete for the same segment id under the new turnId with empty
// text. The body must stay, and turnId must stay on the completed turn so
// Discord/web do not swallow the final under the next Working tip.
const createdAt = "2026-07-27T05:54:00.000Z";
const completeAt = "2026-07-27T05:57:21.174Z";
const restampAt = "2026-07-27T05:57:32.206Z";
const model = createEmptyReadModel(createdAt);

// Single runPromise so we stay within the file's LEGACY_BASELINE for
// t3code/no-manual-effect-runtime-in-tests.
const afterRestamp = await Effect.runPromise(
Effect.gen(function* () {
const afterCreate = yield* projectEvent(
model,
makeEvent({
sequence: 1,
type: "thread.created",
aggregateKind: "thread",
aggregateId: "thread-1",
occurredAt: createdAt,
commandId: "cmd-create",
payload: {
threadId: "thread-1",
projectId: "project-1",
title: "demo",
modelSelection: {
provider: ProviderDriverKind.make("codex"),
model: "gpt-5.3-codex",
},
runtimeMode: "full-access",
branch: null,
worktreePath: null,
createdAt,
updatedAt: createdAt,
},
}),
);

const afterFinal = yield* projectEvent(
afterCreate,
makeEvent({
sequence: 2,
type: "thread.message-sent",
aggregateKind: "thread",
aggregateId: "thread-1",
occurredAt: completeAt,
commandId: "cmd-final-delta",
payload: {
threadId: "thread-1",
messageId: "assistant:run:segment:5",
role: "assistant",
text: "**Yes — the bug is almost entirely a naming/dual-use problem.**",
turnId: "turn-prior",
streaming: true,
createdAt: completeAt,
updatedAt: completeAt,
},
}),
);

const afterComplete = yield* projectEvent(
afterFinal,
makeEvent({
sequence: 3,
type: "thread.message-sent",
aggregateKind: "thread",
aggregateId: "thread-1",
occurredAt: completeAt,
commandId: "cmd-final-complete",
payload: {
threadId: "thread-1",
messageId: "assistant:run:segment:5",
role: "assistant",
text: "",
turnId: "turn-prior",
streaming: false,
createdAt: completeAt,
updatedAt: completeAt,
},
}),
);

return yield* projectEvent(
afterComplete,
makeEvent({
sequence: 4,
type: "thread.message-sent",
aggregateKind: "thread",
aggregateId: "thread-1",
occurredAt: restampAt,
commandId: "cmd-restamp-complete",
payload: {
threadId: "thread-1",
messageId: "assistant:run:segment:5",
role: "assistant",
text: "",
turnId: "turn-next",
streaming: false,
createdAt: restampAt,
updatedAt: restampAt,
},
}),
);
}),
);

const message = firstThread(afterRestamp)?.messages[0];
expect(message?.id).toBe("assistant:run:segment:5");
expect(message?.text).toBe("**Yes — the bug is almost entirely a naming/dual-use problem.**");
expect(message?.turnId).toBe("turn-prior");
expect(message?.streaming).toBe(false);
});

it("prunes reverted turn messages from in-memory thread snapshot", async () => {
const createdAt = "2026-02-23T10:00:00.000Z";
const model = createEmptyReadModel(createdAt);
Expand Down
14 changes: 13 additions & 1 deletion apps/server/src/orchestration/projector.ts
Original file line number Diff line number Diff line change
Expand Up @@ -503,7 +503,19 @@ export function projectEvent(
: entry.text,
streaming: message.streaming,
updatedAt: message.updatedAt,
turnId: message.turnId,
// Once an assistant bubble is bound to a turn, never rebind it
// to a different turn. Queue drain races re-emit
// assistant.complete for the same messageId under the next
// activeTurnId (empty text + streaming:false), which would
// otherwise orphan the real final from its completed turn and
// hide it under the next Working tip (Discord/web fold).
turnId:
entry.role === "assistant" &&
entry.turnId !== null &&
message.turnId !== null &&
entry.turnId !== message.turnId
? entry.turnId
: message.turnId,
...(message.attachments !== undefined
? { attachments: message.attachments }
: {}),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ const LEGACY_BASELINE = new Map<string, number>([
["apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts", 70],
["apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts", 31],
["apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts", 2],
["apps/server/src/orchestration/projector.test.ts", 20],
["apps/server/src/orchestration/projector.test.ts", 21],
["apps/server/src/project/Layers/ProjectSetupScriptRunner.test.ts", 4],
["apps/server/src/provider/acp/CursorAcpSupport.test.ts", 1],
["apps/server/src/provider/Layers/ClaudeAdapter.test.ts", 2],
Expand Down
Loading