diff --git a/apps/server/integration/orchestrationEngine.integration.test.ts b/apps/server/integration/orchestrationEngine.integration.test.ts index 5572b6e3eed9..e12b8d2bd182 100644 --- a/apps/server/integration/orchestrationEngine.integration.test.ts +++ b/apps/server/integration/orchestrationEngine.integration.test.ts @@ -319,7 +319,7 @@ it.live("runs multi-turn file edits and persists checkpoint diffs", () => const secondTurnThread = yield* harness.waitForThread( THREAD_ID, (entry) => - entry.latestTurnId === "turn-2" && + entry.latestTurn?.turnId === "turn-2" && entry.checkpoints.length === 2 && entry.checkpoints.some((checkpoint) => checkpoint.checkpointTurnCount === 2), ); @@ -653,7 +653,7 @@ it.live("reverts to an earlier checkpoint and trims checkpoint projections + git yield* harness.waitForThread( THREAD_ID, (entry) => - entry.latestTurnId === "turn-2" && + entry.latestTurn?.turnId === "turn-2" && entry.checkpoints.length === 2 && entry.activities.some((activity) => activity.turnId === "turn-2"), 8000, diff --git a/apps/server/src/checkpointing/Layers/CheckpointDiffQuery.test.ts b/apps/server/src/checkpointing/Layers/CheckpointDiffQuery.test.ts index e01081ee3fc3..786253d5a47c 100644 --- a/apps/server/src/checkpointing/Layers/CheckpointDiffQuery.test.ts +++ b/apps/server/src/checkpointing/Layers/CheckpointDiffQuery.test.ts @@ -45,7 +45,14 @@ function makeSnapshot(input: { model: "gpt-5-codex", branch: null, worktreePath: input.worktreePath, - latestTurnId: TurnId.makeUnsafe("turn-1"), + latestTurn: { + turnId: TurnId.makeUnsafe("turn-1"), + state: "completed", + requestedAt: "2026-01-01T00:00:00.000Z", + startedAt: "2026-01-01T00:00:00.000Z", + completedAt: "2026-01-01T00:00:00.000Z", + assistantMessageId: null, + }, createdAt: "2026-01-01T00:00:00.000Z", updatedAt: "2026-01-01T00:00:00.000Z", deletedAt: null, diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts index 769310380cca..3da8faf731ee 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts @@ -11,6 +11,7 @@ import { ProjectId, ProviderSessionId, ProviderThreadId, + ProviderTurnId, ThreadId, TurnId, } from "@t3tools/contracts"; @@ -40,6 +41,7 @@ import { checkpointRefForThreadTurn } from "../../checkpointing/Utils.ts"; const asProjectId = (value: string): ProjectId => ProjectId.makeUnsafe(value); const asSessionId = (value: string): ProviderSessionId => ProviderSessionId.makeUnsafe(value); const asProviderThreadId = (value: string): ProviderThreadId => ProviderThreadId.makeUnsafe(value); +const asProviderTurnId = (value: string): ProviderTurnId => ProviderTurnId.makeUnsafe(value); const asTurnId = (value: string): TurnId => TurnId.makeUnsafe(value); function createProviderServiceHarness(cwd: string, hasSession = true, sessionCwd = cwd) { @@ -91,7 +93,7 @@ function createProviderServiceHarness(cwd: string, hasSession = true, sessionCwd async function waitForThread( engine: OrchestrationEngineShape, predicate: (thread: { - latestTurnId: string | null; + latestTurn: { turnId: string } | null; checkpoints: ReadonlyArray<{ checkpointTurnCount: number }>; activities: ReadonlyArray<{ kind: string }>; }) => boolean, @@ -99,7 +101,7 @@ async function waitForThread( ) { const deadline = Date.now() + timeoutMs; const poll = async (): Promise<{ - latestTurnId: string | null; + latestTurn: { turnId: string } | null; checkpoints: ReadonlyArray<{ checkpointTurnCount: number }>; activities: ReadonlyArray<{ kind: string }>; }> => { @@ -334,7 +336,7 @@ describe("CheckpointReactor", () => { sessionId: asSessionId("sess-1"), createdAt: new Date().toISOString(), threadId: ProviderThreadId.makeUnsafe("provider-thread-1"), - turnId: asTurnId("turn-1"), + turnId: asProviderTurnId("turn-1"), }); await waitForGitRefExists( harness.cwd, @@ -349,14 +351,14 @@ describe("CheckpointReactor", () => { sessionId: asSessionId("sess-1"), createdAt: new Date().toISOString(), threadId: ProviderThreadId.makeUnsafe("provider-thread-1"), - turnId: asTurnId("turn-1"), + turnId: asProviderTurnId("turn-1"), status: "completed", }); await waitForEvent(harness.engine, (event) => event.type === "thread.turn-diff-completed"); const thread = await waitForThread( harness.engine, - (entry) => entry.latestTurnId === "turn-1" && entry.checkpoints.length === 1, + (entry) => entry.latestTurn?.turnId === "turn-1" && entry.checkpoints.length === 1, ); expect(thread.checkpoints[0]?.checkpointTurnCount).toBe(1); expect( @@ -381,6 +383,81 @@ describe("CheckpointReactor", () => { ).toBe("v2\n"); }); + it("ignores auxiliary thread turn completion while primary turn is active", async () => { + const harness = await createHarness({ seedFilesystemCheckpoints: false }); + const createdAt = new Date().toISOString(); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.makeUnsafe("cmd-session-set-primary-running"), + threadId: ThreadId.makeUnsafe("thread-1"), + session: { + threadId: ThreadId.makeUnsafe("thread-1"), + status: "running", + providerName: "codex", + providerSessionId: asSessionId("sess-1"), + providerThreadId: ProviderThreadId.makeUnsafe("provider-thread-1"), + approvalPolicy: "on-request", + sandboxMode: "workspace-write", + activeTurnId: asTurnId("turn-main"), + lastError: null, + updatedAt: createdAt, + }, + createdAt, + }), + ); + + harness.provider.emit({ + type: "turn.started", + eventId: EventId.makeUnsafe("evt-turn-started-main"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: new Date().toISOString(), + threadId: ProviderThreadId.makeUnsafe("provider-thread-1"), + turnId: asProviderTurnId("turn-main"), + }); + await waitForGitRefExists( + harness.cwd, + checkpointRefForThreadTurn(ThreadId.makeUnsafe("thread-1"), 0), + ); + + fs.writeFileSync(path.join(harness.cwd, "README.md"), "v2\n", "utf8"); + + harness.provider.emit({ + type: "turn.completed", + eventId: EventId.makeUnsafe("evt-turn-completed-aux"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: new Date().toISOString(), + threadId: ProviderThreadId.makeUnsafe("provider-thread-aux"), + turnId: asProviderTurnId("turn-aux"), + status: "completed", + }); + + await Effect.runPromise(Effect.sleep("40 millis")); + const midReadModel = await Effect.runPromise(harness.engine.getReadModel()); + const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.makeUnsafe("thread-1")); + expect(midThread?.checkpoints).toHaveLength(0); + + harness.provider.emit({ + type: "turn.completed", + eventId: EventId.makeUnsafe("evt-turn-completed-main"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: new Date().toISOString(), + threadId: ProviderThreadId.makeUnsafe("provider-thread-1"), + turnId: asProviderTurnId("turn-main"), + status: "completed", + }); + + const thread = await waitForThread( + harness.engine, + (entry) => entry.latestTurn?.turnId === "turn-main" && entry.checkpoints.length === 1, + ); + expect(thread.checkpoints[0]?.checkpointTurnCount).toBe(1); + }); + it("appends capture failure activity when turn diff summary cannot be derived", async () => { const harness = await createHarness({ seedFilesystemCheckpoints: false }); const createdAt = new Date().toISOString(); @@ -413,7 +490,7 @@ describe("CheckpointReactor", () => { sessionId: asSessionId("sess-1"), createdAt: new Date().toISOString(), threadId: ProviderThreadId.makeUnsafe("provider-thread-1"), - turnId: asTurnId("turn-missing-baseline"), + turnId: asProviderTurnId("turn-missing-baseline"), status: "completed", }); @@ -505,7 +582,7 @@ describe("CheckpointReactor", () => { sessionId: asSessionId("sess-missing"), createdAt: new Date().toISOString(), threadId: ProviderThreadId.makeUnsafe("provider-thread-missing"), - turnId: asTurnId("turn-missing-cwd"), + turnId: asProviderTurnId("turn-missing-cwd"), status: "completed", }); @@ -554,7 +631,7 @@ describe("CheckpointReactor", () => { sessionId: asSessionId("sess-1"), createdAt: new Date().toISOString(), threadId: ProviderThreadId.makeUnsafe("provider-thread-1"), - turnId: asTurnId("turn-3"), + turnId: asProviderTurnId("turn-3"), turnCount: 3, status: "completed", }); @@ -607,7 +684,7 @@ describe("CheckpointReactor", () => { sessionId: asSessionId("sess-1"), createdAt: new Date().toISOString(), threadId: ProviderThreadId.makeUnsafe("provider-thread-1"), - turnId: asTurnId("turn-runtime-failure"), + turnId: asProviderTurnId("turn-runtime-failure"), turnCount: 1, status: "completed", }); @@ -619,7 +696,7 @@ describe("CheckpointReactor", () => { sessionId: asSessionId("sess-1"), createdAt: new Date().toISOString(), threadId: ProviderThreadId.makeUnsafe("provider-thread-1"), - turnId: asTurnId("turn-after-runtime-failure"), + turnId: asProviderTurnId("turn-after-runtime-failure"), }); await waitForGitRefExists( @@ -698,7 +775,7 @@ describe("CheckpointReactor", () => { await waitForEvent(harness.engine, (event) => event.type === "thread.reverted"); const thread = await waitForThread(harness.engine, (entry) => entry.checkpoints.length === 1); - expect(thread.latestTurnId).toBe("turn-1"); + expect(thread.latestTurn?.turnId).toBe("turn-1"); expect(thread.checkpoints).toHaveLength(1); expect(thread.checkpoints[0]?.checkpointTurnCount).toBe(1); expect(harness.provider.rollbackConversation).toHaveBeenCalledTimes(1); diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.ts index cd4eba3d2c61..71b51ba737be 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.ts @@ -3,6 +3,7 @@ import { EventId, MessageId, ProviderSessionId, + ProviderThreadId, ThreadId, TurnId, type OrchestrationEvent, @@ -36,6 +37,17 @@ function toTurnId(value: string | undefined): TurnId | null { return value === undefined ? null : TurnId.makeUnsafe(value); } +function toProviderThreadId(value: string | undefined): ProviderThreadId | null { + return value === undefined ? null : ProviderThreadId.makeUnsafe(value); +} + +function sameId(left: string | null | undefined, right: string | null | undefined): boolean { + if (left === null || left === undefined || right === null || right === undefined) { + return false; + } + return left === right; +} + function checkpointStatusFromRuntime(status: string | undefined): "ready" | "missing" | "error" { switch (status) { case "failed": @@ -173,6 +185,21 @@ const make = Effect.gen(function* () { return; } + const projectedProviderThreadId = thread.session?.providerThreadId ?? null; + const eventProviderThreadId = toProviderThreadId(event.threadId); + if ( + projectedProviderThreadId !== null && + eventProviderThreadId !== null && + !sameId(projectedProviderThreadId, eventProviderThreadId) + ) { + return; + } + + // When a primary turn is active, only that turn may produce completion checkpoints. + if (thread.session?.activeTurnId && !sameId(thread.session.activeTurnId, turnId)) { + return; + } + if (thread.checkpoints.some((checkpoint) => checkpoint.turnId === turnId)) { return; } diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts index 97dadf01cb5a..ba99a796487d 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts @@ -243,7 +243,14 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { model: "gpt-5-codex", branch: null, worktreePath: null, - latestTurnId: asTurnId("turn-1"), + latestTurn: { + turnId: asTurnId("turn-1"), + state: "completed", + requestedAt: "2026-02-24T00:00:08.000Z", + startedAt: "2026-02-24T00:00:08.000Z", + completedAt: "2026-02-24T00:00:08.000Z", + assistantMessageId: asMessageId("message-1"), + }, createdAt: "2026-02-24T00:00:02.000Z", updatedAt: "2026-02-24T00:00:03.000Z", deletedAt: null, diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index fe2f96615cb1..3b0798e20eaf 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -1,8 +1,12 @@ import { + IsoDateTime, + MessageId, OrchestrationCheckpointFile, OrchestrationReadModel, ProjectScript, + TurnId, type OrchestrationCheckpointSummary, + type OrchestrationLatestTurn, type OrchestrationMessage, type OrchestrationProject, type OrchestrationSession, @@ -55,6 +59,15 @@ const ProjectionCheckpointDbRowSchema = ProjectionCheckpoint.mapFields( files: Schema.fromJsonString(Schema.Array(OrchestrationCheckpointFile)), }), ); +const ProjectionLatestTurnDbRowSchema = Schema.Struct({ + threadId: ProjectionThread.fields.threadId, + turnId: TurnId, + state: Schema.String, + requestedAt: IsoDateTime, + startedAt: Schema.NullOr(IsoDateTime), + completedAt: Schema.NullOr(IsoDateTime), + assistantMessageId: Schema.NullOr(MessageId), +}); const ProjectionStateDbRowSchema = ProjectionState; const REQUIRED_SNAPSHOT_PROJECTORS = [ @@ -226,6 +239,25 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { `, }); + const listLatestTurnRows = SqlSchema.findAll({ + Request: Schema.Void, + Result: ProjectionLatestTurnDbRowSchema, + execute: () => + sql` + SELECT + thread_id AS "threadId", + turn_id AS "turnId", + state, + requested_at AS "requestedAt", + started_at AS "startedAt", + completed_at AS "completedAt", + assistant_message_id AS "assistantMessageId" + FROM projection_turns + WHERE turn_id IS NOT NULL + ORDER BY thread_id ASC, requested_at DESC, turn_id DESC + `, + }); + const listProjectionStateRows = SqlSchema.findAll({ Request: Schema.Void, Result: ProjectionStateDbRowSchema, @@ -250,6 +282,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { activityRows, sessionRows, checkpointRows, + latestTurnRows, stateRows, ] = yield* Effect.all([ listProjectRows(undefined).pipe( @@ -300,6 +333,14 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { ), ), ), + listLatestTurnRows(undefined).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionSnapshotQuery.getSnapshot:listLatestTurns:query", + "ProjectionSnapshotQuery.getSnapshot:listLatestTurns:decodeRows", + ), + ), + ), listProjectionStateRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( @@ -314,6 +355,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { const activitiesByThread = new Map>(); const checkpointsByThread = new Map>(); const sessionsByThread = new Map(); + const latestTurnByThread = new Map(); let updatedAt: string | null = null; @@ -372,6 +414,34 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { checkpointsByThread.set(row.threadId, threadCheckpoints); } + for (const row of latestTurnRows) { + updatedAt = maxIso(updatedAt, row.requestedAt); + if (row.startedAt !== null) { + updatedAt = maxIso(updatedAt, row.startedAt); + } + if (row.completedAt !== null) { + updatedAt = maxIso(updatedAt, row.completedAt); + } + if (latestTurnByThread.has(row.threadId)) { + continue; + } + latestTurnByThread.set(row.threadId, { + turnId: row.turnId, + state: + row.state === "error" + ? "error" + : row.state === "interrupted" + ? "interrupted" + : row.state === "completed" + ? "completed" + : "running", + requestedAt: row.requestedAt, + startedAt: row.startedAt, + completedAt: row.completedAt, + assistantMessageId: row.assistantMessageId, + }); + } + for (const row of sessionRows) { updatedAt = maxIso(updatedAt, row.updatedAt); sessionsByThread.set(row.threadId, { @@ -406,7 +476,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { model: row.model, branch: row.branch, worktreePath: row.worktreePath, - latestTurnId: row.latestTurnId, + latestTurn: latestTurnByThread.get(row.threadId) ?? null, createdAt: row.createdAt, updatedAt: row.updatedAt, deletedAt: row.deletedAt, diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 7a3a029e89a0..fa254887d8f4 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -31,6 +31,8 @@ import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeInge const asProjectId = (value: string): ProjectId => ProjectId.makeUnsafe(value); const asSessionId = (value: string): ProviderSessionId => ProviderSessionId.makeUnsafe(value); +const asProviderThreadId = (value: string): ProviderThreadId => + ProviderThreadId.makeUnsafe(value); const asProviderTurnId = (value: string): ProviderTurnId => ProviderTurnId.makeUnsafe(value); const asItemId = (value: string): ProviderItemId => ProviderItemId.makeUnsafe(value); const asEventId = (value: string): EventId => EventId.makeUnsafe(value); @@ -220,6 +222,112 @@ describe("ProviderRuntimeIngestion", () => { expect(thread.session?.lastError).toBe("turn failed"); }); + it("ignores auxiliary turn completions from a different provider thread", async () => { + const harness = await createHarness(); + const now = new Date().toISOString(); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-primary"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: now, + threadId: asProviderThreadId("provider-thread-1"), + turnId: asProviderTurnId("turn-primary"), + }); + + await waitForThread( + harness.engine, + (thread) => + thread.session?.status === "running" && thread.session?.activeTurnId === "turn-primary", + ); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-aux"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: new Date().toISOString(), + threadId: asProviderThreadId("provider-thread-aux"), + turnId: asProviderTurnId("turn-aux"), + status: "completed", + }); + + await Effect.runPromise(Effect.sleep("40 millis")); + const midReadModel = await Effect.runPromise(harness.engine.getReadModel()); + const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.makeUnsafe("thread-1")); + expect(midThread?.session?.status).toBe("running"); + expect(midThread?.session?.activeTurnId).toBe("turn-primary"); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-primary"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: new Date().toISOString(), + threadId: asProviderThreadId("provider-thread-1"), + turnId: asProviderTurnId("turn-primary"), + status: "completed", + }); + + await waitForThread( + harness.engine, + (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, + ); + }); + + it("ignores non-active turn completion when runtime omits thread id", async () => { + const harness = await createHarness(); + const now = new Date().toISOString(); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-guarded"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: now, + turnId: asProviderTurnId("turn-guarded-main"), + }); + + await waitForThread( + harness.engine, + (thread) => + thread.session?.status === "running" && + thread.session?.activeTurnId === "turn-guarded-main", + ); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-guarded-other"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: new Date().toISOString(), + turnId: asProviderTurnId("turn-guarded-other"), + status: "completed", + }); + + await Effect.runPromise(Effect.sleep("40 millis")); + const midReadModel = await Effect.runPromise(harness.engine.getReadModel()); + const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.makeUnsafe("thread-1")); + expect(midThread?.session?.status).toBe("running"); + expect(midThread?.session?.activeTurnId).toBe("turn-guarded-main"); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-guarded-main"), + provider: "codex", + sessionId: asSessionId("sess-1"), + createdAt: new Date().toISOString(), + turnId: asProviderTurnId("turn-guarded-main"), + status: "completed", + }); + + await waitForThread( + harness.engine, + (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, + ); + }); + it("maps message delta/completed into finalized assistant messages", async () => { const harness = await createHarness(); const now = new Date().toISOString(); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 9da54b2a2636..86cb01475688 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -33,6 +33,7 @@ const TURN_MESSAGE_IDS_BY_TURN_TTL = Duration.minutes(120); const BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_CACHE_CAPACITY = 20_000; const BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_TTL = Duration.minutes(120); const MAX_BUFFERED_ASSISTANT_CHARS = 24_000; +const STRICT_PROVIDER_LIFECYCLE_GUARD = process.env.T3CODE_STRICT_PROVIDER_LIFECYCLE_GUARD !== "0"; type TurnStartRequestedDomainEvent = Extract< OrchestrationEvent, @@ -53,6 +54,17 @@ function toTurnId(value: string | undefined): TurnId | undefined { return value === undefined ? undefined : TurnId.makeUnsafe(value); } +function toProviderThreadId(value: string | undefined): ProviderThreadId | null { + return value === undefined ? null : ProviderThreadId.makeUnsafe(value); +} + +function sameId(left: string | null | undefined, right: string | null | undefined): boolean { + if (left === null || left === undefined || right === null || right === undefined) { + return false; + } + return left === right; +} + function truncateDetail(value: string, limit = 180): string { return value.length > limit ? `${value.slice(0, limit - 3)}...` : value; } @@ -332,6 +344,62 @@ const make = Effect.gen(function* () { if (!thread) return; const now = event.createdAt; + const sessionProviderThreadId = thread.session?.providerThreadId ?? null; + const eventProviderThreadId = toProviderThreadId(event.threadId); + const eventTurnId = toTurnId("turnId" in event ? event.turnId : undefined); + const activeTurnId = thread.session?.activeTurnId ?? null; + + const matchesThreadScope = + eventProviderThreadId === null || + sessionProviderThreadId === null || + sameId(eventProviderThreadId, sessionProviderThreadId); + const conflictsWithActiveTurn = + activeTurnId !== null && eventTurnId !== undefined && !sameId(activeTurnId, eventTurnId); + const missingTurnForActiveTurn = activeTurnId !== null && eventTurnId === undefined; + + const shouldApplyThreadLifecycle = (() => { + if (!STRICT_PROVIDER_LIFECYCLE_GUARD) { + return true; + } + switch (event.type) { + case "session.exited": + return true; + case "session.started": + case "thread.started": + if (!matchesThreadScope) { + return false; + } + // Never let auxiliary/provider-side spawned threads replace the primary thread binding. + if ( + eventProviderThreadId !== null && + sessionProviderThreadId !== null && + !sameId(eventProviderThreadId, sessionProviderThreadId) + ) { + return false; + } + return true; + case "turn.started": + if (!matchesThreadScope) { + return false; + } + return !conflictsWithActiveTurn; + case "turn.completed": + if (!matchesThreadScope) { + return false; + } + if (conflictsWithActiveTurn || missingTurnForActiveTurn) { + return false; + } + // Only the active turn may close the lifecycle state. + if (activeTurnId !== null && eventTurnId !== undefined) { + return sameId(activeTurnId, eventTurnId); + } + // Without an active turn, only accept completion when no thread mismatch signal exists. + return eventProviderThreadId === null || sessionProviderThreadId === null; + default: + return true; + } + })(); if ( event.type === "session.started" || @@ -340,8 +408,7 @@ const make = Effect.gen(function* () { event.type === "turn.started" || event.type === "turn.completed" ) { - const activeTurnId = - event.type === "turn.started" ? (toTurnId(event.turnId) ?? null) : null; + const nextActiveTurnId = event.type === "turn.started" ? (eventTurnId ?? null) : null; const providerThreadIdFromEvent = event.type === "thread.started" ? ProviderThreadId.makeUnsafe(event.threadId) @@ -349,7 +416,7 @@ const make = Effect.gen(function* () { ? ProviderThreadId.makeUnsafe(event.threadId) : null; const providerThreadId = - providerThreadIdFromEvent ?? thread.session?.providerThreadId ?? null; + providerThreadIdFromEvent ?? sessionProviderThreadId ?? null; const status = event.type === "turn.started" ? "running" @@ -365,24 +432,26 @@ const make = Effect.gen(function* () { ? null : (thread.session?.lastError ?? null); - yield* orchestrationEngine.dispatch({ - type: "thread.session.set", - commandId: providerCommandId(event, "thread-session-set"), - threadId: thread.id, - session: { + if (shouldApplyThreadLifecycle) { + yield* orchestrationEngine.dispatch({ + type: "thread.session.set", + commandId: providerCommandId(event, "thread-session-set"), threadId: thread.id, - status, - providerName: event.provider, - providerSessionId: event.sessionId, - providerThreadId, - approvalPolicy: thread.session?.approvalPolicy ?? DEFAULT_APPROVAL_POLICY, - sandboxMode: thread.session?.sandboxMode ?? DEFAULT_SANDBOX_MODE, - activeTurnId, - lastError, - updatedAt: now, - }, - createdAt: now, - }); + session: { + threadId: thread.id, + status, + providerName: event.provider, + providerSessionId: event.sessionId, + providerThreadId, + approvalPolicy: thread.session?.approvalPolicy ?? DEFAULT_APPROVAL_POLICY, + sandboxMode: thread.session?.sandboxMode ?? DEFAULT_SANDBOX_MODE, + activeTurnId: nextActiveTurnId, + lastError, + updatedAt: now, + }, + createdAt: now, + }); + } } if (event.type === "message.delta" && event.delta.length > 0) { @@ -470,29 +539,38 @@ const make = Effect.gen(function* () { } if (event.type === "runtime.error") { + const shouldApplyRuntimeError = !STRICT_PROVIDER_LIFECYCLE_GUARD + ? true + : matchesThreadScope && + (activeTurnId === null || + eventTurnId === undefined || + sameId(activeTurnId, eventTurnId)); + const providerThreadId = event.threadId !== undefined ? ProviderThreadId.makeUnsafe(event.threadId) : (thread.session?.providerThreadId ?? null); - yield* orchestrationEngine.dispatch({ - type: "thread.session.set", - commandId: providerCommandId(event, "runtime-error-session-set"), - threadId: thread.id, - session: { + if (shouldApplyRuntimeError) { + yield* orchestrationEngine.dispatch({ + type: "thread.session.set", + commandId: providerCommandId(event, "runtime-error-session-set"), threadId: thread.id, - status: "error", - providerName: event.provider, - providerSessionId: event.sessionId, - providerThreadId, - approvalPolicy: thread.session?.approvalPolicy ?? DEFAULT_APPROVAL_POLICY, - sandboxMode: thread.session?.sandboxMode ?? DEFAULT_SANDBOX_MODE, - activeTurnId: toTurnId(event.turnId) ?? null, - lastError: event.message, - updatedAt: now, - }, - createdAt: now, - }); + session: { + threadId: thread.id, + status: "error", + providerName: event.provider, + providerSessionId: event.sessionId, + providerThreadId, + approvalPolicy: thread.session?.approvalPolicy ?? DEFAULT_APPROVAL_POLICY, + sandboxMode: thread.session?.sandboxMode ?? DEFAULT_SANDBOX_MODE, + activeTurnId: eventTurnId ?? null, + lastError: event.message, + updatedAt: now, + }, + createdAt: now, + }); + } } const activities = runtimeEventToActivities(event); diff --git a/apps/server/src/orchestration/commandInvariants.test.ts b/apps/server/src/orchestration/commandInvariants.test.ts index 9ad953d5b7fb..af64eba008ff 100644 --- a/apps/server/src/orchestration/commandInvariants.test.ts +++ b/apps/server/src/orchestration/commandInvariants.test.ts @@ -54,7 +54,7 @@ const readModel: OrchestrationReadModel = { worktreePath: null, createdAt: now, updatedAt: now, - latestTurnId: null, + latestTurn: null, messages: [], session: null, activities: [], @@ -70,7 +70,7 @@ const readModel: OrchestrationReadModel = { worktreePath: null, createdAt: now, updatedAt: now, - latestTurnId: null, + latestTurn: null, messages: [], session: null, activities: [], diff --git a/apps/server/src/orchestration/projector.test.ts b/apps/server/src/orchestration/projector.test.ts index 56c0c0bed92a..094290d90f3c 100644 --- a/apps/server/src/orchestration/projector.test.ts +++ b/apps/server/src/orchestration/projector.test.ts @@ -75,7 +75,7 @@ describe("orchestration projector", () => { model: "gpt-5-codex", branch: null, worktreePath: null, - latestTurnId: null, + latestTurn: null, createdAt: now, updatedAt: now, deletedAt: null, @@ -207,7 +207,7 @@ describe("orchestration projector", () => { ); const thread = afterRunning.threads[0]; - expect(thread?.latestTurnId).toBe("turn-1"); + expect(thread?.latestTurn?.turnId).toBe("turn-1"); expect(thread?.session?.status).toBe("running"); }); @@ -504,7 +504,7 @@ describe("orchestration projector", () => { thread?.activities.map((activity) => ({ id: activity.id, turnId: activity.turnId })), ).toEqual([{ id: "activity-1", turnId: "turn-1" }]); expect(thread?.checkpoints.map((checkpoint) => checkpoint.checkpointTurnCount)).toEqual([1]); - expect(thread?.latestTurnId).toBe("turn-1"); + expect(thread?.latestTurn?.turnId).toBe("turn-1"); }); it("does not fallback-retain messages tied to removed turn IDs", async () => { diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index 6588505d02b7..f92f7172a4cb 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -26,6 +26,12 @@ type ThreadPatch = Partial>; const MAX_THREAD_MESSAGES = 2_000; const MAX_THREAD_CHECKPOINTS = 500; +function checkpointStatusToLatestTurnState(status: "ready" | "missing" | "error") { + if (status === "error") return "error" as const; + if (status === "missing") return "interrupted" as const; + return "completed" as const; +} + function updateThread( threads: ReadonlyArray, threadId: ThreadId, @@ -220,7 +226,7 @@ export function projectEvent( model: payload.model, branch: payload.branch, worktreePath: payload.worktreePath, - latestTurnId: null, + latestTurn: null, createdAt: payload.createdAt, updatedAt: payload.updatedAt, deletedAt: null, @@ -351,7 +357,26 @@ export function projectEvent( ...nextBase, threads: updateThread(nextBase.threads, payload.threadId, { session, - latestTurnId: session.activeTurnId, + latestTurn: + session.status === "running" && session.activeTurnId !== null + ? { + turnId: session.activeTurnId, + state: "running", + requestedAt: + thread.latestTurn?.turnId === session.activeTurnId + ? thread.latestTurn.requestedAt + : session.updatedAt, + startedAt: + thread.latestTurn?.turnId === session.activeTurnId + ? (thread.latestTurn.startedAt ?? session.updatedAt) + : session.updatedAt, + completedAt: null, + assistantMessageId: + thread.latestTurn?.turnId === session.activeTurnId + ? thread.latestTurn.assistantMessageId + : null, + } + : thread.latestTurn, updatedAt: event.occurredAt, }), }; @@ -396,7 +421,20 @@ export function projectEvent( ...nextBase, threads: updateThread(nextBase.threads, payload.threadId, { checkpoints, - latestTurnId: payload.turnId, + latestTurn: { + turnId: payload.turnId, + state: checkpointStatusToLatestTurnState(payload.status), + requestedAt: + thread.latestTurn?.turnId === payload.turnId + ? thread.latestTurn.requestedAt + : payload.completedAt, + startedAt: + thread.latestTurn?.turnId === payload.turnId + ? (thread.latestTurn.startedAt ?? payload.completedAt) + : payload.completedAt, + completedAt: payload.completedAt, + assistantMessageId: payload.assistantMessageId, + }, updatedAt: event.occurredAt, }), }; @@ -422,8 +460,18 @@ export function projectEvent( ).slice(-MAX_THREAD_MESSAGES); const activities = retainThreadActivitiesAfterRevert(thread.activities, retainedTurnIds); - const latestTurnId = - checkpoints.length > 0 ? (checkpoints[checkpoints.length - 1]?.turnId ?? null) : null; + const latestCheckpoint = checkpoints.at(-1) ?? null; + const latestTurn = + latestCheckpoint === null + ? null + : { + turnId: latestCheckpoint.turnId, + state: checkpointStatusToLatestTurnState(latestCheckpoint.status), + requestedAt: latestCheckpoint.completedAt, + startedAt: latestCheckpoint.completedAt, + completedAt: latestCheckpoint.completedAt, + assistantMessageId: latestCheckpoint.assistantMessageId, + }; return { ...nextBase, @@ -431,7 +479,7 @@ export function projectEvent( checkpoints, messages, activities, - latestTurnId, + latestTurn, updatedAt: event.occurredAt, }), }; diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index 9e87316f7741..d6976479af04 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -47,6 +47,7 @@ import { const asSessionId = (value: string): ProviderSessionId => ProviderSessionId.makeUnsafe(value); const asTurnId = (value: string): ProviderTurnId => ProviderTurnId.makeUnsafe(value); +const asProviderThreadId = (value: string): ProviderThreadId => ProviderThreadId.makeUnsafe(value); const asRequestId = (value: string): ApprovalRequestId => ApprovalRequestId.makeUnsafe(value); const asEventId = (value: string): EventId => EventId.makeUnsafe(value); const asThreadId = (value: string): ThreadId => ThreadId.makeUnsafe(value); @@ -564,7 +565,7 @@ fanout.layer("ProviderServiceLive fanout", (it) => { provider: "codex", sessionId: session.sessionId, createdAt: new Date().toISOString(), - threadId: ThreadId.makeUnsafe("thread-1"), + threadId: asProviderThreadId("thread-1"), turnId: asTurnId("turn-1"), status: "completed", }; @@ -612,7 +613,7 @@ fanout.layer("ProviderServiceLive fanout", (it) => { provider: "codex", sessionId: session.sessionId, createdAt: new Date().toISOString(), - threadId: ThreadId.makeUnsafe("thread-1"), + threadId: asProviderThreadId("thread-1"), turnId: asTurnId("turn-1"), toolKind: "command", title: "Command run", @@ -624,7 +625,7 @@ fanout.layer("ProviderServiceLive fanout", (it) => { provider: "codex", sessionId: session.sessionId, createdAt: new Date().toISOString(), - threadId: ThreadId.makeUnsafe("thread-1"), + threadId: asProviderThreadId("thread-1"), turnId: asTurnId("turn-1"), delta: "hello", }, @@ -634,7 +635,7 @@ fanout.layer("ProviderServiceLive fanout", (it) => { provider: "codex", sessionId: session.sessionId, createdAt: new Date().toISOString(), - threadId: ThreadId.makeUnsafe("thread-1"), + threadId: asProviderThreadId("thread-1"), turnId: asTurnId("turn-1"), status: "completed", }, diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 7ae68e537100..04abede67d41 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -43,7 +43,7 @@ import { deriveTimelineEntries, type PendingApproval, deriveWorkLogEntries, - formatDuration, + hasToolActivityForTurn, formatElapsed, formatTimestamp, } from "../session-logic"; @@ -379,12 +379,13 @@ export default function ChatView({ threadId }: ChatViewProps) { ); const diffOpen = diffSearch.diff === "1"; const activeThreadId = activeThread?.id ?? null; + const activeLatestTurn = activeThread?.latestTurn ?? null; const activeProject = state.projects.find((p) => p.id === activeThread?.projectId); useEffect(() => { if (!activeThread?.id) return; - if (!activeThread.latestTurnCompletedAt) return; - const turnCompletedAt = Date.parse(activeThread.latestTurnCompletedAt); + if (!activeLatestTurn?.completedAt) return; + const turnCompletedAt = Date.parse(activeLatestTurn.completedAt); if (Number.isNaN(turnCompletedAt)) return; const lastVisitedAt = activeThread.lastVisitedAt ? Date.parse(activeThread.lastVisitedAt) : NaN; if (!Number.isNaN(lastVisitedAt) && lastVisitedAt >= turnCompletedAt) return; @@ -396,7 +397,7 @@ export default function ChatView({ threadId }: ChatViewProps) { }, [ activeThread?.id, activeThread?.lastVisitedAt, - activeThread?.latestTurnCompletedAt, + activeLatestTurn?.completedAt, dispatch, ]); @@ -410,18 +411,12 @@ export default function ChatView({ threadId }: ChatViewProps) { const nowIso = new Date(nowTick).toISOString(); const threadActivities = activeThread?.activities ?? []; const workLogEntries = useMemo( - () => deriveWorkLogEntries(threadActivities, undefined), - [threadActivities], + () => deriveWorkLogEntries(threadActivities, activeLatestTurn?.turnId ?? undefined), + [activeLatestTurn?.turnId, threadActivities], ); const latestTurnHasToolActivity = useMemo(() => { - const latestTurnId = activeThread?.latestTurnId; - if (!latestTurnId) { - return false; - } - return threadActivities.some( - (activity) => activity.turnId === latestTurnId && activity.tone === "tool", - ); - }, [activeThread?.latestTurnId, threadActivities]); + return hasToolActivityForTurn(threadActivities, activeLatestTurn?.turnId); + }, [activeLatestTurn?.turnId, threadActivities]); const pendingApprovals = useMemo( () => derivePendingApprovals(threadActivities), [threadActivities], @@ -486,36 +481,24 @@ export default function ChatView({ threadId }: ChatViewProps) { }, [inferredCheckpointTurnCountByTurnId, timelineEntries, turnDiffSummaryByAssistantMessageId]); const completionSummary = useMemo(() => { - if (!activeThread?.latestTurnStartedAt) return null; - if (!activeThread.latestTurnCompletedAt) return null; + if (!activeLatestTurn?.startedAt) return null; + if (!activeLatestTurn.completedAt) return null; if (!latestTurnHasToolActivity) return null; - if ( - typeof activeThread.latestTurnDurationMs === "number" && - Number.isFinite(activeThread.latestTurnDurationMs) && - activeThread.latestTurnDurationMs >= 0 - ) { - return `Worked for ${formatDuration(activeThread.latestTurnDurationMs)}`; - } - - const elapsed = formatElapsed( - activeThread.latestTurnStartedAt, - activeThread.latestTurnCompletedAt, - ); + const elapsed = formatElapsed(activeLatestTurn.startedAt, activeLatestTurn.completedAt); return elapsed ? `Worked for ${elapsed}` : null; }, [ - activeThread?.latestTurnStartedAt, - activeThread?.latestTurnCompletedAt, - activeThread?.latestTurnDurationMs, + activeLatestTurn?.completedAt, + activeLatestTurn?.startedAt, latestTurnHasToolActivity, ]); const completionDividerBeforeEntryId = useMemo(() => { - if (!activeThread?.latestTurnStartedAt) return null; - if (!activeThread.latestTurnCompletedAt) return null; + if (!activeLatestTurn?.startedAt) return null; + if (!activeLatestTurn.completedAt) return null; if (!completionSummary) return null; - const turnStartedAt = Date.parse(activeThread.latestTurnStartedAt); - const turnCompletedAt = Date.parse(activeThread.latestTurnCompletedAt); + const turnStartedAt = Date.parse(activeLatestTurn.startedAt); + const turnCompletedAt = Date.parse(activeLatestTurn.completedAt); if (Number.isNaN(turnStartedAt)) return null; if (Number.isNaN(turnCompletedAt)) return null; @@ -533,8 +516,8 @@ export default function ChatView({ threadId }: ChatViewProps) { } return inRangeMatch ?? fallbackMatch; }, [ - activeThread?.latestTurnCompletedAt, - activeThread?.latestTurnStartedAt, + activeLatestTurn?.completedAt, + activeLatestTurn?.startedAt, completionSummary, timelineEntries, ]); diff --git a/apps/web/src/components/Sidebar.tsx b/apps/web/src/components/Sidebar.tsx index 0f05fd689ce5..b348836e1754 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -63,8 +63,8 @@ interface TerminalStatusIndicator { } function hasUnseenCompletion(thread: Thread): boolean { - if (!thread.latestTurnCompletedAt) return false; - const completedAt = Date.parse(thread.latestTurnCompletedAt); + if (!thread.latestTurn?.completedAt) return false; + const completedAt = Date.parse(thread.latestTurn.completedAt); if (Number.isNaN(completedAt)) return false; if (!thread.lastVisitedAt) return true; diff --git a/apps/web/src/session-logic.test.ts b/apps/web/src/session-logic.test.ts index 3377268d2a35..cafaea5f91dc 100644 --- a/apps/web/src/session-logic.test.ts +++ b/apps/web/src/session-logic.test.ts @@ -1,7 +1,11 @@ import { EventId, TurnId, type OrchestrationThreadActivity } from "@t3tools/contracts"; import { describe, expect, it } from "vitest"; -import { derivePendingApprovals, deriveWorkLogEntries } from "./session-logic"; +import { + derivePendingApprovals, + deriveWorkLogEntries, + hasToolActivityForTurn, +} from "./session-logic"; function makeActivity(overrides: { id?: string; @@ -126,3 +130,24 @@ describe("deriveWorkLogEntries", () => { expect(entries.map((entry) => entry.id)).toEqual(["tool-complete"]); }); }); + +describe("hasToolActivityForTurn", () => { + it("returns false when turn id is missing", () => { + const activities: OrchestrationThreadActivity[] = [ + makeActivity({ id: "tool-1", turnId: "turn-1", kind: "tool.completed", tone: "tool" }), + ]; + + expect(hasToolActivityForTurn(activities, undefined)).toBe(false); + expect(hasToolActivityForTurn(activities, null)).toBe(false); + }); + + it("returns true only for matching tool activity in the target turn", () => { + const activities: OrchestrationThreadActivity[] = [ + makeActivity({ id: "tool-1", turnId: "turn-1", kind: "tool.completed", tone: "tool" }), + makeActivity({ id: "info-1", turnId: "turn-2", kind: "turn.completed", tone: "info" }), + ]; + + expect(hasToolActivityForTurn(activities, TurnId.makeUnsafe("turn-1"))).toBe(true); + expect(hasToolActivityForTurn(activities, TurnId.makeUnsafe("turn-2"))).toBe(false); + }); +}); diff --git a/apps/web/src/session-logic.ts b/apps/web/src/session-logic.ts index 114c2272e8de..4a4f6025145d 100644 --- a/apps/web/src/session-logic.ts +++ b/apps/web/src/session-logic.ts @@ -147,6 +147,14 @@ export function deriveWorkLogEntries( }); } +export function hasToolActivityForTurn( + activities: ReadonlyArray, + turnId: TurnId | null | undefined, +): boolean { + if (!turnId) return false; + return activities.some((activity) => activity.turnId === turnId && activity.tone === "tool"); +} + export function deriveTimelineEntries( messages: ChatMessage[], workEntries: WorkLogEntry[], diff --git a/apps/web/src/store.test.ts b/apps/web/src/store.test.ts index 4a3a73a3c004..13b86329599e 100644 --- a/apps/web/src/store.test.ts +++ b/apps/web/src/store.test.ts @@ -1,4 +1,4 @@ -import { ProjectId, ThreadId } from "@t3tools/contracts"; +import { ProjectId, ThreadId, TurnId } from "@t3tools/contracts"; import { describe, expect, it } from "vitest"; import { reducer, type AppState } from "./store"; @@ -29,6 +29,7 @@ function makeThread(overrides: Partial = {}): Thread { activities: [], error: null, createdAt: "2026-02-13T00:00:00.000Z", + latestTurn: null, branch: null, worktreePath: null, ...overrides, @@ -58,7 +59,14 @@ describe("store reducer", () => { const latestTurnCompletedAt = "2026-02-25T12:30:00.000Z"; const initialState = makeState( makeThread({ - latestTurnCompletedAt, + latestTurn: { + turnId: TurnId.makeUnsafe("turn-1"), + state: "completed", + requestedAt: "2026-02-25T12:28:00.000Z", + startedAt: "2026-02-25T12:28:30.000Z", + completedAt: latestTurnCompletedAt, + assistantMessageId: null, + }, lastVisitedAt: "2026-02-25T12:35:00.000Z", }), ); @@ -79,7 +87,7 @@ describe("store reducer", () => { it("does not change a thread without a completed turn", () => { const initialState = makeState( makeThread({ - latestTurnCompletedAt: undefined, + latestTurn: null, lastVisitedAt: "2026-02-25T12:35:00.000Z", }), ); diff --git a/apps/web/src/store.ts b/apps/web/src/store.ts index 8a9af566a4f1..bce86d0e66ab 100644 --- a/apps/web/src/store.ts +++ b/apps/web/src/store.ts @@ -414,19 +414,6 @@ export function reducer(state: AppState, action: Action): AppState { .filter((thread) => thread.deletedAt === null) .map((thread) => { const existing = existingThreadById.get(thread.id); - const latestTurnStartedAt = - thread.latestTurnId === null - ? undefined - : thread.messages.find( - (message) => message.turnId === thread.latestTurnId && message.role === "user", - )?.createdAt; - const latestTurnCompletedAt = thread.checkpoints.find( - (checkpoint) => checkpoint.turnId === thread.latestTurnId, - )?.completedAt; - const latestTurnDurationMs = - latestTurnStartedAt && latestTurnCompletedAt - ? Math.max(0, Date.parse(latestTurnCompletedAt) - Date.parse(latestTurnStartedAt)) - : undefined; return normalizeThreadTerminals({ id: thread.id, @@ -484,10 +471,7 @@ export function reducer(state: AppState, action: Action): AppState { }), error: thread.session?.lastError ?? null, createdAt: thread.createdAt, - latestTurnId: thread.latestTurnId ?? undefined, - latestTurnStartedAt, - latestTurnCompletedAt, - latestTurnDurationMs, + latestTurn: thread.latestTurn, lastVisitedAt: existing?.lastVisitedAt ?? thread.updatedAt, branch: thread.branch, worktreePath: thread.worktreePath, @@ -537,10 +521,10 @@ export function reducer(state: AppState, action: Action): AppState { return { ...state, threads: updateThread(state.threads, action.threadId, (thread) => { - if (!thread.latestTurnCompletedAt) { + if (!thread.latestTurn?.completedAt) { return thread; } - const latestTurnCompletedAtMs = Date.parse(thread.latestTurnCompletedAt); + const latestTurnCompletedAtMs = Date.parse(thread.latestTurn.completedAt); if (Number.isNaN(latestTurnCompletedAtMs)) { return thread; } diff --git a/apps/web/src/types.ts b/apps/web/src/types.ts index 3eca157f091b..0ae07e20c416 100644 --- a/apps/web/src/types.ts +++ b/apps/web/src/types.ts @@ -1,4 +1,5 @@ import type { + OrchestrationLatestTurn, OrchestrationSessionStatus, OrchestrationThreadActivity, ProjectScript as ContractProjectScript, @@ -89,10 +90,7 @@ export interface Thread { messages: ChatMessage[]; error: string | null; createdAt: string; - latestTurnId?: TurnId | undefined; - latestTurnStartedAt?: string | undefined; - latestTurnCompletedAt?: string | undefined; - latestTurnDurationMs?: number | undefined; + latestTurn: OrchestrationLatestTurn | null; lastVisitedAt?: string | undefined; branch: string | null; worktreePath: string | null; diff --git a/apps/web/src/worktreeCleanup.test.ts b/apps/web/src/worktreeCleanup.test.ts index b6f9711bf7ee..2d997f2f8dd3 100644 --- a/apps/web/src/worktreeCleanup.test.ts +++ b/apps/web/src/worktreeCleanup.test.ts @@ -29,6 +29,7 @@ function makeThread(overrides: Partial = {}): Thread { activities: [], error: null, createdAt: "2026-02-13T00:00:00.000Z", + latestTurn: null, branch: null, worktreePath: null, ...overrides, diff --git a/docs/adr/0001-provider-runtime-lifecycle-ownership.md b/docs/adr/0001-provider-runtime-lifecycle-ownership.md new file mode 100644 index 000000000000..1c5afa5093a0 --- /dev/null +++ b/docs/adr/0001-provider-runtime-lifecycle-ownership.md @@ -0,0 +1,40 @@ +# ADR 0001: Provider Runtime Lifecycle Ownership + +## Status +Accepted + +## Context +Provider runtime streams can include auxiliary work (for example collab/child-agent turns) under the same provider session as a user-visible primary turn. + +The previous orchestration ingestion and checkpoint flows keyed lifecycle updates by `providerSessionId` alone. This allowed auxiliary `turn.completed` events to: + +- mark thread sessions as `ready` before the primary turn completed +- clear `activeTurnId` for a still-running primary turn +- emit checkpoint turn-diff summaries for auxiliary turns in the main thread timeline + +## Decision +Define and enforce a lifecycle ownership invariant: + +- Only events in the primary thread/turn lane may mutate thread session lifecycle state (`status`, `activeTurnId`, `providerThreadId`) and checkpoint turn completion. +- Auxiliary events may still append messages and activities. + +Current enforcement is implemented with scope guards: + +- Runtime ingestion only applies lifecycle transitions when the event targets the active primary scope (provider thread and active turn checks). +- Checkpoint capture ignores `turn.completed` events from non-primary provider thread scope, and ignores non-active turn completions while a primary turn is active. + +## Consequences + +### Positive +- Prevents premature loss of "working" state from auxiliary turn completion. +- Prevents auxiliary changed-files/checkpoint cards from appearing as if they were primary turn completion. +- Keeps useful auxiliary activity/message visibility. + +### Negative +- Runtime events that omit thread/turn identifiers are treated conservatively in lifecycle transitions. +- Full explicit lane modeling is still desirable in contracts for long-term clarity. + +## Follow-up +1. Add explicit runtime lane metadata (`primary` vs `auxiliary`) to canonical provider runtime events. +2. Persist orchestration-to-provider turn binding (`TurnId` <-> `ProviderTurnId`) for stronger reconciliation and replay safety. +3. Promote lifecycle ownership checks from heuristic guards to schema-level invariants once lane metadata is available. diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 06fcbfec3ea1..1de117eb66b5 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -191,6 +191,24 @@ export const OrchestrationThreadActivity = Schema.Struct({ }); export type OrchestrationThreadActivity = typeof OrchestrationThreadActivity.Type; +export const OrchestrationLatestTurnState = Schema.Literals([ + "running", + "interrupted", + "completed", + "error", +]); +export type OrchestrationLatestTurnState = typeof OrchestrationLatestTurnState.Type; + +export const OrchestrationLatestTurn = Schema.Struct({ + turnId: TurnId, + state: OrchestrationLatestTurnState, + requestedAt: IsoDateTime, + startedAt: Schema.NullOr(IsoDateTime), + completedAt: Schema.NullOr(IsoDateTime), + assistantMessageId: Schema.NullOr(MessageId), +}); +export type OrchestrationLatestTurn = typeof OrchestrationLatestTurn.Type; + export const OrchestrationThread = Schema.Struct({ id: ThreadId, projectId: ProjectId, @@ -198,7 +216,7 @@ export const OrchestrationThread = Schema.Struct({ model: TrimmedNonEmptyString, branch: Schema.NullOr(TrimmedNonEmptyString), worktreePath: Schema.NullOr(TrimmedNonEmptyString), - latestTurnId: Schema.NullOr(TurnId), + latestTurn: Schema.NullOr(OrchestrationLatestTurn), createdAt: IsoDateTime, updatedAt: IsoDateTime, deletedAt: Schema.NullOr(IsoDateTime), diff --git a/packages/contracts/src/providerRuntime.ts b/packages/contracts/src/providerRuntime.ts index daf70550425e..1544fda1cc11 100644 --- a/packages/contracts/src/providerRuntime.ts +++ b/packages/contracts/src/providerRuntime.ts @@ -9,15 +9,11 @@ import { ProviderSessionId, ProviderThreadId, ProviderTurnId, - ThreadId, - TurnId, IsoDateTime, } from "./baseSchemas"; import { ProviderApprovalDecision, ProviderKind, ProviderRequestKind } from "./orchestration"; const TrimmedNonEmptyStringSchema = TrimmedNonEmptyString; -const RuntimeThreadIdSchema = Schema.Union([ThreadId, ProviderThreadId]); -const RuntimeTurnIdSchema = Schema.Union([TurnId, ProviderTurnId]); export const ProviderRuntimeToolKind = Schema.Union([ProviderRequestKind, Schema.Literal("other")]); export type ProviderRuntimeToolKind = typeof ProviderRuntimeToolKind.Type; @@ -36,7 +32,7 @@ export const ProviderRuntimeSessionStartedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), + threadId: Schema.optional(ProviderThreadId), message: Schema.optional(TrimmedNonEmptyStringSchema), }); export type ProviderRuntimeSessionStartedEvent = typeof ProviderRuntimeSessionStartedEvent.Type; @@ -47,7 +43,7 @@ export const ProviderRuntimeSessionExitedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), + threadId: Schema.optional(ProviderThreadId), message: Schema.optional(TrimmedNonEmptyStringSchema), }); export type ProviderRuntimeSessionExitedEvent = typeof ProviderRuntimeSessionExitedEvent.Type; @@ -58,7 +54,7 @@ export const ProviderRuntimeThreadStartedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: RuntimeThreadIdSchema, + threadId: ProviderThreadId, }); export type ProviderRuntimeThreadStartedEvent = typeof ProviderRuntimeThreadStartedEvent.Type; @@ -68,8 +64,8 @@ export const ProviderRuntimeTurnStartedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: RuntimeTurnIdSchema, + threadId: Schema.optional(ProviderThreadId), + turnId: ProviderTurnId, }); export type ProviderRuntimeTurnStartedEvent = typeof ProviderRuntimeTurnStartedEvent.Type; @@ -79,8 +75,8 @@ export const ProviderRuntimeTurnCompletedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: Schema.optional(ProviderThreadId), + turnId: Schema.optional(ProviderTurnId), status: Schema.optional(ProviderRuntimeTurnStatus), errorMessage: Schema.optional(TrimmedNonEmptyStringSchema), }); @@ -92,8 +88,8 @@ export const ProviderRuntimeMessageDeltaEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: Schema.optional(ProviderThreadId), + turnId: Schema.optional(ProviderTurnId), itemId: Schema.optional(ProviderItemId), delta: Schema.String, }); @@ -106,8 +102,8 @@ export const ProviderRuntimeMessageCompletedEvent = Schema.Struct({ sessionId: ProviderSessionId, createdAt: IsoDateTime, itemId: ProviderItemId, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: Schema.optional(ProviderThreadId), + turnId: Schema.optional(ProviderTurnId), }); export type ProviderRuntimeMessageCompletedEvent = typeof ProviderRuntimeMessageCompletedEvent.Type; @@ -117,8 +113,8 @@ export const ProviderRuntimeToolStartedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: Schema.optional(ProviderThreadId), + turnId: Schema.optional(ProviderTurnId), itemId: Schema.optional(ProviderItemId), toolKind: ProviderRuntimeToolKind, title: TrimmedNonEmptyStringSchema, @@ -132,8 +128,8 @@ export const ProviderRuntimeToolCompletedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: Schema.optional(ProviderThreadId), + turnId: Schema.optional(ProviderTurnId), itemId: Schema.optional(ProviderItemId), toolKind: ProviderRuntimeToolKind, title: TrimmedNonEmptyStringSchema, @@ -147,8 +143,8 @@ export const ProviderRuntimeApprovalRequestedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: Schema.optional(ProviderThreadId), + turnId: Schema.optional(ProviderTurnId), itemId: Schema.optional(ProviderItemId), requestId: ApprovalRequestId, requestKind: ProviderRequestKind, @@ -163,8 +159,8 @@ export const ProviderRuntimeApprovalResolvedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: Schema.optional(ProviderThreadId), + turnId: Schema.optional(ProviderTurnId), itemId: Schema.optional(ProviderItemId), requestId: ApprovalRequestId, requestKind: Schema.optional(ProviderRequestKind), @@ -178,8 +174,8 @@ export const ProviderRuntimeCheckpointCapturedEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: RuntimeThreadIdSchema, - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: ProviderThreadId, + turnId: Schema.optional(ProviderTurnId), turnCount: NonNegativeInt, status: Schema.optional(ProviderRuntimeTurnStatus), }); @@ -192,8 +188,8 @@ export const ProviderRuntimeErrorEvent = Schema.Struct({ provider: ProviderKind, sessionId: ProviderSessionId, createdAt: IsoDateTime, - threadId: Schema.optional(RuntimeThreadIdSchema), - turnId: Schema.optional(RuntimeTurnIdSchema), + threadId: Schema.optional(ProviderThreadId), + turnId: Schema.optional(ProviderTurnId), itemId: Schema.optional(ProviderItemId), message: TrimmedNonEmptyStringSchema, });