diff --git a/apps/mobile/src/lib/threadActivity.test.ts b/apps/mobile/src/lib/threadActivity.test.ts index 73080b3dfcc6..67e6dd7261e5 100644 --- a/apps/mobile/src/lib/threadActivity.test.ts +++ b/apps/mobile/src/lib/threadActivity.test.ts @@ -1090,4 +1090,26 @@ describe("quiet timeline: nested agents", () => { { type: "activity-group", id: "nested-done" }, ]); }); + + it("filters reasoning messages at thread-feed derivation entry", () => { + const thread = makeThread({ + id: ThreadId.make("thread-reasoning"), + projectId: ProjectId.make("project-1"), + title: "Reasoning", + messages: [ + { + id: MessageId.make("reasoning-1"), + role: "assistant", + channel: "reasoning", + text: "Private reasoning", + turnId: TurnId.make("turn-1"), + streaming: false, + createdAt: "2026-04-01T00:00:01.000Z", + updatedAt: "2026-04-01T00:00:01.000Z", + }, + ], + }); + + expect(buildThreadFeed(thread)).toEqual([]); + }); }); diff --git a/apps/mobile/src/lib/threadActivity.ts b/apps/mobile/src/lib/threadActivity.ts index f30081dd35b4..cc04c615aee4 100644 --- a/apps/mobile/src/lib/threadActivity.ts +++ b/apps/mobile/src/lib/threadActivity.ts @@ -1,5 +1,6 @@ import { ApprovalRequestId, + isReasoningMessage, isToolLifecycleItemType, ProviderApprovalOption, ProviderRequestKind, @@ -1730,12 +1731,13 @@ export function buildThreadFeed( readonly localMessages?: ReadonlyArray; }, ): ThreadFeedEntry[] { - const loadedMessages = options?.loadedMessages ?? thread.messages; + const sourceMessages = options?.loadedMessages ?? thread.messages; + const loadedMessages = sourceMessages.filter((message) => !isReasoningMessage(message)); const messages = options?.localMessages ? [...loadedMessages, ...options.localMessages] : loadedMessages; const oldestLoadedMessageCreatedAt = - options?.loadedMessages !== undefined ? (loadedMessages[0]?.createdAt ?? null) : null; + options?.loadedMessages !== undefined ? (sourceMessages[0]?.createdAt ?? null) : null; const workLogEntries = deriveWorkLogEntries(thread.activities); const entries = Arr.sortWith( [ diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts index ca4cb7afd9ab..b2261c530c09 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts @@ -28,7 +28,7 @@ import * as ManagedRuntime from "effect/ManagedRuntime"; import * as PubSub from "effect/PubSub"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; -import { afterEach, describe, expect, it, vi } from "vite-plus/test"; +import { afterEach, describe, expect, it, vi } from "@effect/vitest"; import * as CheckpointStore from "../../checkpointing/CheckpointStore.ts"; import * as VcsDriverRegistry from "../../vcs/VcsDriverRegistry.ts"; @@ -490,6 +490,28 @@ describe("CheckpointReactor", () => { checkpointRefForThreadTurn(ThreadId.make("thread-1"), 0), ); + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.message.assistant.complete", + commandId: CommandId.make("cmd-checkpoint-answer-complete"), + threadId: ThreadId.make("thread-1"), + messageId: MessageId.make("answer-1"), + turnId: asTurnId("turn-1"), + createdAt: "2026-01-01T00:00:00.200Z", + }), + ); + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.message.assistant.complete", + commandId: CommandId.make("cmd-checkpoint-reasoning-complete"), + threadId: ThreadId.make("thread-1"), + messageId: MessageId.make("reasoning-1"), + channel: "reasoning", + turnId: asTurnId("turn-1"), + createdAt: "2026-01-01T00:00:00.400Z", + }), + ); + NodeFS.writeFileSync(NodePath.join(harness.cwd, "README.md"), "v2\n", "utf8"); harness.provider.emit({ type: "turn.completed", @@ -502,7 +524,14 @@ describe("CheckpointReactor", () => { payload: { state: "completed" }, }); - await waitForEvent(harness.engine, (event) => event.type === "thread.turn-diff-completed"); + const events = await waitForEvent( + harness.engine, + (event) => event.type === "thread.turn-diff-completed", + ); + const checkpointEvent = events.find((event) => event.type === "thread.turn-diff-completed"); + if (checkpointEvent?.type === "thread.turn-diff-completed") { + expect(checkpointEvent.payload.assistantMessageId).toBe(MessageId.make("answer-1")); + } const thread = await waitForThread( harness.readModel, (entry) => entry.latestTurn?.turnId === "turn-1" && entry.checkpoints.length === 1, @@ -1152,12 +1181,12 @@ describe("CheckpointReactor", () => { }); }); - it("processes consecutive revert requests with deterministic rollback sequencing", async () => { - const harness = await createHarness(); - const createdAt = "2026-01-01T00:00:00.000Z"; + it.effect("processes consecutive revert requests with deterministic rollback sequencing", () => + Effect.gen(function* () { + const harness = yield* Effect.promise(() => createHarness()); + const createdAt = "2026-01-01T00:00:00.000Z"; - await Effect.runPromise( - harness.engine.dispatch({ + yield* harness.engine.dispatch({ type: "thread.session.set", commandId: CommandId.make("cmd-session-set-inline-revert"), threadId: ThreadId.make("thread-1"), @@ -1171,11 +1200,9 @@ describe("CheckpointReactor", () => { updatedAt: createdAt, }, createdAt, - }), - ); + }); - await Effect.runPromise( - harness.engine.dispatch({ + yield* harness.engine.dispatch({ type: "thread.turn.diff.complete", commandId: CommandId.make("cmd-inline-revert-diff-1"), threadId: ThreadId.make("thread-1"), @@ -1186,10 +1213,8 @@ describe("CheckpointReactor", () => { files: [], checkpointTurnCount: 1, createdAt, - }), - ); - await Effect.runPromise( - harness.engine.dispatch({ + }); + yield* harness.engine.dispatch({ type: "thread.turn.diff.complete", commandId: CommandId.make("cmd-inline-revert-diff-2"), threadId: ThreadId.make("thread-1"), @@ -1200,62 +1225,60 @@ describe("CheckpointReactor", () => { files: [], checkpointTurnCount: 2, createdAt, - }), - ); + }); - await Effect.runPromise( - harness.engine.dispatch({ + yield* harness.engine.dispatch({ type: "thread.checkpoint.revert", commandId: CommandId.make("cmd-sequenced-revert-request-1"), threadId: ThreadId.make("thread-1"), turnCount: 1, createdAt, - }), - ); - await Effect.runPromise( - harness.engine.dispatch({ + }); + yield* harness.engine.dispatch({ type: "thread.checkpoint.revert", commandId: CommandId.make("cmd-sequenced-revert-request-0"), threadId: ThreadId.make("thread-1"), turnCount: 0, createdAt, - }), - ); + }); - await harness.drain(); + yield* Effect.promise(() => harness.drain()); - expect(harness.provider.rollbackConversation).toHaveBeenCalledTimes(2); - expect(harness.provider.rollbackConversation.mock.calls[0]?.[0]).toEqual({ - threadId: ThreadId.make("thread-1"), - numTurns: 1, - }); - expect(harness.provider.rollbackConversation.mock.calls[1]?.[0]).toEqual({ - threadId: ThreadId.make("thread-1"), - numTurns: 1, - }); - }); + expect(harness.provider.rollbackConversation).toHaveBeenCalledTimes(2); + expect(harness.provider.rollbackConversation.mock.calls[0]?.[0]).toEqual({ + threadId: ThreadId.make("thread-1"), + numTurns: 1, + }); + expect(harness.provider.rollbackConversation.mock.calls[1]?.[0]).toEqual({ + threadId: ThreadId.make("thread-1"), + numTurns: 1, + }); + }), + ); - it("appends an error activity when revert is requested without an active session", async () => { - const harness = await createHarness({ hasSession: false }); - const createdAt = "2026-01-01T00:00:00.000Z"; + it.effect("appends an error activity when revert is requested without an active session", () => + Effect.gen(function* () { + const harness = yield* Effect.promise(() => createHarness({ hasSession: false })); + const createdAt = "2026-01-01T00:00:00.000Z"; - await Effect.runPromise( - harness.engine.dispatch({ + yield* harness.engine.dispatch({ type: "thread.checkpoint.revert", commandId: CommandId.make("cmd-revert-no-session"), threadId: ThreadId.make("thread-1"), turnCount: 1, createdAt, - }), - ); + }); - const thread = await waitForThread(harness.readModel, (entry) => - entry.activities.some((activity) => activity.kind === "checkpoint.revert.failed"), - ); + const thread = yield* Effect.promise(() => + waitForThread(harness.readModel, (entry) => + entry.activities.some((activity) => activity.kind === "checkpoint.revert.failed"), + ), + ); - expect(thread.activities.some((activity) => activity.kind === "checkpoint.revert.failed")).toBe( - true, - ); - expect(harness.provider.rollbackConversation).not.toHaveBeenCalled(); - }); + expect( + thread.activities.some((activity) => activity.kind === "checkpoint.revert.failed"), + ).toBe(true); + expect(harness.provider.rollbackConversation).not.toHaveBeenCalled(); + }), + ); }); diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.ts index 95adee0cf7f8..b1288ac16bc8 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.ts @@ -3,10 +3,12 @@ import { type CheckpointRef, EventId, MessageId, + isReasoningMessage, type ProjectId, ThreadId, TurnId, type OrchestrationEvent, + type OrchestrationMessageChannel, type ProviderRuntimeEvent, type VcsStatusLocalResult, } from "@t3tools/contracts"; @@ -225,6 +227,7 @@ const make = Effect.gen(function* () { readonly messages: ReadonlyArray<{ readonly id: MessageId; readonly role: string; + readonly channel?: OrchestrationMessageChannel | undefined; readonly turnId: TurnId | null; }>; }; @@ -298,7 +301,12 @@ const make = Effect.gen(function* () { input.assistantMessageId ?? input.thread.messages .toReversed() - .find((entry) => entry.role === "assistant" && entry.turnId === input.turnId)?.id ?? + .find( + (entry) => + entry.role === "assistant" && + !isReasoningMessage(entry) && + entry.turnId === input.turnId, + )?.id ?? MessageId.make(`assistant:${input.turnId}`); yield* orchestrationEngine.dispatch({ diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index 57215803f4b9..05f842b89f83 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -125,8 +125,9 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { threadId: ThreadId.make("thread-1"), messageId: MessageId.make("message-1"), role: "assistant", + channel: "reasoning", text: "hello", - turnId: null, + turnId: TurnId.make("turn-1"), streaming: false, createdAt: now, updatedAt: now, @@ -153,13 +154,21 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { const messageRows = yield* sql<{ readonly messageId: string; readonly text: string; + readonly channel: string | null; }>` SELECT message_id AS "messageId", - text + text, + channel FROM projection_thread_messages `; - assert.deepEqual(messageRows, [{ messageId: "message-1", text: "hello" }]); + assert.deepEqual(messageRows, [ + { messageId: "message-1", text: "hello", channel: "reasoning" }, + ]); + const turnRows = yield* sql<{ readonly assistantMessageId: string | null }>` + SELECT assistant_message_id AS "assistantMessageId" FROM projection_turns + `; + assert.deepEqual(turnRows, []); const stateRows = yield* sql<{ readonly projector: string; @@ -2551,7 +2560,7 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { }), ); - it.effect("does not fallback-retain messages whose turnId is removed by revert", () => + it.effect("removes reasoning messages whose turnId is removed by revert", () => Effect.gen(function* () { const projectionPipeline = yield* OrchestrationProjectionPipeline; const eventStore = yield* OrchestrationEventStore; @@ -2710,6 +2719,7 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { threadId: ThreadId.make("thread-revert"), messageId: MessageId.make("assistant-remove"), role: "assistant", + channel: "reasoning", text: "removed", turnId: TurnId.make("turn-2"), streaming: false, @@ -2758,6 +2768,485 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { ); }); +it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-projection-revert-reasoning-")))( + "OrchestrationProjectionPipeline revert reasoning retention", + (it) => { + // Regression: retainProjectionMessagesAfterRevert counts a kept turn's + // channel=reasoning row as the turn's assistant, so the turnless + // substantive assistant is never recovered through the fallback scan and + // is dropped by the revert. + it.effect( + "reasoning on a kept turn does not satisfy the substantive assistant count on revert", + () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const appendAndProject = (event: Parameters[0]) => + eventStore + .append(event) + .pipe(Effect.flatMap((savedEvent) => projectionPipeline.projectEvent(savedEvent))); + + yield* appendAndProject({ + type: "project.created", + eventId: EventId.make("evt-rrcount-1"), + aggregateKind: "project", + aggregateId: ProjectId.make("project-reason-count"), + occurredAt: "2026-02-27T12:00:00.000Z", + commandId: CommandId.make("cmd-rrcount-1"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-1"), + metadata: {}, + payload: { + projectId: ProjectId.make("project-reason-count"), + title: "Project Reason Count", + workspaceRoot: "/tmp/project-reason-count", + defaultModelSelection: null, + scripts: [], + createdAt: "2026-02-27T12:00:00.000Z", + updatedAt: "2026-02-27T12:00:00.000Z", + }, + }); + + yield* appendAndProject({ + type: "thread.created", + eventId: EventId.make("evt-rrcount-2"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:01.000Z", + commandId: CommandId.make("cmd-rrcount-2"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-2"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + projectId: ProjectId.make("project-reason-count"), + title: "Thread Reason Count", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt: "2026-02-27T12:00:01.000Z", + updatedAt: "2026-02-27T12:00:01.000Z", + }, + }); + + yield* appendAndProject({ + type: "thread.message-sent", + eventId: EventId.make("evt-rrcount-3"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:01.500Z", + commandId: CommandId.make("cmd-rrcount-3"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-3"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + messageId: MessageId.make("user-count-1"), + role: "user", + text: "first prompt", + turnId: TurnId.make("turn-count-1"), + streaming: false, + createdAt: "2026-02-27T12:00:01.500Z", + updatedAt: "2026-02-27T12:00:01.500Z", + }, + }); + + yield* appendAndProject({ + type: "thread.message-sent", + eventId: EventId.make("evt-rrcount-4"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:01.600Z", + commandId: CommandId.make("cmd-rrcount-4"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-4"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + messageId: MessageId.make("reasoning-count-1"), + role: "assistant", + channel: "reasoning", + text: "thinking about the first prompt", + turnId: TurnId.make("turn-count-1"), + streaming: false, + createdAt: "2026-02-27T12:00:01.600Z", + updatedAt: "2026-02-27T12:00:01.600Z", + }, + }); + + yield* appendAndProject({ + type: "thread.message-sent", + eventId: EventId.make("evt-rrcount-5"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:01.700Z", + commandId: CommandId.make("cmd-rrcount-5"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-5"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + messageId: MessageId.make("assistant-count-1"), + role: "assistant", + text: "substantive first answer", + turnId: null, + streaming: false, + createdAt: "2026-02-27T12:00:01.700Z", + updatedAt: "2026-02-27T12:00:01.700Z", + }, + }); + + yield* appendAndProject({ + type: "thread.turn-diff-completed", + eventId: EventId.make("evt-rrcount-6"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:02.000Z", + commandId: CommandId.make("cmd-rrcount-6"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-6"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + turnId: TurnId.make("turn-count-1"), + checkpointTurnCount: 1, + checkpointRef: CheckpointRef.make("refs/t3/checkpoints/thread-reason-count/turn/1"), + status: "ready", + files: [], + assistantMessageId: null, + completedAt: "2026-02-27T12:00:02.000Z", + }, + }); + + yield* appendAndProject({ + type: "thread.message-sent", + eventId: EventId.make("evt-rrcount-7"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:02.500Z", + commandId: CommandId.make("cmd-rrcount-7"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-7"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + messageId: MessageId.make("user-count-2"), + role: "user", + text: "second prompt", + turnId: TurnId.make("turn-count-2"), + streaming: false, + createdAt: "2026-02-27T12:00:02.500Z", + updatedAt: "2026-02-27T12:00:02.500Z", + }, + }); + + yield* appendAndProject({ + type: "thread.message-sent", + eventId: EventId.make("evt-rrcount-8"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:02.600Z", + commandId: CommandId.make("cmd-rrcount-8"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-8"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + messageId: MessageId.make("assistant-count-2"), + role: "assistant", + text: "second answer", + turnId: TurnId.make("turn-count-2"), + streaming: false, + createdAt: "2026-02-27T12:00:02.600Z", + updatedAt: "2026-02-27T12:00:02.600Z", + }, + }); + + yield* appendAndProject({ + type: "thread.turn-diff-completed", + eventId: EventId.make("evt-rrcount-9"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:03.000Z", + commandId: CommandId.make("cmd-rrcount-9"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-9"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + turnId: TurnId.make("turn-count-2"), + checkpointTurnCount: 2, + checkpointRef: CheckpointRef.make("refs/t3/checkpoints/thread-reason-count/turn/2"), + status: "ready", + files: [], + assistantMessageId: MessageId.make("assistant-count-2"), + completedAt: "2026-02-27T12:00:03.000Z", + }, + }); + + yield* appendAndProject({ + type: "thread.reverted", + eventId: EventId.make("evt-rrcount-10"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-count"), + occurredAt: "2026-02-27T12:00:04.000Z", + commandId: CommandId.make("cmd-rrcount-10"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrcount-10"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-count"), + turnCount: 1, + }, + }); + + const messageRows = yield* sql<{ + readonly messageId: string; + readonly turnId: string | null; + readonly role: string; + readonly channel: string | null; + }>` + SELECT + message_id AS "messageId", + turn_id AS "turnId", + role, + channel + FROM projection_thread_messages + WHERE thread_id = 'thread-reason-count' + ORDER BY created_at ASC, message_id ASC + `; + assert.deepEqual(messageRows, [ + { + messageId: "user-count-1", + turnId: "turn-count-1", + role: "user", + channel: null, + }, + { + messageId: "reasoning-count-1", + turnId: "turn-count-1", + role: "assistant", + channel: "reasoning", + }, + { + messageId: "assistant-count-1", + turnId: null, + role: "assistant", + channel: null, + }, + ]); + }), + ); + + // Regression: the fallback assistant scan treats every role=assistant row + // as eligible, so a turnless channel=reasoning row left over from + // discarded later work is selected as the kept turn's assistant and + // survives the revert while the substantive assistant is dropped. + it.effect( + "turnless reasoning from discarded work is never selected as a fallback assistant on revert", + () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const appendAndProject = (event: Parameters[0]) => + eventStore + .append(event) + .pipe(Effect.flatMap((savedEvent) => projectionPipeline.projectEvent(savedEvent))); + + yield* appendAndProject({ + type: "project.created", + eventId: EventId.make("evt-rrfall-1"), + aggregateKind: "project", + aggregateId: ProjectId.make("project-reason-fallback"), + occurredAt: "2026-02-27T12:10:00.000Z", + commandId: CommandId.make("cmd-rrfall-1"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrfall-1"), + metadata: {}, + payload: { + projectId: ProjectId.make("project-reason-fallback"), + title: "Project Reason Fallback", + workspaceRoot: "/tmp/project-reason-fallback", + defaultModelSelection: null, + scripts: [], + createdAt: "2026-02-27T12:10:00.000Z", + updatedAt: "2026-02-27T12:10:00.000Z", + }, + }); + + yield* appendAndProject({ + type: "thread.created", + eventId: EventId.make("evt-rrfall-2"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-fallback"), + occurredAt: "2026-02-27T12:10:01.000Z", + commandId: CommandId.make("cmd-rrfall-2"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrfall-2"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-fallback"), + projectId: ProjectId.make("project-reason-fallback"), + title: "Thread Reason Fallback", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt: "2026-02-27T12:10:01.000Z", + updatedAt: "2026-02-27T12:10:01.000Z", + }, + }); + + yield* appendAndProject({ + type: "thread.message-sent", + eventId: EventId.make("evt-rrfall-3"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-fallback"), + occurredAt: "2026-02-27T12:10:01.500Z", + commandId: CommandId.make("cmd-rrfall-3"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrfall-3"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-fallback"), + messageId: MessageId.make("user-fallback-1"), + role: "user", + text: "kept prompt", + turnId: TurnId.make("turn-fallback-1"), + streaming: false, + createdAt: "2026-02-27T12:10:01.500Z", + updatedAt: "2026-02-27T12:10:01.500Z", + }, + }); + + yield* appendAndProject({ + type: "thread.turn-diff-completed", + eventId: EventId.make("evt-rrfall-4"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-fallback"), + occurredAt: "2026-02-27T12:10:02.000Z", + commandId: CommandId.make("cmd-rrfall-4"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrfall-4"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-fallback"), + turnId: TurnId.make("turn-fallback-1"), + checkpointTurnCount: 1, + checkpointRef: CheckpointRef.make( + "refs/t3/checkpoints/thread-reason-fallback/turn/1", + ), + status: "ready", + files: [], + assistantMessageId: null, + completedAt: "2026-02-27T12:10:02.000Z", + }, + }); + + yield* appendAndProject({ + type: "thread.message-sent", + eventId: EventId.make("evt-rrfall-5"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-fallback"), + occurredAt: "2026-02-27T12:10:02.100Z", + commandId: CommandId.make("cmd-rrfall-5"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrfall-5"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-fallback"), + messageId: MessageId.make("reasoning-discarded"), + role: "assistant", + channel: "reasoning", + text: "planning work that gets reverted away", + turnId: null, + streaming: false, + createdAt: "2026-02-27T12:10:02.100Z", + updatedAt: "2026-02-27T12:10:02.100Z", + }, + }); + + yield* appendAndProject({ + type: "thread.message-sent", + eventId: EventId.make("evt-rrfall-6"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-fallback"), + occurredAt: "2026-02-27T12:10:02.200Z", + commandId: CommandId.make("cmd-rrfall-6"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrfall-6"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-fallback"), + messageId: MessageId.make("assistant-recovered"), + role: "assistant", + text: "substantive answer for the kept turn", + turnId: null, + streaming: false, + createdAt: "2026-02-27T12:10:02.200Z", + updatedAt: "2026-02-27T12:10:02.200Z", + }, + }); + + yield* appendAndProject({ + type: "thread.reverted", + eventId: EventId.make("evt-rrfall-7"), + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-reason-fallback"), + occurredAt: "2026-02-27T12:10:03.000Z", + commandId: CommandId.make("cmd-rrfall-7"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-rrfall-7"), + metadata: {}, + payload: { + threadId: ThreadId.make("thread-reason-fallback"), + turnCount: 1, + }, + }); + + const messageRows = yield* sql<{ + readonly messageId: string; + readonly turnId: string | null; + readonly role: string; + readonly channel: string | null; + }>` + SELECT + message_id AS "messageId", + turn_id AS "turnId", + role, + channel + FROM projection_thread_messages + WHERE thread_id = 'thread-reason-fallback' + ORDER BY created_at ASC, message_id ASC + `; + assert.deepEqual(messageRows, [ + { + messageId: "user-fallback-1", + turnId: "turn-fallback-1", + role: "user", + channel: null, + }, + { + messageId: "assistant-recovered", + turnId: null, + role: "assistant", + channel: null, + }, + ]); + }), + ); + }, +); + it.layer(makeProjectionPipelinePrefixedTestLayer("t3-pending-turn-terminal-test-"))( "OrchestrationProjectionPipeline pending turn cleanup", (it) => { diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index ac514d7eb282..cd89d357fe00 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -3,6 +3,7 @@ import { type ChatAttachment, type OrchestrationEvent, type OrchestrationSessionStatus, + isReasoningMessage, ThreadId, } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; @@ -285,7 +286,10 @@ function retainProjectionMessagesAfterRevert( } const retainedAssistantCount = messages.filter( - (message) => message.role === "assistant" && retainedMessageIds.has(message.messageId), + (message) => + message.role === "assistant" && + message.channel !== "reasoning" && + retainedMessageIds.has(message.messageId), ).length; const missingAssistantCount = Math.max(0, turnCount - retainedAssistantCount); if (missingAssistantCount > 0) { @@ -293,6 +297,7 @@ function retainProjectionMessagesAfterRevert( .filter( (message) => message.role === "assistant" && + message.channel !== "reasoning" && !retainedMessageIds.has(message.messageId) && (message.turnId === null || retainedTurnIds.has(message.turnId)), ) @@ -1035,6 +1040,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti threadId: event.payload.threadId, turnId: event.payload.turnId, role: event.payload.role, + ...(event.payload.channel !== undefined ? { channel: event.payload.channel } : {}), text: nextText, ...(nextAttachments !== undefined ? { attachments: [...nextAttachments] } : {}), isStreaming: event.payload.streaming, @@ -1380,7 +1386,11 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti } case "thread.message-sent": { - if (event.payload.turnId === null || event.payload.role !== "assistant") { + if ( + event.payload.turnId === null || + event.payload.role !== "assistant" || + isReasoningMessage(event.payload) + ) { return; } // A completed assistant message only settles the turn once the diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts index 30892c760e77..0733e0a55344 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts @@ -123,21 +123,35 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { thread_id, turn_id, role, + channel, text, is_streaming, created_at, updated_at ) - VALUES ( - 'message-1', - 'thread-1', - 'turn-1', - 'assistant', - 'hello from projection', - 0, - '2026-02-24T00:00:04.000Z', - '2026-02-24T00:00:05.000Z' - ) + VALUES + ( + 'reasoning-1', + 'thread-1', + 'turn-1', + 'assistant', + 'reasoning', + 'thinking before answer', + 0, + '2026-02-24T00:00:03.500Z', + '2026-02-24T00:00:03.600Z' + ), + ( + 'message-1', + 'thread-1', + 'turn-1', + 'assistant', + NULL, + 'hello from projection', + 0, + '2026-02-24T00:00:04.000Z', + '2026-02-24T00:00:05.000Z' + ) `; yield* sql` @@ -337,6 +351,16 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { titleRegeneration: null, deletedAt: null, messages: [ + { + id: asMessageId("reasoning-1"), + role: "assistant", + channel: "reasoning", + text: "thinking before answer", + turnId: asTurnId("turn-1"), + streaming: false, + createdAt: "2026-02-24T00:00:03.500Z", + updatedAt: "2026-02-24T00:00:03.600Z", + }, { id: asMessageId("message-1"), role: "assistant", diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index 0b9698eaf9cf..94080806e0c7 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -5,6 +5,7 @@ import { MessageId, NonNegativeInt, OrchestrationCheckpointFile, + OrchestrationMessageChannel, OrchestrationProposedPlanId, OrchestrationReadModel, OrchestrationThreadSearchSource, @@ -83,6 +84,7 @@ const ProjectionProjectDbRowSchema = ProjectionProject.mapFields( const ProjectionThreadMessageDbRowSchema = ProjectionThreadMessage.mapFields( Struct.assign({ isStreaming: Schema.Number, + channel: Schema.NullOr(OrchestrationMessageChannel), attachments: Schema.NullOr(Schema.fromJsonString(Schema.Array(ChatAttachment))), }), ); @@ -538,6 +540,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { thread_id AS "threadId", turn_id AS "turnId", role, + channel, text, attachments_json AS "attachments", is_streaming AS "isStreaming", @@ -983,6 +986,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { thread_id AS "threadId", turn_id AS "turnId", role, + channel, text, attachments_json AS "attachments", is_streaming AS "isStreaming", @@ -1226,6 +1230,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { thread_id AS "threadId", turn_id AS "turnId", role, + channel, text, attachments_json AS "attachments", is_streaming AS "isStreaming", @@ -1568,6 +1573,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { threadMessages.push({ id: row.messageId, role: row.role, + ...(row.channel !== null ? { channel: row.channel } : {}), text: row.text, ...(row.attachments !== null ? { attachments: row.attachments } : {}), turnId: row.turnId, @@ -2652,6 +2658,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { const message = { id: row.messageId, role: row.role, + ...(row.channel !== null ? { channel: row.channel } : {}), text: row.text, turnId: row.turnId, streaming: row.isStreaming === 1, diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 4ad35e54cca8..655b3b21f1c0 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -773,6 +773,17 @@ describe("ProviderCommandReactor", () => { createdAt: now, }), ); + await harness.runEffect( + harness.engine.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-reasoning-before-title-regeneration"), + threadId: ThreadId.make("thread-1"), + messageId: asMessageId("reasoning-before-title-regeneration"), + delta: "Internal reasoning must not appear in title context.", + channel: "reasoning", + createdAt: "2026-01-01T00:00:00.500Z", + }), + ); await harness.runEffect( harness.engine.dispatch({ type: "thread.message.assistant.delta", diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 22c70094ce0e..46bf9e3bb04d 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -1,5 +1,7 @@ import { type ChatAttachment, + type OrchestrationMessageChannel, + isReasoningMessage, CommandId, EventId, type ModelSelection, @@ -102,12 +104,13 @@ const FIRST_USER_CONTEXT_TRUNCATION_MARKER = "\n[First user message truncated]"; type ThreadTitleMessage = { readonly role: "user" | "assistant" | "system"; + readonly channel?: OrchestrationMessageChannel | undefined; readonly text: string; readonly attachments?: ReadonlyArray | undefined; }; function formatThreadTitleSection(message: ThreadTitleMessage): string | undefined { - if (message.role === "system") { + if (message.role === "system" || isReasoningMessage(message)) { return undefined; } const text = message.text.trim(); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 84858b6affe9..904b79df3c03 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -20,6 +20,7 @@ import { ProjectId, ProviderItemId, type ServerSettings, + type ServerSettingsPatch, ThreadId, TurnId, } from "@t3tools/contracts"; @@ -197,7 +198,10 @@ async function waitForThread( describe("ProviderRuntimeIngestion", () => { let runtime: ManagedRuntime.ManagedRuntime< - OrchestrationEngineService | ProviderRuntimeIngestionService | ProjectionSnapshotQuery, + | OrchestrationEngineService + | ProviderRuntimeIngestionService + | ProjectionSnapshotQuery + | ServerSettingsService, unknown > | null = null; let scope: Scope.Closeable | null = null; @@ -242,7 +246,8 @@ describe("ProviderRuntimeIngestion", () => { Layer.provide(RepositoryIdentityResolver.layer), Layer.provide(SqlitePersistenceMemory), ); - const layer = ProviderRuntimeIngestionLive.pipe( + const serverSettingsLayer = makeTestServerSettingsLayer(options?.serverSettings); + const ingestionLayer = ProviderRuntimeIngestionLive.pipe( Layer.provideMerge(orchestrationLayer), Layer.provideMerge(projectionSnapshotLayer), // Single shared liveness instance across ingestion (writer), the @@ -251,14 +256,16 @@ describe("ProviderRuntimeIngestion", () => { Layer.provideMerge(ThreadPlanProgress.layer), Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(Layer.succeed(ProviderService, provider.service)), - Layer.provideMerge(makeTestServerSettingsLayer(options?.serverSettings)), + Layer.provideMerge(serverSettingsLayer), Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(NodeServices.layer), ); + const layer = Layer.merge(ingestionLayer, serverSettingsLayer); runtime = ManagedRuntime.make(layer); const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery)); const ingestion = await runtime.runPromise(Effect.service(ProviderRuntimeIngestionService)); + const serverSettings = await runtime.runPromise(Effect.service(ServerSettingsService)); scope = await Effect.runPromise(Scope.make("sequential")); await Effect.runPromise(ingestion.start().pipe(Scope.provide(scope))); const drain = () => Effect.runPromise(ingestion.drain); @@ -323,6 +330,8 @@ describe("ProviderRuntimeIngestion", () => { readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), emit: provider.emit, setProviderSession: provider.setSession, + updateServerSettings: (patch: ServerSettingsPatch) => + runtime!.runPromise(serverSettings.updateSettings(patch)), drain, }; } @@ -3624,4 +3633,870 @@ describe("ProviderRuntimeIngestion", () => { expect(thread.session?.status).toBe("error"); expect(thread.session?.lastError).toBe("runtime still processed"); }); + + it("segments buffered reasoning around answers with monotonic message stamps", async () => { + const harness = await createHarness(); + const turnId = asTurnId("turn-reasoning-buffered"); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-summary"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_summary_text", delta: "Inspecting files" }, + }); + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-answer-after-reasoning"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + itemId: asItemId("answer-after-reasoning"), + payload: { streamKind: "assistant_text", delta: "Answer" }, + }); + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-answer-after-reasoning-complete"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + itemId: asItemId("answer-after-reasoning"), + payload: { itemType: "assistant_message", status: "completed" }, + }); + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-second-segment"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:02.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "Checking results" }, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-reasoning-aborted"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:03.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { reason: "cancelled" }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => { + const reasoning = entry.messages.filter((message) => message.channel === "reasoning"); + return reasoning.length === 2 && reasoning.every((message) => !message.streaming); + }); + expect(thread.messages.map((message) => [message.channel, message.text])).toEqual([ + ["reasoning", "Inspecting files"], + [undefined, "Answer"], + ["reasoning", "Checking results"], + ]); + + const events = await Effect.runPromise( + Stream.runCollect(harness.engine.readEvents(0)).pipe( + Effect.map((chunk) => Array.from(chunk)), + ), + ); + const messageEvents = events.filter( + (event): event is Extract<(typeof events)[number], { type: "thread.message-sent" }> => + event.type === "thread.message-sent" && event.payload.threadId === "thread-1", + ); + const messageTimes = messageEvents.map((event) => Date.parse(event.payload.createdAt)); + expect( + messageTimes.every((time, index) => index === 0 || time > messageTimes[index - 1]!), + ).toBe(true); + }); + + it("preserves buffered reasoning duration through finalization", async () => { + const harness = await createHarness(); + const turnId = asTurnId("turn-reasoning-duration"); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-duration-start"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "Thinking for a while" }, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-reasoning-duration-end"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:03.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { reason: "cancelled" }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const reasoning = readModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.find((message) => message.channel === "reasoning"); + + expect(reasoning?.createdAt).toBe("2026-05-01T00:00:00.000Z"); + expect(reasoning?.updatedAt).toBe("2026-05-01T00:00:03.000Z"); + }); + + it("finalizes turnless reasoning when a named turn completes", async () => { + const harness = await createHarness(); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-turnless"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + payload: { streamKind: "reasoning_text", delta: "Reasoning before turn binding" }, + }); + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-reasoning-turnless-complete"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:01.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-reasoning-late-bound"), + status: "completed", + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const reasoning = readModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.filter((message) => message.channel === "reasoning"); + + expect(reasoning).toEqual([ + expect.objectContaining({ + id: "reasoning:thread-1:turnless:segment:0", + streaming: false, + text: "Reasoning before turn binding", + }), + ]); + }); + + it("continues reasoning segment numbering after session cleanup", async () => { + const harness = await createHarness(); + const turnId = asTurnId("turn-reasoning-resumed"); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-before-cleanup"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "Before cleanup" }, + }); + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-reasoning-session-cleanup"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:01.000Z", + threadId: asThreadId("thread-1"), + payload: {}, + }); + await harness.drain(); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-after-cleanup"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:02.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "After cleanup" }, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-reasoning-after-cleanup-complete"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:03.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { reason: "cancelled" }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const reasoning = readModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.filter((message) => message.channel === "reasoning"); + + expect(reasoning?.map((message) => [message.id, message.text])).toEqual([ + ["reasoning:thread-1:turn-reasoning-resumed:segment:0", "Before cleanup"], + ["reasoning:thread-1:turn-reasoning-resumed:segment:1", "After cleanup"], + ]); + }); + + it("preserves reasoning chunk order when delivery mode changes mid-segment", async () => { + const harness = await createHarness(); + const turnId = asTurnId("turn-reasoning-mode-change"); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-buffered-before-mode-change"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "A" }, + }); + await harness.drain(); + await harness.updateServerSettings({ enableLegacyTokenStreaming: true }); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-streaming-after-mode-change"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:01.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "B" }, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-reasoning-mode-change-complete"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:02.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { reason: "cancelled" }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const reasoning = readModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.find((message) => message.channel === "reasoning"); + + expect(reasoning?.text).toBe("AB"); + expect(reasoning?.streaming).toBe(false); + }); + + it("keeps active-turn reasoning open after a stale runtime error", async () => { + const harness = await createHarness(); + const activeTurnId = asTurnId("turn-reasoning-active"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-reasoning-active-started"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + }); + await harness.drain(); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-before-stale-error"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:01.000Z", + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + payload: { streamKind: "reasoning_text", delta: "Before " }, + }); + await harness.drain(); + + harness.emit({ + type: "runtime.error", + eventId: asEventId("evt-reasoning-stale-error"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:02.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-reasoning-superseded"), + payload: { message: "stale turn failed" }, + }); + await harness.drain(); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-after-stale-error"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:03.000Z", + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + payload: { streamKind: "reasoning_text", delta: "after" }, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-reasoning-active-aborted"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:04.000Z", + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + payload: { reason: "cancelled" }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const reasoning = readModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.filter((message) => message.channel === "reasoning"); + + expect(reasoning?.map((message) => message.text)).toEqual(["Before after"]); + expect(reasoning?.every((message) => !message.streaming)).toBe(true); + }); + + it("keeps active-turn reasoning open after a rejected conflicting turn start", async () => { + const harness = await createHarness(); + const activeTurnId = asTurnId("turn-reasoning-active-start"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-reasoning-current-turn-started"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + }); + await harness.drain(); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-before-stale-start"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:01.000Z", + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + payload: { streamKind: "reasoning_text", delta: "Before " }, + }); + await harness.drain(); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-reasoning-conflicting-turn-started"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:02.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-reasoning-stale-start"), + }); + await harness.drain(); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-after-stale-start"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:03.000Z", + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + payload: { streamKind: "reasoning_text", delta: "after" }, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-reasoning-current-turn-aborted"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-01T00:00:04.000Z", + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + payload: { reason: "cancelled" }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const reasoning = readModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.filter((message) => message.channel === "reasoning"); + + expect(reasoning?.map((message) => message.text)).toEqual(["Before after"]); + expect(reasoning?.every((message) => !message.streaming)).toBe(true); + }); + + it("streams both reasoning delta kinds live and finalizes them on session exit", async () => { + const harness = await createHarness({ serverSettings: { enableLegacyTokenStreaming: true } }); + const turnId = asTurnId("turn-reasoning-streaming"); + for (const [index, streamKind, delta] of [ + [1, "reasoning_summary_text", "Plan "], + [2, "reasoning_text", "details"], + ] as const) { + harness.emit({ + type: "content.delta", + eventId: asEventId(`evt-reasoning-live-${index}`), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-02T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind, delta }, + }); + } + + await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message) => + message.channel === "reasoning" && message.streaming && message.text === "Plan details", + ), + ); + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-reasoning-session-exit"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-02T00:00:01.000Z", + threadId: asThreadId("thread-1"), + payload: {}, + }); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message) => + message.channel === "reasoning" && !message.streaming && message.text === "Plan details", + ), + ); + expect(thread.messages.filter((message) => message.channel === "reasoning")).toHaveLength(1); + }); + + it("finalizes reasoning before a different turn, superseding turn, and visible activity", async () => { + const harness = await createHarness(); + const emitReasoning = (turnId: TurnId, eventId: string, createdAt: string, delta: string) => + harness.emit({ + type: "content.delta", + eventId: asEventId(eventId), + provider: ProviderDriverKind.make("opencode"), + createdAt, + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta }, + }); + + emitReasoning( + asTurnId("turn-reasoning-old"), + "evt-reasoning-old", + "2026-05-03T00:00:00.000Z", + "old turn", + ); + emitReasoning( + asTurnId("turn-reasoning-next"), + "evt-reasoning-next", + "2026-05-03T00:00:01.000Z", + "next turn", + ); + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-reasoning-superseded"), + provider: ProviderDriverKind.make("opencode"), + createdAt: "2026-05-03T00:00:02.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-reasoning-final"), + }); + emitReasoning( + asTurnId("turn-reasoning-final"), + "evt-reasoning-final", + "2026-05-03T00:00:03.000Z", + "before tool", + ); + harness.emit({ + type: "task.started", + eventId: asEventId("evt-task-after-reasoning"), + provider: ProviderDriverKind.make("opencode"), + createdAt: "2026-05-03T00:00:04.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-reasoning-final"), + payload: { taskId: "task-after-reasoning", description: "Run checks" }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => { + const reasoning = entry.messages.filter((message) => message.channel === "reasoning"); + return ( + reasoning.length === 3 && + reasoning.every((message) => !message.streaming) && + entry.activities.some((activity) => activity.id === "evt-task-after-reasoning") + ); + }); + expect( + thread.messages + .filter((message) => message.channel === "reasoning") + .map((message) => message.text), + ).toEqual(["old turn", "next turn", "before tool"]); + + const events = await Effect.runPromise( + Stream.runCollect(harness.engine.readEvents(0)).pipe( + Effect.map((chunk) => Array.from(chunk)), + ), + ); + const finalReasoningComplete = events.find( + (event) => + event.type === "thread.message-sent" && + event.payload.channel === "reasoning" && + event.payload.text === "" && + event.payload.turnId === "turn-reasoning-final", + ); + const activity = events.find( + (event) => + event.type === "thread.activity-appended" && + event.payload.activity.id === "evt-task-after-reasoning", + ); + expect(finalReasoningComplete?.sequence).toBeLessThan(activity?.sequence ?? 0); + }); + + it("completes a persisted streaming reasoning message after in-memory segment loss", async () => { + const harness = await createHarness(); + const recoveredTurnId = asTurnId("turn-reasoning-recovered"); + + // Seed projected streaming rows directly through the engine so the + // ingestion caches never learn about them, modeling a server restart or + // cache eviction between streaming and turn settlement. + await harness.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-seed-recovered-reasoning"), + threadId: asThreadId("thread-1"), + messageId: asMessageId("reasoning:thread-1:turn-reasoning-recovered:segment:0"), + delta: "Recovered thinking", + channel: "reasoning", + turnId: recoveredTurnId, + createdAt: "2026-05-05T00:00:00.000Z", + }); + await harness.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-seed-unrelated-reasoning"), + threadId: asThreadId("thread-1"), + messageId: asMessageId("reasoning:thread-1:turn-reasoning-unrelated:segment:0"), + delta: "Unrelated thinking", + channel: "reasoning", + turnId: asTurnId("turn-reasoning-unrelated"), + createdAt: "2026-05-05T00:00:01.000Z", + }); + await harness.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-seed-substantive-assistant"), + threadId: asThreadId("thread-1"), + messageId: asMessageId("assistant:item-reasoning-recovered"), + delta: "Streaming answer", + turnId: recoveredTurnId, + createdAt: "2026-05-05T00:00:02.000Z", + }); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-reasoning-recovered-complete"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-05T00:00:03.000Z", + threadId: asThreadId("thread-1"), + turnId: recoveredTurnId, + status: "completed", + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const messages = readModel.threads.find((entry) => entry.id === "thread-1")?.messages ?? []; + const byId = new Map(messages.map((message) => [message.id, message])); + + const recovered = byId.get( + asMessageId("reasoning:thread-1:turn-reasoning-recovered:segment:0"), + ); + expect(recovered?.streaming).toBe(false); + expect(recovered?.text).toBe("Recovered thinking"); + expect(recovered?.updatedAt).toBe("2026-05-05T00:00:03.000Z"); + // Reasoning owned by another turn and substantive assistant output stay + // untouched by the recovery. + expect( + byId.get(asMessageId("reasoning:thread-1:turn-reasoning-unrelated:segment:0"))?.streaming, + ).toBe(true); + expect(byId.get(asMessageId("assistant:item-reasoning-recovered"))?.streaming).toBe(true); + }); + + it("recovers an orphaned streaming reasoning segment when a delta creates its replacement", async () => { + const harness = await createHarness(); + const turnId = asTurnId("turn-reasoning-replaced"); + + // Seed projected streaming rows directly through the engine so the + // ingestion caches never learn about them, modeling a server restart or + // cache eviction between streaming and the next reasoning delta. + await harness.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-seed-orphaned-reasoning"), + threadId: asThreadId("thread-1"), + messageId: asMessageId("reasoning:thread-1:turn-reasoning-replaced:segment:0"), + delta: "Orphaned thinking", + channel: "reasoning", + turnId, + createdAt: "2026-05-08T00:00:00.000Z", + }); + await harness.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-seed-other-turn-reasoning"), + threadId: asThreadId("thread-1"), + messageId: asMessageId("reasoning:thread-1:turn-reasoning-other:segment:0"), + delta: "Unrelated thinking", + channel: "reasoning", + turnId: asTurnId("turn-reasoning-other"), + createdAt: "2026-05-08T00:00:01.000Z", + }); + await harness.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-seed-substantive-assistant-replaced"), + threadId: asThreadId("thread-1"), + messageId: asMessageId("assistant:item-reasoning-replaced"), + delta: "Streaming answer", + turnId, + createdAt: "2026-05-08T00:00:02.000Z", + }); + + // The new delta creates a replacement segment and repopulates the cache, + // so the terminal boundary no longer sees a wholly absent state; the + // orphaned row must still settle. + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-replacement-delta"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-08T00:00:03.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "Replacement thinking" }, + }); + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-reasoning-replaced-turn-complete"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-08T00:00:04.000Z", + threadId: asThreadId("thread-1"), + turnId, + status: "completed", + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const messages = readModel.threads.find((entry) => entry.id === "thread-1")?.messages ?? []; + const byId = new Map(messages.map((message) => [message.id, message])); + + // The orphan settles at the replacement delta, not the turn boundary, so + // its duration does not absorb the rest of the turn. + const orphan = byId.get(asMessageId("reasoning:thread-1:turn-reasoning-replaced:segment:0")); + expect(orphan?.streaming).toBe(false); + expect(orphan?.text).toBe("Orphaned thinking"); + expect(orphan?.updatedAt).toBe("2026-05-08T00:00:03.000Z"); + + const replacement = byId.get( + asMessageId("reasoning:thread-1:turn-reasoning-replaced:segment:1"), + ); + expect(replacement?.streaming).toBe(false); + expect(replacement?.text).toBe("Replacement thinking"); + + expect( + messages.filter((message) => message.channel === "reasoning" && message.turnId === turnId), + ).toHaveLength(2); + // Reasoning owned by another turn and substantive assistant output stay + // untouched by the recovery. + expect( + byId.get(asMessageId("reasoning:thread-1:turn-reasoning-other:segment:0"))?.streaming, + ).toBe(true); + expect(byId.get(asMessageId("assistant:item-reasoning-replaced"))?.streaming).toBe(true); + }); + + it("scopes turnless replacement-delta recovery to turnless reasoning rows", async () => { + // Streaming mode projects the replacement delta immediately, so the + // mid-stream state right after recovery is observable. + const harness = await createHarness({ serverSettings: { enableLegacyTokenStreaming: true } }); + + // Seed projected streaming rows the ingestion caches never learn about: a + // turnless orphan and a named-turn row that must survive the recovery. + await harness.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-seed-turnless-orphan-reasoning"), + threadId: asThreadId("thread-1"), + messageId: asMessageId("reasoning:thread-1:turnless:segment:0"), + delta: "Orphaned turnless thinking", + channel: "reasoning", + createdAt: "2026-05-09T00:00:00.000Z", + }); + await harness.dispatch({ + type: "thread.message.assistant.delta", + commandId: CommandId.make("cmd-seed-named-turn-reasoning"), + threadId: asThreadId("thread-1"), + messageId: asMessageId("reasoning:thread-1:turn-reasoning-named:segment:0"), + delta: "Named-turn thinking", + channel: "reasoning", + turnId: asTurnId("turn-reasoning-named"), + createdAt: "2026-05-09T00:00:01.000Z", + }); + + // A turnless delta on the cold path must settle only the turnless orphan, + // never sweep the named turn's row. + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-turnless-replacement-delta"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-09T00:00:02.000Z", + threadId: asThreadId("thread-1"), + payload: { streamKind: "reasoning_text", delta: "Turnless replacement thinking" }, + }); + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turnless-replacement-turn-complete"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-09T00:00:03.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-turnless-late-bound"), + status: "completed", + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const messages = readModel.threads.find((entry) => entry.id === "thread-1")?.messages ?? []; + const byId = new Map(messages.map((message) => [message.id, message])); + + const orphan = byId.get(asMessageId("reasoning:thread-1:turnless:segment:0")); + expect(orphan?.streaming).toBe(false); + expect(orphan?.text).toBe("Orphaned turnless thinking"); + expect(orphan?.updatedAt).toBe("2026-05-09T00:00:02.000Z"); + + const replacement = byId.get(asMessageId("reasoning:thread-1:turnless:segment:1")); + expect(replacement?.streaming).toBe(false); + expect(replacement?.text).toBe("Turnless replacement thinking"); + + expect( + byId.get(asMessageId("reasoning:thread-1:turn-reasoning-named:segment:0"))?.streaming, + ).toBe(true); + }); + + it("finalizes active reasoning when a tool user-input request opens", async () => { + const harness = await createHarness(); + const turnId = asTurnId("turn-reasoning-tool-input"); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-before-tool-input"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-06T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "Considering the question" }, + }); + harness.emit({ + type: "request.opened", + eventId: asEventId("evt-tool-input-request-opened"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-06T00:00:01.000Z", + threadId: asThreadId("thread-1"), + turnId, + requestId: ApprovalRequestId.make("req-tool-input"), + payload: { requestType: "tool_user_input", detail: "Pick one" }, + }); + await harness.drain(); + + const pausedModel = await harness.readModel(); + const pausedReasoning = pausedModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.filter((message) => message.channel === "reasoning"); + expect(pausedReasoning).toEqual([ + expect.objectContaining({ + id: "reasoning:thread-1:turn-reasoning-tool-input:segment:0", + streaming: false, + text: "Considering the question", + updatedAt: "2026-05-06T00:00:01.000Z", + }), + ]); + + // Reasoning after the human answers lands in a new segment instead of + // extending the paused one with wait time. + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-after-tool-input"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-06T00:05:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "After the answer" }, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-reasoning-tool-input-aborted"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-06T00:05:01.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { reason: "cancelled" }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const reasoning = readModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.filter((message) => message.channel === "reasoning"); + expect(reasoning?.map((message) => [message.id, message.text])).toEqual([ + ["reasoning:thread-1:turn-reasoning-tool-input:segment:0", "Considering the question"], + ["reasoning:thread-1:turn-reasoning-tool-input:segment:1", "After the answer"], + ]); + }); + + it("finalizes active reasoning when user input is requested", async () => { + const harness = await createHarness(); + const turnId = asTurnId("turn-reasoning-user-input"); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-before-user-input"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-07T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "Weighing options" }, + }); + harness.emit({ + type: "user-input.requested", + eventId: asEventId("evt-reasoning-user-input-requested"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-07T00:00:01.000Z", + threadId: asThreadId("thread-1"), + turnId, + requestId: ApprovalRequestId.make("req-reasoning-user-input"), + payload: { + questions: [ + { + id: "choice", + header: "Choice", + question: "Pick one", + options: [{ label: "A", description: "Option A" }], + }, + ], + }, + }); + await harness.drain(); + + const pausedModel = await harness.readModel(); + const pausedReasoning = pausedModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.filter((message) => message.channel === "reasoning"); + expect(pausedReasoning).toEqual([ + expect.objectContaining({ + id: "reasoning:thread-1:turn-reasoning-user-input:segment:0", + streaming: false, + text: "Weighing options", + updatedAt: "2026-05-07T00:00:01.000Z", + }), + ]); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-reasoning-after-user-input"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-07T00:10:00.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { streamKind: "reasoning_text", delta: "After the reply" }, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-reasoning-user-input-aborted"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-05-07T00:10:01.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { reason: "cancelled" }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const reasoning = readModel.threads + .find((entry) => entry.id === "thread-1") + ?.messages.filter((message) => message.channel === "reasoning"); + expect(reasoning?.map((message) => [message.id, message.text])).toEqual([ + ["reasoning:thread-1:turn-reasoning-user-input:segment:0", "Weighing options"], + ["reasoning:thread-1:turn-reasoning-user-input:segment:1", "After the reply"], + ]); + }); }); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 7ec3a7e64243..1dddbe4ebe7e 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -9,6 +9,7 @@ import { CheckpointRef, classifyTaskAgentKind, EventId, + isReasoningMessage, isToolLifecycleItemType, ThreadId, type ThreadTokenUsageSnapshot, @@ -23,6 +24,7 @@ import * as Cache from "effect/Cache"; import * as Cause from "effect/Cause"; import * as Crypto from "effect/Crypto"; import * as Duration from "effect/Duration"; +import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -90,6 +92,17 @@ interface AssistantSegmentState { activeMessageId: MessageId | null; } +interface ReasoningSegmentState { + turnId: TurnId | null; + nextSegmentIndex: number; + active: { + messageId: MessageId; + firstDeltaAt: string; + hasProjectedMessage: boolean; + deliveryMode: AssistantDeliveryMode; + } | null; +} + const TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY = 10_000; const TURN_MESSAGE_IDS_BY_TURN_TTL = Duration.minutes(120); const BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_CACHE_CAPACITY = 20_000; @@ -141,7 +154,7 @@ function hasAssistantMessageForTurn( if (!message) { continue; } - if (message.role !== "assistant" || message.turnId !== turnId) { + if (message.role !== "assistant" || isReasoningMessage(message) || message.turnId !== turnId) { continue; } if (options?.streamingOnly === true && !message.streaming) { @@ -165,6 +178,25 @@ function findMessageById( return undefined; } +function nextReasoningSegmentIndex( + messages: ReadonlyArray, + threadId: ThreadId, + turnId: TurnId | null, +): number { + const prefix = `reasoning:${threadId}:${turnId ?? "turnless"}:segment:`; + let nextIndex = 0; + for (const message of messages) { + if (!isReasoningMessage(message) || !message.id.startsWith(prefix)) { + continue; + } + const segmentIndex = Number(message.id.slice(prefix.length)); + if (Number.isSafeInteger(segmentIndex) && segmentIndex >= nextIndex) { + nextIndex = segmentIndex + 1; + } + } + return nextIndex; +} + function findProposedPlanById( proposedPlans: ReadonlyArray< Pick @@ -901,6 +933,11 @@ const make = Effect.gen(function* () { Effect.map((uuid) => CommandId.make(`provider:${event.eventId}:${tag}:${uuid}`)), ); + // Read per event: the setting can change between turns. + const getAssistantDeliveryMode = Effect.map(serverSettingsService.getSettings, (settings) => + settings.enableLegacyTokenStreaming ? "streaming" : "buffered", + ); + const turnMessageIdsByTurnKey = yield* Cache.make>({ capacity: TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY, timeToLive: TURN_MESSAGE_IDS_BY_TURN_TTL, @@ -922,6 +959,16 @@ const make = Effect.gen(function* () { ), }); + const reasoningSegmentByThreadId = yield* Cache.make({ + capacity: TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY, + timeToLive: TURN_MESSAGE_IDS_BY_TURN_TTL, + lookup: () => + Effect.die( + new Error("reasoning segment state should be read through getOption before initialization"), + ), + }); + const lastMessageStampMsByThreadId = new Map(); + const bufferedProposedPlanById = yield* Cache.make({ capacity: BUFFERED_PROPOSED_PLAN_BY_ID_CACHE_CAPACITY, timeToLive: BUFFERED_PROPOSED_PLAN_BY_ID_TTL, @@ -961,6 +1008,25 @@ const make = Effect.gen(function* () { .pipe(Effect.map(Option.getOrUndefined)); }); + const nextMessageStamp = Effect.fn("nextMessageStamp")(function* ( + threadId: ThreadId, + intendedAt: string, + ) { + let lastStampMs = lastMessageStampMsByThreadId.get(threadId); + if (lastStampMs === undefined) { + const thread = yield* resolveThreadDetail(threadId); + lastStampMs = thread?.messages.reduce( + (latest, message) => + Math.max(latest, DateTime.toEpochMillis(DateTime.makeUnsafe(message.createdAt))), + Number.NEGATIVE_INFINITY, + ); + } + const intendedMs = DateTime.toEpochMillis(DateTime.makeUnsafe(intendedAt)); + const nextStampMs = Math.max(intendedMs, (lastStampMs ?? Number.NEGATIVE_INFINITY) + 1); + lastMessageStampMsByThreadId.set(threadId, nextStampMs); + return DateTime.formatIso(DateTime.makeUnsafe(nextStampMs)); + }); + const rememberAssistantMessageId = (threadId: ThreadId, turnId: TurnId, messageId: MessageId) => Cache.getOption(turnMessageIdsByTurnKey, providerTurnKey(threadId, turnId)).pipe( Effect.flatMap((existingIds) => @@ -1162,7 +1228,7 @@ const make = Effect.gen(function* () { messageId: input.messageId, delta: bufferedText, ...(input.turnId ? { turnId: input.turnId } : {}), - createdAt: input.createdAt, + createdAt: yield* nextMessageStamp(input.threadId, input.createdAt), }); return true; }); @@ -1229,7 +1295,7 @@ const make = Effect.gen(function* () { messageId: input.messageId, delta: text, ...(input.turnId ? { turnId: input.turnId } : {}), - createdAt: input.createdAt, + createdAt: yield* nextMessageStamp(input.threadId, input.createdAt), }); } @@ -1240,7 +1306,7 @@ const make = Effect.gen(function* () { threadId: input.threadId, messageId: input.messageId, ...(input.turnId ? { turnId: input.turnId } : {}), - createdAt: input.createdAt, + createdAt: yield* nextMessageStamp(input.threadId, input.createdAt), }); } yield* clearAssistantMessageState(input.messageId); @@ -1288,6 +1354,139 @@ const make = Effect.gen(function* () { } }); + const getReasoningSegmentState = (threadId: ThreadId) => + Cache.getOption(reasoningSegmentByThreadId, threadId).pipe(Effect.map(Option.getOrUndefined)); + + // A thread has at most one open reasoning segment. The active path stays in + // memory; a cold/new turn derives its next index from projected messages so + // cache eviction and session cleanup cannot reuse a durable message id. + const getOrCreateReasoningSegment = (input: { + threadId: ThreadId; + turnId: TurnId | null; + firstDeltaAt: string; + deliveryMode: AssistantDeliveryMode; + }) => + Effect.gen(function* () { + const existing = yield* getReasoningSegmentState(input.threadId); + const sameTurn = existing !== undefined && existing.turnId === input.turnId; + if (sameTurn && existing.active !== null) { + return { ...existing, active: existing.active }; + } + + const segmentIndex = sameTurn + ? existing.nextSegmentIndex + : nextReasoningSegmentIndex( + (yield* resolveThreadDetail(input.threadId))?.messages ?? [], + input.threadId, + input.turnId, + ); + const state = { + turnId: input.turnId, + nextSegmentIndex: segmentIndex + 1, + active: { + messageId: MessageId.make( + `reasoning:${input.threadId}:${input.turnId ?? "turnless"}:segment:${segmentIndex}`, + ), + firstDeltaAt: input.firstDeltaAt, + hasProjectedMessage: false, + deliveryMode: input.deliveryMode, + }, + } satisfies ReasoningSegmentState; + yield* Cache.set(reasoningSegmentByThreadId, input.threadId, state); + return state; + }); + + // Durable complement to the in-memory reasoning segment: after a restart or + // cache eviction a projected reasoning row can still be streaming with no + // state left to close it, so settle it straight from the projection. The + // turn guard mirrors finalizeActiveReasoningSegment: a named turn closes + // its own and turnless rows, never another turn's; `null` closes turnless + // rows only; omitting `turnId` is the deliberate unscoped terminal/session + // sweep that closes everything. + const completeProjectedStreamingReasoning = (input: { + event: ProviderRuntimeEvent; + threadId: ThreadId; + turnId?: TurnId | null; + }) => + Effect.gen(function* () { + const messages = (yield* resolveThreadDetail(input.threadId))?.messages ?? []; + for (const message of messages) { + if (!isReasoningMessage(message) || !message.streaming) { + continue; + } + if ( + input.turnId !== undefined && + message.turnId !== null && + message.turnId !== input.turnId + ) { + continue; + } + yield* orchestrationEngine.dispatch({ + type: "thread.message.assistant.complete", + commandId: yield* providerCommandId(input.event, "reasoning-recover-complete"), + threadId: input.threadId, + messageId: message.id, + channel: "reasoning", + ...(message.turnId !== null ? { turnId: message.turnId } : {}), + createdAt: yield* nextMessageStamp(input.threadId, input.event.createdAt), + }); + } + }); + + // Closes the thread's open reasoning segment; `turnId` limits the close to a + // segment opened by that turn. `recoverProjected` also settles reasoning + // rows left streaming in the projection when the in-memory state was lost + // (restart, TTL, capacity eviction); only terminal/pause boundaries pass it + // so hot paths never load thread detail. + const finalizeActiveReasoningSegment = (input: { + event: ProviderRuntimeEvent; + threadId: ThreadId; + turnId?: TurnId; + recoverProjected?: boolean; + }) => + Effect.gen(function* () { + const state = yield* getReasoningSegmentState(input.threadId); + if (state === undefined && input.recoverProjected === true) { + yield* completeProjectedStreamingReasoning(input); + return; + } + const active = state?.active ?? null; + if ( + state === undefined || + active === null || + (input.turnId !== undefined && state.turnId !== null && state.turnId !== input.turnId) + ) { + return; + } + + const bufferedText = yield* takeBufferedAssistantText(active.messageId); + const hasText = hasRenderableAssistantText(bufferedText); + if (hasText) { + yield* orchestrationEngine.dispatch({ + type: "thread.message.assistant.delta", + commandId: yield* providerCommandId(input.event, "reasoning-finalize-delta"), + threadId: input.threadId, + messageId: active.messageId, + delta: bufferedText, + channel: "reasoning", + ...(state.turnId !== null ? { turnId: state.turnId } : {}), + createdAt: yield* nextMessageStamp(input.threadId, active.firstDeltaAt), + }); + } + if (active.hasProjectedMessage || hasText) { + yield* orchestrationEngine.dispatch({ + type: "thread.message.assistant.complete", + commandId: yield* providerCommandId(input.event, "reasoning-finalize-complete"), + threadId: input.threadId, + messageId: active.messageId, + channel: "reasoning", + ...(state.turnId !== null ? { turnId: state.turnId } : {}), + createdAt: yield* nextMessageStamp(input.threadId, input.event.createdAt), + }); + } + yield* Cache.set(reasoningSegmentByThreadId, input.threadId, { ...state, active: null }); + }); + const upsertProposedPlan = (input: { event: ProviderRuntimeEvent; threadId: ThreadId; @@ -1401,6 +1600,9 @@ const make = Effect.gen(function* () { : Effect.void, { concurrency: 1 }, ).pipe(Effect.asVoid); + // The reasoning segment is already finalized by the session.exited trigger. + yield* Cache.invalidate(reasoningSegmentByThreadId, threadId); + lastMessageStampMsByThreadId.delete(threadId); yield* Effect.forEach( proposedPlanKeys, (key) => @@ -1570,6 +1772,31 @@ const make = Effect.gen(function* () { ? yield* getSourceProposedPlanReferenceForAcceptedTurnStart(thread.id, eventTurnId) : null; + const terminalSessionState = + event.type === "session.state.changed" && + !sessionStatusAllowsActiveTurn( + orchestrationSessionStatusFromRuntimeState(event.payload.state), + ); + if ( + (event.type === "turn.started" && shouldApplyThreadLifecycle) || + (event.type === "runtime.error" && !conflictsWithActiveTurn) || + event.type === "session.exited" || + terminalSessionState + ) { + yield* finalizeActiveReasoningSegment({ + event, + threadId: thread.id, + recoverProjected: true, + }); + } else if (event.type === "turn.completed" || event.type === "turn.aborted") { + yield* finalizeActiveReasoningSegment({ + event, + threadId: thread.id, + ...(eventTurnId !== undefined ? { turnId: eventTurnId } : {}), + recoverProjected: true, + }); + } + if ( event.type === "session.started" || event.type === "session.state.changed" || @@ -1602,14 +1829,11 @@ const make = Effect.gen(function* () { const nextActiveTurnId = event.type === "turn.started" ? (eventTurnId ?? null) - : event.type === "turn.completed" || event.type === "session.exited" + : event.type === "turn.completed" || + event.type === "session.exited" || + terminalSessionState ? null - : event.type === "session.state.changed" && - !sessionStatusAllowsActiveTurn( - orchestrationSessionStatusFromRuntimeState(event.payload.state), - ) - ? null - : activeTurnId; + : activeTurnId; const lastError = event.type === "session.state.changed" && event.payload.state === "error" ? (event.payload.reason ?? thread.session?.lastError ?? "Provider session error") @@ -1666,11 +1890,22 @@ const make = Effect.gen(function* () { event.type === "content.delta" && event.payload.streamKind === "assistant_text" ? event.payload.delta : undefined; + const reasoningDelta = + event.type === "content.delta" && + (event.payload.streamKind === "reasoning_text" || + event.payload.streamKind === "reasoning_summary_text") + ? event.payload.delta + : undefined; const proposedPlanDelta = event.type === "turn.proposed.delta" ? event.payload.delta : undefined; if (assistantDelta && assistantDelta.length > 0) { const turnId = toTurnId(event.turnId); + yield* finalizeActiveReasoningSegment({ + event, + threadId: thread.id, + ...(turnId !== undefined ? { turnId } : {}), + }); const assistantMessageId = yield* getOrCreateAssistantMessageId({ threadId: thread.id, event, @@ -1680,10 +1915,7 @@ const make = Effect.gen(function* () { yield* rememberAssistantMessageId(thread.id, turnId, assistantMessageId); } - const assistantDeliveryMode: AssistantDeliveryMode = yield* Effect.map( - serverSettingsService.getSettings, - (settings) => (settings.enableLegacyTokenStreaming ? "streaming" : "buffered"), - ); + const assistantDeliveryMode = yield* getAssistantDeliveryMode; if (assistantDeliveryMode === "buffered") { const spillChunk = yield* appendBufferedAssistantText(assistantMessageId, assistantDelta); if (spillChunk.length > 0) { @@ -1694,7 +1926,7 @@ const make = Effect.gen(function* () { messageId: assistantMessageId, delta: spillChunk, ...(turnId ? { turnId } : {}), - createdAt: now, + createdAt: yield* nextMessageStamp(thread.id, now), }); } } else { @@ -1705,21 +1937,83 @@ const make = Effect.gen(function* () { messageId: assistantMessageId, delta: assistantDelta, ...(turnId ? { turnId } : {}), - createdAt: now, + createdAt: yield* nextMessageStamp(thread.id, now), }); } } + if (reasoningDelta && reasoningDelta.length > 0) { + const turnId = toTurnId(event.turnId) ?? activeTurnId; + const openSegment = yield* getReasoningSegmentState(thread.id); + // Reasoning from another turn closes the segment the previous turn opened. + if (openSegment?.active && openSegment.turnId !== turnId) { + yield* finalizeActiveReasoningSegment({ event, threadId: thread.id }); + } + // Wholly absent state means a projected row left streaming is an + // orphan (restart, TTL, capacity eviction). Settle it before the + // replacement below repopulates the cache — after that the terminal + // recoverProjected path can no longer see the loss. This scans the + // projection only on the cold path, which loads thread detail anyway. + if (openSegment === undefined) { + yield* completeProjectedStreamingReasoning({ + event, + threadId: thread.id, + turnId, + }); + } + const deliveryMode = yield* getAssistantDeliveryMode; + const reasoningState = yield* getOrCreateReasoningSegment({ + threadId: thread.id, + turnId, + firstDeltaAt: now, + deliveryMode, + }); + const active = reasoningState.active; + const buffered = active.deliveryMode === "buffered"; + // Buffered mode only projects when the buffer spills; streaming projects every delta. + const projectedDelta = buffered + ? yield* appendBufferedAssistantText(active.messageId, reasoningDelta) + : reasoningDelta; + if (projectedDelta.length > 0) { + yield* orchestrationEngine.dispatch({ + type: "thread.message.assistant.delta", + commandId: yield* providerCommandId( + event, + buffered ? "reasoning-delta-buffer-spill" : "reasoning-delta", + ), + threadId: thread.id, + messageId: active.messageId, + delta: projectedDelta, + channel: "reasoning", + ...(turnId !== null ? { turnId } : {}), + createdAt: yield* nextMessageStamp(thread.id, buffered ? active.firstDeltaAt : now), + }); + // A projected segment gets a completion even when its text is whitespace only. + if (!active.hasProjectedMessage) { + yield* Cache.set(reasoningSegmentByThreadId, thread.id, { + ...reasoningState, + active: { ...active, hasProjectedMessage: true }, + }); + } + } + } + const pauseForUserTurnId = event.type === "request.opened" || event.type === "user-input.requested" ? toTurnId(event.turnId) : undefined; if (pauseForUserTurnId) { + // A user-facing pause also settles in-flight reasoning so it does not + // stay "Thinking" or absorb human wait time; tool_user_input requests + // emit no activity, so nothing later would close it. + yield* finalizeActiveReasoningSegment({ + event, + threadId: thread.id, + turnId: pauseForUserTurnId, + recoverProjected: true, + }); const detailedThread = yield* getLoadedThreadDetail(); - const assistantDeliveryMode: AssistantDeliveryMode = yield* Effect.map( - serverSettingsService.getSettings, - (settings) => (settings.enableLegacyTokenStreaming ? "streaming" : "buffered"), - ); + const assistantDeliveryMode = yield* getAssistantDeliveryMode; const flushedMessageIds = assistantDeliveryMode === "buffered" ? yield* flushBufferedAssistantMessagesForTurn({ @@ -1779,6 +2073,12 @@ const make = Effect.gen(function* () { : undefined; if (assistantCompletion) { + const completionTurnId = toTurnId(event.turnId); + yield* finalizeActiveReasoningSegment({ + event, + threadId: thread.id, + ...(completionTurnId !== undefined ? { turnId: completionTurnId } : {}), + }); const detailedThread = yield* getLoadedThreadDetail(); const messages = detailedThread?.messages ?? []; const turnId = toTurnId(event.turnId); @@ -2028,6 +2328,13 @@ const make = Effect.gen(function* () { } const activities = runtimeEventToActivities(event, taskTitle); + if (activities.length > 0) { + yield* finalizeActiveReasoningSegment({ + event, + threadId: thread.id, + ...(eventTurnId !== undefined ? { turnId: eventTurnId } : {}), + }); + } yield* Effect.forEach(activities, (activity) => providerCommandId(event, "thread-activity-append").pipe( Effect.flatMap((commandId) => @@ -2041,6 +2348,10 @@ const make = Effect.gen(function* () { ), ), ).pipe(Effect.asVoid); + + if (event.type === "turn.completed" || event.type === "turn.aborted") { + lastMessageStampMsByThreadId.delete(thread.id); + } }); const processDomainEvent = (_event: TurnStartRequestedDomainEvent) => Effect.void; diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index f3fdd462f437..cb532ae6ad71 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -1246,6 +1246,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" threadId: command.threadId, messageId: command.messageId, role: "assistant", + ...(command.channel !== undefined ? { channel: command.channel } : {}), text: command.delta, turnId: command.turnId ?? null, streaming: true, @@ -1273,6 +1274,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" threadId: command.threadId, messageId: command.messageId, role: "assistant", + ...(command.channel !== undefined ? { channel: command.channel } : {}), text: "", turnId: command.turnId ?? null, streaming: false, diff --git a/apps/server/src/orchestration/projector.test.ts b/apps/server/src/orchestration/projector.test.ts index dad3d07370f9..1b69ad86669b 100644 --- a/apps/server/src/orchestration/projector.test.ts +++ b/apps/server/src/orchestration/projector.test.ts @@ -1,3 +1,4 @@ +import { it as effectIt } from "@effect/vitest"; import { CommandId, EventId, @@ -448,6 +449,7 @@ describe("orchestration projector", () => { threadId: "thread-1", messageId: "assistant:msg-1", role: "assistant", + channel: "reasoning", text: "hello", turnId: "turn-1", streaming: true, @@ -472,6 +474,7 @@ describe("orchestration projector", () => { threadId: "thread-1", messageId: "assistant:msg-1", role: "assistant", + channel: "reasoning", text: "", turnId: "turn-1", streaming: false, @@ -484,12 +487,13 @@ describe("orchestration projector", () => { const message = afterComplete.threads[0]?.messages[0]; expect(message?.id).toBe("assistant:msg-1"); + expect(message?.channel).toBe("reasoning"); expect(message?.text).toBe("hello"); expect(message?.streaming).toBe(false); expect(message?.updatedAt).toBe(completeAt); }); - it("prunes reverted turn messages from in-memory thread snapshot", async () => { + it("prunes reverted reasoning messages from in-memory thread snapshot", async () => { const createdAt = "2026-02-23T10:00:00.000Z"; const model = createEmptyReadModel(createdAt); @@ -625,6 +629,7 @@ describe("orchestration projector", () => { threadId: "thread-1", messageId: "assistant-msg-2", role: "assistant", + channel: "reasoning", text: "Updated README to v3.\n", turnId: "turn-2", streaming: false, @@ -857,6 +862,356 @@ describe("orchestration projector", () => { ).toEqual([{ id: "assistant-keep", role: "assistant", turnId: "turn-1" }]); }); + // Regression: retainThreadMessagesAfterRevert counts a kept turn's + // channel=reasoning message as the turn's assistant, so the turnless + // substantive assistant is never recovered through the fallback scan and + // is dropped by the revert. + effectIt.effect( + "reasoning on a kept turn does not satisfy the substantive assistant count on revert", + () => + Effect.gen(function* () { + const createdAt = "2026-02-27T12:00:00.000Z"; + const model = createEmptyReadModel(createdAt); + + const afterCreate = yield* projectEvent( + model, + makeEvent({ + sequence: 1, + type: "thread.created", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: createdAt, + commandId: "cmd-rrcount-create", + payload: { + threadId: "thread-reason-count", + projectId: "project-1", + title: "demo", + modelSelection: { + provider: ProviderDriverKind.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt, + updatedAt: createdAt, + }, + }), + ); + + const events: ReadonlyArray = [ + makeEvent({ + sequence: 2, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: "2026-02-27T12:00:01.500Z", + commandId: "cmd-rrcount-user-1", + payload: { + threadId: "thread-reason-count", + messageId: "user-count-1", + role: "user", + text: "first prompt", + turnId: "turn-count-1", + streaming: false, + createdAt: "2026-02-27T12:00:01.500Z", + updatedAt: "2026-02-27T12:00:01.500Z", + }, + }), + makeEvent({ + sequence: 3, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: "2026-02-27T12:00:01.600Z", + commandId: "cmd-rrcount-reasoning-1", + payload: { + threadId: "thread-reason-count", + messageId: "reasoning-count-1", + role: "assistant", + channel: "reasoning", + text: "thinking about the first prompt", + turnId: "turn-count-1", + streaming: false, + createdAt: "2026-02-27T12:00:01.600Z", + updatedAt: "2026-02-27T12:00:01.600Z", + }, + }), + makeEvent({ + sequence: 4, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: "2026-02-27T12:00:01.700Z", + commandId: "cmd-rrcount-assistant-1", + payload: { + threadId: "thread-reason-count", + messageId: "assistant-count-1", + role: "assistant", + text: "substantive first answer", + turnId: null, + streaming: false, + createdAt: "2026-02-27T12:00:01.700Z", + updatedAt: "2026-02-27T12:00:01.700Z", + }, + }), + makeEvent({ + sequence: 5, + type: "thread.turn-diff-completed", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: "2026-02-27T12:00:02.000Z", + commandId: "cmd-rrcount-turn-1", + payload: { + threadId: "thread-reason-count", + turnId: "turn-count-1", + checkpointTurnCount: 1, + checkpointRef: "refs/t3/checkpoints/thread-reason-count/turn/1", + status: "ready", + files: [], + assistantMessageId: null, + completedAt: "2026-02-27T12:00:02.000Z", + }, + }), + makeEvent({ + sequence: 6, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: "2026-02-27T12:00:02.500Z", + commandId: "cmd-rrcount-user-2", + payload: { + threadId: "thread-reason-count", + messageId: "user-count-2", + role: "user", + text: "second prompt", + turnId: "turn-count-2", + streaming: false, + createdAt: "2026-02-27T12:00:02.500Z", + updatedAt: "2026-02-27T12:00:02.500Z", + }, + }), + makeEvent({ + sequence: 7, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: "2026-02-27T12:00:02.600Z", + commandId: "cmd-rrcount-assistant-2", + payload: { + threadId: "thread-reason-count", + messageId: "assistant-count-2", + role: "assistant", + text: "second answer", + turnId: "turn-count-2", + streaming: false, + createdAt: "2026-02-27T12:00:02.600Z", + updatedAt: "2026-02-27T12:00:02.600Z", + }, + }), + makeEvent({ + sequence: 8, + type: "thread.turn-diff-completed", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: "2026-02-27T12:00:03.000Z", + commandId: "cmd-rrcount-turn-2", + payload: { + threadId: "thread-reason-count", + turnId: "turn-count-2", + checkpointTurnCount: 2, + checkpointRef: "refs/t3/checkpoints/thread-reason-count/turn/2", + status: "ready", + files: [], + assistantMessageId: "assistant-count-2", + completedAt: "2026-02-27T12:00:03.000Z", + }, + }), + makeEvent({ + sequence: 9, + type: "thread.reverted", + aggregateKind: "thread", + aggregateId: "thread-reason-count", + occurredAt: "2026-02-27T12:00:04.000Z", + commandId: "cmd-rrcount-revert", + payload: { + threadId: "thread-reason-count", + turnCount: 1, + }, + }), + ]; + + let afterRevert = afterCreate; + for (const event of events) { + afterRevert = yield* projectEvent(afterRevert, event); + } + + const thread = afterRevert.threads[0]; + expect( + thread?.messages.map((message) => ({ + id: message.id, + role: message.role, + channel: message.channel ?? null, + turnId: message.turnId, + })), + ).toEqual([ + { id: "user-count-1", role: "user", channel: null, turnId: "turn-count-1" }, + { + id: "reasoning-count-1", + role: "assistant", + channel: "reasoning", + turnId: "turn-count-1", + }, + { id: "assistant-count-1", role: "assistant", channel: null, turnId: null }, + ]); + }), + ); + + // Regression: the fallback assistant scan treats every role=assistant + // message as eligible, so a turnless channel=reasoning message left over + // from discarded later work is selected as the kept turn's assistant and + // survives the revert while the substantive assistant is dropped. + effectIt.effect( + "turnless reasoning from discarded work is never selected as a fallback assistant on revert", + () => + Effect.gen(function* () { + const createdAt = "2026-02-27T12:10:00.000Z"; + const model = createEmptyReadModel(createdAt); + + const afterCreate = yield* projectEvent( + model, + makeEvent({ + sequence: 1, + type: "thread.created", + aggregateKind: "thread", + aggregateId: "thread-reason-fallback", + occurredAt: createdAt, + commandId: "cmd-rrfall-create", + payload: { + threadId: "thread-reason-fallback", + projectId: "project-1", + title: "demo", + modelSelection: { + provider: ProviderDriverKind.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + branch: null, + worktreePath: null, + createdAt, + updatedAt: createdAt, + }, + }), + ); + + const events: ReadonlyArray = [ + makeEvent({ + sequence: 2, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-reason-fallback", + occurredAt: "2026-02-27T12:10:01.500Z", + commandId: "cmd-rrfall-user-1", + payload: { + threadId: "thread-reason-fallback", + messageId: "user-fallback-1", + role: "user", + text: "kept prompt", + turnId: "turn-fallback-1", + streaming: false, + createdAt: "2026-02-27T12:10:01.500Z", + updatedAt: "2026-02-27T12:10:01.500Z", + }, + }), + makeEvent({ + sequence: 3, + type: "thread.turn-diff-completed", + aggregateKind: "thread", + aggregateId: "thread-reason-fallback", + occurredAt: "2026-02-27T12:10:02.000Z", + commandId: "cmd-rrfall-turn-1", + payload: { + threadId: "thread-reason-fallback", + turnId: "turn-fallback-1", + checkpointTurnCount: 1, + checkpointRef: "refs/t3/checkpoints/thread-reason-fallback/turn/1", + status: "ready", + files: [], + assistantMessageId: null, + completedAt: "2026-02-27T12:10:02.000Z", + }, + }), + makeEvent({ + sequence: 4, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-reason-fallback", + occurredAt: "2026-02-27T12:10:02.100Z", + commandId: "cmd-rrfall-reasoning", + payload: { + threadId: "thread-reason-fallback", + messageId: "reasoning-discarded", + role: "assistant", + channel: "reasoning", + text: "planning work that gets reverted away", + turnId: null, + streaming: false, + createdAt: "2026-02-27T12:10:02.100Z", + updatedAt: "2026-02-27T12:10:02.100Z", + }, + }), + makeEvent({ + sequence: 5, + type: "thread.message-sent", + aggregateKind: "thread", + aggregateId: "thread-reason-fallback", + occurredAt: "2026-02-27T12:10:02.200Z", + commandId: "cmd-rrfall-assistant", + payload: { + threadId: "thread-reason-fallback", + messageId: "assistant-recovered", + role: "assistant", + text: "substantive answer for the kept turn", + turnId: null, + streaming: false, + createdAt: "2026-02-27T12:10:02.200Z", + updatedAt: "2026-02-27T12:10:02.200Z", + }, + }), + makeEvent({ + sequence: 6, + type: "thread.reverted", + aggregateKind: "thread", + aggregateId: "thread-reason-fallback", + occurredAt: "2026-02-27T12:10:03.000Z", + commandId: "cmd-rrfall-revert", + payload: { + threadId: "thread-reason-fallback", + turnCount: 1, + }, + }), + ]; + + let afterRevert = afterCreate; + for (const event of events) { + afterRevert = yield* projectEvent(afterRevert, event); + } + + const thread = afterRevert.threads[0]; + expect( + thread?.messages.map((message) => ({ + id: message.id, + role: message.role, + channel: message.channel ?? null, + turnId: message.turnId, + })), + ).toEqual([ + { id: "user-fallback-1", role: "user", channel: null, turnId: "turn-fallback-1" }, + { id: "assistant-recovered", role: "assistant", channel: null, turnId: null }, + ]); + }), + ); + 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 1c4cd65d5123..af1482073422 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -126,7 +126,10 @@ function retainThreadMessagesAfterRevert( } const retainedAssistantCount = messages.filter( - (message) => message.role === "assistant" && retainedMessageIds.has(message.id), + (message) => + message.role === "assistant" && + message.channel !== "reasoning" && + retainedMessageIds.has(message.id), ).length; const missingAssistantCount = Math.max(0, turnCount - retainedAssistantCount); if (missingAssistantCount > 0) { @@ -134,6 +137,7 @@ function retainThreadMessagesAfterRevert( .filter( (message) => message.role === "assistant" && + message.channel !== "reasoning" && !retainedMessageIds.has(message.id) && (message.turnId === null || retainedTurnIds.has(message.turnId)), ) @@ -521,6 +525,7 @@ export function projectEvent( { id: payload.messageId, role: payload.role, + ...(payload.channel !== undefined ? { channel: payload.channel } : {}), text: payload.text, ...(payload.attachments !== undefined ? { attachments: payload.attachments } : {}), turnId: payload.turnId, @@ -546,6 +551,7 @@ export function projectEvent( streaming: message.streaming, updatedAt: message.updatedAt, turnId: message.turnId, + ...(message.channel !== undefined ? { channel: message.channel } : {}), ...(message.attachments !== undefined ? { attachments: message.attachments } : {}), diff --git a/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts b/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts index 719191668869..5ff17491c147 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts @@ -5,7 +5,7 @@ import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; import * as Struct from "effect/Struct"; -import { ChatAttachment } from "@t3tools/contracts"; +import { ChatAttachment, OrchestrationMessageChannel } from "@t3tools/contracts"; import { toPersistenceSqlError } from "../Errors.ts"; import { @@ -20,6 +20,7 @@ import { const ProjectionThreadMessageDbRowSchema = ProjectionThreadMessage.mapFields( Struct.assign({ isStreaming: Schema.Number, + channel: Schema.NullOr(OrchestrationMessageChannel), attachments: Schema.NullOr(Schema.fromJsonString(Schema.Array(ChatAttachment))), }), ); @@ -32,6 +33,7 @@ function toProjectionThreadMessage( threadId: row.threadId, turnId: row.turnId, role: row.role, + ...(row.channel !== null ? { channel: row.channel } : {}), text: row.text, isStreaming: row.isStreaming === 1, createdAt: row.createdAt, @@ -54,6 +56,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { thread_id, turn_id, role, + channel, text, attachments_json, is_streaming, @@ -65,6 +68,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { ${row.threadId}, ${row.turnId}, ${row.role}, + ${row.channel ?? null}, ${row.text}, COALESCE( ${nextAttachmentsJson}, @@ -83,6 +87,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { thread_id = excluded.thread_id, turn_id = excluded.turn_id, role = excluded.role, + channel = COALESCE(excluded.channel, projection_thread_messages.channel), text = excluded.text, attachments_json = COALESCE( excluded.attachments_json, @@ -105,6 +110,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { thread_id AS "threadId", turn_id AS "turnId", role, + channel, text, attachments_json AS "attachments", is_streaming AS "isStreaming", @@ -126,6 +132,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { thread_id AS "threadId", turn_id AS "turnId", role, + channel, text, attachments_json AS "attachments", is_streaming AS "isStreaming", diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index 8abbe87fce3e..287b15dbd63b 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_AuthSessionClientConnection.ts"; import Migration0042 from "./Migrations/042_ProjectionThreadLinkedPullRequest.ts"; import Migration0043 from "./Migrations/043_ProjectionThreadsUnsettledAt.ts"; +import Migration0044 from "./Migrations/044_ProjectionThreadMessagesChannel.ts"; /** * Migration loader with all migrations defined inline. @@ -111,6 +112,7 @@ export const migrationEntries = [ [41, "AuthSessionClientConnection", Migration0041], [42, "ProjectionThreadLinkedPullRequest", Migration0042], [43, "ProjectionThreadsUnsettledAt", Migration0043], + [44, "ProjectionThreadMessagesChannel", Migration0044], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/044_ProjectionThreadMessagesChannel.test.ts b/apps/server/src/persistence/Migrations/044_ProjectionThreadMessagesChannel.test.ts new file mode 100644 index 000000000000..c8990f65502a --- /dev/null +++ b/apps/server/src/persistence/Migrations/044_ProjectionThreadMessagesChannel.test.ts @@ -0,0 +1,43 @@ +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"; +import Migration0044 from "./044_ProjectionThreadMessagesChannel.ts"; + +const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layerMemory())); + +const channelColumns = Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const columns = yield* sql<{ readonly name: string; readonly notnull: number }>` + PRAGMA table_info(projection_thread_messages) + `; + return columns.filter((column) => column.name === "channel"); +}); + +layer("044_ProjectionThreadMessagesChannel", (it) => { + it.effect("adds the nullable channel column to message projections", () => + Effect.gen(function* () { + yield* runMigrations({ toMigrationInclusive: 43 }); + yield* runMigrations({ toMigrationInclusive: 44 }); + + const columns = yield* channelColumns; + + assert.equal(columns.length, 1); + assert.equal(columns[0]?.notnull, 0); + }), + ); + + it.effect("is a no-op when the column already exists", () => + Effect.gen(function* () { + yield* runMigrations({ toMigrationInclusive: 44 }); + yield* Migration0044; + + const columns = yield* channelColumns; + + assert.equal(columns.length, 1); + }), + ); +}); diff --git a/apps/server/src/persistence/Migrations/044_ProjectionThreadMessagesChannel.ts b/apps/server/src/persistence/Migrations/044_ProjectionThreadMessagesChannel.ts new file mode 100644 index 000000000000..478fc21f2bec --- /dev/null +++ b/apps/server/src/persistence/Migrations/044_ProjectionThreadMessagesChannel.ts @@ -0,0 +1,16 @@ +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 columns = yield* sql<{ readonly name: string }>` + PRAGMA table_info(projection_thread_messages) + `; + + if (!columns.some((column) => column.name === "channel")) { + yield* sql` + ALTER TABLE projection_thread_messages + ADD COLUMN channel TEXT + `; + } +}); diff --git a/apps/server/src/persistence/Services/ProjectionThreadMessages.ts b/apps/server/src/persistence/Services/ProjectionThreadMessages.ts index d50ff3202563..5c401d4af03a 100644 --- a/apps/server/src/persistence/Services/ProjectionThreadMessages.ts +++ b/apps/server/src/persistence/Services/ProjectionThreadMessages.ts @@ -10,6 +10,7 @@ import { ChatAttachment, MessageId, OrchestrationMessageRole, + OrchestrationMessageChannel, ThreadId, TurnId, IsoDateTime, @@ -26,6 +27,7 @@ export const ProjectionThreadMessage = Schema.Struct({ threadId: ThreadId, turnId: Schema.NullOr(TurnId), role: OrchestrationMessageRole, + channel: Schema.optional(OrchestrationMessageChannel), text: Schema.String, attachments: Schema.optional(Schema.Array(ChatAttachment)), isStreaming: Schema.Boolean, diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index 2f0efeac5f53..5a70132a4499 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -665,6 +665,50 @@ describe("ClaudeAdapterLive", () => { ); }); + it.effect("requests summarized thinking display by default", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + const createInput = harness.getLastCreateQueryInput(); + assert.deepEqual(createInput?.options.thinking, { + type: "adaptive", + display: "summarized", + }); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("omits thinking query options when the Claude thinking toggle is off", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + modelSelection: createModelSelection( + ProviderInstanceId.make("claudeAgent"), + "claude-haiku-4-5", + [{ id: "thinking", value: false }], + ), + runtimeMode: "full-access", + }); + + const createInput = harness.getLastCreateQueryInput(); + assert.equal(Object.hasOwn(createInput?.options ?? {}, "thinking"), false); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + it.effect("forwards claude fast mode into SDK settings", () => { const harness = makeHarness(); return Effect.gen(function* () { diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index 6989378d8287..3476d0ebdf45 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -4309,6 +4309,8 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( ...(existingResumeSessionId ? { resume: existingResumeSessionId } : {}), ...(newSessionId ? { sessionId: newSessionId } : {}), includePartialMessages: true, + // Claude otherwise streams only redacted thinking token estimates. + ...(thinking === false ? {} : { thinking: { type: "adaptive", display: "summarized" } }), canUseTool, onUserDialog, supportedDialogKinds: ["resume_return"], diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts index 6a6cec5b1e61..7a1148692a5c 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts @@ -120,6 +120,7 @@ describe("buildTurnStartParams", () => { ], model: "gpt-5.3-codex", effort: "medium", + summary: "detailed", collaborationMode: { mode: "plan", settings: { @@ -169,6 +170,7 @@ describe("buildTurnStartParams", () => { }, ], model: "gpt-5.3-codex", + summary: "detailed", collaborationMode: { mode: "default", settings: { @@ -220,6 +222,7 @@ describe("buildTurnStartParams", () => { text: "Ship it", }, ], + summary: "detailed", }); }), ); @@ -246,6 +249,7 @@ describe("buildTurnStartParams", () => { text: "Review", }, ], + summary: "detailed", }); }); }); diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index d83489763f5c..340f190fca70 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -635,6 +635,8 @@ export function buildTurnStartParams(input: { ...(input.model ? { model: input.model } : {}), ...(input.serviceTier ? { serviceTier: input.serviceTier } : {}), ...(input.effort ? { effort: input.effort } : {}), + // Readable reasoning summaries; without this Codex emits empty reasoning items. + summary: "detailed", ...(collaborationMode ? { collaborationMode } : {}), }).pipe( Effect.mapError((cause) => diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts index 3cd8d338d361..150f9ec96f29 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts @@ -4,10 +4,35 @@ import { computeMessageDurationStart, deriveMessagesTimelineRows, normalizeCompactToolLabel, + resolveReasoningDisclosureExpanded, resolveAssistantMessageCopyState, shouldPreserveAssistantLineBreaks, + toggleReasoningDisclosureExpansion, } from "./MessagesTimeline.logic"; +describe("reasoning disclosure expansion", () => { + it("defaults to the message streaming state when the user has not toggled it", () => { + const overrides = new Map(); + const messageId = "reasoning-1" as never; + + expect(resolveReasoningDisclosureExpanded(overrides, messageId, true)).toBe(true); + expect(resolveReasoningDisclosureExpanded(overrides, messageId, false)).toBe(false); + }); + + it("persists an explicit toggle across virtualization and stream completion", () => { + const messageId = "reasoning-1" as never; + const initial = new Map(); + const collapsed = toggleReasoningDisclosureExpansion(initial, messageId, true); + + expect(collapsed).not.toBe(initial); + expect(resolveReasoningDisclosureExpanded(initial, messageId, true)).toBe(true); + expect(resolveReasoningDisclosureExpanded(collapsed, messageId, true)).toBe(false); + + const expanded = toggleReasoningDisclosureExpansion(collapsed, messageId, true); + expect(resolveReasoningDisclosureExpanded(expanded, messageId, false)).toBe(true); + }); +}); + describe("shouldPreserveAssistantLineBreaks", () => { it("preserves Claude insight formatting without changing regular markdown", () => { expect( @@ -93,7 +118,10 @@ describe("computeMessageDurationStart", () => { ); }); - it("does not advance the boundary for a streaming message", () => { + it.each([ + { label: "streaming", streaming: true }, + { label: "completed reasoning", streaming: false, channel: "reasoning" as const }, + ])("does not advance the boundary for a $label message", ({ streaming, channel }) => { const result = computeMessageDurationStart([ { id: "u1", @@ -105,9 +133,10 @@ describe("computeMessageDurationStart", () => { { id: "a1", role: "assistant", + ...(channel ? { channel } : {}), createdAt: "2026-01-01T00:00:30Z", updatedAt: "2026-01-01T00:00:40Z", - streaming: true, + streaming, }, { id: "a2", @@ -627,6 +656,493 @@ describe("deriveMessagesTimelineRows", () => { expect(rows.map((row) => row.id)).toEqual(["turn-fold:turn-1", "assistant-final-entry"]); }); + it("anchors the turn fold at leading reasoning and folds it with the rest of the turn", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "reasoning-entry", + kind: "message", + createdAt: "2026-01-01T00:00:01Z", + message: { + id: "reasoning" as never, + role: "assistant", + channel: "reasoning", + text: "Thinking", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:01Z", + updatedAt: "2026-01-01T00:00:02Z", + streaming: false, + }, + }, + { + id: "assistant-first-entry", + kind: "message", + createdAt: "2026-01-01T00:00:03Z", + message: { + id: "assistant-first" as never, + role: "assistant", + text: "I am checking the implementation.", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:03Z", + updatedAt: "2026-01-01T00:00:04Z", + streaming: false, + }, + }, + { + id: "work-entry", + kind: "work", + createdAt: "2026-01-01T00:00:05Z", + entry: { + id: "work-1", + createdAt: "2026-01-01T00:00:05Z", + turnId: "turn-1" as never, + label: "Read files", + tone: "tool", + }, + }, + { + id: "assistant-final-entry", + kind: "message", + createdAt: "2026-01-01T00:00:06Z", + message: { + id: "assistant-final" as never, + role: "assistant", + text: "Done", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:06Z", + updatedAt: "2026-01-01T00:00:07Z", + streaming: false, + }, + }, + ], + isWorking: false, + activeTurnStartedAt: null, + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows.map((row) => row.id)).toEqual(["turn-fold:turn-1", "assistant-final-entry"]); + expect(rows.find((row) => row.kind === "turn-fold")).toMatchObject({ + createdAt: "2026-01-01T00:00:01Z", + }); + }); + + it("does not create a turn fold for an empty completed reasoning message", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "assistant-final-entry", + kind: "message", + createdAt: "2026-01-01T00:00:05Z", + message: { + id: "assistant-final" as never, + role: "assistant", + text: "Done", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:05Z", + updatedAt: "2026-01-01T00:00:06Z", + streaming: false, + }, + }, + { + id: "reasoning-empty-entry", + kind: "message", + createdAt: "2026-01-01T00:00:07Z", + message: { + id: "reasoning-empty" as never, + role: "assistant", + channel: "reasoning", + text: "", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:07Z", + updatedAt: "2026-01-01T00:00:08Z", + streaming: false, + }, + }, + ], + isWorking: false, + activeTurnStartedAt: null, + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows.map((row) => row.id)).toEqual(["assistant-final-entry"]); + }); + + it("groups work across an invisible empty completed reasoning entry", () => { + // A whitespace-only completed reasoning message renders no row, so it must + // not split adjacent work entries into two work-toggle groups either. + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "work-entry-1", + kind: "work", + createdAt: "2026-01-01T00:00:01Z", + entry: { + id: "work-1", + createdAt: "2026-01-01T00:00:01Z", + label: "Ran command", + tone: "tool" as const, + itemType: "command_execution" as const, + toolLifecycleStatus: "completed" as const, + }, + }, + { + id: "reasoning-empty-entry", + kind: "message", + createdAt: "2026-01-01T00:00:02Z", + message: { + id: "reasoning-empty" as never, + role: "assistant", + channel: "reasoning", + text: " ", + turnId: null, + createdAt: "2026-01-01T00:00:02Z", + updatedAt: "2026-01-01T00:00:02Z", + streaming: false, + }, + }, + { + id: "work-entry-2", + kind: "work", + createdAt: "2026-01-01T00:00:03Z", + entry: { + id: "work-2", + createdAt: "2026-01-01T00:00:03Z", + label: "Ran command", + tone: "tool" as const, + itemType: "command_execution" as const, + toolLifecycleStatus: "completed" as const, + }, + }, + ], + isWorking: false, + activeTurnStartedAt: null, + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows.map((row) => row.id)).toEqual(["work-toggle:work-entry-1"]); + expect(rows.find((row) => row.kind === "work-toggle")).toMatchObject({ + hiddenCount: 2, + summary: "Ran 2 commands", + }); + }); + + it("keeps the live work group when an empty completed reasoning entry trails it", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "latest-command-entry", + kind: "work", + createdAt: "2026-01-01T00:00:05Z", + entry: { + id: "latest-command", + createdAt: "2026-01-01T00:00:05Z", + turnId: "turn-1" as never, + label: "Ran rg", + command: "rg toolCall", + requestKind: "command", + tone: "tool" as const, + toolLifecycleStatus: "completed" as const, + }, + }, + { + id: "reasoning-empty-entry", + kind: "message", + createdAt: "2026-01-01T00:00:06Z", + message: { + id: "reasoning-empty" as never, + role: "assistant", + channel: "reasoning", + text: " ", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:06Z", + updatedAt: "2026-01-01T00:00:06Z", + streaming: false, + }, + }, + ], + latestTurn: { + turnId: "turn-1" as never, + state: "running", + startedAt: "2026-01-01T00:00:00Z", + completedAt: null, + }, + isWorking: true, + activeTurnStartedAt: "2026-01-01T00:00:00Z", + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows.map((row) => row.kind)).toEqual(["working", "work-live"]); + expect(rows.find((row) => row.kind === "work-live")).toMatchObject({ + entry: { id: "latest-command" }, + groupedEntries: [{ id: "latest-command" }], + }); + }); + + it("keeps the Thinking label when filtered empty reasoning precedes the active-turn header", () => { + // Empty completed reasoning entries are filtered before indices are + // computed, so they must not shift active-turn membership onto prior-turn + // entries and hide the Thinking label. + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "reasoning-empty-entry-1", + kind: "message", + createdAt: "2026-01-01T00:00:01Z", + message: { + id: "reasoning-empty-1" as never, + role: "assistant", + channel: "reasoning", + text: "", + turnId: "turn-0" as never, + createdAt: "2026-01-01T00:00:01Z", + updatedAt: "2026-01-01T00:00:01Z", + streaming: false, + }, + }, + { + id: "reasoning-empty-entry-2", + kind: "message", + createdAt: "2026-01-01T00:00:02Z", + message: { + id: "reasoning-empty-2" as never, + role: "assistant", + channel: "reasoning", + text: " ", + turnId: "turn-0" as never, + createdAt: "2026-01-01T00:00:02Z", + updatedAt: "2026-01-01T00:00:02Z", + streaming: false, + }, + }, + { + id: "assistant-prior-entry", + kind: "message", + createdAt: "2026-01-01T00:00:03Z", + message: { + id: "assistant-prior" as never, + role: "assistant", + text: "Earlier answer", + turnId: "turn-0" as never, + createdAt: "2026-01-01T00:00:03Z", + updatedAt: "2026-01-01T00:00:04Z", + streaming: false, + }, + }, + { + id: "user-entry", + kind: "message", + createdAt: "2026-01-01T00:00:05Z", + message: { + id: "user-1" as never, + role: "user", + text: "Next question", + turnId: null, + createdAt: "2026-01-01T00:00:05Z", + updatedAt: "2026-01-01T00:00:05Z", + streaming: false, + }, + }, + { + id: "reasoning-streaming-entry", + kind: "message", + createdAt: "2026-01-01T00:00:06Z", + message: { + id: "reasoning-streaming" as never, + role: "assistant", + channel: "reasoning", + text: "", + turnId: null, + createdAt: "2026-01-01T00:00:06Z", + updatedAt: "2026-01-01T00:00:06Z", + streaming: true, + }, + }, + ], + isWorking: true, + activeTurnStartedAt: "2026-01-01T00:00:05Z", + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + const workingRow = rows.find((row) => row.kind === "working"); + expect(workingRow).toMatchObject({ showThinking: true }); + }); + + it("shows the working indicator when trailing filtered empty reasoning ends the timeline", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "user-entry", + kind: "message", + createdAt: "2026-01-01T00:00:00Z", + message: { + id: "user-1" as never, + role: "user", + text: "Go", + turnId: null, + createdAt: "2026-01-01T00:00:00Z", + updatedAt: "2026-01-01T00:00:00Z", + streaming: false, + }, + }, + { + id: "reasoning-empty-entry", + kind: "message", + createdAt: "2026-01-01T00:00:01Z", + message: { + id: "reasoning-empty" as never, + role: "assistant", + channel: "reasoning", + text: "", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:01Z", + updatedAt: "2026-01-01T00:00:01Z", + streaming: false, + }, + }, + ], + latestTurn: { + turnId: "turn-1" as never, + state: "running", + startedAt: "2026-01-01T00:00:00Z", + completedAt: null, + }, + isWorking: true, + activeTurnStartedAt: "2026-01-01T00:00:00Z", + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows.map((row) => row.id)).toEqual(["user-entry", "working-indicator-row"]); + expect(rows.find((row) => row.kind === "working")).toMatchObject({ showThinking: true }); + }); + + it("restores the Thinking label once the active turn's only reasoning message completes", () => { + // Regression: activeTurnHasVisibleContent treats any non-empty assistant + // message as live visible content, including a *completed* reasoning + // message. Once the thought finishes and the turn has no substantive + // assistant text and no in-progress tool, the working row is the only + // signal the agent is still busy, so it must show the generic Thinking + // label again. + const makeInput = (reasoningStreaming: boolean) => ({ + timelineEntries: [ + { + id: "user-entry", + kind: "message" as const, + createdAt: "2026-01-01T00:00:00Z", + message: { + id: "user-1" as never, + role: "user" as const, + text: "Go", + turnId: null, + createdAt: "2026-01-01T00:00:00Z", + updatedAt: "2026-01-01T00:00:00Z", + streaming: false, + }, + }, + { + id: "reasoning-entry", + kind: "message" as const, + createdAt: "2026-01-01T00:00:01Z", + message: { + id: "reasoning-1" as never, + role: "assistant" as const, + channel: "reasoning" as const, + text: "Weighing the options", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:01Z", + updatedAt: "2026-01-01T00:00:01Z", + streaming: reasoningStreaming, + }, + }, + ], + latestTurn: { + turnId: "turn-1" as never, + state: "running" as const, + startedAt: "2026-01-01T00:00:00Z", + completedAt: null, + }, + isWorking: true, + activeTurnStartedAt: "2026-01-01T00:00:00Z", + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + // While the reasoning text still streams it is itself a live Thinking + // disclosure, so the redundant generic label stays suppressed. + const streamingRows = deriveMessagesTimelineRows(makeInput(true)); + expect(streamingRows.find((row) => row.kind === "working")).toMatchObject({ + showThinking: false, + }); + + const completedRows = deriveMessagesTimelineRows(makeInput(false)); + expect(completedRows.find((row) => row.kind === "working")).toMatchObject({ + showThinking: true, + }); + }); + + it("still splits work groups on a non-empty completed reasoning entry", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "work-entry-1", + kind: "work", + createdAt: "2026-01-01T00:00:01Z", + entry: { + id: "work-1", + createdAt: "2026-01-01T00:00:01Z", + label: "Ran command", + tone: "tool" as const, + itemType: "command_execution" as const, + toolLifecycleStatus: "completed" as const, + }, + }, + { + id: "reasoning-entry", + kind: "message", + createdAt: "2026-01-01T00:00:02Z", + message: { + id: "reasoning" as never, + role: "assistant", + channel: "reasoning", + text: "Weighing the next step.", + turnId: null, + createdAt: "2026-01-01T00:00:02Z", + updatedAt: "2026-01-01T00:00:02Z", + streaming: false, + }, + }, + { + id: "work-entry-2", + kind: "work", + createdAt: "2026-01-01T00:00:03Z", + entry: { + id: "work-2", + createdAt: "2026-01-01T00:00:03Z", + label: "Ran command", + tone: "tool" as const, + itemType: "command_execution" as const, + toolLifecycleStatus: "completed" as const, + }, + }, + ], + isWorking: false, + activeTurnStartedAt: null, + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows.map((row) => row.id)).toEqual([ + "work-toggle:work-entry-1", + "reasoning-entry", + "work-toggle:work-entry-2", + ]); + }); + it("derives a sane duration for a steer-superseded turn with one instant commentary message", () => { // A steer ends the previous turn early: its only message completes the // instant it is created, and trailing work entries land after it. The @@ -1322,6 +1838,21 @@ describe("deriveMessagesTimelineRows", () => { streaming: false, }, }, + { + id: "reasoning-tail-entry", + kind: "message", + createdAt: "2026-01-01T00:00:40Z", + message: { + id: "reasoning-tail" as never, + role: "assistant", + channel: "reasoning", + text: "Considering a follow-up.", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:40Z", + updatedAt: "2026-01-01T00:00:45Z", + streaming: false, + }, + }, ], expandedTurnIds: new Set(["turn-1" as never]), isWorking: false, @@ -1335,7 +1866,8 @@ describe("deriveMessagesTimelineRows", () => { row.kind === "message" && row.message.role === "assistant", ); - expect(assistantRows.map((row) => row.showAssistantMeta)).toEqual([false, true]); + expect(assistantRows.map((row) => row.showAssistantMeta)).toEqual([false, true, false]); + expect(assistantRows.map((row) => row.showAssistantCopyButton)).toEqual([false, true, false]); }); it("withholds assistant metadata while the active turn is still in progress", () => { diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.ts b/apps/web/src/components/chat/MessagesTimeline.logic.ts index c190643ee7f3..914ebb4ad8c8 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.ts @@ -9,7 +9,13 @@ import { type WorkLogEntry, } from "../../session-logic"; import { type ChatMessage, type ProposedPlan, type TurnDiffSummary } from "../../types"; -import { type MessageId, type OrchestrationLatestTurn, type TurnId } from "@t3tools/contracts"; +import { + isReasoningMessage, + type MessageId, + type OrchestrationLatestTurn, + type OrchestrationMessageChannel, + type TurnId, +} from "@t3tools/contracts"; export const TIMELINE_MINIMAP_ITEM_SPACING = 8; export const TIMELINE_MINIMAP_MIN_ITEMS = 2; @@ -68,6 +74,24 @@ export function shouldPreserveAssistantLineBreaks(text: string): boolean { return /^★ Insight(?:\s|─)/mu.test(text); } +export function resolveReasoningDisclosureExpanded( + overrides: ReadonlyMap, + messageId: MessageId, + streaming: boolean, +): boolean { + return overrides.get(messageId) ?? streaming; +} + +export function toggleReasoningDisclosureExpansion( + overrides: ReadonlyMap, + messageId: MessageId, + defaultExpanded: boolean, +): ReadonlyMap { + const next = new Map(overrides); + next.set(messageId, !(overrides.get(messageId) ?? defaultExpanded)); + return next; +} + export function resolveTimelineMinimapHeightStyle(itemCount: number): string { const naturalHeight = Math.max(1, (itemCount - 1) * TIMELINE_MINIMAP_ITEM_SPACING); return `min(${naturalHeight}px, ${TIMELINE_MINIMAP_MAX_HEIGHT_CSS})`; @@ -166,6 +190,7 @@ function maxIsoTimestamp(a: string | null, b: string | null): string | null { export interface TimelineDurationMessage { id: string; role: "user" | "assistant" | "system"; + channel?: OrchestrationMessageChannel | undefined; createdAt: string; updatedAt: string; streaming: boolean; @@ -254,7 +279,7 @@ export function computeMessageDurationStart( lastBoundary = message.createdAt; } result.set(message.id, lastBoundary ?? message.createdAt); - if (message.role === "assistant" && !message.streaming) { + if (message.role === "assistant" && !isReasoningMessage(message) && !message.streaming) { lastBoundary = message.updatedAt; } } @@ -453,7 +478,7 @@ function deriveTerminalAssistantMessageIds(timelineEntries: ReadonlyArray; }): MessagesTimelineRow[] { const nextRows: MessagesTimelineRow[] = []; + // A whitespace-only completed reasoning message renders no row, so every + // derivation below (folds, forward/backward work-group scans, the row loop) + // must see the same list without it — otherwise an invisible entry breaks + // work-group adjacency. + const timelineEntries = input.timelineEntries.filter( + (entry) => !isEmptyCompletedReasoningEntry(entry), + ); const durationStartByMessageId = computeMessageDurationStart( - input.timelineEntries.flatMap((entry) => (entry.kind === "message" ? [entry.message] : [])), + timelineEntries.flatMap((entry) => (entry.kind === "message" ? [entry.message] : [])), ); - const terminalAssistantMessageIds = deriveTerminalAssistantMessageIds(input.timelineEntries); + const terminalAssistantMessageIds = deriveTerminalAssistantMessageIds(timelineEntries); const unsettledTurnId = deriveUnsettledTurnId( input.latestTurn ?? null, input.runningTurnId ?? null, ); const foldsByAnchorEntryId = deriveTurnFolds({ - timelineEntries: input.timelineEntries, + timelineEntries, terminalAssistantMessageIds, latestTurn: input.latestTurn ?? null, unsettledTurnId, @@ -680,13 +721,13 @@ export function deriveMessagesTimelineRows(input: { } } - let activeTurnHeaderIndex = input.timelineEntries.length; + let activeTurnHeaderIndex = timelineEntries.length; if (input.isWorking) { - const latestUserMessageIndex = lastUserMessageIndex(input.timelineEntries); + const latestUserMessageIndex = lastUserMessageIndex(timelineEntries); const firstOwnedAfterUser = unsettledTurnId === null ? -1 - : input.timelineEntries.findIndex( + : timelineEntries.findIndex( (entry, index) => index > latestUserMessageIndex && timelineEntryTurnId(entry) === unsettledTurnId, ); @@ -703,11 +744,16 @@ export function deriveMessagesTimelineRows(input: { entry.toolLifecycleStatus === "inProgress" && entry.turnId === unsettledTurnId; const activeEntries = input.isWorking - ? input.timelineEntries.filter((entry, index) => entryBelongsToActiveTurn(entry, index)) + ? timelineEntries.filter((entry, index) => entryBelongsToActiveTurn(entry, index)) : []; const activeTurnHasVisibleContent = activeEntries.some((entry) => { if (entry.kind === "message") { - return entry.message.role === "assistant" && (entry.message.text?.trim().length ?? 0) > 0; + if (entry.message.role !== "assistant" || (entry.message.text?.trim().length ?? 0) === 0) { + return false; + } + // Completed reasoning is a transcript artifact, not live output; it only + // counts as visible content while it is still streaming. + return isReasoningMessage(entry.message) ? entry.message.streaming : true; } if (entry.kind === "work") { return ( @@ -721,8 +767,8 @@ export function deriveMessagesTimelineRows(input: { }); const activeToolEntries: Array> = []; - for (let index = input.timelineEntries.length - 1; index >= activeTurnHeaderIndex; index -= 1) { - const entry = input.timelineEntries[index]!; + for (let index = timelineEntries.length - 1; index >= activeTurnHeaderIndex; index -= 1) { + const entry = timelineEntries[index]!; if ( !entryBelongsToActiveTurn(entry, index) || entry.kind !== "work" || @@ -780,8 +826,8 @@ export function deriveMessagesTimelineRows(input: { } }; - for (let index = 0; index < input.timelineEntries.length; index += 1) { - const timelineEntry = input.timelineEntries[index]; + for (let index = 0; index < timelineEntries.length; index += 1) { + const timelineEntry = timelineEntries[index]; if (!timelineEntry) { continue; } @@ -828,8 +874,8 @@ export function deriveMessagesTimelineRows(input: { } const groupedEntries = [timelineEntry.entry]; let cursor = index + 1; - while (cursor < input.timelineEntries.length) { - const nextEntry = input.timelineEntries[cursor]; + while (cursor < timelineEntries.length) { + const nextEntry = timelineEntries[cursor]; if ( !nextEntry || nextEntry.kind !== "work" || @@ -963,7 +1009,7 @@ export function deriveMessagesTimelineRows(input: { }); } - if (input.isWorking && activeTurnHeaderIndex === input.timelineEntries.length) { + if (input.isWorking && activeTurnHeaderIndex === timelineEntries.length) { appendWorkingRow(); } diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index c654b8707d9b..e2f253e59ff7 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -1,4 +1,5 @@ import { + isReasoningMessage, type ChatFileAttachment, type EnvironmentId, type MessageId, @@ -39,6 +40,7 @@ import { LegendList, type LegendListRef } from "@legendapp/list/react"; import { FileDiff } from "@pierre/diffs/react"; import { deriveTimelineEntries, + formatDuration, workEntryDisplayIndicatesToolFailure, workEntrySignalsSevereFailure, workLogEntryIsToolLike, @@ -91,6 +93,7 @@ import { deriveMessagesTimelineRows, normalizeCompactToolLabel, resolveAssistantMessageCopyState, + resolveReasoningDisclosureExpanded, resolveTimelineIsAtEnd, resolveTimelineMinimapHasPersistentGutter, resolveTimelineMinimapHeightStyle, @@ -99,6 +102,7 @@ import { resolveTimelineMinimapInteractiveWidth, resolveTimelineMinimapTopPercent, shouldPreserveAssistantLineBreaks, + toggleReasoningDisclosureExpansion, toolGroupAction, workEntryIsVisibleInGroup, type StableMessagesTimelineRowsState, @@ -163,6 +167,12 @@ interface TimelineRowSharedState { onOpenTurnDiff: (turnId: TurnId, filePath?: string) => void; onToggleTurnFold: (turnId: TurnId) => void; onToggleWorkGroup: (groupId: string, anchorKey: string) => void; + reasoningExpansionOverrides: ReadonlyMap; + onToggleReasoningDisclosure: ( + messageId: MessageId, + defaultExpanded: boolean, + anchorKey: string, + ) => void; agentPanelModel: AgentPanelModel; onOpenAgents: () => void; } @@ -309,6 +319,9 @@ export const MessagesTimeline = memo(function MessagesTimeline({ }: MessagesTimelineProps) { const [expandedTurnIds, setExpandedTurnIds] = useState>(new Set()); const [expandedWorkGroupIds, setExpandedWorkGroupIds] = useState>(new Set()); + const [reasoningExpansionOverrides, setReasoningExpansionOverrides] = useState< + ReadonlyMap + >(new Map()); const [disclosureToggleSettling, setDisclosureToggleSettling] = useState(false); const [minimapStripMap] = useState(() => new Map()); const disclosureAnchorKeyRef = useRef(null); @@ -400,6 +413,15 @@ export const MessagesTimeline = memo(function MessagesTimeline({ }, [suspendEndScrollMaintenanceForDisclosure], ); + const onToggleReasoningDisclosure = useCallback( + (messageId: MessageId, defaultExpanded: boolean, anchorKey: string) => { + suspendEndScrollMaintenanceForDisclosure(anchorKey); + setReasoningExpansionOverrides((existing) => + toggleReasoningDisclosureExpansion(existing, messageId, defaultExpanded), + ); + }, + [suspendEndScrollMaintenanceForDisclosure], + ); // An in-session interrupt leaves its turn expanded so the user keeps their // place; the next turn (or a reload, since this is local state) folds it. @@ -555,6 +577,8 @@ export const MessagesTimeline = memo(function MessagesTimeline({ onOpenTurnDiff, onToggleTurnFold, onToggleWorkGroup, + reasoningExpansionOverrides, + onToggleReasoningDisclosure, agentPanelModel, onOpenAgents, }), @@ -574,6 +598,8 @@ export const MessagesTimeline = memo(function MessagesTimeline({ onOpenTurnDiff, onToggleTurnFold, onToggleWorkGroup, + reasoningExpansionOverrides, + onToggleReasoningDisclosure, agentPanelModel, onOpenAgents, ], @@ -675,7 +701,9 @@ function keyExtractor(item: MessagesTimelineRow) { } function getItemType(item: MessagesTimelineRow) { - return item.kind === "message" ? `message:${item.message.role}` : item.kind; + return item.kind === "message" + ? `message:${isReasoningMessage(item.message) ? "reasoning" : item.message.role}` + : item.kind; } interface TimelineMinimapItem { @@ -726,7 +754,7 @@ function resolveFinalAssistantTextForTurn( if (row.message.role === "user") { break; } - if (row.message.role === "assistant") { + if (row.message.role === "assistant" && !isReasoningMessage(row.message)) { finalAssistantText = row.message.text ?? null; } } @@ -1004,7 +1032,11 @@ const TimelineRowContent = memo(function TimelineRowContent({ row }: { row: Time {row.kind === "turn-fold" ? : null} {row.kind === "message" && row.message.role === "user" ? : null} {row.kind === "message" && row.message.role === "assistant" ? ( - + isReasoningMessage(row.message) ? ( + + ) : ( + + ) ) : null} {row.kind === "proposed-plan" ? : null} {row.kind === "working" ? : null} @@ -1243,6 +1275,44 @@ function TurnFoldTimelineRow({ row }: { row: Extract }) { + const ctx = use(TimelineRowCtx); + const streaming = row.message.streaming; + const expanded = resolveReasoningDisclosureExpanded( + ctx.reasoningExpansionOverrides, + row.message.id, + streaming, + ); + const Icon = expanded ? ChevronDownIcon : ChevronRightIcon; + const durationMs = Date.parse(row.message.updatedAt) - Date.parse(row.message.createdAt); + const label = streaming + ? "Thinking..." + : Number.isFinite(durationMs) && durationMs >= 1_000 + ? `Thought for ${formatDuration(durationMs)}` + : "Thought"; + const text = row.message.text; + + return ( +
+ + {expanded ? ( +
+ {text} +
+ ) : null} +
+ ); +} + function AssistantTimelineRow({ row }: { row: Extract }) { const ctx = use(TimelineRowCtx); const messageText = row.message.text || (row.message.streaming ? "" : "(empty response)"); diff --git a/docs/README.md b/docs/README.md index 2e2e55fbbadc..efce553d8fb7 100644 --- a/docs/README.md +++ b/docs/README.md @@ -6,6 +6,7 @@ - [Permission modes](./user/permission-modes.md) - [Keyboard shortcuts](./user/keybindings.md) - [Organizing threads](./user/thread-sidebar.md) +- [Reasoning in a thread](./user/reasoning.md) - [Review usage](./user/usage.md) - [Customize a project icon](./user/project-settings.md) - [Mobile appearance](./user/mobile-appearance.md) diff --git a/docs/internals/glossary.md b/docs/internals/glossary.md index c1b4251f91d0..bb1e0b47e920 100644 --- a/docs/internals/glossary.md +++ b/docs/internals/glossary.md @@ -43,6 +43,14 @@ A single user-to-assistant work cycle inside a thread. It starts with user input A user-visible log item attached to a thread. In [the contracts][1], activities cover important non-message events like approvals, tool actions, and failures. They are projected into thread state in [projector.ts][4]. +#### Message channel + +An optional discriminator on a message in [the contracts][1]. No channel means ordinary conversation text; `reasoning` marks provider thinking. Read it through the `isReasoningMessage` helper rather than the raw field, so server and clients branch the same way. + +#### Reasoning message + +An assistant message with `channel: "reasoning"`, holding one burst of provider thinking. It is deliberately inert: [ProjectionPipeline.ts][11] does not let it settle a turn, [CheckpointReactor.ts][6] leaves it out of checkpoints, and [ProviderCommandReactor.ts][12] drops it from thread-title context. Delivery reuses the [assistant delivery mode](#assistant-delivery-mode) chosen in [ProviderRuntimeIngestion.ts][5]. The web timeline renders each one as a collapsible "Thought" row in [MessagesTimeline.tsx][27], excluded from timeline search and the minimap by [MessagesTimeline.logic.ts][28]; mobile filters them out entirely. See [reasoning][29] for the shipped behavior. + ### Orchestration Orchestration is the server-side domain layer that turns runtime activity into stable app state. The main entry point is [OrchestrationEngine.ts][7], with core logic in [decider.ts][8] and [projector.ts][4]. @@ -201,3 +209,6 @@ ships T3 Code already matching it. [24]: ./overview.md [25]: ../../apps/server/src/environmentTheme.ts [26]: ../user/environment-theme.md +[27]: ../../apps/web/src/components/chat/MessagesTimeline.tsx +[28]: ../../apps/web/src/components/chat/MessagesTimeline.logic.ts +[29]: ../user/reasoning.md diff --git a/docs/user/providers-claude.md b/docs/user/providers-claude.md index 2442d3315ab4..ab1f03dfb5b1 100644 --- a/docs/user/providers-claude.md +++ b/docs/user/providers-claude.md @@ -53,6 +53,12 @@ T3 Code looks for Claude skills in the Claude config directory's `skills` folder If the same skill name exists in more than one folder, the later folder wins. +## What Claude's Thoughts Show + +Claude's thinking appears in the thread timeline as **Thought** rows. Recent Claude Code versions +do not hand out raw thinking, so T3 Code asks for the summarized form and that summary is what you +read. See [Reasoning in a thread](./reasoning.md). + ## I Want Work And Personal Claude Accounts Use a different Claude config directory for each account. diff --git a/docs/user/reasoning.md b/docs/user/reasoning.md new file mode 100644 index 000000000000..85fa39320a7c --- /dev/null +++ b/docs/user/reasoning.md @@ -0,0 +1,23 @@ +# Reasoning in a thread + +Some providers report the thinking they do before answering. Each burst of it becomes one +collapsible row in the thread timeline, tucked under the turn's **Worked for …** group with the +tools the agent ran. The row reads **Thinking…** while it arrives and **Thought for 4s** once it +finishes. It collapses on its own when the thinking ends, so the answer stays the thing you see. +Click the row to read it. + +Delivery follows the same setting as assistant text. With **Stream token by token (legacy)** on in +**Settings → General**, thinking streams in as it is produced. With it off, which is the default, +each thought appears as one finished block. + +Reasoning never changes the outcome of a turn. It does not keep a turn from settling, it is not +part of checkpoints or turn diffs, it never becomes a thread title, and thread search and the +timeline's jump markers skip it. + +## Where it shows up + +Claude, Codex, and OpenCode report reasoning. Cursor and Grok do not, and their threads look the +same as before. Thoughts appear on web and desktop. The mobile app does not show them yet. + +For Claude, what you read is a summary of the thinking rather than the model's raw words. Recent +Claude Code versions do not hand out raw thinking, so T3 Code asks for the summarized form. diff --git a/packages/client-runtime/src/state/threadReducer.test.ts b/packages/client-runtime/src/state/threadReducer.test.ts index 2042b2168c88..734d9683b9ab 100644 --- a/packages/client-runtime/src/state/threadReducer.test.ts +++ b/packages/client-runtime/src/state/threadReducer.test.ts @@ -497,6 +497,58 @@ describe("applyThreadDetailEvent", () => { expect(result.thread.latestTurn?.completedAt).toBeNull(); } }); + + it("stores reasoning messages without rebinding answer or checkpoint semantics", () => { + const turnId = TurnId.make("turn-1"); + const thread: OrchestrationThread = { + ...baseThread, + latestTurn: { + turnId, + state: "running", + requestedAt: "2026-04-01T06:59:00.000Z", + startedAt: "2026-04-01T06:59:00.000Z", + completedAt: null, + assistantMessageId: MessageId.make("answer-1"), + }, + checkpoints: [ + { + turnId, + checkpointTurnCount: 1, + checkpointRef: CheckpointRef.make("ref-1"), + status: "ready", + files: [], + assistantMessageId: MessageId.make("answer-1"), + completedAt: "2026-04-01T07:00:00.000Z", + }, + ], + }; + const result = applyThreadDetailEvent(thread, { + ...baseEventFields, + sequence: 9, + occurredAt: "2026-04-01T07:01:00.000Z", + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-1"), + type: "thread.message-sent", + payload: { + threadId: ThreadId.make("thread-1"), + messageId: MessageId.make("reasoning-1"), + role: "assistant", + channel: "reasoning", + text: "Thinking", + turnId, + streaming: false, + createdAt: "2026-04-01T07:01:00.000Z", + updatedAt: "2026-04-01T07:01:00.000Z", + }, + }); + + expect(result.kind).toBe("updated"); + if (result.kind === "updated") { + expect(result.thread.messages.at(-1)?.channel).toBe("reasoning"); + expect(result.thread.latestTurn?.assistantMessageId).toBe("answer-1"); + expect(result.thread.checkpoints[0]?.assistantMessageId).toBe("answer-1"); + } + }); }); describe("thread.session-set", () => { diff --git a/packages/client-runtime/src/state/threadReducer.ts b/packages/client-runtime/src/state/threadReducer.ts index 10d5898c8fe8..8746a33a69dc 100644 --- a/packages/client-runtime/src/state/threadReducer.ts +++ b/packages/client-runtime/src/state/threadReducer.ts @@ -1,16 +1,17 @@ import { pipe } from "effect/Function"; import * as Arr from "effect/Array"; import * as O from "effect/Order"; -import type { - MessageId, - OrchestrationCheckpointSummary, - OrchestrationEvent, - OrchestrationLatestTurn, - OrchestrationMessage, - OrchestrationSession, - OrchestrationThread, - OrchestrationThreadActivity, - TurnId, +import { + isReasoningMessage, + type MessageId, + type OrchestrationCheckpointSummary, + type OrchestrationEvent, + type OrchestrationLatestTurn, + type OrchestrationMessage, + type OrchestrationSession, + type OrchestrationThread, + type OrchestrationThreadActivity, + type TurnId, } from "@t3tools/contracts"; export type ThreadDetailReducerResult = @@ -296,6 +297,7 @@ export function applyThreadDetailEvent( const message: OrchestrationMessage = { id: event.payload.messageId, role: event.payload.role, + ...(event.payload.channel !== undefined ? { channel: event.payload.channel } : {}), text: event.payload.text, ...(event.payload.attachments !== undefined ? { attachments: event.payload.attachments } @@ -319,6 +321,7 @@ export function applyThreadDetailEvent( ? message.text : entry.text, streaming: message.streaming, + ...(message.channel !== undefined ? { channel: message.channel } : {}), ...(message.turnId !== undefined ? { turnId: message.turnId } : {}), ...(message.streaming ? {} : { updatedAt: message.updatedAt }), ...(message.attachments !== undefined @@ -339,6 +342,7 @@ export function applyThreadDetailEvent( const settlesTurn = !event.payload.streaming && !turnStillRunning; const latestTurn: OrchestrationThread["latestTurn"] = event.payload.role === "assistant" && + !isReasoningMessage(event.payload) && event.payload.turnId !== null && (thread.latestTurn === null || thread.latestTurn.turnId === event.payload.turnId) ? { @@ -369,7 +373,9 @@ export function applyThreadDetailEvent( // Rebind checkpoint assistant message IDs for assistant messages. const checkpoints = - event.payload.role === "assistant" && event.payload.turnId !== null + event.payload.role === "assistant" && + !isReasoningMessage(event.payload) && + event.payload.turnId !== null ? rebindCheckpointAssistantMessage( thread.checkpoints, event.payload.turnId, diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 7ca7175ad175..87d4657cac21 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -303,9 +303,13 @@ export type OrchestrationProject = typeof OrchestrationProject.Type; export const OrchestrationMessageRole = Schema.Literals(["user", "assistant", "system"]); export type OrchestrationMessageRole = typeof OrchestrationMessageRole.Type; +export const OrchestrationMessageChannel = Schema.Literal("reasoning"); +export type OrchestrationMessageChannel = typeof OrchestrationMessageChannel.Type; + export const OrchestrationMessage = Schema.Struct({ id: MessageId, role: OrchestrationMessageRole, + channel: Schema.optional(OrchestrationMessageChannel), text: Schema.String, attachments: Schema.optional(Schema.Array(ChatAttachment)), turnId: Schema.NullOr(TurnId), @@ -315,6 +319,12 @@ export const OrchestrationMessage = Schema.Struct({ }); export type OrchestrationMessage = typeof OrchestrationMessage.Type; +export function isReasoningMessage(message: { + readonly channel?: OrchestrationMessageChannel | undefined; +}): boolean { + return message.channel === "reasoning"; +} + export const OrchestrationProposedPlanId = TrimmedNonEmptyString; export type OrchestrationProposedPlanId = typeof OrchestrationProposedPlanId.Type; @@ -1047,6 +1057,7 @@ const ThreadMessageAssistantDeltaCommand = Schema.Struct({ threadId: ThreadId, messageId: MessageId, delta: Schema.String, + channel: Schema.optional(OrchestrationMessageChannel), turnId: Schema.optional(TurnId), createdAt: IsoDateTime, }); @@ -1056,6 +1067,7 @@ const ThreadMessageAssistantCompleteCommand = Schema.Struct({ commandId: CommandId, threadId: ThreadId, messageId: MessageId, + channel: Schema.optional(OrchestrationMessageChannel), turnId: Schema.optional(TurnId), createdAt: IsoDateTime, }); @@ -1306,6 +1318,7 @@ export const ThreadMessageSentPayload = Schema.Struct({ threadId: ThreadId, messageId: MessageId, role: OrchestrationMessageRole, + channel: Schema.optional(OrchestrationMessageChannel), text: Schema.String, attachments: Schema.optional(Schema.Array(ChatAttachment)), turnId: Schema.NullOr(TurnId),