Compare commits

...
10 changed files with 674 additions and 342 deletions
@@ -1,5 +1,6 @@
import type { IntegrationOAuthMethodRegistration } from "@opencode/plugin/effect/integration"
import { define } from "@opencode/plugin/effect/plugin"
import type { SessionRequest } from "@opencode/plugin/effect/session"
import { Deferred, Effect, Option, Schema, Semaphore, Stream } from "effect"
import type { Server } from "node:http"
import { App } from "../../app.js"
@@ -307,6 +308,13 @@ export const OpenAIPlugin = define({
}),
{ providerID: Provider.ID.openai },
)
// The ChatGPT backend rejects a requested output limit, and OpenAI counts one against rate limits.
const omitOutputLimit = (evt: SessionRequest) =>
Effect.sync(() => {
delete evt.options.maxTokens
})
for (const name of ["context", "compaction"] as const)
yield* ctx.session.hook(name, omitOutputLimit, { providerID: Provider.ID.openai })
const refresh = () => loading.withPermit(load().pipe(Effect.andThen(ctx.provider.reload())))
yield* bus.subscribe(Credential.Event.Switched).pipe(
Stream.filter((event) => event.data.integrationID === Integration.ID.make("openai")),
+400 -324
View File
@@ -10,12 +10,13 @@ import {
LLMRequest,
Message,
type ContentPart,
type ToolResultPart,
type Usage,
} from "@opencode/ai"
import type { StreamOptions } from "@opencode/ai/route"
import type { SessionCompactionResult } from "@opencode/plugin/effect/session"
import { SessionError } from "@opencode/schema/session-error"
import { Context, Effect, Layer, Stream } from "effect"
import { Context, Effect, Layer, Ref, Stream } from "effect"
import { Bus } from "../bus.js"
import { Database } from "../database/database.js"
import { makeLocationNode } from "@opencode/util/effect/app-node"
@@ -41,8 +42,16 @@ const DEFAULT_BUFFER = 20_000
const DEFAULT_KEEP_TOKENS = 15_000
const OUTPUT_TOKEN_MAX = 32_000
const TOOL_OUTPUT_MAX_CHARS = 2_000
const FALLBACK_TOOL_CHARS = 1_000
const FALLBACK_OVERFLOW_RETRIES = 3
const IMAGE_TOKEN_ESTIMATE = 1_500
const PDF_TOKEN_ESTIMATE = 2_000
/** The prompt size that starts automatic compaction and bounds a reduced summary request. */
const promptCeiling = (limit: SessionRunnerModel.Resolved["limit"], buffer: number) =>
Math.min(
limit.input === undefined ? Number.POSITIVE_INFINITY : limit.input - buffer,
limit.context - Math.max(Math.min(limit.output, OUTPUT_TOKEN_MAX), buffer),
)
const SUMMARY_TEMPLATE = `You MUST use this format for your response (you may omit sections that aren't applicable). Do not include the <template> tags in your response.
<template>
## Objective
@@ -165,13 +174,46 @@ export interface Interface extends State.Transformable<Editor> {
export class Service extends Context.Service<Service, Interface>()("@opencode/SessionCompaction") {}
const inputTokens = (tokens: NonNullable<SessionMessage.Assistant["tokens"]>) =>
tokens.input + tokens.cache.read + tokens.cache.write
const hasInputUsage = (message: SessionMessage.Info) =>
message.type === "assistant" &&
!message.error &&
message.tokens !== undefined &&
message.tokens.input + message.tokens.cache.read + message.tokens.cache.write > 0
message.type === "assistant" && !message.error && message.tokens !== undefined && inputTokens(message.tokens) > 0
const lastCheckpoint = (messages: readonly SessionMessage.Info[]) =>
messages.findLast(
(message): message is SessionMessage.CompactionCompleted =>
message.type === "compaction" && message.status === "completed",
)
/** Index of the oldest item in the newest run whose sizes total at most `budget`. */
const fitNewest = <T>(items: readonly T[], size: (item: T) => number, budget: number) => {
let total = 0
let start = items.length
while (start > 0) {
const next = total + size(items[start - 1])
if (next > budget) break
total = next
start--
}
return start
}
/** System prompt and tool definitions: sent with every request but outside the message history. */
const estimateFixed = (system: ReadonlyArray<{ readonly text: string }>, tools: SessionContext.Loaded["tools"]) =>
system.reduce((sum, part) => sum + Token.estimate(part.text), 0) +
tools.definitions.reduce(
(sum, tool) => sum + Token.estimate(tool.name + tool.description + JSON.stringify(tool.inputSchema)),
0,
)
export const estimateTokens = (input: RequiredInput) => {
const prompt = estimatePrompt(input)
return prompt.measured + prompt.estimated
}
/** The prompt size: `measured` is what the provider reported at the latest response, `estimated` is the text since. */
export const estimatePrompt = (input: RequiredInput) => {
const index = input.messages.findLastIndex(hasInputUsage)
const last = input.messages[index]
// Keep the anchor's local tool results: they are not covered by its provider usage.
@@ -182,14 +224,7 @@ export const estimateTokens = (input: RequiredInput) => {
.filter((message) => message.role !== "assistant" || message.id !== last?.id)
.reduce((sum, message) => sum + message.content.reduce((sum, part) => sum + estimatePart(part), 0), 0)
if (last?.type === "assistant" && last.tokens)
return (
added +
last.tokens.input +
last.tokens.cache.read +
last.tokens.cache.write +
last.tokens.output +
last.tokens.reasoning
)
return { measured: inputTokens(last.tokens) + last.tokens.output + last.tokens.reasoning, estimated: added }
const transcript = SessionModelRequest.baseTranscript({
agent: input.context.agent.info,
model: input.resolved,
@@ -197,14 +232,7 @@ export const estimateTokens = (input: RequiredInput) => {
initial: input.context.initial,
messages: [],
})
return (
added +
transcript.system.reduce((sum, part) => sum + Token.estimate(part.text), 0) +
input.context.tools.definitions.reduce(
(sum, tool) => sum + Token.estimate(tool.name + tool.description + JSON.stringify(tool.inputSchema)),
0,
)
)
return { measured: 0, estimated: added + estimateFixed(transcript.system, input.context.tools) }
}
const estimateMedia = (mime: string) => {
@@ -224,9 +252,12 @@ const estimatePart = (part: ContentPart): number => {
(sum, content) => sum + (content.type === "text" ? Token.estimate(content.text) : estimateMedia(content.mime)),
0,
)
return Token.estimate(
typeof part.result.value === "string" ? part.result.value : (JSON.stringify(part.result.value) ?? ""),
)
return Token.estimate(toolResultText(part.result))
}
const toolResultText = (result: ToolResultPart["result"]) => {
if (result.type === "content") return serializeToolContent(result.value)
return typeof result.value === "string" ? result.value : (JSON.stringify(result.value) ?? "")
}
/** Keep whole, real user messages, never synthetic guidance or half an attachment/tool exchange. */
@@ -244,21 +275,14 @@ export const retainUsers = (
model.capabilities,
),
)
let tokens = 0
let start = users.length
for (let index = users.length - 1; index >= 0; index--) {
const size = users[index].content.reduce((sum, part) => sum + estimatePart(part), 0)
if (tokens + size > keepTokens) break
tokens += size
start = index
}
return users.slice(start)
const size = (message: Message) => message.content.reduce((sum, part) => sum + estimatePart(part), 0)
return users.slice(fitNewest(users, size, keepTokens))
}
export const truncateToolOutput = (value: string) => {
if (value.length <= TOOL_OUTPUT_MAX_CHARS) return value
export const truncateToolOutput = (value: string, maxChars = TOOL_OUTPUT_MAX_CHARS) => {
if (value.length <= maxChars) return value
let end = 0
for (let count = 0; count < TOOL_OUTPUT_MAX_CHARS && end < value.length; count++) {
for (let count = 0; count < maxChars && end < value.length; count++) {
const code = value.charCodeAt(end)
end +=
code >= 0xd800 && code <= 0xdbff && value.charCodeAt(end + 1) >= 0xdc00 && value.charCodeAt(end + 1) <= 0xdfff
@@ -269,7 +293,7 @@ export const truncateToolOutput = (value: string) => {
return `${value.slice(0, end)}\n[truncated]`
}
export const serializeToolContent = (content: SessionMessage.ToolStateCompleted["content"]) =>
export const serializeToolContent = (content: ReadonlyArray<SessionMessage.ToolStateCompleted["content"][number]>) =>
content
.map((item) =>
item.type === "text" ? item.text : `[Attached ${item.mime}${item.name === undefined ? "" : `: ${item.name}`}]`,
@@ -299,14 +323,11 @@ const serializeRecentMessage = (message: SessionMessage.Info) => {
if (part.type === "text") return [`[Assistant]: ${part.text}`]
if (part.type === "reasoning") return part.text ? [`[Assistant reasoning]: ${part.text}`] : []
const input = typeof part.state.input === "string" ? part.state.input : JSON.stringify(part.state.input)
const call = `[Assistant tool call]: ${part.name}(${input})`
if (part.state.status === "completed")
return [
`[Assistant tool call]: ${part.name}(${input})`,
`[Tool result]: ${truncateToolOutput(serializeToolContent(part.state.content))}`,
]
if (part.state.status === "error")
return [`[Assistant tool call]: ${part.name}(${input})`, `[Tool error]: ${part.state.error.message}`]
return [`[Assistant tool call]: ${part.name}(${input})`]
return [call, `[Tool result]: ${truncateToolOutput(serializeToolContent(part.state.content))}`]
if (part.state.status === "error") return [call, `[Tool error]: ${part.state.error.message}`]
return [call]
})
.join("\n")
}
@@ -319,6 +340,39 @@ const serializeRecentMessage = (message: SessionMessage.Info) => {
return ""
}
/** Flatten provider-bound history into text so a reduced summary request carries no tool pairs or reasoning signatures. */
const serializeFallback = (messages: readonly Message[]) =>
messages.flatMap((message) => {
const parts = message.content.flatMap((part) => {
if (part.type === "text" || part.type === "reasoning") return part.text ? [part.text] : []
if (part.type === "media")
return [`[Attached ${part.media.mediaType}${part.filename ? `: ${part.filename}` : ""}; content omitted]`]
if (part.type === "tool-call") return [`[Tool call ${part.name}(${JSON.stringify(part.input)})]`]
if (part.type === "tool-result")
return [`[Tool result ${part.name}]: ${truncateToolOutput(toolResultText(part.result), FALLBACK_TOOL_CHARS)}`]
if (part.type === "compaction" && part.text) return [part.text]
return []
})
return parts.length ? [{ role: message.role, text: `[${message.role}]: ${parts.join("\n")}` }] : []
})
/** Retain a prior checkpoint and the newest complete exchanges that fit the summary input budget, or all of them. */
const fitFallback = (entries: ReturnType<typeof serializeFallback>, budget = Number.POSITIVE_INFINITY) => {
const previous = entries[0]?.text.includes("<conversation-checkpoint>") ? entries[0] : undefined
const rest = previous ? entries.slice(1) : entries
const groups = rest.reduce<Array<string>>((groups, entry) => {
if (entry.role === "user" || groups.length === 0) return [...groups, entry.text]
return [...groups.slice(0, -1), `${groups[groups.length - 1]}\n\n${entry.text}`]
}, [])
const header = previous ? `${previous.text}\n\n` : ""
const allowance = budget - Token.estimate(header)
if (allowance <= 0) return
const start = fitNewest(groups, Token.estimate, allowance)
if (groups.length && start === groups.length) return
const text = `${header}${start ? `[${start} older exchanges omitted from this summary input]\n\n` : ""}${groups.slice(start).join("\n\n")}`
return { text, omitted: start, tokens: Token.estimate(text) }
}
const splitHistory = (messages: readonly SessionMessage.Info[], keepTokens: number) => {
const tailStart = findTailStart(messages, keepTokens)
if (tailStart === undefined) return
@@ -336,29 +390,20 @@ const findTailStart = (messages: readonly SessionMessage.Info[], keepTokens: num
if (conversation.length === 0) return undefined
// Keep at least the newest entry, even if it exceeds the allowance.
let total = 0
let start = conversation.length
for (let index = conversation.length - 1; index >= 0; index--) {
const next = total + Token.estimate(conversation[index].text)
if (start < conversation.length && next > keepTokens) break
total = next
start = index
}
const fitted = Math.min(
fitNewest(conversation, (item) => Token.estimate(item.text), keepTokens),
conversation.length - 1,
)
// Start at a user boundary so an assistant's tool calls and results stay together.
while (start > 0 && conversation[start].message.type !== "user") start--
const start = conversation.findLastIndex((item, index) => index <= fitted && item.message.type === "user")
if (start > 0) return conversation[start].index
// If everything fits, retain only the latest exchange to leave an older prefix to summarize.
const latestUser = conversation.findLastIndex((item) => item.message.type === "user")
if (latestUser > 0) return conversation[latestUser].index
const previousSummary = messages.findLast(
(message): message is SessionMessage.CompactionCompleted =>
message.type === "compaction" && message.status === "completed",
)
// Without an older retained tail to summarize, summarize everything and retain nothing.
return previousSummary?.recent ? conversation[0].index : messages.length
return lastCheckpoint(messages)?.recent ? conversation[0].index : messages.length
}
export const buildPrompt = (update: boolean, legacy = false) => {
@@ -392,6 +437,19 @@ export const buildPrompt = (update: boolean, legacy = false) => {
const hasSummarySection = (summary: string) =>
summary.split("\n").some((line) => SUMMARY_HEADINGS.includes(line.trim()))
type Envelope = Pick<SessionEvent.Compaction.Failed["data"], "sessionID" | "reason" | "inputID">
/** One summary attempt: text so far plus the first failure, if any, and whether it was a context overflow. */
type Summary = {
readonly summary: string
readonly overflow: boolean
readonly failure?: SessionError.Error
readonly providerState?: SessionMessage.ProviderState
}
const NUDGE =
"The previous response did not fill in the required summary template. Do not call tools. Return the summary as text using the exact section headings from the template."
export const layer = Layer.effect(
Service,
Effect.gen(function* () {
@@ -413,34 +471,35 @@ export const layer = Layer.effect(
},
}),
})
const failed = Effect.fnUntraced(function* (input: SessionEvent.Compaction.Failed["data"]) {
yield* bus.publish(SessionEvent.Compaction.Failed, input)
return { status: "failed" as const, error: input.error }
const envelope = (input: ExecuteInput): Envelope => ({
sessionID: input.context.session.id,
reason: input.reason,
inputID: input.inputID,
})
const started = (input: ExecuteInput, recent: string) =>
input.started
? Effect.void
: bus.publish(SessionEvent.Compaction.Started, {
sessionID: input.context.session.id,
reason: input.reason,
recent,
inputID: input.inputID,
})
const supplied = Effect.fn("SessionCompaction.supplied")(function* (
const recordUsage = (sessionID: SessionSchema.ID, usage: SessionUsage.Recorded | undefined) =>
usage ? bus.publish(SessionEvent.UsageRecorded, { sessionID, source: "compaction", ...usage }) : Effect.void
const failed = Effect.fnUntraced(function* (
target: Envelope,
error: SessionError.Error,
usage?: SessionUsage.Recorded,
) {
yield* recordUsage(target.sessionID, usage)
yield* bus.publish(SessionEvent.Compaction.Failed, { ...target, error, ...usage })
return { status: "failed" as const, error }
})
const ended = Effect.fnUntraced(function* (
input: ExecuteInput,
result: SessionCompactionResult,
recent: string,
result: {
readonly text: string
readonly recent: string
readonly providerState?: SessionMessage.ProviderState
readonly providerContext?: SessionProviderContext.Info
readonly usage?: SessionUsage.Recorded
readonly metadata?: Record<string, unknown>
},
) {
const context = input.context
const usage = result.tokens
? { tokens: result.tokens, cost: SessionUsage.calculateCost(context.model.cost, result.tokens) }
: undefined
if (usage)
yield* bus.publish(SessionEvent.UsageRecorded, {
sessionID: context.session.id,
source: "compaction",
...usage,
})
yield* recordUsage(context.session.id, result.usage)
yield* bus.publish(
SessionEvent.Compaction.Ended,
{
@@ -448,65 +507,106 @@ export const layer = Layer.effect(
reason: input.reason,
model: context.model.ref,
providerState: result.providerState,
text: result.summary,
recent,
...usage,
providerContext: result.providerContext,
text: result.text,
recent: result.recent,
...result.usage,
},
{ metadata: result.metadata },
)
return { status: "completed" as const }
})
const started = (input: ExecuteInput, recent: string) =>
input.started ? Effect.void : bus.publish(SessionEvent.Compaction.Started, { ...envelope(input), recent })
/** A hook answered the request itself. */
const supplied = (
input: ExecuteInput,
result: SessionCompactionResult,
recent: string,
prior?: SessionUsage.Recorded,
) => {
const own = result.tokens && {
tokens: result.tokens,
cost: SessionUsage.calculateCost(input.context.model.cost, result.tokens),
}
return ended(input, {
text: result.summary,
recent,
providerState: result.providerState,
usage: prior && own ? SessionUsage.add(prior, own) : (prior ?? own),
metadata: result.metadata,
})
}
// Manual controls settle through the inbox; only automatic work needs a durable interruption record.
const interrupted = (input: ExecuteInput) =>
input.reason === "auto"
? failed({
sessionID: input.context.session.id,
reason: input.reason,
inputID: input.inputID,
error: { type: "compaction.interrupted", message: "Compaction was interrupted" },
}).pipe(Effect.asVoid)
: Effect.void
const interrupted = (input: ExecuteInput, usage?: SessionUsage.Recorded) =>
Effect.gen(function* () {
yield* recordUsage(input.context.session.id, usage)
if (input.reason !== "auto") return
yield* failed(envelope(input), { type: "compaction.interrupted", message: "Compaction was interrupted" })
})
const prepare = (
input: ExecuteInput,
transcript: Pick<SessionModelRequest.Input, "system" | "messages">,
inputTokens: SessionModelRequest.Input["inputTokens"],
webSocket?: "session",
) =>
input.prepare({
session: input.context.session,
agent: input.context.agent.id,
model: input.context.model,
tools: input.context.tools,
system: transcript.system,
messages: transcript.messages,
webSocket,
inputTokens,
})
const transcript = (input: ExecuteInput, messages: readonly SessionMessage.Info[]) =>
SessionModelRequest.baseTranscript({
agent: input.context.agent.info,
model: input.context.model,
tools: input.context.tools,
initial: input.context.initial,
messages,
})
const compactionRequest = (
input: ExecuteInput,
messages: readonly SessionMessage.Info[],
webSocket?: "session",
summaryPrompt?: string,
) => {
const context = input.context
const transcript = SessionModelRequest.baseTranscript({
agent: context.agent.info,
model: context.model,
tools: context.tools,
initial: context.initial,
messages,
})
return input.prepare({
session: context.session,
agent: context.agent.id,
model: context.model,
tools: context.tools,
system: transcript.system,
messages: [
...transcript.messages,
...(input.instructionUpdate ? [Message.system(input.instructionUpdate)] : []),
],
const base = transcript(input, messages)
const prompt = estimatePrompt({ messages, resolved: input.context.model, context: input.context })
return prepare(
input,
{
system: base.system,
messages: [...base.messages, ...(input.instructionUpdate ? [Message.system(input.instructionUpdate)] : [])],
},
// The instruction update and summary prompt are sent outside the history, so count them too.
{
measured: prompt.measured,
estimated: prompt.estimated + Token.estimate((input.instructionUpdate ?? "") + (summaryPrompt ?? "")),
},
webSocket,
})
)
}
const retry = Effect.fnUntraced(function* (input: ExecuteInput, hook: SessionModelRequest.Prepared["retry"]) {
return SessionRunnerRetry.transient(yield* SessionRunnerRetry.policy(input.context.session.id), {
agent: input.context.agent.id,
model: input.context.model.ref,
hook,
})
})
/** The durable transcript since the last local summary, re-expanding every native window. */
const original = (sessionID: SessionSchema.ID) => SessionHistory.load(db, sessionID, "local").pipe(Effect.orDie)
const recoverLocally = (input: ExecuteInput) =>
original(input.context.session.id).pipe(
Effect.flatMap((messages) => execute({ ...input, context: { ...input.context, messages } })),
Effect.flatMap((messages) => summarize({ ...input, context: { ...input.context, messages } })),
)
const executeProvider = Effect.fn("SessionCompaction.executeProvider")(function* (input: ExecuteInput) {
const native = Effect.fn("SessionCompaction.native")(function* (input: ExecuteInput) {
const context = input.context
const reject = (message: string) =>
failed({
sessionID: context.session.id,
reason: input.reason,
inputID: input.inputID,
error: { type: "provider.unsupported-operation", message },
})
const reject = (message: string) => failed(envelope(input), { type: "provider.unsupported-operation", message })
const prepared = yield* compactionRequest(input, context.messages, "session")
if (prepared.event.result) {
yield* started(input, "")
@@ -517,16 +617,12 @@ export const layer = Layer.effect(
if (!provenance) return yield* reject("Provider compaction requires a stable, configured endpoint")
// History is selected before request hooks. Until that interface can select on the final route,
// require routing in the catalog; never install a checkpoint that the next request would skip.
if (
!SessionProviderContext.compatible(
provenance,
SessionProviderContext.provenance({ model: request.model, ref: context.model.ref }),
)
)
const routed = SessionProviderContext.provenance({ model: request.model, ref: context.model.ref })
if (!SessionProviderContext.compatible(provenance, routed))
return yield* reject(
"Provider compaction requires the endpoint in provider/model settings, not a model.request rewrite",
)
const native = state
const strategy = state
.get()
.native.toReversed()
.map((strategy) =>
@@ -539,39 +635,25 @@ export const layer = Layer.effect(
}),
)
.find((effect) => effect !== undefined)
if (!native)
if (!strategy)
return yield* reject(
`No plugin provides native compaction for ${request.model.provider}/${request.model.route.id}`,
)
const transient = SessionRunnerRetry.transient(yield* SessionRunnerRetry.policy(context.session.id), {
agent: context.agent.id,
model: context.model.ref,
hook: prepared.retry,
})
const transient = yield* retry(input, prepared.retry)
yield* started(input, "")
// Transient provider failures retry like any other request; only a known automatic overflow permits
// local recovery, and nothing is installed until the provider returns a checkpoint.
return yield* Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
// Transient provider failures retry like any other request; only a known automatic overflow permits
// local recovery, and nothing is installed until the provider returns a checkpoint.
const result = yield* restore(native.pipe(transient))
const usage = result.usage ? SessionUsage.record(result.usage, context.model.cost) : undefined
if (usage)
yield* bus.publish(SessionEvent.UsageRecorded, {
sessionID: context.session.id,
source: "compaction" as const,
...usage,
})
yield* bus.publish(SessionEvent.Compaction.Ended, {
sessionID: context.session.id,
reason: input.reason,
model: context.model.ref,
text: "",
recent: "",
providerContext: SessionProviderContext.encode(provenance, result.replacement),
...usage,
})
return { status: "completed" as const }
}),
restore(strategy.pipe(transient)).pipe(
Effect.flatMap((result) =>
ended(input, {
text: "",
recent: "",
providerContext: SessionProviderContext.encode(provenance, result.replacement),
usage: result.usage && SessionUsage.record(result.usage, context.model.cost),
}),
),
),
).pipe(
Effect.onInterrupt(() => interrupted(input)),
Effect.catchTag(
@@ -583,161 +665,173 @@ export const layer = Layer.effect(
result.status === "completed" ? { ...result, recoveredOverflow: true } : result,
),
)
: failed({
sessionID: context.session.id,
reason: input.reason,
inputID: input.inputID,
error: toSessionError(cause),
}),
: failed(envelope(input), toSessionError(cause)),
),
)
})
const execute = Effect.fn("SessionCompaction.execute")(function* (input: ExecuteInput) {
/** One summary request. Usage accumulates in `usage` so an interruption can still account for it. */
const stream = (
input: ExecuteInput,
request: LLMRequest,
options: StreamOptions,
transient: ReturnType<typeof SessionRunnerRetry.transient>,
usage: Ref.Ref<SessionUsage.Recorded | undefined>,
) => {
const context = input.context
const key = context.model.model.route.providerMetadataKey ?? context.model.model.provider
return llm.stream(request, options).pipe(
Stream.runFoldEffect(
(): Summary => ({ summary: "", overflow: false }),
(acc, event): Effect.Effect<Summary, AIError> => {
if (LLMEvent.is.providerError(event)) {
const overflow = event.classification === "context-overflow"
return Effect.succeed({
...acc,
overflow,
failure: { type: overflow ? "provider.invalid-request" : "provider.error", message: event.message },
})
}
if (LLMEvent.is.textDelta(event))
return bus
.publish(SessionEvent.Compaction.Delta, { sessionID: context.session.id, text: event.text })
.pipe(Effect.as({ ...acc, summary: acc.summary + event.text }))
if (LLMEvent.is.stepFinish(event))
return Ref.update(usage, (total) => {
const step = SessionUsage.record(event.usage, context.model.cost)
return total ? SessionUsage.add(total, step) : step
}).pipe(Effect.as({ ...acc, providerState: event.providerMetadata?.[key] }))
if (!LLMEvent.is.finish(event)) return Effect.succeed(acc)
const reason = event.reason.normalized
if (reason === "unknown")
return Effect.fail(
new AIError({
reason: new InvalidProviderOutputError({
message: "The provider response ended with an unknown finish reason.",
classification: "incomplete-stream",
}),
}),
)
if (reason === "error")
return Effect.fail(
new AIError({ reason: new UnknownProviderError({ message: "Compaction generation failed" }) }),
)
if (reason === "length")
return Effect.succeed({
...acc,
failure: { type: "compaction.failed", message: "Compaction summary reached the output token limit" },
})
if (reason === "content-filter")
return Effect.succeed({
...acc,
failure: { type: "provider.content-filter", message: "Compaction summary was blocked by the provider" },
})
return Effect.succeed(acc)
},
),
transient,
Effect.catchTag("AI.Error", (error) =>
Effect.succeed<Summary>({
summary: "",
overflow: isContextOverflowFailure(error),
failure: toSessionError(error),
}),
),
Effect.onInterrupt(() => Ref.get(usage).pipe(Effect.flatMap((total) => interrupted(input, total)))),
)
}
const summarize = Effect.fn("SessionCompaction.summarize")(function* (input: ExecuteInput) {
const context = input.context
const history = splitHistory(context.messages, state.get().tokens)
if (!history)
return yield* failed({
sessionID: context.session.id,
reason: input.reason,
error: { type: "compaction.unavailable", message: "Nothing to compact yet" },
inputID: input.inputID,
})
return yield* failed(envelope(input), { type: "compaction.unavailable", message: "Nothing to compact yet" })
yield* started(input, history.recent)
const chunks: string[] = []
let failure: SessionError.Error | undefined
let usage: SessionUsage.Recorded | undefined
let providerState: SessionMessage.ProviderState | undefined
const recordUsage = Effect.suspend(() =>
usage
? bus.publish(SessionEvent.UsageRecorded, {
sessionID: context.session.id,
source: "compaction",
...usage,
})
: Effect.void,
)
const previous = history.messages.findLast(
(message): message is SessionMessage.CompactionCompleted =>
message.type === "compaction" && message.status === "completed",
)
const previous = lastCheckpoint(history.messages)
// Checkpoints from the previous template ran far longer than this one asks for; its catch-all heading identifies them.
const legacy = previous?.summary.includes(LEGACY_HEADING) ?? false
const prepared = yield* compactionRequest(input, history.messages)
const prompt = buildPrompt(previous !== undefined, previous?.summary.includes(LEGACY_HEADING) ?? false)
const prepared = yield* compactionRequest(input, history.messages, undefined, prompt)
if (prepared.event.result) return yield* supplied(input, prepared.event.result, history.recent)
// Hooks see the transcript alone; the summary prompt is appended after they run.
const first = LLMRequest.update(prepared.request, {
messages: [...prepared.request.messages, Message.user(buildPrompt(previous !== undefined, legacy))],
})
// Both requests share the retry allowance; rejected output never enters the reminder request.
const transient = SessionRunnerRetry.transient(yield* SessionRunnerRetry.policy(context.session.id), {
agent: context.agent.id,
model: context.model.ref,
hook: prepared.retry,
const transient = yield* retry(input, prepared.retry)
const usage = yield* Ref.make<SessionUsage.Recorded | undefined>(undefined)
// Hooks see the transcript alone; the summary prompt is appended after they run.
const generate = Effect.fnUntraced(function* (request: LLMRequest, options: StreamOptions) {
const prompted = LLMRequest.update(request, { messages: [...request.messages, Message.user(prompt)] })
const first = yield* stream(input, prompted, options, transient, usage)
if (first.failure || hasSummarySection(first.summary)) return first
const nudged = LLMRequest.update(prompted, { messages: [...prompted.messages, Message.user(NUDGE)] })
return yield* stream(input, nudged, options, transient, usage)
})
for (const request of [
first,
LLMRequest.update(first, {
messages: [
...first.messages,
Message.user(
"The previous response did not fill in the required summary template. Do not call tools. Return the summary as text using the exact section headings from the template.",
),
],
}),
]) {
yield* Stream.suspend(() => {
chunks.length = 0
providerState = undefined
failure = undefined
return llm.stream(request, prepared.options)
}).pipe(
Stream.runForEach((event) => {
if (LLMEvent.is.providerError(event))
failure = {
type: event.classification === "context-overflow" ? "provider.invalid-request" : "provider.error",
message: event.message,
}
if (LLMEvent.is.textDelta(event)) {
chunks.push(event.text)
return bus.publish(SessionEvent.Compaction.Delta, {
sessionID: context.session.id,
text: event.text,
})
}
if (LLMEvent.is.stepFinish(event)) {
providerState =
event.providerMetadata?.[context.model.model.route.providerMetadataKey ?? context.model.model.provider]
const step = SessionUsage.record(event.usage, context.model.cost)
usage = usage ? SessionUsage.add(usage, step) : step
}
if (LLMEvent.is.finish(event)) {
if (event.reason.normalized === "length")
failure = { type: "compaction.failed", message: "Compaction summary reached the output token limit" }
if (event.reason.normalized === "content-filter")
failure = {
type: "provider.content-filter",
message: "Compaction summary was blocked by the provider",
}
if (event.reason.normalized === "unknown")
return Effect.fail(
new AIError({
reason: new InvalidProviderOutputError({
message: "The provider response ended with an unknown finish reason.",
classification: "incomplete-stream",
}),
}),
)
if (event.reason.normalized === "error")
return Effect.fail(
new AIError({ reason: new UnknownProviderError({ message: "Compaction generation failed" }) }),
)
}
return Effect.void
}),
transient,
Effect.catchTag("AI.Error", (error) =>
Effect.sync(() => {
failure = toSessionError(error)
}),
),
Effect.onInterrupt(() => recordUsage.pipe(Effect.andThen(interrupted(input)))),
)
if (failure || hasSummarySection(chunks.join(""))) break
}
yield* recordUsage
const summary = chunks.join("")
if (failure || !hasSummarySection(summary)) {
const error = failure ?? {
type: "compaction.failed" as const,
message: summary.trim()
? "Compaction summary did not match the required template"
: "Compaction produced no summary",
}
return yield* failed({
sessionID: context.session.id,
reason: input.reason,
error,
inputID: input.inputID,
...usage,
const finish = Effect.fnUntraced(function* (result: Summary, omitted: number) {
const total = yield* Ref.get(usage)
if (result.failure || !hasSummarySection(result.summary))
return yield* failed(
envelope(input),
result.failure ?? {
type: "compaction.failed",
message: result.summary.trim()
? "Compaction summary did not match the required template"
: "Compaction produced no summary",
},
total,
)
return yield* ended(input, {
text: omitted
? `${result.summary}\n\n[${omitted} older exchanges were omitted from the summary input; original session history is retained.]`
: result.summary,
recent: history.recent,
providerState: result.providerState,
usage: total,
})
}
yield* bus.publish(SessionEvent.Compaction.Ended, {
sessionID: context.session.id,
reason: input.reason,
model: context.model.ref,
providerState,
text: summary,
recent: history.recent,
...usage,
})
return { status: "completed" as const }
const limit = context.model.limit
const fixed = estimateFixed(prepared.request.system, context.tools) + Token.estimate(prompt)
const ceiling = limit.context > 0 ? promptCeiling(limit, state.get().buffer) - fixed : undefined
// A prefix that cannot fit the window at all skips the normal request and starts from the reduced text.
const oversized =
(ceiling ?? 0) > 0 &&
estimateTokens({ messages: history.messages, resolved: context.model, context }) + Token.estimate(prompt) >
Math.min(limit.context, limit.input ?? Number.POSITIVE_INFINITY)
const normal = oversized ? undefined : yield* generate(prepared.request, prepared.options)
if (normal && !normal.overflow) return yield* finish(normal, 0)
// Flatten the history to text and drop the oldest exchanges until the provider accepts the request.
const entries = serializeFallback(prepared.request.messages)
const system = transcript(input, []).system
let budget = ceiling
let last: string | undefined
for (let attempt = 1; ; attempt++) {
const fitted = fitFallback(entries, budget)
if (!fitted || fitted.text === last)
return yield* failed(
envelope(input),
{
type: "compaction.failed",
message: "The summary input cannot be reduced further without losing the latest exchange or checkpoint",
},
yield* Ref.get(usage),
)
const reduced = yield* prepare(
input,
{ system, messages: [Message.user(fitted.text)] },
{ measured: 0, estimated: fixed + fitted.tokens },
)
if (reduced.event.result)
return yield* supplied(input, reduced.event.result, history.recent, yield* Ref.get(usage))
const result = yield* generate(reduced.request, reduced.options)
if (!result.overflow || attempt === FALLBACK_OVERFLOW_RETRIES) return yield* finish(result, fitted.omitted)
budget = Math.floor(fitted.tokens / 2)
last = fitted.text
}
})
const run = (input: ExecuteInput) =>
input.context.model.compaction?.type === "native" ? native(input) : summarize(input)
const compact = Effect.fn("SessionCompaction.compact")(function* (input: AutoInput): Effect.fn.Return<Outcome> {
const request = { ...input, reason: "auto" as const }
if (input.overflow) return yield* recoverLocally(request)
if (input.context.model.compaction?.type !== "native") return yield* execute(request)
return yield* executeProvider(request)
return yield* input.overflow ? recoverLocally(request) : run(request)
})
const required = (input: RequiredInput) => {
const config = state.get()
@@ -752,43 +846,25 @@ export const layer = Layer.effect(
)
return false
const limit = input.resolved.limit
const context = limit.context
if (context <= 0) return false
const output = Math.min(limit.output, OUTPUT_TOKEN_MAX)
const promptCeiling = Math.min(
limit.input === undefined ? Number.POSITIVE_INFINITY : limit.input - config.buffer,
context - Math.max(output, config.buffer),
)
return estimateTokens(input) >= promptCeiling
if (limit.context <= 0) return false
return estimateTokens(input) >= promptCeiling(limit, config.buffer)
}
const compactManual = Effect.fn("SessionCompaction.compactManual")(function* (input: ManualInput) {
const target: Envelope = { sessionID: input.session.id, reason: "manual", inputID: input.inputID }
if (findTailStart(input.messages, state.get().tokens) === undefined)
return yield* failed({
sessionID: input.session.id,
reason: "manual",
error: { type: "compaction.unavailable", message: "Nothing to compact yet" },
inputID: input.inputID,
})
return yield* failed(target, { type: "compaction.unavailable", message: "Nothing to compact yet" })
return yield* input.resolveContext(input.session).pipe(
Effect.matchEffect({
onFailure: (cause) =>
failed({
sessionID: input.session.id,
reason: "manual",
error: toSessionError(cause),
inputID: input.inputID,
}),
onSuccess: (context) => {
const request = {
onFailure: (cause) => failed(target, toSessionError(cause)),
onSuccess: (context) =>
run({
context,
instructionUpdate: context.instructionUpdate,
prepare: input.prepare,
reason: "manual" as const,
reason: "manual",
inputID: input.inputID,
started: input.started,
}
return context.model.compaction?.type === "native" ? executeProvider(request) : execute(request)
},
}),
}),
)
})
+27 -1
View File
@@ -44,6 +44,14 @@ const IMAGE_BYTES_TARGET = 15 * 1024 * 1024 // 15 MiB
const IMAGE_REMOVED =
"[This image was removed to reduce the request size and is no longer visible. Do not make claims about its contents from memory. If needed, retrieve it again with an available tool or ask the user to attach it again.]"
const GENERATION_KEYS = new Set(Object.keys(GenerationOptions.fields))
// Default output limit caps per request kind. Titles and generate have none and keep the provider default, because
// their reasoning is hard to budget.
const OUTPUT_TOKEN_CAPS: Partial<Record<SessionRequestKind, number>> = { primary: 256_000, compaction: 32_000 }
// Used when the catalog has no output limit for the model.
const OUTPUT_TOKEN_FALLBACK = 32_000
// Prompt text is estimated at about 4 characters per token, which can run low on dense text such as code.
const ESTIMATE_ERROR = 0.15
const OUTPUT_TOKEN_MIN = 1_024
/** Tool errors, plus the user declining a permission or dismissing a question. */
export type ExecuteError = Tool.Error | Permission.DeclinedError | QuestionTool.CancelledError
@@ -69,6 +77,16 @@ export interface Input {
readonly toolChoice?: LLM.RequestInput["toolChoice"]
/** Only the durable runner may use a stateful WebSocket. */
readonly webSocket?: "session"
/** Prompt size, measured by the provider or estimated. The default output limit leaves room for it. */
readonly inputTokens?: { readonly measured: number; readonly estimated: number }
}
/** The default output limit: the catalog limit, capped, and fitted to the room the prompt leaves in the context window. */
export const outputLimit = (limit: Model.Info["limit"], cap: number, inputTokens?: Input["inputTokens"]) => {
const requested = Math.min(limit.output > 0 ? limit.output : OUTPUT_TOKEN_FALLBACK, cap)
if (inputTokens === undefined || limit.context <= 0) return requested
const room = limit.context - inputTokens.measured - Math.ceil(inputTokens.estimated * (1 + ESTIMATE_ERROR))
return Math.min(requested, Math.max(OUTPUT_TOKEN_MIN, room))
}
export const baseTranscript = (input: {
@@ -218,8 +236,16 @@ export const layer = Layer.effect(
const given = new Map(
tools.definitions.map((t) => [{ description: t.description, input: { ...t.inputSchema } }, t] as const),
)
// Hooks see the default output limit and may change or remove it.
const cap = OUTPUT_TOKEN_CAPS[kind]
const shaped = yield* shape(
{ sessionID: session.id, model: model.ref, system: input.system, messages: input.messages, options: {} },
{
sessionID: session.id,
model: model.ref,
system: input.system,
messages: input.messages,
options: cap === undefined ? {} : { maxTokens: outputLimit(model.limit, cap, input.inputTokens) },
},
Object.fromEntries(Array.from(given, ([d, t]) => [t.name, d])),
)
// Match by identity first, then by key. Entries matching neither were invented by a
+5
View File
@@ -244,6 +244,11 @@ const layer = Layer.effect(
// Keep tool definitions on the final Step to preserve the provider's cached prefix.
toolChoice: stepLimitReached ? "none" : undefined,
webSocket: "session",
inputTokens: SessionCompaction.estimatePrompt({
messages: loaded.messages,
resolved: loaded.model,
context: loaded,
}),
})
const outcome = yield* steps.attempt({
isLocationClosed: lifecycle.isClosed,
@@ -216,6 +216,31 @@ describe("OpenAIPlugin", () => {
}),
)
it.effect("omits the default output limit from OpenAI steps and compaction", () =>
Effect.gen(function* () {
yield* addPlugin()
const hooks = yield* PluginHooks.Service
const maxTokens = (providerID: Provider.ID) =>
Effect.gen(function* () {
const draft = {
sessionID: Session.ID.make("ses_test"),
model: Model.Ref.make({ providerID, id: Model.ID.make("gpt-5.5") }),
system: [],
messages: [],
options: { maxTokens: 128_000 },
}
const events = [
yield* hooks.trigger("session", "context", { ...draft, agent: Agent.ID.make("build"), tools: {} }),
yield* hooks.trigger("session", "compaction", { ...draft, agent: Agent.ID.make("build"), tools: {} }),
]
return events.map((event) => event.options.maxTokens)
})
expect(yield* maxTokens(Provider.ID.openai)).toEqual([undefined, undefined])
expect(yield* maxTokens(Provider.ID.azure)).toEqual([128_000, 128_000])
}),
)
it.effect("selects WebSocket only from explicit policy", () =>
Effect.gen(function* () {
const credentials = yield* Credential.Service
@@ -1,5 +1,5 @@
import { expect, test } from "bun:test"
import { LLMClient, LLMEvent, LanguageModel, ToolDefinition, type LLMRequest } from "@opencode/ai"
import { GenerationOptions, LLMClient, LLMEvent, LanguageModel, ToolDefinition, type LLMRequest } from "@opencode/ai"
import { OpenAIChat } from "@opencode/ai/protocols"
import { Database } from "@opencode/core/database/database"
import { AppNodeBuilder } from "@opencode/core/effect/app-node-builder"
@@ -407,7 +407,7 @@ it.effect("manual compaction summarizes short context instead of no-op", () =>
"x-opencode-session": sessionID,
"x-opencode-client": "opencode",
})
expect(requests[0]?.generation).toBeUndefined()
expect(requests[0]?.generation).toEqual(GenerationOptions.make({ maxTokens: 32_000 }))
expect(JSON.stringify(requests[0]?.messages)).toContain("Manual compaction should include this short conversation.")
expect(JSON.stringify(requests[0]?.messages)).toContain("Use Effect services and generators.")
expect(JSON.stringify(requests[0]?.messages)).toContain("User shell pwd completed: /project")
@@ -154,3 +154,61 @@ describe("SessionModelRequest HTTP hooks", () => {
}),
)
})
describe("SessionModelRequest output limit", () => {
const input = { session, agent: Agent.ID.make("build"), model, system: [], messages: [] }
it.effect("caps the default output limit per request kind", () =>
Effect.gen(function* () {
const requests = yield* SessionModelRequest.Service.pipe(Effect.provide(SessionModelRequest.layer))
const large = {
...input,
model: SessionRunnerModel.resolved(OpenAIChat.route.model({ id: "large-output", provider: "test" }), {
capabilities: { tools: true, input: ["text"], output: ["text"] },
cost: [],
limit: { context: 1_000_000, output: 384_000 },
}),
}
const maxTokens = (prepared: SessionModelRequest.Prepared<unknown>) => prepared.request.generation?.maxTokens
expect(maxTokens(yield* requests.primary(large))).toBe(256_000)
expect(maxTokens(yield* requests.compaction(large))).toBe(32_000)
const inputTokens = { measured: 170_000, estimated: 8_000 }
expect(maxTokens(yield* requests.primary({ ...input, inputTokens }))).toBe(20_800)
expect(maxTokens(yield* requests.compaction({ ...input, inputTokens }))).toBe(20_800)
}).pipe(Effect.provideService(SessionModelTransport.Service, transport)),
)
// Provider plugins that remove the default limit only hook `context` and `compaction`. If titles or generate get a
// default, also hook `title` and `generate` in: the OpenAI plugin (`omitOutputLimit`), whose ChatGPT backend
// rejects any requested limit.
it.effect("sends no output limit for titles and generate by default", () =>
Effect.gen(function* () {
const requests = yield* SessionModelRequest.Service.pipe(Effect.provide(SessionModelRequest.layer))
expect((yield* requests.title(input)).request.generation?.maxTokens).toBeUndefined()
expect((yield* requests.generate(input)).request.generation?.maxTokens).toBeUndefined()
}).pipe(Effect.provideService(SessionModelTransport.Service, transport)),
)
it.effect("lets hooks change or remove the default output limit", () =>
Effect.gen(function* () {
const hooks = yield* PluginHooks.Service
const seen: Array<number | undefined> = []
yield* hooks.register("session", "context", (event) =>
Effect.sync(() => {
seen.push(event.options.maxTokens)
delete event.options.maxTokens
}),
)
yield* hooks.register("session", "title", (event) =>
Effect.sync(() => {
event.options.maxTokens = 100
}),
)
const requests = yield* SessionModelRequest.Service.pipe(Effect.provide(SessionModelRequest.layer))
expect((yield* requests.primary(input)).request.generation).toBeUndefined()
expect((yield* requests.title(input)).request.generation?.maxTokens).toBe(100)
expect(seen).toEqual([32_000])
}).pipe(Effect.provideService(SessionModelTransport.Service, transport)),
)
})
@@ -1,9 +1,40 @@
import { describe, expect, test } from "bun:test"
import { Message, ToolResultPart, Media } from "@opencode/ai"
import { boundImages, unsupportedParts } from "@opencode/core/session/model-request"
import { boundImages, outputLimit, unsupportedParts } from "@opencode/core/session/model-request"
const capabilities = (input: string[]) => ({ tools: true, input, output: ["text"] })
describe("SessionModelRequest.outputLimit", () => {
test("requests the catalog output limit up to the cap", () => {
expect(outputLimit({ context: 1_000_000, output: 128_000 }, 256_000)).toBe(128_000)
expect(outputLimit({ context: 200_000, output: 64_000 }, 256_000)).toBe(64_000)
expect(outputLimit({ context: 1_048_576, output: 1_048_576 }, 256_000)).toBe(256_000)
expect(outputLimit({ context: 200_000, output: 64_000 }, 32_000)).toBe(32_000)
})
test("falls back to 32k when the catalog has no output limit", () => {
expect(outputLimit({ context: 200_000, output: 0 }, 256_000)).toBe(32_000)
})
test("fits the limit to the room the prompt leaves in the context window", () => {
const limit = { context: 1_000_000, output: 128_000 }
expect(outputLimit(limit, 256_000, { measured: 50_000, estimated: 0 })).toBe(128_000)
expect(outputLimit(limit, 256_000, { measured: 900_000, estimated: 0 })).toBe(100_000)
// Estimated text counts 15% extra, so 40k estimated takes 46k of the room.
expect(outputLimit(limit, 256_000, { measured: 900_000, estimated: 40_000 })).toBe(54_000)
})
test("keeps a minimum limit when the prompt nearly fills the context window", () => {
const prompt = { measured: 199_000, estimated: 0 }
expect(outputLimit({ context: 200_000, output: 64_000 }, 256_000, prompt)).toBe(1_024)
expect(outputLimit({ context: 200_000, output: 512 }, 256_000, prompt)).toBe(512)
})
test("ignores the prompt size when the context window is unknown", () => {
expect(outputLimit({ context: 0, output: 32_000 }, 256_000, { measured: 500_000, estimated: 0 })).toBe(32_000)
})
})
describe("SessionModelRequest.unsupportedParts", () => {
test("replaces unsupported user media with a visible error", () => {
const messages = unsupportedParts(
@@ -393,13 +393,15 @@ it.live("only known automatic native overflow falls back locally and failed reco
fixture.state.overflow = true
fixture.state.localFailure = true
expect(yield* fixture.automatic).toMatchObject({ status: "failed" })
expect(fixture.state.calls).toBe(5)
expect(fixture.state.calls).toBe(6)
expect(yield* fixture.checkpoint).toEqual(installed)
expect(JSON.stringify(fixture.bodies[4])).toContain("Original durable request")
expect(JSON.stringify(fixture.bodies[4])).not.toContain("encrypted_1")
expect(JSON.stringify(fixture.bodies[5])).toContain("Original durable request")
expect(JSON.stringify(fixture.bodies[5])).not.toContain("encrypted_1")
fixture.state.localFailure = false
expect(yield* fixture.automatic).toEqual({ status: "completed", recoveredOverflow: true })
expect(fixture.state.calls).toBe(7)
expect(fixture.state.calls).toBe(8)
expect((yield* fixture.load).messages).toContainEqual(
expect.objectContaining({ type: "compaction", summary: "## Objective\n- Recovered locally" }),
)
+113 -12
View File
@@ -127,6 +127,7 @@ const fullOutputModel = testModel("full-output", { context: 262_144, output: 262
const unknownContextModel = testModel("unknown-context", { context: 0, output: 32_000 })
const undersizedContextModel = testModel("undersized-context", { context: 1, output: 1_000 })
const recoveryModel = testModel("recovery", { context: 200_000, output: 1_000 })
const fittedOutputModel = testModel("fitted-output", { context: 100_000, output: 64_000 })
test("calculates step cost using the matching context tier", () => {
expect(
@@ -2749,22 +2750,16 @@ describe("SessionRunnerLLM", () => {
})
}
for (const response of ["length", "content-filter", "context overflow"] as const) {
for (const response of ["length", "content-filter"] as const) {
scenario(`rejects compaction ${response} without retrying or committing its draft`, function* (s) {
yield* s.llm.push(TestLLM.text("Earlier answer", "history"))
yield* s.runPrompt("Earlier question")
s.requests.length = 0
yield* s.llm.push(
response === "context overflow"
? Stream.fail(
new AIError({
reason: new InvalidRequestError({ message: "Too long", classification: "context-overflow" }),
}),
)
: TestLLM.complete(
{ reason: { normalized: response } },
LLMEvent.textDelta({ id: "truncated", text: "## Objective\n- Incomplete summary" }),
),
TestLLM.complete(
{ reason: { normalized: response } },
LLMEvent.textDelta({ id: "truncated", text: "## Objective\n- Incomplete summary" }),
),
)
const compaction = yield* s.session.compact({ sessionID })
yield* s.resume
@@ -2778,6 +2773,100 @@ describe("SessionRunnerLLM", () => {
})
}
scenario("stops after three smaller compaction inputs overflow", function* (s) {
yield* s.llm.push(...Array.from({ length: 8 }, (_, index) => TestLLM.text(`Answer ${index}`, `answer-${index}`)))
yield* Effect.forEach(
Array.from({ length: 8 }, (_, index) => index),
(index) => s.runPrompt(`Request ${index}: ${"context ".repeat(30)}`),
)
s.currentModel = unknownContextModel
s.requests.length = 0
const overflow = () =>
Stream.fail(
new AIError({ reason: new InvalidRequestError({ message: "Too long", classification: "context-overflow" }) }),
)
yield* s.llm.push(overflow(), overflow(), overflow(), overflow())
const compaction = yield* s.session.compact({ sessionID })
yield* s.resume
expect(s.requests).toHaveLength(4)
expect(userTexts(s.requests[2])[0].length).toBeLessThan(userTexts(s.requests[1])[0].length)
expect(userTexts(s.requests[3])[0].length).toBeLessThan(userTexts(s.requests[2])[0].length)
expect((yield* s.messages).find((message) => message.id === compaction.id)).toMatchObject({ status: "failed" })
yield* s.llm.push(TestLLM.text("Continued", "continued"))
yield* s.runPrompt("Continue")
expect(userTexts(s.requests[4])).toContain("Request 0: " + "context ".repeat(30))
})
scenario("serializes history and omits media after a summary input overflow", function* (s) {
const image = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII="
yield* s.session.prompt({
sessionID,
text: "Earlier question",
files: [{ uri: `data:image/png;base64,${image}` }],
resume: false,
})
yield* s.llm.push(
TestLLM.stop(
LLMEvent.toolCall({ id: "hosted", name: "web_search", input: { query: "earlier" }, providerExecuted: true }),
LLMEvent.toolResult({
id: "hosted",
name: "web_search",
result: { type: "text", value: "x".repeat(5_000) },
providerExecuted: true,
}),
LLMEvent.textStart({ id: "history" }),
LLMEvent.textDelta({ id: "history", text: "Earlier answer" }),
LLMEvent.textEnd({ id: "history" }),
),
)
yield* s.resume
s.requests.length = 0
yield* s.llm.push(
[LLMEvent.providerError({ message: "Too long", classification: "context-overflow" })],
TestLLM.text("## Objective\n- Recovered", "summary"),
)
const compaction = yield* s.session.compact({ sessionID })
yield* s.resume
expect(s.requests).toHaveLength(2)
expect(s.requests[0]?.messages.some((message) => message.content.some((part) => part.type === "media"))).toBeTrue()
expect(s.requests[1]?.messages.every((message) => message.role === "user")).toBeTrue()
expect(userTexts(s.requests[1])[0]).toContain("[Attached image/png; content omitted]")
expect(userTexts(s.requests[1])[0]).toContain("[Tool result web_search]:")
expect(userTexts(s.requests[1])[0]).toContain("[truncated]")
expect(userTexts(s.requests[1])[0]).not.toContain("x".repeat(2_000))
expect(s.requests[1]?.system).toEqual(s.requests[0]?.system)
expect(s.requests[1]?.tools).toEqual(s.requests[0]?.tools)
expect((yield* s.messages).find((message) => message.id === compaction.id)).toMatchObject({ status: "completed" })
})
scenario("fits serialized compaction history by omitting oldest exchanges", function* (s) {
const service = yield* SessionCompaction.Service
yield* service.transform((editor) => editor.configure({ buffer: 3_000 }))
s.currentModel = testModel("large-history", { context: 1_000_000, output: 32_000 })
yield* s.llm.push(...Array.from({ length: 4 }, (_, index) => TestLLM.text(`Answer ${index}`, `answer-${index}`)))
yield* Effect.forEach(
Array.from({ length: 4 }, (_, index) => index),
(index) => s.runPrompt(`Request ${index}: ${"x".repeat(8_000)}`),
)
s.currentModel = testModel("smaller-history", { context: 7_000, output: 1_000 })
s.requests.length = 0
yield* s.llm.push(TestLLM.text("## Objective\n- Recovered", "summary"))
const compaction = yield* s.session.compact({ sessionID })
yield* s.resume
expect(s.requests).toHaveLength(1)
expect(userTexts(s.requests[0])[0]).toContain("older exchanges omitted")
expect(userTexts(s.requests[0])[0]).not.toContain("Request 0:")
expect(userTexts(s.requests[0])[0]).not.toContain("Request 1:")
expect(userTexts(s.requests[0])[0]).toContain("Request 2:")
expect((yield* s.messages).find((message) => message.id === compaction.id)).toMatchObject({
status: "completed",
summary: expect.stringContaining("older exchanges were omitted from the summary input"),
})
})
scenario("records cancelled manual compaction without surfacing an internal failure", function* (s) {
yield* s.llm.push(TestLLM.text("Earlier answer", "text-manual-interrupt-history"))
yield* s.runPrompt("Earlier question")
@@ -2985,7 +3074,7 @@ describe("SessionRunnerLLM", () => {
expect(yield* Effect.exit(s.resume)).toMatchObject({ _tag: "Failure" })
expect(s.requests).toHaveLength(1)
expect(s.requests[0]?.generation).toBeUndefined()
expect(s.requests[0]?.generation?.maxTokens).toBe(50)
expect(yield* s.context).toContainEqual(
expect.objectContaining({
type: "compaction",
@@ -3190,6 +3279,18 @@ describe("SessionRunnerLLM", () => {
])
})
scenario("fits the output limit to the prompt size", function* (s) {
s.currentModel = fittedOutputModel
yield* s.llm.push(TestLLM.textWithUsage("Earlier answer", "text-fitted-first", 50_000))
yield* s.runPrompt("Earlier question")
yield* s.llm.push(TestLLM.text("Continued", "text-fitted-final"))
yield* s.runPrompt("Continue")
expect(s.requests[0]?.generation?.maxTokens).toBe(64_000)
expect(s.requests[1]?.generation?.maxTokens).toBeLessThan(100_000 - 50_000)
expect(s.requests[1]?.generation?.maxTokens).toBeGreaterThan(100_000 - 50_000 - 100)
})
scenario("publishes the original overflow when recovery summarization fails", function* (s) {
yield* setupOverflowRecovery(s)
yield* s.llm.push(