diff --git a/apps/server/src/provider/Layers/AntigravityAdapter.ts b/apps/server/src/provider/Layers/AntigravityAdapter.ts index aa7d4c6a6755..97925cd4912d 100644 --- a/apps/server/src/provider/Layers/AntigravityAdapter.ts +++ b/apps/server/src/provider/Layers/AntigravityAdapter.ts @@ -1245,6 +1245,7 @@ export const makeAntigravityAdapter = Effect.fn("makeAntigravityAdapter")(functi return { provider: PROVIDER, capabilities: { sessionModelSwitch: "in-session", supportsConversationRollback: false }, + compaction: { type: "slash-command", command: "/compact" }, startSession, sendTurn, interruptTurn, diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index 2f0d15e281df..1ecde618e191 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -5075,6 +5075,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( capabilities: { sessionModelSwitch: "in-session", }, + compaction: { type: "slash-command", command: "/compact" }, startSession, sendTurn, interruptTurn, diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 4676d780a530..ef6e97d8993d 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -353,7 +353,8 @@ sessionErrorLayer("CodexAdapterLive session errors", (it) => { Stream.runHead, Effect.forkChild, ); - yield* adapter.compactThread!(threadId); + NodeAssert.ok(adapter.compaction?.type === "native"); + yield* adapter.compaction.start(threadId); yield* runtime.emit({ id: asEventId("evt-compaction-item-completed"), kind: "notification", diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index d1981b33d47d..5e2244336afc 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -2504,14 +2504,12 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( ), ); - const compactThread: NonNullable = Effect.fn("compactThread")( - function* (threadId) { - const session = yield* requireSession(threadId); - yield* session.runtime.compactThread.pipe( - Effect.mapError((cause) => mapCodexRuntimeError(threadId, "thread/compact/start", cause)), - ); - }, - ); + const compactThread = Effect.fn("compactThread")(function* (threadId: ThreadId) { + const session = yield* requireSession(threadId); + yield* session.runtime.compactThread.pipe( + Effect.mapError((cause) => mapCodexRuntimeError(threadId, "thread/compact/start", cause)), + ); + }); const readThread: CodexAdapterShape["readThread"] = (threadId) => requireSession(threadId).pipe( @@ -2658,7 +2656,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( }, startSession, sendTurn, - compactThread, + compaction: { type: "native", start: compactThread }, interruptTurn, readThread, rollbackThread, diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index a8aea90e8972..1ed9648a50b6 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -1212,6 +1212,7 @@ export function makeCursorAdapter( return { provider: PROVIDER, capabilities: { sessionModelSwitch: "in-session" }, + compaction: { type: "slash-command", command: "/compress" }, startSession, sendTurn, interruptTurn, diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index dae7c2ca08f6..25188adcffcc 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -2128,6 +2128,7 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte return { provider: PROVIDER, capabilities: { sessionModelSwitch: "in-session" }, + compaction: { type: "slash-command", command: "/compact" }, startSession, sendTurn, interruptTurn, diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index c8a7d12af276..ee5767f9d356 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -1080,7 +1080,8 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { threadId, runtimeMode: "full-access", }); - yield* adapter.compactThread!( + NodeAssert.ok(adapter.compaction?.type === "native"); + yield* adapter.compaction.start( threadId, createModelSelection(ProviderInstanceId.make("opencode"), "openai/gpt-5"), ); diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 38444be4ad19..742ee9b86d6a 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -4,6 +4,7 @@ import { ProviderDriverKind, ProviderInstanceId, type ProviderRuntimeEvent, + type ProviderSendTurnInput, type ProviderSession, RuntimeItemId, RuntimeRequestId, @@ -3426,9 +3427,10 @@ export function makeOpenCodeAdapter( ); }); - const compactThread: NonNullable = Effect.fn( - "compactThread", - )(function* (threadId, requestedModelSelection) { + const compactThread = Effect.fn("compactThread")(function* ( + threadId: ThreadId, + requestedModelSelection?: ProviderSendTurnInput["modelSelection"], + ) { const context = yield* ensureSessionContext(sessions, threadId); yield* awaitOpenCodeContextReady(context); const modelSelection = @@ -3843,7 +3845,7 @@ export function makeOpenCodeAdapter( }, startSession, sendTurn, - compactThread, + compaction: { type: "native", start: compactThread }, interruptTurn, respondToRequest, respondToUserInput, diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index 0342f3ca79e3..4d17aabaa5f5 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -190,7 +190,7 @@ function makeFakeCodexAdapter( Effect.void, ); - const compactThread = vi.fn((threadId: ThreadId) => + const compactThread = vi.fn((threadId: ThreadId): Effect.Effect => Effect.sync(() => emit({ type: "thread.state.changed", @@ -278,7 +278,13 @@ function makeFakeCodexAdapter( }, startSession, sendTurn, - ...(provider === CODEX_DRIVER ? { compactThread } : {}), + ...(provider === CODEX_DRIVER + ? { compaction: { type: "native", start: compactThread } } + : provider === CURSOR_DRIVER + ? { compaction: { type: "slash-command", command: "/compress" } } + : provider === CLAUDE_AGENT_DRIVER + ? { compaction: { type: "slash-command", command: "/compact" } } + : {}), interruptTurn, respondToRequest, respondToUserInput, @@ -981,6 +987,125 @@ it.effect("ProviderServiceLive rejects new sessions for disabled custom instance const routing = makeProviderServiceLayer(); +const customCompactionDriver = ProviderDriverKind.make("custom-compaction-provider"); +const nativeCompactionInstanceId = ProviderInstanceId.make("native-compaction"); +const slashCompactionInstanceId = ProviderInstanceId.make("slash-compaction"); +const unsupportedCompactionInstanceId = ProviderInstanceId.make("unsupported-compaction"); +const customNativeCompaction = makeFakeCodexAdapter(customCompactionDriver); +const customSlashCompaction = makeFakeCodexAdapter(customCompactionDriver); +const unsupportedCompaction = makeFakeCodexAdapter(customCompactionDriver); +const declaredCompaction = makeProviderServiceLayer({ + registry: makeStaticInstanceRegistry([ + [ + nativeCompactionInstanceId, + { + ...customNativeCompaction.adapter, + compaction: { type: "native", start: customNativeCompaction.compactThread }, + }, + ], + [ + slashCompactionInstanceId, + { + ...customSlashCompaction.adapter, + compaction: { type: "slash-command", command: "/reduce-context" }, + }, + ], + [unsupportedCompactionInstanceId, unsupportedCompaction.adapter], + ]), +}); + +declaredCompaction.layer("ProviderService declared compaction", (it) => { + it.effect("starts declared native compaction instead of sending a prompt", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("custom-native-compaction"); + const requestId = MessageId.make("custom-native-request"); + yield* provider.startSession(threadId, { + providerInstanceId: nativeCompactionInstanceId, + threadId, + runtimeMode: "full-access", + }); + const compactedEventFiber = yield* provider.streamEvents.pipe( + Stream.filter( + (event) => event.threadId === threadId && event.type === "thread.state.changed", + ), + Stream.runHead, + Effect.forkChild({ startImmediately: true }), + ); + yield* advanceTestClock(50); + yield* provider.compactThread(threadId, undefined, requestId); + const compacted = Option.getOrThrow(yield* Fiber.join(compactedEventFiber)); + assert.equal(compacted.requestId, String(requestId)); + assert.equal(customNativeCompaction.compactThread.mock.calls.length, 1); + assert.equal(customNativeCompaction.sendTurn.mock.calls.length, 0); + yield* provider.stopSession({ threadId }); + }), + ); + + it.effect("sends the declared slash command as the compaction turn", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("custom-slash-compaction"); + const requestId = MessageId.make("custom-slash-request"); + const modelSelection = createModelSelection(slashCompactionInstanceId, "custom-model"); + yield* provider.startSession(threadId, { + providerInstanceId: slashCompactionInstanceId, + threadId, + runtimeMode: "full-access", + }); + const compactedEventFiber = yield* provider.streamEvents.pipe( + Stream.filter( + (event) => event.threadId === threadId && event.type === "thread.state.changed", + ), + Stream.runHead, + Effect.forkChild({ startImmediately: true }), + ); + const compactFiber = yield* provider + .compactThread(threadId, modelSelection, requestId) + .pipe(Effect.forkChild); + yield* advanceTestClock(50); + customSlashCompaction.emit({ + type: "turn.completed", + eventId: asEventId("custom-slash-completed"), + provider: customCompactionDriver, + createdAt: "2026-01-01T00:00:01.000Z", + threadId, + turnId: asTurnId(`turn-${threadId}`), + payload: { state: "completed" }, + }); + yield* Fiber.join(compactFiber); + const compacted = Option.getOrThrow(yield* Fiber.join(compactedEventFiber)); + assert.equal(compacted.requestId, String(requestId)); + assert.equal(customSlashCompaction.compactThread.mock.calls.length, 0); + assert.equal(customSlashCompaction.sendTurn.mock.calls.length, 1); + assert.equal(customSlashCompaction.sendTurn.mock.calls[0]?.[0].input, "/reduce-context"); + assert.deepEqual( + customSlashCompaction.sendTurn.mock.calls[0]?.[0].modelSelection, + modelSelection, + ); + yield* provider.stopSession({ threadId }); + }), + ); + + it.effect("rejects compaction for adapters without a declared strategy", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const threadId = asThreadId("custom-unsupported-compaction"); + yield* provider.startSession(threadId, { + providerInstanceId: unsupportedCompactionInstanceId, + threadId, + runtimeMode: "full-access", + }); + const failure = yield* provider.compactThread(threadId).pipe(Effect.flip); + assert.instanceOf(failure, ProviderValidationError); + assert.include(failure.message, "does not support context compaction"); + assert.equal(unsupportedCompaction.sendTurn.mock.calls.length, 0); + assert.equal(unsupportedCompaction.compactThread.mock.calls.length, 0); + yield* provider.stopSession({ threadId }); + }), + ); +}); + const antigravityDriver = ProviderDriverKind.make("antigravity"); const replacementAntigravity = makeFakeCodexAdapter(antigravityDriver); const originalAntigravityInstanceId = ProviderInstanceId.make("antigravity-personal"); diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index d9cac46ec4d9..2b2719faabd3 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -76,6 +76,9 @@ import * as ServerSettings from "../../serverSettings.ts"; import * as ProjectionSnapshotQuery from "../../orchestration/Services/ProjectionSnapshotQuery.ts"; const isModelSelection = Schema.is(ModelSelection); +/** How long a manual context compaction may run before ProviderService gives up on it. */ +const COMPACTION_COMPLETION_TIMEOUT = "10 minutes"; + interface PendingCompaction { readonly completion: Deferred.Deferred; readonly native: boolean; @@ -1522,18 +1525,24 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( "provider.thread_id": threadId, }); yield* McpSessionRegistry.touchActiveMcpThread(threadId); - const nativeCompaction = routed.adapter.compactThread; + const compaction = routed.adapter.compaction; + if (compaction === undefined) { + return yield* toValidationError( + "ProviderService.compactThread", + `Provider '${routed.adapter.provider}' does not support context compaction.`, + ); + } const completion = yield* Deferred.make(); const pending: PendingCompaction = { completion, - native: nativeCompaction !== undefined, + native: compaction.type === "native", providerInstanceId: routed.instanceId, requestId, earlyEvents: [], compactedEventObserved: false, expectedTurnId: undefined, }; - if (nativeCompaction !== undefined && timedOutNativeCompactions.has(threadId)) { + if (compaction.type === "native" && timedOutNativeCompactions.has(threadId)) { return yield* new ProviderAdapterRequestError({ provider: routed.adapter.provider, method: "thread/compact", @@ -1558,14 +1567,10 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( pendingCompactions.delete(threadId); } }); - const nativeCompletionTimeout = - routed.adapter.provider === "codex" || routed.adapter.provider === "opencode" - ? "10 minutes" - : "30 seconds"; const awaitNativeCompaction = (start: Effect.Effect) => start.pipe( Effect.andThen(Deferred.await(completion)), - Effect.timeout(nativeCompletionTimeout), + Effect.timeout(COMPACTION_COMPLETION_TIMEOUT), Effect.catchTag("TimeoutError", (cause) => Effect.sync(() => { timedOutNativeCompactions.add(threadId); @@ -1575,7 +1580,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( new ProviderAdapterRequestError({ provider: routed.adapter.provider, method: "thread/compact", - detail: `Provider did not report completed context compaction within ${nativeCompletionTimeout}.`, + detail: `Provider did not report completed context compaction within ${COMPACTION_COMPLETION_TIMEOUT}.`, cause, }), ), @@ -1584,24 +1589,24 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ), ); const awaitFallbackCompaction = Deferred.await(completion).pipe( - Effect.timeout("10 minutes"), + Effect.timeout(COMPACTION_COMPLETION_TIMEOUT), Effect.mapError( (cause) => new ProviderAdapterRequestError({ provider: routed.adapter.provider, method: "turn/start", - detail: "Provider did not finish context compaction within 10 minutes.", + detail: `Provider did not finish context compaction within ${COMPACTION_COMPLETION_TIMEOUT}.`, cause, }), ), ); const terminal = yield* ( - nativeCompaction - ? awaitNativeCompaction(nativeCompaction(routed.threadId, modelSelection)) + compaction.type === "native" + ? awaitNativeCompaction(compaction.start(routed.threadId, modelSelection)) : Effect.gen(function* () { const turn = yield* sendTurn({ threadId, - input: routed.adapter.provider === "cursor" ? "/compress" : "/compact", + input: compaction.command, ...(modelSelection !== undefined ? { modelSelection } : {}), }).pipe( Effect.onError(() => @@ -1621,7 +1626,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( if (terminal !== "completed") { return yield* new ProviderAdapterRequestError({ provider: routed.adapter.provider, - method: nativeCompaction ? "thread/compact" : "turn/start", + method: compaction.type === "native" ? "thread/compact" : "turn/start", detail: `Context compaction ended with ${terminal}.`, }); } diff --git a/apps/server/src/provider/Services/ProviderAdapter.ts b/apps/server/src/provider/Services/ProviderAdapter.ts index 0e4d696335b0..c9b62fd79525 100644 --- a/apps/server/src/provider/Services/ProviderAdapter.ts +++ b/apps/server/src/provider/Services/ProviderAdapter.ts @@ -27,6 +27,21 @@ import type * as Stream from "effect/Stream"; export type ProviderSessionModelSwitchMode = "in-session" | "unsupported"; +/** + * How ProviderService runs manual context compaction for an adapter. + * Native adapters expose a start call and must emit a compacted thread state + * when they finish. Slash-command adapters get the command sent as a turn. + */ +export type ProviderCompaction = + | { + readonly type: "native"; + readonly start: ( + threadId: ThreadId, + modelSelection?: ProviderSendTurnInput["modelSelection"], + ) => Effect.Effect; + } + | { readonly type: "slash-command"; readonly command: `/${string}` }; + export interface ProviderAdapterCapabilities { /** * Declares whether changing the model on an existing session is supported. @@ -70,10 +85,8 @@ export interface ProviderAdapterShape { input: ProviderSendTurnInput, ) => Effect.Effect; - readonly compactThread?: ( - threadId: ThreadId, - modelSelection?: ProviderSendTurnInput["modelSelection"], - ) => Effect.Effect; + /** Omitted when this adapter does not support manual context compaction. */ + readonly compaction?: ProviderCompaction; /** * Interrupt an active turn.