diff --git a/.changeset/preserve-openai-compaction-state.md b/.changeset/preserve-openai-compaction-state.md new file mode 100644 index 00000000000..b38cb9f6cbe --- /dev/null +++ b/.changeset/preserve-openai-compaction-state.md @@ -0,0 +1,10 @@ +--- +"@moonshot-ai/kimi-code": patch +"@moonshot-ai/agent-core": patch +"@moonshot-ai/agent-core-v2": patch +"@moonshot-ai/kosong": patch +--- + +Preserve opaque OpenAI Responses compaction state across turns and automatically +use `/responses/compact` when the active provider exposes that capability, +falling back to Kimi's existing local summarizer when it does not. diff --git a/apps/kimi-code/src/tui/utils/message-replay.ts b/apps/kimi-code/src/tui/utils/message-replay.ts index d778b9a4721..8edb209d498 100644 --- a/apps/kimi-code/src/tui/utils/message-replay.ts +++ b/apps/kimi-code/src/tui/utils/message-replay.ts @@ -171,6 +171,7 @@ export function collectReplayMessageContent( break; case 'audio_url': case 'image_url': + case 'openai_compaction': case 'video_url': break; } @@ -285,6 +286,8 @@ function contentPartToText(part: ContentPart): string { return mediaUrlPartToText('video', part.videoUrl.url); case 'audio_url': return mediaUrlPartToText('audio', part.audioUrl.url); + case 'openai_compaction': + return ''; } } diff --git a/apps/vscode/src/runtime/replay-adapter.ts b/apps/vscode/src/runtime/replay-adapter.ts index facd978ac68..bdaa995f9d2 100644 --- a/apps/vscode/src/runtime/replay-adapter.ts +++ b/apps/vscode/src/runtime/replay-adapter.ts @@ -511,6 +511,8 @@ function toLegacyContent(content: readonly ContentPart[]): LegacyContentPart[] { case "video_url": result.push({ type: "video_url", video_url: { ...part.videoUrl } }); break; + case "openai_compaction": + break; } } return result; diff --git a/apps/vscode/src/utils/session-context.ts b/apps/vscode/src/utils/session-context.ts index e4ff4ce5b63..e2e21d0473b 100644 --- a/apps/vscode/src/utils/session-context.ts +++ b/apps/vscode/src/utils/session-context.ts @@ -196,6 +196,7 @@ function formatPartMarkdown(part: ContentPart): string { case "image_url": return "[image]"; case "audio_url": return "[audio]"; case "video_url": return "[video]"; + case "openai_compaction": return ""; } } @@ -205,7 +206,8 @@ function stringifyParts(parts: readonly ContentPart[]): string { if (part.type === "think") return part.think.trim() ? `\n${part.think}\n` : ""; if (part.type === "image_url") return "[image]"; if (part.type === "audio_url") return "[audio]"; - return "[video]"; + if (part.type === "video_url") return "[video]"; + return ""; }).filter(Boolean).join("\n"); } diff --git a/docs/en/guides/sessions.md b/docs/en/guides/sessions.md index d5010c3f3f5..9f79c10f5da 100644 --- a/docs/en/guides/sessions.md +++ b/docs/en/guides/sessions.md @@ -68,6 +68,12 @@ You can manage sessions without leaving the terminal. The following slash comman As a conversation grows, Kimi Code CLI automatically compresses the message history when the context approaches the window limit, freeing up token space. You can also trigger compression manually at any time: +For OpenAI Responses-compatible providers, Kimi Code automatically tries the +provider's native `/responses/compact` endpoint and preserves its opaque +replacement state. If that capability is unavailable, or when `/compact` +includes a custom instruction, it uses Kimi Code's existing summary compaction +instead. No separate setting is required. + ``` /compact ``` diff --git a/docs/zh/guides/sessions.md b/docs/zh/guides/sessions.md index 63a79e1abe6..bf1b5aca04a 100644 --- a/docs/zh/guides/sessions.md +++ b/docs/zh/guides/sessions.md @@ -68,6 +68,10 @@ kimi --session 对话变长时,Kimi Code CLI 会在上下文接近窗口上限时自动压缩历史消息,释放 token 空间。也可以随时手动触发: +对于兼容 OpenAI Responses 的提供商,Kimi Code 会自动尝试其原生 +`/responses/compact` 接口,并保留接口返回的不透明替换状态。如果该能力不可用, +或 `/compact` 带有自定义指引,则回退到 Kimi Code 原有的摘要压缩。无需额外配置。 + ``` /compact ``` diff --git a/packages/agent-core-v2/docs/wire-manifest.d.ts b/packages/agent-core-v2/docs/wire-manifest.d.ts index bd572a7c8cb..ce5b58134f0 100644 --- a/packages/agent-core-v2/docs/wire-manifest.d.ts +++ b/packages/agent-core-v2/docs/wire-manifest.d.ts @@ -114,7 +114,7 @@ interface ContextAppendMessagePayload { message: { role: 'system' | 'user' | 'assistant' | 'tool'; name?: string; - content: ('text' | 'think' | 'image_url' | 'audio_url' | 'video_url')[]; + content: ('text' | 'think' | 'image_url' | 'audio_url' | 'video_url' | 'openai_compaction')[]; toolCalls: { type: 'function'; id: string; diff --git a/packages/agent-core-v2/src/agent/contextMemory/compactionHandoff.ts b/packages/agent-core-v2/src/agent/contextMemory/compactionHandoff.ts index 966c533ffe9..9748a28778b 100644 --- a/packages/agent-core-v2/src/agent/contextMemory/compactionHandoff.ts +++ b/packages/agent-core-v2/src/agent/contextMemory/compactionHandoff.ts @@ -24,6 +24,7 @@ export interface CompactionUserSelection { export interface ContextCompactionShapeInput { readonly summary: string; + readonly replacementMessages?: readonly ContextMessage[]; readonly legacySummaryMessage?: ContextMessage; readonly contextSummary?: string; readonly compactedCount: number; @@ -51,6 +52,22 @@ export function buildContextCompactionShape( history: readonly ContextMessage[], input: ContextCompactionShapeInput, ): ContextCompactionShape { + if (input.replacementMessages !== undefined) { + const messages = [...input.replacementMessages]; + const contextSummary = input.contextSummary ?? input.summary; + return { + summary: input.summary, + contextSummary, + compactedCount: input.compactedCount, + tokensBefore: input.tokensBefore, + tokensAfter: input.tokensAfter ?? estimateTokensForMessages(messages), + keptUserMessageCount: + input.keptUserMessageCount ?? messages.filter((message) => message.role === 'user').length, + keptHeadUserMessageCount: input.keptHeadUserMessageCount, + droppedCount: input.droppedCount, + messages, + }; + } if (usesLegacyTailShape(input)) { const contextSummary = input.contextSummary ?? input.summary; const messages = [ diff --git a/packages/agent-core-v2/src/agent/contextMemory/contextMemory.ts b/packages/agent-core-v2/src/agent/contextMemory/contextMemory.ts index 3334d76ab4f..476c69660f9 100644 --- a/packages/agent-core-v2/src/agent/contextMemory/contextMemory.ts +++ b/packages/agent-core-v2/src/agent/contextMemory/contextMemory.ts @@ -6,6 +6,8 @@ import type { ContextMessage } from './types'; export interface ContextCompactionInput { readonly summary: string; + /** Provider-owned canonical replacement window, when native compaction is used. */ + readonly replacementMessages?: readonly ContextMessage[]; readonly contextSummary?: string; readonly compactedCount: number; readonly tokensBefore: number; diff --git a/packages/agent-core-v2/src/agent/contextMemory/contextMemoryService.ts b/packages/agent-core-v2/src/agent/contextMemory/contextMemoryService.ts index 53485e99e41..aae03853021 100644 --- a/packages/agent-core-v2/src/agent/contextMemory/contextMemoryService.ts +++ b/packages/agent-core-v2/src/agent/contextMemory/contextMemoryService.ts @@ -113,6 +113,8 @@ export class AgentContextMemoryService extends Disposable implements IAgentConte this.wire.dispatch( contextApplyCompaction({ summary: result.summary, + replacementMessages: + input.replacementMessages === undefined ? undefined : [...input.replacementMessages], contextSummary: result.contextSummary, compactedCount: result.compactedCount, tokensBefore: result.tokensBefore, diff --git a/packages/agent-core-v2/src/agent/contextMemory/contextOps.ts b/packages/agent-core-v2/src/agent/contextMemory/contextOps.ts index 8446efadd0a..edcc7c00bff 100644 --- a/packages/agent-core-v2/src/agent/contextMemory/contextOps.ts +++ b/packages/agent-core-v2/src/agent/contextMemory/contextOps.ts @@ -156,6 +156,7 @@ const contextCompactionBaseShape = { keptHeadUserMessageCount: z.number().optional(), droppedCount: z.number().optional(), legacyTail: z.boolean().optional(), + replacementMessages: z.array(contextMessageSchema).optional(), }; const contextApplyCompactionSchema = z.union([ @@ -210,6 +211,7 @@ export function readContextCompactionShapeInput( const keptUserMessageCount = readOptionalNumber(fields, 'keptUserMessageCount'); return { summary: readContextCompactionRawSummary(fields), + replacementMessages: readReplacementMessages(fields), legacySummaryMessage: readLegacySummaryMessage(fields), contextSummary: readOptionalString(fields, 'contextSummary'), compactedCount: readContextCompactedCount(fields), @@ -222,6 +224,11 @@ export function readContextCompactionShapeInput( }; } +function readReplacementMessages(record: UnknownRecord): readonly ContextMessage[] | undefined { + const value = record['replacementMessages']; + return Array.isArray(value) ? (value as ContextMessage[]) : undefined; +} + export function readContextCompactedCount(record: ContextCompactionRecord): number { const fields = record as UnknownRecord; const compactedCount = fields['compactedCount']; diff --git a/packages/agent-core-v2/src/agent/contextMemory/messageProjection.ts b/packages/agent-core-v2/src/agent/contextMemory/messageProjection.ts index 7e48c659951..d97664bdcc6 100644 --- a/packages/agent-core-v2/src/agent/contextMemory/messageProjection.ts +++ b/packages/agent-core-v2/src/agent/contextMemory/messageProjection.ts @@ -36,7 +36,7 @@ function toProtocolRole(role: ContextMessage['role']): MessageRole { return role as MessageRole; } -function mapContentPart(part: ContextMessage['content'][number]): MessageContent { +function mapContentPart(part: ContextMessage['content'][number]): MessageContent | undefined { switch (part.type) { case 'text': return { type: 'text', text: part.text }; @@ -59,13 +59,22 @@ function mapContentPart(part: ContextMessage['content'][number]): MessageContent ? { type: 'video', source: { kind: 'file', file_id: ref.fileId } } : { type: 'video', source: { kind: 'url', url: part.videoUrl.url, id: part.videoUrl.id } }; } + case 'openai_compaction': + return undefined; } } +function mapContentParts(parts: ContextMessage['content']): MessageContent[] { + return parts.flatMap((part) => { + const mapped = mapContentPart(part); + return mapped === undefined ? [] : [mapped]; + }); +} + function buildProtocolContent(msg: ContextMessage): MessageContent[] { if (msg.role === 'tool') { if (msg.toolCallId === undefined) { - return msg.content.map((p) => mapContentPart(p)); + return mapContentParts(msg.content); } const hasMediaPart = msg.content.some( (p) => p.type === 'image_url' || p.type === 'video_url' || p.type === 'audio_url', @@ -89,7 +98,7 @@ function buildProtocolContent(msg: ContextMessage): MessageContent[] { return [part]; } - const base = msg.content.map((p) => mapContentPart(p)); + const base = mapContentParts(msg.content); if (msg.role === 'assistant' && msg.toolCalls.length > 0) { for (const call of msg.toolCalls) { diff --git a/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts b/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts index 4c87fedd0b9..e3e871f4aa1 100644 --- a/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts +++ b/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts @@ -74,6 +74,7 @@ const OVERFLOW_CONTEXT_SAFETY_RATIO = 0.85; const OVERFLOW_STATUS_RECOVERY_RATIO = 0.5; const MAX_COMPACTION_OVERFLOW_SHRINK_ATTEMPTS = 3; const COMPACTION_OVERFLOW_SHRINK_RATIOS = [0.7, 0.5, 0.35] as const; +const REMOTE_COMPACTION_SUMMARY = '[OpenAI server compaction checkpoint]'; const EMPTY_TOOL_PARAMETERS: Record = { type: 'object', properties: {}, @@ -115,6 +116,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull private compactionCountInTurn = 0; private _compacting: ActiveCompaction | null = null; private readonly observedMaxContextTokensByModel = new Map(); + private readonly remoteCompactionUnavailableModels = new Set(); private lastCompactedTokenCount: number | null = null; private consecutiveOverflowCompactions = 0; private activeTurnId: number | undefined; @@ -537,6 +539,71 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull const resolvedModel = this.profile.resolveModelContext(); thinkingEffort = resolvedModel.thinkingLevel; + + const modelAlias = resolvedModel.modelAlias; + const hasCustomInstruction = (data.instruction?.trim().length ?? 0) > 0; + // Provider-owned compaction has no portable custom-instruction field. + if (!hasCustomInstruction && !this.remoteCompactionUnavailableModels.has(modelAlias)) { + try { + const remote = await this.llmRequester.compact?.( + { + messages: stripDynamicToolContext(originalHistory), + source: { + type: 'operation', + turnId: active.originTurnId, + requestKind: 'remote_compaction', + }, + }, + signal, + ); + if (remote !== undefined) { + if (!historySafeToCompact(this.context.get(), originalHistory)) { + this.cancelActive(active); + throw compactionCancelledReason(active); + } + const result = this.context.applyCompaction({ + summary: REMOTE_COMPACTION_SUMMARY, + contextSummary: REMOTE_COMPACTION_SUMMARY, + replacementMessages: remote.messages as readonly ContextMessage[], + compactedCount: originalHistory.length, + tokensBefore, + }); + const properties: CompactionFinishedEvent = { + turn_id: active.originTurnId, + source: data.source, + tokens_before: result.tokensBefore, + tokens_after: result.tokensAfter, + duration_ms: Date.now() - startedAt, + compacted_count: result.compactedCount, + retry_count: 0, + round: 1, + thinking_effort: thinkingEffort, + ...usageTelemetry(remote.usage ?? null), + }; + this.telemetry.track2('compaction_finished', properties); + return result; + } + this.remoteCompactionUnavailableModels.add(modelAlias); + } catch (error) { + if (isAbortError(error)) throw error; + if ( + isError2(error) && + (error.code === ErrorCodes.AUTH_LOGIN_REQUIRED || + error.code === ErrorCodes.PROVIDER_AUTH_ERROR) + ) { + throw error; + } + const status = findAPIStatusError(error)?.statusCode; + if (status === 400 || status === 404 || status === 405 || status === 501) { + this.remoteCompactionUnavailableModels.add(modelAlias); + } + this.log.warn('remote compaction unavailable; falling back to local summary', { + model: modelAlias, + status, + }); + } + } + const maxContextTokens = resolvedModel.modelCapabilities.max_context_tokens; const defaultCompactionCap = maxContextTokens > 0 diff --git a/packages/agent-core-v2/src/agent/llmRequester/llmRequester.ts b/packages/agent-core-v2/src/agent/llmRequester/llmRequester.ts index a7a9f694d15..169ebceec17 100644 --- a/packages/agent-core-v2/src/agent/llmRequester/llmRequester.ts +++ b/packages/agent-core-v2/src/agent/llmRequester/llmRequester.ts @@ -1,5 +1,9 @@ import { createDecorator } from '#/_base/di/instantiation'; -import type { FinishReason, ThinkingEffort } from '#/kosong/contract/provider'; +import type { + FinishReason, + ProviderCompactionResult, + ThinkingEffort, +} from '#/kosong/contract/provider'; import type { Message, StreamedMessagePart } from '#/kosong/contract/message'; import type { Tool } from '#/kosong/contract/tool'; import type { TokenUsage } from '#/kosong/contract/usage'; @@ -70,6 +74,11 @@ export interface IAgentLLMRequesterService { onPart?: AgentLLMRequestPartHandler, signal?: AbortSignal, ): AgentLLMRequestTask; + + compact?( + overrides?: AgentLLMRequestOverrides, + signal?: AbortSignal, + ): Promise; } export const IAgentLLMRequesterService = createDecorator( diff --git a/packages/agent-core-v2/src/agent/llmRequester/llmRequesterService.ts b/packages/agent-core-v2/src/agent/llmRequester/llmRequesterService.ts index 3543e437e1d..f5487978312 100644 --- a/packages/agent-core-v2/src/agent/llmRequester/llmRequesterService.ts +++ b/packages/agent-core-v2/src/agent/llmRequester/llmRequesterService.ts @@ -55,7 +55,10 @@ import { isRetryableGenerateError, } from '#/kosong/contract/errors'; import { type Message } from '#/kosong/contract/message'; -import { type ThinkingEffort } from '#/kosong/contract/provider'; +import type { + ProviderCompactionResult, + ThinkingEffort, +} from '#/kosong/contract/provider'; import { type Tool } from '#/kosong/contract/tool'; import { emptyUsage, inputTotal, type TokenUsage } from '#/kosong/contract/usage'; import { ILogService, type LogContext } from '#/_base/log/log'; @@ -194,6 +197,23 @@ export class AgentLLMRequesterService implements IAgentLLMRequesterService { }; } + async compact( + overrides: AgentLLMRequestOverrides = {}, + signal?: AbortSignal, + ): Promise { + const request = this.resolveRequest(overrides); + const input = { + systemPrompt: request.systemPrompt, + tools: request.tools, + messages: request.messages, + }; + const result = await request.requester.compact?.(input, signal, request.params); + if (result?.usage !== undefined) { + this.usage.record(request.modelAlias, result.usage, request.source); + } + return result; + } + private async requestWithTrace( trace: MutableLLMRequestTrace, overrides: AgentLLMRequestOverrides, diff --git a/packages/agent-core-v2/src/agent/loop/loopService.ts b/packages/agent-core-v2/src/agent/loop/loopService.ts index 3b16df8ca6b..02b95c688f7 100644 --- a/packages/agent-core-v2/src/agent/loop/loopService.ts +++ b/packages/agent-core-v2/src/agent/loop/loopService.ts @@ -948,6 +948,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService { case 'image_url': case 'audio_url': case 'video_url': + case 'openai_compaction': return; case 'function': { onResponseEvent(); diff --git a/packages/agent-core-v2/src/agent/mcp/output.ts b/packages/agent-core-v2/src/agent/mcp/output.ts index 137930e5403..65fdfef93ae 100644 --- a/packages/agent-core-v2/src/agent/mcp/output.ts +++ b/packages/agent-core-v2/src/agent/mcp/output.ts @@ -246,7 +246,11 @@ function applyBinaryPartCap(parts: readonly ContentPart[]): { const out: ContentPart[] = []; for (const part of parts) { - if (part.type === 'text' || part.type === 'think') { + if ( + part.type === 'text' || + part.type === 'think' || + part.type === 'openai_compaction' + ) { out.push(part); continue; } diff --git a/packages/agent-core-v2/src/kosong/contract/message.ts b/packages/agent-core-v2/src/kosong/contract/message.ts index 6ea2ff05e49..c1366006213 100644 --- a/packages/agent-core-v2/src/kosong/contract/message.ts +++ b/packages/agent-core-v2/src/kosong/contract/message.ts @@ -40,7 +40,25 @@ export interface VideoURLPart { videoUrl: { url: string; id?: string | undefined }; } -export type ContentPart = TextPart | ThinkPart | ImageURLPart | AudioURLPart | VideoURLPart; +/** Opaque state emitted by and replayed only to the OpenAI Responses API. */ +export interface OpenAICompactionPart { + type: 'openai_compaction'; + encryptedContent: string; + id?: string | undefined; + /** Origin identity used to prevent replaying opaque state to another backend. */ + source?: { + readonly model: string; + readonly baseUrl?: string | undefined; + } | undefined; +} + +export type ContentPart = + | TextPart + | ThinkPart + | ImageURLPart + | AudioURLPart + | VideoURLPart + | OpenAICompactionPart; export interface ToolCall { type: 'function'; @@ -72,7 +90,22 @@ export interface Message { export function isContentPart(part: StreamedMessagePart): part is ContentPart { const t = part.type; return ( - t === 'text' || t === 'think' || t === 'image_url' || t === 'audio_url' || t === 'video_url' + t === 'text' || + t === 'think' || + t === 'image_url' || + t === 'audio_url' || + t === 'video_url' || + t === 'openai_compaction' + ); +} + +/** True when an assistant turn contains only provider-owned opaque state. */ +export function isOpaqueAssistantMessage(message: Message): boolean { + return ( + message.role === 'assistant' && + message.toolCalls.length === 0 && + message.content.length > 0 && + message.content.every((part) => part.type === 'openai_compaction') ); } diff --git a/packages/agent-core-v2/src/kosong/contract/provider.ts b/packages/agent-core-v2/src/kosong/contract/provider.ts index 753bd808920..17828b67dd1 100644 --- a/packages/agent-core-v2/src/kosong/contract/provider.ts +++ b/packages/agent-core-v2/src/kosong/contract/provider.ts @@ -159,6 +159,16 @@ export interface GenerateOptions { onTraceId?: (traceId: string | null) => void; } +/** + * Canonical replacement window returned by a provider-owned compaction API. + * The caller must install `messages` atomically and replay them unchanged. + */ +export interface ProviderCompactionResult { + readonly messages: readonly Message[]; + readonly usage?: TokenUsage; + readonly id?: string; +} + /** * A constructed, immutable wire adapter. All construction-time variation * (endpoint, credentials, dialect hooks) is baked in by the factory; all @@ -177,5 +187,15 @@ export interface ChatProvider { history: Message[], options?: GenerateOptions, ): Promise; + /** + * Compact a conversation using the provider's native protocol. Presence is + * the capability signal; providers without a native endpoint omit it. + */ + compact?( + systemPrompt: string, + tools: Tool[], + history: Message[], + options?: GenerateOptions, + ): Promise; uploadVideo?(input: string | VideoUploadInput, options?: GenerateOptions): Promise; } diff --git a/packages/agent-core-v2/src/kosong/contract/tokens.ts b/packages/agent-core-v2/src/kosong/contract/tokens.ts index 9e17577691a..a0621abc4d3 100644 --- a/packages/agent-core-v2/src/kosong/contract/tokens.ts +++ b/packages/agent-core-v2/src/kosong/contract/tokens.ts @@ -82,6 +82,8 @@ export function estimateTokensForContentPart(part: ContentPart): number { case 'audio_url': case 'video_url': return MEDIA_TOKEN_ESTIMATE; + case 'openai_compaction': + return estimateTokens(part.encryptedContent); default: { const exhaustive: never = part; void exhaustive; diff --git a/packages/agent-core-v2/src/kosong/model/modelRequester.ts b/packages/agent-core-v2/src/kosong/model/modelRequester.ts index b4b2de2a5f0..07f23d7be2b 100644 --- a/packages/agent-core-v2/src/kosong/model/modelRequester.ts +++ b/packages/agent-core-v2/src/kosong/model/modelRequester.ts @@ -14,6 +14,7 @@ import type { Message, StreamedMessagePart, VideoURLPart } from '#/kosong/contract/message'; import type { FinishReason, + ProviderCompactionResult, ResponseFormat, SamplingOptions, ThinkingEffort, @@ -102,6 +103,12 @@ export interface ModelRequester { params?: ModelRequestParams, ): AsyncIterable; + compact?( + input: ModelRequestInput, + signal?: AbortSignal, + params?: ModelRequestParams, + ): Promise; + uploadVideo?( input: string | VideoUploadInput, options?: { readonly signal?: AbortSignal }, diff --git a/packages/agent-core-v2/src/kosong/model/modelRequesterImpl.ts b/packages/agent-core-v2/src/kosong/model/modelRequesterImpl.ts index 6e7d52e2600..afa33cb1797 100644 --- a/packages/agent-core-v2/src/kosong/model/modelRequesterImpl.ts +++ b/packages/agent-core-v2/src/kosong/model/modelRequesterImpl.ts @@ -28,6 +28,7 @@ import { generate, type GenerateResult } from '#/kosong/contract/generate'; import type { ChatProvider, GenerateOptions, + ProviderCompactionResult, ProviderRequestAuth, StreamDecodeStats, VideoUploadInput, @@ -79,6 +80,29 @@ export class ModelRequesterImpl implements ModelRequester { return queue; } + async compact( + input: ModelRequestInput, + signal?: AbortSignal, + params?: ModelRequestParams, + ): Promise { + signal?.throwIfAborted(); + const provider = this.resolveChatProvider(); + if (provider.compact === undefined) return undefined; + const compact = provider.compact.bind(provider); + try { + return await this.runWithAuthRefresh((auth) => + compact(input.systemPrompt, [...input.tools], [...input.messages], { + signal, + auth, + cacheKey: params?.cacheKey, + }), + ); + } catch (error) { + if (isAbortError(error) || signal?.aborted === true) throw error; + throw translateProviderError(error); + } + } + async uploadVideo( input: string | VideoUploadInput, options?: { readonly signal?: AbortSignal }, diff --git a/packages/agent-core-v2/src/kosong/provider/bases/google-genai/google-genai.ts b/packages/agent-core-v2/src/kosong/provider/bases/google-genai/google-genai.ts index 07819eed18b..5ce40b61ea8 100644 --- a/packages/agent-core-v2/src/kosong/provider/bases/google-genai/google-genai.ts +++ b/packages/agent-core-v2/src/kosong/provider/bases/google-genai/google-genai.ts @@ -22,7 +22,10 @@ import { normalizeAPIStatusError, } from '#/kosong/contract/errors'; import type { Message, StreamedMessagePart, ThinkPart, ToolCall } from '#/kosong/contract/message'; -import { isToolDeclarationOnlyMessage } from '#/kosong/contract/message'; +import { + isOpaqueAssistantMessage, + isToolDeclarationOnlyMessage, +} from '#/kosong/contract/message'; import type { ChatProvider, FinishReason, @@ -258,6 +261,8 @@ function messageToGoogleGenAI(message: Message): GoogleContent { case 'video_url': parts.push(convertMediaUrl(part.videoUrl.url, 'video/mp4')); break; + case 'openai_compaction': + break; } } @@ -323,6 +328,8 @@ function toolMessageToFunctionResponseParts( break; case 'think': break; + case 'openai_compaction': + break; } } @@ -350,6 +357,10 @@ export function messagesToGoogleGenAIContents(messages: Message[]): GoogleConten i += 1; continue; } + if (isOpaqueAssistantMessage(message)) { + i += 1; + continue; + } if (message.role === 'system') { const text = message.content diff --git a/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-common.ts b/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-common.ts index 68eed46d317..2ac17fceb1d 100644 --- a/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-common.ts +++ b/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-common.ts @@ -73,6 +73,9 @@ export function convertContentPart(part: ContentPart): OpenAIContentPart | null ? { url: part.videoUrl.url } : { url: part.videoUrl.url, id: part.videoUrl.id }, }; + case 'openai_compaction': + // Only the Responses API understands this opaque item. + return null; default: throw new Error(`Unknown content part type: ${(part as ContentPart).type}`); } diff --git a/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-legacy.ts b/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-legacy.ts index 1afe600788a..35828d58342 100644 --- a/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-legacy.ts +++ b/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-legacy.ts @@ -34,7 +34,10 @@ import type { ToolCall, VideoURLPart, } from '#/kosong/contract/message'; -import { isToolDeclarationOnlyMessage } from '#/kosong/contract/message'; +import { + isOpaqueAssistantMessage, + isToolDeclarationOnlyMessage, +} from '#/kosong/contract/message'; import type { ChatProvider, FinishReason, @@ -349,6 +352,7 @@ function convertHistoryMessages( for (const msg of history) { if (isToolDeclarationOnlyMessage(msg)) continue; + if (isOpaqueAssistantMessage(msg)) continue; if (msg.role !== 'tool') { appendToolResultMediaMessage(messages, pendingToolResultMedia); } @@ -619,6 +623,7 @@ export class OpenAILegacyChatProvider implements ChatProvider { // Trait mode: the tool-declaration-only skip and the tool-result media // extraction are handed over to the trait wholesale. for (const msg of normalizedHistory) { + if (isOpaqueAssistantMessage(msg)) continue; const converted = convertMessage(msg, reasoningKey, null, preserveThinking, false); const shaped = convertMessageHook(msg, converted); if (shaped !== null) { diff --git a/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-responses.ts b/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-responses.ts index 13e2b866306..037cec8683c 100644 --- a/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-responses.ts +++ b/packages/agent-core-v2/src/kosong/provider/bases/openai/openai-responses.ts @@ -28,6 +28,7 @@ import type { ChatProvider, FinishReason, GenerateOptions, + ProviderCompactionResult, ProviderRequestAuth, ResponseFormat, StreamedMessage, @@ -107,6 +108,11 @@ type ResponseOutputItemView = encryptedContent?: string; summary: RawObject[]; } + | { + type: 'compaction'; + id?: string; + encryptedContent: string; + } | { type: 'other'; }; @@ -204,6 +210,14 @@ function readResponseOutputItem(value: unknown, context: string): ResponseOutput }; } + if (type === 'compaction') { + return { + type, + id: readStringField(item, 'id'), + encryptedContent: requireStringField(item, 'encrypted_content', context), + }; + } + return { type: 'other' }; } @@ -370,6 +384,26 @@ interface ResponseInputItem { [key: string]: unknown; } +interface CompactionSource { + readonly model: string; + readonly baseUrl?: string | undefined; +} + +function normalizeBaseUrl(value: string | undefined): string | undefined { + return value?.replace(/\/+$/, ''); +} + +function canReplayCompaction( + source: CompactionSource | undefined, + target: CompactionSource, +): boolean { + return ( + source === undefined || + (source.model === target.model && + normalizeBaseUrl(source.baseUrl) === normalizeBaseUrl(target.baseUrl)) + ); +} + interface ResponseToolParam { type: string; name: string; @@ -422,6 +456,9 @@ function contentPartsToInputItems(parts: ContentPart[]): unknown[] { break; case 'think': break; + case 'openai_compaction': + // Emitted as a top-level item by convertMessage. + break; } } return items; @@ -459,6 +496,9 @@ function messageContentToFunctionOutputItems(content: ContentPart[]): unknown[] break; case 'think': break; + case 'openai_compaction': + // Compaction items cannot appear inside function outputs. + break; } } return items; @@ -516,6 +556,7 @@ function convertMessage( message: Message, modelName: string, toolMessageConversion: ToolMessageConversion, + compactionTarget: CompactionSource, ): ResponseInputItem[] { let role: string = message.role; if (usesOpenAIResponsesDeveloperRole(modelName) && role === 'system') { @@ -589,6 +630,16 @@ function convertMessage( type: 'reasoning', encrypted_content: encryptedValue, }); + } else if (part.type === 'openai_compaction') { + flushPendingParts(); + if (canReplayCompaction(part.source, compactionTarget)) { + result.push({ + type: 'compaction', + encrypted_content: part.encryptedContent, + id: part.id, + }); + } + i += 1; } else { pendingParts.push(part); i += 1; @@ -624,6 +675,7 @@ function convertHistoryMessages( history: readonly Message[], modelName: string, toolMessageConversion: ToolMessageConversion, + compactionTarget: CompactionSource, ): unknown[] { const input: unknown[] = []; const pendingToolResultMedia: unknown[] = []; @@ -643,7 +695,7 @@ function convertHistoryMessages( if (msg.role !== 'tool') { flushPendingMedia(); } - input.push(...convertMessage(msg, modelName, toolMessageConversion)); + input.push(...convertMessage(msg, modelName, toolMessageConversion, compactionTarget)); if (msg.role === 'tool' && toolMessageConversion === 'extract_text') { pendingToolResultMedia.push( ...messageContentToFunctionOutputItems(msg.content.filter(isMediaPart)), @@ -655,6 +707,108 @@ function convertHistoryMessages( return input; } +function convertCompactedResponse( + response: unknown, + source: CompactionSource, +): ProviderCompactionResult { + const raw = asRawObject(response); + if (raw === null) { + throw new ChatProviderError('OpenAI Responses compact returned an invalid response.'); + } + const output = readObjectArrayField(raw, 'output'); + if (output === undefined || output.length === 0) { + throw new ChatProviderError('OpenAI Responses compact returned an empty replacement window.'); + } + + const messages: Message[] = []; + let compactionCount = 0; + for (const item of output) { + const type = readStringField(item, 'type'); + if (type === 'message') { + const roleValue = readStringField(item, 'role'); + const role = + roleValue === 'developer' + ? 'system' + : roleValue === 'user' || roleValue === 'system' + ? roleValue + : undefined; + const content = readObjectArrayField(item, 'content'); + if (role === undefined || content === undefined) { + throw new ChatProviderError( + 'OpenAI Responses compact returned an unsupported message item.', + ); + } + const parts: ContentPart[] = content.map((part) => { + const partType = readStringField(part, 'type'); + if (partType === 'input_text' || partType === 'output_text') { + const text = readStringField(part, 'text'); + if (text !== undefined) return { type: 'text', text }; + } + if (partType === 'input_image') { + const imageUrl = readStringField(part, 'image_url'); + if (imageUrl !== undefined) return { type: 'image_url', imageUrl: { url: imageUrl } }; + } + throw new ChatProviderError( + `OpenAI Responses compact returned unsupported message content: ${partType ?? 'unknown'}.`, + ); + }); + messages.push({ role, content: parts, toolCalls: [] }); + continue; + } + if (type === 'compaction') { + const encryptedContent = readStringField(item, 'encrypted_content'); + if (encryptedContent === undefined) { + throw new ChatProviderError( + 'OpenAI Responses compact returned a compaction item without encrypted content.', + ); + } + compactionCount += 1; + messages.push({ + role: 'assistant', + content: [{ + type: 'openai_compaction', + encryptedContent, + id: readStringField(item, 'id'), + source, + }], + toolCalls: [], + }); + continue; + } + throw new ChatProviderError( + `OpenAI Responses compact returned an unsupported replacement item: ${type ?? 'unknown'}.`, + ); + } + if (compactionCount !== 1 || messages.at(-1)?.content[0]?.type !== 'openai_compaction') { + throw new ChatProviderError( + 'OpenAI Responses compact must return user messages followed by one compaction item.', + ); + } + + const usageRaw = readObjectField(raw, 'usage'); + let usage: TokenUsage | undefined; + if (usageRaw !== undefined && usageRaw !== null) { + const inputTokens = readNumberField(usageRaw, 'input_tokens') ?? 0; + const outputTokens = readNumberField(usageRaw, 'output_tokens') ?? 0; + const details = readObjectField(usageRaw, 'input_tokens_details'); + const cached = + details === undefined || details === null + ? 0 + : (readNumberField(details, 'cached_tokens') ?? 0); + usage = { + inputOther: Math.max(0, inputTokens - cached), + output: outputTokens, + inputCacheRead: cached, + inputCacheCreation: 0, + }; + } + return { + messages, + usage, + id: readStringField(raw, 'id'), + }; +} + export class OpenAIResponsesStreamedMessage implements StreamedMessage { private _id: string | null = null; private _usage: TokenUsage | null = null; @@ -662,7 +816,11 @@ export class OpenAIResponsesStreamedMessage implements StreamedMessage { private _rawFinishReason: string | null = null; private readonly _iter: AsyncGenerator; - constructor(response: unknown, isStream: boolean) { + constructor( + response: unknown, + isStream: boolean, + private readonly compactionSource?: CompactionSource, + ) { if (isStream) { this._iter = this._convertStreamResponse(response as AsyncIterable); } else { @@ -766,6 +924,13 @@ export class OpenAIResponsesStreamedMessage implements StreamedMessage { } yield thinkPart; } + } else if (outputItem.type === 'compaction') { + yield { + type: 'openai_compaction', + encryptedContent: outputItem.encryptedContent, + id: outputItem.id, + ...(this.compactionSource === undefined ? {} : { source: this.compactionSource }), + }; } } } @@ -906,6 +1071,13 @@ export class OpenAIResponsesStreamedMessage implements StreamedMessage { (thinkPart as { encrypted: string }).encrypted = item.encryptedContent; } yield thinkPart; + } else if (item.type === 'compaction') { + yield { + type: 'openai_compaction', + encryptedContent: item.encryptedContent, + id: item.id, + ...(this.compactionSource === undefined ? {} : { source: this.compactionSource }), + }; } else if (item.type === 'function_call' && typeof item.arguments === 'string') { const streamIndex = responseStreamIndex(item.itemId, outputIndex); yield* yieldFinalArgumentsSuffix(streamIndex, item.arguments, type); @@ -1053,8 +1225,17 @@ export class OpenAIResponsesChatProvider implements ChatProvider { history, OPENAI_RESPONSES_TOOL_CALL_ID_POLICY, ); + const compactionSource = { + model: this._model, + baseUrl: options?.auth?.baseUrl ?? this._baseUrl, + }; input.push( - ...convertHistoryMessages(normalizedHistory, this._model, this._toolMessageConversion), + ...convertHistoryMessages( + normalizedHistory, + this._model, + this._toolMessageConversion, + compactionSource, + ), ); let kwargs: Record = { ...this._generationKwargs }; @@ -1149,7 +1330,55 @@ export class OpenAIResponsesChatProvider implements ChatProvider { create(params: unknown, opts?: unknown): Promise; } ).create(createParams, options?.signal ? { signal: options.signal } : undefined); - return new OpenAIResponsesStreamedMessage(response, this._stream); + return new OpenAIResponsesStreamedMessage(response, this._stream, compactionSource); + } catch (error: unknown) { + throw convertOpenAIError(error); + } + } + + async compact( + systemPrompt: string, + _tools: Tool[], + history: Message[], + options?: GenerateOptions, + ): Promise { + const normalizedHistory = normalizeToolCallIdsForProvider( + history, + OPENAI_RESPONSES_TOOL_CALL_ID_POLICY, + ); + const params: Record = { + model: this._model, + input: convertHistoryMessages( + normalizedHistory, + this._model, + this._toolMessageConversion, + { + model: this._model, + baseUrl: options?.auth?.baseUrl ?? this._baseUrl, + }, + ), + }; + if (systemPrompt) params['instructions'] = systemPrompt; + if (options?.cacheKey !== undefined) params['prompt_cache_key'] = options.cacheKey; + + try { + const client = this._createClient(options?.auth); + const compact = (client as { responses?: { compact?: unknown } }).responses?.compact; + if (typeof compact !== 'function') { + throw new TypeError( + 'OpenAI SDK version does not support Responses compaction. Upgrade to a version with responses.compact.', + ); + } + options?.onRequestSent?.(); + const response = await ( + client.responses as { + compact(params: unknown, opts?: unknown): Promise; + } + ).compact(params, options?.signal ? { signal: options.signal } : undefined); + return convertCompactedResponse(response, { + model: this._model, + baseUrl: options?.auth?.baseUrl ?? this._baseUrl, + }); } catch (error: unknown) { throw convertOpenAIError(error); } diff --git a/packages/agent-core-v2/test/agent/contextMemory/context.test.ts b/packages/agent-core-v2/test/agent/contextMemory/context.test.ts index 2189dac79b4..32366bd1815 100644 --- a/packages/agent-core-v2/test/agent/contextMemory/context.test.ts +++ b/packages/agent-core-v2/test/agent/contextMemory/context.test.ts @@ -63,6 +63,43 @@ describe('Agent context', () => { expect(ctx.project().some((message) => 'origin' in message)).toBe(false); }); + it('installs and replays a provider-owned compaction window atomically', () => { + context.append({ + role: 'user', + content: [{ type: 'text', text: 'old request' }], + toolCalls: [], + }); + const replacement: ContextMessage[] = [ + { + role: 'user', + content: [{ type: 'text', text: 'old request' }], + toolCalls: [], + }, + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'opaque-state', + id: 'cmp_1', + source: { model: 'gpt-5.4', baseUrl: 'https://api.openai.com/v1' }, + }, + ], + toolCalls: [], + }, + ]; + + const result = context.applyCompaction({ + summary: '[OpenAI server compaction checkpoint]', + replacementMessages: replacement, + compactedCount: 1, + tokensBefore: 50, + }); + + expect(context.get()).toEqual(replacement); + expect(result.keptUserMessageCount).toBe(1); + }); + it('renders tool error and empty-output status as model-visible text', () => { context.append( { diff --git a/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts b/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts index 68fe0d460c2..26c1381cfc7 100644 --- a/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts +++ b/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts @@ -36,7 +36,7 @@ import { MASTER_ENV } from '#/app/flag/flagService'; import { estimateTokensForMessages } from '#/kosong/contract/tokens'; import { recordingTelemetry, type TelemetryRecord } from '../../app/telemetry/stubs'; import type { TestAgentContext, TestAgentOptions, TestAgentServiceOverride } from '../../harness'; -import { appServices, createCommandRunner, execEnvServices, hostEnvironmentServices, sessionServices, testAgent } from '../../harness'; +import { agentService, appServices, createCommandRunner, execEnvServices, hostEnvironmentServices, sessionServices, testAgent } from '../../harness'; import { IAgentFullCompactionService, IModelOAuthTokens, @@ -50,6 +50,7 @@ import { } from '#/index'; import { IAgentLoopService } from '#/agent/loop/loop'; import { IAgentContextSizeService } from '#/agent/contextSize/contextSize'; +import { IAgentLLMRequesterService } from '#/agent/llmRequester/llmRequester'; import { IAgentGoalService } from '#/agent/goal/goal'; import { IAgentTelemetryContextService } from '#/app/telemetry/agentTelemetryContext'; import { HostFileSystem } from '#/os/backends/node-local/hostFsService'; @@ -316,6 +317,68 @@ describe('FullCompaction', () => { await ctx.expectResumeMatches(); }); + it('prefers provider-owned compaction and installs its replacement window', async () => { + let compactCalls = 0; + const requester: IAgentLLMRequesterService = { + _serviceBrand: undefined, + prepareTurnConfig: () => undefined, + request: async () => { + throw new Error('local summarizer should not run'); + }, + start: () => { + throw new Error('local summarizer should not run'); + }, + compact: async () => { + compactCalls += 1; + return { + messages: [ + { + role: 'user', + content: [{ type: 'text', text: 'old request' }], + toolCalls: [], + }, + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'opaque-state', + id: 'cmp_1', + }, + ], + toolCalls: [], + }, + ], + usage: { + inputOther: 80, + output: 20, + inputCacheRead: 0, + inputCacheCreation: 0, + }, + }; + }, + }; + const ctx = testAgent(agentService(IAgentLLMRequesterService, requester)); + ctx.configure({ + provider: CATALOGUED_PROVIDER, + modelCapabilities: CATALOGUED_MODEL_CAPABILITIES, + }); + ctx.appendExchange(1, 'old request', 'old answer', 20); + const completed = ctx.once('compaction.completed'); + + await ctx.rpc.beginCompaction({}); + await completed; + + expect(compactCalls).toBe(1); + expect( + ctx.context.get().some((message) => + message.content.some((part) => part.type === 'openai_compaction'), + ), + ).toBe(true); + expect(ctx.llmCalls).toHaveLength(0); + await ctx.expectResumeMatches(); + }); + it('refreshes the active profile system prompt after compaction without resetting active tools', async () => { const homeDir = mkdtempSync(join(tmpdir(), 'kimi-compact-refresh-home-')); const workDir = mkdtempSync(join(tmpdir(), 'kimi-compact-refresh-work-')); diff --git a/packages/agent-core-v2/test/agent/llmRequester/llmRequesterService.test.ts b/packages/agent-core-v2/test/agent/llmRequester/llmRequesterService.test.ts index 6afd2e982c0..6b118cd8716 100644 --- a/packages/agent-core-v2/test/agent/llmRequester/llmRequesterService.test.ts +++ b/packages/agent-core-v2/test/agent/llmRequester/llmRequesterService.test.ts @@ -171,7 +171,13 @@ function createService( get: () => ({ size: 0, measured: 0, estimated: 0 }), measured: () => undefined, }; - const usage = { record: () => undefined, status: () => ({}) }; + const usageRecords: Array<{ model: string; usage: ReturnType }> = []; + const usage = { + record: (model: string, recordedUsage: ReturnType) => { + usageRecords.push({ model, usage: recordedUsage }); + }, + status: () => ({}), + }; const context = { get: () => history }; const tools = { list: () => [] }; const config: Partial = { @@ -242,9 +248,47 @@ function createService( records, events, telemetryRecords, + usageRecords, }; } +describe('AgentLLMRequesterService native compaction', () => { + it('records provider-reported compaction usage', async () => { + const requester = createRequester({ value: 0 }, null); + requester.compact = async () => ({ + messages: [ + { + role: 'assistant', + content: [{ type: 'openai_compaction', encryptedContent: 'opaque-state' }], + toolCalls: [], + }, + ], + usage: { + inputOther: 80, + output: 20, + inputCacheRead: 10, + inputCacheCreation: 0, + }, + }); + const { service, usageRecords } = createService(requester, undefined); + + const result = await service.compact?.(); + + expect(result?.messages[0]?.content[0]?.type).toBe('openai_compaction'); + expect(usageRecords).toEqual([ + { + model: 'm', + usage: { + inputOther: 80, + output: 20, + inputCacheRead: 10, + inputCacheCreation: 0, + }, + }, + ]); + }); +}); + describe('AgentLLMRequesterService Anthropic effort diagnostics', () => { it('warns and sends when the effort is not listed by the model', async () => { const calls = { value: 0 }; diff --git a/packages/agent-core-v2/test/kosong/provider/composition.test.ts b/packages/agent-core-v2/test/kosong/provider/composition.test.ts index 331281ef5ab..2fcb8b6fd1a 100644 --- a/packages/agent-core-v2/test/kosong/provider/composition.test.ts +++ b/packages/agent-core-v2/test/kosong/provider/composition.test.ts @@ -57,7 +57,10 @@ import { import '#/kosong/provider/bases/google-genai/index'; import { GoogleGenAIChatProvider } from '#/kosong/provider/bases/google-genai/google-genai'; import '#/kosong/provider/bases/openai/index'; -import { OpenAIResponsesChatProvider } from '#/kosong/provider/bases/openai/openai-responses'; +import { + OpenAIResponsesChatProvider, + OpenAIResponsesStreamedMessage, +} from '#/kosong/provider/bases/openai/openai-responses'; import { OpenAILegacyChatProvider } from '#/kosong/provider/bases/openai/openai-legacy'; import { ProtocolAdapterRegistry } from '#/kosong/provider/protocolAdapterRegistry'; import { @@ -629,6 +632,7 @@ async function captureAnthropicBody( async function captureGoogleBody( provider: ChatProvider, options?: GenerateOptions, + history: Message[] = PROBE_HISTORY, ): Promise> { let captured: Record | undefined; const client = sdkClient(provider) as { models: { generateContent: unknown } }; @@ -642,7 +646,7 @@ async function captureGoogleBody( modelVersion: 'probe', }); }); - await drain(await provider.generate('', [], PROBE_HISTORY, options)); + await drain(await provider.generate('', [], history, options)); if (captured === undefined) throw new Error('expected models.generateContent to be called'); return captured; } @@ -650,6 +654,7 @@ async function captureGoogleBody( async function captureResponsesBody( provider: ChatProvider, options?: GenerateOptions, + history: Message[] = PROBE_HISTORY, ): Promise> { let captured: Record | undefined; const client = sdkClient(provider) as { responses: { create: unknown } }; @@ -657,7 +662,7 @@ async function captureResponsesBody( captured = params as Record; return Promise.resolve(responsesEventStream()); }); - await drain(await provider.generate('', [], PROBE_HISTORY, options)); + await drain(await provider.generate('', [], history, options)); if (captured === undefined) throw new Error('expected responses.create to be called'); return captured; } @@ -986,6 +991,182 @@ describe('responseFormat wire encoding (per base)', () => { // unreachable now: no channel seeds `text.verbosity` (the per-request // merge in the base still stands, but only per-turn formats reach it). }); + + it('round-trips opaque OpenAI compaction items on the Responses wire', async () => { + const provider = new OpenAIResponsesChatProvider({ model: 'gpt-4.1', apiKey: 'sk-probe' }); + const history: Message[] = [ + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'encrypted-summary', + id: 'cmp_123', + }, + ], + toolCalls: [], + }, + ]; + + const body = await captureResponsesBody(provider, undefined, history); + expect(body['input']).toEqual([ + { + type: 'compaction', + encrypted_content: 'encrypted-summary', + id: 'cmp_123', + }, + ]); + + async function* stream(): AsyncIterable { + yield { + type: 'response.output_item.done', + item: { + type: 'compaction', + id: 'cmp_456', + encrypted_content: 'encrypted-stream-summary', + }, + }; + } + + const parts = []; + for await (const part of new OpenAIResponsesStreamedMessage(stream(), true)) { + parts.push(part); + } + expect(parts).toEqual([ + { + type: 'openai_compaction', + encryptedContent: 'encrypted-stream-summary', + id: 'cmp_456', + }, + ]); + }); + + it('uses the native compact endpoint and returns its canonical replacement window', async () => { + const provider = new OpenAIResponsesChatProvider({ model: 'gpt-5.4', apiKey: 'sk-probe' }); + let captured: Record | undefined; + const client = sdkClient(provider) as { responses: { compact: unknown } }; + client.responses.compact = vi.fn().mockImplementation((params: unknown) => { + captured = params as Record; + return Promise.resolve({ + id: 'cmp_response_1', + output: [ + { + type: 'message', + role: 'user', + content: [{ type: 'input_text', text: 'Keep this request.' }], + }, + { type: 'compaction', id: 'cmp_1', encrypted_content: 'opaque-state' }, + ], + usage: { + input_tokens: 100, + output_tokens: 20, + input_tokens_details: { cached_tokens: 40 }, + }, + }); + }); + + const result = await provider.compact('', [], [ + { role: 'user', content: [{ type: 'text', text: 'Keep this request.' }], toolCalls: [] }, + { role: 'assistant', content: [{ type: 'text', text: 'Old answer.' }], toolCalls: [] }, + ], { cacheKey: 'session-probe' }); + + expect(captured).toEqual({ + model: 'gpt-5.4', + input: [ + { + type: 'message', + role: 'user', + content: [{ type: 'input_text', text: 'Keep this request.' }], + }, + { + type: 'message', + role: 'assistant', + content: [{ type: 'output_text', text: 'Old answer.', annotations: [] }], + }, + ], + prompt_cache_key: 'session-probe', + }); + expect(result).toEqual({ + id: 'cmp_response_1', + messages: [ + { + role: 'user', + content: [{ type: 'text', text: 'Keep this request.' }], + toolCalls: [], + }, + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'opaque-state', + id: 'cmp_1', + source: { model: 'gpt-5.4', baseUrl: 'https://api.openai.com/v1' }, + }, + ], + toolCalls: [], + }, + ], + usage: { + inputOther: 60, + output: 20, + inputCacheRead: 40, + inputCacheCreation: 0, + }, + }); + }); + + it('does not replay opaque compaction state to another Responses model', async () => { + const provider = new OpenAIResponsesChatProvider({ model: 'gpt-4.1', apiKey: 'sk-probe' }); + const body = await captureResponsesBody(provider, undefined, [ + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'other-model-state', + source: { model: 'gpt-5.4', baseUrl: 'https://api.openai.com/v1' }, + }, + ], + toolCalls: [], + }, + ]); + + expect(body['input']).toEqual([]); + }); + + it('drops compaction-only assistant shells when switching provider protocols', async () => { + const history: Message[] = [ + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'encrypted-summary', + id: 'cmp_123', + }, + ], + toolCalls: [], + }, + { role: 'user', content: [{ type: 'text', text: 'Continue' }], toolCalls: [] }, + ]; + + const legacy = new OpenAILegacyChatProvider({ + model: 'gpt-4.1', + apiKey: 'sk-probe', + stream: false, + }); + const legacyBody = await captureOpenAIBody(legacy, undefined, history); + expect(legacyBody['messages']).toEqual([{ role: 'user', content: 'Continue' }]); + + const google = new GoogleGenAIChatProvider({ + model: 'gemini-2.5-flash', + apiKey: 'sk-probe', + stream: false, + }); + const googleBody = await captureGoogleBody(google, undefined, history); + expect(googleBody['contents']).toEqual([{ role: 'user', parts: [{ text: 'Continue' }] }]); + }); }); describe('Anthropic thinking keep (context-management overlay)', () => { diff --git a/packages/agent-core/src/agent/compaction/full.ts b/packages/agent-core/src/agent/compaction/full.ts index 7f0fc234916..e4d292595a8 100644 --- a/packages/agent-core/src/agent/compaction/full.ts +++ b/packages/agent-core/src/agent/compaction/full.ts @@ -61,6 +61,7 @@ export const MAX_COMPACTION_RETRY_ATTEMPTS = 5; const DEFAULT_COMPACTION_MAX_COMPLETION_TOKENS = 128 * 1024; const OVERFLOW_CONTEXT_SAFETY_RATIO = 0.85; const OVERFLOW_STATUS_RECOVERY_RATIO = 0.5; +const REMOTE_COMPACTION_SUMMARY = '[OpenAI server compaction checkpoint]'; class CompactionTruncatedError extends Error { constructor() { @@ -77,6 +78,7 @@ export class FullCompaction { blockedByTurn: boolean; } | null = null; private readonly observedMaxContextTokensByModel = new Map(); + private readonly remoteCompactionUnavailableModels = new Set(); // Token count right after the last successful compaction. While no new // content has been appended (tokenCountWithPending <= this value), the // history is already in its minimal compacted form ([kept user prompts @@ -435,6 +437,85 @@ export class FullCompaction { }), capability, }); + const hasCustomInstruction = (data.instruction?.trim().length ?? 0) > 0; + // The native endpoint does not expose a portable custom-instruction field. + if (!hasCustomInstruction && !this.remoteCompactionUnavailableModels.has(model)) { + try { + const remote = await this.agent.compactProvider( + provider, + this.agent.config.systemPrompt, + [...this.agent.tools.loopTools], + this.agent.context.project(stripDynamicToolContext(originalHistory), { + synthesizeMissing: true, + dropOrphanResults: true, + }), + { signal }, + ); + if (remote !== undefined) { + const newHistory = this.agent.context.history; + for (let i = 0; i < originalHistory.length; i++) { + if (newHistory[i] !== originalHistory[i]) { + this.cancel(); + return undefined; + } + } + if ( + newHistory + .slice(originalHistory.length) + .some((message) => !isRealUserInput(message)) + ) { + this.cancel(); + return undefined; + } + if (remote.usage !== undefined) { + this.agent.usage.record(model, remote.usage); + } + const result = this.agent.context.applyCompaction({ + summary: REMOTE_COMPACTION_SUMMARY, + contextSummary: REMOTE_COMPACTION_SUMMARY, + replacementMessages: remote.messages, + compactedCount: originalHistory.length, + tokensBefore, + }); + this.agent.telemetry.track('compaction_finished', { + source: data.source, + tokens_before: result.tokensBefore, + tokens_after: result.tokensAfter, + duration_ms: Date.now() - startedAt, + compacted_count: result.compactedCount, + retry_count: 0, + round: 1, + thinking_effort: this.agent.config.thinkingEffort, + ...(remote.usage === undefined + ? {} + : { + input_tokens: inputTotal(remote.usage), + output_tokens: remote.usage.output, + }), + }); + this.lastCompactedTokenCount = this.tokenCountWithPending; + return result; + } + this.remoteCompactionUnavailableModels.add(model); + } catch (error) { + if (isAbortError(error)) return undefined; + if ( + isKimiError(error) && + (error.code === ErrorCodes.AUTH_LOGIN_REQUIRED || + error.code === ErrorCodes.PROVIDER_AUTH_ERROR) + ) { + throw error; + } + const status = findAPIStatusError(error)?.statusCode; + if (status === 400 || status === 404 || status === 405 || status === 501) { + this.remoteCompactionUnavailableModels.add(model); + } + this.agent.log.warn('remote compaction unavailable; falling back to local summary', { + model, + status, + }); + } + } const instruction = this.buildInstruction(data.instruction); const delays = retryBackoffDelays(MAX_COMPACTION_RETRY_ATTEMPTS); diff --git a/packages/agent-core/src/agent/compaction/types.ts b/packages/agent-core/src/agent/compaction/types.ts index ff80e5385fa..6f7653715cc 100644 --- a/packages/agent-core/src/agent/compaction/types.ts +++ b/packages/agent-core/src/agent/compaction/types.ts @@ -40,6 +40,8 @@ export interface CompactionResult { * compatibility with older wire records. */ droppedCount?: number; + /** Canonical provider-owned replacement window for native compaction. */ + replacementMessages?: readonly import('@moonshot-ai/kosong').Message[]; } /** @@ -52,7 +54,7 @@ export type CompactionInput = Pick >; diff --git a/packages/agent-core/src/agent/context/index.ts b/packages/agent-core/src/agent/context/index.ts index 52f8c6c4fb5..947c97dd6f1 100644 --- a/packages/agent-core/src/agent/context/index.ts +++ b/packages/agent-core/src/agent/context/index.ts @@ -312,6 +312,39 @@ export class ContextMemory { } applyCompaction(input: CompactionInput): CompactionResult { + if (input.replacementMessages !== undefined) { + const replacementMessages = input.replacementMessages as readonly ContextMessage[]; + const contextSummary = input.contextSummary ?? input.summary; + const result: CompactionResult = { + summary: input.summary, + contextSummary, + compactedCount: input.compactedCount, + tokensBefore: input.tokensBefore, + tokensAfter: input.tokensAfter ?? estimateTokensForMessages(replacementMessages), + keptUserMessageCount: + input.keptUserMessageCount ?? + replacementMessages.filter((message) => message.role === 'user').length, + keptHeadUserMessageCount: input.keptHeadUserMessageCount, + droppedCount: input.droppedCount, + replacementMessages: [...replacementMessages], + }; + this.agent.records.logRecord({ + type: 'context.apply_compaction', + ...result, + }); + this.agent.replayBuilder.patchLast('compaction', { result }); + this._history = [...replacementMessages]; + this.openSteps.clear(); + this.pendingToolResultIds.clear(); + this.deferredMessages = []; + this._tokenCount = result.tokensAfter; + this.tokenCountCoveredMessageCount = this._history.length; + this.agent.microCompaction.reset(); + this.agent.injection.onContextCompacted(); + this.agent.tools.onContextCompacted(); + this.agent.emitStatusUpdated(); + return result; + } // Single derivation point for the post-compaction shape: the kept user // messages (verbatim, within the token budget — the oldest head plus the // most recent tail, with an elision marker between them when the pool diff --git a/packages/agent-core/src/agent/index.ts b/packages/agent-core/src/agent/index.ts index 7786c5eb04a..143100049b0 100644 --- a/packages/agent-core/src/agent/index.ts +++ b/packages/agent-core/src/agent/index.ts @@ -6,7 +6,14 @@ import { ErrorCodes, KimiError, makeErrorPayload } from '#/errors'; import { log } from '#/logging/logger'; import type { Logger } from '#/logging/types'; import type { AgentAPI, AgentEvent, KimiConfig, SDKAgentRPC, UsageStatus } from '#/rpc'; -import { generate, type ChatProvider } from '@moonshot-ai/kosong'; +import { + generate, + type ChatProvider, + type GenerateOptions, + type Message, + type ProviderCompactionResult, + type Tool, +} from '@moonshot-ai/kosong'; import type { EnabledPluginSessionStart, PluginCommandDef } from '#/plugin'; import { expandCommandArguments } from '../plugin/commands'; @@ -319,6 +326,31 @@ export class Agent { }; } + async compactProvider( + provider: ChatProvider, + systemPrompt: string, + tools: Tool[], + history: Message[], + options?: GenerateOptions, + ): Promise { + if (provider.compact === undefined) return undefined; + const compact = provider.compact.bind(provider); + const modelAlias = this.config.modelAlias; + if (options?.auth !== undefined) { + return compact(systemPrompt, tools, history, options); + } + const withAuth = + modelAlias === undefined + ? undefined + : this.modelProvider?.resolveAuth?.(modelAlias, { log: this.log }); + if (withAuth === undefined) { + return compact(systemPrompt, tools, history, options); + } + return withAuth((auth) => + compact(systemPrompt, tools, history, { ...options, auth }), + ); + } + private warnAboutAnthropicThinkingEffort( provider: ChatProvider, modelAlias: string | undefined, diff --git a/packages/agent-core/src/mcp/output.ts b/packages/agent-core/src/mcp/output.ts index 08fe82e9a1d..d5b09ffa6fd 100644 --- a/packages/agent-core/src/mcp/output.ts +++ b/packages/agent-core/src/mcp/output.ts @@ -322,7 +322,11 @@ function applyBinaryPartCap(parts: readonly ContentPart[]): { const out: ContentPart[] = []; for (const part of parts) { - if (part.type === 'text' || part.type === 'think') { + if ( + part.type === 'text' || + part.type === 'think' || + part.type === 'openai_compaction' + ) { out.push(part); continue; } diff --git a/packages/agent-core/src/services/message/message.ts b/packages/agent-core/src/services/message/message.ts index 1d25ec626f1..0815e14ad57 100644 --- a/packages/agent-core/src/services/message/message.ts +++ b/packages/agent-core/src/services/message/message.ts @@ -150,7 +150,7 @@ function toProtocolRole(role: ContextMessage['role']): MessageRole { * Translate kosong content parts to SCHEMAS §3 content parts. See header * for the full mapping table. */ -function mapContentPart(part: ContextMessage['content'][number]): MessageContent { +function mapContentPart(part: ContextMessage['content'][number]): MessageContent | undefined { switch (part.type) { case 'text': return { type: 'text', text: part.text }; @@ -177,9 +177,18 @@ function mapContentPart(part: ContextMessage['content'][number]): MessageContent type: 'text', text: `[video:${part.videoUrl.url}]`, }; + case 'openai_compaction': + return undefined; } } +function mapContentParts(parts: ContextMessage['content']): MessageContent[] { + return parts.flatMap((part) => { + const mapped = mapContentPart(part); + return mapped === undefined ? [] : [mapped]; + }); +} + /** * Build the protocol-shaped `Message.content[]` for one ContextMessage. * @@ -199,7 +208,7 @@ function buildProtocolContent(msg: ContextMessage): MessageContent[] { if (msg.toolCallId === undefined) { // Defensive — kosong tool messages always carry toolCallId. If absent, // fall back to text passthrough so we don't lose user-visible content. - return msg.content.map((p) => mapContentPart(p)); + return mapContentParts(msg.content); } const hasMediaPart = msg.content.some( (p) => p.type === 'image_url' || p.type === 'video_url' || p.type === 'audio_url', @@ -222,7 +231,7 @@ function buildProtocolContent(msg: ContextMessage): MessageContent[] { return [part]; } - const base = msg.content.map((p) => mapContentPart(p)); + const base = mapContentParts(msg.content); if (msg.role === 'assistant' && msg.toolCalls.length > 0) { for (const call of msg.toolCalls) { diff --git a/packages/agent-core/src/services/message/transcript.ts b/packages/agent-core/src/services/message/transcript.ts index c3d2e25abaf..8d2fdf536a2 100644 --- a/packages/agent-core/src/services/message/transcript.ts +++ b/packages/agent-core/src/services/message/transcript.ts @@ -243,6 +243,20 @@ export function reduceWireRecords(records: Iterable): { applyLoopEvent(record.event, record.time); break; case 'context.apply_compaction': { + if (record.replacementMessages !== undefined) { + transcript.push({ + message: { + role: 'user', + content: [{ type: 'text', text: record.summary }], + toolCalls: [], + origin: { kind: 'compaction_summary' }, + }, + time: record.time, + }); + foldedLength = record.replacementMessages.length; + resetOpenState(); + break; + } // Mirrors ContextMemory.applyCompaction: the live context becomes the // kept user messages (head + tail, possibly separated by an elision // marker) followed by a user-role summary. The transcript keeps the diff --git a/packages/agent-core/src/utils/tokens.ts b/packages/agent-core/src/utils/tokens.ts index cb4a366627a..95ab0f1ad15 100644 --- a/packages/agent-core/src/utils/tokens.ts +++ b/packages/agent-core/src/utils/tokens.ts @@ -110,6 +110,8 @@ export function estimateTokensForContentPart(part: ContentPart): number { case 'audio_url': case 'video_url': return MEDIA_TOKEN_ESTIMATE; + case 'openai_compaction': + return estimateTokens(part.encryptedContent); default: { // Exhaustiveness guard: a new ContentPart kind must declare its estimate // here rather than silently counting as 0 (the CMP-03 defect). diff --git a/packages/agent-core/test/agent/context.test.ts b/packages/agent-core/test/agent/context.test.ts index 432187410ac..9f4eeb377af 100644 --- a/packages/agent-core/test/agent/context.test.ts +++ b/packages/agent-core/test/agent/context.test.ts @@ -1034,6 +1034,42 @@ describe('Agent context', () => { expect(result.keptUserMessageCount).toBe(1); }); + it('applyCompaction installs and restores a provider-owned replacement window', async () => { + const ctx = testAgent(); + ctx.configure(); + ctx.agent.context.appendUserMessage([{ type: 'text', text: 'old request' }]); + const replacement = [ + { + role: 'user' as const, + content: [{ type: 'text' as const, text: 'old request' }], + toolCalls: [], + }, + { + role: 'assistant' as const, + content: [ + { + type: 'openai_compaction' as const, + encryptedContent: 'opaque-state', + id: 'cmp_1', + source: { model: 'gpt-5.4', baseUrl: 'https://api.openai.com/v1' }, + }, + ], + toolCalls: [], + }, + ]; + + const result = ctx.agent.context.applyCompaction({ + summary: '[OpenAI server compaction checkpoint]', + replacementMessages: replacement, + compactedCount: 1, + tokensBefore: 50, + }); + + expect(ctx.agent.context.history).toEqual(replacement); + expect(result.keptUserMessageCount).toBe(1); + await ctx.expectResumeMatches(); + }); + it('clears context before the next LLM request', async () => { const ctx = testAgent(); ctx.configure(); diff --git a/packages/agent-core/test/services/message-transcript.test.ts b/packages/agent-core/test/services/message-transcript.test.ts index c2655bf9d17..e1fa6e70cbc 100644 --- a/packages/agent-core/test/services/message-transcript.test.ts +++ b/packages/agent-core/test/services/message-transcript.test.ts @@ -214,6 +214,48 @@ describe('reduceWireRecords', () => { expect(foldedLength).toBe(3); }); + it('keeps the readable transcript while tracking a provider-owned replacement window', () => { + const replacementMessages: ContextMessage[] = [ + userMessage('u1'), + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'opaque-state', + source: { model: 'gpt-5.4', baseUrl: 'https://api.openai.com/v1' }, + }, + ], + toolCalls: [], + }, + ]; + const { entries, foldedLength } = reduceWireRecords([ + appendMessage(userMessage('u1')), + ...assistantStep('s1', 'a1'), + { + type: 'context.apply_compaction', + summary: '[OpenAI server compaction checkpoint]', + replacementMessages, + compactedCount: 2, + tokensBefore: 100, + tokensAfter: 20, + } as AgentRecord, + appendMessage(userMessage('u2')), + ]); + + // REST consumers retain the readable pre-compaction conversation and a + // checkpoint marker. The opaque replacement belongs only to live model + // context, whose length is replacementMessages + the subsequently + // appended user message. + expect(entries.map((entry) => textOf(entry.message))).toEqual([ + 'u1', + 'a1', + '[OpenAI server compaction checkpoint]', + 'u2', + ]); + expect(foldedLength).toBe(3); + }); + it('drops a late tool result after compaction closes an open exchange', () => { const { entries, foldedLength } = reduceWireRecords([ appendMessage(userMessage('u1')), diff --git a/packages/kosong/src/message.ts b/packages/kosong/src/message.ts index 12db25c84ed..5d0894cbbc2 100644 --- a/packages/kosong/src/message.ts +++ b/packages/kosong/src/message.ts @@ -28,6 +28,18 @@ export interface VideoURLPart { videoUrl: { url: string; id?: string | undefined }; } +/** Opaque state emitted by and replayed only to the OpenAI Responses API. */ +export interface OpenAICompactionPart { + type: 'openai_compaction'; + encryptedContent: string; + id?: string | undefined; + /** Origin identity used to prevent replaying opaque state to another backend. */ + source?: { + readonly model: string; + readonly baseUrl?: string | undefined; + } | undefined; +} + /** * A single piece of content within a {@link Message}. * @@ -35,7 +47,13 @@ export interface VideoURLPart { * Providers convert these to their native content-block format during * {@link ChatProvider.generate}. */ -export type ContentPart = TextPart | ThinkPart | ImageURLPart | AudioURLPart | VideoURLPart; +export type ContentPart = + | TextPart + | ThinkPart + | ImageURLPart + | AudioURLPart + | VideoURLPart + | OpenAICompactionPart; export interface ToolCall { type: 'function'; @@ -114,11 +132,26 @@ export interface Message { readonly tools?: readonly Tool[] | undefined; } -/** Check if a streamed part is a ContentPart (text, think, image_url, audio_url, video_url). */ +/** Check if a streamed part is a ContentPart. */ export function isContentPart(part: StreamedMessagePart): part is ContentPart { const t = part.type; return ( - t === 'text' || t === 'think' || t === 'image_url' || t === 'audio_url' || t === 'video_url' + t === 'text' || + t === 'think' || + t === 'image_url' || + t === 'audio_url' || + t === 'video_url' || + t === 'openai_compaction' + ); +} + +/** True when an assistant turn contains only provider-owned opaque state. */ +export function isOpaqueAssistantMessage(message: Message): boolean { + return ( + message.role === 'assistant' && + message.toolCalls.length === 0 && + message.content.length > 0 && + message.content.every((part) => part.type === 'openai_compaction') ); } diff --git a/packages/kosong/src/provider.ts b/packages/kosong/src/provider.ts index 3103ea6619f..4326db83371 100644 --- a/packages/kosong/src/provider.ts +++ b/packages/kosong/src/provider.ts @@ -201,6 +201,16 @@ export interface StreamDecodeStats { readonly clientConsumeMs: number; } +/** + * Canonical replacement window returned by a provider-owned compaction API. + * The caller must install `messages` atomically and replay them unchanged. + */ +export interface ProviderCompactionResult { + readonly messages: readonly Message[]; + readonly usage?: TokenUsage; + readonly id?: string; +} + /** * In-memory video bytes for providers that require an uploaded file * reference instead of an inline data URL. @@ -251,6 +261,16 @@ export interface ChatProvider { history: Message[], options?: GenerateOptions, ): Promise; + /** + * Compact a conversation using the provider's native protocol. Presence is + * the capability signal; providers without a native endpoint omit it. + */ + compact?( + systemPrompt: string, + tools: Tool[], + history: Message[], + options?: GenerateOptions, + ): Promise; /** Return a shallow copy of this provider with the given thinking effort. */ withThinking(effort: ThinkingEffort): ChatProvider; /** diff --git a/packages/kosong/src/providers/google-genai.ts b/packages/kosong/src/providers/google-genai.ts index 88a26e6e0ea..ec2e5222f7e 100644 --- a/packages/kosong/src/providers/google-genai.ts +++ b/packages/kosong/src/providers/google-genai.ts @@ -5,7 +5,7 @@ import { normalizeAPIStatusError, } from '#/errors'; import type { Message, StreamedMessagePart, ThinkPart, ToolCall } from '#/message'; -import { isToolDeclarationOnlyMessage } from '#/message'; +import { isOpaqueAssistantMessage, isToolDeclarationOnlyMessage } from '#/message'; import type { ChatProvider, FinishReason, @@ -278,6 +278,8 @@ function messageToGoogleGenAI(message: Message): GoogleContent { case 'video_url': parts.push(convertMediaUrl(part.videoUrl.url, 'video/mp4')); break; + case 'openai_compaction': + break; } } @@ -356,6 +358,8 @@ function toolMessageToFunctionResponseParts( case 'think': // Skip — handled separately via reasoning channel. break; + case 'openai_compaction': + break; } } @@ -387,6 +391,10 @@ export function messagesToGoogleGenAIContents(messages: Message[]): GoogleConten i += 1; continue; } + if (isOpaqueAssistantMessage(message)) { + i += 1; + continue; + } if (message.role === 'system') { // Google GenAI's `Content.role` only accepts "user" or "model", so a diff --git a/packages/kosong/src/providers/kimi.ts b/packages/kosong/src/providers/kimi.ts index c228eba91db..39882275c13 100644 --- a/packages/kosong/src/providers/kimi.ts +++ b/packages/kosong/src/providers/kimi.ts @@ -1,6 +1,7 @@ import { normalizeKimiToolSchema } from './kimi-schema'; import { parseTraceId } from '#/errors'; import type { ContentPart, Message, StreamedMessagePart, ToolCall } from '#/message'; +import { isOpaqueAssistantMessage } from '#/message'; import type { ChatProvider, FinishReason, @@ -506,6 +507,7 @@ export class KimiChatProvider implements ChatProvider { const reasoningKey = this._reasoningKeyDialect.outboundKey() as ReasoningKey; const normalizedHistory = normalizeToolCallIdsForProvider(history, KIMI_TOOL_CALL_ID_POLICY); for (const msg of normalizedHistory) { + if (isOpaqueAssistantMessage(msg)) continue; messages.push(convertMessage(msg, preservedThinkingEnabled, reasoningKey)); } diff --git a/packages/kosong/src/providers/openai-common.ts b/packages/kosong/src/providers/openai-common.ts index 5f94a9fc661..ad0cb86018d 100644 --- a/packages/kosong/src/providers/openai-common.ts +++ b/packages/kosong/src/providers/openai-common.ts @@ -62,6 +62,9 @@ export function convertContentPart(part: ContentPart): OpenAIContentPart | null ? { url: part.videoUrl.url } : { url: part.videoUrl.url, id: part.videoUrl.id }, }; + case 'openai_compaction': + // Only the Responses API understands this opaque item. + return null; default: throw new Error(`Unknown content part type: ${(part as ContentPart).type}`); } diff --git a/packages/kosong/src/providers/openai-legacy.ts b/packages/kosong/src/providers/openai-legacy.ts index 3f5cf3df40a..1f307d9a64e 100644 --- a/packages/kosong/src/providers/openai-legacy.ts +++ b/packages/kosong/src/providers/openai-legacy.ts @@ -1,5 +1,5 @@ import type { ContentPart, Message, StreamedMessagePart, ToolCall } from '#/message'; -import { isToolDeclarationOnlyMessage } from '#/message'; +import { isOpaqueAssistantMessage, isToolDeclarationOnlyMessage } from '#/message'; import type { ChatProvider, FinishReason, @@ -308,6 +308,7 @@ function convertHistoryMessages( // because the leftover `{role:"system"}` without content is rejected by // the Chat Completions API. See isToolDeclarationOnlyMessage. if (isToolDeclarationOnlyMessage(msg)) continue; + if (isOpaqueAssistantMessage(msg)) continue; if (msg.role !== 'tool') { appendToolResultMediaMessage(messages, pendingToolResultMedia); } diff --git a/packages/kosong/src/providers/openai-responses.ts b/packages/kosong/src/providers/openai-responses.ts index e546c4f8458..373205f3f3f 100644 --- a/packages/kosong/src/providers/openai-responses.ts +++ b/packages/kosong/src/providers/openai-responses.ts @@ -10,6 +10,7 @@ import type { ChatProvider, FinishReason, GenerateOptions, + ProviderCompactionResult, ProviderRequestAuth, ResponseFormat, StreamedMessage, @@ -98,6 +99,11 @@ type ResponseOutputItemView = encryptedContent?: string; summary: RawObject[]; } + | { + type: 'compaction'; + id?: string; + encryptedContent: string; + } | { type: 'other'; }; @@ -198,6 +204,14 @@ function readResponseOutputItem( }; } + if (type === 'compaction') { + return { + type, + id: readStringField(item, 'id'), + encryptedContent: requireStringField(item, 'encrypted_content', context), + }; + } + return { type: 'other' }; } @@ -375,6 +389,26 @@ interface ResponseInputItem { [key: string]: unknown; } +interface CompactionSource { + readonly model: string; + readonly baseUrl?: string | undefined; +} + +function normalizeBaseUrl(value: string | undefined): string | undefined { + return value?.replace(/\/+$/, ''); +} + +function canReplayCompaction( + source: CompactionSource | undefined, + target: CompactionSource, +): boolean { + return ( + source === undefined || + (source.model === target.model && + normalizeBaseUrl(source.baseUrl) === normalizeBaseUrl(target.baseUrl)) + ); +} + interface ResponseToolParam { type: string; name: string; @@ -431,6 +465,9 @@ function contentPartsToInputItems(parts: ContentPart[]): unknown[] { case 'think': // Handled separately as reasoning items. break; + case 'openai_compaction': + // Emitted as a top-level item by convertMessage. + break; } } return items; @@ -474,6 +511,9 @@ function messageContentToFunctionOutputItems(content: ContentPart[]): unknown[] case 'think': // Handled separately as reasoning items. break; + case 'openai_compaction': + // Compaction items cannot appear inside function outputs. + break; } } return items; @@ -508,6 +548,7 @@ function convertMessage( message: Message, modelName: string, toolMessageConversion: ToolMessageConversion, + compactionTarget: CompactionSource, ): ResponseInputItem[] { let role: string = message.role; if (usesOpenAIResponsesDeveloperRole(modelName) && role === 'system') { @@ -589,6 +630,16 @@ function convertMessage( type: 'reasoning', encrypted_content: encryptedValue, }); + } else if (part.type === 'openai_compaction') { + flushPendingParts(); + if (canReplayCompaction(part.source, compactionTarget)) { + result.push({ + type: 'compaction', + encrypted_content: part.encryptedContent, + id: part.id, + }); + } + i += 1; } else { pendingParts.push(part); i += 1; @@ -632,6 +683,7 @@ function convertHistoryMessages( history: readonly Message[], modelName: string, toolMessageConversion: ToolMessageConversion, + compactionTarget: CompactionSource, ): unknown[] { const input: unknown[] = []; const pendingToolResultMedia: unknown[] = []; @@ -657,7 +709,7 @@ function convertHistoryMessages( if (msg.role !== 'tool') { flushPendingMedia(); } - input.push(...convertMessage(msg, modelName, toolMessageConversion)); + input.push(...convertMessage(msg, modelName, toolMessageConversion, compactionTarget)); if (msg.role === 'tool' && toolMessageConversion === 'extract_text') { pendingToolResultMedia.push( ...messageContentToFunctionOutputItems(msg.content.filter(isMediaPart)), @@ -668,6 +720,109 @@ function convertHistoryMessages( flushPendingMedia(); return input; } + +function convertCompactedResponse( + response: unknown, + source: CompactionSource, +): ProviderCompactionResult { + const raw = asRawObject(response); + if (raw === null) { + throw new ChatProviderError('OpenAI Responses compact returned an invalid response.'); + } + const output = readObjectArrayField(raw, 'output'); + if (output === undefined || output.length === 0) { + throw new ChatProviderError('OpenAI Responses compact returned an empty replacement window.'); + } + + const messages: Message[] = []; + let compactionCount = 0; + for (const item of output) { + const type = readStringField(item, 'type'); + if (type === 'message') { + const roleValue = readStringField(item, 'role'); + const role = + roleValue === 'developer' + ? 'system' + : roleValue === 'user' || roleValue === 'system' + ? roleValue + : undefined; + const content = readObjectArrayField(item, 'content'); + if (role === undefined || content === undefined) { + throw new ChatProviderError( + 'OpenAI Responses compact returned an unsupported message item.', + ); + } + const parts: ContentPart[] = content.map((part) => { + const partType = readStringField(part, 'type'); + if (partType === 'input_text' || partType === 'output_text') { + const text = readStringField(part, 'text'); + if (text !== undefined) return { type: 'text', text }; + } + if (partType === 'input_image') { + const imageUrl = readStringField(part, 'image_url'); + if (imageUrl !== undefined) return { type: 'image_url', imageUrl: { url: imageUrl } }; + } + throw new ChatProviderError( + `OpenAI Responses compact returned unsupported message content: ${partType ?? 'unknown'}.`, + ); + }); + messages.push({ role, content: parts, toolCalls: [] }); + continue; + } + if (type === 'compaction') { + const encryptedContent = readStringField(item, 'encrypted_content'); + if (encryptedContent === undefined) { + throw new ChatProviderError( + 'OpenAI Responses compact returned a compaction item without encrypted content.', + ); + } + compactionCount += 1; + messages.push({ + role: 'assistant', + content: [{ + type: 'openai_compaction', + encryptedContent, + id: readStringField(item, 'id'), + source, + }], + toolCalls: [], + }); + continue; + } + throw new ChatProviderError( + `OpenAI Responses compact returned an unsupported replacement item: ${type ?? 'unknown'}.`, + ); + } + if (compactionCount !== 1 || messages.at(-1)?.content[0]?.type !== 'openai_compaction') { + throw new ChatProviderError( + 'OpenAI Responses compact must return user messages followed by one compaction item.', + ); + } + + const usageRaw = readObjectField(raw, 'usage'); + let usage: TokenUsage | undefined; + if (usageRaw !== undefined && usageRaw !== null) { + const inputTokens = readNumberField(usageRaw, 'input_tokens') ?? 0; + const outputTokens = readNumberField(usageRaw, 'output_tokens') ?? 0; + const details = readObjectField(usageRaw, 'input_tokens_details'); + const cached = + details === undefined || details === null + ? 0 + : (readNumberField(details, 'cached_tokens') ?? 0); + usage = { + inputOther: Math.max(0, inputTokens - cached), + output: outputTokens, + inputCacheRead: cached, + inputCacheCreation: 0, + }; + } + return { + messages, + usage, + id: readStringField(raw, 'id'), + }; +} + export class OpenAIResponsesStreamedMessage implements StreamedMessage { private _id: string | null = null; private _usage: TokenUsage | null = null; @@ -675,7 +830,11 @@ export class OpenAIResponsesStreamedMessage implements StreamedMessage { private _rawFinishReason: string | null = null; private readonly _iter: AsyncGenerator; - constructor(response: unknown, isStream: boolean) { + constructor( + response: unknown, + isStream: boolean, + private readonly compactionSource?: CompactionSource, + ) { if (isStream) { this._iter = this._convertStreamResponse(response as AsyncIterable); } else { @@ -779,6 +938,13 @@ export class OpenAIResponsesStreamedMessage implements StreamedMessage { } yield thinkPart; } + } else if (outputItem.type === 'compaction') { + yield { + type: 'openai_compaction', + encryptedContent: outputItem.encryptedContent, + id: outputItem.id, + ...(this.compactionSource === undefined ? {} : { source: this.compactionSource }), + }; } } } @@ -935,6 +1101,13 @@ export class OpenAIResponsesStreamedMessage implements StreamedMessage { (thinkPart as { encrypted: string }).encrypted = item.encryptedContent; } yield thinkPart; + } else if (item.type === 'compaction') { + yield { + type: 'openai_compaction', + encryptedContent: item.encryptedContent, + id: item.id, + ...(this.compactionSource === undefined ? {} : { source: this.compactionSource }), + }; } else if (item.type === 'function_call' && typeof item.arguments === 'string') { const streamIndex = responseStreamIndex(item.itemId, outputIndex); yield* yieldFinalArgumentsSuffix(streamIndex, item.arguments, type); @@ -1095,8 +1268,17 @@ export class OpenAIResponsesChatProvider implements ChatProvider { history, OPENAI_RESPONSES_TOOL_CALL_ID_POLICY, ); + const compactionSource = { + model: this._model, + baseUrl: options?.auth?.baseUrl ?? this._baseUrl, + }; input.push( - ...convertHistoryMessages(normalizedHistory, this._model, this._toolMessageConversion), + ...convertHistoryMessages( + normalizedHistory, + this._model, + this._toolMessageConversion, + compactionSource, + ), ); const kwargs: Record = { ...this._generationKwargs }; @@ -1154,7 +1336,61 @@ export class OpenAIResponsesChatProvider implements ChatProvider { create(params: unknown, opts?: unknown): Promise; } ).create(createParams, options?.signal ? { signal: options.signal } : undefined); - return new OpenAIResponsesStreamedMessage(response, this._stream); + return new OpenAIResponsesStreamedMessage(response, this._stream, compactionSource); + } catch (error: unknown) { + throw convertOpenAIError(error); + } + } + + async compact( + systemPrompt: string, + _tools: Tool[], + history: Message[], + options?: GenerateOptions, + ): Promise { + const normalizedHistory = normalizeToolCallIdsForProvider( + history, + OPENAI_RESPONSES_TOOL_CALL_ID_POLICY, + ); + const params: Record = { + model: this._model, + input: convertHistoryMessages( + normalizedHistory, + this._model, + this._toolMessageConversion, + { + model: this._model, + baseUrl: options?.auth?.baseUrl ?? this._baseUrl, + }, + ), + }; + if (systemPrompt) params['instructions'] = systemPrompt; + const cacheKey = this._generationKwargs['prompt_cache_key']; + if (typeof cacheKey === 'string') params['prompt_cache_key'] = cacheKey; + // Preserve the encrypted reasoning chain across the compaction boundary. + // Codex sets `include: ["reasoning.encrypted_content"]` unconditionally on + // every Responses request (compaction included); without it the model + // retains message continuity but loses its accumulated reasoning state. + params['include'] = ['reasoning.encrypted_content']; + + try { + const client = this._createClient(options?.auth); + const compact = (client as { responses?: { compact?: unknown } }).responses?.compact; + if (typeof compact !== 'function') { + throw new TypeError( + 'OpenAI SDK version does not support Responses compaction. Upgrade to a version with responses.compact.', + ); + } + options?.onRequestSent?.(); + const response = await ( + client.responses as { + compact(params: unknown, opts?: unknown): Promise; + } + ).compact(params, options?.signal ? { signal: options.signal } : undefined); + return convertCompactedResponse(response, { + model: this._model, + baseUrl: options?.auth?.baseUrl ?? this._baseUrl, + }); } catch (error: unknown) { throw convertOpenAIError(error); } diff --git a/packages/kosong/test/kimi.test.ts b/packages/kosong/test/kimi.test.ts index 23f3fe24753..c7716ff83c8 100644 --- a/packages/kosong/test/kimi.test.ts +++ b/packages/kosong/test/kimi.test.ts @@ -213,6 +213,21 @@ describe('KimiChatProvider', () => { ]); }); + it('drops compaction-only assistant shells from Kimi history', async () => { + const provider = createProvider(); + const history: Message[] = [ + { + role: 'assistant', + content: [{ type: 'openai_compaction', encryptedContent: 'encrypted-summary' }], + toolCalls: [], + }, + { role: 'user', content: [{ type: 'text', text: 'Continue' }], toolCalls: [] }, + ]; + const body = await captureRequestBody(provider, '', [], history); + + expect(body['messages']).toEqual([{ role: 'user', content: 'Continue' }]); + }); + it('multi-turn with system prompt', async () => { const provider = createProvider(); const history: Message[] = [ diff --git a/packages/kosong/test/openai-legacy.test.ts b/packages/kosong/test/openai-legacy.test.ts index 3f5d8b5bf83..e5e57b3c8e8 100644 --- a/packages/kosong/test/openai-legacy.test.ts +++ b/packages/kosong/test/openai-legacy.test.ts @@ -125,6 +125,21 @@ describe('OpenAILegacyChatProvider', () => { ]); }); + it('drops compaction-only assistant shells from Chat Completions history', async () => { + const provider = createProvider(); + const history: Message[] = [ + { + role: 'assistant', + content: [{ type: 'openai_compaction', encryptedContent: 'encrypted-summary' }], + toolCalls: [], + }, + { role: 'user', content: [{ type: 'text', text: 'Continue' }], toolCalls: [] }, + ]; + const body = await captureRequestBody(provider, '', [], history); + + expect(body['messages']).toEqual([{ role: 'user', content: 'Continue' }]); + }); + it('multi-turn with system prompt', async () => { const provider = createProvider(); const history: Message[] = [ diff --git a/packages/kosong/test/openai-responses.test.ts b/packages/kosong/test/openai-responses.test.ts index c7f9b3ee97d..9a1a0dd6188 100644 --- a/packages/kosong/test/openai-responses.test.ts +++ b/packages/kosong/test/openai-responses.test.ts @@ -97,6 +97,80 @@ const MUL_TOOL: Tool = { }; describe('OpenAIResponsesChatProvider', () => { + it('uses responses.compact and returns a validated replacement window', async () => { + const provider = new OpenAIResponsesChatProvider({ + model: 'gpt-4.1', + apiKey: 'test-key', + generationKwargs: { prompt_cache_key: 'session-probe' }, + }); + let captured: Record | undefined; + ((provider as any)._client.responses as Record)['compact'] = vi + .fn() + .mockImplementation((params: unknown) => { + captured = params as Record; + return Promise.resolve({ + id: 'cmp_response_1', + output: [ + { + type: 'message', + role: 'user', + content: [{ type: 'input_text', text: 'Keep this request.' }], + }, + { type: 'compaction', id: 'cmp_1', encrypted_content: 'opaque-state' }, + ], + usage: { + input_tokens: 100, + output_tokens: 20, + input_tokens_details: { cached_tokens: 40 }, + }, + }); + }); + + const result = await provider.compact('', [], [ + { role: 'user', content: [{ type: 'text', text: 'Keep this request.' }], toolCalls: [] }, + ]); + + expect(captured).toEqual({ + model: 'gpt-4.1', + input: [ + { + type: 'message', + role: 'user', + content: [{ type: 'input_text', text: 'Keep this request.' }], + }, + ], + prompt_cache_key: 'session-probe', + // Compaction preserves the encrypted reasoning chain across the boundary + // (matches Codex's unconditional include on every Responses request). + include: ['reasoning.encrypted_content'], + }); + expect(result.messages).toEqual([ + { + role: 'user', + content: [{ type: 'text', text: 'Keep this request.' }], + toolCalls: [], + }, + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'opaque-state', + id: 'cmp_1', + source: { model: 'gpt-4.1', baseUrl: 'https://api.openai.com/v1' }, + }, + ], + toolCalls: [], + }, + ]); + expect(result.usage).toEqual({ + inputOther: 60, + output: 20, + inputCacheRead: 40, + inputCacheCreation: 0, + }); + }); + describe('message conversion', () => { it('sends system prompt as top-level instructions', async () => { const provider = createProvider(); @@ -145,6 +219,77 @@ describe('OpenAIResponsesChatProvider', () => { ]); }); + it('replays OpenAI compaction items as top-level response input items', async () => { + const provider = createProvider(); + const history: Message[] = [ + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'encrypted-summary', + id: 'cmp_123', + }, + { type: 'text', text: 'Continuing after compaction.' }, + ], + toolCalls: [], + }, + ]; + const body = await captureRequestBody(provider, '', [], history); + + expect(body['input']).toEqual([ + { + type: 'compaction', + encrypted_content: 'encrypted-summary', + id: 'cmp_123', + }, + { + content: [ + { + type: 'output_text', + text: 'Continuing after compaction.', + annotations: [], + }, + ], + role: 'assistant', + type: 'message', + }, + ]); + }); + + it('does not replay compaction state from another model or endpoint', async () => { + const provider = createProvider(); + const history: Message[] = [ + { + role: 'assistant', + content: [ + { + type: 'openai_compaction', + encryptedContent: 'other-model', + source: { model: 'gpt-5', baseUrl: 'https://api.openai.com/v1' }, + }, + { + type: 'openai_compaction', + encryptedContent: 'other-endpoint', + source: { model: 'gpt-4.1', baseUrl: 'https://example.com/v1' }, + }, + { type: 'text', text: 'portable output' }, + ], + toolCalls: [], + }, + ]; + + const body = await captureRequestBody(provider, '', [], history); + + expect(body['input']).toEqual([ + { + content: [{ type: 'output_text', text: 'portable output', annotations: [] }], + role: 'assistant', + type: 'message', + }, + ]); + }); + it('multi-turn with system prompt', async () => { const provider = createProvider(); const history: Message[] = [ @@ -1316,6 +1461,37 @@ describe('OpenAIResponsesChatProvider', () => { ]); }); + it('yields opaque compaction items from non-stream responses', async () => { + const provider = createProvider(); + (provider as any)._stream = false; + ((provider as any)._client.responses as unknown as Record)['create'] = vi + .fn() + .mockResolvedValue({ + id: 'resp_compaction', + output: [ + { + type: 'compaction', + id: 'cmp_123', + encrypted_content: 'encrypted-summary', + }, + ], + usage: { input_tokens: 5, output_tokens: 3, total_tokens: 8 }, + }); + + const stream = await provider.generate('', [], []); + const parts: StreamedMessagePart[] = []; + for await (const part of stream) parts.push(part); + + expect(parts).toEqual([ + { + type: 'openai_compaction', + encryptedContent: 'encrypted-summary', + id: 'cmp_123', + source: { model: 'gpt-4.1', baseUrl: 'https://api.openai.com/v1' }, + }, + ]); + }); + it('yields an empty ThinkPart from a non-stream reasoning item with no summaries', async () => { const provider = createProvider(); (provider as any)._stream = false; @@ -1877,6 +2053,31 @@ describe('OpenAIResponsesChatProvider', () => { expect(parts).toEqual([{ type: 'think', think: '', encrypted: 'enc_done' }]); }); + it('yields opaque compaction items from response.output_item.done', async () => { + const events = [ + { + type: 'response.output_item.done', + item: { + type: 'compaction', + id: 'cmp_stream', + encrypted_content: 'encrypted-stream-summary', + }, + }, + ]; + + const stream = new OpenAIResponsesStreamedMessage(makeAsyncIterable(events), true); + const parts: StreamedMessagePart[] = []; + for await (const part of stream) parts.push(part); + + expect(parts).toEqual([ + { + type: 'openai_compaction', + encryptedContent: 'encrypted-stream-summary', + id: 'cmp_stream', + }, + ]); + }); + it('yields ThinkPart from response.output_item.done reasoning item without encrypted_content', async () => { const events = [ { diff --git a/packages/kosong/test/type-safety.test.ts b/packages/kosong/test/type-safety.test.ts index 5c3dec1a293..ea773348910 100644 --- a/packages/kosong/test/type-safety.test.ts +++ b/packages/kosong/test/type-safety.test.ts @@ -26,6 +26,8 @@ function processPartSafely(part: StreamedMessagePart): string { return part.audioUrl.url; // AudioURLPart.audioUrl.url -> string case 'video_url': return part.videoUrl.url; // VideoURLPart.videoUrl.url -> string + case 'openai_compaction': + return part.encryptedContent; // OpenAICompactionPart.encryptedContent -> string case 'function': return part.name; // ToolCall.name -> string case 'tool_call_part':