diff --git a/apps/mobile/src/components/AppSymbol.tsx b/apps/mobile/src/components/AppSymbol.tsx index e421316e9..56ff46ed3 100644 --- a/apps/mobile/src/components/AppSymbol.tsx +++ b/apps/mobile/src/components/AppSymbol.tsx @@ -14,6 +14,7 @@ import IconArrowUpCircle from "@tabler/icons-react-native/IconArrowUpCircle"; import IconArrowUpRight from "@tabler/icons-react-native/IconArrowUpRight"; import IconArrowUpRightCircle from "@tabler/icons-react-native/IconArrowUpRightCircle"; import IconArrowsMaximize from "@tabler/icons-react-native/IconArrowsMaximize"; +import IconArrowsMinimize from "@tabler/icons-react-native/IconArrowsMinimize"; import IconBellRinging from "@tabler/icons-react-native/IconBellRinging"; import IconBolt from "@tabler/icons-react-native/IconBolt"; import IconBox from "@tabler/icons-react-native/IconBox"; @@ -97,6 +98,7 @@ const ANDROID_ICON_BY_SF_SYMBOL: Partial> = { "arrow.up": IconArrowUp, "arrow.up.circle": IconArrowUpCircle, "arrow.up.left.and.arrow.down.right": IconArrowsMaximize, + "arrow.down.right.and.arrow.up.left": IconArrowsMinimize, "arrow.up.right": IconArrowUpRight, "arrow.up.right.circle": IconArrowUpRightCircle, "arrow.uturn.backward": IconArrowBackUp, diff --git a/apps/mobile/src/features/threads/NewTaskDraftScreen.tsx b/apps/mobile/src/features/threads/NewTaskDraftScreen.tsx index c88a798b8..a11851876 100644 --- a/apps/mobile/src/features/threads/NewTaskDraftScreen.tsx +++ b/apps/mobile/src/features/threads/NewTaskDraftScreen.tsx @@ -330,6 +330,7 @@ export function NewTaskDraftScreen(props: { sessionResources: null, showInteractionModeToggle: flow.showInteractionModeToggle, hasThread: false, + hasCompactableConversation: false, enabled: isComposerFocused && !isComposerInteractionLocked, onChangeDraftMessage: flow.setPrompt, onUpdateInteractionMode: flow.setInteractionMode, diff --git a/apps/mobile/src/features/threads/ThreadComposer.tsx b/apps/mobile/src/features/threads/ThreadComposer.tsx index 88bc9db01..0a9fd2343 100644 --- a/apps/mobile/src/features/threads/ThreadComposer.tsx +++ b/apps/mobile/src/features/threads/ThreadComposer.tsx @@ -213,6 +213,7 @@ export interface ThreadComposerProps { readonly connectionState: RemoteClientConnectionState; readonly environmentLabel: string | null; readonly selectedThread: OrchestrationThreadShell; + readonly hasCompactableConversation: boolean; readonly serverConfig: T3ServerConfig | null; readonly localOutboxCount: number; readonly onManagePendingSends: () => void; @@ -1243,6 +1244,7 @@ export const ThreadComposer = memo(function ThreadComposer(props: ThreadComposer sessionResources: props.sessionResources, showInteractionModeToggle, hasThread: true, + hasCompactableConversation: props.hasCompactableConversation, enabled: !props.sessionInputBlocked, onChangeDraftMessage: props.onChangeDraftMessage, onUpdateInteractionMode: props.onUpdateInteractionMode, diff --git a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx index 71de73d61..69dd1ccec 100644 --- a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx @@ -136,6 +136,7 @@ export interface ThreadDetailScreenProps { readonly sessionCompactionScopeKey: string | null; readonly sessionCompactionPendingAction: SessionCompactionMenuAction | null; readonly activeWorkStartedAt: string | null; + readonly isCompacting: boolean; readonly activePendingApproval: PendingApproval | null; readonly respondingApprovalId: ApprovalRequestId | null; readonly activePendingUserInput: PendingUserInput | null; @@ -468,8 +469,10 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread if (threadSyncLabel !== null) { return { kind: "syncing", label: threadSyncLabel }; } + // Prime reports compaction through its own control snapshot; every other + // provider reports it as a compacting turn. if ( - isSessionCompactionInProgress(props.sessionCompaction) && + (isSessionCompactionInProgress(props.sessionCompaction) || props.isCompacting) && contentPresentationKind === "ready" ) { return { kind: "compacting" }; @@ -481,6 +484,15 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread })(); const showWorkingControl = floatingStatus !== null; const selectedThreadFeed = props.selectedThreadFeed; + const hasCompactableConversation = + selectedThreadFeed.some( + (entry) => + entry.type === "message" && + entry.message.role === "user" && + ((entry.message.attachments?.length ?? 0) > 0 || + entry.message.text.trim().toLowerCase() !== "/compact"), + ) || + (Boolean(props.loadEarlier) && props.selectedThread.latestUserMessageAt !== null); const composerChrome = composerExpanded ? COMPOSER_EXPANDED_CHROME : COMPOSER_COLLAPSED_CHROME; const composerOverlapHeight = composerChrome + composerBottomInset; // While a user-input request is pending, the questionnaire owns the @@ -1052,6 +1064,7 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread connectionState={props.connectionStateLabel} environmentLabel={props.environmentLabel} selectedThread={props.selectedThread} + hasCompactableConversation={hasCompactableConversation && !props.isCompacting} serverConfig={props.serverConfig} localOutboxCount={props.localOutboxCount} onManagePendingSends={props.onManagePendingSends} diff --git a/apps/mobile/src/features/threads/ThreadFeed.tsx b/apps/mobile/src/features/threads/ThreadFeed.tsx index b633bc6b9..080e35435 100644 --- a/apps/mobile/src/features/threads/ThreadFeed.tsx +++ b/apps/mobile/src/features/threads/ThreadFeed.tsx @@ -144,6 +144,7 @@ import { } from "@t3tools/mobile-markdown-text/links"; import { deriveThreadFeedPresentation, + isContextCompactionActivityGroup, type ThreadFeedEntry, type ThreadFeedLatestTurn, } from "../../lib/threadActivity"; @@ -1417,6 +1418,29 @@ function renderFeedEntry( ); } + if (entry.type === "activity-group" && isContextCompactionActivityGroup(entry)) { + const label = entry.activities[0]!.summary; + return ( + + + + + {label} + + + + ); + } + if (entry.type === "message") { const { message } = entry; const isUser = message.role === "user"; @@ -2615,6 +2639,9 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { case "thinking": return WORK_GROUP_TOGGLE_HEIGHT; case "activity-group": + if (isContextCompactionActivityGroup(entry)) { + return undefined; + } // Expanded rows append a variable detail block — fall back to // measurement for those groups. return entry.activities.some((activity) => expandedWorkRows[activity.id]) diff --git a/apps/mobile/src/features/threads/ThreadRouteScreen.tsx b/apps/mobile/src/features/threads/ThreadRouteScreen.tsx index bfc8516ee..56c55801f 100644 --- a/apps/mobile/src/features/threads/ThreadRouteScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadRouteScreen.tsx @@ -982,6 +982,7 @@ function ThreadRouteContent( sessionCompactionScopeKey={composer.sessionCompactionScopeKey} sessionCompactionPendingAction={composer.sessionCompactionPendingAction} activeWorkStartedAt={composer.activeWorkStartedAt} + isCompacting={composer.isCompacting} activePendingApproval={requests.activePendingApproval} respondingApprovalId={requests.respondingApprovalId} activePendingUserInput={requests.activePendingUserInput} diff --git a/apps/mobile/src/features/threads/use-composer-command-menu.test.ts b/apps/mobile/src/features/threads/use-composer-command-menu.test.ts index 20a94b509..04f805661 100644 --- a/apps/mobile/src/features/threads/use-composer-command-menu.test.ts +++ b/apps/mobile/src/features/threads/use-composer-command-menu.test.ts @@ -49,6 +49,7 @@ const baseInput = { providerSlashCommands: [], showInteractionModeToggle: false, hasThread: true, + hasCompactableConversation: true, pathEntries: [], } as const; diff --git a/apps/mobile/src/features/threads/use-composer-command-menu.ts b/apps/mobile/src/features/threads/use-composer-command-menu.ts index b831309e4..ffc38497e 100644 --- a/apps/mobile/src/features/threads/use-composer-command-menu.ts +++ b/apps/mobile/src/features/threads/use-composer-command-menu.ts @@ -84,6 +84,7 @@ export function buildComposerCommandItems({ providerSlashCommands, showInteractionModeToggle, hasThread, + hasCompactableConversation, pathEntries, }: { readonly trigger: ComposerTrigger | null; @@ -91,6 +92,7 @@ export function buildComposerCommandItems({ readonly providerSlashCommands: ReadonlyArray; readonly showInteractionModeToggle: boolean; readonly hasThread: boolean; + readonly hasCompactableConversation: boolean; readonly pathEntries: ReadonlyArray; }): ComposerCommandItem[] { if (!trigger) return []; @@ -146,6 +148,8 @@ export function buildComposerCommandItems({ for (const cmd of expandableCommands) { if (!cmd.name.toLowerCase().includes(q)) continue; if (!hasThread && cmd.localAction === "usage-limits") continue; + // Nothing to summarize before the thread has a compactable exchange. + if (cmd.name === "compact" && !hasCompactableConversation) continue; // Codex `/feedback` uploads an existing thread's session and logs, so it // has nothing to send before the thread exists. if (!hasThread && selectedProviderStatus?.driver === "codex" && cmd.name === "feedback") { @@ -307,6 +311,7 @@ export function useComposerCommandMenu({ sessionResources, showInteractionModeToggle, hasThread, + hasCompactableConversation, enabled = true, onChangeDraftMessage, onUpdateInteractionMode, @@ -319,6 +324,7 @@ export function useComposerCommandMenu({ readonly sessionResources: SessionResourcesSnapshot | null; readonly showInteractionModeToggle: boolean; readonly hasThread: boolean; + readonly hasCompactableConversation: boolean; readonly enabled?: boolean; readonly onChangeDraftMessage: (value: string) => void; readonly onUpdateInteractionMode?: (mode: ProviderInteractionMode) => void; @@ -441,10 +447,12 @@ export function useComposerCommandMenu({ showInteractionModeToggle: showInteractionModeToggle && onUpdateInteractionMode !== undefined, hasThread, + hasCompactableConversation, pathEntries: pathSearch.entries, }), [ hasThread, + hasCompactableConversation, onUpdateInteractionMode, pathSearch.entries, providerSlashCommands, diff --git a/apps/mobile/src/lib/threadActivity.test.ts b/apps/mobile/src/lib/threadActivity.test.ts index e5ffdb1e6..6d0ad6dfa 100644 --- a/apps/mobile/src/lib/threadActivity.test.ts +++ b/apps/mobile/src/lib/threadActivity.test.ts @@ -470,6 +470,33 @@ describe("buildThreadFeed", () => { expect(prepended.at(-1)).toBe(page.at(-1)); }); + it("keeps context compaction as a standalone timeline row", () => { + const thread = makeThread({ + id: ThreadId.make("thread-context-compaction"), + projectId: ProjectId.make("project-1"), + title: "Context compaction", + activities: [ + makeActivity({ + id: EventId.make("context-compaction"), + kind: "context-compaction", + tone: "info", + summary: "Compacted context 899K → 19K tokens", + createdAt: "2026-09-01T00:00:00.000Z", + turnId: TurnId.make("turn-context-compaction"), + }), + ], + }); + + const presented = deriveThreadFeedPresentation(buildThreadFeed(thread), null, new Set()); + expect(presented).toMatchObject([ + { + type: "activity-group", + id: "context-compaction", + activities: [{ summary: "Compacted context 899K → 19K tokens" }], + }, + ]); + }); + it("keeps long Claude commands expandable without repeating them in full detail", () => { const command = `printf 'first line\nsecond line'\n&& printf done`; const thread = makeThread({ diff --git a/apps/mobile/src/lib/threadActivity.ts b/apps/mobile/src/lib/threadActivity.ts index c79ee62d6..f7c7e029d 100644 --- a/apps/mobile/src/lib/threadActivity.ts +++ b/apps/mobile/src/lib/threadActivity.ts @@ -222,6 +222,15 @@ export type ThreadFeedLatestTurn = Pick< "turnId" | "state" | "startedAt" | "completedAt" >; +export function isContextCompactionActivityGroup( + entry: Extract, +): boolean { + return ( + entry.activities.length === 1 && + entry.activities[0]?.workEntry.sourceActivityKind === "context-compaction" + ); +} + type ThreadFeedActivityGroup = Extract; // These keys are immutable inputs. Weak caches release old histories with their source data. @@ -1528,7 +1537,11 @@ function groupAdjacentActivities(entries: ReadonlyArray): Th continue; } - const isStandalone = entry.activity.terminalResponseNotice === true; + // Terminal response notices and context compaction each own their card, so + // neither joins the surrounding tool group. + const isStandalone = + entry.activity.terminalResponseNotice === true || + entry.activity.workEntry.sourceActivityKind === "context-compaction"; if (isStandalone || firstActivityEntry?.turnId !== entry.turnId) { flushGroup(); } @@ -1652,6 +1665,16 @@ function deriveThreadFeedTurnFolds( if (hiddenEntryIds.size === 0) { continue; } + // A lone compaction row stays visible on its own; it only folds away as + // part of a turn that already folds other work. + const hidesNonCompactionWork = entries.some( + (entry) => + hiddenEntryIds.has(entry.id) && + !(entry.type === "activity-group" && isContextCompactionActivityGroup(entry)), + ); + if (!hidesNonCompactionWork) { + continue; + } const firstEntry = entries[0]; const firstHiddenEntry = entries.find((entry) => hiddenEntryIds.has(entry.id)); @@ -1816,6 +1839,10 @@ function appendPresentedFeedEntry( result.push(entry); return; } + if (isContextCompactionActivityGroup(entry)) { + result.push(entry); + return; + } let cached = presentedActivityGroupsCache.get(entry); if ( diff --git a/apps/mobile/src/state/use-thread-composer-state.ts b/apps/mobile/src/state/use-thread-composer-state.ts index 5ce1862ae..e95d77951 100644 --- a/apps/mobile/src/state/use-thread-composer-state.ts +++ b/apps/mobile/src/state/use-thread-composer-state.ts @@ -102,7 +102,7 @@ import { } from "./thread-outbox"; import { removeThreadOutboxMessageIfCurrent } from "./thread-outbox-removal"; import { recoverPendingSendToComposer } from "./thread-outbox-recovery"; -import { useThreadOutboxMessages } from "./use-thread-outbox"; +import { dispatchingQueuedMessageIdAtom, useThreadOutboxMessages } from "./use-thread-outbox"; import { useAtomCommand } from "./use-atom-command"; import { composerAttachmentUploadBlockReason, @@ -171,6 +171,7 @@ export function useThreadComposerState() { ) ?? null); const composerDrafts = useAtomValue(composerDraftsAtom); const queuedMessagesByThreadKey = useThreadOutboxMessages(); + const dispatchingQueuedMessageId = useAtomValue(dispatchingQueuedMessageIdAtom); const [feedbackSubmissionsByThreadKey, setFeedbackSubmissionsByThreadKey] = useState< Record> >({}); @@ -549,6 +550,48 @@ export function useThreadComposerState() { }; }, [selectedThreadDetail, selectedThreadShell]); + const isCompacting = useMemo(() => { + const queuedMessage = selectedThreadQueuedMessages.findLast( + (message) => + message.messageId === dispatchingQueuedMessageId && + message.text.trim().toLowerCase() === "/compact" && + message.attachments.length === 0, + ); + const latestCompactMessage = selectedThreadDetail?.messages.findLast( + (message) => + message.role === "user" && + message.text.trim().toLowerCase() === "/compact" && + !message.attachments?.length, + ); + const compactRequestIsActive = + latestCompactMessage !== undefined && + (latestCompactMessage.createdAt > + (selectedThread?.latestTurn?.requestedAt ?? latestCompactMessage.createdAt) || + (selectedThread?.latestTurn?.state === "running" && + latestCompactMessage.createdAt === selectedThread.latestTurn.requestedAt)); + const compactionSettled = selectedThreadDetail?.activities.some((activity) => { + if (!["context-compaction", "provider.turn.start.failed"].includes(activity.kind)) + return false; + const payload = + typeof activity.payload === "object" && activity.payload !== null + ? (activity.payload as { readonly requestId?: unknown }) + : null; + return payload?.requestId === latestCompactMessage?.id; + }); + return ( + queuedMessage !== undefined || + ((selectedThread?.session?.status === "starting" || + selectedThread?.session?.status === "running") && + compactRequestIsActive && + !compactionSettled) + ); + }, [ + dispatchingQueuedMessageId, + selectedThread, + selectedThreadDetail, + selectedThreadQueuedMessages, + ]); + const activeWorkStartedAt = useMemo(() => { const selectedThread = selectedThreadDetail ?? selectedThreadShell; if (!selectedThread) { @@ -1551,6 +1594,7 @@ export function useThreadComposerState() { selectedThreadQueueCount, selectedThreadQueueHold: selectedThreadQueuedMessages[0]?.deliveryHold ?? null, activeWorkStartedAt, + isCompacting, draftMessage, draftAttachments, modelSelection, diff --git a/apps/mobile/src/state/use-thread-outbox-drain.ts b/apps/mobile/src/state/use-thread-outbox-drain.ts index 08c277134..ea2756e46 100644 --- a/apps/mobile/src/state/use-thread-outbox-drain.ts +++ b/apps/mobile/src/state/use-thread-outbox-drain.ts @@ -12,7 +12,7 @@ import { } from "@t3tools/contracts"; import { buildTemporaryWorktreeBranchName } from "@t3tools/shared/git"; import * as Cause from "effect/Cause"; -import { AsyncResult, Atom } from "effect/unstable/reactivity"; +import { AsyncResult } from "effect/unstable/reactivity"; import { useCallback, useEffect, useRef, useState } from "react"; import { scopedProjectKey, scopedThreadKey } from "../lib/scopedEntities"; @@ -64,6 +64,7 @@ import { } from "./use-composer-drafts"; import { useAtomCommand } from "./use-atom-command"; import { + dispatchingQueuedMessageIdAtom, editingQueuedMessageIdsAtom, useThreadOutboxMessages, useThreadOutboxShellStatuses, @@ -73,11 +74,6 @@ import { useRemoteConnectionStatus, } from "./use-remote-environment-registry"; -export const dispatchingQueuedMessageIdAtom = Atom.make(null).pipe( - Atom.keepAlive, - Atom.withLabel("mobile:thread-outbox:dispatching-message-id"), -); - function beginDispatchingQueuedMessage(queuedMessageId: MessageId): void { appAtomRegistry.set(dispatchingQueuedMessageIdAtom, queuedMessageId); } diff --git a/apps/mobile/src/state/use-thread-outbox.ts b/apps/mobile/src/state/use-thread-outbox.ts index 542c2d401..6ed00e2a0 100644 --- a/apps/mobile/src/state/use-thread-outbox.ts +++ b/apps/mobile/src/state/use-thread-outbox.ts @@ -31,6 +31,11 @@ export const editingQueuedMessageIdsAtom = Atom.make(null).pipe( + Atom.keepAlive, + Atom.withLabel("mobile:thread-outbox:dispatching-message-id"), +); + export function holdEditingQueuedMessage(messageId: MessageId): void { const current = appAtomRegistry.get(editingQueuedMessageIdsAtom); if (current[messageId]) { diff --git a/apps/server/integration/orphanedProviderSessionStartup.integration.test.ts b/apps/server/integration/orphanedProviderSessionStartup.integration.test.ts index 0f6d6fe80..a7ce6dae8 100644 --- a/apps/server/integration/orphanedProviderSessionStartup.integration.test.ts +++ b/apps/server/integration/orphanedProviderSessionStartup.integration.test.ts @@ -113,6 +113,7 @@ const startupDependencies = Layer.mergeAll( Layer.succeed(ProviderService.ProviderService, { startSession: () => Effect.die("unused"), sendTurn: () => Effect.die("unused"), + compactThread: () => Effect.die("unused"), interruptTurn: () => Effect.die("unused"), respondToRequest: () => Effect.die("unused"), respondToUserInput: () => Effect.die("unused"), diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts index 47d8ac8db..c85af28f5 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts @@ -115,6 +115,7 @@ function createProviderServiceHarness( const service: ProviderServiceShape = { startSession: () => unsupported(), sendTurn: () => unsupported(), + compactThread: () => unsupported(), interruptTurn: () => unsupported(), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index 96bda5472..820018e3c 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -3775,6 +3775,69 @@ it.layer(makeProjectionPipelinePrefixedTestLayer("t3-pending-turn-terminal-test- assert.deepEqual(pendingRows, []); }), ); + + it.effect("only clears the compact request that produced the compaction activity", () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const threadId = ThreadId.make("thread-compaction-correlation"); + + for (const [index, messageId] of ["compact-request", "new-message"].entries()) { + const createdAt = `2026-02-26T15:00:0${index}.000Z`; + yield* eventStore.append({ + type: "thread.turn-start-requested", + eventId: EventId.make(`evt-compaction-pending-${index}`), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: createdAt, + commandId: CommandId.make(`cmd-compaction-pending-${index}`), + causationEventId: null, + correlationId: CorrelationId.make(`cmd-compaction-pending-${index}`), + metadata: {}, + payload: { + threadId, + messageId: MessageId.make(messageId), + runtimeMode: "full-access", + createdAt, + }, + }); + } + yield* eventStore.append({ + type: "thread.activity-appended", + eventId: EventId.make("evt-compaction-stale"), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: "2026-02-26T15:00:02.000Z", + commandId: CommandId.make("cmd-compaction-stale"), + causationEventId: null, + correlationId: CorrelationId.make("cmd-compaction-stale"), + metadata: {}, + payload: { + threadId, + activity: { + id: EventId.make("activity-compaction-stale"), + tone: "info", + kind: "context-compaction", + summary: "Context compacted", + payload: { requestId: "compact-request" }, + turnId: null, + createdAt: "2026-02-26T15:00:02.000Z", + }, + }, + }); + yield* projectionPipeline.bootstrap; + + const pendingRows = yield* sql<{ readonly messageId: string }>` + SELECT pending_message_id AS "messageId" + FROM projection_turns + WHERE thread_id = ${threadId} + AND turn_id IS NULL + AND state = 'pending' + `; + assert.deepEqual(pendingRows, [{ messageId: "new-message" }]); + }), + ); }, ); diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index fc4496bd8..2f42126f8 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -1377,6 +1377,22 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti return; case "thread.turn-start-requested": { + const pendingTurnStart = yield* projectionTurnRepository.getPendingTurnStartByThreadId({ + threadId: event.payload.threadId, + }); + if (Option.isSome(pendingTurnStart)) { + const pendingMessage = yield* projectionThreadMessageRepository.getByMessageId({ + messageId: pendingTurnStart.value.messageId, + }); + if ( + Option.isSome(pendingMessage) && + pendingMessage.value.role === "user" && + (pendingMessage.value.attachments?.length ?? 0) === 0 && + pendingMessage.value.text.trim().toLowerCase() === "/compact" + ) { + return; + } + } yield* projectionTurnRepository.replacePendingTurnStart({ threadId: event.payload.threadId, messageId: event.payload.messageId, @@ -1387,10 +1403,42 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti return; } + case "thread.activity-appended": { + if (event.payload.activity.kind === "context-compaction") { + const pendingTurnStart = yield* projectionTurnRepository.getPendingTurnStartByThreadId( + event.payload, + ); + if ( + Option.isNone(pendingTurnStart) || + String(pendingTurnStart.value.messageId) !== + extractActivityRequestId(event.payload.activity.payload) + ) { + return; + } + yield* projectionTurnRepository.deletePendingTurnStartByThreadId(event.payload); + return; + } + if (event.payload.activity.kind !== "provider.turn.start.failed") return; + const pendingTurnStart = yield* projectionTurnRepository.getPendingTurnStartByThreadId( + event.payload, + ); + if ( + Option.isNone(pendingTurnStart) || + String(pendingTurnStart.value.messageId) !== + extractActivityRequestId(event.payload.activity.payload) + ) { + return; + } + yield* projectionTurnRepository.deletePendingTurnStartByThreadId(event.payload); + return; + } + case "thread.session-set": { const turnId = event.payload.session.activeTurnId; if (turnId === null || event.payload.session.status !== "running") { if ( + (event.payload.session.status === "ready" && + event.commandId?.startsWith("server:provider-session-set:") === true) || event.payload.session.status === "error" || event.payload.session.status === "stopped" || event.payload.session.status === "interrupted" diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 06d68ca80..8e36a053e 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -196,6 +196,8 @@ describe("ProviderCommandReactor", () => { readonly titleRegenerationCompletionDispatchFailures?: number; readonly titleRegenerationBeforeStart?: "one" | "two"; readonly serverActivation?: Effect.Effect; + readonly beforeReadySessionDispatch?: () => Effect.Effect; + readonly compactThreadEffect?: () => Effect.Effect; readonly interruptTurnEffect?: () => Effect.Effect; readonly stopSessionEffect?: () => Effect.Effect; readonly startSessionEffect?: ( @@ -341,6 +343,7 @@ describe("ProviderCommandReactor", () => { Effect.andThen(result), ); }); + const compactThread = vi.fn((_: ThreadId) => input?.compactThreadEffect?.() ?? Effect.void); const interruptTurn = vi.fn((_: unknown) => input?.interruptTurnEffect?.() ?? Effect.void); const respondToRequest = vi.fn(() => Effect.void); const respondToUserInput = vi.fn(() => Effect.void); @@ -443,6 +446,7 @@ describe("ProviderCommandReactor", () => { const service: ProviderServiceShape = { startSession: startSession as ProviderServiceShape["startSession"], sendTurn: sendTurn as ProviderServiceShape["sendTurn"], + compactThread, interruptTurn: interruptTurn as ProviderServiceShape["interruptTurn"], respondToRequest: respondToRequest as ProviderServiceShape["respondToRequest"], respondToUserInput: respondToUserInput as ProviderServiceShape["respondToUserInput"], @@ -570,7 +574,11 @@ describe("ProviderCommandReactor", () => { return Effect.die(new Error("Injected title regeneration completion failure")); } } - return engine.dispatch(command); + return ( + command.type === "thread.session.set" && command.session.status === "ready" + ? (input?.beforeReadySessionDispatch?.() ?? Effect.void) + : Effect.void + ).pipe(Effect.andThen(engine.dispatch(command))); }, get streamDomainEvents() { return engine.streamDomainEvents; @@ -800,8 +808,20 @@ describe("ProviderCommandReactor", () => { engine, snapshotQuery, readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), + readPendingTurnStarts: () => + runtime!.runPromise( + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + return yield* sql<{ readonly threadId: string }>` + SELECT thread_id AS "threadId" + FROM projection_turns + WHERE turn_id IS NULL AND state = 'pending' + `; + }), + ), startSession, sendTurn, + compactThread, interruptTurn, respondToRequest, respondToUserInput, @@ -1066,6 +1086,286 @@ describe("ProviderCommandReactor", () => { expect(harness.sendTurn).not.toHaveBeenCalled(); }); + effectIt.effect("rejects /compact without conversation context", () => + Effect.gen(function* () { + const harness = yield* Effect.promise(() => createHarness()); + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-empty-compact"), + threadId: ThreadId.make("thread-1"), + message: { + messageId: asMessageId("user-message-empty-compact"), + role: "user", + text: "/compact", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:00.000Z", + }); + yield* Effect.promise(() => harness.drain()); + expect(harness.compactThread).not.toHaveBeenCalled(); + }), + ); + + effectIt.effect("keeps turns blocked until compaction restores the session", () => + Effect.gen(function* () { + const readyDispatchStarted = yield* Deferred.make(); + const releaseReadyDispatch = yield* Deferred.make(); + let blockReadyDispatch = false; + const harness = yield* Effect.promise(() => + createHarness({ + beforeReadySessionDispatch: () => + blockReadyDispatch + ? Deferred.succeed(readyDispatchStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseReadyDispatch)), + ) + : Effect.void, + }), + ); + const threadId = ThreadId.make("thread-1"); + const now = "2026-01-01T00:00:00.000Z"; + const dispatchTurn = (id: string, text: string, createdAt: string) => + harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make(`cmd-${id}`), + threadId, + message: { + messageId: asMessageId(`user-message-${id}`), + role: "user", + text, + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt, + }); + + yield* dispatchTurn("before-blocked-compact", "hello", now); + yield* Effect.promise(() => waitFor(() => harness.sendTurn.mock.calls.length === 1)); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-ready-before-blocked-compact"), + threadId, + session: { + threadId, + status: "ready", + providerName: "codex", + providerInstanceId: ProviderInstanceId.make("codex"), + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + + blockReadyDispatch = true; + yield* dispatchTurn("blocked-compact", "/compact", "2026-01-01T00:00:01.000Z"); + yield* Deferred.await(readyDispatchStarted); + + // Upstream blocks this in the reactor and records a failure activity. + // Pylon's decider owns turn exclusivity: the compaction still holds the + // thread's pending admission, so the command is refused outright and + // never reaches the reactor's own `compactingThreadIds` guard. + const blockedTurn = yield* dispatchTurn( + "during-compact-recovery", + "too soon", + "2026-01-01T00:00:02.000Z", + ).pipe(Effect.result); + expect(blockedTurn._tag).toBe("Failure"); + expect(harness.sendTurn).toHaveBeenCalledTimes(1); + expect(yield* Effect.promise(() => harness.readPendingTurnStarts())).toEqual([ + { threadId: "thread-1" }, + ]); + + yield* Deferred.succeed(releaseReadyDispatch, undefined); + yield* Effect.promise(() => + waitFor(async () => { + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + return thread?.session?.status === "ready"; + }), + ); + }), + ); + + effectIt.effect("does not overwrite concurrent session state after compaction failure", () => + Effect.gen(function* () { + const releaseCompaction = yield* Deferred.make(); + const releaseRunningCompaction = yield* Deferred.make(); + const releaseFailedStop = yield* Deferred.make(); + let compactionCount = 0; + const harness = yield* Effect.promise(() => + createHarness({ + compactThreadEffect: () => + Deferred.await( + compactionCount++ === 0 ? releaseCompaction : releaseRunningCompaction, + ).pipe(Effect.andThen(Effect.die("Compaction stopped"))), + stopSessionEffect: () => + Deferred.await(releaseFailedStop).pipe( + Effect.andThen( + Effect.fail( + new ProviderAdapterRequestError({ + provider: "codex", + method: "session.stop", + detail: "provider stop failed", + }), + ), + ), + ), + }), + ); + const threadId = ThreadId.make("thread-1"); + const now = "2026-01-01T00:00:00.000Z"; + const dispatchCompact = (suffix: string, createdAt: string) => + harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make(`cmd-compact-${suffix}`), + threadId, + message: { + messageId: asMessageId(`user-message-compact-${suffix}`), + role: "user", + text: "/compact", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt, + }); + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-message-before-compact"), + threadId, + message: { + messageId: asMessageId("user-message-before-compact"), + role: "user", + text: "hello", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + }); + yield* Effect.promise(() => waitFor(() => harness.sendTurn.mock.calls.length === 1)); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-ready-before-compact"), + threadId, + session: { + threadId, + status: "ready", + providerName: "codex", + providerInstanceId: ProviderInstanceId.make("codex"), + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + yield* dispatchCompact("before-stop", now); + yield* Effect.promise(() => waitFor(() => harness.compactThread.mock.calls.length === 1)); + const compactingThread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + expect(compactingThread?.session?.status).toBe("starting"); + yield* harness.engine.dispatch({ + type: "thread.session.stop", + commandId: CommandId.make("cmd-stop-during-compact"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + }); + yield* Effect.promise(() => waitFor(() => harness.stopSession.mock.calls.length === 1)); + yield* Deferred.succeed(releaseCompaction, undefined); + yield* Effect.promise(() => + waitFor(async () => { + const compactingThread = (await harness.readModel()).threads.find( + (entry) => entry.id === threadId, + ); + return ( + compactingThread?.activities.some( + (activity) => activity.kind === "provider.turn.start.failed", + ) === true + ); + }), + ); + // Pylon's decider persists the `stopped` transition in the same + // transaction as the stop intent, so the session reads stopped here where + // upstream still reads starting. What matters is the same either way: the + // in-flight compaction does not overwrite it, and the failed stop is + // recorded and recovered below. + const stoppingThread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + expect(stoppingThread?.session?.status).toBe("stopped"); + yield* Deferred.succeed(releaseFailedStop, undefined); + yield* Effect.promise(() => harness.drain()); + + const recoveredThread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + // Upstream restores the session to ready here because its stop had not + // committed yet. Pylon's stop is authoritative once the decider records + // it, so a failed provider stop leaves the thread stopped and the next + // turn starts a fresh session. The stop failure is still recorded. + expect(recoveredThread?.session?.status).toBe("stopped"); + expect( + recoveredThread?.activities.find( + (activity) => activity.kind === "provider.session.stop.failed", + ), + ).toMatchObject({ + summary: "Provider session stop failed", + payload: { detail: "provider stop failed" }, + }); + + yield* dispatchCompact("before-running", "2026-01-01T00:00:02.000Z"); + yield* Effect.promise(() => waitFor(() => harness.compactThread.mock.calls.length === 2)); + yield* harness.engine.dispatch({ + type: "thread.session.stop", + commandId: CommandId.make("cmd-failed-stop-before-compaction-settles"), + threadId, + createdAt: "2026-01-01T00:00:02.500Z", + }); + yield* Effect.promise(() => + waitFor(async () => { + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + return ( + thread?.activities.filter( + (activity) => activity.kind === "provider.session.stop.failed", + ).length === 2 + ); + }), + ); + const restartedThread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + // Same eager reverse transition as above. What phase two actually checks + // is below: a compaction that fails after the session moved to running + // must not overwrite it back to ready. + expect(restartedThread?.session?.status).toBe("stopped"); + const restartedSession = restartedThread?.session; + if (!restartedSession) return yield* Effect.die("Compaction session missing"); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-running-during-compact"), + threadId, + session: { + ...restartedSession, + status: "running", + activeTurnId: asTurnId("compaction-turn"), + updatedAt: "2026-01-01T00:00:03.000Z", + }, + createdAt: "2026-01-01T00:00:03.000Z", + }); + yield* Deferred.succeed(releaseRunningCompaction, undefined); + yield* Effect.promise(() => harness.drain()); + const runningThread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + expect(runningThread?.session?.status).toBe("running"); + }), + ); effectIt.effect("projects starting before a slow provider session finishes", () => Effect.gen(function* () { const releaseStart = yield* Deferred.make(); diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 5abee4b6c..4852281cb 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -39,6 +39,7 @@ import { resolveThreadWorkspaceCwd } from "../../checkpointing/Utils.ts"; import { increment, orchestrationEventsProcessedTotal } from "../../observability/Metrics.ts"; import { ProviderAdapterRequestError, + ProviderAdapterValidationError, ProviderWorkspaceMissingError, } from "../../provider/Errors.ts"; import type { ProviderServiceError } from "../../provider/Errors.ts"; @@ -65,6 +66,7 @@ import { VcsStatusBroadcaster } from "../../vcs/VcsStatusBroadcaster.ts"; import { GitWorkflowService } from "../../git/GitWorkflowService.ts"; const isProviderAdapterRequestError = Schema.is(ProviderAdapterRequestError); const isProviderWorkspaceMissingError = Schema.is(ProviderWorkspaceMissingError); +const isProviderAdapterValidationError = Schema.is(ProviderAdapterValidationError); const isProviderDriverKind = Schema.is(ProviderDriverKind); type TurnAdmissionIntent = NonNullable< @@ -102,6 +104,16 @@ export function providerFollowUpInputFromMessage(message: { }; } +/** `/compact` is a client-side command, not text the provider should ever see. */ +const isCompactCommandMessage = (message: { + readonly role: string; + readonly text: string; + readonly attachments?: ReadonlyArray | undefined; +}): boolean => + message.role === "user" && + (message.attachments?.length ?? 0) === 0 && + message.text.trim().toLowerCase() === "/compact"; + function mapProviderSessionStatusToOrchestrationStatus( status: "connecting" | "ready" | "running" | "error" | "closed", ): OrchestrationSession["status"] { @@ -369,6 +381,8 @@ const make = Effect.gen(function* () { ); const threadModelSelections = new Map(); + const compactingThreadIds = new Set(); + const stoppingThreadIds = new Set(); const admissionFibers = new Map>(); const admissionFiberThreads = new Map(); type AdmissionStopToken = object; @@ -474,11 +488,11 @@ const make = Effect.gen(function* () { const formatFailureDetail = (cause: Cause.Cause): string => { const failReason = cause.reasons.find(Cause.isFailReason); - const providerError = isProviderAdapterRequestError(failReason?.error) - ? failReason.error - : undefined; - if (providerError) { - return providerError.detail; + if (isProviderAdapterRequestError(failReason?.error)) { + return failReason.error.detail; + } + if (isProviderAdapterValidationError(failReason?.error)) { + return failReason.error.issue; } if (isProviderWorkspaceMissingError(failReason?.error)) { return failReason.error.message; @@ -628,6 +642,101 @@ const make = Effect.gen(function* () { .pipe(Effect.map(Option.getOrUndefined)); }); + const setThreadSession = (input: { + readonly threadId: ThreadId; + readonly session: OrchestrationSession; + readonly createdAt: string; + }) => + serverCommandId("provider-session-set").pipe( + Effect.flatMap((commandId) => + orchestrationEngine.dispatch({ + type: "thread.session.set", + commandId, + threadId: input.threadId, + session: input.session, + createdAt: input.createdAt, + }), + ), + Effect.asVoid, + ); + + const restoreCompaction = Effect.fnUntraced(function* (threadId: ThreadId, fromRunning = false) { + if (stoppingThreadIds.has(threadId)) { + compactingThreadIds.delete(threadId); + return; + } + const thread = yield* resolveThreadShell(threadId); + if (!thread?.session) return; + if ( + thread.session.status !== "starting" && + thread.session.status !== "ready" && + (!fromRunning || thread.session.status !== "running") + ) + return; + const completedAt = DateTime.formatIso(yield* DateTime.now); + if (stoppingThreadIds.has(threadId)) { + compactingThreadIds.delete(threadId); + return; + } + yield* setThreadSession({ + threadId, + session: { + ...thread.session, + status: "ready", + activeTurnId: null, + // Pylon's decider reserves a turn admission for every `thread.turn.start`, + // including the `/compact` one, but compaction never becomes a provider + // turn that could accept it. Retire it here or the spread above carries + // it forward and the decider rejects every later turn on this thread. + pendingTurnRequestId: undefined, + pendingTurnMessageId: undefined, + pendingTurnRequestedAt: undefined, + pendingTurnDeadlineAt: undefined, + pendingTurnSessionId: undefined, + activeTurnRequestId: undefined, + lastError: null, + updatedAt: completedAt, + }, + createdAt: completedAt, + }); + }); + + /** + * Marks a stop in flight so a compaction finishing underneath it cannot flip + * the session back to ready. Pylon's decider records `stopped` in the same + * transaction as the stop intent, so there is nothing to restore afterwards + * even when the provider stop itself fails. + */ + const withSessionStopTracking = (threadId: ThreadId, stop: Effect.Effect) => + Effect.suspend(() => { + stoppingThreadIds.add(threadId); + return stop.pipe( + Effect.onError((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.void + : DateTime.now.pipe( + Effect.flatMap((failedAt) => + appendProviderFailureActivity({ + threadId, + kind: "provider.session.stop.failed", + summary: "Provider session stop failed", + detail: formatFailureDetail(cause), + turnId: null, + createdAt: DateTime.formatIso(failedAt), + }), + ), + Effect.catchCause((appendCause) => + Effect.logWarning("failed to record a provider session stop failure", { + threadId, + cause: Cause.pretty(appendCause), + }), + ), + ), + ), + Effect.ensuring(Effect.sync(() => void stoppingThreadIds.delete(threadId))), + ); + }); + const rejectStartedThreadModelChangeIfRequired = Effect.fnUntraced(function* (input: { readonly threadId: ThreadId; readonly currentModelSelection: ModelSelection; @@ -1515,9 +1624,20 @@ const make = Effect.gen(function* () { detail: `User message '${event.payload.messageId}' was not found for turn start request.`, turnId: null, createdAt: event.payload.createdAt, + requestId: event.payload.messageId, }); return; } + const appendTurnStartFailure = (summary: string, detail: string) => + appendProviderFailureActivity({ + threadId: event.payload.threadId, + kind: "provider.turn.start.failed", + summary, + detail, + turnId: null, + createdAt: event.payload.createdAt, + requestId: event.payload.messageId, + }); const { message, hasOtherUserMessages } = turnStart.value; @@ -1563,6 +1683,100 @@ const make = Effect.gen(function* () { return; } + // `/compact` never reaches the provider as a turn. It runs through + // ProviderService.compactThread, which either calls the adapter's native + // compaction or sends the adapter's own slash command. + if (isCompactCommandMessage(message)) { + if (!hasOtherUserMessages) { + return yield* appendTurnStartFailure( + "Context compaction failed", + "Context compaction requires an existing conversation.", + ); + } + const latestSession = (yield* resolveThreadShell(event.payload.threadId))?.session; + // Pylon admits the turn before the reactor observes its intent event, so + // this thread already reads as "starting" for the compaction's own + // request. Upstream's status check would reject every compaction here; + // reject only when a different turn owns the session. + const otherTurnOwnsSession = + latestSession?.status === "running" || + latestSession?.activeTurnId != null || + (latestSession?.pendingTurnRequestId != null && + latestSession.pendingTurnRequestId !== requestId); + if (compactingThreadIds.has(event.payload.threadId) || otherTurnOwnsSession) { + yield* appendTurnStartFailure( + "Context compaction failed", + "Context compaction is unavailable while a provider turn is running.", + ); + return; + } + const handleCompactionFailure = (cause: Cause.Cause) => { + if (Cause.hasInterruptsOnly(cause)) return Effect.void; + const detail = formatFailureDetail(cause); + return appendTurnStartFailure("Context compaction failed", detail).pipe( + Effect.ensuring( + // A no-op unless the session actually reached a restorable state, + // so this covers both a failed ensure and a failed compaction. + restoreCompaction(event.payload.threadId).pipe( + Effect.catchCause((restoreCause) => + Effect.logWarning("failed to restore provider session after compaction failure", { + threadId: event.payload.threadId, + cause: Cause.pretty(restoreCause), + }), + ), + ), + ), + Effect.asVoid, + ); + }; + compactingThreadIds.add(event.payload.threadId); + yield* Effect.gen(function* () { + // Deliberately no pending turn admission: `restoreCompaction` spreads + // the existing session forward, so an admission recorded here would + // survive it and the decider would reject every later turn on the + // thread. `compactingThreadIds` is what serialises compaction instead. + yield* ensureSessionForThread(event.payload.threadId, event.payload.createdAt, { + ...(event.payload.modelSelection !== undefined + ? { modelSelection: event.payload.modelSelection } + : {}), + runtimeMode: event.payload.runtimeMode, + interactionMode: event.payload.interactionMode, + pendingTurnStart: true, + }); + if (event.payload.modelSelection !== undefined) { + threadModelSelections.set(event.payload.threadId, event.payload.modelSelection); + } + yield* providerService.compactThread( + event.payload.threadId, + event.payload.modelSelection, + event.payload.messageId, + ); + }).pipe( + Effect.andThen(restoreCompaction(event.payload.threadId, true)), + Effect.catchCause((cause) => + handleCompactionFailure(cause).pipe( + Effect.catchCause((recoveryCause) => + Effect.logWarning("provider command reactor failed to recover compaction failure", { + eventType: event.type, + threadId: event.payload.threadId, + cause: Cause.pretty(recoveryCause), + originalCause: Cause.pretty(cause), + }), + ), + ), + ), + Effect.ensuring(Effect.sync(() => void compactingThreadIds.delete(event.payload.threadId))), + Effect.forkScoped, + ); + return; + } + if (compactingThreadIds.has(event.payload.threadId)) { + return yield* appendTurnStartFailure( + "Provider turn start failed", + "Wait for context compaction to finish before sending another message.", + ); + } + // Subscribe before provider admission starts. Only the persisted exact // admission CAS disarms the watchdog; observing a raw provider start is not // sufficient because ingestion or projection can still reject it. @@ -2056,13 +2270,16 @@ const make = Effect.gen(function* () { // runtime. A still-projected stopped target owns the original reservation. const invalidateStartReservation = currentSession?.status === "stopped" && pendingStopTargetIsCurrent(currentSession, target); - yield* providerService.stopSession({ - threadId: target.threadId, - expectedProviderInstanceId: target.providerInstanceId, - expectedSessionIncarnationId: target.sessionIncarnationId, - expectedAdmissionRequestId: target.turnRequestId, - invalidateStartReservation, - }); + yield* withSessionStopTracking( + target.threadId, + providerService.stopSession({ + threadId: target.threadId, + expectedProviderInstanceId: target.providerInstanceId, + expectedSessionIncarnationId: target.sessionIncarnationId, + expectedAdmissionRequestId: target.turnRequestId, + invalidateStartReservation, + }), + ); yield* clearPendingSessionStop(target); releaseAdmissionLaneIfIdle(target.threadId); }); @@ -2087,18 +2304,21 @@ const make = Effect.gen(function* () { if (event.commandId === null) { // Legacy stop events have no durable pending-stop marker. Keep their // prior exact-stop behavior, but never manufacture a cleanup target. - yield* providerService.stopSession({ - threadId: event.payload.threadId, - ...(event.payload.targetProviderInstanceId !== undefined || - event.payload.targetSessionIncarnationId !== undefined || - event.payload.targetTurnRequestId !== undefined - ? { - expectedProviderInstanceId: event.payload.targetProviderInstanceId ?? null, - expectedSessionIncarnationId: event.payload.targetSessionIncarnationId ?? null, - expectedAdmissionRequestId: targetRequestId, - } - : {}), - }); + yield* withSessionStopTracking( + event.payload.threadId, + providerService.stopSession({ + threadId: event.payload.threadId, + ...(event.payload.targetProviderInstanceId !== undefined || + event.payload.targetSessionIncarnationId !== undefined || + event.payload.targetTurnRequestId !== undefined + ? { + expectedProviderInstanceId: event.payload.targetProviderInstanceId ?? null, + expectedSessionIncarnationId: event.payload.targetSessionIncarnationId ?? null, + expectedAdmissionRequestId: targetRequestId, + } + : {}), + }), + ); return; } yield* stopPendingSessionTarget({ diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index fa267022d..7aac79269 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -129,6 +129,7 @@ function createProviderServiceHarness() { const service: ProviderServiceShape = { startSession: () => unsupported(), sendTurn: () => unsupported(), + compactThread: () => unsupported(), interruptTurn: () => unsupported(), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), @@ -5748,10 +5749,57 @@ describe("ProviderRuntimeIngestion", () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; + const compactCommand = { + type: "thread.turn.start", + commandId: CommandId.make("cmd-thread-compact"), + threadId: asThreadId("thread-1"), + message: { + messageId: asMessageId("message-compact"), + role: "user", + text: "/compact", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + } satisfies OrchestrationCommand; + await harness.dispatch(compactCommand); + harness.emit({ + type: "session.state.changed", + eventId: asEventId("evt-session-starting-compact"), + provider: ProviderDriverKind.make("codex"), + providerInstanceId: ProviderInstanceId.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + payload: { state: "starting" }, + }); + await waitForThread(harness.readModel, (entry) => entry.session?.status === "starting"); + + for (const [index, usedTokens] of [899_000, 0].entries()) { + harness.emit({ + type: "thread.token-usage.updated", + eventId: asEventId(`evt-thread-token-usage-${index}`), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + payload: { usage: { usedTokens } }, + }); + } + // Pylon upserts one context-window row per thread, so both updates land on + // the same activity; the compaction label reads its retained previous total. + await waitForThread(harness.readModel, (entry) => + entry.activities.some( + (activity: ProviderRuntimeTestActivity) => + activity.kind === "context-window.updated" && + (activity.payload as { readonly usedTokens?: number } | undefined)?.usedTokens === 0, + ), + ); + harness.emit({ type: "thread.state.changed", eventId: asEventId("evt-thread-compacted"), provider: ProviderDriverKind.make("codex"), + providerInstanceId: ProviderInstanceId.make("codex"), createdAt: now, threadId: asThreadId("thread-1"), turnId: asTurnId("turn-1"), @@ -5763,15 +5811,16 @@ describe("ProviderRuntimeIngestion", () => { const thread = await waitForThread(harness.readModel, (entry) => entry.activities.some( - (activity: ProviderRuntimeTestActivity) => activity.kind === "context-compaction", + (activity: ProviderRuntimeTestActivity) => activity.id === "evt-thread-compacted", ), ); const activity = thread.activities.find( - (candidate: ProviderRuntimeTestActivity) => candidate.kind === "context-compaction", + (candidate: ProviderRuntimeTestActivity) => candidate.id === "evt-thread-compacted", ); - expect(activity?.summary).toBe("Context compacted"); + expect(activity?.summary).toBe("Compacted context 899K → 0 tokens"); expect(activity?.tone).toBe("info"); + expect(activity?.payload).toMatchObject({ requestId: "message-compact" }); }); it("projects Codex task lifecycle chunks into thread activities", async () => { diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 6a07bdd4f..9bd2bbdb2 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -19,6 +19,7 @@ import { type OrchestrationSession, type OrchestrationThreadActivity, type ProviderRuntimeEvent, + RuntimeRequestId, type SessionInteractionRequest, type SessionInteractionResponse, } from "@t3tools/contracts"; @@ -27,12 +28,14 @@ import * as Cause from "effect/Cause"; import * as Crypto from "effect/Crypto"; import * as Context from "effect/Context"; import * as Data from "effect/Data"; +import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Stream from "effect/Stream"; import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker"; +import { formatTokens } from "@t3tools/shared/usageFormat"; import { ProviderRegistry } from "../../provider/Services/ProviderRegistry.ts"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; @@ -139,6 +142,8 @@ interface AssistantSegmentState { activeMessageId: MessageId | null; } +const CONTEXT_WINDOW_HISTORY_BY_THREAD_CACHE_CAPACITY = 2_000; +const CONTEXT_WINDOW_HISTORY_BY_THREAD_TTL = Duration.minutes(120); const TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY = 10_000; const TURN_MESSAGE_IDS_BY_TURN_TTL = Duration.minutes(120); const BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_CACHE_CAPACITY = 20_000; @@ -283,6 +288,29 @@ function buildContextWindowActivityPayload( return event.payload.usage; } +/** + * The two most recent context-window totals seen for a thread. Pylon upserts a + * single `context-window.updated` activity per thread so a streaming turn does + * not append a row per token tick, which means the pre-compaction total is no + * longer in the projection by the time the compacted event lands. Retaining the + * previous value here keeps the "899K -> 19K" label without reintroducing a row + * per update. + */ +interface ContextWindowHistory { + readonly previousUsedTokens: number | undefined; + readonly latestUsedTokens: number; +} + +function compactedTokenCounts( + history: ContextWindowHistory | undefined, +): { readonly beforeTokens: number; readonly afterTokens: number } | undefined { + if (history?.previousUsedTokens === undefined) return undefined; + const beforeTokens = history.previousUsedTokens; + const afterTokens = history.latestUsedTokens; + if (afterTokens >= beforeTokens) return undefined; + return { beforeTokens, afterTokens }; +} + function normalizeRuntimeTurnState( value: string | undefined, ): "completed" | "failed" | "interrupted" | "cancelled" { @@ -1165,14 +1193,26 @@ export function runtimeEventToActivities( return []; } + const beforeTokens = event.payload.beforeTokens; + const afterTokens = event.payload.afterTokens; + const summary = + beforeTokens !== undefined && afterTokens !== undefined + ? `Compacted context ${formatTokens(beforeTokens)} → ${formatTokens(afterTokens)} tokens` + : "Context compacted"; return [ { id: event.eventId, createdAt: event.createdAt, tone: "info", kind: "context-compaction", - summary: "Context compacted", - payload: { state: event.payload.state }, + summary, + payload: { + state: event.payload.state, + ...(beforeTokens !== undefined ? { beforeTokens } : {}), + ...(afterTokens !== undefined ? { afterTokens } : {}), + ...(event.requestId !== undefined ? { requestId: event.requestId } : {}), + ...(event.payload.detail !== undefined ? { detail: event.payload.detail } : {}), + }, turnId: toTurnId(event.turnId) ?? null, ...maybeSequence, }, @@ -1524,6 +1564,15 @@ const make = Effect.gen(function* () { Effect.map((uuid) => CommandId.make(`provider:${event.eventId}:${tag}:${uuid}`)), ); + const contextWindowHistoryByThreadId = yield* Cache.make({ + capacity: CONTEXT_WINDOW_HISTORY_BY_THREAD_CACHE_CAPACITY, + timeToLive: CONTEXT_WINDOW_HISTORY_BY_THREAD_TTL, + lookup: () => + Effect.die( + new Error("context window history should be read through getOption before initialization"), + ), + }); + const turnMessageIdsByTurnKey = yield* Cache.make>({ capacity: TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY, timeToLive: TURN_MESSAGE_IDS_BY_TURN_TTL, @@ -3071,7 +3120,66 @@ const make = Effect.gen(function* () { } yield* requireRuntimeEventCurrent; - const activities = runtimeEventToActivities(event, taskTitle); + if (event.type === "thread.token-usage.updated" && event.payload.usage.usedTokens >= 0) { + const previous = Option.getOrUndefined( + yield* Cache.getOption(contextWindowHistoryByThreadId, thread.id), + ); + yield* Cache.set(contextWindowHistoryByThreadId, thread.id, { + previousUsedTokens: previous?.latestUsedTokens, + latestUsedTokens: event.payload.usage.usedTokens, + }); + } + let activityEvent = event; + if ( + isCompactedThreadState && + event.requestId === undefined && + Option.isSome(pendingTurnStart) && + thread.session?.status === "starting" && + activeTurnId === null && + sameId(thread.session.providerName, event.provider) && + sameId(thread.session.providerInstanceId, event.providerInstanceId) && + DateTime.isGreaterThanOrEqualTo( + DateTime.makeUnsafe(event.createdAt), + DateTime.makeUnsafe(pendingTurnStart.value.requestedAt), + ) + ) { + const pendingMessage = yield* getThreadMessageById( + thread.id, + pendingTurnStart.value.messageId, + ); + if ( + pendingMessage?.role === "user" && + (pendingMessage.attachments?.length ?? 0) === 0 && + pendingMessage.text.trim().toLowerCase() === "/compact" + ) { + activityEvent = { + ...event, + requestId: RuntimeRequestId.make(String(pendingTurnStart.value.messageId)), + }; + } + } + if ( + activityEvent.type === "thread.state.changed" && + activityEvent.payload.state === "compacted" && + (activityEvent.payload.beforeTokens === undefined || + activityEvent.payload.afterTokens === undefined) + ) { + const tokenCounts = compactedTokenCounts( + Option.getOrUndefined(yield* Cache.getOption(contextWindowHistoryByThreadId, thread.id)), + ); + if (tokenCounts) { + activityEvent = { + ...activityEvent, + payload: { + ...activityEvent.payload, + beforeTokens: activityEvent.payload.beforeTokens ?? tokenCounts.beforeTokens, + afterTokens: activityEvent.payload.afterTokens ?? tokenCounts.afterTokens, + }, + }; + } + } + + const activities = runtimeEventToActivities(activityEvent, taskTitle); yield* Effect.forEach(activities, (activity) => providerCommandId(event, "thread-activity-append").pipe( Effect.flatMap((commandId) => diff --git a/apps/server/src/provider/Layers/AntigravityAdapter.ts b/apps/server/src/provider/Layers/AntigravityAdapter.ts index 957ec5fa3..4c2e7f48a 100644 --- a/apps/server/src/provider/Layers/AntigravityAdapter.ts +++ b/apps/server/src/provider/Layers/AntigravityAdapter.ts @@ -1249,9 +1249,7 @@ export const makeAntigravityAdapter = Effect.fn("makeAntigravityAdapter")(functi sessionModelSwitch: "in-session", conversationRollback: BUILT_IN_ADAPTER_CONVERSATION_ROLLBACK_MODES.antigravity, }, - // Antigravity declares `compaction: { type: "slash-command", command: "/compact" }` - // upstream. `ProviderAdapterShape` only grows that field with #10112, which is a - // separate concern; restore this line when that lands. + compaction: { type: "slash-command", command: "/compact" }, startSession, sendTurn, interruptTurn, diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index f16033f4c..bf69865e5 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -3348,7 +3348,7 @@ describe("ClaudeAdapterLive", () => { const harness = makeHarness(); return Effect.gen(function* () { const adapter = yield* ClaudeAdapter; - const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 9).pipe( + const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 11).pipe( Stream.runCollect, Effect.forkChild, ); @@ -3388,6 +3388,13 @@ describe("ClaudeAdapterLive", () => { session_id: "sdk-session-compacted-usage", uuid: "compact-boundary-usage", } as unknown as SDKMessage); + harness.query.emit({ + type: "system", + subtype: "compact_boundary", + compact_metadata: { post_tokens: 40 }, + session_id: "sdk-session-compacted-usage", + uuid: "compact-boundary-post-usage", + } as unknown as SDKMessage); harness.query.emit({ type: "result", subtype: "success", @@ -3411,6 +3418,14 @@ describe("ClaudeAdapterLive", () => { } as unknown as SDKMessage); const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); + const compactionEvents = runtimeEvents.filter( + (event): event is Extract => + event.type === "thread.state.changed" && event.payload.state === "compacted", + ); + assert.equal(compactionEvents[0]?.payload.beforeTokens, 200); + assert.equal(compactionEvents[0]?.payload.afterTokens, 40); + assert.equal(compactionEvents[1]?.payload.beforeTokens, undefined); + assert.equal(compactionEvents[1]?.payload.afterTokens, 40); const finalUsageEvent = runtimeEvents.findLast( (event) => event.type === "thread.token-usage.updated", ); @@ -3418,7 +3433,6 @@ describe("ClaudeAdapterLive", () => { if (finalUsageEvent?.type === "thread.token-usage.updated") { assert.deepEqual(finalUsageEvent.payload.usage, { usedTokens: 40, - lastUsedTokens: 200, totalProcessedTokens: 450, maxTokens: 200000, }); diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index fd9c3dd93..bafeb4834 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -724,12 +724,17 @@ function compactBoundaryTokenUsageSnapshot( } const preTokens = finiteNonNegativeInteger(compactMetadata.pre_tokens); - return makeClaudeTokenUsageSnapshot({ + const snapshot = makeClaudeTokenUsageSnapshot({ activeTokens: postTokens, ...(preTokens !== undefined ? { lastUsedTokens: preTokens } : {}), ...(contextWindow !== undefined ? { contextWindow } : {}), ...(totalProcessedTokens !== undefined ? { totalProcessedTokens } : {}), }); + if (snapshot === undefined || preTokens !== undefined) { + return snapshot; + } + const { lastUsedTokens: _lastUsedTokens, ...snapshotWithoutBeforeTokens } = snapshot; + return snapshotWithoutBeforeTokens; } function normalizeClaudeTaskProgressTokenUsage( @@ -3425,32 +3430,36 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }, }); return; - case "compact_boundary": + case "compact_boundary": { if (context.turnState) { context.turnState.latestAssistantUsage = undefined; context.turnState.compactedSinceLatestAssistantUsage = true; } - yield* emitThreadTokenUsage( - context, - compactBoundaryTokenUsageSnapshot( - message as unknown as Record, - context.lastKnownContextWindow, - context.lastKnownTotalProcessedTokens, - ), - { - rawMethod: "claude/system/compact_boundary", - rawPayload: message, - }, + const compactedUsage = compactBoundaryTokenUsageSnapshot( + message as unknown as Record, + context.lastKnownContextWindow, + context.lastKnownTotalProcessedTokens, ); + yield* emitThreadTokenUsage(context, compactedUsage, { + rawMethod: "claude/system/compact_boundary", + rawPayload: message, + }); yield* offerRuntimeEvent(context.sessionIncarnationId, { ...base, type: "thread.state.changed", payload: { state: "compacted", + ...(compactedUsage?.lastUsedTokens !== undefined + ? { beforeTokens: compactedUsage.lastUsedTokens } + : {}), + ...(compactedUsage?.usedTokens !== undefined + ? { afterTokens: compactedUsage.usedTokens } + : {}), detail: message, }, }); return; + } case "hook_started": yield* offerRuntimeEvent(context.sessionIncarnationId, { ...base, @@ -5162,6 +5171,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( sessionModelSwitch: "in-session", conversationRollback: BUILT_IN_ADAPTER_CONVERSATION_ROLLBACK_MODES.claude, }, + compaction: { type: "slash-command", command: "/compact" }, startSession, sendTurn, interruptTurn, diff --git a/apps/server/src/provider/Layers/ClaudeProvider.ts b/apps/server/src/provider/Layers/ClaudeProvider.ts index 811226b8c..e1a8638bc 100644 --- a/apps/server/src/provider/Layers/ClaudeProvider.ts +++ b/apps/server/src/provider/Layers/ClaudeProvider.ts @@ -24,6 +24,7 @@ import { import { buildServerProvider, + COMPACT_SLASH_COMMAND, DEFAULT_TIMEOUT_MS, isCommandMissingCause, parseGenericCliVersion, @@ -616,13 +617,7 @@ export const checkClaudeProviderStatus = Effect.fn("checkClaudeProviderStatus")( ? yield* resolveCapabilities(claudeSettings).pipe(Effect.orElseSucceed(() => undefined)) : undefined; const skills = yield* discoverClaudeSkills(claudeSettings, cwd, resolvedEnvironment); - const slashCommands = [ - { - name: "compact", - description: "Summarize the conversation and reduce context usage", - }, - ...(capabilities?.slashCommands ?? []), - ]; + const slashCommands = [COMPACT_SLASH_COMMAND, ...(capabilities?.slashCommands ?? [])]; const dedupedSlashCommands = dedupeSlashCommands(slashCommands); if (!capabilities) { diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 3778d5561..9a9349dd4 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -86,6 +86,8 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape { }), ); + public readonly compactThread = Effect.void; + public readonly interruptTurnImpl = vi.fn( (_turnId?: TurnId): Promise => Promise.resolve(undefined), ); @@ -339,6 +341,50 @@ sessionErrorLayer("CodexAdapterLive session errors", (it) => { }), ); + it.effect("compacts the active Codex thread and emits compacted state", () => + Effect.gen(function* () { + const adapter = yield* CodexAdapter; + const threadId = asThreadId("thread-compact"); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("codex"), + threadId, + runtimeMode: "full-access", + }); + const runtime = sessionRuntimeFactory.lastRuntime; + NodeAssert.ok(runtime); + const compactedEventFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.type === "thread.state.changed"), + Stream.runHead, + Effect.forkChild, + ); + const codexCompaction = adapter.compaction; + NodeAssert.equal(codexCompaction?.type, "native"); + if (codexCompaction?.type !== "native") throw new Error("expected native compaction"); + yield* codexCompaction.start(threadId); + yield* runtime.emit({ + id: asEventId("evt-compaction-item-completed"), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + method: "item/completed", + threadId, + payload: { + completedAtMs: 1_778_000_000_000, + threadId: "provider-thread-1", + turnId: "provider-compact-turn", + item: { + id: "provider-compact-item", + type: "contextCompaction", + }, + }, + }); + const event = Option.getOrThrow(yield* Fiber.join(compactedEventFiber)); + NodeAssert.ok(event.type === "thread.state.changed"); + NodeAssert.equal(event.payload.state, "compacted"); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("uploads feedback for the active Codex thread", () => Effect.gen(function* () { const adapter = yield* CodexAdapter; diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index 19038499c..73bc049ab 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -1220,7 +1220,18 @@ function mapToRuntimeEvents( ]; } const completed = mapItemLifecycle(event, canonicalThreadId, "item.completed"); - return completed ? [completed] : []; + if (!completed || itemType !== "context_compaction") { + return completed ? [completed] : []; + } + return [ + completed, + { + ...runtimeEventBase(event, canonicalThreadId), + eventId: EventId.make(`${event.id}:thread-compacted`), + type: "thread.state.changed", + payload: { state: "compacted" }, + }, + ]; } if ( @@ -1998,6 +2009,13 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( ), ); + const compactThread = Effect.fn("compactThread")(function* (threadId: ThreadId) { + const session = yield* requireSession(threadId); + yield* session.runtime.compactThread.pipe( + Effect.mapError((cause) => mapCodexRuntimeError(threadId, "thread/compact/start", cause)), + ); + }); + const readThread: CodexAdapterShape["readThread"] = (threadId) => requireSession(threadId).pipe( Effect.flatMap((session) => session.runtime.readThread), @@ -2134,6 +2152,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( }, startSession, sendTurn, + compaction: { type: "native", start: compactThread }, interruptTurn, readThread, rollbackThread, diff --git a/apps/server/src/provider/Layers/CodexProvider.ts b/apps/server/src/provider/Layers/CodexProvider.ts index e5882c585..85dc8c028 100644 --- a/apps/server/src/provider/Layers/CodexProvider.ts +++ b/apps/server/src/provider/Layers/CodexProvider.ts @@ -34,6 +34,7 @@ import { codexAppServerArgs, resolveCodexLaunchArgs } from "./codexLaunchArgs.ts import { AUTH_PROBE_TIMEOUT_MS, buildServerProvider, + COMPACT_SLASH_COMMAND, type ServerProviderDraft, } from "../providerSnapshot.ts"; import { resolveProviderHomePath } from "../../pathExpansion.ts"; @@ -790,6 +791,7 @@ export const checkCodexProviderStatus = Effect.fn("checkCodexProviderStatus")(fu models: snapshot.models, skills: snapshot.skills, slashCommands: [ + COMPACT_SLASH_COMMAND, { name: "feedback", description: "Send this thread and Codex logs to OpenAI", diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index 5712e932a..79204b63d 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -196,6 +196,7 @@ export interface CodexSessionRuntimeShape { readonly sendTurn: ( input: CodexSessionRuntimeSendTurnInput, ) => Effect.Effect; + readonly compactThread: Effect.Effect; readonly interruptTurn: (turnId?: TurnId) => Effect.Effect; readonly readThread: Effect.Effect; readonly rollbackThread: ( @@ -2333,6 +2334,10 @@ export const makeCodexSessionRuntime = ( return { start, getSession: Ref.get(sessionRef), + compactThread: Effect.gen(function* () { + const providerThreadId = yield* readProviderThreadId; + yield* client.request("thread/compact/start", { threadId: providerThreadId }); + }), sendTurn: (input) => Effect.gen(function* () { const providerThreadId = yield* readProviderThreadId; diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index 7bb61bbe5..248cbf9bd 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -1276,6 +1276,7 @@ export function makeCursorAdapter( sessionModelSwitch: "in-session", conversationRollback: BUILT_IN_ADAPTER_CONVERSATION_ROLLBACK_MODES.cursor, }, + compaction: { type: "slash-command", command: "/compress" }, startSession, sendTurn, interruptTurn, diff --git a/apps/server/src/provider/Layers/CursorProvider.ts b/apps/server/src/provider/Layers/CursorProvider.ts index ca23c01b4..1d80fa20b 100644 --- a/apps/server/src/provider/Layers/CursorProvider.ts +++ b/apps/server/src/provider/Layers/CursorProvider.ts @@ -36,6 +36,7 @@ import { buildBooleanOptionDescriptor, buildSelectOptionDescriptor, buildServerProvider, + COMPACT_SLASH_COMMAND, collectStreamAsString, isCommandMissingCause, providerModelsFromSettings, @@ -659,6 +660,7 @@ export function buildCursorProviderSnapshot(input: { input.cursorSettings.customModels, EMPTY_CAPABILITIES, ), + slashCommands: [COMPACT_SLASH_COMMAND], probe: { installed: true, version: input.parsed.version, diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index b2ebe59d8..ec9c7bc53 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -2195,6 +2195,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte sessionModelSwitch: "in-session", conversationRollback: BUILT_IN_ADAPTER_CONVERSATION_ROLLBACK_MODES.grok, }, + compaction: { type: "slash-command", command: "/compact" }, startSession, sendTurn, interruptTurn, diff --git a/apps/server/src/provider/Layers/GrokProvider.ts b/apps/server/src/provider/Layers/GrokProvider.ts index af6ea6e56..bd47035f6 100644 --- a/apps/server/src/provider/Layers/GrokProvider.ts +++ b/apps/server/src/provider/Layers/GrokProvider.ts @@ -22,6 +22,7 @@ import { resolveSpawnCommand } from "@t3tools/shared/shell"; import { AUTH_PROBE_TIMEOUT_MS, buildServerProvider, + COMPACT_SLASH_COMMAND, isCommandMissingCause, parseGenericCliVersion, providerModelsFromSettings, @@ -501,6 +502,7 @@ export const checkGrokProviderStatus = Effect.fn("checkGrokProviderStatus")(func checkedAt, models, skills, + slashCommands: [COMPACT_SLASH_COMMAND], probe: { installed: true, version, diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index 900b19f7d..cbf3a96a9 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -88,6 +88,7 @@ const runtimeMock = { messageCalls: [] as Array<{ sessionID: string; messageID: string }>, messageFailures: 0, promptCalls: [] as Array, + summarizeCalls: [] as Array, promptAsyncError: null as Error | null, promptAsyncImplementation: null as (() => Promise) | null, autoPromptEcho: true, @@ -147,6 +148,7 @@ const runtimeMock = { this.state.messageCalls.length = 0; this.state.messageFailures = 0; this.state.promptCalls.length = 0; + this.state.summarizeCalls.length = 0; this.state.promptAsyncError = null; this.state.promptAsyncImplementation = null; this.state.autoPromptEcho = true; @@ -349,6 +351,10 @@ const OpenCodeRuntimeTestDouble: OpenCodeRuntimeShape = { }); } }, + summarize: async (input: unknown) => { + runtimeMock.state.summarizeCalls.push(input); + return { data: true }; + }, messages: async () => ({ data: runtimeMock.state.messages }), message: async ({ sessionID, messageID }: { sessionID: string; messageID: string }) => { runtimeMock.state.messageCalls.push({ sessionID, messageID }); @@ -1055,6 +1061,42 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("compacts through the native OpenCode session API", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-compact"); + runtimeMock.state.subscribedEvents.push({ + type: "session.compacted", + properties: { sessionID: "http://127.0.0.1:9999/session" }, + }); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId), + Stream.take(3), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const openCodeCompaction = adapter.compaction; + NodeAssert.equal(openCodeCompaction?.type, "native"); + if (openCodeCompaction?.type !== "native") throw new Error("expected native compaction"); + yield* openCodeCompaction.start( + threadId, + createModelSelection(ProviderInstanceId.make("opencode"), "openai/gpt-5"), + ); + const summarizeCall = runtimeMock.state.summarizeCalls[0] as Record; + NodeAssert.equal(summarizeCall.modelID, "gpt-5"); + const events = Array.from(yield* Fiber.join(eventsFiber)); + yield* adapter.stopSession(threadId); + const compacted = events.some( + (event) => event.type === "thread.state.changed" && event.payload.state === "compacted", + ); + NodeAssert.equal(compacted, true); + }), + ); it.effect("falls back to a fresh session when the persisted session is gone", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 29e794f01..6252cbeb0 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -4,6 +4,7 @@ import { ProviderDriverKind, ProviderInstanceId, type ProviderRuntimeEvent, + type ProviderSendTurnInput, type ProviderSession, RuntimeItemId, RuntimeRequestId, @@ -2214,7 +2215,6 @@ export function makeOpenCodeAdapter( } break; } - case "session.compacted": { // Surfaces OpenCode context compaction the same way Claude's // compact_boundary does: ingestion turns it into a "Context @@ -3310,6 +3310,72 @@ export function makeOpenCodeAdapter( ); }); + const compactThread = Effect.fn("compactThread")(function* ( + threadId: ThreadId, + requestedModelSelection?: ProviderSendTurnInput["modelSelection"], + ) { + const context = yield* ensureSessionContext(sessions, threadId); + yield* awaitOpenCodeContextReady(context); + const modelSelection = + requestedModelSelection ?? + (context.session.model + ? { instanceId: boundInstanceId, model: context.session.model } + : undefined); + if (modelSelection !== undefined && modelSelection.instanceId !== boundInstanceId) { + return yield* new ProviderAdapterValidationError({ + provider: PROVIDER, + operation: "compactThread", + issue: `OpenCode model selection is bound to instance '${modelSelection.instanceId}', expected '${boundInstanceId}'.`, + }); + } + const parsedModel = parseOpenCodeModelSlug(modelSelection?.model); + if (!parsedModel) { + return yield* new ProviderAdapterValidationError({ + provider: PROVIDER, + operation: "compactThread", + issue: "OpenCode compaction requires an active 'provider/model' selection.", + }); + } + yield* context.promptSemaphore.withPermit( + Effect.gen(function* () { + if (sessions.get(threadId) !== context || (yield* Ref.get(context.stopped))) { + return yield* Effect.interrupt; + } + if (context.activeTurnId !== undefined) { + return yield* new ProviderAdapterValidationError({ + provider: PROVIDER, + operation: "compactThread", + issue: "OpenCode cannot compact while a turn is running.", + }); + } + yield* runOpenCodeSdk("session.summarize", (signal) => + context.client.session.summarize( + { + sessionID: context.openCodeSessionId, + ...parsedModel, + auto: false, + }, + { signal }, + ), + ).pipe( + Effect.timeout("10 minutes"), + Effect.catchTags({ + OpenCodeRuntimeError: (cause) => Effect.fail(toRequestError(cause)), + TimeoutError: (cause) => + Effect.fail( + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "session.summarize", + detail: "OpenCode session compaction did not complete within 10 minutes.", + cause, + }), + ), + }), + Effect.asVoid, + ); + }), + ); + }); const interruptTurn: OpenCodeAdapterShape["interruptTurn"] = Effect.fn("interruptTurn")( function* (threadId, turnId) { const context = yield* ensureSessionContext(sessions, threadId); @@ -3666,6 +3732,7 @@ export function makeOpenCodeAdapter( }, startSession, sendTurn, + compaction: { type: "native", start: compactThread }, interruptTurn, respondToRequest, respondToUserInput, diff --git a/apps/server/src/provider/Layers/OpenCodeProvider.ts b/apps/server/src/provider/Layers/OpenCodeProvider.ts index 8a1ab4c7f..aae28639e 100644 --- a/apps/server/src/provider/Layers/OpenCodeProvider.ts +++ b/apps/server/src/provider/Layers/OpenCodeProvider.ts @@ -13,6 +13,7 @@ import { createModelCapabilities } from "@t3tools/shared/model"; import { compareSemverVersions } from "@t3tools/shared/semver"; import { buildServerProvider, + COMPACT_SLASH_COMMAND, nonEmptyTrimmed, parseGenericCliVersion, providerModelsFromSettings, @@ -526,6 +527,7 @@ export const checkOpenCodeProviderStatus = Effect.fn("checkOpenCodeProviderStatu checkedAt, models, skills, + slashCommands: [COMPACT_SLASH_COMMAND], probe: { installed: true, version, diff --git a/apps/server/src/provider/Layers/ProviderRegistry.test.ts b/apps/server/src/provider/Layers/ProviderRegistry.test.ts index 0ec631c13..a54254f45 100644 --- a/apps/server/src/provider/Layers/ProviderRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderRegistry.test.ts @@ -55,6 +55,7 @@ import { resolveProviderStatusCachePath, writeProviderStatusCache, } from "../providerStatusCache.ts"; +import { COMPACT_SLASH_COMMAND } from "../providerSnapshot.ts"; import { providerBackendsFromCapacityRefresh, type ProviderInstance } from "../ProviderDriver.ts"; import * as ProviderInstanceRegistry from "../Services/ProviderInstanceRegistry.ts"; import * as ProviderRegistry from "../Services/ProviderRegistry.ts"; @@ -571,7 +572,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te shortDescription: "Debug failing GitHub Actions checks", }, ]); - assert.deepStrictEqual(status.slashCommands, [ + assert.deepStrictEqual(status.slashCommands.slice(1), [ { name: "feedback", description: "Send this thread and Codex logs to OpenAI", @@ -3248,11 +3249,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te }), ); - assert.deepStrictEqual(status.slashCommands, [ - { - name: "compact", - description: "Summarize the conversation and reduce context usage", - }, + assert.deepStrictEqual(status.slashCommands.slice(1), [ { name: "review", description: "Review a pull request", @@ -3296,10 +3293,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ); assert.deepStrictEqual(status.slashCommands, [ - { - name: "compact", - description: "Summarize the conversation and reduce context usage", - }, + COMPACT_SLASH_COMMAND, { name: "ui", description: "Explore and refine UI", diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index bc18a271d..587d651cc 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -88,7 +88,10 @@ import { import * as ServerConfig from "../../config.ts"; import * as ServerSettings from "../../serverSettings.ts"; import * as AnalyticsService from "../../telemetry/AnalyticsService.ts"; -import { makeAdapterRegistryMock } from "../testUtils/providerAdapterRegistryMock.ts"; +import { + makeAdapterRegistryMock, + makeInstanceAdapterRegistryMock, +} from "../testUtils/providerAdapterRegistryMock.ts"; const defaultServerSettingsLayer = ServerSettings.ServerSettingsService.layerTest(); const serverConfigTestLayer = ServerConfig.layerTest(process.cwd(), process.cwd()).pipe( @@ -205,6 +208,19 @@ function makeFakeCodexAdapter( Effect.void, ); + const compactThread = vi.fn( + (threadId: ThreadId): Effect.Effect => + Effect.sync(() => + emit({ + type: "thread.state.changed", + eventId: asEventId("evt-native-compact"), + provider, + createdAt: "2026-01-01T00:00:00.000Z", + threadId, + payload: { state: "compacted" }, + }), + ), + ); const respondToRequest = vi.fn( ( _threadId: ThreadId, @@ -441,6 +457,13 @@ function makeFakeCodexAdapter( startSession, sendTurn, followUp, + ...(provider === CODEX_DRIVER + ? { compaction: { type: "native", start: compactThread } } + : provider === CURSOR_DRIVER + ? { compaction: { type: "slash-command", command: "/compress" } } + : provider === CLAUDE_AGENT_DRIVER + ? { compaction: { type: "slash-command", command: "/compact" } } + : {}), interruptTurn, respondToRequest, respondToUserInput, @@ -475,8 +498,17 @@ function makeFakeCodexAdapter( }, }; + // Real adapters stamp every event with the incarnation that produced it, and + // ProviderService drops unstamped events so a replaced adapter cannot write + // into the new session. Stamp here too, or these fakes would exercise a path + // production never takes. const emit = (event: LegacyProviderRuntimeEvent): void => { - Effect.runSync(PubSub.publish(runtimeEventPubSub, event as unknown as ProviderRuntimeEvent)); + const incarnationId = sessions.get(event.threadId as ThreadId)?.sessionIncarnationId; + const stamped = + incarnationId === undefined || event.sessionIncarnationId !== undefined + ? event + : { ...event, sessionIncarnationId: incarnationId }; + Effect.runSync(PubSub.publish(runtimeEventPubSub, stamped as unknown as ProviderRuntimeEvent)); }; const removeSession = (threadId: ThreadId): void => { @@ -535,6 +567,7 @@ function makeFakeCodexAdapter( abortSessionCompaction, setSessionAutoCompaction, refineSessionHarness, + compactThread, interruptTurn, respondToRequest, respondToUserInput, @@ -565,17 +598,20 @@ const hasMetricSnapshot = ( function makeProviderServiceLayer( input: { readonly directory?: ProviderSessionDirectory.ProviderSessionDirectory["Service"]; + readonly registry?: ProviderAdapterRegistry.ProviderAdapterRegistry["Service"]; } = {}, ) { const startReservationCounts: number[] = []; const codex = makeFakeCodexAdapter(); const claude = makeFakeCodexAdapter(CLAUDE_AGENT_DRIVER); const cursor = makeFakeCodexAdapter(CURSOR_DRIVER, { supportsSideQuestions: false }); - const registry = makeAdapterRegistryMock({ - [ProviderDriverKind.make("codex")]: codex.adapter, - [ProviderDriverKind.make("claudeAgent")]: claude.adapter, - [ProviderDriverKind.make("cursor")]: cursor.adapter, - }); + const registry = + input.registry ?? + makeAdapterRegistryMock({ + [ProviderDriverKind.make("codex")]: codex.adapter, + [ProviderDriverKind.make("claudeAgent")]: claude.adapter, + [ProviderDriverKind.make("cursor")]: cursor.adapter, + }); const providerAdapterLayer = Layer.succeed( ProviderAdapterRegistry.ProviderAdapterRegistry, @@ -1030,6 +1066,125 @@ it.effect("ProviderServiceLive rejects new sessions for disabled custom instance const routing = makeProviderServiceLayer(); +const customCompactionDriver = ProviderDriverKind.make("custom-compaction-provider"); +const nativeCompactionInstanceId = ProviderInstanceId.make("native-compaction"); +const slashCompactionInstanceId = ProviderInstanceId.make("slash-compaction"); +const unsupportedCompactionInstanceId = ProviderInstanceId.make("unsupported-compaction"); +const customNativeCompaction = makeFakeCodexAdapter(customCompactionDriver); +const customSlashCompaction = makeFakeCodexAdapter(customCompactionDriver); +const unsupportedCompaction = makeFakeCodexAdapter(customCompactionDriver); +const declaredCompaction = makeProviderServiceLayer({ + registry: makeInstanceAdapterRegistryMock([ + [ + nativeCompactionInstanceId, + { + ...customNativeCompaction.adapter, + compaction: { type: "native", start: customNativeCompaction.compactThread }, + }, + ], + [ + slashCompactionInstanceId, + { + ...customSlashCompaction.adapter, + compaction: { type: "slash-command", command: "/reduce-context" }, + }, + ], + [unsupportedCompactionInstanceId, unsupportedCompaction.adapter], + ]), +}); + +declaredCompaction.layer("ProviderService declared compaction", (it) => { + it.effect("starts declared native compaction instead of sending a prompt", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("custom-native-compaction"); + const requestId = MessageId.make("custom-native-request"); + yield* provider.startSession(threadId, { + providerInstanceId: nativeCompactionInstanceId, + threadId, + runtimeMode: "full-access", + }); + const compactedEventFiber = yield* provider.streamEvents.pipe( + Stream.filter( + (event) => event.threadId === threadId && event.type === "thread.state.changed", + ), + Stream.runHead, + Effect.forkChild({ startImmediately: true }), + ); + yield* advanceTestClock(50); + yield* provider.compactThread(threadId, undefined, requestId); + const compacted = Option.getOrThrow(yield* Fiber.join(compactedEventFiber)); + assert.equal(compacted.requestId, String(requestId)); + assert.equal(customNativeCompaction.compactThread.mock.calls.length, 1); + assert.equal(customNativeCompaction.sendTurn.mock.calls.length, 0); + yield* provider.stopSession({ threadId }); + }), + ); + + it.effect("sends the declared slash command as the compaction turn", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("custom-slash-compaction"); + const requestId = MessageId.make("custom-slash-request"); + const modelSelection = createModelSelection(slashCompactionInstanceId, "custom-model"); + yield* provider.startSession(threadId, { + providerInstanceId: slashCompactionInstanceId, + threadId, + runtimeMode: "full-access", + }); + const compactedEventFiber = yield* provider.streamEvents.pipe( + Stream.filter( + (event) => event.threadId === threadId && event.type === "thread.state.changed", + ), + Stream.runHead, + Effect.forkChild({ startImmediately: true }), + ); + const compactFiber = yield* provider + .compactThread(threadId, modelSelection, requestId) + .pipe(Effect.forkChild); + yield* advanceTestClock(50); + customSlashCompaction.emit({ + type: "turn.completed", + eventId: asEventId("custom-slash-completed"), + provider: customCompactionDriver, + createdAt: "2026-01-01T00:00:01.000Z", + threadId, + turnId: asTurnId(`turn-${threadId}`), + payload: { state: "completed" }, + }); + yield* Fiber.join(compactFiber); + const compacted = Option.getOrThrow(yield* Fiber.join(compactedEventFiber)); + assert.equal(compacted.requestId, String(requestId)); + assert.equal(customSlashCompaction.compactThread.mock.calls.length, 0); + assert.equal(customSlashCompaction.sendTurn.mock.calls.length, 1); + assert.equal(customSlashCompaction.sendTurn.mock.calls[0]?.[0].input, "/reduce-context"); + assert.deepEqual( + customSlashCompaction.sendTurn.mock.calls[0]?.[0].modelSelection, + modelSelection, + ); + yield* provider.stopSession({ threadId }); + }), + ); + + it.effect("rejects compaction for adapters without a declared strategy", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("custom-unsupported-compaction"); + yield* provider.startSession(threadId, { + providerInstanceId: unsupportedCompactionInstanceId, + threadId, + runtimeMode: "full-access", + }); + const failure = yield* provider.compactThread(threadId).pipe(Effect.flip); + assert.instanceOf(failure, ProviderValidationError); + assert.include(failure.message, "does not support context compaction"); + assert.equal(unsupportedCompaction.sendTurn.mock.calls.length, 0); + assert.equal(unsupportedCompaction.compactThread.mock.calls.length, 0); + yield* provider.stopSession({ threadId }); + }), + ); +}); + it.effect( "ProviderServiceLive uploads feedback through the adapter that recovered the session", () => @@ -2025,6 +2180,243 @@ routing.layer("ProviderServiceLive routing", (it) => { }), ); + it.effect("marks a successful fallback compaction as compacted", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("thread-compact-cursor"); + yield* provider.startSession(threadId, { + provider: CURSOR_DRIVER, + providerInstanceId: ProviderInstanceId.make("cursor"), + threadId, + runtimeMode: "full-access", + }); + const compactedEventFiber = yield* provider.streamEvents.pipe( + Stream.filter((event) => event.type === "thread.state.changed"), + Stream.runHead, + Effect.forkChild, + ); + const requestId = MessageId.make("message-compact-cursor"); + const compactFiber = yield* provider + .compactThread(threadId, undefined, requestId) + .pipe(Effect.forkChild); + yield* advanceTestClock(50); + routing.cursor.emit({ + type: "turn.completed", + eventId: asEventId("evt-cursor-stale-turn-completed"), + provider: CURSOR_DRIVER, + createdAt: "2026-01-01T00:00:00.500Z", + threadId, + turnId: asTurnId("turn-before-compaction"), + payload: { state: "completed" }, + }); + yield* Effect.yieldNow; + assert.equal(compactFiber.pollUnsafe(), undefined); + routing.cursor.emit({ + type: "turn.completed", + eventId: asEventId("evt-cursor-compact-completed"), + provider: CURSOR_DRIVER, + createdAt: "2026-01-01T00:00:01.000Z", + threadId, + turnId: asTurnId(`turn-${threadId}`), + payload: { state: "completed" }, + }); + yield* Fiber.join(compactFiber); + + const compacted = yield* Fiber.join(compactedEventFiber); + assert.equal(compacted._tag, "Some"); + if (Option.isSome(compacted)) { + assert.equal(compacted.value.requestId, String(requestId)); + } + + const observedEvents = yield* Ref.make>([]); + const observedEventsFiber = yield* provider.streamEvents.pipe( + Stream.filter( + (event) => + event.threadId === threadId && + event.type === "thread.state.changed" && + event.payload.state === "compacted", + ), + Stream.runForEach((event) => Ref.update(observedEvents, (events) => [...events, event])), + Effect.forkChild, + ); + const observedRequestId = MessageId.make("message-observed-compact-cursor"); + const observedCompactFiber = yield* provider + .compactThread(threadId, undefined, observedRequestId) + .pipe(Effect.forkChild); + yield* advanceTestClock(50); + routing.cursor.emit({ + type: "thread.state.changed", + eventId: asEventId("evt-cursor-provider-compacted"), + provider: CURSOR_DRIVER, + createdAt: "2026-01-01T00:00:02.000Z", + threadId, + turnId: asTurnId(`turn-${threadId}`), + payload: { state: "compacted" }, + }); + routing.cursor.emit({ + type: "turn.completed", + eventId: asEventId("evt-cursor-observed-compact-completed"), + provider: CURSOR_DRIVER, + createdAt: "2026-01-01T00:00:03.000Z", + threadId, + turnId: asTurnId(`turn-${threadId}`), + payload: { state: "completed" }, + }); + yield* Fiber.join(observedCompactFiber); + yield* Effect.yieldNow; + const observed = yield* Ref.get(observedEvents); + assert.equal(observed.length, 1); + assert.equal(observed[0]?.requestId, String(observedRequestId)); + yield* Fiber.interrupt(observedEventsFiber); + + const failedStartEventId = asEventId("evt-cursor-failed-compact-start"); + const failedStartEventFiber = yield* provider.streamEvents.pipe( + Stream.filter((event) => event.eventId === failedStartEventId), + Stream.runHead, + Effect.forkChild, + ); + routing.cursor.sendTurn.mockImplementationOnce((input) => + Effect.gen(function* () { + routing.cursor.emit({ + type: "turn.completed", + eventId: failedStartEventId, + provider: CURSOR_DRIVER, + createdAt: "2026-01-01T00:00:04.000Z", + threadId: input.threadId, + turnId: asTurnId("turn-cursor-failed-compact-start"), + payload: { state: "failed" }, + }); + yield* Effect.yieldNow; + return yield* new ProviderAdapterRequestError({ + provider: String(CURSOR_DRIVER), + method: "turn/start", + detail: "Failed after emitting a terminal event.", + }); + }), + ); + const failedStart = yield* provider.compactThread(threadId).pipe(Effect.result); + assert.equal(failedStart._tag, "Failure"); + assert.equal(Option.isSome(yield* Fiber.join(failedStartEventFiber)), true); + yield* provider.stopSession({ threadId }); + }), + ); + + it.effect("serializes native compaction and quarantines timed-out completions", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("thread-compact-timeout"); + yield* provider.startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }); + routing.codex.compactThread.mockClear(); + routing.codex.compactThread.mockImplementationOnce(() => Effect.never); + + const resultFiber = yield* provider + .compactThread(threadId) + .pipe(Effect.result, Effect.forkChild); + yield* advanceTestClock(50); + const concurrent = yield* provider.compactThread(threadId).pipe(Effect.result); + assert.equal(concurrent._tag, "Failure"); + assert.equal(routing.codex.compactThread.mock.calls.length, 1); + + routing.cursor.emit({ + type: "thread.state.changed", + eventId: asEventId("evt-stale-provider-compact"), + provider: CURSOR_DRIVER, + createdAt: "2026-01-01T00:00:00.100Z", + threadId, + payload: { state: "compacted" }, + }); + yield* Effect.yieldNow; + assert.equal(resultFiber.pollUnsafe(), undefined); + + yield* advanceTestClock(600_001); + const result = yield* Fiber.join(resultFiber); + assert.equal(result._tag, "Failure"); + if (result._tag === "Failure") { + assert.equal(result.failure._tag, "ProviderAdapterRequestError"); + } + + const blockedRetry = yield* provider.compactThread(threadId).pipe(Effect.result); + assert.equal(blockedRetry._tag, "Failure"); + assert.equal(routing.codex.compactThread.mock.calls.length, 1); + + routing.codex.emit({ + type: "thread.state.changed", + eventId: asEventId("evt-native-compact-late"), + provider: CODEX_DRIVER, + createdAt: "2026-01-01T00:10:01.000Z", + threadId, + payload: { state: "compacted" }, + }); + yield* Effect.yieldNow; + yield* provider.compactThread(threadId); + assert.equal(routing.codex.compactThread.mock.calls.length, 2); + + routing.codex.compactThread.mockImplementationOnce(() => Effect.void); + const stoppedResultFiber = yield* provider + .compactThread(threadId) + .pipe(Effect.result, Effect.forkChild); + yield* advanceTestClock(50); + yield* provider.stopSession({ threadId }); + const stoppedResult = yield* Fiber.join(stoppedResultFiber); + assert.equal(stoppedResult._tag, "Failure"); + + yield* provider.startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }); + yield* provider.compactThread(threadId); + assert.equal(routing.codex.compactThread.mock.calls.length, 4); + yield* provider.stopSession({ threadId }); + }), + ); + + it.effect("times out fallback compaction when its turn never settles", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("thread-compact-fallback-timeout"); + yield* provider.startSession(threadId, { + provider: CURSOR_DRIVER, + providerInstanceId: ProviderInstanceId.make("cursor"), + threadId, + runtimeMode: "full-access", + }); + + const resultFiber = yield* provider + .compactThread(threadId) + .pipe(Effect.result, Effect.forkChild); + yield* advanceTestClock(600_001); + const result = yield* Fiber.join(resultFiber); + assert.equal(result._tag, "Failure"); + + routing.cursor.sendTurn.mockImplementationOnce((input) => + Effect.succeed({ + threadId: input.threadId, + turnId: asTurnId("turn-compact-fallback-retry"), + }), + ); + const retryFiber = yield* provider.compactThread(threadId).pipe(Effect.forkChild); + yield* advanceTestClock(50); + routing.cursor.emit({ + type: "turn.completed", + eventId: asEventId("evt-compact-fallback-retry-completed"), + provider: CURSOR_DRIVER, + createdAt: "2026-01-01T00:10:02.000Z", + threadId, + turnId: asTurnId("turn-compact-fallback-retry"), + payload: { state: "completed" }, + }); + yield* Fiber.join(retryFiber); + yield* provider.stopSession({ threadId }); + }), + ); + it.effect("routes feedback to the Codex adapter and returns its feedback ID", () => 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 f8d50400d..27ea9aaa7 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -13,6 +13,8 @@ import * as NodeCrypto from "node:crypto"; */ import { CommandId, + EventId, + MessageId, ModelSelection, NonNegativeInt, PROVIDER_SESSION_AGENT_DEPTH_MAX_SETTABLE, @@ -42,6 +44,7 @@ import { ProviderRespondToRequestInput, ProviderRespondToUserInputInput, ProviderRespondToInteractionInput, + RuntimeRequestId, ProviderSendTurnInput, ProviderSessionStartInput, ProviderStopSessionInput, @@ -54,6 +57,7 @@ import { import { expandAssistantCitationsForProvider } from "@t3tools/shared/assistantCitations"; 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 FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; @@ -102,6 +106,19 @@ import * as ServerSettings from "../../serverSettings.ts"; import { RollbackSagaRepository } from "../../persistence/Services/RollbackSagas.ts"; const isModelSelection = Schema.is(ModelSelection); +/** How long a manual context compaction may run before ProviderService gives up on it. */ +const COMPACTION_COMPLETION_TIMEOUT = "10 minutes"; + +interface PendingCompaction { + readonly completion: Deferred.Deferred; + readonly native: boolean; + readonly providerInstanceId: ProviderInstanceId; + readonly requestId: MessageId | undefined; + readonly earlyEvents: ProviderRuntimeEvent[]; + compactedEventObserved: boolean; + expectedTurnId: TurnId | undefined; +} + /** * Hook for tests that want to override the canonical event logger pulled * from `ProviderEventLoggers`. Production wiring leaves this undefined and @@ -462,6 +479,15 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( }); return next; }).pipe(Effect.tap(() => reportStartReservationCount)); + const pendingCompactions = new Map(); + const timedOutNativeCompactions = new Set(); + const settleCompaction = (threadId: ThreadId, pending: PendingCompaction, terminal: string) => + Effect.gen(function* () { + if (pendingCompactions.get(threadId) !== pending) return false; + pendingCompactions.delete(threadId); + yield* Deferred.succeed(pending.completion, terminal); + return true; + }); const nowIso = Effect.map(DateTime.now, DateTime.formatIso); const requireAdapterGenerationCurrent = Effect.fnUntraced(function* ( adapter: ProviderAdapterShape, @@ -575,6 +601,107 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( yield* PubSub.publish(runtimeEventPubSub, event); }); + const isCompactedEvent = ( + event: ProviderRuntimeEvent, + ): event is Extract => + event.type === "thread.state.changed" && event.payload.state === "compacted"; + const withCompactionRequestId = ( + event: ProviderRuntimeEvent, + pending: PendingCompaction, + ): ProviderRuntimeEvent => + pending.requestId === undefined + ? event + : { + ...event, + requestId: RuntimeRequestId.make(String(pending.requestId)), + }; + const compactionTerminal = (event: ProviderRuntimeEvent): string | null => + event.type === "turn.completed" + ? event.payload.state + : event.type === "runtime.error" || event.type === "turn.aborted" + ? event.type + : null; + const processFallbackCompactionEvent = ( + pending: PendingCompaction, + event: ProviderRuntimeEvent, + ): Effect.Effect => + Effect.gen(function* () { + if (pendingCompactions.get(event.threadId) !== pending) { + yield* publishRuntimeEvent(event); + return; + } + const matchesTurn = event.turnId !== undefined && event.turnId === pending.expectedTurnId; + if (matchesTurn && isCompactedEvent(event)) { + pending.compactedEventObserved = true; + yield* publishRuntimeEvent(withCompactionRequestId(event, pending)); + return; + } + yield* publishRuntimeEvent(event); + const terminal = compactionTerminal(event); + if (!matchesTurn || terminal === null) return; + const settled = yield* settleCompaction(event.threadId, pending, terminal); + if (!settled || terminal !== "completed" || pending.compactedEventObserved) return; + const compactedEvent = { + ...event, + eventId: EventId.make(`${event.eventId}:context-compaction`), + type: "thread.state.changed", + payload: { + state: "compacted", + detail: { source: "provider-native-command" }, + }, + ...(pending.requestId !== undefined + ? { requestId: RuntimeRequestId.make(String(pending.requestId)) } + : {}), + } satisfies ProviderRuntimeEvent; + yield* increment(providerRuntimeEventsTotal, { + provider: compactedEvent.provider, + eventType: compactedEvent.type, + }); + yield* publishRuntimeEvent(compactedEvent); + }); + + /** + * Publishes a runtime event, letting a manual compaction claim the events for + * the turn it started. A native run settles on its own compacted or terminal + * event; the slash-command fallback publishes through its own path so it can + * attach the request id and synthesize the compacted event. Callers keep their + * own post-publish bookkeeping either way. + */ + const publishCompactionAwareRuntimeEvent = ( + instanceId: ProviderInstanceId, + event: ProviderRuntimeEvent, + ): Effect.Effect => + Effect.gen(function* () { + if (isCompactedEvent(event) && timedOutNativeCompactions.delete(event.threadId)) { + yield* publishRuntimeEvent(event); + return; + } + const pending = pendingCompactions.get(event.threadId); + if (pending === undefined || pending.providerInstanceId !== instanceId) { + yield* publishRuntimeEvent(event); + return; + } + if (pending.native) { + const compacted = isCompactedEvent(event); + const terminal = compacted ? "completed" : compactionTerminal(event); + yield* publishRuntimeEvent(compacted ? withCompactionRequestId(event, pending) : event); + if (terminal !== null) yield* settleCompaction(event.threadId, pending, terminal); + return; + } + // The compaction turn's id is only known once turn.started lands. Hold + // terminal events from other turns until then rather than mistaking one + // for this compaction's outcome. + if ( + pending.expectedTurnId === undefined && + event.turnId !== undefined && + (isCompactedEvent(event) || compactionTerminal(event) !== null) + ) { + pending.earlyEvents.push(event); + return; + } + yield* processFallbackCompactionEvent(pending, event); + }); + const requireBindingInstanceId = ( operation: string, payload: { @@ -802,7 +929,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( provider: canonicalEvent.provider, eventType: canonicalEvent.type, }); - yield* publishRuntimeEvent(canonicalEvent); + yield* publishCompactionAwareRuntimeEvent(source.instanceId, canonicalEvent); if ( canonicalEvent.type === "turn.completed" || @@ -1365,6 +1492,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( input.modelSelection.model.trim().length > 0, }), ); + timedOutNativeCompactions.delete(threadId); // Changing runtime mode restarts the session, so the transition is only // observable here, by diffing against the mode the previous session for @@ -1784,6 +1912,130 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( } }); + const compactThread: ProviderServiceMethod<"compactThread"> = Effect.fn("compactThread")( + function* (threadId, modelSelection, requestId) { + const routed = yield* resolveRoutableSession({ + threadId, + operation: "ProviderService.compactThread", + allowRecovery: true, + }); + yield* Effect.annotateCurrentSpan({ + "provider.operation": "compact-thread", + "provider.kind": routed.adapter.provider, + "provider.thread_id": threadId, + }); + yield* McpSessionRegistry.touchActiveMcpThread(threadId); + const compaction = routed.adapter.compaction; + if (compaction === undefined) { + return yield* toValidationError( + "ProviderService.compactThread", + `Provider '${routed.adapter.provider}' does not support context compaction.`, + ); + } + const completion = yield* Deferred.make(); + const pending: PendingCompaction = { + completion, + native: compaction.type === "native", + providerInstanceId: routed.instanceId, + requestId, + earlyEvents: [], + compactedEventObserved: false, + expectedTurnId: undefined, + }; + if (compaction.type === "native" && timedOutNativeCompactions.has(threadId)) { + return yield* new ProviderAdapterRequestError({ + provider: routed.adapter.provider, + method: "thread/compact", + detail: + "The previous context compaction may still be running. Restart the provider session before retrying.", + }); + } + const claimed = yield* Effect.sync(() => { + if (pendingCompactions.has(threadId)) return false; + pendingCompactions.set(threadId, pending); + return true; + }); + if (!claimed) { + return yield* new ProviderAdapterRequestError({ + provider: routed.adapter.provider, + method: "thread/compact", + detail: "Context compaction is already in progress.", + }); + } + const clearPending = Effect.sync(() => { + if (pendingCompactions.get(threadId) === pending) { + pendingCompactions.delete(threadId); + } + }); + const awaitNativeCompaction = (start: Effect.Effect) => + start.pipe( + Effect.andThen(Deferred.await(completion)), + Effect.timeout(COMPACTION_COMPLETION_TIMEOUT), + Effect.catchTag("TimeoutError", (cause) => + Effect.sync(() => { + timedOutNativeCompactions.add(threadId); + }).pipe( + Effect.andThen( + Effect.fail( + new ProviderAdapterRequestError({ + provider: routed.adapter.provider, + method: "thread/compact", + detail: `Provider did not report completed context compaction within ${COMPACTION_COMPLETION_TIMEOUT}.`, + cause, + }), + ), + ), + ), + ), + ); + const awaitFallbackCompaction = Deferred.await(completion).pipe( + Effect.timeout(COMPACTION_COMPLETION_TIMEOUT), + Effect.mapError( + (cause) => + new ProviderAdapterRequestError({ + provider: routed.adapter.provider, + method: "turn/start", + detail: `Provider did not finish context compaction within ${COMPACTION_COMPLETION_TIMEOUT}.`, + cause, + }), + ), + ); + const terminal = yield* ( + compaction.type === "native" + ? awaitNativeCompaction(compaction.start(routed.threadId, modelSelection)) + : Effect.gen(function* () { + const turn = yield* sendTurn({ + threadId, + input: compaction.command, + ...(modelSelection !== undefined ? { modelSelection } : {}), + }).pipe( + Effect.onError(() => + Effect.forEach(pending.earlyEvents.splice(0), publishRuntimeEvent, { + discard: true, + }), + ), + ); + pending.expectedTurnId = turn.turnId; + const earlyEvents = pending.earlyEvents.splice(0); + for (const earlyEvent of earlyEvents) { + yield* processFallbackCompactionEvent(pending, earlyEvent); + } + return yield* awaitFallbackCompaction; + }) + ).pipe(Effect.ensuring(clearPending)); + if (terminal !== "completed") { + return yield* new ProviderAdapterRequestError({ + provider: routed.adapter.provider, + method: compaction.type === "native" ? "thread/compact" : "turn/start", + detail: `Context compaction ended with ${terminal}.`, + }); + } + yield* analytics.record("provider.thread.compacted", { + provider: routed.adapter.provider, + }); + }, + ); + const interruptTurn: ProviderServiceMethod<"interruptTurn"> = Effect.fn("interruptTurn")( function* (rawInput) { const input = yield* decodeInputOrValidationError({ @@ -2688,6 +2940,12 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ) { currentSessionIncarnations.delete(input.threadId); } + // A stop is a terminal outcome for any compaction this session owed. + const pendingCompaction = pendingCompactions.get(input.threadId); + if (pendingCompaction !== undefined) { + yield* settleCompaction(input.threadId, pendingCompaction, "turn.aborted"); + } + timedOutNativeCompactions.delete(input.threadId); yield* clearMcpSession(input.threadId, routed.adapter.runtimeFence); const latestBinding = Option.getOrUndefined(yield* directory.getBinding(input.threadId)); @@ -3335,6 +3593,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( startSession, sendTurn, recoverRestartSessions, + compactThread, interruptTurn, respondToRequest, respondToUserInput, diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index 907aba64a..de5530dea 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -164,6 +164,7 @@ describe("ProviderSessionReaper", () => { const providerService: ProviderServiceShape = { startSession: () => unsupported(), sendTurn: () => unsupported(), + compactThread: () => unsupported(), interruptTurn: () => unsupported(), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), diff --git a/apps/server/src/provider/Services/ProviderAdapter.ts b/apps/server/src/provider/Services/ProviderAdapter.ts index ad5792385..5e9f5c906 100644 --- a/apps/server/src/provider/Services/ProviderAdapter.ts +++ b/apps/server/src/provider/Services/ProviderAdapter.ts @@ -47,6 +47,21 @@ import type { Json } from "effect/Schema"; import type { ProviderRuntimeFence } from "../ProviderDriver.ts"; export type ProviderSessionModelSwitchMode = "in-session" | "unsupported"; + +/** + * How ProviderService runs manual context compaction for an adapter. + * Native adapters expose a start call and must emit a compacted thread state + * when they finish. Slash-command adapters get the command sent as a turn. + */ +export type ProviderCompaction = + | { + readonly type: "native"; + readonly start: ( + threadId: ThreadId, + modelSelection?: ProviderSendTurnInput["modelSelection"], + ) => Effect.Effect; + } + | { readonly type: "slash-command"; readonly command: `/${string}` }; export type ProviderConversationRollbackMode = "absolute" | "relative" | "unsupported"; export const BUILT_IN_ADAPTER_CONVERSATION_ROLLBACK_MODES = { @@ -160,6 +175,9 @@ export interface ProviderAdapterShape { input: ProviderSendTurnInput, ) => Effect.Effect; + /** Omitted when this adapter does not support manual context compaction. */ + readonly compaction?: ProviderCompaction; + /** Server-private hook. Prime uses it to replace an idle ordinary owner with a recoverable one. */ readonly prepareTurnRecovery?: (input: ProviderSendTurnInput) => Effect.Effect; diff --git a/apps/server/src/provider/Services/ProviderService.ts b/apps/server/src/provider/Services/ProviderService.ts index 55dedb23c..37c9ebe07 100644 --- a/apps/server/src/provider/Services/ProviderService.ts +++ b/apps/server/src/provider/Services/ProviderService.ts @@ -52,6 +52,7 @@ import type { ProviderStopSessionInput, ProviderUploadFeedbackInput, ProviderUploadFeedbackResult, + MessageId, ThreadId, ProviderTurnStartResult, } from "@t3tools/contracts"; @@ -95,6 +96,12 @@ export interface ProviderServiceShape { input: ProviderSendTurnInput, ) => Effect.Effect; + readonly compactThread: ( + threadId: ThreadId, + modelSelection?: ProviderSendTurnInput["modelSelection"], + requestId?: MessageId, + ) => Effect.Effect; + /** Adopt eligible surviving Prime executions before startup orphan reconciliation. */ readonly recoverRestartSessions?: (input?: { readonly unrecoverableAbsoluteRollbacks?: ReadonlySet; diff --git a/apps/server/src/provider/providerSnapshot.ts b/apps/server/src/provider/providerSnapshot.ts index 08fb7bd36..8c3552b98 100644 --- a/apps/server/src/provider/providerSnapshot.ts +++ b/apps/server/src/provider/providerSnapshot.ts @@ -25,6 +25,11 @@ export const DEFAULT_TIMEOUT_MS = 4_000; // Auth status checks involve disk/network lookups and can be slow on first run (especially Windows) export const AUTH_PROBE_TIMEOUT_MS = 10_000; +export const COMPACT_SLASH_COMMAND = { + name: "compact", + description: "Summarize the conversation and reduce context usage", +} satisfies ServerProviderSlashCommand; + export interface CommandResult { readonly stdout: string; readonly stderr: string; diff --git a/apps/server/src/provider/testUtils/providerAdapterRegistryMock.ts b/apps/server/src/provider/testUtils/providerAdapterRegistryMock.ts index ae1a12081..1e814c19b 100644 --- a/apps/server/src/provider/testUtils/providerAdapterRegistryMock.ts +++ b/apps/server/src/provider/testUtils/providerAdapterRegistryMock.ts @@ -72,3 +72,41 @@ export const makeAdapterRegistryMock = (adapters: KindAdapterMap): ProviderAdapt ), }; }; + +/** + * Build a `ProviderAdapterRegistryShape` from explicit instance ids, for tests + * that need several instances of one driver kind with different adapters. + */ +export const makeInstanceAdapterRegistryMock = ( + entries: ReadonlyArray]>, +): ProviderAdapterRegistryShape => { + const byInstanceId = new Map(entries); + const unsupported = (instanceId: ProviderInstanceId) => + new ProviderUnsupportedError({ provider: ProviderDriverKind.make(instanceId) }); + + return { + getByInstance: (instanceId) => { + const adapter = byInstanceId.get(instanceId); + return adapter ? Effect.succeed(adapter) : Effect.fail(unsupported(instanceId)); + }, + getInstanceInfo: (instanceId) => { + const adapter = byInstanceId.get(instanceId); + if (!adapter) return Effect.fail(unsupported(instanceId)); + const driverKind = ProviderDriverKind.make(adapter.provider); + return Effect.succeed({ + instanceId, + driverKind, + displayName: undefined, + enabled: true, + continuationIdentity: { + driverKind, + continuationKey: `${adapter.provider}:instance:${instanceId}`, + }, + }); + }, + listInstances: () => Effect.succeed(Array.from(byInstanceId.keys())), + subscribeChanges: Effect.flatMap(PubSub.unbounded(), (pubsub) => + PubSub.subscribe(pubsub), + ), + }; +}; diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 8f53ee77f..6f592e0a3 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -63,6 +63,7 @@ const makeProviderService = ( startSession: () => Effect.die("unused"), sendTurn: () => Effect.die("unused"), recoverRestartSessions, + compactThread: () => Effect.die("unused"), interruptTurn: () => Effect.die("unused"), respondToRequest: () => Effect.die("unused"), respondToUserInput: () => Effect.die("unused"), diff --git a/apps/web/src/components/ChatView.logic.test.ts b/apps/web/src/components/ChatView.logic.test.ts index 229bff6dd..c02b3d6ea 100644 --- a/apps/web/src/components/ChatView.logic.test.ts +++ b/apps/web/src/components/ChatView.logic.test.ts @@ -1509,8 +1509,44 @@ describe("hasServerAcknowledgedLocalDispatch", () => { expect(hasServerAcknowledgedLocalDispatch({ ...common, hasPendingApproval: true })).toBe(true); expect(hasServerAcknowledgedLocalDispatch({ ...common, hasPendingUserInput: true })).toBe(true); + expect( + hasServerAcknowledgedLocalDispatch({ + ...common, + latestTurnStartFailureId: "turn-start-failure-1", + }), + ).toBe(true); expect(hasServerAcknowledgedLocalDispatch({ ...common, threadError: "failed" })).toBe(true); }); + + it("acknowledges only a new turn-start failure", () => { + const localDispatch = { + ...createLocalDispatchSnapshot(makeThread()), + latestTurnStartFailureId: "turn-start-failure-old", + }; + const common = { + localDispatch, + phase: "ready" as const, + latestTurn: null, + latestUserMessageId: localDispatch.latestUserMessageId, + session: null, + hasPendingApproval: false, + hasPendingUserInput: false, + threadError: null, + }; + + expect( + hasServerAcknowledgedLocalDispatch({ + ...common, + latestTurnStartFailureId: "turn-start-failure-old", + }), + ).toBe(false); + expect( + hasServerAcknowledgedLocalDispatch({ + ...common, + latestTurnStartFailureId: "turn-start-failure-new", + }), + ).toBe(true); + }); }); describe("server pull request link changes", () => { diff --git a/apps/web/src/components/ChatView.logic.ts b/apps/web/src/components/ChatView.logic.ts index 9729e2631..1e7bfc150 100644 --- a/apps/web/src/components/ChatView.logic.ts +++ b/apps/web/src/components/ChatView.logic.ts @@ -868,6 +868,24 @@ export interface LocalDispatchSnapshot { latestTurnCompletedAt: string | null; sessionStatus: NonNullable["status"] | null; sessionUpdatedAt: string | null; + latestTurnStartFailureId: string | null; +} + +export function latestTurnStartFailureId( + activeThread: Thread | undefined, + latestUserMessageId: ChatMessage["id"] | null, +): string | null { + if (latestUserMessageId === null) return null; + return ( + activeThread?.activities.findLast((activity) => { + if (activity.kind !== "provider.turn.start.failed") return false; + const payload = + typeof activity.payload === "object" && activity.payload !== null + ? (activity.payload as { readonly requestId?: unknown }) + : null; + return payload?.requestId === latestUserMessageId; + })?.id ?? null + ); } export function createLocalDispatchSnapshot( @@ -891,6 +909,7 @@ export function createLocalDispatchSnapshot( latestTurnCompletedAt: latestTurn?.completedAt ?? null, sessionStatus: session?.status ?? null, sessionUpdatedAt: session?.updatedAt ?? null, + latestTurnStartFailureId: latestTurnStartFailureId(activeThread, latestUserMessage?.id ?? null), }; } @@ -902,6 +921,7 @@ export function hasServerAcknowledgedLocalDispatch(input: { session: Thread["session"] | null; hasPendingApproval: boolean; hasPendingUserInput: boolean; + latestTurnStartFailureId?: string | null; threadError: string | null | undefined; }): boolean { if (!input.localDispatch) { @@ -910,6 +930,13 @@ export function hasServerAcknowledgedLocalDispatch(input: { if (input.hasPendingApproval || input.hasPendingUserInput || Boolean(input.threadError)) { return true; } + if ( + input.latestTurnStartFailureId !== undefined && + input.latestTurnStartFailureId !== null && + input.latestTurnStartFailureId !== input.localDispatch.latestTurnStartFailureId + ) { + return true; + } if (input.phase === "connecting") { return false; } diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 2f09b4874..8cf2c05a9 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -399,7 +399,7 @@ import { } from "@t3tools/shared/usageLimits"; import { ComposerSurface } from "./chat/ComposerSurface"; import { - hasAvailableClaudeCompactionProvider, + hasAvailableCompactionProvider, hasDismissedResumeCompaction, shouldOfferResumeCompaction, } from "./chat/ContextWindowMeter.logic"; @@ -425,6 +425,7 @@ import { deriveComposerSendState, dismissBranchMismatchForSession, hasEnvironmentReconnectWarningGraceElapsed, + latestTurnStartFailureId, scheduleEnvironmentReconnectWarning, hasServerAcknowledgedLocalDispatch, isBranchMismatchDismissedForSession, @@ -701,6 +702,11 @@ function formatOutgoingPrompt(params: { const SCRIPT_TERMINAL_COLS = 120; const SCRIPT_TERMINAL_ROWS = 30; +function isCompactCommandMessage(message: ChatMessage): boolean { + const text = message.text.trim().toLowerCase(); + return message.role === "user" && text === "/compact" && !message.attachments?.length; +} + type ChatViewProps = | { environmentId: EnvironmentId; @@ -744,6 +750,10 @@ function useLocalDispatchState(input: { (message) => message.role === "user", ); const latestUserMessageId = latestUserMessage?.id ?? null; + const currentTurnStartFailureId = + localDispatch === null + ? null + : latestTurnStartFailureId(input.activeThread, latestUserMessageId); const resetLocalDispatch = useCallback(() => { setLocalDispatch(null); @@ -759,6 +769,7 @@ function useLocalDispatchState(input: { session: input.activeThread?.session ?? null, hasPendingApproval: input.activePendingApproval !== null, hasPendingUserInput: input.activePendingUserInput !== null, + latestTurnStartFailureId: currentTurnStartFailureId, threadError: input.threadError, }), [ @@ -769,6 +780,7 @@ function useLocalDispatchState(input: { input.phase, input.threadError, latestUserMessageId, + currentTurnStartFailureId, localDispatch, ], ); @@ -2911,7 +2923,33 @@ export default function ChatView(props: ChatViewProps) { activePendingUserInput: activePendingUserInput?.requestId ?? null, threadError, }); - const isWorking = phase === "running" || isSendBusy || isConnecting || isRevertingCheckpoint; + const optimisticCompactionMessage = optimisticUserMessages.at(-1); + const pendingCompactionMessage = + isSendBusy && + optimisticCompactionMessage !== undefined && + isCompactCommandMessage(optimisticCompactionMessage) + ? optimisticCompactionMessage + : activeThread?.messages.findLast(isCompactCommandMessage); + const compactRequestIsActive = + pendingCompactionMessage !== undefined && + (pendingCompactionMessage.createdAt > + (activeLatestTurn?.requestedAt ?? pendingCompactionMessage.createdAt) || + (activeLatestTurn?.state === "running" && + pendingCompactionMessage.createdAt === activeLatestTurn.requestedAt)); + const compactionSettled = + pendingCompactionMessage !== undefined && + (latestTurnStartFailureId(activeThread, pendingCompactionMessage.id) !== null || + activeThread?.activities.some((activity) => { + if (activity.kind !== "context-compaction") return false; + const payload = activity.payload as { readonly requestId?: unknown } | null | undefined; + return payload?.requestId === pendingCompactionMessage.id; + })); + const isCompacting = + (isSendBusy || phase === "connecting" || phase === "running") && + compactRequestIsActive && + !compactionSettled; + const isWorking = + phase === "running" || isSendBusy || isConnecting || isRevertingCheckpoint || isCompacting; const activeWorkStartedAt = deriveActiveWorkStartedAt( activeLatestTurn, activeThread?.session ?? null, @@ -3306,13 +3344,14 @@ export default function ChatView(props: ChatViewProps) { activeThread?.modelSelection.instanceId ?? activeProject?.defaultModelSelection?.instanceId ?? null; - const compactionProviderAvailable = useMemo( + const manualCompactionProviderAvailable = useMemo( () => - hasAvailableClaudeCompactionProvider({ + hasAvailableCompactionProvider({ providers: applyProviderInstanceSettings( deriveProviderInstanceEntries(providerStatuses), settings, ), + driverKind: selectedProvider, instanceId: activeProviderInstanceId, lockedInstanceId: lockedProvider ? (activeThread?.session?.providerInstanceId ?? @@ -3326,6 +3365,7 @@ export default function ChatView(props: ChatViewProps) { activeThread?.session?.providerInstanceId, lockedProvider, providerStatuses, + selectedProvider, settings, ], ); @@ -5814,12 +5854,16 @@ export default function ChatView(props: ChatViewProps) { activeThread && activeContextWindow ? `${activeThread.id}:${activeContextWindow.updatedAt}` : null; - const compactDisabled = + const activeThreadHasCompactableConversation = + activeThread?.messages.some( + (message) => message.role === "user" && !isCompactCommandMessage(message), + ) ?? false; + const compactThreadUnavailable = !activeThread || + !activeThreadHasCompactableConversation || !activeProject || !isServerThread || - selectedProvider !== "claudeAgent" || - !compactionProviderAvailable || + !manualCompactionProviderAvailable || isWorking || threadDetailLoading || isPreparingWorktree || @@ -5827,15 +5871,15 @@ export default function ChatView(props: ChatViewProps) { feedbackUploading || pendingApprovals.length > 0 || pendingUserInputs.length > 0 || - showPlanFollowUpPrompt || - composerHasUnsentContent; + showPlanFollowUpPrompt; + const compactDisabled = compactThreadUnavailable || composerHasUnsentContent; const compactDisabledReason = compactDisabled ? composerHasUnsentContent ? "Send or clear your draft before compacting" : !activeProject ? "Choose a project before compacting" - : !compactionProviderAvailable - ? "Enable a Claude provider before compacting" + : !manualCompactionProviderAvailable + ? "Compaction is unavailable for this provider" : "Compacting is unavailable right now" : null; const resumeCompactionBannerItem = useMemo(() => { @@ -8958,6 +9002,7 @@ export default function ChatView(props: ChatViewProps) { workingStepLabel={workingStepLabel} activeTurnInProgress={isWorking || !latestTurnSettled} isPreparingWorktree={isPreparingWorktree} + isCompacting={isCompacting} activeTurnStartedAt={activeWorkStartedAt} listRef={legendListRef} timelineEntries={timelineEntries} @@ -9151,6 +9196,7 @@ export default function ChatView(props: ChatViewProps) { activeProject?.defaultModelSelection } activeThreadModelSelection={activeThread?.modelSelection} + compactThreadUnavailable={compactThreadUnavailable} compactDisabled={compactDisabled} activeThreadActivities={activeThread?.activities} quickQuestionAvailable={quickQuestionAvailable} diff --git a/apps/web/src/components/chat/ChatComposer.tsx b/apps/web/src/components/chat/ChatComposer.tsx index 6480abe35..cfefc8e9e 100644 --- a/apps/web/src/components/chat/ChatComposer.tsx +++ b/apps/web/src/components/chat/ChatComposer.tsx @@ -233,7 +233,10 @@ import { ContextWindowMeter } from "./ContextWindowMeter"; import { SessionGoalControl } from "./SessionGoalControl"; import { SessionHarnessControl } from "./SessionHarnessControl"; import { SessionInputQueueControl } from "./SessionInputQueueControl"; -import { resolveContextWindowModelDisplayName } from "./ContextWindowMeter.logic"; +import { + providerSupportsManualCompaction, + resolveContextWindowModelDisplayName, +} from "./ContextWindowMeter.logic"; import { attachVideoThumbnail, buildExpandedImagePreview, @@ -1365,6 +1368,7 @@ export interface ChatComposerProps { activeThreadActivities: Thread["activities"] | undefined; // Context window + compactThreadUnavailable: boolean; compactDisabled: boolean; // Ephemeral session side question @@ -1514,6 +1518,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) activeProjectDefaultModelSelection, activeThreadModelSelection, activeThreadActivities, + compactThreadUnavailable, compactDisabled, quickQuestionAvailable, quickQuestionIdentity, @@ -1884,6 +1889,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) () => selectedProviderEntry?.snapshot ?? null, [selectedProviderEntry], ); + const compactCommandAvailable = providerSupportsManualCompaction(selectedProviderEntry); const activeSessionProviderStatus = useMemo(() => { const instanceId = activeThread?.session?.providerInstanceId; if (instanceId === undefined) return null; @@ -2709,7 +2715,6 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) prompt, ], ); - // ------------------------------------------------------------------ // Derived: composer trigger / menu // ------------------------------------------------------------------ @@ -2721,6 +2726,17 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) cwd: isPathTrigger ? gitCwd : null, query: isPathTrigger ? pathTriggerQuery : null, }); + const compactSlashCommandAvailable = + composerTrigger?.kind === "slash-command" && + prompt.slice(0, composerTrigger.rangeStart).trim() === "" && + !compactThreadUnavailable && + prompt.slice(composerTrigger.rangeEnd).trim() === "" && + composerImages.length + composerFiles.length === 0 && + composerDraft.persistedAttachments.length === 0 && + composerTerminalContexts.length === 0 && + composerElementContexts.length === 0 && + composerPreviewAnnotations.length === 0 && + composerReviewComments.length === 0; const composerMenuItems = useMemo(() => { if (!composerTrigger) return []; @@ -2789,8 +2805,11 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) skill.description ?? (skill.scope ? `${skill.scope} skill` : ""), })); + const visibleProviderSlashCommandItems = providerSlashCommandItems.filter( + (item) => item.command.name !== "compact" || compactSlashCommandAvailable, + ); const slashCommandItems = slashCommandItemsForPromptPosition( - [...builtInSlashCommandItems, ...providerSlashCommandItems, ...skillItems], + [...builtInSlashCommandItems, ...visibleProviderSlashCommandItems, ...skillItems], composerTrigger.rangeStart === 0, ); return searchSlashCommandItems(slashCommandItems, query); @@ -2810,6 +2829,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) } return []; }, [ + compactSlashCommandAvailable, composerProviderControls.showInteractionModeToggle, composerTrigger, providerSlashCommands, @@ -3698,6 +3718,11 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) if ( compactDisabled || providerTurnUnavailable || + noProviderAvailable || + // The composer can hold a draft for a provider other than the running + // session's, so re-check the selected entry rather than trusting the + // thread-level gate alone. + !compactCommandAvailable || composerSendState.hasSendableContent || activePendingApproval !== null || pendingUserInputs.length > 0 || @@ -3734,11 +3759,13 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) }, [ activePendingApproval, activeThreadId, + compactCommandAvailable, compactDisabled, composerDraftTarget, composerSendState.hasSendableContent, isConnecting, isSendBusy, + noProviderAvailable, providerTurnUnavailable, pendingUserInputs.length, phase, diff --git a/apps/web/src/components/chat/ContextWindowMeter.logic.test.ts b/apps/web/src/components/chat/ContextWindowMeter.logic.test.ts index 032076c0b..b4c4d2217 100644 --- a/apps/web/src/components/chat/ContextWindowMeter.logic.test.ts +++ b/apps/web/src/components/chat/ContextWindowMeter.logic.test.ts @@ -3,7 +3,7 @@ import { describe, expect, it } from "vite-plus/test"; import { deriveProviderInstanceEntries } from "../../providerInstances"; import { formatContextWindowCompactionMessage, - hasAvailableClaudeCompactionProvider, + hasAvailableCompactionProvider, hasDismissedResumeCompaction, resolveContextWindowModelDisplayName, shouldOfferResumeCompaction, @@ -25,12 +25,12 @@ function claudeProvider(input: { auth: { status: "authenticated" }, checkedAt: "2026-08-24T12:00:00.000Z", models: [], - slashCommands: [], + slashCommands: [{ name: "compact", description: "" }], skills: [], }; } -describe("hasAvailableClaudeCompactionProvider", () => { +describe("hasAvailableCompactionProvider", () => { const originalInstanceId = ProviderInstanceId.make("claude_original"); it("rejects a fallback in a different locked continuation group", () => { @@ -47,8 +47,9 @@ describe("hasAvailableClaudeCompactionProvider", () => { ]); expect( - hasAvailableClaudeCompactionProvider({ + hasAvailableCompactionProvider({ providers, + driverKind: ProviderDriverKind.make("claudeAgent"), instanceId: originalInstanceId, lockedInstanceId: originalInstanceId, }), @@ -69,8 +70,9 @@ describe("hasAvailableClaudeCompactionProvider", () => { ]); expect( - hasAvailableClaudeCompactionProvider({ + hasAvailableCompactionProvider({ providers, + driverKind: ProviderDriverKind.make("claudeAgent"), instanceId: originalInstanceId, lockedInstanceId: originalInstanceId, }), diff --git a/apps/web/src/components/chat/ContextWindowMeter.logic.ts b/apps/web/src/components/chat/ContextWindowMeter.logic.ts index 8664c593a..6ff2b6e0a 100644 --- a/apps/web/src/components/chat/ContextWindowMeter.logic.ts +++ b/apps/web/src/components/chat/ContextWindowMeter.logic.ts @@ -1,4 +1,4 @@ -import type { ModelSelection, ProviderInstanceId } from "@t3tools/contracts"; +import type { ModelSelection, ProviderDriverKind, ProviderInstanceId } from "@t3tools/contracts"; import { CLAUDE_RESUME_COMPACTION_NEVER_ANSWER, isClaudeResumeCompactionQuestion, @@ -12,27 +12,33 @@ import { getTriggerDisplayModelName, type ModelEsque } from "./providerIconUtils const CLAUDE_RESUME_COMPACTION_MINUTES = 70; const CLAUDE_RESUME_COMPACTION_TOKENS = 100_000; -export function hasAvailableClaudeCompactionProvider(input: { +export function providerSupportsManualCompaction( + provider: ProviderInstanceEntry | null | undefined, +): boolean { + return provider?.snapshot.slashCommands.some((command) => command.name === "compact") ?? false; +} + +export function hasAvailableCompactionProvider(input: { readonly providers: ReadonlyArray; + readonly driverKind: ProviderDriverKind; readonly instanceId: ProviderInstanceId | null; readonly lockedInstanceId: ProviderInstanceId | null; }): boolean { - const claudeProviders = input.providers.filter( - (provider) => provider.driverKind === "claudeAgent", + const driverProviders = input.providers.filter( + (provider) => provider.driverKind === input.driverKind, ); const lockedContinuationGroupKey = input.lockedInstanceId - ? claudeProviders.find((provider) => provider.instanceId === input.lockedInstanceId) + ? driverProviders.find((provider) => provider.instanceId === input.lockedInstanceId) ?.continuationGroupKey : undefined; const compatibleProviders = lockedContinuationGroupKey - ? claudeProviders.filter( + ? driverProviders.filter( (provider) => provider.continuationGroupKey === lockedContinuationGroupKey, ) - : claudeProviders; + : driverProviders; - return ( - resolveSelectableProviderInstanceEntry(compatibleProviders, input.instanceId ?? undefined) !== - undefined + return providerSupportsManualCompaction( + resolveSelectableProviderInstanceEntry(compatibleProviders, input.instanceId ?? undefined), ); } diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts index 2b2a2c10b..00fc21bca 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts @@ -886,6 +886,38 @@ describe("resolveAssistantMessageCopyState", () => { }); describe("deriveMessagesTimelineRows", () => { + it("keeps context compaction visible outside folded work", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "compaction-entry", + kind: "work", + createdAt: "2026-01-01T00:00:00Z", + entry: { + id: "compaction", + createdAt: "2026-01-01T00:00:00Z", + label: "Compacted context 899K → 19K tokens", + tone: "info", + sourceActivityKind: "context-compaction", + }, + }, + ], + isWorking: false, + activeTurnStartedAt: null, + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows).toEqual([ + { + kind: "context-compaction", + id: "compaction-entry", + createdAt: "2026-01-01T00:00:00Z", + label: "Compacted context 899K → 19K tokens", + }, + ]); + }); + it("only enables assistant copy for the terminal assistant message in a turn", () => { const rows = deriveMessagesTimelineRows({ timelineEntries: [ diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.ts b/apps/web/src/components/chat/MessagesTimeline.logic.ts index 37c9dd915..15783c976 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.ts @@ -306,6 +306,12 @@ export type MessagesTimelineRow = label: string; expanded: boolean; } + | { + kind: "context-compaction"; + id: string; + createdAt: string; + label: string; + } | { kind: "message"; id: string; @@ -802,6 +808,16 @@ function deriveTurnFolds(input: { if (hiddenEntryIds.size === 0) { continue; } + // A lone compaction row stays visible on its own; it only folds away as + // part of a turn that already folds other work. + const hidesNonCompactionWork = group.entries.some( + (entry) => + hiddenEntryIds.has(entry.id) && + !(entry.kind === "work" && entry.entry.sourceActivityKind === "context-compaction"), + ); + if (!hidesNonCompactionWork) { + continue; + } const firstEntry = group.entries[0]; const firstHiddenEntry = group.entries.find((entry) => hiddenEntryIds.has(entry.id)); @@ -992,6 +1008,7 @@ export function deriveMessagesTimelineRows(input: { !entryBelongsToActiveTurn(entry, index) || entry.kind !== "work" || entry.entry.agentSpawn !== undefined || + entry.entry.sourceActivityKind === "context-compaction" || entry.entry.tone === "error" ) { break; @@ -1101,6 +1118,19 @@ export function deriveMessagesTimelineRows(input: { continue; } + if ( + timelineEntry.kind === "work" && + timelineEntry.entry.sourceActivityKind === "context-compaction" + ) { + nextRows.push({ + kind: "context-compaction", + id: timelineEntry.id, + createdAt: timelineEntry.createdAt, + label: timelineEntry.entry.label, + }); + continue; + } + if (timelineEntry.kind === "work") { if (timelineEntry.entry.agentSpawn !== undefined || timelineEntry.entry.tone === "error") { nextRows.push({ @@ -1122,6 +1152,7 @@ export function deriveMessagesTimelineRows(input: { nextEntry.kind !== "work" || workLogEntryIsMissingResponse(nextEntry.entry) || nextEntry.entry.agentSpawn !== undefined || + nextEntry.entry.sourceActivityKind === "context-compaction" || nextEntry.entry.tone === "error" || activeWorkEntryIds.has(nextEntry.id) || collapsedEntryIds.has(nextEntry.id) || @@ -1385,6 +1416,11 @@ function isRowUnchanged(a: MessagesTimelineRow, b: MessagesTimelineRow): boolean return a.createdAt === bf.createdAt && a.label === bf.label && a.expanded === bf.expanded; } + case "context-compaction": { + const bc = b as typeof a; + return a.createdAt === bc.createdAt && a.label === bc.label; + } + case "proposed-plan": return a.proposedPlan === (b as typeof a).proposedPlan; diff --git a/apps/web/src/components/chat/MessagesTimeline.test.tsx b/apps/web/src/components/chat/MessagesTimeline.test.tsx index d2b4ada49..d42722230 100644 --- a/apps/web/src/components/chat/MessagesTimeline.test.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.test.tsx @@ -1245,7 +1245,7 @@ describe("MessagesTimeline", () => { entry: { id: "work-1", createdAt: "2026-03-17T19:12:28.000Z", - label: "Context compacted", + label: "Compacted context 899K → 19K tokens", tone: "info", }, }, @@ -1253,7 +1253,7 @@ describe("MessagesTimeline", () => { />, ); - expect(markup).toContain("Context compacted"); + expect(markup).toContain("Compacted context 899K → 19K tokens"); }); it("reports a coordinator-only pending workflow as queued", () => { diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index a855e6327..5b6f4e7fb 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -92,6 +92,7 @@ import { GlobeIcon, HammerIcon, MessageCircleIcon, + Minimize2Icon, MinusIcon, MousePointerClickIcon, PaintbrushIcon, @@ -222,6 +223,7 @@ interface TimelineRowSharedState { interface TimelineRowActivityState { isWorking: boolean; isPreparingWorktree: boolean; + isCompacting: boolean; isRevertingCheckpoint: boolean; activeTurnInProgress: boolean; latestTurnId: TurnId | null; @@ -311,6 +313,7 @@ interface MessagesTimelineProps { workingStepLabel?: string | null; activeTurnInProgress: boolean; isPreparingWorktree?: boolean; + isCompacting?: boolean; activeTurnStartedAt: string | null; listRef: React.RefObject; timelineEntries: ReturnType; @@ -371,6 +374,7 @@ export const MessagesTimeline = memo(function MessagesTimeline({ workingStepLabel = null, activeTurnInProgress, isPreparingWorktree = false, + isCompacting = false, activeTurnStartedAt, agentPanelModel = EMPTY_AGENT_PANEL_MODEL, onOpenAgents = NOOP_OPEN_AGENTS, @@ -787,6 +791,7 @@ export const MessagesTimeline = memo(function MessagesTimeline({ () => ({ isWorking, isPreparingWorktree, + isCompacting, isRevertingCheckpoint, activeTurnInProgress, latestTurnId: latestTurn?.turnId ?? null, @@ -794,6 +799,7 @@ export const MessagesTimeline = memo(function MessagesTimeline({ }), [ activeTurnInProgress, + isCompacting, isRevertingCheckpoint, isWorking, isPreparingWorktree, @@ -1244,6 +1250,7 @@ const TimelineRowContent = memo(function TimelineRowContent({ row }: { row: Time {row.kind === "work-live" ? : null} {row.kind === "work-toggle" ? : null} {row.kind === "turn-fold" ? : null} + {row.kind === "context-compaction" ? : null} {row.kind === "message" && row.message.role === "user" ? : null} {row.kind === "message" && row.message.role === "assistant" ? ( @@ -1256,6 +1263,27 @@ const TimelineRowContent = memo(function TimelineRowContent({ row }: { row: Time ); }); +function ContextCompactionTimelineRow({ + row, +}: { + row: Extract; +}) { + return ( +
+ + + + +
+ ); +} + function UserTimelineRow({ row }: { row: Extract }) { const ctx = use(TimelineRowCtx); const resources = useMemo( @@ -1678,16 +1706,18 @@ function ProposedPlanTimelineRow({ } function WorkingTimelineRow({ row }: { row: Extract }) { - const { workingStepLabel, isPreparingWorktree } = use(TimelineRowActivityCtx); + const { workingStepLabel, isCompacting, isPreparingWorktree } = use(TimelineRowActivityCtx); return (
{isPreparingWorktree ? ( "Setting up worktree…" + ) : isCompacting ? ( + ) : row.createdAt ? ( <> Working for @@ -1707,15 +1737,26 @@ function WorkingTimelineRow({ row }: { row: Extract - {isPreparingWorktree ? null : } + {isPreparingWorktree || isCompacting ? null : ( + + )}
); } +function CompactingLabel() { + return ( + + + ); +} + // --------------------------------------------------------------------------- // Self-ticking labels — update their own text nodes so elapsed-time display // does not create a React commit every second while a response is streaming. diff --git a/docs/user/composer.md b/docs/user/composer.md index f12e533d3..fb370e926 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -262,6 +262,8 @@ project or worktree. Web, desktop, and mobile request this list from the connect so remote projects use their remote skills. Discovery can take a moment after you switch providers or workspaces. If discovery fails, later typing retries after a short cooldown. +In a thread with prior conversation context, send `/compact` to reduce context usage. Web and desktop also offer this action from the context meter, and the work log records token counts when the provider reports them. + By default, the `/` menu includes skills. To keep this menu command-only, turn off **Show skills in slash menu** in **Settings → General** in the web or desktop app. Skill results use the `/skill:Skill Name` label and add the same `$name` skill token to your message. The original skill diff --git a/packages/contracts/src/providerRuntime.ts b/packages/contracts/src/providerRuntime.ts index 5551cf878..f3ccc7018 100644 --- a/packages/contracts/src/providerRuntime.ts +++ b/packages/contracts/src/providerRuntime.ts @@ -341,6 +341,8 @@ export type ThreadStartedPayload = typeof ThreadStartedPayload.Type; const ThreadStateChangedPayload = Schema.Struct({ state: RuntimeThreadState, + beforeTokens: Schema.optional(NonNegativeInt), + afterTokens: Schema.optional(NonNegativeInt), detail: Schema.optional(Schema.Unknown), }); export type ThreadStateChangedPayload = typeof ThreadStateChangedPayload.Type;