diff --git a/apps/mobile/src/features/threads/ThreadFeed.tsx b/apps/mobile/src/features/threads/ThreadFeed.tsx index 4e7a89f77d0e..be5b0c4a8895 100644 --- a/apps/mobile/src/features/threads/ThreadFeed.tsx +++ b/apps/mobile/src/features/threads/ThreadFeed.tsx @@ -107,7 +107,7 @@ import { useAppearanceCodeSurface } from "../settings/appearance/useAppearanceCo import { markdownFileIconSource } from "@t3tools/mobile-markdown-text/file-icons"; import { resolveMarkdownLinkPresentation } from "@t3tools/mobile-markdown-text/links"; import { - deriveThreadFeedPresentation, + deriveThreadFeedPresentationState, type ThreadFeedEntry, type ThreadFeedLatestTurn, } from "../../lib/threadActivity"; @@ -190,6 +190,18 @@ function isFreshTimestamp(input: string): boolean { return Number.isFinite(timestamp) && Date.now() - timestamp < FRESH_ENTRY_WINDOW_MS; } +function haveSameStringSet(left: ReadonlySet, right: ReadonlySet): boolean { + if (left.size !== right.size) { + return false; + } + for (const value of left) { + if (!right.has(value)) { + return false; + } + } + return true; +} + export interface ThreadFeedProps { readonly environmentId: EnvironmentId; readonly threadId: ThreadId; @@ -1042,6 +1054,7 @@ function renderFeedEntry( > & { readonly copiedRowId: string | null; readonly expandedWorkRows: Record; + readonly settledTurnOpeningAssistantMessageIds: ReadonlySet; readonly terminalAssistantMessageIds: ReadonlySet; readonly unsettledTurnId: TurnId | null; readonly onCopyWorkRow: (rowId: string, value: string) => void; @@ -1119,7 +1132,8 @@ function renderFeedEntry( message.turnId === props.unsettledTurnId; const showAssistantMeta = message.role === "assistant" && - props.terminalAssistantMessageIds.has(message.id) && + (props.terminalAssistantMessageIds.has(message.id) || + props.settledTurnOpeningAssistantMessageIds.has(message.id)) && !assistantTurnStillInProgress && !message.streaming; @@ -2039,6 +2053,7 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { const headerMaterialVisibleRef = useRef(false); const lastContentInsetReportRef = useRef(null); const previousLatestTurnRef = useRef(props.latestTurn); + const settledTurnOpeningAssistantMessageIdsRef = useRef>(new Set()); const userScrollSettleTimerRef = useRef | null>(null); const { width: windowWidth } = useWindowDimensions(); const { appearance } = useAppearancePreferences(); @@ -2449,9 +2464,9 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { } return ids; }, [expandedWorkGroups]); - const presentedFeed = useMemo( + const presentationState = useMemo( () => - deriveThreadFeedPresentation( + deriveThreadFeedPresentationState( props.feed, props.latestTurn, expandedTurnIds, @@ -2466,6 +2481,22 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { props.latestTurn, ], ); + const presentedFeed = presentationState.entries; + const derivedSettledTurnOpeningAssistantMessageIds = + presentationState.settledTurnOpeningAssistantMessageIds; + if ( + !haveSameStringSet( + settledTurnOpeningAssistantMessageIdsRef.current, + derivedSettledTurnOpeningAssistantMessageIds, + ) + ) { + settledTurnOpeningAssistantMessageIdsRef.current = derivedSettledTurnOpeningAssistantMessageIds; + } + const settledTurnOpeningAssistantMessageIds = settledTurnOpeningAssistantMessageIdsRef.current; + const feedAppearanceData = useMemo( + () => ({ listAppearanceData, settledTurnOpeningAssistantMessageIds }), + [listAppearanceData, settledTurnOpeningAssistantMessageIds], + ); // The empty↔filled key below remounts the list, which resets its imperative // content-inset override. Re-report the animated inset before the fresh @@ -2731,6 +2762,7 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { steerPendingMessageIds: props.steerPendingMessageIds, copiedRowId, expandedWorkRows, + settledTurnOpeningAssistantMessageIds, terminalAssistantMessageIds, unsettledTurnId, onCopyWorkRow, @@ -2751,6 +2783,7 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { [ copiedRowId, expandedWorkRows, + settledTurnOpeningAssistantMessageIds, terminalAssistantMessageIds, unsettledTurnId, iconSubtleColor, @@ -2856,7 +2889,7 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { } maintainVisibleContentPosition={maintainVisibleContentPosition} data={presentedFeed} - extraData={listAppearanceData} + extraData={feedAppearanceData} renderItem={renderItem} keyExtractor={(entry) => entry.id} getItemType={(entry) => diff --git a/apps/mobile/src/lib/threadActivity.test.ts b/apps/mobile/src/lib/threadActivity.test.ts index c5910d237950..10b31c4e0ce2 100644 --- a/apps/mobile/src/lib/threadActivity.test.ts +++ b/apps/mobile/src/lib/threadActivity.test.ts @@ -15,6 +15,7 @@ import { buildPendingUserInputAnswers, buildThreadFeed, deriveThreadFeedPresentation, + deriveThreadFeedPresentationState, isPendingUserInputOptionSelected, setPendingUserInputCustomAnswer, togglePendingUserInputOptionSelection, @@ -445,6 +446,185 @@ describe("buildThreadFeed", () => { ]); }); + it("keeps a settled turn fold anchored when an older page prepends entries", () => { + const turnId = TurnId.make("turn-windowed-fold"); + const thread = makeThread({ + id: ThreadId.make("thread-windowed-fold"), + projectId: ProjectId.make("project-1"), + title: "Windowed fold", + latestTurn: { + turnId, + state: "completed", + requestedAt: "2026-04-01T00:00:00.000Z", + startedAt: "2026-04-01T00:00:01.000Z", + completedAt: "2026-04-01T00:00:10.000Z", + assistantMessageId: MessageId.make("assistant-final"), + }, + messages: [ + { + id: MessageId.make("assistant-first"), + role: "assistant", + text: "Starting.", + turnId, + streaming: false, + createdAt: "2026-04-01T00:00:02.000Z", + updatedAt: "2026-04-01T00:00:02.000Z", + }, + { + id: MessageId.make("assistant-middle"), + role: "assistant", + text: "Still working.", + turnId, + streaming: false, + createdAt: "2026-04-01T00:00:06.000Z", + updatedAt: "2026-04-01T00:00:06.000Z", + }, + { + id: MessageId.make("assistant-final"), + role: "assistant", + text: "Done.", + turnId, + streaming: false, + createdAt: "2026-04-01T00:00:10.000Z", + updatedAt: "2026-04-01T00:00:10.000Z", + }, + ], + activities: [ + makeActivity({ + id: EventId.make("work-early"), + kind: "tool.completed", + tone: "tool", + summary: "Early work", + createdAt: "2026-04-01T00:00:04.000Z", + turnId, + payload: { title: "Early work", itemType: "file_read", status: "completed" }, + }), + makeActivity({ + id: EventId.make("work-late"), + kind: "tool.completed", + tone: "tool", + summary: "Late work", + createdAt: "2026-04-01T00:00:08.000Z", + turnId, + payload: { title: "Late work", itemType: "file_read", status: "completed" }, + }), + ], + }); + + const initialWindow = buildThreadFeed(thread, { loadedMessages: thread.messages.slice(1) }); + const prependedWindow = buildThreadFeed(thread, { loadedMessages: thread.messages }); + const initialFold = deriveThreadFeedPresentation( + initialWindow, + thread.latestTurn, + new Set(), + ).find((entry) => entry.type === "turn-fold"); + const prependedFold = deriveThreadFeedPresentation( + prependedWindow, + thread.latestTurn, + new Set(), + ).find((entry) => entry.type === "turn-fold"); + + expect(initialFold).toMatchObject({ + id: "turn-fold:turn-windowed-fold", + }); + expect(prependedFold).toMatchObject({ + id: initialFold?.id, + createdAt: initialFold?.createdAt, + }); + }); + + it("derives metadata controls for settled opening assistant messages", () => { + const settledTurnId = TurnId.make("turn-settled"); + const streamingTurnId = TurnId.make("turn-streaming"); + const unsettledTurnId = TurnId.make("turn-unsettled"); + const thread = makeThread({ + id: ThreadId.make("thread-opening-meta"), + projectId: ProjectId.make("project-1"), + title: "Opening response metadata", + latestTurn: { + turnId: unsettledTurnId, + state: "running", + requestedAt: "2026-04-01T00:00:10.000Z", + startedAt: "2026-04-01T00:00:10.000Z", + completedAt: null, + assistantMessageId: MessageId.make("unsettled-final"), + }, + messages: [ + { + id: MessageId.make("settled-opening"), + role: "assistant", + text: "Substantive opening.", + turnId: settledTurnId, + streaming: false, + createdAt: "2026-04-01T00:00:01.000Z", + updatedAt: "2026-04-01T00:00:01.000Z", + }, + { + id: MessageId.make("settled-final"), + role: "assistant", + text: "Done.", + turnId: settledTurnId, + streaming: false, + createdAt: "2026-04-01T00:00:02.000Z", + updatedAt: "2026-04-01T00:00:02.000Z", + }, + { + id: MessageId.make("single-response"), + role: "assistant", + text: "Only response.", + turnId: TurnId.make("turn-single"), + streaming: false, + createdAt: "2026-04-01T00:00:03.000Z", + updatedAt: "2026-04-01T00:00:03.000Z", + }, + { + id: MessageId.make("streaming-opening"), + role: "assistant", + text: "Streaming opening.", + turnId: streamingTurnId, + streaming: true, + createdAt: "2026-04-01T00:00:04.000Z", + updatedAt: "2026-04-01T00:00:04.000Z", + }, + { + id: MessageId.make("streaming-final"), + role: "assistant", + text: "Streaming final.", + turnId: streamingTurnId, + streaming: false, + createdAt: "2026-04-01T00:00:05.000Z", + updatedAt: "2026-04-01T00:00:05.000Z", + }, + { + id: MessageId.make("unsettled-opening"), + role: "assistant", + text: "Unsettled opening.", + turnId: unsettledTurnId, + streaming: false, + createdAt: "2026-04-01T00:00:11.000Z", + updatedAt: "2026-04-01T00:00:11.000Z", + }, + { + id: MessageId.make("unsettled-final"), + role: "assistant", + text: "Unsettled final.", + turnId: unsettledTurnId, + streaming: false, + createdAt: "2026-04-01T00:00:12.000Z", + updatedAt: "2026-04-01T00:00:12.000Z", + }, + ], + }); + + const state = deriveThreadFeedPresentationState( + buildThreadFeed(thread), + thread.latestTurn, + new Set(), + ); + + expect([...state.settledTurnOpeningAssistantMessageIds]).toEqual(["settled-opening"]); + }); + it("folds assistant messages between the first and terminal messages", () => { const turnId = TurnId.make("turn-1"); const thread = makeThread({ @@ -555,7 +735,20 @@ describe("buildThreadFeed", () => { ], activities: [ makeActivity({ - id: EventId.make("work-1"), + id: EventId.make("work-before-response"), + kind: "tool.completed", + tone: "tool", + summary: "Read inputs", + createdAt: "2026-04-01T00:00:05.000Z", + turnId: firstTurnId, + payload: { + title: "Read inputs", + itemType: "file_read", + status: "completed", + }, + }), + makeActivity({ + id: EventId.make("work-after-response"), kind: "tool.completed", tone: "tool", summary: "Ran command", @@ -576,6 +769,24 @@ describe("buildThreadFeed", () => { turnId: firstTurnId, label: "Worked for 12s", }); + expect(collapsed.map((entry) => entry.id)).toEqual([ + "user-1", + "turn-fold:turn-1", + "assistant-commentary", + "user-2", + "assistant-next", + ]); + + const expanded = deriveThreadFeedPresentation(feed, thread.latestTurn, new Set([firstTurnId])); + expect(expanded.map((entry) => entry.id)).toEqual([ + "user-1", + "turn-fold:turn-1", + "work-before-response", + "assistant-commentary", + "work-after-response", + "user-2", + "assistant-next", + ]); }); it("keeps an active turn expanded and classifies error-shaped tool output", () => { diff --git a/apps/mobile/src/lib/threadActivity.ts b/apps/mobile/src/lib/threadActivity.ts index 8e3f40ffde61..63f9de57d372 100644 --- a/apps/mobile/src/lib/threadActivity.ts +++ b/apps/mobile/src/lib/threadActivity.ts @@ -1142,18 +1142,23 @@ interface ThreadFeedTurnFold { readonly label: string; } +interface ThreadFeedTurnDerivation { + readonly foldsByAnchorId: ReadonlyMap; + readonly settledTurnOpeningAssistantMessageIds: ReadonlySet; +} + function deriveThreadFeedTurnFolds( feed: ReadonlyArray, latestTurn: ThreadFeedLatestTurn | null, -): ReadonlyMap { +): ThreadFeedTurnDerivation { const firstAssistantMessageIdByTurn = new Map(); const terminalAssistantMessageIdByTurn = new Map(); for (const entry of feed) { if (entry.type === "message" && entry.message.role === "assistant" && entry.message.turnId) { if (!firstAssistantMessageIdByTurn.has(entry.message.turnId)) { - firstAssistantMessageIdByTurn.set(entry.message.turnId, entry.id); + firstAssistantMessageIdByTurn.set(entry.message.turnId, entry.message.id); } - terminalAssistantMessageIdByTurn.set(entry.message.turnId, entry.id); + terminalAssistantMessageIdByTurn.set(entry.message.turnId, entry.message.id); } } @@ -1179,10 +1184,7 @@ function deriveThreadFeedTurnFolds( } let group = groupsByTurnId.get(turnId); if (!group) { - group = { - entries: [], - startBoundary: pendingUserBoundary, - }; + group = { entries: [], startBoundary: pendingUserBoundary }; pendingUserBoundary = null; groupsByTurnId.set(turnId, group); } @@ -1191,17 +1193,31 @@ function deriveThreadFeedTurnFolds( const unsettledTurnId = deriveUnsettledTurnId(latestTurn); const foldsByAnchorId = new Map(); + const settledTurnOpeningAssistantMessageIds = new Set(); for (const [turnId, group] of groupsByTurnId) { const { entries } = group; if (turnId === unsettledTurnId) { continue; } - if (entries.some((entry) => entry.type === "message" && entry.message.streaming)) { + if ( + entries.some( + (entry) => + entry.type === "message" && entry.message.role === "assistant" && entry.message.streaming, + ) + ) { continue; } const firstAssistantMessageId = firstAssistantMessageIdByTurn.get(turnId); const terminalAssistantMessageId = terminalAssistantMessageIdByTurn.get(turnId); + if ( + firstAssistantMessageId && + terminalAssistantMessageId && + firstAssistantMessageId !== terminalAssistantMessageId + ) { + settledTurnOpeningAssistantMessageIds.add(firstAssistantMessageId); + } + const hiddenEntryIds = new Set( entries .filter( @@ -1248,26 +1264,36 @@ function deriveThreadFeedTurnFolds( foldsByAnchorId.set(firstHiddenEntry.id, { turnId, - createdAt: firstHiddenEntry.createdAt, + // Keep insertion in source order, but do not let page prepends change the + // row's timestamp and invalidate the virtualized anchor. + createdAt: terminalEntry?.createdAt ?? lastEntry.createdAt, hiddenEntryIds, label, }); } - return foldsByAnchorId; + return { foldsByAnchorId, settledTurnOpeningAssistantMessageIds }; } -export function deriveThreadFeedPresentation( +export interface ThreadFeedPresentationState { + readonly entries: ThreadFeedEntry[]; + readonly settledTurnOpeningAssistantMessageIds: ReadonlySet; +} + +export function deriveThreadFeedPresentationState( feed: ReadonlyArray, latestTurn: ThreadFeedLatestTurn | null, expandedTurnIds: ReadonlySet, expandedWorkGroupIds: ReadonlySet = new Set(), activeWorkStartedAt: string | null = null, -): ThreadFeedEntry[] { +): ThreadFeedPresentationState { const sourceFeed = feed.filter( - (entry) => + (entry): entry is PresentableThreadFeedEntry => entry.type !== "turn-fold" && entry.type !== "work-toggle" && entry.type !== "working", ); - const foldsByAnchorId = deriveThreadFeedTurnFolds(sourceFeed, latestTurn); + const { foldsByAnchorId, settledTurnOpeningAssistantMessageIds } = deriveThreadFeedTurnFolds( + sourceFeed, + latestTurn, + ); const collapsedEntryIds = new Set(); for (const fold of foldsByAnchorId.values()) { if (!expandedTurnIds.has(fold.turnId)) { @@ -1277,11 +1303,11 @@ export function deriveThreadFeedPresentation( } } - const result: ThreadFeedEntry[] = []; + const entries: ThreadFeedEntry[] = []; for (const entry of sourceFeed) { const fold = foldsByAnchorId.get(entry.id); if (fold) { - result.push({ + entries.push({ type: "turn-fold", id: `turn-fold:${fold.turnId}`, createdAt: fold.createdAt, @@ -1291,22 +1317,43 @@ export function deriveThreadFeedPresentation( }); } if (!collapsedEntryIds.has(entry.id)) { - appendPresentedFeedEntry(result, entry, expandedWorkGroupIds); + appendPresentedFeedEntry(entries, entry, expandedWorkGroupIds); } } if (activeWorkStartedAt !== null) { - result.push({ + entries.push({ type: "working", id: "working-indicator-row", createdAt: activeWorkStartedAt, }); } - return result; + return { entries, settledTurnOpeningAssistantMessageIds }; } +export function deriveThreadFeedPresentation( + feed: ReadonlyArray, + latestTurn: ThreadFeedLatestTurn | null, + expandedTurnIds: ReadonlySet, + expandedWorkGroupIds: ReadonlySet = new Set(), + activeWorkStartedAt: string | null = null, +): ThreadFeedEntry[] { + return deriveThreadFeedPresentationState( + feed, + latestTurn, + expandedTurnIds, + expandedWorkGroupIds, + activeWorkStartedAt, + ).entries; +} + +type PresentableThreadFeedEntry = Exclude< + ThreadFeedEntry, + { readonly type: "turn-fold" | "work-toggle" | "working" } +>; + function appendPresentedFeedEntry( result: ThreadFeedEntry[], - entry: Exclude, + entry: PresentableThreadFeedEntry, expandedWorkGroupIds: ReadonlySet, ): void { if (entry.type !== "activity-group") { diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 2fbb8aaf5174..6164ef17e9bc 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -46,6 +46,7 @@ import { OrchestrationCommandReceiptRepositoryLive } from "../../persistence/Lay import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import { ProviderService, + shouldApplySessionScopedRuntimeEvent, type ProviderServiceShape, } from "../../provider/Services/ProviderService.ts"; import * as RepositoryIdentityResolver from "../../project/RepositoryIdentityResolver.ts"; @@ -105,7 +106,12 @@ function isLegacyTurnCompletedEvent( } function createProviderServiceHarness() { - const runtimeEventPubSub = Effect.runSync(PubSub.unbounded()); + const runtimeEventPubSub = Effect.runSync( + PubSub.unbounded<{ + readonly event: ProviderRuntimeEvent; + readonly eventIsStaleGeneration: boolean; + }>(), + ); const runtimeSessions: ProviderSession[] = []; const unsupported = () => Effect.die(new Error("Unsupported provider call in test")) as never; @@ -135,7 +141,12 @@ function createProviderServiceHarness() { }, rollbackConversation: () => unsupported(), get streamEvents() { - return Stream.fromPubSub(runtimeEventPubSub); + return Stream.fromPubSub(runtimeEventPubSub).pipe( + Stream.filter(({ eventIsStaleGeneration }) => + shouldApplySessionScopedRuntimeEvent(eventIsStaleGeneration, true), + ), + Stream.map(({ event }) => event), + ); }, }; @@ -148,6 +159,10 @@ function createProviderServiceHarness() { runtimeSessions.push(session); }; + const clearSessions = (): void => { + runtimeSessions.length = 0; + }; + const normalizeLegacyEvent = (event: LegacyProviderRuntimeEvent): ProviderRuntimeEvent => { if (isLegacyTurnCompletedEvent(event)) { const normalized: Extract = { @@ -163,14 +178,23 @@ function createProviderServiceHarness() { return event as ProviderRuntimeEvent; }; - const emit = (event: LegacyProviderRuntimeEvent): void => { - Effect.runSync(PubSub.publish(runtimeEventPubSub, normalizeLegacyEvent(event))); + const emit = ( + event: LegacyProviderRuntimeEvent, + options?: { readonly eventIsStaleGeneration?: boolean }, + ): void => { + Effect.runSync( + PubSub.publish(runtimeEventPubSub, { + event: normalizeLegacyEvent(event), + eventIsStaleGeneration: options?.eventIsStaleGeneration ?? false, + }), + ); }; return { service, emit, setSession, + clearSessions, }; } @@ -384,6 +408,7 @@ describe("ProviderRuntimeIngestion", () => { readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), emit: provider.emit, setProviderSession: provider.setSession, + clearProviderSessions: provider.clearSessions, drain, getRateLimitState, listUsageSnapshots, @@ -391,6 +416,38 @@ describe("ProviderRuntimeIngestion", () => { }; } + async function setRunningLiveSession( + harness: Awaited>, + sessionGenerationId: string, + ) { + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-live-generation"); + const createdAt = "2026-01-01T00:00:00.000Z"; + harness.setProviderSession({ + provider: ProviderDriverKind.make("codex"), + status: "running", + runtimeMode: "approval-required", + threadId, + activeTurnId: turnId, + sessionGenerationId, + createdAt, + updatedAt: createdAt, + }); + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-live-generation"), + provider: ProviderDriverKind.make("codex"), + threadId, + turnId, + createdAt, + }); + await waitForThread( + harness.readModel, + (thread) => thread.session?.status === "running" && thread.session.activeTurnId === turnId, + ); + return { threadId, turnId }; + } + it("maps turn started/completed events into thread session updates", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; @@ -745,6 +802,159 @@ describe("ProviderRuntimeIngestion", () => { }), ); + it("ignores a session exit from a superseded session generation", async () => { + const harness = await createHarness(); + const { threadId, turnId } = await setRunningLiveSession(harness, "generation-2"); + + harness.emit( + { + type: "session.exited", + eventId: asEventId("evt-session-exited-stale-generation"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + payload: { sessionGenerationId: "generation-1" }, + }, + { eventIsStaleGeneration: true }, + ); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(turnId); + }); + + it("does not apply a stale exit over a replacement that is still starting", async () => { + const harness = await createHarness(); + const threadId = asThreadId("thread-1"); + harness.clearProviderSessions(); + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-replacement"), + threadId, + session: { + threadId, + status: "starting", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: "2026-01-01T00:00:01.000Z", + }, + createdAt: "2026-01-01T00:00:01.000Z", + }); + + harness.emit( + { + type: "session.exited", + eventId: asEventId("evt-session-exited-during-replacement-start"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + payload: { sessionGenerationId: "generation-1" }, + }, + { eventIsStaleGeneration: true }, + ); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("starting"); + }); + + it("ignores a turn completion from a superseded session generation", async () => { + const harness = await createHarness(); + const { threadId, turnId } = await setRunningLiveSession(harness, "generation-2"); + + harness.emit( + { + type: "turn.completed", + eventId: asEventId("evt-turn-completed-stale-generation"), + provider: ProviderDriverKind.make("codex"), + threadId, + turnId, + createdAt: "2026-01-01T00:00:01.000Z", + payload: { state: "completed", sessionGenerationId: "generation-1" }, + }, + { eventIsStaleGeneration: true }, + ); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(turnId); + }); + + it("ignores a runtime error from a superseded session generation", async () => { + const harness = await createHarness(); + const { threadId, turnId } = await setRunningLiveSession(harness, "generation-2"); + + harness.emit( + { + type: "runtime.error", + eventId: asEventId("evt-runtime-error-stale-generation"), + provider: ProviderDriverKind.make("opencode"), + threadId, + turnId, + createdAt: "2026-01-01T00:00:01.000Z", + payload: { + message: "old transport failed", + class: "transport_error", + sessionGenerationId: "generation-1", + }, + }, + { eventIsStaleGeneration: true }, + ); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(turnId); + }); + + it("uses the strict lifecycle flag for the shared generation decision", () => { + expect(shouldApplySessionScopedRuntimeEvent(true, true)).toBe(false); + expect(shouldApplySessionScopedRuntimeEvent(true, false)).toBe(true); + expect(shouldApplySessionScopedRuntimeEvent(false, true)).toBe(true); + }); + + it("accepts a session exit from the live session generation", async () => { + const harness = await createHarness(); + const { threadId } = await setRunningLiveSession(harness, "generation-2"); + + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-session-exited-live-generation"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + payload: { sessionGenerationId: "generation-2" }, + }); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("stopped"); + expect(thread?.session?.activeTurnId).toBeNull(); + }); + + it("accepts an unstamped legacy session exit", async () => { + const harness = await createHarness(); + const { threadId } = await setRunningLiveSession(harness, "generation-2"); + + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-session-exited-legacy"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + payload: {}, + }); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("stopped"); + expect(thread?.session?.activeTurnId).toBeNull(); + }); + it("does not clear active turn when session/thread started arrives mid-turn", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; @@ -3042,17 +3252,20 @@ describe("ProviderRuntimeIngestion", () => { updatedAt: "2026-01-01T00:00:00.000Z", }); - harness.emit({ - type: "session.cwd.changed", - eventId: asEventId("evt-session-cwd-stale"), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - payload: { - cwd: worktreePath, - sessionGenerationId: "generation-1", + harness.emit( + { + type: "session.cwd.changed", + eventId: asEventId("evt-session-cwd-stale"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + payload: { + cwd: worktreePath, + sessionGenerationId: "generation-1", + }, }, - }); + { eventIsStaleGeneration: true }, + ); await harness.drain(); const thread = (await harness.readModel()).threads.find( diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 1308c69674d2..fd7cd7f8009d 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -1491,30 +1491,6 @@ const make = Effect.gen(function* () { }) { const { event, thread } = input; - // Events are drained asynchronously, so one emitted by a session generation - // that has since been replaced can arrive after the thread was rebound. The - // provider binding rejects those; the thread's metadata must too, or the UI - // ends up describing a directory no live session is using. - const sessions = yield* providerService.listSessions(); - const liveSession = sessions.find((entry) => entry.threadId === thread.id); - // Reject only against a *different* live generation. Requiring a live - // session would discard the most important case of all: the agent leaves a - // worktree as its last act, the session exits, and the event drains after - // the adapter has already dropped it — leaving the thread pointing at a - // directory nothing uses. Events drain in order, so when no session is live - // the newest observation is still the correct one. - if ( - liveSession !== undefined && - liveSession.sessionGenerationId !== event.payload.sessionGenerationId - ) { - yield* Effect.logDebug("provider.session.cwd-changed-stale-generation", { - threadId: thread.id, - eventGenerationId: event.payload.sessionGenerationId, - liveGenerationId: liveSession.sessionGenerationId, - }); - return; - } - const workspace = yield* projectionSnapshotQuery .getThreadCheckpointContext(thread.id) .pipe(Effect.map(Option.getOrUndefined)); @@ -1721,7 +1697,6 @@ const make = Effect.gen(function* () { const conflictsWithActiveTurn = activeTurnId !== null && eventTurnId !== undefined && !sameId(activeTurnId, eventTurnId); const missingTurnForActiveTurn = activeTurnId !== null && eventTurnId === undefined; - // A turn.started that conflicts with the active turn is legitimate when // the server itself has a turn start pending for this thread AND the // provider session already tracks the event's turn as its active turn: @@ -2088,7 +2063,7 @@ const make = Effect.gen(function* () { } } - if (event.type === "session.exited") { + if (event.type === "session.exited" && shouldApplyThreadLifecycle) { yield* clearTurnStateForSession(thread.id); } @@ -2225,28 +2200,6 @@ const make = Effect.gen(function* () { break; } case "session.exited": { - // Events drain asynchronously, so an exit emitted by a session - // generation that has since been replaced can arrive after the - // replacement session already recorded fresh liveness. Clearing on - // that stale exit would wipe the live session's state, mirroring the - // generation guard in mirrorSessionCwdOntoThread: reject only - // against a *different* live generation, and treat a missing event - // generation (older adapters) or no live session as clearable. - const exitedGenerationId = event.payload.sessionGenerationId; - const sessions = yield* providerService.listSessions(); - const liveSession = sessions.find((entry) => entry.threadId === thread.id); - if ( - liveSession !== undefined && - exitedGenerationId !== undefined && - liveSession.sessionGenerationId !== exitedGenerationId - ) { - yield* Effect.logDebug("provider.session.exited-stale-generation-liveness", { - threadId: thread.id, - eventGenerationId: exitedGenerationId, - liveGenerationId: liveSession.sessionGenerationId, - }); - break; - } threadBackgroundLiveness.clearThreadLiveness(thread.id); threadPlanProgress.clearThreadPlanProgress(thread.id); break; diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index b8cbee38ded2..930ff5a5523d 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -378,7 +378,14 @@ describe("ClaudeAdapterLive", () => { // Project settings are loaded for skills, so discovery must opt out of // the workspace's `.mcp.json` servers rather than booting them. assert.equal(harness.getLastCreateQueryInput()?.options.strictMcpConfig, true); - assert.equal(harness.getLastCreateQueryInput()?.options.mcpServers, undefined); + assert.deepEqual(harness.getLastCreateQueryInput()?.options.mcpServers, {}); + assert.deepEqual(harness.getLastCreateQueryInput()?.options.allowedTools, []); + // Connected claude.ai MCP servers live outside filesystem config, so the + // shared probe options must disable them independently. + assert.equal( + harness.getLastCreateQueryInput()?.options.env?.ENABLE_CLAUDEAI_MCP_SERVERS, + "false", + ); assert.equal(harness.query.closeCalls, 1); }).pipe( Effect.provide(harness.layer), diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index c247ee4de8e1..80715ac41fe6 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -96,6 +96,7 @@ import { collectClaudeSkillReferenceNames, } from "../Drivers/ClaudeSkillReferences.ts"; import { + buildClaudeCapabilitiesProbeQueryOptions, CLAUDE_SDK_INITIALIZATION_TIMEOUT_MS, getClaudeModelCapabilities, gateClaudeSkillsByUserInvocation, @@ -5201,25 +5202,16 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( ); } })(), - options: { - cwd: canonicalCwd, - persistSession: false, - pathToClaudeCodeExecutable: claudeSdkExecutablePath, + // Same isolation as the capability probe, from the same builder: + // project settings load so project skills appear, but a metadata + // lookup must not run hooks or boot filesystem/claude.ai MCP + // servers every time the picker opens. + options: buildClaudeCapabilitiesProbeQueryOptions({ + executablePath: claudeSdkExecutablePath, abortController, - settingSources: ["user", "project", "local"], - // Same reason the capability probe disables them: this runs every - // time the picker opens or refreshes, and a SessionStart hook has - // no business firing for a metadata lookup. - settings: { disableAllHooks: true }, - // Project settings must be loaded to see project skills, but a - // metadata lookup must not boot the workspace's `.mcp.json` - // servers. Strict mode keeps discovery to the skills we asked - // for instead of spawning arbitrary project subprocesses. - strictMcpConfig: true, - allowedTools: [], - env: claudeEnvironment, - stderr: () => {}, - }, + environment: claudeEnvironment, + cwd: canonicalCwd, + }), }); const initialization = await claudeQuery.initializationResult(); const commandNames = new Set( diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 4e7fd376a1eb..59abb4e9ed1d 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -789,14 +789,14 @@ const lifecycleLayer = it.layer( function startLifecycleRuntime() { return Effect.gen(function* () { const adapter = yield* CodexAdapter; - yield* adapter.startSession({ + const session = yield* adapter.startSession({ provider: ProviderDriverKind.make("codex"), threadId: asThreadId("thread-1"), runtimeMode: "full-access", }); const runtime = lifecycleRuntimeFactory.lastRuntime; NodeAssert.ok(runtime); - return { adapter, runtime }; + return { adapter, runtime, session }; }); } @@ -1102,11 +1102,41 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { }), ); + it.effect("stamps turn completion events with the session generation", () => + Effect.gen(function* () { + const { adapter, runtime, session } = yield* startLifecycleRuntime(); + const firstEventFiber = yield* Stream.runHead(adapter.streamEvents).pipe(Effect.forkChild); + + NodeAssert.ok(session.sessionGenerationId); + yield* runtime.emit({ + id: asEventId("evt-turn-completed-generation"), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-1"), + createdAt: "2026-01-01T00:00:00.000Z", + method: "turn/completed", + payload: { + threadId: "thread-1", + turn: { id: "turn-1", status: "completed", items: [] }, + }, + } satisfies ProviderEvent); + + const firstEvent = yield* Fiber.join(firstEventFiber); + NodeAssert.equal(firstEvent._tag, "Some"); + if (firstEvent._tag !== "Some" || firstEvent.value.type !== "turn.completed") { + return; + } + NodeAssert.equal(firstEvent.value.payload.sessionGenerationId, session.sessionGenerationId); + }), + ); + it.effect("maps session/closed lifecycle events to canonical session.exited runtime events", () => Effect.gen(function* () { - const { adapter, runtime } = yield* startLifecycleRuntime(); + const { adapter, runtime, session } = yield* startLifecycleRuntime(); const firstEventFiber = yield* Stream.runHead(adapter.streamEvents).pipe(Effect.forkChild); + NodeAssert.ok(session.sessionGenerationId); const event: ProviderEvent = { id: asEventId("evt-session-closed"), kind: "session", @@ -1130,6 +1160,7 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { } NodeAssert.equal(firstEvent.value.threadId, "thread-1"); NodeAssert.equal(firstEvent.value.payload.reason, "Session stopped"); + NodeAssert.equal(firstEvent.value.payload.sessionGenerationId, session.sessionGenerationId); }), ); diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index 82bce6e1a28d..8575c9306aa3 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -96,12 +96,33 @@ export interface CodexAdapterLiveOptions { interface CodexAdapterSessionContext { readonly threadId: ThreadId; + readonly sessionGenerationId: string; readonly scope: Scope.Closeable; readonly runtime: CodexSessionRuntimeShape; readonly eventFiber: Fiber.Fiber; stopped: boolean; } +function stampSessionGeneration( + event: ProviderRuntimeEvent, + sessionGenerationId: string, +): ProviderRuntimeEvent { + switch (event.type) { + case "session.exited": + return { + ...event, + payload: { ...event.payload, sessionGenerationId }, + }; + case "turn.completed": + return { + ...event, + payload: { ...event.payload, sessionGenerationId }, + }; + default: + return event; + } +} + function mapCodexRuntimeError( threadId: ThreadId, method: string, @@ -1645,6 +1666,17 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( options?.nativeEventLogger === undefined ? nativeEventLogger : undefined; const runtimeEventQueue = yield* Queue.unbounded(); const sessions = new Map(); + const randomUUIDv4 = crypto.randomUUIDv4.pipe( + Effect.mapError( + (cause) => + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "crypto/randomUUIDv4", + detail: "Failed to generate Codex runtime identifier.", + cause, + }), + ), + ); const startSession: CodexAdapterShape["startSession"] = (input) => Effect.scoped( @@ -1666,6 +1698,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( input.modelSelection?.instanceId === boundInstanceId ? getCodexServiceTierOptionValue(input.modelSelection) : undefined; + const sessionGenerationId = yield* randomUUIDv4; const mcpSession = McpProviderSession.readMcpProviderSession(input.threadId); const runtimeInput: CodexSessionRuntimeOptions = { threadId: input.threadId, @@ -1726,7 +1759,9 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( const eventFiber = yield* Stream.runForEach(runtime.events, (event) => Effect.gen(function* () { yield* writeNativeEvent(event); - const runtimeEvents = mapToRuntimeEvents(event, event.threadId); + const runtimeEvents = mapToRuntimeEvents(event, event.threadId).map((runtimeEvent) => + stampSessionGeneration(runtimeEvent, sessionGenerationId), + ); if (runtimeEvents.length === 0) { yield* Effect.logDebug("ignoring unhandled Codex provider event", { method: event.method, @@ -1759,8 +1794,10 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( ), ); + const session = { ...started, sessionGenerationId }; sessions.set(input.threadId, { threadId: input.threadId, + sessionGenerationId, scope: sessionScope, runtime, eventFiber, @@ -1768,7 +1805,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( }); sessionScopeTransferred = true; - return started; + return session; }), ); @@ -2032,7 +2069,13 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( const listSessions: CodexAdapterShape["listSessions"] = () => Effect.forEach( Array.from(sessions.values()).filter((session) => !session.stopped), - (session) => session.runtime.getSession, + (session) => + session.runtime.getSession.pipe( + Effect.map((runtimeSession) => ({ + ...runtimeSession, + sessionGenerationId: session.sessionGenerationId, + })), + ), { concurrency: 1 }, ); diff --git a/apps/server/src/provider/Layers/CursorAdapter.test.ts b/apps/server/src/provider/Layers/CursorAdapter.test.ts index cd5cdb7f01aa..f6394b5ae87e 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.test.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.test.ts @@ -195,6 +195,7 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { schemaVersion: 1, sessionId: "mock-session-1", }); + assert.isString(session.sessionGenerationId); yield* adapter.sendTurn({ threadId, @@ -236,6 +237,13 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { event.type === "item.completed" && event.payload.itemType === "assistant_message", ); assert.isDefined(assistantCompleted); + const turnCompleted = runtimeEvents.find((event) => event.type === "turn.completed"); + assert.equal( + turnCompleted?.type === "turn.completed" + ? turnCompleted.payload.sessionGenerationId + : undefined, + session.sessionGenerationId, + ); const planUpdate = runtimeEvents.find((event) => event.type === "turn.plan.updated"); assert.isDefined(planUpdate); @@ -250,6 +258,38 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { }), ); + it.effect("stamps session exit events with the session generation", () => + Effect.gen(function* () { + const adapter = yield* CursorAdapter; + const settings = yield* ServerSettingsService; + const threadId = ThreadId.make("cursor-exit-generation"); + const wrapperPath = yield* Effect.promise(() => makeMockAgentWrapper()); + yield* settings.updateSettings({ providers: { cursor: { binaryPath: wrapperPath } } }); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.take(4), + Stream.runCollect, + Effect.forkChild, + ); + + const session = yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("cursor"), + cwd: process.cwd(), + runtimeMode: "full-access", + }); + yield* adapter.stopSession(threadId); + + const events = Array.from(yield* Fiber.join(eventsFiber).pipe(Effect.timeout("2 seconds"))); + const exited = events.find((event) => event.type === "session.exited"); + assert.isString(session.sessionGenerationId); + assert.equal( + exited?.type === "session.exited" ? exited.payload.sessionGenerationId : undefined, + session.sessionGenerationId, + ); + }), + ); + it.effect("steers a running turn instead of opening a new one on mid-turn sendTurn", () => Effect.gen(function* () { const adapter = yield* CursorAdapter; diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index 30c173d8fae8..033b9bc8389b 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -124,6 +124,7 @@ interface PendingUserInput { interface CursorSessionContext { readonly threadId: ThreadId; + readonly sessionGenerationId: string; session: ProviderSession; readonly scope: Scope.Closeable; readonly acp: AcpSessionRuntime.AcpSessionRuntime["Service"]; @@ -472,7 +473,10 @@ export function makeCursorAdapter( ...(yield* makeEventStamp()), provider: PROVIDER, threadId: ctx.threadId, - payload: { exitKind: "graceful" }, + payload: { + exitKind: "graceful", + sessionGenerationId: ctx.sessionGenerationId, + }, }); }); @@ -751,6 +755,7 @@ export function makeCursorAdapter( }); const now = yield* nowIso; + const sessionGenerationId = yield* randomUUIDv4; const session: ProviderSession = { provider: PROVIDER, providerInstanceId: boundInstanceId, @@ -763,12 +768,14 @@ export function makeCursorAdapter( schemaVersion: CURSOR_RESUME_VERSION, sessionId: started.sessionId, }, + sessionGenerationId, createdAt: now, updatedAt: now, }; ctx = { threadId: input.threadId, + sessionGenerationId, session, scope: sessionScope, acp, @@ -1046,6 +1053,7 @@ export function makeCursorAdapter( payload: { state: result.stopReason === "cancelled" ? "cancelled" : "completed", stopReason: result.stopReason ?? null, + sessionGenerationId: ctx.sessionGenerationId, }, }); } diff --git a/apps/server/src/provider/Layers/GrokAdapter.test.ts b/apps/server/src/provider/Layers/GrokAdapter.test.ts index 6cb71660a74c..a9cc11d65ba3 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.test.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.test.ts @@ -97,6 +97,8 @@ it("requires a settlement to match the live Grok turn", () => { grokPromptSettlementBelongsToContext({ liveAcpSessionId: "session-1", expectedAcpSessionId: "session-1", + liveSessionGenerationId: "generation-1", + originatingSessionGenerationId: "generation-1", liveActiveTurnId: replacementTurnId, liveSessionActiveTurnId: replacementTurnId, turnId: staleTurnId, @@ -106,6 +108,19 @@ it("requires a settlement to match the live Grok turn", () => { grokPromptSettlementBelongsToContext({ liveAcpSessionId: "replacement-session", expectedAcpSessionId: "stale-session", + liveSessionGenerationId: "generation-1", + originatingSessionGenerationId: "generation-1", + liveActiveTurnId: staleTurnId, + liveSessionActiveTurnId: staleTurnId, + turnId: staleTurnId, + }), + ); + assert.isFalse( + grokPromptSettlementBelongsToContext({ + liveAcpSessionId: "session-1", + expectedAcpSessionId: "session-1", + liveSessionGenerationId: "generation-2", + originatingSessionGenerationId: "generation-1", liveActiveTurnId: staleTurnId, liveSessionActiveTurnId: staleTurnId, turnId: staleTurnId, @@ -115,6 +130,8 @@ it("requires a settlement to match the live Grok turn", () => { grokPromptSettlementBelongsToContext({ liveAcpSessionId: "session-1", expectedAcpSessionId: "session-1", + liveSessionGenerationId: "generation-1", + originatingSessionGenerationId: "generation-1", liveActiveTurnId: staleTurnId, liveSessionActiveTurnId: staleTurnId, turnId: staleTurnId, @@ -157,6 +174,7 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { schemaVersion: 1, sessionId: "mock-session-1", }); + assert.isString(session.sessionGenerationId); yield* adapter.sendTurn({ threadId, @@ -183,6 +201,11 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { if (delta?.type === "content.delta") { assert.equal(delta.payload.delta, "hello from mock"); } + const completed = runtimeEvents.find((event) => event.type === "turn.completed"); + assert.equal( + completed?.type === "turn.completed" ? completed.payload.sessionGenerationId : undefined, + session.sessionGenerationId, + ); yield* adapter.stopSession(threadId); }), @@ -202,8 +225,14 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { }), ); const adapter = yield* makeTestAdapter(wrapperPath); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.take(4), + Stream.runCollect, + Effect.forkChild, + ); - yield* adapter.startSession({ + const session = yield* adapter.startSession({ threadId, provider: ProviderDriverKind.make("grok"), cwd: process.cwd(), @@ -215,6 +244,12 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { const exitLog = yield* waitForFileContent(exitLogPath); assert.include(exitLog, "SIGTERM"); + const events = Array.from(yield* Fiber.join(eventsFiber).pipe(Effect.timeout("2 seconds"))); + const exited = events.find((event) => event.type === "session.exited"); + assert.equal( + exited?.type === "session.exited" ? exited.payload.sessionGenerationId : undefined, + session.sessionGenerationId, + ); }), ); diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index 858d862e6d5f..263230b4c4f2 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -101,6 +101,7 @@ interface PendingUserInput { interface GrokSessionContext { readonly threadId: ThreadId; readonly acpSessionId: string; + readonly sessionGenerationId: string; session: ProviderSession; readonly scope: Scope.Closeable; readonly acp: AcpSessionRuntime.AcpSessionRuntime["Service"]; @@ -214,12 +215,15 @@ function completedStopReasonFromPromptResponse( export function grokPromptSettlementBelongsToContext(input: { readonly liveAcpSessionId: string; readonly expectedAcpSessionId: string; + readonly liveSessionGenerationId: string; + readonly originatingSessionGenerationId: string; readonly liveActiveTurnId: TurnId | undefined; readonly liveSessionActiveTurnId: TurnId | undefined; readonly turnId: TurnId; }): boolean { return ( input.liveAcpSessionId === input.expectedAcpSessionId && + input.liveSessionGenerationId === input.originatingSessionGenerationId && (input.liveActiveTurnId === input.turnId || input.liveSessionActiveTurnId === input.turnId) ); } @@ -298,6 +302,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte threadId: ThreadId, turnId: TurnId, expectedAcpSessionId: string, + originatingSessionGenerationId: string, options?: { readonly errorMessage?: string; readonly completedStopReason?: EffectAcpSchema.StopReason | null; @@ -314,6 +319,8 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const settlementBelongsToLiveContext = grokPromptSettlementBelongsToContext({ liveAcpSessionId: liveCtx.acpSessionId, expectedAcpSessionId, + liveSessionGenerationId: liveCtx.sessionGenerationId, + originatingSessionGenerationId, liveActiveTurnId: liveCtx.activeTurnId, liveSessionActiveTurnId: liveCtx.session.activeTurnId, turnId, @@ -339,6 +346,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte payload: { state: "failed", errorMessage: options.errorMessage, + sessionGenerationId: originatingSessionGenerationId, }, }); } else if (options?.completedStopReason !== undefined) { @@ -351,6 +359,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte payload: { state: options.completedStopReason === "cancelled" ? "cancelled" : "completed", stopReason: options.completedStopReason ?? null, + sessionGenerationId: originatingSessionGenerationId, }, }); } @@ -415,6 +424,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte payload: { state: "failed", errorMessage: options.errorMessage, + sessionGenerationId: originatingSessionGenerationId, }, }); } else if (shouldEmitCompletedTurn) { @@ -427,6 +437,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte payload: { state: options.completedStopReason === "cancelled" ? "cancelled" : "completed", stopReason: options.completedStopReason ?? null, + sessionGenerationId: originatingSessionGenerationId, }, }); } @@ -523,7 +534,10 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ...(yield* makeEventStamp()), provider: PROVIDER, threadId: ctx.threadId, - payload: { exitKind: "graceful" }, + payload: { + exitKind: "graceful", + sessionGenerationId: ctx.sessionGenerationId, + }, }); }); @@ -747,6 +761,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte }); const now = yield* nowIso; + const sessionGenerationId = yield* randomUUIDv4; const session: ProviderSession = { provider: PROVIDER, providerInstanceId: boundInstanceId, @@ -759,6 +774,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte schemaVersion: GROK_RESUME_VERSION, sessionId: started.sessionId, }, + sessionGenerationId, createdAt: now, updatedAt: now, }; @@ -766,6 +782,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const ctx: GrokSessionContext = { threadId: input.threadId, acpSessionId: started.sessionId, + sessionGenerationId, session, scope: sessionScope, acp, @@ -1011,11 +1028,17 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte yield* Effect.yieldNow; } if (ctx.interruptedTurnIds.has(turnId)) { - yield* settlePromptInFlight(input.threadId, turnId, ctx.acpSessionId, { - completedStopReason: "cancelled", - emitTurnCompletion: false, - settleAllPrompts: true, - }); + yield* settlePromptInFlight( + input.threadId, + turnId, + ctx.acpSessionId, + ctx.sessionGenerationId, + { + completedStopReason: "cancelled", + emitTurnCompletion: false, + settleAllPrompts: true, + }, + ); return yield* new ProviderAdapterRequestError({ provider: PROVIDER, method: "session/prompt", @@ -1047,6 +1070,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte return { acp: ctx.acp, acpSessionId: ctx.acpSessionId, + sessionGenerationId: ctx.sessionGenerationId, displayModel, promptParts, turnId, @@ -1058,10 +1082,16 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte if (!liveCtx) { return; } - yield* settlePromptInFlight(input.threadId, turnId, liveCtx.acpSessionId, { - errorMessage: "Grok prompt preparation failed.", - emitTurnCompletion: false, - }); + yield* settlePromptInFlight( + input.threadId, + turnId, + liveCtx.acpSessionId, + liveCtx.sessionGenerationId, + { + errorMessage: "Grok prompt preparation failed.", + emitTurnCompletion: false, + }, + ); }), ), ); @@ -1102,11 +1132,15 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte input.threadId, Effect.gen(function* () { const ctx = yield* requireSession(input.threadId); - if (ctx.acpSessionId !== prepared.acpSessionId) { + if ( + ctx.acpSessionId !== prepared.acpSessionId || + ctx.sessionGenerationId !== prepared.sessionGenerationId + ) { yield* settlePromptInFlight( input.threadId, prepared.turnId, prepared.acpSessionId, + prepared.sessionGenerationId, { errorMessage: "Grok session changed before the turn completed.", settleAllPrompts: true, @@ -1194,6 +1228,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte payload: { state: result.stopReason === "cancelled" ? "cancelled" : "completed", stopReason: completedStopReason, + sessionGenerationId: prepared.sessionGenerationId, }, }); ctx.interruptedTurnIds.delete(prepared.turnId); @@ -1225,11 +1260,15 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte input.threadId, Effect.gen(function* () { const ctx = yield* requireSession(input.threadId); - if (ctx.acpSessionId !== prepared.acpSessionId) { + if ( + ctx.acpSessionId !== prepared.acpSessionId || + ctx.sessionGenerationId !== prepared.sessionGenerationId + ) { yield* settlePromptInFlight( input.threadId, prepared.turnId, prepared.acpSessionId, + prepared.sessionGenerationId, { errorMessage: "Grok session changed before the turn completed.", settleAllPrompts: true, @@ -1257,6 +1296,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte input.threadId, prepared.turnId, prepared.acpSessionId, + prepared.sessionGenerationId, { completedStopReason: completedStopReasonFromPromptResponse(promptResult), }, @@ -1269,9 +1309,15 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const errorMessage = yield* Ref.get(promptFailureMessageRef); yield* withThreadLock( input.threadId, - settlePromptInFlight(input.threadId, prepared.turnId, prepared.acpSessionId, { - errorMessage: errorMessage ?? "Grok prompt request failed.", - }), + settlePromptInFlight( + input.threadId, + prepared.turnId, + prepared.acpSessionId, + prepared.sessionGenerationId, + { + errorMessage: errorMessage ?? "Grok prompt request failed.", + }, + ), ); }).pipe(Effect.catch(() => Effect.void)), ), @@ -1338,10 +1384,16 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte ); if (interruptedTurnId) { ctx.interruptedTurnIds.add(interruptedTurnId); - yield* settlePromptInFlight(threadId, interruptedTurnId, ctx.acpSessionId, { - completedStopReason: "cancelled", - settleAllPrompts: true, - }); + yield* settlePromptInFlight( + threadId, + interruptedTurnId, + ctx.acpSessionId, + ctx.sessionGenerationId, + { + completedStopReason: "cancelled", + settleAllPrompts: true, + }, + ); } else if ( ctx.promptsInFlight > 0 || ctx.session.status === "running" || diff --git a/apps/server/src/provider/Layers/HermesAdapter.test.ts b/apps/server/src/provider/Layers/HermesAdapter.test.ts index 3124fae0e966..799c1fdc2662 100644 --- a/apps/server/src/provider/Layers/HermesAdapter.test.ts +++ b/apps/server/src/provider/Layers/HermesAdapter.test.ts @@ -12,6 +12,7 @@ import { ProviderDriverKind, ProviderInstanceId, ThreadId, + TurnId, type ProviderRuntimeEvent, } from "@t3tools/contracts"; import * as Deferred from "effect/Deferred"; @@ -23,7 +24,7 @@ import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; import { ServerConfig } from "../../config.ts"; -import { makeHermesAdapter } from "./HermesAdapter.ts"; +import { hermesPromptSettlementBelongsToContext, makeHermesAdapter } from "./HermesAdapter.ts"; const decodeHermesSettings = Schema.decodeSync(HermesSettings); const __dirname = NodePath.dirname(NodeURL.fileURLToPath(import.meta.url)); @@ -83,6 +84,32 @@ const makeTestAdapter = (binaryPath: string) => decodeHermesSettings({ enabled: true, binaryPath, requireGateway: false }), ).pipe(Effect.orDie); +it("requires a settlement to match the originating Hermes generation", () => { + const turnId = TurnId.make("turn-1"); + assert.isFalse( + hermesPromptSettlementBelongsToContext({ + liveAcpSessionId: "shared-session", + expectedAcpSessionId: "shared-session", + liveSessionGenerationId: "generation-2", + originatingSessionGenerationId: "generation-1", + liveActiveTurnId: turnId, + liveSessionActiveTurnId: turnId, + turnId, + }), + ); + assert.isTrue( + hermesPromptSettlementBelongsToContext({ + liveAcpSessionId: "shared-session", + expectedAcpSessionId: "shared-session", + liveSessionGenerationId: "generation-1", + originatingSessionGenerationId: "generation-1", + liveActiveTurnId: turnId, + liveSessionActiveTurnId: turnId, + turnId, + }), + ); +}); + it.layer(hermesAdapterTestLayer)("HermesAdapter", (it) => { it.effect("surfaces terminal-only Hermes authentication during session startup", () => Effect.gen(function* () { @@ -155,6 +182,14 @@ it.layer(hermesAdapterTestLayer)("HermesAdapter", (it) => { turnStarted?.type === "turn.started" ? turnStarted.payload.model : undefined, "anthropic:claude-sonnet-5", ); + const turnCompleted = events.find((event) => event.type === "turn.completed"); + assert.isString(session.sessionGenerationId); + assert.equal( + turnCompleted?.type === "turn.completed" + ? turnCompleted.payload.sessionGenerationId + : undefined, + session.sessionGenerationId, + ); assert.includeMembers( events.map((event) => event.type), ["session.started", "thread.started", "turn.started", "content.delta", "turn.completed"], @@ -164,6 +199,35 @@ it.layer(hermesAdapterTestLayer)("HermesAdapter", (it) => { }), ); + it.effect("stamps session exit events with the session generation", () => + Effect.gen(function* () { + const threadId = ThreadId.make("hermes-exit-generation"); + const adapter = yield* makeTestAdapter(yield* Effect.promise(() => makeMockHermesWrapper())); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.take(4), + Stream.runCollect, + Effect.forkChild, + ); + + const session = yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("hermes"), + cwd: process.cwd(), + runtimeMode: "full-access", + }); + yield* adapter.stopSession(threadId); + + const events = Array.from(yield* Fiber.join(eventsFiber).pipe(Effect.timeout("2 seconds"))); + const exited = events.find((event) => event.type === "session.exited"); + assert.isString(session.sessionGenerationId); + assert.equal( + exited?.type === "session.exited" ? exited.payload.sessionGenerationId : undefined, + session.sessionGenerationId, + ); + }), + ); + it.effect("echoes the default sentinel at session start and after sendTurn", () => Effect.gen(function* () { const threadId = ThreadId.make("hermes-default-model"); diff --git a/apps/server/src/provider/Layers/HermesAdapter.ts b/apps/server/src/provider/Layers/HermesAdapter.ts index 373720d2eaa6..ac9956c86145 100644 --- a/apps/server/src/provider/Layers/HermesAdapter.ts +++ b/apps/server/src/provider/Layers/HermesAdapter.ts @@ -168,6 +168,7 @@ interface HermesTurnItem { interface HermesSessionContext { readonly threadId: ThreadId; readonly acpSessionId: string; + readonly sessionGenerationId: string; session: ProviderSession; readonly scope: Scope.Closeable; readonly acp: AcpSessionRuntime.AcpSessionRuntime["Service"]; @@ -305,12 +306,15 @@ function completedStopReasonFromPromptResponse( export function hermesPromptSettlementBelongsToContext(input: { readonly liveAcpSessionId: string; readonly expectedAcpSessionId: string; + readonly liveSessionGenerationId: string; + readonly originatingSessionGenerationId: string; readonly liveActiveTurnId: TurnId | undefined; readonly liveSessionActiveTurnId: TurnId | undefined; readonly turnId: TurnId; }): boolean { return ( input.liveAcpSessionId === input.expectedAcpSessionId && + input.liveSessionGenerationId === input.originatingSessionGenerationId && (input.liveActiveTurnId === input.turnId || input.liveSessionActiveTurnId === input.turnId) ); } @@ -421,6 +425,7 @@ export function makeHermesAdapter( threadId: ThreadId, turnId: TurnId, expectedAcpSessionId: string, + originatingSessionGenerationId: string, options?: { readonly errorMessage?: string; readonly completedStopReason?: EffectAcpSchema.StopReason | null; @@ -437,6 +442,8 @@ export function makeHermesAdapter( const settlementBelongsToLiveContext = hermesPromptSettlementBelongsToContext({ liveAcpSessionId: liveCtx.acpSessionId, expectedAcpSessionId, + liveSessionGenerationId: liveCtx.sessionGenerationId, + originatingSessionGenerationId, liveActiveTurnId: liveCtx.activeTurnId, liveSessionActiveTurnId: liveCtx.session.activeTurnId, turnId, @@ -462,6 +469,7 @@ export function makeHermesAdapter( payload: { state: "failed", errorMessage: options.errorMessage, + sessionGenerationId: originatingSessionGenerationId, }, }); } else if (options?.completedStopReason !== undefined) { @@ -474,6 +482,7 @@ export function makeHermesAdapter( payload: { state: options.completedStopReason === "cancelled" ? "cancelled" : "completed", stopReason: options.completedStopReason ?? null, + sessionGenerationId: originatingSessionGenerationId, }, }); } @@ -560,6 +569,7 @@ export function makeHermesAdapter( payload: { state: "failed", errorMessage: promptFailureMessage, + sessionGenerationId: originatingSessionGenerationId, }, }); } else if (shouldEmitCompletedTurn) { @@ -572,6 +582,7 @@ export function makeHermesAdapter( payload: { state: completedStopReason === "cancelled" ? "cancelled" : "completed", stopReason: completedStopReason ?? null, + sessionGenerationId: originatingSessionGenerationId, }, }); } @@ -667,7 +678,10 @@ export function makeHermesAdapter( ...(yield* makeEventStamp()), provider: PROVIDER, threadId: ctx.threadId, - payload: { exitKind: "graceful" }, + payload: { + exitKind: "graceful", + sessionGenerationId: ctx.sessionGenerationId, + }, }); }); @@ -898,6 +912,7 @@ export function makeHermesAdapter( }); const now = yield* nowIso; + const sessionGenerationId = yield* randomUUIDv4; const session: ProviderSession = { provider: PROVIDER, providerInstanceId: boundInstanceId, @@ -911,6 +926,7 @@ export function makeHermesAdapter( sessionId: started.sessionId, defaultModelId: defaultModelId ?? null, }, + sessionGenerationId, createdAt: now, updatedAt: now, }; @@ -918,6 +934,7 @@ export function makeHermesAdapter( const ctx: HermesSessionContext = { threadId: input.threadId, acpSessionId: started.sessionId, + sessionGenerationId, session, scope: sessionScope, acp, @@ -1278,6 +1295,7 @@ export function makeHermesAdapter( return { acp: ctx.acp, acpSessionId: ctx.acpSessionId, + sessionGenerationId: ctx.sessionGenerationId, promptHandle, promptParts, promptSequence, @@ -1298,7 +1316,30 @@ export function makeHermesAdapter( input.threadId, Effect.gen(function* () { const ctx = yield* requireSession(input.threadId); - if (ctx.acpSessionId !== prepared.acpSessionId) { + if ( + ctx.acpSessionId !== prepared.acpSessionId || + ctx.sessionGenerationId !== prepared.sessionGenerationId + ) { + yield* settlePromptInFlight( + input.threadId, + prepared.turnId, + prepared.acpSessionId, + prepared.sessionGenerationId, + outcome._tag === "Failure" + ? { + errorMessage: mapAcpToAdapterError( + PROVIDER, + input.threadId, + "session/prompt", + outcome.error, + ).message, + } + : { + completedStopReason: completedStopReasonFromPromptResponse( + outcome.result, + ), + }, + ); return; } if (ctx.interruptedTurnIds.has(prepared.turnId)) { @@ -1316,6 +1357,7 @@ export function makeHermesAdapter( input.threadId, prepared.turnId, prepared.acpSessionId, + prepared.sessionGenerationId, { errorMessage: mapAcpToAdapterError( PROVIDER, @@ -1391,6 +1433,7 @@ export function makeHermesAdapter( payload: { state: "failed", errorMessage: promptFailureMessage, + sessionGenerationId: prepared.sessionGenerationId, }, } : { @@ -1404,6 +1447,7 @@ export function makeHermesAdapter( stopReason: promptWasCancelled ? "cancelled" : completedStopReason, + sessionGenerationId: prepared.sessionGenerationId, }, }, ); @@ -1499,10 +1543,16 @@ export function makeHermesAdapter( ); if (interruptedTurnId) { ctx.interruptedTurnIds.add(interruptedTurnId); - yield* settlePromptInFlight(threadId, interruptedTurnId, ctx.acpSessionId, { - completedStopReason: "cancelled", - settleAllPrompts: true, - }).pipe( + yield* settlePromptInFlight( + threadId, + interruptedTurnId, + ctx.acpSessionId, + ctx.sessionGenerationId, + { + completedStopReason: "cancelled", + settleAllPrompts: true, + }, + ).pipe( Effect.ensuring( Effect.sync(() => { ctx.interruptedTurnIds.delete(interruptedTurnId); diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index e90ad28aac53..9a51462387d9 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -36,6 +36,7 @@ import { isOpenCodeNotFound, isSameOpenCodeDirectory, makeOpenCodeAdapter, + makeOpenCodeUnexpectedRuntimeErrorPayload, mergeOpenCodeAssistantText, } from "./OpenCodeAdapter.ts"; @@ -247,7 +248,8 @@ const providerSessionDirectoryTestLayer = Layer.succeed(ProviderSessionDirectory // the layer graph reach for it — but the routing values the assertions // probe (serverUrl, serverPassword) must be threaded directly through the // decoded `OpenCodeSettings`. -const openCodeAdapterTestSettings = Schema.decodeSync(OpenCodeSettings)({ +const decodeOpenCodeSettings = Schema.decodeSync(OpenCodeSettings); +const openCodeAdapterTestSettings = decodeOpenCodeSettings({ binaryPath: "fake-opencode", serverUrl: "http://127.0.0.1:9999", serverPassword: "secret-password", @@ -294,6 +296,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { NodeAssert.equal(session.provider, "opencode"); NodeAssert.equal(session.threadId, "thread-opencode"); + NodeAssert.ok(session.sessionGenerationId); NodeAssert.deepEqual(runtimeMock.state.startCalls, []); NodeAssert.deepEqual(runtimeMock.state.sessionCreateUrls, ["http://127.0.0.1:9999"]); NodeAssert.deepEqual(runtimeMock.state.authHeaders, [ @@ -614,7 +617,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { Effect.forkChild, ); - yield* adapter.startSession({ + const session = yield* adapter.startSession({ provider: ProviderDriverKind.make("opencode"), threadId, runtimeMode: "full-access", @@ -626,6 +629,12 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { events.map((event) => event.type), ["session.started", "thread.started", "session.exited"], ); + const exited = events.at(-1); + NodeAssert.ok(session.sessionGenerationId); + NodeAssert.equal( + exited?.type === "session.exited" ? exited.payload.sessionGenerationId : undefined, + session.sessionGenerationId, + ); }), ); @@ -660,6 +669,19 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("stamps unexpected runtime errors with the session generation", () => + Effect.sync(() => { + NodeAssert.deepEqual( + makeOpenCodeUnexpectedRuntimeErrorPayload("server exited", "generation-1"), + { + message: "server exited", + class: "transport_error", + sessionGenerationId: "generation-1", + }, + ); + }), + ); + it.effect("completes streamEvents when the adapter scope closes", () => Effect.gen(function* () { const scope = yield* Scope.make("sequential"); diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 8f7e42c11d7c..653768833302 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -221,7 +221,19 @@ function isOpenCodeDefaultTitle(title: string): boolean { return OPENCODE_DEFAULT_TITLE_PATTERN.test(title); } +export function makeOpenCodeUnexpectedRuntimeErrorPayload( + message: string, + sessionGenerationId: string, +): Extract["payload"] { + return { + message, + class: "transport_error", + sessionGenerationId, + }; +} + interface OpenCodeSessionContext { + readonly sessionGenerationId: string; session: ProviderSession; readonly client: OpencodeClient; readonly server: OpenCodeServerConnection; @@ -706,10 +718,7 @@ export function makeOpenCodeAdapter( turnId, })), type: "runtime.error", - payload: { - message, - class: "transport_error", - }, + payload: makeOpenCodeUnexpectedRuntimeErrorPayload(message, context.sessionGenerationId), }).pipe(Effect.ignore); yield* emit({ ...(yield* buildEventBase({ @@ -721,6 +730,7 @@ export function makeOpenCodeAdapter( reason: message, recoverable: false, exitKind: "error", + sessionGenerationId: context.sessionGenerationId, }, }).pipe(Effect.ignore); // Inline the teardown that `stopOpenCodeContext` would do; we can't @@ -1085,6 +1095,7 @@ export function makeOpenCodeAdapter( type: "turn.completed", payload: { state: "completed", + sessionGenerationId: context.sessionGenerationId, }, }); } @@ -1114,6 +1125,7 @@ export function makeOpenCodeAdapter( payload: { state: "failed", errorMessage: message, + sessionGenerationId: context.sessionGenerationId, }, }); } @@ -1208,6 +1220,7 @@ export function makeOpenCodeAdapter( const serverPassword = openCodeSettings.serverPassword; const directory = input.cwd ?? serverConfig.cwd; const resumeSessionId = parseOpenCodeResume(input.resumeCursor)?.sessionId; + const sessionGenerationId = yield* randomUUIDv4; const existing = sessions.get(input.threadId); if (existing) { yield* stopOpenCodeContext(existing); @@ -1382,11 +1395,13 @@ export function makeOpenCodeAdapter( schemaVersion: OPENCODE_RESUME_VERSION, sessionId: started.openCodeSession.id, }, + sessionGenerationId, createdAt, updatedAt: createdAt, }; const context: OpenCodeSessionContext = { + sessionGenerationId, session, client: started.client, server: started.server, @@ -1638,6 +1653,7 @@ export function makeOpenCodeAdapter( reason: "Session stopped.", recoverable: false, exitKind: "graceful", + sessionGenerationId: context.sessionGenerationId, }, }); }, diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index b067240c3efe..d0707a4de891 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -12,7 +12,6 @@ import type { } from "@t3tools/contracts"; import { ApprovalRequestId, - EnvironmentId, EventId, ProviderDriverKind, ProviderInstanceId, @@ -23,6 +22,7 @@ import { import { createModelSelection } from "@t3tools/shared/model"; import { it, assert, describe, vi } from "@effect/vitest"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; @@ -134,9 +134,10 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) { const sessions = new Map(); const runtimeEventPubSub = Effect.runSync(PubSub.unbounded()); let sessionGeneration = 0; + let beforeSessionStored: ((session: ProviderSession) => Effect.Effect) | undefined; const startSession = vi.fn((input: ProviderSessionStartInput) => - Effect.sync(() => { + Effect.gen(function* () { const now = "2026-01-01T00:00:00.000Z"; const session: ProviderSession = { provider, @@ -154,6 +155,9 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) { createdAt: now, updatedAt: now, }; + if (beforeSessionStored !== undefined) { + yield* beforeSessionStored(session); + } sessions.set(session.threadId, session); return session; }), @@ -293,6 +297,11 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) { adapter, emit, updateSession, + setBeforeSessionStored: ( + hook: ((session: ProviderSession) => Effect.Effect) | undefined, + ) => { + beforeSessionStored = hook; + }, startSession, forkSession, sendTurn, @@ -2567,6 +2576,170 @@ fanout.layer("ProviderServiceLive fanout", (it) => { }), ); + it.effect("rejects a superseded generation while its replacement is starting", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("thread-replacement-generation-decision"); + const initialSession = yield* provider.startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }); + assert.isString(initialSession.sessionGenerationId); + + const staleExitEventId = asEventId("evt-exit-during-replacement-start"); + const replacementEventId = asEventId("evt-replacement-generation-during-start"); + const guardedEventFiber = yield* provider.streamEvents.pipe( + Stream.filter( + (event) => event.eventId === staleExitEventId || event.eventId === replacementEventId, + ), + Stream.runHead, + Effect.forkChild, + ); + yield* advanceTestClock(50); + fanout.codex.setBeforeSessionStored((replacementSession) => + Effect.sync(() => { + fanout.codex.emit({ + type: "session.exited", + eventId: staleExitEventId, + provider: CODEX_DRIVER, + createdAt: "2026-01-01T00:00:01.000Z", + threadId, + payload: { + exitKind: "graceful", + sessionGenerationId: initialSession.sessionGenerationId, + }, + }); + fanout.codex.emit({ + type: "runtime.error", + eventId: replacementEventId, + provider: CODEX_DRIVER, + createdAt: "2026-01-01T00:00:02.000Z", + threadId, + payload: { + message: "replacement generation sentinel", + sessionGenerationId: replacementSession.sessionGenerationId, + }, + }); + }).pipe(Effect.andThen(Effect.yieldNow)), + ); + + const replacementSession = yield* provider.startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }); + fanout.codex.setBeforeSessionStored(undefined); + yield* advanceTestClock(50); + + const guardedEvent = Option.getOrThrow(yield* Fiber.join(guardedEventFiber)); + assert.notEqual(replacementSession.sessionGenerationId, initialSession.sessionGenerationId); + assert.equal(guardedEvent.eventId, replacementEventId); + }), + ); + + it.effect("clears the pending replacement marker when startup is interrupted", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("thread-interrupted-replacement-marker"); + const initialSession = yield* provider.startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }); + const replacementReached = yield* Deferred.make(); + const holdReplacement = yield* Deferred.make(); + fanout.codex.setBeforeSessionStored((replacementSession) => + Deferred.succeed(replacementReached, replacementSession).pipe( + Effect.andThen(Deferred.await(holdReplacement)), + ), + ); + + const replacementStartFiber = yield* provider + .startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }) + .pipe(Effect.forkChild); + const interruptedReplacement = yield* Deferred.await(replacementReached); + yield* Fiber.interrupt(replacementStartFiber); + fanout.codex.setBeforeSessionStored(undefined); + + const currentEventId = asEventId("evt-current-after-interrupted-replacement"); + const currentEventFiber = yield* provider.streamEvents.pipe( + Stream.filter((event) => event.eventId === currentEventId), + Stream.runHead, + Effect.forkChild, + ); + yield* Effect.yieldNow; + fanout.codex.emit({ + type: "runtime.error", + eventId: asEventId("evt-stale-after-interrupted-replacement"), + provider: CODEX_DRIVER, + createdAt: "2026-01-01T00:00:03.000Z", + threadId, + payload: { + message: "interrupted replacement sentinel", + sessionGenerationId: interruptedReplacement.sessionGenerationId, + }, + }); + fanout.codex.emit({ + type: "runtime.error", + eventId: currentEventId, + provider: CODEX_DRIVER, + createdAt: "2026-01-01T00:00:04.000Z", + threadId, + payload: { + message: "current generation sentinel", + sessionGenerationId: initialSession.sessionGenerationId, + }, + }); + + assert.equal(Option.getOrThrow(yield* Fiber.join(currentEventFiber)).eventId, currentEventId); + }), + ); + + it.effect("accepts a stamped exit when no session is live", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("thread-stopped-generation-compatibility"); + yield* provider.startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }); + yield* provider.stopSession({ threadId }); + + const eventId = asEventId("evt-exit-after-session-stopped"); + const eventFiber = yield* provider.streamEvents.pipe( + Stream.filter((event) => event.eventId === eventId), + Stream.runHead, + Effect.forkChild, + ); + yield* advanceTestClock(50); + fanout.codex.emit({ + type: "session.exited", + eventId, + provider: CODEX_DRIVER, + createdAt: "2026-01-01T00:00:03.000Z", + threadId, + payload: { + exitKind: "graceful", + sessionGenerationId: "superseded-generation", + }, + }); + yield* advanceTestClock(50); + + assert.equal(Option.getOrThrow(yield* Fiber.join(eventFiber)).eventId, eventId); + }), + ); + it.effect("refreshes the session binding and persists pending work on turn completion", () => Effect.gen(function* () { const provider = yield* ProviderService.ProviderService; diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index 0a7e7ea4a9cb..6d4c938fe3e5 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -26,6 +26,7 @@ import { } from "@t3tools/contracts"; import { causeErrorTag } from "@t3tools/shared/observability"; import * as DateTime from "effect/DateTime"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -65,7 +66,9 @@ import { isExistingDirectory } from "../../pathExpansion.ts"; import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import * as McpSessionRegistry from "../../mcp/McpSessionRegistry.ts"; import * as ServerSettings from "../../serverSettings.ts"; + const isModelSelection = Schema.is(ModelSelection); +const STRICT_PROVIDER_LIFECYCLE_GUARD = process.env.T3CODE_STRICT_PROVIDER_LIFECYCLE_GUARD !== "0"; /** * Hook for tests that want to override the canonical event logger pulled @@ -183,6 +186,18 @@ function readSessionGenerationId(runtimePayload: unknown | null | undefined): st return typeof value === "string" && value.length > 0 ? value : undefined; } +function readRuntimeEventSessionGenerationId(event: ProviderRuntimeEvent): string | undefined { + const value = + "sessionGenerationId" in event.payload ? event.payload.sessionGenerationId : undefined; + return typeof value === "string" && value.length > 0 ? value : undefined; +} + +interface PendingSessionReplacement { + readonly targetInstanceId: ProviderInstanceId; + readonly previousGenerationId: string; + readonly settled: Deferred.Deferred; +} + function readPersistedModelSelection( runtimePayload: ProviderSessionDirectory.ProviderRuntimeBinding["runtimePayload"], ): ModelSelection | undefined { @@ -353,6 +368,9 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( const revokeMcpCredential = options?.revokeMcpCredential ?? McpSessionRegistry.revokeActiveMcpThread; const runtimeEventPubSub = yield* PubSub.unbounded(); + const pendingSessionReplacements = yield* Ref.make( + new Map(), + ); const threadLocksRef = yield* Ref.make>(new Map()); const getThreadLock = Effect.fn("ProviderService.getThreadLock")(function* (threadId: ThreadId) { const existing = (yield* Ref.get(threadLocksRef)).get(threadId); @@ -455,14 +473,24 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( Effect.tap(() => Effect.sync(() => McpProviderSession.clearMcpProviderSession(threadId))), ); - const publishRuntimeEvent = (event: ProviderRuntimeEvent): Effect.Effect => + const publishRuntimeEvent = ( + event: ProviderRuntimeEvent, + eventIsStaleGeneration: boolean, + ): Effect.Effect => Effect.succeed(event).pipe( Effect.tap((canonicalEvent) => canonicalEventLogger ? canonicalEventLogger.write(canonicalEvent, canonicalEvent.threadId) : Effect.void, ), - Effect.flatMap((canonicalEvent) => PubSub.publish(runtimeEventPubSub, canonicalEvent)), + Effect.flatMap((canonicalEvent) => + ProviderService.shouldApplySessionScopedRuntimeEvent( + eventIsStaleGeneration, + STRICT_PROVIDER_LIFECYCLE_GUARD, + ) + ? PubSub.publish(runtimeEventPubSub, canonicalEvent) + : Effect.void, + ), Effect.asVoid, ); @@ -475,33 +503,64 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( }, event: ProviderRuntimeEvent, ) { - if ( - event.type !== "turn.completed" && - event.type !== "session.exited" && - event.type !== "session.cwd.changed" - ) { - return; + const eventGenerationId = readRuntimeEventSessionGenerationId(event); + const refreshesBinding = + event.type === "turn.completed" || + event.type === "session.exited" || + event.type === "session.cwd.changed"; + if (!refreshesBinding && eventGenerationId === undefined) { + return false; } // A cwd observation is one-shot and cannot be reconstructed from a later - // event: if it is dropped, the binding keeps a directory the session has - // already left. `refreshIfUnchanged` reports a lost compare-and-swap by - // returning false rather than failing, so `Effect.retry` never sees it — - // re-read the binding and reapply instead. - const conflictAttempts = event.type === "session.cwd.changed" ? 3 : 1; + // event. Terminal events get one re-read as well so a replacement binding + // that wins the compare-and-swap race can still mark the event stale before + // it is published. + const conflictAttempts = event.type === "session.cwd.changed" ? 3 : 2; for (let attempt = 1; attempt <= conflictAttempts; attempt += 1) { + const pendingReplacement = (yield* Ref.get(pendingSessionReplacements)).get(event.threadId); + if ( + eventGenerationId !== undefined && + pendingReplacement !== undefined && + eventGenerationId === pendingReplacement.previousGenerationId + ) { + yield* Effect.logDebug("provider.session.runtime-event-stale-replacement-generation", { + threadId: event.threadId, + eventProvider: source.provider, + eventProviderInstanceId: source.instanceId, + eventGenerationId, + targetProviderInstanceId: pendingReplacement.targetInstanceId, + }); + return true; + } + const binding = Option.getOrUndefined(yield* directory.getBinding(event.threadId)); if (!binding) { - return; + return false; + } + if (binding.status === "stopped" && pendingReplacement === undefined) { + return false; } - const eventGenerationId = event.payload.sessionGenerationId; const bindingGenerationId = readSessionGenerationId(binding.runtimePayload); - if ( - binding.provider !== source.provider || - binding.providerInstanceId !== source.instanceId || - bindingGenerationId !== eventGenerationId - ) { + const bindingMatchesEvent = + binding.provider === source.provider && + binding.providerInstanceId === source.instanceId && + bindingGenerationId === eventGenerationId; + if (!bindingMatchesEvent) { + const eventMayBeFromPendingReplacement = + eventGenerationId !== undefined && + pendingReplacement !== undefined && + source.instanceId === pendingReplacement.targetInstanceId; + if (eventMayBeFromPendingReplacement) { + yield* Deferred.await(pendingReplacement.settled); + // Waiting for replacement settlement does not consume a binding + // compare-and-swap retry. + attempt -= 1; + continue; + } + + const eventIsStaleGeneration = eventGenerationId !== undefined; yield* Effect.logDebug("provider.session.runtime-event-binding-mismatch", { threadId: event.threadId, eventProvider: source.provider, @@ -510,11 +569,14 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( bindingProviderInstanceId: binding.providerInstanceId, bindingGenerationId, eventGenerationId, + eventIsStaleGeneration, }); - return; + return eventIsStaleGeneration; } - if (binding.status === "stopped") return; + if (!refreshesBinding || binding.status === "stopped") { + return false; + } const hasPendingWork = event.type === "turn.completed" ? event.payload.hasPendingWork : false; const runtimePayloadPatch = @@ -536,7 +598,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( }) .pipe(Effect.retry({ times: 2 })); if (refreshed) { - return; + return false; } yield* Effect.logDebug("provider.session.runtime-event-binding-changed", { @@ -548,6 +610,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( remainingAttempts: conflictAttempts - attempt, }); } + return false; }); const requireBindingInstanceId = ( @@ -650,15 +713,14 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( provider: canonicalEvent.provider, providerInstanceId: source.instanceId, cause, - }), + }).pipe(Effect.as(false)), ), - Effect.andThen( + Effect.flatMap((eventIsStaleGeneration) => increment(providerRuntimeEventsTotal, { provider: canonicalEvent.provider, eventType: canonicalEvent.type, - }), + }).pipe(Effect.andThen(publishRuntimeEvent(canonicalEvent, eventIsStaleGeneration))), ), - Effect.andThen(publishRuntimeEvent(canonicalEvent)), ), ), ); @@ -1070,36 +1132,74 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( "provider.cwd.effective": effectiveCwd ?? "", }); const adapter = yield* registry.getByInstance(resolvedInstanceId); - yield* prepareMcpSession(threadId, resolvedInstanceId); - const session = yield* adapter - .startSession({ - ...input, + const previousGenerationId = readSessionGenerationId(persistedBinding?.runtimePayload); + const pendingReplacement = + previousGenerationId !== undefined + ? { + targetInstanceId: resolvedInstanceId, + previousGenerationId, + settled: yield* Deferred.make(), + } + : undefined; + const sessionWithInstance = yield* Effect.gen(function* () { + if (pendingReplacement !== undefined) { + yield* Ref.update(pendingSessionReplacements, (current) => { + const next = new Map(current); + next.set(threadId, pendingReplacement); + return next; + }); + } + yield* prepareMcpSession(threadId, resolvedInstanceId); + const session = yield* adapter + .startSession({ + ...input, + providerInstanceId: resolvedInstanceId, + ...(effectiveCwd !== undefined ? { cwd: effectiveCwd } : {}), + ...(effectiveResumeCursor !== undefined + ? { resumeCursor: effectiveResumeCursor } + : {}), + }) + .pipe(Effect.onError(() => clearMcpSession(threadId))); + + if (session.provider !== adapter.provider) { + yield* clearMcpSession(threadId); + return yield* toValidationError( + "ProviderService.startSession", + `Adapter/provider mismatch: requested '${adapter.provider}', received '${session.provider}'.`, + ); + } + const boundSession = { + ...session, providerInstanceId: resolvedInstanceId, - ...(effectiveCwd !== undefined ? { cwd: effectiveCwd } : {}), - ...(effectiveResumeCursor !== undefined ? { resumeCursor: effectiveResumeCursor } : {}), - }) - .pipe(Effect.onError(() => clearMcpSession(threadId))); - - if (session.provider !== adapter.provider) { - yield* clearMcpSession(threadId); - return yield* toValidationError( - "ProviderService.startSession", - `Adapter/provider mismatch: requested '${adapter.provider}', received '${session.provider}'.`, - ); - } - const sessionWithInstance = { - ...session, - providerInstanceId: resolvedInstanceId, - }; + }; - yield* stopStaleSessionsForThread({ - threadId, - currentInstanceId: resolvedInstanceId, - }); - yield* upsertSessionBinding(sessionWithInstance, threadId, { - modelSelection: input.modelSelection, - clearHasPendingWork: true, - }); + yield* stopStaleSessionsForThread({ + threadId, + currentInstanceId: resolvedInstanceId, + }); + yield* upsertSessionBinding(boundSession, threadId, { + modelSelection: input.modelSelection, + clearHasPendingWork: true, + }); + return boundSession; + }).pipe( + Effect.ensuring( + pendingReplacement === undefined + ? Effect.void + : Ref.update(pendingSessionReplacements, (current) => { + if (current.get(threadId) !== pendingReplacement) { + return current; + } + const next = new Map(current); + next.delete(threadId); + return next; + }).pipe( + Effect.andThen( + Deferred.succeed(pendingReplacement.settled, undefined).pipe(Effect.ignore), + ), + ), + ), + ); yield* analytics.record("provider.session.started", { provider: sessionWithInstance.provider, runtimeMode: input.runtimeMode, @@ -1814,9 +1914,8 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( getCapabilities, getInstanceInfo, rollbackConversation: (input) => withThreadLock(input.threadId, rollbackConversation(input)), - // Each access creates a fresh PubSub subscription so that multiple - // consumers (ProviderRuntimeIngestion, CheckpointReactor, etc.) each - // independently receive all runtime events. + // Each access creates a fresh PubSub subscription so multiple consumers + // independently receive every generation-accepted runtime event. get streamEvents(): ProviderServiceMethod<"streamEvents"> { return Stream.fromPubSub(runtimeEventPubSub); }, diff --git a/apps/server/src/provider/Services/ProviderService.ts b/apps/server/src/provider/Services/ProviderService.ts index 5b922c5b0109..9e64a507483f 100644 --- a/apps/server/src/provider/Services/ProviderService.ts +++ b/apps/server/src/provider/Services/ProviderService.ts @@ -34,6 +34,13 @@ import type { ProviderForkSessionResult } from "./ProviderAdapter.ts"; import type { ProviderInstanceRoutingInfo } from "./ProviderAdapterRegistry.ts"; import type { ProviderRuntimeBindingWithMetadata } from "./ProviderSessionDirectory.ts"; +export function shouldApplySessionScopedRuntimeEvent( + eventIsStaleGeneration: boolean, + strictProviderLifecycleGuard: boolean, +): boolean { + return !strictProviderLifecycleGuard || !eventIsStaleGeneration; +} + export interface ProviderSessionStartOptions { /** * `fail` makes an existing persisted cursor authoritative: the session may @@ -133,7 +140,7 @@ export interface ProviderServiceShape { }) => Effect.Effect; /** - * Canonical provider runtime event stream. + * Generation-guarded canonical runtime events from all registered adapters. * * Fan-out is owned by ProviderService (not by a standalone event-bus service). */ diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts index e3059cf5fc98..97e2c7eecfff 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts @@ -3,119 +3,11 @@ import { computeStableMessagesTimelineRows, computeMessageDurationStart, deriveMessagesTimelineRows, - deriveTimelineMinimapItems, normalizeCompactToolLabel, resolveAssistantMessageCopyState, - resolveTimelineMinimapAriaLabel, - resolveTimelineMinimapItemIndexFromPointer, - resolveTimelineMinimapTooltipTranslate, shouldPreserveAssistantLineBreaks, } from "./MessagesTimeline.logic"; - -describe("resolveTimelineMinimapAriaLabel", () => { - it("uses a concise role and position label with a capped excerpt", () => { - const longResponse = "A".repeat(180); - - expect( - resolveTimelineMinimapAriaLabel({ - activeIndex: null, - itemCount: 3, - primaryText: null, - side: "right", - }), - ).toBe("Jump to message: Agent response"); - expect( - resolveTimelineMinimapAriaLabel({ - activeIndex: 1, - itemCount: 3, - primaryText: longResponse, - side: "right", - }), - ).toBe(`Jump to agent response 2 of 3: ${"A".repeat(120)}`); - }); -}); - -describe("resolveTimelineMinimapItemIndexFromPointer", () => { - it("selects the nearest rendered reply in the full turn-position space", () => { - const items = [ - { positionIndex: 0, positionCount: 4 }, - { positionIndex: 3, positionCount: 4 }, - ]; - - expect( - resolveTimelineMinimapItemIndexFromPointer({ - items, - railTop: 100, - railHeight: 300, - pointerY: 360, - }), - ).toBe(1); - expect( - resolveTimelineMinimapItemIndexFromPointer({ - items, - railTop: 100, - railHeight: 300, - pointerY: 160, - }), - ).toBe(0); - }); - - it("compares sparse markers against the continuous pointer position", () => { - expect( - resolveTimelineMinimapItemIndexFromPointer({ - items: [ - { positionIndex: 0, positionCount: 4 }, - { positionIndex: 2, positionCount: 4 }, - ], - railTop: 0, - railHeight: 24, - pointerY: 11, - }), - ).toBe(1); - }); -}); - -describe("resolveTimelineMinimapTooltipTranslate", () => { - it("only edge-clamps markers at the endpoints of the full turn space", () => { - expect(resolveTimelineMinimapTooltipTranslate(0, 5)).toBe("0%"); - expect(resolveTimelineMinimapTooltipTranslate(1, 5)).toBe("-50%"); - expect(resolveTimelineMinimapTooltipTranslate(3, 5)).toBe("-50%"); - expect(resolveTimelineMinimapTooltipTranslate(4, 5)).toBe("-100%"); - }); -}); - -describe("deriveTimelineMinimapItems", () => { - it("keeps sparse agent replies aligned with their user-turn positions", () => { - const messageRow = ( - id: string, - role: "user" | "assistant", - isFinalAssistantResponse: boolean, - ) => - ({ - kind: "message", - id, - message: { id, role, text: id }, - isFinalAssistantResponse, - }) as never; - const rows = [ - messageRow("user-1", "user", false), - messageRow("assistant-1", "assistant", true), - messageRow("user-2", "user", false), - messageRow("assistant-commentary", "assistant", false), - messageRow("user-3", "user", false), - messageRow("assistant-3", "assistant", true), - ]; - - expect( - deriveTimelineMinimapItems(rows, "final-assistant").map( - ({ id, positionIndex, positionCount }) => ({ id, positionIndex, positionCount }), - ), - ).toEqual([ - { id: "assistant-1", positionIndex: 0, positionCount: 3 }, - { id: "assistant-3", positionIndex: 2, positionCount: 3 }, - ]); - }); -}); +import { deriveTimelineMinimapItems } from "./MessagesTimeline.minimap"; describe("shouldPreserveAssistantLineBreaks", () => { it("preserves Claude insight formatting without changing regular markdown", () => { @@ -382,7 +274,7 @@ describe("resolveAssistantMessageCopyState", () => { }); describe("deriveMessagesTimelineRows", () => { - it("only enables assistant copy for the terminal assistant message in a turn", () => { + it("enables assistant copy for the opening and terminal assistant messages in a turn", () => { const rows = deriveMessagesTimelineRows({ timelineEntries: [ { @@ -441,7 +333,9 @@ describe("deriveMessagesTimelineRows", () => { ); expect(assistantRows).toHaveLength(2); - expect(assistantRows[0]?.showAssistantCopyButton).toBe(false); + // The first message of a settled two-message turn is the preserved + // opening response, so it carries the copy button like the terminal one. + expect(assistantRows[0]?.showAssistantCopyButton).toBe(true); expect(assistantRows[1]?.showAssistantCopyButton).toBe(true); }); @@ -663,6 +557,112 @@ describe("deriveMessagesTimelineRows", () => { ).toBeDefined(); }); + it("shows the assistant meta row on the preserved opening response of a settled turn", () => { + const message = (id: string, text: string, createdAt: string) => ({ + id: `${id}-entry`, + kind: "message" as const, + createdAt, + message: { + id: id as never, + role: "assistant" as const, + text, + turnId: "turn-1" as never, + createdAt, + updatedAt: createdAt, + streaming: false, + }, + }); + const 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: "Build it", + turnId: null, + createdAt: "2026-01-01T00:00:00Z", + updatedAt: "2026-01-01T00:00:00Z", + streaming: false, + }, + }, + message("assistant-first", "Here is the plan up front.", "2026-01-01T00:00:05Z"), + message("assistant-middle", "Now doing step two.", "2026-01-01T00:00:10Z"), + message("assistant-final", "Done", "2026-01-01T00:00:20Z"), + ]; + + const rows = deriveMessagesTimelineRows({ + timelineEntries, + expandedTurnIds: new Set(["turn-1" as never]), + isWorking: false, + activeTurnStartedAt: null, + completedTurnAssistantMessageIds: new Set(["assistant-final" as never]), + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + const messageRow = (id: string) => + rows.find( + (row): row is Extract<(typeof rows)[number], { kind: "message" }> => + row.kind === "message" && row.message.id === (id as never), + ); + + // The preserved opening earns the meta row (copy/summary/speech), but + // never counts as the final response. + expect(messageRow("assistant-first")?.showAssistantMeta).toBe(true); + expect(messageRow("assistant-first")?.isFinalAssistantResponse).toBe(false); + // Folded-away commentary stays bare even when expanded into view. + expect(messageRow("assistant-middle")?.showAssistantMeta).toBe(false); + expect(messageRow("assistant-final")?.showAssistantMeta).toBe(true); + expect(messageRow("assistant-final")?.isFinalAssistantResponse).toBe(true); + }); + + it("withholds the opening meta row while the turn is still unsettled", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "assistant-first-entry", + kind: "message" as const, + createdAt: "2026-01-01T00:00:05Z", + message: { + id: "assistant-first" as never, + role: "assistant" as const, + text: "Opening answer", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:05Z", + updatedAt: "2026-01-01T00:00:06Z", + streaming: false, + }, + }, + { + id: "assistant-followup-entry", + kind: "message" as const, + createdAt: "2026-01-01T00:00:10Z", + message: { + id: "assistant-followup" as never, + role: "assistant" as const, + text: "Still working", + turnId: "turn-1" as never, + createdAt: "2026-01-01T00:00:10Z", + updatedAt: "2026-01-01T00:00:11Z", + streaming: false, + }, + }, + ], + runningTurnId: "turn-1" as never, + isWorking: true, + activeTurnStartedAt: "2026-01-01T00:00:04Z", + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + const openingRow = rows.find( + (row): row is Extract<(typeof rows)[number], { kind: "message" }> => + row.kind === "message" && row.message.id === ("assistant-first" as never), + ); + expect(openingRow?.showAssistantMeta).toBe(false); + }); + it("folds assistant messages between the first and terminal messages", () => { const timelineEntries = [ { @@ -1389,7 +1389,7 @@ describe("deriveMessagesTimelineRows", () => { expect(rows.map((row) => row.id)).toContain("work-live:running-work-entry"); }); - it("only shows assistant metadata on the terminal assistant message", () => { + it("shows assistant metadata on the opening and terminal assistant messages", () => { const rows = deriveMessagesTimelineRows({ timelineEntries: [ { @@ -1448,7 +1448,7 @@ 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([true, true]); expect(deriveTimelineMinimapItems(rows, "user-turn")).toEqual([ { id: "user-prompt-entry", diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.ts b/apps/web/src/components/chat/MessagesTimeline.logic.ts index acf0511f1b18..5c8da64d0fdd 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.ts @@ -17,7 +17,6 @@ export const TIMELINE_MINIMAP_MIN_ITEMS = 2; export const TIMELINE_MINIMAP_MAX_HEIGHT_CSS = "calc(100vh - 18rem)"; export const TIMELINE_CONTENT_MAX_WIDTH = 768; export const TIMELINE_MINIMAP_PERSISTENT_GUTTER = 48; -export const TIMELINE_MINIMAP_ARIA_EXCERPT_MAX_LENGTH = 120; export function workEntryIsVisibleInGroup( entry: WorkLogEntry, @@ -82,19 +81,6 @@ export function resolveTimelineMinimapTopPercent(index: number, itemCount: numbe return (Math.max(0, Math.min(index, itemCount - 1)) / (itemCount - 1)) * 100; } -export function resolveTimelineMinimapTooltipTranslate( - positionIndex: number, - positionCount: number, -): string { - if (positionIndex === 0) { - return "0%"; - } - if (positionIndex === positionCount - 1) { - return "-100%"; - } - return "-50%"; -} - export function resolveTimelineMinimapIndexFromPointer(input: { readonly itemCount: number; readonly railTop: number; @@ -112,35 +98,6 @@ export function resolveTimelineMinimapIndexFromPointer(input: { return Math.max(0, Math.min(input.itemCount - 1, Math.round(progress * (input.itemCount - 1)))); } -export function resolveTimelineMinimapItemIndexFromPointer(input: { - readonly items: ReadonlyArray>; - readonly railTop: number; - readonly railHeight: number; - readonly pointerY: number; -}): number | null { - const firstItem = input.items[0]; - if (!firstItem || input.railHeight <= 0) { - return null; - } - const progress = Math.max(0, Math.min(1, (input.pointerY - input.railTop) / input.railHeight)); - const targetPosition = progress * Math.max(0, firstItem.positionCount - 1); - - let nearestItemIndex = 0; - let nearestDistance = Number.POSITIVE_INFINITY; - for (let index = 0; index < input.items.length; index += 1) { - const item = input.items[index]; - if (!item) { - continue; - } - const distance = Math.abs(item.positionIndex - targetPosition); - if (distance < nearestDistance) { - nearestItemIndex = index; - nearestDistance = distance; - } - } - return nearestItemIndex; -} - export function resolveTimelineMinimapHasPersistentGutter(viewportWidth: number): boolean { if (!Number.isFinite(viewportWidth) || viewportWidth <= 0) { return false; @@ -295,108 +252,6 @@ export interface StableMessagesTimelineRowsState { result: MessagesTimelineRow[]; } -export type TimelineMinimapKind = "user-turn" | "final-assistant"; - -export interface TimelineMinimapItem { - readonly id: string; - readonly rowIndex: number; - readonly positionIndex: number; - readonly positionCount: number; - readonly primaryText: string | null; - readonly secondaryText: string | null; -} - -export function deriveTimelineMinimapItems( - rows: ReadonlyArray, - kind: TimelineMinimapKind, -): TimelineMinimapItem[] { - const items: TimelineMinimapItem[] = []; - const positionCount = Math.max( - 1, - rows.filter((row) => row.kind === "message" && row.message.role === "user").length, - ); - let positionIndex = -1; - for (let index = 0; index < rows.length; index += 1) { - const row = rows[index]; - if (row?.kind !== "message") { - continue; - } - if (row.message.role === "user") { - positionIndex += 1; - } - - if (kind === "final-assistant") { - if (row.message.role !== "assistant" || !row.isFinalAssistantResponse) { - continue; - } - items.push({ - id: row.id, - rowIndex: index, - positionIndex: Math.max(0, positionIndex), - positionCount, - primaryText: compactMinimapPreview(row.message.text), - secondaryText: null, - }); - continue; - } - - if (row.message.role !== "user") { - continue; - } - items.push({ - id: row.id, - rowIndex: index, - positionIndex: Math.max(0, positionIndex), - positionCount, - primaryText: compactMinimapPreview(row.message.text), - secondaryText: compactMinimapPreview(resolveFinalAssistantTextForTurn(rows, index)), - }); - } - return items; -} - -function resolveFinalAssistantTextForTurn( - rows: ReadonlyArray, - userRowIndex: number, -) { - let finalAssistantText: string | null = null; - for (let index = userRowIndex + 1; index < rows.length; index += 1) { - const row = rows[index]; - if (row?.kind !== "message") { - continue; - } - if (row.message.role === "user") { - break; - } - if (row.message.role === "assistant") { - finalAssistantText = row.message.text ?? null; - } - } - return finalAssistantText; -} - -function compactMinimapPreview(text: string | null | undefined) { - const compact = text?.replace(/\s+/g, " ").trim() ?? ""; - return compact.length > 0 ? compact : null; -} - -export function resolveTimelineMinimapAriaLabel(input: { - readonly activeIndex: number | null; - readonly itemCount: number; - readonly primaryText: string | null; - readonly side: "left" | "right"; -}): string { - const fallback = input.side === "left" ? "User message" : "Agent response"; - if (input.activeIndex === null) { - return `Jump to message: ${fallback}`; - } - - const role = input.side === "left" ? "user message" : "agent response"; - const position = `${input.activeIndex + 1} of ${input.itemCount}`; - const excerpt = input.primaryText?.slice(0, TIMELINE_MINIMAP_ARIA_EXCERPT_MAX_LENGTH) ?? null; - return excerpt ? `Jump to ${role} ${position}: ${excerpt}` : `Jump to ${role} ${position}`; -} - export function computeMessageDurationStart( messages: ReadonlyArray, ): Map { @@ -671,13 +526,20 @@ function timelineEntryTurnId(entry: TimelineEntry): TurnId | null { * Everything between them folds behind a "Worked for ..." row anchored at * the first hidden entry. Keeping both ends prevents a short follow-up from * hiding a substantive opening response while still bounding noisy turns. + * + * Also reports which preserved opening responses differ from their turn's + * terminal message, so they can carry the assistant meta row like the + * first-class answers they are. */ function deriveTurnFolds(input: { timelineEntries: ReadonlyArray; terminalAssistantMessageIds: ReadonlySet; latestTurn: TimelineLatestTurn | null; unsettledTurnId: TurnId | null; -}): ReadonlyMap { +}): { + foldsByAnchorEntryId: ReadonlyMap; + settledTurnOpeningAssistantMessageIds: ReadonlySet; +} { interface TurnGroup { entries: Array; terminalEntry: Extract | null; @@ -733,6 +595,7 @@ function deriveTurnFolds(input: { } const foldsByAnchorEntryId = new Map(); + const settledTurnOpeningAssistantMessageIds = new Set(); for (const [turnId, group] of groupsByTurnId) { if (turnId === input.unsettledTurnId) { continue; @@ -743,6 +606,9 @@ function deriveTurnFolds(input: { const firstAssistantEntry = group.entries.find( (entry): entry is Extract => entry.kind === "message", ); + if (firstAssistantEntry && firstAssistantEntry.message.id !== group.terminalEntry?.message.id) { + settledTurnOpeningAssistantMessageIds.add(firstAssistantEntry.message.id); + } const hiddenEntryIds = new Set(); for (const entry of group.entries) { if (entry.id === firstAssistantEntry?.id || entry.id === group.terminalEntry?.id) { @@ -800,7 +666,7 @@ function deriveTurnFolds(input: { label, }); } - return foldsByAnchorEntryId; + return { foldsByAnchorEntryId, settledTurnOpeningAssistantMessageIds }; } export function deriveMessagesTimelineRows(input: { @@ -824,7 +690,7 @@ export function deriveMessagesTimelineRows(input: { input.latestTurn ?? null, input.runningTurnId ?? null, ); - const foldsByAnchorEntryId = deriveTurnFolds({ + const { foldsByAnchorEntryId, settledTurnOpeningAssistantMessageIds } = deriveTurnFolds({ timelineEntries: input.timelineEntries, terminalAssistantMessageIds, latestTurn: input.latestTurn ?? null, @@ -1151,10 +1017,13 @@ export function deriveMessagesTimelineRows(input: { // While the turn is still running, the latest assistant message is only // provisionally terminal — withhold the metadata row until the turn - // settles so commentary doesn't flash timestamps mid-work. + // settles so commentary doesn't flash timestamps mid-work. Once settled, + // the turn's opening response earns the row too whenever it differs from + // the terminal message — with or without a fold between them. const showAssistantMeta = timelineEntry.message.role === "assistant" && - terminalAssistantMessageIds.has(timelineEntry.message.id) && + (terminalAssistantMessageIds.has(timelineEntry.message.id) || + settledTurnOpeningAssistantMessageIds.has(timelineEntry.message.id)) && !assistantTurnStillInProgress; const assistantTurnDiffSummary = timelineEntry.message.role === "assistant" diff --git a/apps/web/src/components/chat/MessagesTimeline.minimap.test.ts b/apps/web/src/components/chat/MessagesTimeline.minimap.test.ts new file mode 100644 index 000000000000..708a9bc25f65 --- /dev/null +++ b/apps/web/src/components/chat/MessagesTimeline.minimap.test.ts @@ -0,0 +1,372 @@ +import { describe, expect, it } from "vite-plus/test"; +import type { MessagesTimelineRow } from "./MessagesTimeline.logic"; +import { + computeTimelineMinimapState, + deriveTimelineMinimapItems, + EMPTY_TIMELINE_MINIMAP_STATE, + resolveTimelineMinimapAriaLabel, + resolveTimelineMinimapItemIndexFromPointer, + resolveTimelineMinimapTooltipTranslate, + resolveTimelineMinimapVisibleItemIds, + TIMELINE_MINIMAP_PREVIEW_MAX_LENGTH, + type TimelineMinimapItem, +} from "./MessagesTimeline.minimap"; + +function messageRow(input: { + id: string; + role: "user" | "assistant"; + text: string; + final?: boolean; + streaming?: boolean; +}): Extract { + return { + kind: "message", + id: input.id, + message: { + id: input.id, + role: input.role, + text: input.text, + streaming: input.streaming ?? false, + }, + isFinalAssistantResponse: input.final ?? false, + } as never; +} + +describe("resolveTimelineMinimapAriaLabel", () => { + it("uses a concise role and position label with a capped excerpt", () => { + const longResponse = "A".repeat(180); + + expect( + resolveTimelineMinimapAriaLabel({ + activeIndex: null, + itemCount: 3, + primaryText: null, + side: "right", + }), + ).toBe("Jump to message: Agent response"); + expect( + resolveTimelineMinimapAriaLabel({ + activeIndex: 1, + itemCount: 3, + primaryText: longResponse, + side: "right", + }), + ).toBe(`Jump to agent response 2 of 3: ${"A".repeat(120)}`); + }); +}); + +describe("resolveTimelineMinimapTooltipTranslate", () => { + it("only edge-clamps markers at the endpoints of the full turn space", () => { + expect(resolveTimelineMinimapTooltipTranslate(0, 5)).toBe("0%"); + expect(resolveTimelineMinimapTooltipTranslate(1, 5)).toBe("-50%"); + expect(resolveTimelineMinimapTooltipTranslate(3, 5)).toBe("-50%"); + expect(resolveTimelineMinimapTooltipTranslate(4, 5)).toBe("-100%"); + }); +}); + +describe("deriveTimelineMinimapItems", () => { + it("keeps sparse agent replies aligned with their user-turn positions", () => { + const rows = [ + messageRow({ id: "user-1", role: "user", text: "user-1" }), + messageRow({ id: "assistant-1", role: "assistant", text: "assistant-1", final: true }), + messageRow({ id: "user-2", role: "user", text: "user-2" }), + messageRow({ + id: "assistant-commentary", + role: "assistant", + text: "assistant-commentary", + }), + messageRow({ id: "user-3", role: "user", text: "user-3" }), + messageRow({ id: "assistant-3", role: "assistant", text: "assistant-3", final: true }), + ]; + + expect( + deriveTimelineMinimapItems(rows, "final-assistant").map( + ({ id, positionIndex, positionCount }) => ({ id, positionIndex, positionCount }), + ), + ).toEqual([ + { id: "assistant-1", positionIndex: 0, positionCount: 3 }, + { id: "assistant-3", positionIndex: 2, positionCount: 3 }, + ]); + }); +}); + +describe("computeTimelineMinimapState", () => { + it("reuses historical previews and the unaffected rail during streaming", () => { + const rows = [ + messageRow({ id: "user-1", role: "user", text: " First\n prompt " }), + messageRow({ + id: "assistant-1", + role: "assistant", + text: "First final response.", + final: true, + }), + messageRow({ id: "user-2", role: "user", text: "Second prompt." }), + messageRow({ + id: "assistant-stream", + role: "assistant", + text: "Streaming one \n", + streaming: true, + }), + ]; + const initial = computeTimelineMinimapState(rows, EMPTY_TIMELINE_MINIMAP_STATE); + const updated = computeTimelineMinimapState( + [ + ...rows.slice(0, -1), + messageRow({ + id: "assistant-stream", + role: "assistant", + text: "Streaming one \n more chunk", + streaming: true, + }), + ], + initial, + ); + + expect(updated.previewByRowId.get("user-1")).toBe(initial.previewByRowId.get("user-1")); + expect(updated.previewByRowId.get("assistant-1")).toBe( + initial.previewByRowId.get("assistant-1"), + ); + expect(updated.previewByRowId.get("user-2")).toBe(initial.previewByRowId.get("user-2")); + expect(updated.previewByRowId.get("assistant-stream")).not.toBe( + initial.previewByRowId.get("assistant-stream"), + ); + expect(updated.userItems[0]).toBe(initial.userItems[0]); + expect(updated.userItems[1]).not.toBe(initial.userItems[1]); + expect(updated.userItems[1]?.secondaryText).toBe("Streaming one more chunk"); + expect(updated.assistantItems).toBe(initial.assistantItems); + }); + + it("caps compact previews and reuses marker output after the cap", () => { + const prefix = "A".repeat(TIMELINE_MINIMAP_PREVIEW_MAX_LENGTH - 2); + const user = messageRow({ id: "user", role: "user", text: "Prompt" }); + const initial = computeTimelineMinimapState( + [ + user, + messageRow({ + id: "assistant", + role: "assistant", + text: `${prefix} `, + streaming: true, + }), + ], + EMPTY_TIMELINE_MINIMAP_STATE, + ); + const capped = computeTimelineMinimapState( + [ + user, + messageRow({ + id: "assistant", + role: "assistant", + text: `${prefix} word`, + streaming: true, + }), + ], + initial, + ); + const afterCap = computeTimelineMinimapState( + [ + user, + messageRow({ + id: "assistant", + role: "assistant", + text: `${prefix} word more text`, + streaming: true, + }), + ], + capped, + ); + + expect(capped.userItems[0]?.secondaryText).toBe(`${prefix} w`); + expect(capped.userItems[0]?.secondaryText).toHaveLength(TIMELINE_MINIMAP_PREVIEW_MAX_LENGTH); + expect(afterCap.userItems[0]).toBe(capped.userItems[0]); + expect(afterCap.previewByRowId.get("assistant")).not.toBe( + capped.previewByRowId.get("assistant"), + ); + }); + + it("provides a new visibility refresh token for same-length layout changes", () => { + const user = messageRow({ id: "user", role: "user", text: "Prompt" }); + const assistant = messageRow({ + id: "assistant", + role: "assistant", + text: "Response", + final: true, + }); + const workRow = (label: string) => + ({ + kind: "work", + id: "work", + createdAt: "2026-08-23T00:00:00.000Z", + groupedEntries: [{ label }], + isExpandedToolGroupEntry: false, + isLastExpandedToolGroupEntry: false, + }) as never; + const initial = computeTimelineMinimapState( + [user, workRow("Short"), assistant], + EMPTY_TIMELINE_MINIMAP_STATE, + ); + const updated = computeTimelineMinimapState( + [user, workRow("A much longer work row that can change layout"), assistant], + initial, + ); + + expect(updated).not.toBe(initial); + expect(updated.userItems).toBe(initial.userItems); + expect(updated.assistantItems).toBe(initial.assistantItems); + }); + + it("bounds its preview cache to messages represented by the loaded rows", () => { + const initial = computeTimelineMinimapState( + [ + messageRow({ id: "user-1", role: "user", text: "First" }), + messageRow({ id: "assistant-1", role: "assistant", text: "Done" }), + ], + EMPTY_TIMELINE_MINIMAP_STATE, + ); + const replaced = computeTimelineMinimapState( + [messageRow({ id: "user-2", role: "user", text: "Second" })], + initial, + ); + + expect([...replaced.previewByRowId.keys()]).toEqual(["user-2"]); + }); +}); + +describe("resolveTimelineMinimapItemIndexFromPointer", () => { + it("selects the nearest rendered reply in the full turn-position space", () => { + const items = [ + { positionIndex: 0, positionCount: 4 }, + { positionIndex: 3, positionCount: 4 }, + ]; + + expect( + resolveTimelineMinimapItemIndexFromPointer({ + items, + railTop: 100, + railHeight: 300, + pointerY: 360, + }), + ).toBe(1); + expect( + resolveTimelineMinimapItemIndexFromPointer({ + items, + railTop: 100, + railHeight: 300, + pointerY: 160, + }), + ).toBe(0); + }); + + it("compares sparse markers against the continuous pointer position", () => { + expect( + resolveTimelineMinimapItemIndexFromPointer({ + items: [ + { positionIndex: 0, positionCount: 4 }, + { positionIndex: 2, positionCount: 4 }, + ], + railTop: 0, + railHeight: 24, + pointerY: 11, + }), + ).toBe(1); + }); + + it("uses logarithmic lookup in the sorted turn-position space", () => { + let positionReads = 0; + const itemCount = 1_024; + const items = Array.from({ length: itemCount }, (_, index) => ({ + positionCount: itemCount, + get positionIndex() { + positionReads += 1; + return index; + }, + })); + + expect( + resolveTimelineMinimapItemIndexFromPointer({ + items, + railTop: 0, + railHeight: itemCount - 1, + pointerY: 713.4, + }), + ).toBe(713); + expect(positionReads).toBeLessThan(30); + }); + + it("keeps the earlier marker on an exact midpoint tie", () => { + expect( + resolveTimelineMinimapItemIndexFromPointer({ + items: [ + { positionCount: 4, positionIndex: 0 }, + { positionCount: 4, positionIndex: 2 }, + ], + railTop: 0, + railHeight: 2, + pointerY: 2 / 3, + }), + ).toBe(0); + }); +}); + +describe("resolveTimelineMinimapVisibleItemIds", () => { + it("measures only the binary-search path and visible markers", () => { + let positionReads = 0; + const items: TimelineMinimapItem[] = Array.from({ length: 1_024 }, (_, index) => ({ + id: `item-${index}`, + rowIndex: index * 2, + positionIndex: index, + positionCount: 1_024, + primaryText: null, + secondaryText: null, + })); + + expect( + resolveTimelineMinimapVisibleItemIds({ + items, + state: { + positionAtIndex: (rowIndex) => { + positionReads += 1; + return rowIndex * 10; + }, + sizeAtIndex: () => 5, + }, + scrollTop: 10_000, + scrollBottom: 10_025, + }), + ).toEqual(["item-500", "item-501"]); + expect(positionReads).toBeLessThan(30); + }); + + it("falls back to a full scan when layout positions are incomplete", () => { + const items: TimelineMinimapItem[] = [ + { + id: "item-0", + rowIndex: 0, + positionIndex: 0, + positionCount: 2, + primaryText: null, + secondaryText: null, + }, + { + id: "item-1", + rowIndex: 1, + positionIndex: 1, + positionCount: 2, + primaryText: null, + secondaryText: null, + }, + ]; + + expect( + resolveTimelineMinimapVisibleItemIds({ + items, + state: { + positionAtIndex: (rowIndex) => (rowIndex === 0 ? 0 : undefined), + sizeAtIndex: () => 20, + }, + scrollTop: 0, + scrollBottom: 10, + }), + ).toEqual(["item-0"]); + }); +}); diff --git a/apps/web/src/components/chat/MessagesTimeline.minimap.ts b/apps/web/src/components/chat/MessagesTimeline.minimap.ts new file mode 100644 index 000000000000..4b1d5ac7148b --- /dev/null +++ b/apps/web/src/components/chat/MessagesTimeline.minimap.ts @@ -0,0 +1,389 @@ +import type { MessagesTimelineRow } from "./MessagesTimeline.logic"; + +export const TIMELINE_MINIMAP_ARIA_EXCERPT_MAX_LENGTH = 120; +export const TIMELINE_MINIMAP_PREVIEW_MAX_LENGTH = 240; + +export type TimelineMinimapKind = "user-turn" | "final-assistant"; + +export interface TimelineMinimapItem { + readonly id: string; + readonly rowIndex: number; + readonly positionIndex: number; + readonly positionCount: number; + readonly primaryText: string | null; + readonly secondaryText: string | null; +} + +interface TimelineMinimapPreviewCacheEntry { + readonly contentVersion: string | null | undefined; + readonly compactText: string | null; + readonly streaming: boolean; + readonly endsWithWhitespace: boolean; +} + +export interface TimelineMinimapState { + readonly previewByRowId: ReadonlyMap; + readonly userItems: ReadonlyArray; + readonly assistantItems: ReadonlyArray; +} + +export const EMPTY_TIMELINE_MINIMAP_STATE: TimelineMinimapState = { + previewByRowId: new Map(), + userItems: [], + assistantItems: [], +}; + +type MessageRow = Extract; + +function compactMinimapPreview(text: string | null | undefined) { + const compact = text?.replace(/\s+/g, " ").trim() ?? ""; + return compact.length > 0 ? compact.slice(0, TIMELINE_MINIMAP_PREVIEW_MAX_LENGTH) : null; +} + +function textEndsWithWhitespace(text: string | null | undefined) { + if (!text || text.length === 0) { + return false; + } + return /\s/u.test(text.at(-1)!); +} + +function appendCompactMinimapPreview( + cached: TimelineMinimapPreviewCacheEntry, + contentVersion: string, +): string | null { + if ((cached.compactText?.length ?? 0) >= TIMELINE_MINIMAP_PREVIEW_MAX_LENGTH) { + return cached.compactText; + } + + const previousText = cached.contentVersion; + if (typeof previousText !== "string") { + return compactMinimapPreview(contentVersion); + } + const appendedText = contentVersion.slice(previousText.length); + const compactAppend = compactMinimapPreview(appendedText); + if (compactAppend === null) { + return cached.compactText; + } + if (cached.compactText === null) { + return compactAppend; + } + const needsSpace = cached.endsWithWhitespace || /^\s/u.test(appendedText); + return `${cached.compactText}${needsSpace ? " " : ""}${compactAppend}`.slice( + 0, + TIMELINE_MINIMAP_PREVIEW_MAX_LENGTH, + ); +} + +function resolveCachedPreview( + row: MessageRow, + previous: TimelineMinimapState, + nextPreviewByRowId: Map, +): string | null { + const contentVersion = row.message.text; + const streaming = row.message.streaming; + const cached = nextPreviewByRowId.get(row.id) ?? previous.previewByRowId.get(row.id); + let next: TimelineMinimapPreviewCacheEntry; + if (cached?.contentVersion === contentVersion) { + next = cached.streaming === streaming ? cached : { ...cached, streaming }; + } else { + // The message projector constructs a stable streaming message by appending + // each delta, so a longer in-progress version can normalize only its tail. + const canAppendStreamingTail = + cached !== undefined && + cached.streaming && + streaming && + typeof cached.contentVersion === "string" && + typeof contentVersion === "string" && + contentVersion.length >= cached.contentVersion.length; + next = { + contentVersion, + compactText: canAppendStreamingTail + ? appendCompactMinimapPreview(cached, contentVersion) + : compactMinimapPreview(contentVersion), + streaming, + endsWithWhitespace: textEndsWithWhitespace(contentVersion), + }; + } + nextPreviewByRowId.set(row.id, next); + return next.compactText; +} + +function reuseMinimapItem( + previousById: ReadonlyMap, + item: TimelineMinimapItem, +): TimelineMinimapItem { + const previous = previousById.get(item.id); + return previous && + previous.rowIndex === item.rowIndex && + previous.positionIndex === item.positionIndex && + previous.positionCount === item.positionCount && + previous.primaryText === item.primaryText && + previous.secondaryText === item.secondaryText + ? previous + : item; +} + +function reuseMinimapItems( + previous: ReadonlyArray, + next: TimelineMinimapItem[], +): ReadonlyArray { + return previous.length === next.length && next.every((item, index) => previous[index] === item) + ? previous + : next; +} + +/** + * Derives both rails in one shared pass after counting turn positions and + * reuses compact previews for unchanged message text. Streaming normalizes + * only the appended response tail; historical previews and the unaffected rail + * retain their identity. + */ +export function computeTimelineMinimapState( + rows: ReadonlyArray, + previous: TimelineMinimapState, +): TimelineMinimapState { + const positionCount = Math.max( + 1, + rows.reduce( + (count, row) => count + (row.kind === "message" && row.message.role === "user" ? 1 : 0), + 0, + ), + ); + const previousUserById = new Map(previous.userItems.map((item) => [item.id, item])); + const previousAssistantById = new Map(previous.assistantItems.map((item) => [item.id, item])); + const nextPreviewByRowId = new Map(); + const userItems: TimelineMinimapItem[] = []; + const assistantItems: TimelineMinimapItem[] = []; + let positionIndex = -1; + let currentUserRow: MessageRow | null = null; + let currentUserRowIndex = -1; + let currentTurnLastAssistantRow: MessageRow | null = null; + + const appendCurrentUserItem = () => { + if (currentUserRow === null) { + return; + } + userItems.push( + reuseMinimapItem(previousUserById, { + id: currentUserRow.id, + rowIndex: currentUserRowIndex, + positionIndex: Math.max(0, positionIndex), + positionCount, + primaryText: resolveCachedPreview(currentUserRow, previous, nextPreviewByRowId), + secondaryText: + currentTurnLastAssistantRow === null + ? null + : resolveCachedPreview(currentTurnLastAssistantRow, previous, nextPreviewByRowId), + }), + ); + }; + + for (let rowIndex = 0; rowIndex < rows.length; rowIndex += 1) { + const row = rows[rowIndex]; + if (row?.kind !== "message") { + continue; + } + + if (row.message.role === "user") { + appendCurrentUserItem(); + positionIndex += 1; + currentUserRow = row; + currentUserRowIndex = rowIndex; + currentTurnLastAssistantRow = null; + continue; + } + + if (row.message.role !== "assistant") { + continue; + } + + currentTurnLastAssistantRow = row; + if (row.isFinalAssistantResponse) { + assistantItems.push( + reuseMinimapItem(previousAssistantById, { + id: row.id, + rowIndex, + positionIndex: Math.max(0, positionIndex), + positionCount, + primaryText: resolveCachedPreview(row, previous, nextPreviewByRowId), + secondaryText: null, + }), + ); + } + } + appendCurrentUserItem(); + + return { + previewByRowId: nextPreviewByRowId, + userItems: reuseMinimapItems(previous.userItems, userItems), + assistantItems: reuseMinimapItems(previous.assistantItems, assistantItems), + }; +} + +export function deriveTimelineMinimapItems( + rows: ReadonlyArray, + kind: TimelineMinimapKind, +): ReadonlyArray { + const state = computeTimelineMinimapState(rows, EMPTY_TIMELINE_MINIMAP_STATE); + return kind === "user-turn" ? state.userItems : state.assistantItems; +} + +export function resolveTimelineMinimapTooltipTranslate( + positionIndex: number, + positionCount: number, +): string { + if (positionIndex === 0) { + return "0%"; + } + if (positionIndex === positionCount - 1) { + return "-100%"; + } + return "-50%"; +} + +export function resolveTimelineMinimapItemIndexFromPointer(input: { + readonly items: ReadonlyArray>; + readonly railTop: number; + readonly railHeight: number; + readonly pointerY: number; +}): number | null { + const firstItem = input.items[0]; + if (!firstItem || input.railHeight <= 0) { + return null; + } + const progress = Math.max(0, Math.min(1, (input.pointerY - input.railTop) / input.railHeight)); + const targetPosition = progress * Math.max(0, firstItem.positionCount - 1); + + let low = 0; + let high = input.items.length; + while (low < high) { + const middle = low + Math.floor((high - low) / 2); + if (input.items[middle]!.positionIndex < targetPosition) { + low = middle + 1; + } else { + high = middle; + } + } + + if (low === 0) { + return 0; + } + if (low === input.items.length) { + return input.items.length - 1; + } + + const previousIndex = low - 1; + const previousDistance = targetPosition - input.items[previousIndex]!.positionIndex; + const nextDistance = input.items[low]!.positionIndex - targetPosition; + return previousDistance <= nextDistance ? previousIndex : low; +} + +export function resolveTimelineMinimapAriaLabel(input: { + readonly activeIndex: number | null; + readonly itemCount: number; + readonly primaryText: string | null; + readonly side: "left" | "right"; +}): string { + const fallback = input.side === "left" ? "User message" : "Agent response"; + if (input.activeIndex === null) { + return `Jump to message: ${fallback}`; + } + + const role = input.side === "left" ? "user message" : "agent response"; + const position = `${input.activeIndex + 1} of ${input.itemCount}`; + const excerpt = input.primaryText?.slice(0, TIMELINE_MINIMAP_ARIA_EXCERPT_MAX_LENGTH) ?? null; + return excerpt ? `Jump to ${role} ${position}: ${excerpt}` : `Jump to ${role} ${position}`; +} + +export interface TimelineMinimapPositionState { + readonly positionAtIndex?: (index: number) => number | undefined; + readonly sizeAtIndex?: (index: number) => number | undefined; +} + +function resolveTimelineRowTop(state: TimelineMinimapPositionState, rowIndex: number) { + const top = state.positionAtIndex?.(rowIndex); + return typeof top === "number" && Number.isFinite(top) ? top : null; +} + +function resolveTimelineRowHeight(state: TimelineMinimapPositionState, rowIndex: number) { + const height = state.sizeAtIndex?.(rowIndex); + return typeof height === "number" && Number.isFinite(height) ? height : null; +} + +function resolveVisibleItemIdsByScan(input: { + readonly items: ReadonlyArray; + readonly state: TimelineMinimapPositionState; + readonly scrollTop: number; + readonly scrollBottom: number; +}): string[] { + const visibleIds: string[] = []; + for (const item of input.items) { + const rowTop = resolveTimelineRowTop(input.state, item.rowIndex); + const rowHeight = resolveTimelineRowHeight(input.state, item.rowIndex); + if ( + rowTop !== null && + rowTop < input.scrollBottom && + rowTop + Math.max(1, rowHeight ?? 1) > input.scrollTop + ) { + visibleIds.push(item.id); + } + } + return visibleIds; +} + +/** Finds visible markers in O(log n + visible markers), with a scan fallback for incomplete layouts. */ +export function resolveTimelineMinimapVisibleItemIds(input: { + readonly items: ReadonlyArray; + readonly state: TimelineMinimapPositionState; + readonly scrollTop: number; + readonly scrollBottom: number; +}): ReadonlyArray { + if (input.items.length === 0 || input.state.positionAtIndex === undefined) { + return []; + } + + let low = 0; + let high = input.items.length; + while (low < high) { + const middle = low + Math.floor((high - low) / 2); + const rowTop = resolveTimelineRowTop(input.state, input.items[middle]!.rowIndex); + if (rowTop === null) { + return resolveVisibleItemIdsByScan(input); + } + if (rowTop < input.scrollTop) { + low = middle + 1; + } else { + high = middle; + } + } + + const visibleIds: string[] = []; + const precedingIndex = low - 1; + if (precedingIndex >= 0) { + const precedingItem = input.items[precedingIndex]!; + const rowTop = resolveTimelineRowTop(input.state, precedingItem.rowIndex); + const rowHeight = resolveTimelineRowHeight(input.state, precedingItem.rowIndex); + if ( + rowTop === null || + (rowTop < input.scrollBottom && rowTop + Math.max(1, rowHeight ?? 1) > input.scrollTop) + ) { + if (rowTop === null) { + return resolveVisibleItemIdsByScan(input); + } + visibleIds.push(precedingItem.id); + } + } + + for (let index = low; index < input.items.length; index += 1) { + const item = input.items[index]!; + const rowTop = resolveTimelineRowTop(input.state, item.rowIndex); + if (rowTop === null) { + return resolveVisibleItemIdsByScan(input); + } + if (rowTop >= input.scrollBottom) { + break; + } + visibleIds.push(item.id); + } + return visibleIds; +} diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 63ebd9ca517e..869a1c13b80f 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -89,7 +89,6 @@ import { keepTimelineEndVisibleAfterOverlayGrowth } from "./timelineScrollAnchor import { MessageCopyButton } from "./MessageCopyButton"; import { computeStableMessagesTimelineRows, - deriveTimelineMinimapItems, deriveMessagesTimelineRows, normalizeCompactToolLabel, resolveAssistantMessageCopyState, @@ -97,20 +96,26 @@ import { resolveTimelineMinimapHasPersistentGutter, resolveTimelineMinimapHeightStyle, resolveTimelineMinimapHitStripWidth, - resolveTimelineMinimapItemIndexFromPointer, resolveTimelineMinimapInteractiveWidth, resolveTimelineMinimapTopPercent, - resolveTimelineMinimapTooltipTranslate, - resolveTimelineMinimapAriaLabel, shouldPreserveAssistantLineBreaks, + TIMELINE_MINIMAP_MIN_ITEMS, toolGroupAction, workEntryIsVisibleInGroup, type StableMessagesTimelineRowsState, type MessagesTimelineRow, - TIMELINE_MINIMAP_MIN_ITEMS, - type TimelineMinimapItem, type TimelineLatestTurn, } from "./MessagesTimeline.logic"; +import { + computeTimelineMinimapState, + EMPTY_TIMELINE_MINIMAP_STATE, + resolveTimelineMinimapAriaLabel, + resolveTimelineMinimapItemIndexFromPointer, + resolveTimelineMinimapTooltipTranslate, + resolveTimelineMinimapVisibleItemIds, + type TimelineMinimapItem, + type TimelineMinimapPositionState, +} from "./MessagesTimeline.minimap"; import { TerminalContextInlineChip } from "./TerminalContextInlineChip"; import { Tooltip, TooltipPopup, TooltipTrigger } from "../ui/tooltip"; import { @@ -338,6 +343,8 @@ export const MessagesTimeline = memo(function MessagesTimeline({ const [disclosureToggleSettling, setDisclosureToggleSettling] = useState(false); const [userMinimapStripMap] = useState(() => new Map()); const [assistantMinimapStripMap] = useState(() => new Map()); + const userMinimapInViewIdsRef = useRef>(new Set()); + const assistantMinimapInViewIdsRef = useRef>(new Set()); const disclosureAnchorKeyRef = useRef(null); const disclosureSettleFrameRef = useRef(null); const disclosureSettleSecondFrameRef = useRef(null); @@ -488,11 +495,9 @@ export const MessagesTimeline = memo(function MessagesTimeline({ ], ); const rows = useStableRows(rawRows); - const userMinimapItems = useMemo(() => deriveTimelineMinimapItems(rows, "user-turn"), [rows]); - const assistantMinimapItems = useMemo( - () => deriveTimelineMinimapItems(rows, "final-assistant"), - [rows], - ); + const minimapState = useTimelineMinimapState(rows); + const userMinimapItems = minimapState.userItems; + const assistantMinimapItems = minimapState.assistantItems; const [timelineViewportElement, setTimelineViewportElement] = useState( null, ); @@ -526,26 +531,22 @@ export const MessagesTimeline = memo(function MessagesTimeline({ const scrollTop = state.scroll ?? 0; const scrollBottom = scrollTop + (state.scrollLength ?? 0); - for (const [items, stripMap] of [ - [userMinimapItems, userMinimapStripMap], - [assistantMinimapItems, assistantMinimapStripMap], - ] as const) { - for (const item of items) { - const strip = stripMap.get(item.id); - if (!strip) { - continue; - } - - const rowTop = resolveTimelineRowTop(state, item.rowIndex); - const rowHeight = resolveTimelineRowHeight(state, item.rowIndex); - const inView = - rowTop !== null && - rowTop < scrollBottom && - rowTop + Math.max(1, rowHeight ?? 1) > scrollTop; - - strip.dataset.inView = inView ? "true" : "false"; - } - } + userMinimapInViewIdsRef.current = updateTimelineMinimapInView({ + items: userMinimapItems, + previousIds: userMinimapInViewIdsRef.current, + scrollBottom, + scrollTop, + state, + stripMap: userMinimapStripMap, + }); + assistantMinimapInViewIdsRef.current = updateTimelineMinimapInView({ + items: assistantMinimapItems, + previousIds: assistantMinimapInViewIdsRef.current, + scrollBottom, + scrollTop, + state, + stripMap: assistantMinimapStripMap, + }); }, [ assistantMinimapItems, assistantMinimapStripMap, @@ -557,9 +558,13 @@ export const MessagesTimeline = memo(function MessagesTimeline({ ]); useEffect(() => { + // Depends on the state object, not just handleScroll: row content growth + // (streaming text, live work rows) moves markers without changing the + // reference-stable item arrays, and only a fresh derivation signals it. + // rAF coalesces the recompute to once per frame while content grows. const frame = requestAnimationFrame(() => handleScroll()); return () => cancelAnimationFrame(frame); - }, [handleScroll, rows.length]); + }, [handleScroll, minimapState]); useEffect(() => { if (!timelineViewportElement) { @@ -652,6 +657,17 @@ export const MessagesTimeline = memo(function MessagesTimeline({ ), [], ); + const handleMinimapSelect = useCallback( + (item: TimelineMinimapItem) => { + onManualNavigation(); + void listRef.current?.scrollToIndex({ + index: item.rowIndex, + animated: true, + viewOffset: 24, + }); + }, + [listRef, onManualNavigation], + ); if (rows.length === 0 && !isWorking) { if (hideEmptyPlaceholder) { @@ -712,14 +728,7 @@ export const MessagesTimeline = memo(function MessagesTimeline({ side="left" stripMap={userMinimapStripMap} testId="timeline-minimap" - onSelect={(item) => { - onManualNavigation(); - void listRef.current?.scrollToIndex({ - index: item.rowIndex, - animated: true, - viewOffset: 24, - }); - }} + onSelect={handleMinimapSelect} /> { - onManualNavigation(); - void listRef.current?.scrollToIndex({ - index: item.rowIndex, - animated: true, - viewOffset: 24, - }); - }} + onSelect={handleMinimapSelect} /> @@ -751,29 +753,47 @@ function getItemType(item: MessagesTimelineRow) { return item.kind === "message" ? `message:${item.message.role}` : item.kind; } -interface TimelinePositionState { - readonly contentLength?: number; - readonly scroll?: number; - readonly scrollLength?: number; - readonly positionAtIndex?: (index: number) => number | undefined; - readonly sizeAtIndex?: (index: number) => number | undefined; -} - -function resolveTimelineRowTop(state: TimelinePositionState, rowIndex: number) { - const top = state.positionAtIndex?.(rowIndex); - return typeof top === "number" && Number.isFinite(top) ? top : null; -} - -function resolveTimelineRowHeight(state: TimelinePositionState, rowIndex: number) { - const height = state.sizeAtIndex?.(rowIndex); - return typeof height === "number" && Number.isFinite(height) ? height : null; +function updateTimelineMinimapInView(input: { + readonly items: ReadonlyArray; + readonly previousIds: ReadonlySet; + readonly scrollBottom: number; + readonly scrollTop: number; + readonly state: TimelineMinimapPositionState; + readonly stripMap: ReadonlyMap; +}): ReadonlySet { + const nextIds = new Set( + resolveTimelineMinimapVisibleItemIds({ + items: input.items, + scrollBottom: input.scrollBottom, + scrollTop: input.scrollTop, + state: input.state, + }), + ); + for (const id of input.previousIds) { + if (!nextIds.has(id)) { + const strip = input.stripMap.get(id); + if (strip) { + strip.dataset.inView = "false"; + } + } + } + for (const id of nextIds) { + // Idempotent on purpose: an id can enter previousIds before its strip + // mounts (the rail renders only past the minimum item count), so skipping + // already-known ids would leave that strip dim until it left the viewport. + const strip = input.stripMap.get(id); + if (strip && strip.dataset.inView !== "true") { + strip.dataset.inView = "true"; + } + } + return nextIds; } function timelineMinimapEventTargetsPreview(target: EventTarget): boolean { return target instanceof Element && target.closest("[data-minimap-preview]") !== null; } -function TimelineMinimap({ +const TimelineMinimap = memo(function TimelineMinimap({ hasPersistentGutter, hitStripWidth, items, @@ -993,7 +1013,7 @@ function TimelineMinimap({ ); -} +}); // --------------------------------------------------------------------------- // TimelineRowContent — the actual row component @@ -2548,6 +2568,15 @@ function useStableRows(rows: MessagesTimelineRow[]): MessagesTimelineRow[] { }, [rows]); } +function useTimelineMinimapState(rows: ReadonlyArray) { + const previousState = useRef(EMPTY_TIMELINE_MINIMAP_STATE); + return useMemo(() => { + const nextState = computeTimelineMinimapState(rows, previousState.current); + previousState.current = nextState; + return nextState; + }, [rows]); +} + // --------------------------------------------------------------------------- // Pure helpers // --------------------------------------------------------------------------- diff --git a/apps/web/src/components/settings/ProviderModelsSection.tsx b/apps/web/src/components/settings/ProviderModelsSection.tsx index c734ecd3ef36..591bc7c374b3 100644 --- a/apps/web/src/components/settings/ProviderModelsSection.tsx +++ b/apps/web/src/components/settings/ProviderModelsSection.tsx @@ -229,7 +229,19 @@ export function ProviderModelsSection({ const canMoveDown = nextModel !== undefined && favoriteModelSet.has(nextModel.slug) === isFavorite; const descriptors = caps?.optionDescriptors ?? []; - if (descriptors.some((descriptor) => descriptor.id === "fastMode")) { + // Codex models expose fast mode as a "Fast" service tier rather + // than a fastMode boolean. Match TraitsPicker's trigger display: + // codex-only, detected by the option label (tier ids vary). + if ( + descriptors.some( + (descriptor) => + descriptor.id === "fastMode" || + (driverKind === "codex" && + descriptor.id === "serviceTier" && + descriptor.type === "select" && + descriptor.options.some((option) => option.label === "Fast")), + ) + ) { capLabels.push("Fast mode"); } if (descriptors.some((descriptor) => descriptor.id === "thinking")) { diff --git a/apps/web/src/components/thread-split/threadSplitStore.test.ts b/apps/web/src/components/thread-split/threadSplitStore.test.ts index cb3de482726b..35b712ad6333 100644 --- a/apps/web/src/components/thread-split/threadSplitStore.test.ts +++ b/apps/web/src/components/thread-split/threadSplitStore.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it } from "vite-plus/test"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vite-plus/test"; import { EnvironmentId, ThreadId } from "@t3tools/contracts"; import type { ComposerHandleRef } from "../../composerHandleContext"; @@ -129,3 +129,49 @@ describe("clampSplitRatio", () => { expect(clampSplitRatio(Number.NaN)).toBe(0.5); }); }); + +// This suite runs without a DOM; the store only touches `window` through +// guarded accessors, so a plain stub is enough to model restricted hosts. +describe("blocked storage", () => { + const globalWithWindow = globalThis as { window?: unknown }; + + const importFreshStore = async () => { + vi.resetModules(); + return import("./threadSplitStore"); + }; + + afterEach(() => { + delete globalWithWindow.window; + vi.resetModules(); + }); + + it("falls back to the default ratio when the localStorage property throws", async () => { + globalWithWindow.window = { + get localStorage(): Storage { + throw new Error("SecurityError"); + }, + }; + const store = await importFreshStore(); + expect(store.useThreadSplitStore.getState().splitRatio).toBe(0.5); + // Writes must stay silent too. + store.useThreadSplitStore.getState().setSplitRatio(0.6); + expect(store.useThreadSplitStore.getState().splitRatio).toBe(0.6); + }); + + it("survives a storage whose operations throw after resolving", async () => { + globalWithWindow.window = { + localStorage: { + getItem(): string | null { + throw new Error("SecurityError"); + }, + setItem(): void { + throw new Error("SecurityError"); + }, + }, + }; + const store = await importFreshStore(); + expect(store.useThreadSplitStore.getState().splitRatio).toBe(0.5); + store.useThreadSplitStore.getState().setSplitRatio(0.6); + expect(store.useThreadSplitStore.getState().splitRatio).toBe(0.6); + }); +}); diff --git a/apps/web/src/components/thread-split/threadSplitStore.ts b/apps/web/src/components/thread-split/threadSplitStore.ts index 609a4f8c5b6f..11fe76ffd24f 100644 --- a/apps/web/src/components/thread-split/threadSplitStore.ts +++ b/apps/web/src/components/thread-split/threadSplitStore.ts @@ -30,8 +30,21 @@ export function clampSplitRatio(ratio: number): number { // `window` alone is not enough: the test runner and other embedded hosts define // a window without `localStorage`, and this store is constructed at import time. +// Restricted hosts go further: the property access and every Storage operation +// can throw a SecurityError. Dragging the splitter writes at pointer-move rate, +// so the handle resolves once and is dropped on the first failing operation +// instead of throwing (and catching) per move. +let resolvedStorage: Storage | null | undefined; + function splitRatioStorage(): Storage | null { - return typeof window === "undefined" ? null : (window.localStorage ?? null); + if (resolvedStorage === undefined) { + try { + resolvedStorage = typeof window === "undefined" ? null : (window.localStorage ?? null); + } catch { + resolvedStorage = null; + } + } + return resolvedStorage; } function readStoredSplitRatio(): number { @@ -39,7 +52,13 @@ function readStoredSplitRatio(): number { if (storage === null) { return DEFAULT_SPLIT_RATIO; } - const raw = storage.getItem(SPLIT_RATIO_STORAGE_KEY); + let raw: string | null; + try { + raw = storage.getItem(SPLIT_RATIO_STORAGE_KEY); + } catch { + resolvedStorage = null; + return DEFAULT_SPLIT_RATIO; + } // Number("") is 0, which would clamp to the minimum instead of the default. if (raw === null || raw.trim().length === 0) { return DEFAULT_SPLIT_RATIO; @@ -107,7 +126,13 @@ export const useThreadSplitStore = create((set, get) => ({ const next = clampSplitRatio(ratio); if (get().splitRatio === next) return; set({ splitRatio: next }); - splitRatioStorage()?.setItem(SPLIT_RATIO_STORAGE_KEY, String(next)); + // Persistence is best-effort: a full or read-only storage must not break + // the drag that triggered the write, and one failure stops further tries. + try { + splitRatioStorage()?.setItem(SPLIT_RATIO_STORAGE_KEY, String(next)); + } catch { + resolvedStorage = null; + } }, })); diff --git a/packages/contracts/src/providerRuntime.ts b/packages/contracts/src/providerRuntime.ts index 0e7a5d1a0668..7f12d6d9fb31 100644 --- a/packages/contracts/src/providerRuntime.ts +++ b/packages/contracts/src/providerRuntime.ts @@ -789,6 +789,7 @@ const RuntimeErrorPayload = Schema.Struct({ message: TrimmedNonEmptyStringSchema, class: Schema.optional(RuntimeErrorClass), detail: Schema.optional(Schema.Unknown), + sessionGenerationId: Schema.optional(TrimmedNonEmptyStringSchema), }); export type RuntimeErrorPayload = typeof RuntimeErrorPayload.Type;