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"),