From 721b974621e9580275749f6d1f1d6ee4d86c4ecd Mon Sep 17 00:00:00 2001 From: harrydawson Date: Mon, 8 Jun 2026 18:13:25 +1000 Subject: [PATCH 1/2] =?UTF-8?q?feat:=20/compact=20=E2=80=94=20provider=20c?= =?UTF-8?q?ompaction=20+=20summary=20banner=20(#87)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add thread.context.compact + thread.context.summarize command schemas - Add thread.context-compacted + thread.context-summarized event schemas - Add optional summary field to ContextTrimPoint and ThreadContextTrimCommand - Decider: validate and emit events for compact and summarize commands - Decider: include summary in trim point when provided via thread.context.trim - ProviderAdapter: add compactThread method to interface - CodexSessionRuntime: implement compactThread via thread/compact/start RPC - Claude/OpenCode/Cursor adapters: return unsupported error for compactThread - ProviderService: delegate compactThread to adapter - ProviderCommandReactor: process thread.context.compact — call provider, emit summarize + trim - Composer: parse /compact, register in autocomplete, dispatch in ChatView - ContextSummaryBanner component: blue-tinted 'Context compacted' callout - MessagesTimeline: render summary banner above trim divider when summary present - Tests: contracts schemas, decider validation, composer parsing --- .../Layers/ProviderCommandReactor.ts | 68 +++- .../orchestration/decider.contextTrim.test.ts | 309 ++++++++++++++++++ apps/server/src/orchestration/decider.ts | 59 ++++ .../src/provider/Layers/ClaudeAdapter.ts | 10 + .../src/provider/Layers/CodexAdapter.test.ts | 7 + .../src/provider/Layers/CodexAdapter.ts | 11 + .../provider/Layers/CodexSessionRuntime.ts | 58 ++++ .../src/provider/Layers/CursorAdapter.ts | 10 + .../src/provider/Layers/OpenCodeAdapter.ts | 10 + .../Layers/ProviderAdapterRegistry.test.ts | 4 + .../provider/Layers/ProviderService.test.ts | 8 + .../src/provider/Layers/ProviderService.ts | 37 +++ .../src/provider/Services/ProviderAdapter.ts | 10 + .../src/provider/Services/ProviderService.ts | 8 + apps/web/src/components/ChatView.tsx | 15 + apps/web/src/components/chat/ChatComposer.tsx | 8 + .../components/chat/ContextSummaryBanner.tsx | 24 ++ .../src/components/chat/MessagesTimeline.tsx | 14 +- apps/web/src/composer-logic.test.ts | 12 + apps/web/src/composer-logic.ts | 5 +- packages/contracts/src/orchestration.test.ts | 188 +++++++++++ packages/contracts/src/orchestration.ts | 48 +++ 22 files changed, 915 insertions(+), 8 deletions(-) create mode 100644 apps/web/src/components/chat/ContextSummaryBanner.tsx diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index f63b873bc3de..680e7ce2de94 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -16,6 +16,7 @@ import { isTemporaryWorktreeBranch, WORKTREE_BRANCH_PREFIX } from "@t3tools/shar import * as Cache from "effect/Cache"; import * as Cause from "effect/Cause"; import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Equal from "effect/Equal"; @@ -52,7 +53,8 @@ type ProviderIntentEvent = Extract< | "thread.turn-interrupt-requested" | "thread.approval-response-requested" | "thread.user-input-response-requested" - | "thread.session-stop-requested"; + | "thread.session-stop-requested" + | "thread.context-compacted"; } >; @@ -944,6 +946,64 @@ const make = Effect.gen(function* () { }); }); + const processContextCompacted = Effect.fn("processContextCompacted")(function* ( + event: Extract, + ) { + const threadId = event.payload.threadId; + const thread = yield* resolveThread(threadId); + if (!thread) { + return; + } + + const nowIso = yield* Effect.map(DateTime.now, DateTime.formatIso); + + // Attempt provider compaction + const compactResult = yield* providerService.compactThread({ threadId }).pipe( + Effect.map((result) => ({ _tag: "success" as const, ...result })), + Effect.catchAll((cause) => + Effect.gen(function* () { + yield* Effect.logWarning("provider compaction failed, proceeding with trim-only", { + threadId: String(threadId), + cause: Cause.pretty(cause), + }); + return { _tag: "failure" as const }; + }), + ), + ); + + if (compactResult._tag === "success") { + // Dispatch summarize with the compaction summary + const summarizeCommandId = yield* serverCommandId("compact-summarize"); + yield* orchestrationEngine.dispatch({ + type: "thread.context.summarize", + commandId: summarizeCommandId, + threadId, + summary: compactResult.summary, + compactDurationMs: compactResult.durationMs, + createdAt: nowIso, + }); + + // Dispatch trim with the summary attached + const trimCommandId = yield* serverCommandId("compact-trim"); + yield* orchestrationEngine.dispatch({ + type: "thread.context.trim", + commandId: trimCommandId, + threadId, + summary: compactResult.summary, + createdAt: nowIso, + }); + } else { + // Dispatch trim without summary (fallback) + const trimCommandId = yield* serverCommandId("compact-trim-fallback"); + yield* orchestrationEngine.dispatch({ + type: "thread.context.trim", + commandId: trimCommandId, + threadId, + createdAt: nowIso, + }); + } + }); + const processDomainEvent = Effect.fn("processDomainEvent")(function* ( event: ProviderIntentEvent, ) { @@ -984,6 +1044,9 @@ const make = Effect.gen(function* () { case "thread.session-stop-requested": yield* processSessionStopRequested(event); return; + case "thread.context-compacted": + yield* processContextCompacted(event); + return; } }); @@ -1010,7 +1073,8 @@ const make = Effect.gen(function* () { event.type === "thread.turn-interrupt-requested" || event.type === "thread.approval-response-requested" || event.type === "thread.user-input-response-requested" || - event.type === "thread.session-stop-requested" + event.type === "thread.session-stop-requested" || + event.type === "thread.context-compacted" ) { return yield* worker.enqueue(event); } diff --git a/apps/server/src/orchestration/decider.contextTrim.test.ts b/apps/server/src/orchestration/decider.contextTrim.test.ts index 63c36aec44d5..b9bf31fa7180 100644 --- a/apps/server/src/orchestration/decider.contextTrim.test.ts +++ b/apps/server/src/orchestration/decider.contextTrim.test.ts @@ -716,3 +716,312 @@ describe("thread.context.trim decider", () => { ]); }); }); + +describe("thread.context.compact decider", () => { + it("emits thread.context-compacted for a valid thread", async () => { + const now = "2026-01-01T00:00:00.000Z"; + const threadId = asThreadId("thread-compact"); + let model = createEmptyReadModel(now); + + model = await Effect.runPromise( + projectEvent(model, { + sequence: 1, + eventId: asEventId("evt-project-compact"), + aggregateKind: "project", + aggregateId: asProjectId("project-compact"), + type: "project.created", + occurredAt: now, + commandId: asCommandId("cmd-project-compact"), + causationEventId: null, + correlationId: asCommandId("cmd-project-compact"), + metadata: {}, + payload: { + projectId: asProjectId("project-compact"), + title: "Compact Project", + workspaceRoot: "/tmp/compact-project", + defaultModelSelection: null, + scripts: [], + createdAt: now, + updatedAt: now, + }, + }), + ); + + model = await Effect.runPromise( + projectEvent(model, { + sequence: 2, + eventId: asEventId("evt-thread-compact"), + aggregateKind: "thread", + aggregateId: threadId, + type: "thread.created", + occurredAt: now, + commandId: asCommandId("cmd-thread-compact"), + causationEventId: null, + correlationId: asCommandId("cmd-thread-compact"), + metadata: {}, + payload: { + threadId, + projectId: asProjectId("project-compact"), + title: "Compact Thread", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: DEFAULT_RUNTIME_MODE, + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + }, + }), + ); + + const result = await runDecide({ + command: { + type: "thread.context.compact", + commandId: asCommandId("cmd-compact"), + threadId, + createdAt: "2026-01-02T00:00:00.000Z", + } as Extract, + readModel: model, + }); + + const events = Array.isArray(result) ? result : [result]; + expect(events.map((e) => e.type)).toEqual(["thread.context-compacted"]); + }); + + it("rejects compact on non-existent thread", async () => { + const readModel = createEmptyReadModel("2026-01-01T00:00:00.000Z"); + + await expect( + runDecide({ + command: { + type: "thread.context.compact", + commandId: asCommandId("cmd-compact-unknown"), + threadId: asThreadId("thread-unknown"), + createdAt: "2026-01-02T00:00:00.000Z", + } as Extract, + readModel, + }), + ).rejects.toBeDefined(); + }); + + it("rejects compact when thread has an active turn", async () => { + const now = "2026-01-01T00:00:00.000Z"; + const threadId = asThreadId("thread-compact-active"); + let model = createEmptyReadModel(now); + + model = await Effect.runPromise( + projectEvent(model, { + sequence: 1, + eventId: asEventId("evt-project-compact-active"), + aggregateKind: "project", + aggregateId: asProjectId("project-compact-active"), + type: "project.created", + occurredAt: now, + commandId: asCommandId("cmd-project-compact-active"), + causationEventId: null, + correlationId: asCommandId("cmd-project-compact-active"), + metadata: {}, + payload: { + projectId: asProjectId("project-compact-active"), + title: "Compact Active", + workspaceRoot: "/tmp/compact-active", + defaultModelSelection: null, + scripts: [], + createdAt: now, + updatedAt: now, + }, + }), + ); + + model = await Effect.runPromise( + projectEvent(model, { + sequence: 2, + eventId: asEventId("evt-thread-compact-active"), + aggregateKind: "thread", + aggregateId: threadId, + type: "thread.created", + occurredAt: now, + commandId: asCommandId("cmd-thread-compact-active"), + causationEventId: null, + correlationId: asCommandId("cmd-thread-compact-active"), + metadata: {}, + payload: { + threadId, + projectId: asProjectId("project-compact-active"), + title: "Active Thread", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: DEFAULT_RUNTIME_MODE, + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + }, + }), + ); + + model = await Effect.runPromise( + projectEvent(model, { + sequence: 3, + eventId: asEventId("evt-session-compact-active"), + aggregateKind: "thread", + aggregateId: threadId, + type: "thread.session-set", + occurredAt: "2026-01-01T00:05:00.000Z", + commandId: asCommandId("cmd-session-compact-active"), + causationEventId: null, + correlationId: asCommandId("cmd-session-compact-active"), + metadata: {}, + payload: { + threadId, + session: { + threadId, + status: "running" as const, + providerName: "codex", + runtimeMode: DEFAULT_RUNTIME_MODE, + activeTurnId: asTurnId("turn-active-compact"), + lastError: null, + updatedAt: "2026-01-01T00:05:00.000Z", + }, + }, + }), + ); + + await expect( + runDecide({ + command: { + type: "thread.context.compact", + commandId: asCommandId("cmd-compact-active"), + threadId, + createdAt: "2026-01-02T00:00:00.000Z", + } as Extract, + readModel: model, + }), + ).rejects.toBeDefined(); + }); +}); + +describe("thread.context.summarize decider", () => { + it("emits thread.context-summarized for a valid thread", async () => { + const now = "2026-01-01T00:00:00.000Z"; + const threadId = asThreadId("thread-summarize"); + let model = createEmptyReadModel(now); + + model = await Effect.runPromise( + projectEvent(model, { + sequence: 1, + eventId: asEventId("evt-project-summarize"), + aggregateKind: "project", + aggregateId: asProjectId("project-summarize"), + type: "project.created", + occurredAt: now, + commandId: asCommandId("cmd-project-summarize"), + causationEventId: null, + correlationId: asCommandId("cmd-project-summarize"), + metadata: {}, + payload: { + projectId: asProjectId("project-summarize"), + title: "Summarize Project", + workspaceRoot: "/tmp/summarize-project", + defaultModelSelection: null, + scripts: [], + createdAt: now, + updatedAt: now, + }, + }), + ); + + model = await Effect.runPromise( + projectEvent(model, { + sequence: 2, + eventId: asEventId("evt-thread-summarize"), + aggregateKind: "thread", + aggregateId: threadId, + type: "thread.created", + occurredAt: now, + commandId: asCommandId("cmd-thread-summarize"), + causationEventId: null, + correlationId: asCommandId("cmd-thread-summarize"), + metadata: {}, + payload: { + threadId, + projectId: asProjectId("project-summarize"), + title: "Summarize Thread", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: DEFAULT_RUNTIME_MODE, + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + }, + }), + ); + + const result = await runDecide({ + command: { + type: "thread.context.summarize", + commandId: asCommandId("cmd-summarize"), + threadId, + summary: "The user asked about compacting threads.", + compactDurationMs: 1200, + createdAt: "2026-01-02T00:00:00.000Z", + } as Extract, + readModel: model, + }); + + const events = Array.isArray(result) ? result : [result]; + expect(events.map((e) => e.type)).toEqual(["thread.context-summarized"]); + + const summarizeEvent = events.find((e) => e.type === "thread.context-summarized"); + expect(summarizeEvent?.payload.summary).toBe("The user asked about compacting threads."); + expect(summarizeEvent?.payload.compactDurationMs).toBe(1200); + }); + + it("rejects summarize on non-existent thread", async () => { + const readModel = createEmptyReadModel("2026-01-01T00:00:00.000Z"); + + await expect( + runDecide({ + command: { + type: "thread.context.summarize", + commandId: asCommandId("cmd-summarize-unknown"), + threadId: asThreadId("thread-unknown"), + summary: "Test", + createdAt: "2026-01-02T00:00:00.000Z", + } as Extract, + readModel, + }), + ).rejects.toBeDefined(); + }); +}); + +describe("thread.context.trim with summary", () => { + it("includes summary in trim point when provided", async () => { + const readModel = await seedThreadWithMessages(); + + const result = await runDecide({ + command: { + type: "thread.context.trim", + commandId: asCommandId("cmd-trim-summary"), + threadId: asThreadId("thread-trim"), + summary: "Compact summary text", + createdAt: "2026-01-02T00:00:00.000Z", + } as Extract, + readModel, + }); + + const events = Array.isArray(result) ? result : [result]; + const trimEvent = events.find((e) => e.type === "thread.trim-point-created"); + expect(trimEvent?.payload.trimPoint.summary).toBe("Compact summary text"); + }); +}); diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index f075ecf39361..b4432cd38993 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -1149,6 +1149,9 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" beforeEntryId, prunedMessageCount, prunedTurnIds, + ...("summary" in command && command.summary !== undefined + ? { summary: command.summary } + : {}), }; return [ @@ -1181,6 +1184,62 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" ]; } + case "thread.context.compact": { + const thread = yield* requireThread({ + readModel, + command, + threadId: command.threadId, + }); + yield* requireTurnNotActive({ thread, command }); + const occurredAt = yield* nowIso; + return [ + { + ...(yield* withEventBase({ + aggregateKind: "thread", + aggregateId: command.threadId, + occurredAt, + commandId: command.commandId, + })), + type: "thread.context-compacted", + payload: { + threadId: command.threadId, + summary: "", + compactDurationMs: undefined, + }, + }, + ]; + } + + case "thread.context.summarize": { + const thread = yield* requireThread({ + readModel, + command, + threadId: command.threadId, + }); + yield* requireTurnNotActive({ thread, command }); + const occurredAt = yield* nowIso; + const trimPointId = EventId.make(crypto.randomUUID()); + return [ + { + ...(yield* withEventBase({ + aggregateKind: "thread", + aggregateId: command.threadId, + occurredAt, + commandId: command.commandId, + })), + type: "thread.context-summarized", + payload: { + threadId: command.threadId, + trimPointId, + summary: command.summary, + ...("compactDurationMs" in command && command.compactDurationMs !== undefined + ? { compactDurationMs: command.compactDurationMs } + : {}), + }, + }, + ]; + } + case "thread.session.set": { yield* requireThread({ readModel, diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index 1a0f41621953..b291104d1493 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -3371,6 +3371,15 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }, ); + const compactThread: ClaudeAdapterShape["compactThread"] = (_threadId) => + Effect.fail( + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "thread/compact/start", + detail: "Compaction is not supported by this provider.", + }), + ); + const respondToRequest: ClaudeAdapterShape["respondToRequest"] = Effect.fn("respondToRequest")( function* (threadId, requestId, decision) { const context = yield* requireSession(threadId); @@ -3514,6 +3523,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( interruptTurn, readThread, rollbackThread, + compactThread, respondToRequest, respondToUserInput, stopSession, diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index e12fdb9bed53..eb3e1afb11e5 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -103,6 +103,11 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape { }), ); + public readonly compactThreadImpl = vi.fn( + (): Promise<{ summary: string; durationMs: number }> => + Promise.resolve({ summary: "", durationMs: 0 }), + ); + public readonly respondToRequestImpl = vi.fn( (_requestId: ApprovalRequestId, _decision: ProviderApprovalDecision): Promise => Promise.resolve(undefined), @@ -141,6 +146,8 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape { return Effect.promise(() => this.rollbackThreadImpl(numTurns)); } + compactThread = Effect.promise(() => this.compactThreadImpl()); + respondToRequest(requestId: ApprovalRequestId, decision: ProviderApprovalDecision) { return Effect.promise(() => this.respondToRequestImpl(requestId, decision)); } diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index 8ac4b95c4104..a60fac9105ab 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -1595,6 +1595,16 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( ); }; + const compactThread: CodexAdapterShape["compactThread"] = (threadId) => + requireSession(threadId).pipe( + Effect.flatMap((session) => session.runtime.compactThread), + Effect.mapError((cause) => + cause._tag === "ProviderAdapterSessionNotFoundError" + ? cause + : mapCodexRuntimeError(threadId, "thread/compact/start", cause), + ), + ); + const respondToRequest: CodexAdapterShape["respondToRequest"] = (threadId, requestId, decision) => requireSession(threadId).pipe( Effect.flatMap((session) => session.runtime.respondToRequest(requestId, decision)), @@ -1683,6 +1693,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( interruptTurn, readThread, rollbackThread, + compactThread, respondToRequest, respondToUserInput, stopSession, diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index f9b9c6ab4fba..ab1c6c8c2673 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -138,6 +138,7 @@ export interface CodexSessionRuntimeShape { readonly rollbackThread: ( numTurns: number, ) => Effect.Effect; + readonly compactThread: Effect.Effect<{ summary: string; durationMs: number }, CodexSessionRuntimeError>; readonly respondToRequest: ( requestId: ApprovalRequestId, decision: ProviderApprovalDecision, @@ -1323,6 +1324,63 @@ export const makeCodexSessionRuntime = ( }); return parseThreadSnapshot(response); }), + compactThread: Effect.gen(function* () { + const startTime = Date.now(); + const providerThreadId = yield* readProviderThreadId; + const compactedDeferred = yield* Deferred.make< + { summary: string; durationMs: number }, + CodexSessionRuntimeError + >(); + + yield* client + .handleServerNotification("thread/compacted", (_payload) => + Effect.gen(function* () { + const snapshot = yield* client.request("thread/read", { + threadId: providerThreadId, + includeTurns: true, + }); + const parsed = parseThreadSnapshot(snapshot); + let summary = ""; + for (const turn of [...parsed.turns].reverse()) { + for (const item of turn.items) { + if (typeof (item as Record).type === "string") { + const itemType = (item as Record).type as string; + if ( + (itemType === "compaction" || itemType === "context_compaction") && + "encrypted_content" in (item as Record) + ) { + summary = String( + (item as Record).encrypted_content ?? "", + ); + break; + } + } + } + if (summary) break; + } + const durationMs = Date.now() - startTime; + yield* Deferred.succeed(compactedDeferred, { summary, durationMs }); + }), + ) + .pipe(Effect.forkScoped); + + yield* client.request("thread/compact/start", { + threadId: providerThreadId, + }); + + const result = yield* Effect.raceFirst( + Deferred.await(compactedDeferred), + Effect.clockWith((clock) => + clock.sleep("30 seconds").pipe( + Effect.andThen( + Effect.succeed({ summary: "", durationMs: Date.now() - startTime }), + ), + ), + ), + ); + + return result; + }), respondToRequest: (requestId, decision) => Effect.gen(function* () { const pending = (yield* Ref.get(pendingApprovalsRef)).get(requestId); diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index 1f6b88be773b..ab7c564b784c 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -1070,6 +1070,15 @@ export function makeCursorAdapter( return { threadId, turns: ctx.turns }; }); + const compactThread: CursorAdapterShape["compactThread"] = (_threadId) => + Effect.fail( + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "thread/compact/start", + detail: "Compaction is not supported by this provider.", + }), + ); + const stopSession: CursorAdapterShape["stopSession"] = (threadId) => withThreadLock( threadId, @@ -1111,6 +1120,7 @@ export function makeCursorAdapter( interruptTurn, readThread, rollbackThread, + compactThread, respondToRequest, respondToUserInput, stopSession, diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index e893f510aa02..d2039167949d 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -1408,6 +1408,15 @@ export function makeOpenCodeAdapter( }, ); + const compactThread: OpenCodeAdapterShape["compactThread"] = (_threadId) => + Effect.fail( + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "thread/compact/start", + detail: "Compaction is not supported by this provider.", + }), + ); + const stopAll: OpenCodeAdapterShape["stopAll"] = () => Effect.gen(function* () { const contexts = [...sessions.values()]; @@ -1439,6 +1448,7 @@ export function makeOpenCodeAdapter( hasSession, readThread, rollbackThread, + compactThread, stopAll, get streamEvents() { return Stream.fromQueue(runtimeEvents); diff --git a/apps/server/src/provider/Layers/ProviderAdapterRegistry.test.ts b/apps/server/src/provider/Layers/ProviderAdapterRegistry.test.ts index 33b36f6faf82..3980a271f3e4 100644 --- a/apps/server/src/provider/Layers/ProviderAdapterRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderAdapterRegistry.test.ts @@ -40,6 +40,7 @@ const fakeCodexAdapter: CodexAdapterShape = { hasSession: vi.fn(), readThread: vi.fn(), rollbackThread: vi.fn(), + compactThread: vi.fn(), stopAll: vi.fn(), streamEvents: Stream.empty, }; @@ -57,6 +58,7 @@ const fakeClaudeAdapter: ClaudeAdapterShape = { hasSession: vi.fn(), readThread: vi.fn(), rollbackThread: vi.fn(), + compactThread: vi.fn(), toggleMcpServerOnThread: vi.fn(), listMcpServersOnThread: vi.fn(), stopAll: vi.fn(), @@ -76,6 +78,7 @@ const fakeOpenCodeAdapter: OpenCodeAdapterShape = { hasSession: vi.fn(), readThread: vi.fn(), rollbackThread: vi.fn(), + compactThread: vi.fn(), stopAll: vi.fn(), streamEvents: Stream.empty, }; @@ -93,6 +96,7 @@ const fakeCursorAdapter: CursorAdapterShape = { hasSession: vi.fn(), readThread: vi.fn(), rollbackThread: vi.fn(), + compactThread: vi.fn(), stopAll: vi.fn(), streamEvents: Stream.empty, }; diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index ad4238e1264b..32ea14eaec3b 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -196,6 +196,13 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) { Effect.succeed({ threadId, turns: [] }), ); + const compactThread = vi.fn( + ( + _threadId: ThreadId, + ): Effect.Effect<{ summary: string; durationMs: number }, ProviderAdapterError> => + Effect.succeed({ summary: "", durationMs: 0 }), + ); + const stopAll = vi.fn( (): Effect.Effect => Effect.sync(() => { @@ -219,6 +226,7 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) { hasSession, readThread, rollbackThread, + compactThread, stopAll, get streamEvents() { return Stream.fromPubSub(runtimeEventPubSub); diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index 2bce1f483b74..6fc8420842ba 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -976,6 +976,42 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ); }); + const compactThread: ProviderServiceShape["compactThread"] = Effect.fn("compactThread")( + function* (rawInput) { + const input = yield* decodeInputOrValidationError({ + operation: "ProviderService.compactThread", + schema: Schema.Struct({ + threadId: ThreadId, + }), + payload: rawInput, + }); + + return yield* Effect.gen(function* () { + const routed = yield* resolveRoutableSession({ + threadId: input.threadId, + operation: "ProviderService.compactThread", + allowRecovery: false, + }); + yield* Effect.annotateCurrentSpan({ + "provider.operation": "compact-thread", + "provider.kind": routed.adapter.provider, + "provider.thread_id": input.threadId, + }); + return yield* routed.adapter.compactThread(routed.threadId); + }).pipe( + Effect.mapError((cause) => + cause._tag === "ProviderAdapterRequestError" || + cause._tag === "ProviderAdapterSessionNotFoundError" + ? cause + : new ProviderValidationError({ + operation: "ProviderService.compactThread", + issue: Cause.pretty(cause), + }), + ), + ); + }, + ); + const runStopAll = Effect.fn("runStopAll")(function* () { const threadIds = yield* directory.listThreadIds(); const currentAdapters = yield* getAdapterEntries; @@ -1043,6 +1079,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( getCapabilities, getInstanceInfo, rollbackConversation, + compactThread, // Each access creates a fresh PubSub subscription so that multiple // consumers (ProviderRuntimeIngestion, CheckpointReactor, etc.) each // independently receive all runtime events. diff --git a/apps/server/src/provider/Services/ProviderAdapter.ts b/apps/server/src/provider/Services/ProviderAdapter.ts index cbbf378bc3f6..361749111a0b 100644 --- a/apps/server/src/provider/Services/ProviderAdapter.ts +++ b/apps/server/src/provider/Services/ProviderAdapter.ts @@ -118,6 +118,16 @@ export interface ProviderAdapterShape { numTurns: number, ) => Effect.Effect; + /** + * Compact a provider thread context into a summary. + * Returns the compaction summary text and the duration the compaction took. + * Providers that don't support compaction should fail with an error + * (callers fall back to a trim-only flow). + */ + readonly compactThread: ( + threadId: ThreadId, + ) => Effect.Effect<{ summary: string; durationMs: number }, TError>; + /** * Stop all sessions owned by this adapter. */ diff --git a/apps/server/src/provider/Services/ProviderService.ts b/apps/server/src/provider/Services/ProviderService.ts index 4d4cb4fa01a7..1c974580feb5 100644 --- a/apps/server/src/provider/Services/ProviderService.ts +++ b/apps/server/src/provider/Services/ProviderService.ts @@ -105,6 +105,14 @@ export interface ProviderServiceShape { readonly numTurns: number; }) => Effect.Effect; + /** + * Compact provider thread context into a summary. + * Falls back gracefully if provider doesn't support compaction. + */ + readonly compactThread: (input: { + readonly threadId: ThreadId; + }) => Effect.Effect<{ summary: string; durationMs: number }, ProviderServiceError>; + /** * Canonical provider runtime event stream. * diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 6d0c2c450ca4..e6dd15bc1b27 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -2979,6 +2979,21 @@ export default function ChatView(props: ChatViewProps) { to: "/thread/$threadId", params: { threadId: nextThreadId }, }); + } else if (standaloneSlashCommand === "compact") { + api.orchestration.dispatchCommand({ + type: "thread.context.compact", + commandId: newCommandId(), + threadId: activeThread.id, + createdAt: new Date().toISOString(), + }).catch((error) => { + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Failed to compact context", + description: error instanceof Error ? error.message : "An unexpected error occurred.", + }), + ); + }); } else { handleInteractionModeChange(standaloneSlashCommand); } diff --git a/apps/web/src/components/chat/ChatComposer.tsx b/apps/web/src/components/chat/ChatComposer.tsx index e894e68af193..b9f8ea330927 100644 --- a/apps/web/src/components/chat/ChatComposer.tsx +++ b/apps/web/src/components/chat/ChatComposer.tsx @@ -886,6 +886,14 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) label: "/default", description: "Switch this thread back to normal build mode", }, + { + id: "slash:compact", + type: "slash-command", + command: "compact", + label: "/compact", + description: "Compact conversation context with provider-summarize", + disabled: props.activeTurnInProgress, + }, { id: "slash:clear", type: "slash-command", diff --git a/apps/web/src/components/chat/ContextSummaryBanner.tsx b/apps/web/src/components/chat/ContextSummaryBanner.tsx new file mode 100644 index 000000000000..2026c440334e --- /dev/null +++ b/apps/web/src/components/chat/ContextSummaryBanner.tsx @@ -0,0 +1,24 @@ +import { ChevronDownIcon } from "lucide-react"; +import { memo } from "react"; + +interface ContextSummaryBannerProps { + summary: string; +} + +export const ContextSummaryBanner = memo(function ContextSummaryBanner({ + summary, +}: ContextSummaryBannerProps) { + if (!summary) return null; + + return ( +
+
+ + + Context compacted + +
+

{summary}

+
+ ); +}); diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index b2d1bed4c4fb..7ca5e5af6ddc 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -47,6 +47,7 @@ import { ChangedFilesTree } from "./ChangedFilesTree"; import { DiffStatLabel, hasNonZeroStat } from "./DiffStatLabel"; import { MessageCopyButton } from "./MessageCopyButton"; import { ContextTrimPointDivider } from "./ContextTrimPointDivider"; +import { ContextSummaryBanner } from "./ContextSummaryBanner"; import { computeStableMessagesTimelineRows, MAX_VISIBLE_WORK_LOG_ENTRIES, @@ -359,11 +360,14 @@ const TimelineRowContent = memo(function TimelineRowContent({ row }: { row: Time {row.kind === "working" ? : null} {row.kind === "manager-instruction" ? : null} {row.kind === "context-trim" ? ( - ctx.onTrimToggle(row.trimPoint.id)} - /> + <> + {row.trimPoint.summary ? : null} + ctx.onTrimToggle(row.trimPoint.id)} + /> + ) : null} ); diff --git a/apps/web/src/composer-logic.test.ts b/apps/web/src/composer-logic.test.ts index f3be2385f081..bb504c98f76b 100644 --- a/apps/web/src/composer-logic.test.ts +++ b/apps/web/src/composer-logic.test.ts @@ -365,4 +365,16 @@ describe("parseStandaloneComposerSlashCommand", () => { it("ignores /new with extra message text", () => { expect(parseStandaloneComposerSlashCommand("/new start fresh")).toBeNull(); }); + + it("parses standalone /compact command", () => { + expect(parseStandaloneComposerSlashCommand("/compact")).toBe("compact"); + }); + + it("parses /compact with surrounding whitespace", () => { + expect(parseStandaloneComposerSlashCommand(" /compact ")).toBe("compact"); + }); + + it("ignores /compact with extra message text", () => { + expect(parseStandaloneComposerSlashCommand("/compact start fresh")).toBeNull(); + }); }); diff --git a/apps/web/src/composer-logic.ts b/apps/web/src/composer-logic.ts index 16cacf608c36..3b3ec51e8680 100644 --- a/apps/web/src/composer-logic.ts +++ b/apps/web/src/composer-logic.ts @@ -263,7 +263,7 @@ export function detectComposerTrigger(text: string, cursorInput: number): Compos export function parseStandaloneComposerSlashCommand( text: string, -): "plan" | "default" | "new" | ComposerClearSlashCommand | null { +): "plan" | "default" | "new" | "compact" | ComposerClearSlashCommand | null { const trimmed = text.trim(); const clearMatch = /^\/clear(?:\s+(\d+))?\s*$/i.exec(trimmed); if (clearMatch) { @@ -276,6 +276,9 @@ export function parseStandaloneComposerSlashCommand( if (/^\/new\s*$/i.test(trimmed)) { return "new"; } + if (/^\/compact\s*$/i.test(trimmed)) { + return "compact"; + } const match = /^\/(plan|default)\s*$/i.exec(trimmed); if (!match) { return null; diff --git a/packages/contracts/src/orchestration.test.ts b/packages/contracts/src/orchestration.test.ts index 72baca4b1eb7..813a913f1ff9 100644 --- a/packages/contracts/src/orchestration.test.ts +++ b/packages/contracts/src/orchestration.test.ts @@ -1119,6 +1119,194 @@ it.effect("decodes thread.context.trim as part of OrchestrationCommand union", ( }), ); +// ── thread.context.summarize command & event ────────────────────────── + +import { + ThreadContextSummarizeCommand, + ThreadContextSummarizedPayload, + ThreadContextCompactCommand, + ThreadContextCompactedPayload, +} from "./orchestration.ts"; + +const decodeThreadContextSummarizeCommand = Schema.decodeUnknownEffect(ThreadContextSummarizeCommand); +const decodeThreadContextSummarizedPayload = Schema.decodeUnknownEffect(ThreadContextSummarizedPayload); +const decodeThreadContextCompactCommand = Schema.decodeUnknownEffect(ThreadContextCompactCommand); +const decodeThreadContextCompactedPayload = Schema.decodeUnknownEffect(ThreadContextCompactedPayload); + +it.effect("decodes thread.context.summarize command", () => + Effect.gen(function* () { + const parsed = yield* decodeThreadContextSummarizeCommand({ + type: "thread.context.summarize", + commandId: "cmd-summarize-1", + threadId: "thread-1", + summary: "Compacted summary text", + compactDurationMs: 1500, + createdAt: "2026-01-01T00:00:00.000Z", + }); + assert.strictEqual(parsed.type, "thread.context.summarize"); + assert.strictEqual(parsed.summary, "Compacted summary text"); + assert.strictEqual(parsed.compactDurationMs, 1500); + }), +); + +it.effect("decodes thread.context-summarized payload", () => + Effect.gen(function* () { + const parsed = yield* decodeThreadContextSummarizedPayload({ + threadId: "thread-1", + trimPointId: "trim-1", + summary: "Compacted summary text", + compactDurationMs: 1500, + }); + assert.strictEqual(parsed.threadId, "thread-1"); + assert.strictEqual(parsed.trimPointId, "trim-1"); + assert.strictEqual(parsed.summary, "Compacted summary text"); + assert.strictEqual(parsed.compactDurationMs, 1500); + }), +); + +it.effect("decodes thread.context-summarized event", () => + Effect.gen(function* () { + const event = yield* decodeOrchestrationEvent({ + sequence: 200, + eventId: "event-summarized-1", + aggregateKind: "thread", + aggregateId: "thread-1", + type: "thread.context-summarized", + occurredAt: "2026-01-01T00:00:00.000Z", + commandId: "cmd-summarize-1", + causationEventId: null, + correlationId: "cmd-summarize-1", + metadata: {}, + payload: { + threadId: "thread-1", + trimPointId: "trim-1", + summary: "Compacted context summary", + compactDurationMs: 1200, + }, + }); + assert.strictEqual(event.type, "thread.context-summarized"); + if (event.type === "thread.context-summarized") { + assert.strictEqual(event.payload.summary, "Compacted context summary"); + assert.strictEqual(event.payload.compactDurationMs, 1200); + } + }), +); + +// ── thread.context.compact command & event ──────────────────────────── + +it.effect("decodes thread.context.compact command", () => + Effect.gen(function* () { + const parsed = yield* decodeThreadContextCompactCommand({ + type: "thread.context.compact", + commandId: "cmd-compact-1", + threadId: "thread-1", + createdAt: "2026-01-01T00:00:00.000Z", + }); + assert.strictEqual(parsed.type, "thread.context.compact"); + assert.strictEqual(parsed.threadId, "thread-1"); + }), +); + +it.effect("decodes thread.context-compacted payload", () => + Effect.gen(function* () { + const parsed = yield* decodeThreadContextCompactedPayload({ + threadId: "thread-1", + summary: "Provider compacted the conversation.", + compactDurationMs: 800, + }); + assert.strictEqual(parsed.threadId, "thread-1"); + assert.strictEqual(parsed.summary, "Provider compacted the conversation."); + assert.strictEqual(parsed.compactDurationMs, 800); + }), +); + +it.effect("decodes thread.context-compacted event", () => + Effect.gen(function* () { + const event = yield* decodeOrchestrationEvent({ + sequence: 300, + eventId: "event-compacted-1", + aggregateKind: "thread", + aggregateId: "thread-1", + type: "thread.context-compacted", + occurredAt: "2026-01-01T00:00:00.000Z", + commandId: "cmd-compact-1", + causationEventId: null, + correlationId: "cmd-compact-1", + metadata: {}, + payload: { + threadId: "thread-1", + summary: "Provider compacted the conversation.", + compactDurationMs: 800, + }, + }); + assert.strictEqual(event.type, "thread.context-compacted"); + if (event.type === "thread.context-compacted") { + assert.strictEqual(event.payload.summary, "Provider compacted the conversation."); + assert.strictEqual(event.payload.compactDurationMs, 800); + } + }), +); + +// ── ContextTrimPoint with summary ────────────────────────────────────── + +it.effect("decodes ContextTrimPoint with summary", () => + Effect.gen(function* () { + const parsed = yield* decodeContextTrimPoint({ + id: "trim-2", + createdAt: "2026-01-01T00:00:00.000Z", + beforeEntryId: "msg-10", + prunedMessageCount: 20, + prunedTurnIds: ["turn-3", "turn-4"], + summary: "This is a compaction summary.", + }); + assert.strictEqual(parsed.id, "trim-2"); + assert.strictEqual(parsed.prunedMessageCount, 20); + assert.strictEqual(parsed.summary, "This is a compaction summary."); + }), +); + +it.effect("decodes ContextTrimPoint without summary (backward compat)", () => + Effect.gen(function* () { + const parsed = yield* decodeContextTrimPoint({ + id: "trim-3", + createdAt: "2026-01-01T00:00:00.000Z", + beforeEntryId: "msg-5", + prunedMessageCount: 5, + prunedTurnIds: ["turn-1"], + }); + assert.strictEqual(parsed.summary, undefined); + }), +); + +// ── ThreadContextTrimCommand with summary ────────────────────────────── + +it.effect("decodes thread.context.trim command with summary", () => + Effect.gen(function* () { + const parsed = yield* decodeThreadContextTrimCommand({ + type: "thread.context.trim", + commandId: "cmd-trim-summary-1", + threadId: "thread-1", + summary: "Trim with compaction summary", + keepLastNTurns: 2, + createdAt: "2026-01-01T00:00:00.000Z", + }); + assert.strictEqual(parsed.summary, "Trim with compaction summary"); + assert.strictEqual(parsed.keepLastNTurns, 2); + }), +); + +it.effect("decodes thread.context.compact as part of OrchestrationCommand union", () => + Effect.gen(function* () { + const parsed = yield* decodeOrchestrationCommand({ + type: "thread.context.compact", + commandId: "cmd-compact-u", + threadId: "thread-1", + createdAt: "2026-01-01T00:00:00.000Z", + }); + assert.strictEqual(parsed.type, "thread.context.compact"); + }), +); + // ── thread.archive-and-new command & event ───────────────────────────── const decodeThreadArchivedAndNewCreatedPayload = Schema.decodeUnknownEffect( diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 8417554c1728..97caa486df39 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -469,6 +469,7 @@ export const ContextTrimPoint = Schema.Struct({ beforeEntryId: Schema.String, prunedMessageCount: NonNegativeInt, prunedTurnIds: Schema.Array(TurnId), + summary: Schema.optional(Schema.String), }); export type ContextTrimPoint = typeof ContextTrimPoint.Type; @@ -929,10 +930,29 @@ export const ThreadContextTrimCommand = Schema.Struct({ commandId: CommandId, threadId: ThreadId, keepLastNTurns: Schema.optional(NonNegativeInt), + summary: Schema.optional(Schema.String), createdAt: IsoDateTime, }); export type ThreadContextTrimCommand = typeof ThreadContextTrimCommand.Type; +export const ThreadContextSummarizeCommand = Schema.Struct({ + type: Schema.Literal("thread.context.summarize"), + commandId: CommandId, + threadId: ThreadId, + summary: Schema.String, + compactDurationMs: Schema.optional(NonNegativeInt), + createdAt: IsoDateTime, +}); +export type ThreadContextSummarizeCommand = typeof ThreadContextSummarizeCommand.Type; + +export const ThreadContextCompactCommand = Schema.Struct({ + type: Schema.Literal("thread.context.compact"), + commandId: CommandId, + threadId: ThreadId, + createdAt: IsoDateTime, +}); +export type ThreadContextCompactCommand = typeof ThreadContextCompactCommand.Type; + const DispatchableClientOrchestrationCommand = Schema.Union([ ProjectCreateCommand, ManagerBootstrapCommand, @@ -961,6 +981,7 @@ const DispatchableClientOrchestrationCommand = Schema.Union([ ThreadCheckpointRevertCommand, ThreadSessionStopCommand, ThreadContextTrimCommand, + ThreadContextCompactCommand, ]); export type DispatchableClientOrchestrationCommand = typeof DispatchableClientOrchestrationCommand.Type; @@ -993,6 +1014,7 @@ export const ClientOrchestrationCommand = Schema.Union([ ThreadCheckpointRevertCommand, ThreadSessionStopCommand, ThreadContextTrimCommand, + ThreadContextCompactCommand, ]); export type ClientOrchestrationCommand = typeof ClientOrchestrationCommand.Type; @@ -1069,6 +1091,7 @@ const InternalOrchestrationCommand = Schema.Union([ ThreadTurnDiffCompleteCommand, ThreadActivityAppendCommand, ThreadRevertCompleteCommand, + ThreadContextSummarizeCommand, ]); export type InternalOrchestrationCommand = typeof InternalOrchestrationCommand.Type; @@ -1105,6 +1128,8 @@ export const OrchestrationEventType = Schema.Literals([ "thread.activity-appended", "thread.manager-queue-items-upserted", "thread.trim-point-created", + "thread.context-summarized", + "thread.context-compacted", "thread.archived-and-new-created", ]); export type OrchestrationEventType = typeof OrchestrationEventType.Type; @@ -1316,6 +1341,19 @@ export const ThreadTrimPointCreatedPayload = Schema.Struct({ trimPoint: ContextTrimPoint, }); +export const ThreadContextSummarizedPayload = Schema.Struct({ + threadId: ThreadId, + trimPointId: Schema.String, + summary: Schema.String, + compactDurationMs: Schema.optional(NonNegativeInt), +}); + +export const ThreadContextCompactedPayload = Schema.Struct({ + threadId: ThreadId, + summary: Schema.String, + compactDurationMs: Schema.optional(NonNegativeInt), +}); + export const OrchestrationEventMetadata = Schema.Struct({ providerTurnId: Schema.optional(TrimmedNonEmptyString), providerItemId: Schema.optional(ProviderItemId), @@ -1468,6 +1506,16 @@ export const OrchestrationEvent = Schema.Union([ type: Schema.Literal("thread.trim-point-created"), payload: ThreadTrimPointCreatedPayload, }), + Schema.Struct({ + ...EventBaseFields, + type: Schema.Literal("thread.context-summarized"), + payload: ThreadContextSummarizedPayload, + }), + Schema.Struct({ + ...EventBaseFields, + type: Schema.Literal("thread.context-compacted"), + payload: ThreadContextCompactedPayload, + }), Schema.Struct({ ...EventBaseFields, type: Schema.Literal("thread.archived-and-new-created"), From 8955033a19858af09c9423a60c4d300712272a7c Mon Sep 17 00:00:00 2001 From: harrydawson Date: Mon, 8 Jun 2026 18:41:27 +1000 Subject: [PATCH 2/2] Add token usage bug report for ClaudeAdapter --- token-usage-bug-report.html | 296 ++++++++++++++++++++++++++++++++++++ 1 file changed, 296 insertions(+) create mode 100644 token-usage-bug-report.html diff --git a/token-usage-bug-report.html b/token-usage-bug-report.html new file mode 100644 index 000000000000..fe6a489e6bee --- /dev/null +++ b/token-usage-bug-report.html @@ -0,0 +1,296 @@ + + + + + +Token Usage Bug Report — ClaudeAdapter + + + + +

Token Usage Bug Report HIGH

+

ClaudeAdapter — Context Window Occupancy Over-Reported

+ +
+

What: The context-window usage meter reports values that are N times larger than the actual context occupancy, where N = the number of API calls made during a turn.

+

Impact: After a few tool uses, the meter hits 100% even though actual context may only be 20–30% full. Users believe their session is exhausted and start new ones unnecessarily.

+

Root cause: The code reads token counts from sources that accumulate across all API calls in a turn, but treats those accumulated totals as the current snapshot of context window usage.

+
+ +

1. What Should Happen

+

The context window meter should show how many tokens are currently in the model's context window. When Claude makes an API call, the response includes per-request usage: input_tokens + cache_creation_input_tokens + cache_read_input_tokens = the number of tokens occupying the context window for that single request.

+

These per-request numbers arrive in two stream events that the Anthropic SDK sends for every API response:

+ + + + + + + + + + + + + +
Stream EventContainsAccuracy
message_startevent.message.usage (BetaUsage)Per-request input tokens at the start of the stream
message_deltaevent.usage (BetaMessageDeltaUsage)Per-request cumulative usage including output, at the end of the stream
+ +

The last message_delta.usage from the most recent API call reflects the true context window occupancy.

+ +

2. What Actually Happens

+

There are three code paths that all set the token usage to an accumulated total instead of the per-request snapshot:

+ +

Path A: task_progress events (subagents)

+

ClaudeAdapter.ts:2238–2256

+
case "task_progress":
+  if (message.usage) {
+    const normalizedUsage = normalizeClaudeTokenUsage(
+      message.usage,
+      context.lastKnownContextWindow
+    );
+    if (normalizedUsage) {
+      context.lastKnownTokenUsage = normalizedUsage;  // ← BUG: stored as current usage
+      yield* offerRuntimeEvent({
+        type: "thread.token-usage.updated",            // ← BUG: emitted as context window update
+        payload: { usage: normalizedUsage },
+      });
+    }
+  }
+

task_progress.usage.total_tokens is the sum of tokens across all API calls the subagent has made so far. normalizeClaudeTokenUsage treats this accumulated sum as the current context window size, only capping it at maxTokens.

+

The same pattern exists for task_notification at ClaudeAdapter.ts:2270–2289.

+ +

Path B: completeTurn fallback (no subagents)

+

ClaudeAdapter.ts:1506–1533

+
// The SDK result.usage contains *accumulated* totals across all API calls
+// This does NOT represent the current context window size.
+// Instead, use the last known context-window-accurate usage from task_progress
+// events and treat the accumulated total as totalProcessedTokens.    ← comment IS wrong
+
+const accumulatedSnapshot = normalizeClaudeTokenUsage(result?.usage, ...);
+const lastGoodUsage = context.lastKnownTokenUsage;     // may be undefined if no subagents
+
+const usageSnapshot = lastGoodUsage
+  ? { ...lastGoodUsage, ... }   // ← uses task_progress value (also accumulated)
+  : accumulatedSnapshot;        // ← fallback: uses result.usage (also accumulated)
+

When no subagents are used, task_progress events never fire, so lastKnownTokenUsage remains undefined. The fallback uses result.usage, which is accumulated across all API calls in the turn.

+

The comment at lines 1509–1510 says "use the last known context-window-accurate usage from task_progress events" — but task_progress is also accumulated. The comment is misleading.

+ +

Path C: handleStreamEvent ignores the correct source

+

ClaudeAdapter.ts:1670–1917

+

The handleStreamEvent function only processes three event types:

+
if (event.type === "content_block_delta")  { /* handle text/tool deltas */ }
+if (event.type === "content_block_start")  { /* handle text/tool blocks */ }
+if (event.type === "content_block_stop")   { /* handle completion */ }
+// ← message_start and message_delta events are silently ignored
+

message_start and message_delta events — which carry the correct per-request token counts — are completely ignored. The function returns undefined for them and moves on.

+ +

3. Walkthrough Example

+ +
+

Scenario: A turn that makes 3 API calls

+
    +
  1. First API call: reads 20,000 context tokens, produces 500 output tokens
  2. +
  3. Second API call (after tool result): reads 22,000 context tokens, produces 600 output tokens
  4. +
  5. Third API call (after second tool result): reads 23,000 context tokens, produces 700 output tokens
  6. +
+
+ +
+
+

What the code reports

+

result.usage.total_tokens67,400 tokens

+

The meter shows 67,400 / 200,000 = 33% after just 3 calls.
+ After 6 calls: 100% (even though real context is ~30%).

+
+
+

What should be reported

+

Last message_delta.usage: ~23,000 tokens (actual context window occupancy)

+

The meter shows 23,000 / 200,000 = 11.5%

+
+
+ +

4. The Root Cause Function

+

ClaudeAdapter.ts:345–389

+
function normalizeClaudeTokenUsage(value: unknown, contextWindow?: number) {
+  const inputTokens = (usage.input_tokens ?? 0) + (usage.cache_creation_input_tokens ?? 0) + (usage.cache_read_input_tokens ?? 0);
+  const outputTokens = usage.output_tokens ?? 0;
+  const totalProcessedTokens = usage.total_tokens ?? (inputTokens + outputTokens);
+  const usedTokens = Math.min(totalProcessedTokens, maxTokens ?? Infinity); // cap prevents overflow but doesn't fix
+
+  return { usedTokens, lastUsedTokens: usedTokens, ... };
+}
+

The function uses total_tokens as usedTokens without any awareness of whether the value represents a per-request count or an accumulated total. It treats both identically.

+ +

5. Summary of Affected Locations

+ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +
FileLinesProblem
ClaudeAdapter.ts1670–1917handleStreamEvent ignores message_start/message_delta events that carry per-request usage
ClaudeAdapter.ts2238–2256task_progress handler stores accumulated total as current context window usage
ClaudeAdapter.ts2270–2289task_notification handler — same bug as task_progress
ClaudeAdapter.ts1506–1533completeTurn falls back to accumulated result.usage or accumulated lastKnownTokenUsage
ClaudeAdapter.ts1509–1510Misleading comment claiming task_progress is "context-window-accurate"
ClaudeAdapter.ts345–389normalizeClaudeTokenUsage doesn't distinguish per-request vs. accumulated totals
ClaudeAdapter.test.ts1729Test validates the buggy behavior: emitting thread.token-usage.updated from task_progress
+ +

6. How to Fix

+ +

Step 1: Capture per-request usage from stream events

+

In handleStreamEvent, add handling for message_start and message_delta:

+
if (event.type === "message_start") {
+  const usage = event.message?.usage;
+  if (usage && typeof usage === "object") {
+    const normalized = normalizeClaudeTokenUsage(usage, context.lastKnownContextWindow);
+    if (normalized) context.lastKnownTokenUsage = normalized;
+  }
+  return;
+}
+
+if (event.type === "message_delta") {
+  const usage = event.usage;
+  if (usage && typeof usage === "object") {
+    const normalized = normalizeClaudeTokenUsage(usage, context.lastKnownContextWindow);
+    if (normalized) context.lastKnownTokenUsage = normalized;
+  }
+  return;
+}
+

This populates lastKnownTokenUsage with the correct per-request numbers during normal streaming.

+ +

Step 2: Remove context-window side effects from task progress/notification

+

In both task_progress and task_notification handlers, remove:

+
    +
  • The context.lastKnownTokenUsage = normalizedUsage assignment
  • +
  • The thread.token-usage.updated event emission
  • +
+

Keep the task.progress and task.completed events needed for subagent UI.

+ +

Step 3: Fix the misleading comment

+

Replace lines 1509–1510:

+
// Instead, use per-request usage captured from message_start/message_delta
+// stream events. task_progress.total_tokens is also accumulated and not suitable.
+ +

Step 4: Update tests

+

Update the test at ClaudeAdapter.test.ts:1729 to verify that task_progress no longer emits thread.token-usage.updated. Add a new test verifying that message_delta stream events correctly set lastKnownTokenUsage.

+ +

Conclusion

+
+

The bug is confirmed and present in the current codebase.

+

The core problem is that handleStreamEvent ignores message_start and message_delta events — the only source of per-API-call token usage. Instead, the code reads token counts from accumulated sources (task_progress.total_tokens, result.usage) and presents them as current context window occupancy. This causes the context meter to grow at N× the real rate, where N is the number of API calls per turn.

+

The fix is small and surgical: add message_start/message_delta handling to handleStreamEvent, remove the incorrect token-usage updates from task_progress/task_notification, and fix the comment.

+
+ + +