Skip to content

Commit ca14d0f

Browse files
committed
fix(session): harden suffix compaction
1 parent 05f3598 commit ca14d0f

10 files changed

Lines changed: 378 additions & 146 deletions

File tree

‎packages/core/src/session/compaction.ts‎

Lines changed: 97 additions & 69 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,15 @@
11
export * as SessionCompaction from "./compaction"
22

3-
import { LLM, LLMError, LLMEvent, Message, type LLMRequest, type Model, type Usage } from "@opencode-ai/llm"
3+
import {
4+
LLM,
5+
LLMError,
6+
LLMEvent,
7+
Message,
8+
mergeGenerationOptions,
9+
type LLMRequest,
10+
type Model,
11+
type Usage,
12+
} from "@opencode-ai/llm"
413
import { SessionCompaction } from "@opencode-ai/schema/session-compaction"
514
import { DateTime, Effect, Stream } from "effect"
615
import type { Config } from "../config"
@@ -197,13 +206,34 @@ export const make = (dependencies: Dependencies) => {
197206
reason: "auto",
198207
requested,
199208
})
209+
let attempted: SessionCompaction.Mode | undefined
210+
let attemptedUsage: Usage | undefined
211+
let terminal = false
212+
const diagnostics = (
213+
input: { readonly fallback?: SessionCompaction.FallbackReason; readonly usage?: Usage } = {},
214+
) =>
215+
Effect.gen(function* () {
216+
const diagnosticTokens = tokens(input.usage ?? attemptedUsage)
217+
return {
218+
requested,
219+
...(attempted ? { used: attempted } : {}),
220+
...(input.fallback ? { fallback: input.fallback } : {}),
221+
durationMs: Math.max(
222+
0,
223+
Math.round(DateTime.toEpochMillis(yield* DateTime.now) - DateTime.toEpochMillis(started)),
224+
),
225+
...(diagnosticTokens ? { tokens: diagnosticTokens } : {}),
226+
}
227+
})
200228
const run = (
201229
request: LLMRequest,
202-
validate: boolean,
230+
mode: SessionCompaction.Mode,
203231
): Effect.Effect<
204232
{ summary: string; usage: Usage | undefined } | { failure: SessionCompaction.FallbackReason; usage?: Usage }
205233
> =>
206234
Effect.gen(function* () {
235+
attempted = mode
236+
attemptedUsage = undefined
207237
const chunks: string[] = []
208238
let failure: SessionCompaction.FallbackReason | undefined
209239
let usage: Usage | undefined
@@ -212,37 +242,19 @@ export const make = (dependencies: Dependencies) => {
212242
if (LLMEvent.is.providerError(event)) failure = "provider_error"
213243
if (LLMEvent.is.toolCall(event)) failure = "tool_call"
214244
if (LLMEvent.is.textDelta(event)) chunks.push(event.text)
215-
if (LLMEvent.is.stepFinish(event) || LLMEvent.is.finish(event)) usage = event.usage
245+
if (LLMEvent.is.stepFinish(event) || LLMEvent.is.finish(event)) {
246+
usage = event.usage
247+
attemptedUsage = usage
248+
}
216249
return Effect.void
217250
}),
218251
Effect.as(true),
219252
Effect.catchTag("LLM.Error", () => Effect.succeed(false)),
220-
Effect.onInterrupt(() =>
221-
Effect.gen(function* () {
222-
const diagnosticTokens = tokens(usage)
223-
yield* dependencies.events.publish(SessionEvent.Compaction.Failed, {
224-
sessionID: input.sessionID,
225-
messageID,
226-
timestamp: yield* DateTime.now,
227-
reason: "auto",
228-
failure: "interrupted",
229-
diagnostics: {
230-
requested,
231-
used: validate ? "suffix" : "prepend",
232-
durationMs: Math.max(
233-
0,
234-
Math.round(DateTime.toEpochMillis(yield* DateTime.now) - DateTime.toEpochMillis(started)),
235-
),
236-
...(diagnosticTokens ? { tokens: diagnosticTokens } : {}),
237-
},
238-
})
239-
}),
240-
),
241253
)
242254
const summary = chunks.join("")
243255
if (!streamed || failure) return { failure: failure ?? "provider_error", usage }
244256
if (!summary.trim()) return { failure: "empty_summary" as const, usage }
245-
if (validate && !SessionCompactionSuffix.validateSummary(summary))
257+
if (mode === "suffix" && !SessionCompactionSuffix.validateSummary(summary))
246258
return { failure: "invalid_summary" as const, usage }
247259
return { summary, usage }
248260
})
@@ -256,7 +268,7 @@ export const make = (dependencies: Dependencies) => {
256268
tools: [],
257269
generation: { maxTokens: summaryOutput },
258270
}),
259-
false,
271+
"prepend",
260272
)
261273
: Effect.succeed({ failure: "context" as const })
262274
const suffixPrompt = SessionCompactionSuffix.buildPrompt({
@@ -271,7 +283,7 @@ export const make = (dependencies: Dependencies) => {
271283
})
272284
const suffixRequest = LLM.updateRequest(input.request, {
273285
messages: [...input.request.messages, Message.user(suffixPrompt)],
274-
generation: { ...input.request.generation, maxTokens: summaryOutput },
286+
generation: mergeGenerationOptions(input.request.generation, { maxTokens: summaryOutput }),
275287
})
276288
const suffixUnavailable =
277289
input.request.toolChoice && !["auto", "none"].includes(input.request.toolChoice.type)
@@ -287,49 +299,65 @@ export const make = (dependencies: Dependencies) => {
287299
toolChoice: suffixRequest.toolChoice,
288300
}) <=
289301
context - summaryOutput
290-
const suffix =
291-
requested === "suffix" && !suffixUnavailable && suffixFits ? yield* run(suffixRequest, true) : undefined
292-
const suffixSuccess = suffix && "summary" in suffix
293-
const final = suffixSuccess ? suffix : yield* prepend()
294-
const fallback =
295-
(suffix && "failure" in suffix ? suffix.failure : undefined) ??
296-
(requested === "suffix" ? (suffixUnavailable ?? (!suffixFits ? "context" : undefined)) : undefined)
297-
const suffixUsage = suffix && "usage" in suffix ? suffix.usage : undefined
298-
const finalUsage = "usage" in final ? final.usage : undefined
299-
const diagnosticTokens = tokens(suffixUsage ?? finalUsage)
300-
const diagnostics = {
301-
requested,
302-
used: suffixSuccess ? ("suffix" as const) : ("prepend" as const),
303-
...(fallback ? { fallback } : {}),
304-
durationMs: Math.max(
305-
0,
306-
Math.round(DateTime.toEpochMillis(yield* DateTime.now) - DateTime.toEpochMillis(started)),
307-
),
308-
...(diagnosticTokens ? { tokens: diagnosticTokens } : {}),
309-
}
310-
if (!("summary" in final)) {
311-
yield* Effect.logWarning("session compaction failed", { failure: final.failure, ...diagnostics })
312-
yield* dependencies.events.publish(SessionEvent.Compaction.Failed, {
313-
sessionID: input.sessionID,
314-
messageID,
315-
timestamp: yield* DateTime.now,
316-
reason: "auto",
317-
failure: final.failure,
318-
diagnostics,
319-
})
320-
return false
321-
}
322-
yield* Effect.logInfo("session compaction completed", diagnostics)
323-
yield* dependencies.events.publish(SessionEvent.Compaction.Ended, {
324-
sessionID: input.sessionID,
325-
messageID,
326-
timestamp: yield* DateTime.now,
327-
reason: "auto",
328-
text: final.summary,
329-
recent: selected.recent,
330-
diagnostics,
302+
const execute = Effect.gen(function* () {
303+
const suffix =
304+
requested === "suffix" && !suffixUnavailable && suffixFits ? yield* run(suffixRequest, "suffix") : undefined
305+
const suffixSuccess = suffix && "summary" in suffix
306+
const final = suffixSuccess ? suffix : yield* prepend()
307+
const fallback =
308+
(suffix && "failure" in suffix ? suffix.failure : undefined) ??
309+
(requested === "suffix" ? (suffixUnavailable ?? (!suffixFits ? "context" : undefined)) : undefined)
310+
const suffixUsage = suffix && "usage" in suffix ? suffix.usage : undefined
311+
const finalUsage = "usage" in final ? final.usage : undefined
312+
const completedDiagnostics = yield* diagnostics({ fallback, usage: suffixUsage ?? finalUsage })
313+
if (!("summary" in final)) {
314+
yield* Effect.logWarning("session compaction failed", { failure: final.failure, ...completedDiagnostics })
315+
terminal = true
316+
yield* Effect.uninterruptible(
317+
dependencies.events.publish(SessionEvent.Compaction.Failed, {
318+
sessionID: input.sessionID,
319+
messageID,
320+
timestamp: yield* DateTime.now,
321+
reason: "auto",
322+
failure: final.failure,
323+
diagnostics: completedDiagnostics,
324+
}),
325+
)
326+
return false
327+
}
328+
yield* Effect.logInfo("session compaction completed", completedDiagnostics)
329+
terminal = true
330+
yield* Effect.uninterruptible(
331+
dependencies.events.publish(SessionEvent.Compaction.Ended, {
332+
sessionID: input.sessionID,
333+
messageID,
334+
timestamp: yield* DateTime.now,
335+
reason: "auto",
336+
text: final.summary,
337+
recent: selected.recent,
338+
diagnostics: completedDiagnostics,
339+
}),
340+
)
341+
return true
331342
})
332-
return true
343+
return yield* execute.pipe(
344+
Effect.onInterrupt(() =>
345+
terminal
346+
? Effect.void
347+
: Effect.uninterruptible(
348+
Effect.gen(function* () {
349+
yield* dependencies.events.publish(SessionEvent.Compaction.Failed, {
350+
sessionID: input.sessionID,
351+
messageID,
352+
timestamp: yield* DateTime.now,
353+
reason: "auto",
354+
failure: "interrupted",
355+
diagnostics: yield* diagnostics(),
356+
})
357+
}),
358+
),
359+
),
360+
)
333361
})
334362
const compactIfNeeded = Effect.fn("SessionCompaction.compactIfNeeded")(function* (input: Input) {
335363
if (!config.auto) return false

‎packages/core/test/session-compaction.test.ts‎

Lines changed: 102 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { expect, test } from "bun:test"
22
import { DateTime, Effect, Fiber, Stream } from "effect"
3-
import { Finish, LLM, Message, Model, TextDelta, ToolDefinition, Usage } from "@opencode-ai/llm"
3+
import { Finish, LLM, mergeGenerationOptions, Message, Model, TextDelta, ToolDefinition, Usage } from "@opencode-ai/llm"
44
import { route } from "@opencode-ai/llm/protocols/openai-chat"
55
import { Config } from "@opencode-ai/core/config"
66
import { ConfigCompaction } from "@opencode-ai/core/config/compaction"
@@ -265,20 +265,77 @@ test("suffix compaction preserves the request prefix and reports failed suffix u
265265
),
266266
).toBe(true)
267267
expect(requests).toHaveLength(2)
268-
expect([...requests[0]!.messages.slice(0, -1)]).toEqual([...request.messages])
268+
expect(requests[0].messages.slice(0, -1)).toEqual([...request.messages])
269269
expect(requests[0]?.messages).toHaveLength(request.messages.length + 1)
270270
expect(requests[0]?.system).toEqual(request.system)
271271
expect(requests[0]?.tools).toEqual(request.tools)
272272
expect(requests[0]?.toolChoice).toEqual(request.toolChoice)
273273
expect(requests[0]?.http).toEqual(request.http)
274274
expect(requests[0]?.providerOptions).toEqual(request.providerOptions)
275-
expect(requests[0]?.generation).toEqual({ ...request.generation, maxTokens: 4_096 })
275+
expect(requests[0]?.generation).toEqual(mergeGenerationOptions(request.generation, { maxTokens: 4_096 }))
276276
expect(published.at(-1)).toMatchObject({
277277
type: "session.next.compaction.ended",
278278
data: { diagnostics: { requested: "suffix", used: "prepend", fallback: "invalid_summary", tokens: { input: 7 } } },
279279
})
280280
})
281281

282+
test("valid suffix compaction streams once and retains the request prefix", async () => {
283+
const requests: Array<ReturnType<typeof LLM.request>> = []
284+
const published: Array<{ type: string; data: unknown }> = []
285+
const compaction = SessionCompaction.make({
286+
config: [
287+
new Config.Document({
288+
type: "document",
289+
info: new Config.Info({
290+
compaction: new ConfigCompaction.Info({ mode: "suffix", keep: new ConfigCompaction.Keep({ tokens: 0 }) }),
291+
}),
292+
}),
293+
],
294+
events: {
295+
publish: (definition, data) =>
296+
Effect.sync(() => published.push({ type: definition.type, data })).pipe(Effect.as(data as never)),
297+
},
298+
llm: {
299+
stream(request) {
300+
requests.push(request)
301+
return Stream.fromIterable([
302+
TextDelta.make({ type: "text-delta", id: "txt_1", text: validSummary }),
303+
Finish.make({ type: "finish", reason: "stop" }),
304+
])
305+
},
306+
},
307+
})
308+
const model = Model.make({ id: "model", provider: "provider", route: route.with({ limits: { context: 100_000 } }) })
309+
const request = LLM.request({ model, messages: [Message.user("old context")], generation: { maxTokens: 8_192 } })
310+
311+
expect(
312+
await Effect.runPromise(
313+
compaction.compactAfterOverflow({
314+
sessionID: SessionSchema.ID.make("ses_compaction"),
315+
model,
316+
request,
317+
entries: [
318+
{
319+
seq: 1,
320+
message: SessionMessage.User.make({
321+
id: SessionMessage.ID.make("msg_old"),
322+
type: "user",
323+
text: "old context",
324+
time: { created: DateTime.makeUnsafe(0) },
325+
}),
326+
},
327+
],
328+
}),
329+
),
330+
).toBe(true)
331+
expect(requests).toHaveLength(1)
332+
expect([...(requests[0]?.messages.slice(0, -1) ?? [])]).toEqual([...request.messages])
333+
expect(published.at(-1)).toMatchObject({
334+
type: "session.next.compaction.ended",
335+
data: { text: validSummary, diagnostics: { requested: "suffix", used: "suffix" } },
336+
})
337+
})
338+
282339
test("both strategies failing after started emits one durable compaction failure", async () => {
283340
const published: Array<{ type: string; data: unknown }> = []
284341
const compaction = SessionCompaction.make({
@@ -322,7 +379,8 @@ test("both strategies failing after started emits one durable compaction failure
322379
"session.next.compaction.started",
323380
"session.next.compaction.failed",
324381
])
325-
expect(published.at(-1)).toMatchObject({ data: { diagnostics: { requested: "suffix", used: "prepend" } } })
382+
expect(published.at(-1)).toMatchObject({ data: { diagnostics: { requested: "suffix", fallback: "context" } } })
383+
expect(published.at(-1)).not.toMatchObject({ data: { diagnostics: { used: expect.anything() } } })
326384
})
327385

328386
test("forced tool choice skips suffix and reports prepend fallback", async () => {
@@ -419,3 +477,43 @@ test("interruption publishes a terminal compaction failure", async () => {
419477
)
420478
expect(published).toEqual(["session.next.compaction.started", "session.next.compaction.failed"])
421479
})
480+
481+
test("interruption during terminal publication does not publish a second terminal outcome", async () => {
482+
const published: string[] = []
483+
const compaction = SessionCompaction.make({
484+
config: [],
485+
events: {
486+
publish: (definition, data) =>
487+
Effect.sync(() => published.push(definition.type)).pipe(
488+
Effect.andThen(definition.type === "session.next.compaction.ended" ? Effect.interrupt : Effect.void),
489+
Effect.as(data as never),
490+
),
491+
},
492+
llm: {
493+
stream: () => Stream.fromIterable([TextDelta.make({ type: "text-delta", id: "txt_1", text: "summary" })]),
494+
},
495+
})
496+
const model = Model.make({ id: "model", provider: "provider", route: route.with({ limits: { context: 100_000 } }) })
497+
498+
await Effect.runPromise(
499+
Effect.exit(
500+
compaction.compactAfterOverflow({
501+
sessionID: SessionSchema.ID.make("ses_compaction"),
502+
model,
503+
request: LLM.request({ model, prompt: "old context" }),
504+
entries: [
505+
{
506+
seq: 1,
507+
message: SessionMessage.User.make({
508+
id: SessionMessage.ID.make("msg_old"),
509+
type: "user",
510+
text: "old ".repeat(10_000),
511+
time: { created: DateTime.makeUnsafe(0) },
512+
}),
513+
},
514+
],
515+
}),
516+
),
517+
)
518+
expect(published).toEqual(["session.next.compaction.started", "session.next.compaction.ended"])
519+
})

0 commit comments

Comments
 (0)