From b0284d2e8a8b54ecf7c896f8582fdda7bfb90776 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 27 Aug 2026 22:59:20 -0400 Subject: [PATCH 1/2] refactor(core): centralize manual compaction settlement --- packages/core/src/session/compaction.ts | 77 +++++++++++------- packages/core/src/session/runner/llm.ts | 33 +++----- packages/core/test/config/compaction.test.ts | 15 +++- packages/core/test/session-compaction.test.ts | 52 ++++++++---- packages/core/test/session-runner.test.ts | 81 +++++++++++++++++++ 5 files changed, 185 insertions(+), 73 deletions(-) diff --git a/packages/core/src/session/compaction.ts b/packages/core/src/session/compaction.ts index c04e13bca4b4..3542653d4d99 100644 --- a/packages/core/src/session/compaction.ts +++ b/packages/core/src/session/compaction.ts @@ -3,12 +3,13 @@ export * as SessionCompaction from "./compaction.js" import { LLMClient, AIError, LLMEvent, Message, type LLMRequest } from "@opencode-ai/ai" import type { StreamOptions } from "@opencode-ai/ai/route" import { SessionError } from "@opencode-ai/schema/session-error" -import { Context, Effect, Layer, Stream } from "effect" +import { Cause, Context, Effect, Exit, Layer, Stream } from "effect" import { Bus } from "../bus.js" import { makeLocationNode } from "@opencode-ai/util/effect/app-node" import { llmClient } from "../effect/app-node-platform.js" import { SessionEvent } from "./event.js" import type { SessionContext } from "./context.js" +import type { MessageDecodeError } from "./error.js" import type { SessionMessage } from "./message.js" import type { SessionModelRequest } from "./model-request.js" import type { SessionRunnerModel } from "./runner/model.js" @@ -83,9 +84,9 @@ type RequiredInput = Pick export type ManualInput = { readonly session: SessionSchema.Info - readonly messages: readonly SessionMessage.Info[] + readonly messages: Effect.Effect readonly inputID: SessionMessage.ID - readonly started?: boolean + readonly restore: (effect: Effect.Effect) => Effect.Effect /** Invoked after content planning, not when the caller captures the operation. */ readonly resolveModel: SessionContext.Interface["resolveModel"] readonly prepare: SessionModelRequest.Interface["prepare"] @@ -98,7 +99,6 @@ type Plan = { readonly prompt: string readonly recent: string readonly inputID?: SessionMessage.ID - readonly started?: boolean readonly prepare: SessionModelRequest.Interface["prepare"] } @@ -110,7 +110,8 @@ export interface Interface extends State.Transformable { readonly enabled: () => boolean readonly required: (input: RequiredInput) => boolean readonly compact: (input: AutoInput) => Effect.Effect - readonly compactManual: (input: ManualInput) => Effect.Effect + /** Runs an already-started control under its caller's interruption mask. */ + readonly compactManual: (input: ManualInput) => Effect.Effect } export class Service extends Context.Service()("@opencode/SessionCompaction") {} @@ -262,7 +263,7 @@ const make = (dependencies: Dependencies) => { return { status: "failed" as const, error: input.error } }) const execute = Effect.fn("SessionCompaction.execute")(function* (plan: Plan) { - if (!plan.started) + if (plan.reason === "auto") yield* dependencies.bus.publish(SessionEvent.Compaction.Started, { sessionID: plan.session.id, reason: plan.reason, @@ -384,34 +385,50 @@ const make = (dependencies: Dependencies) => { return used >= promptCeiling } const compactManual = Effect.fn("SessionCompaction.compactManual")(function* (input: ManualInput) { - const content = planContent(input.messages, state.get().tokens) - if (!content) - return yield* failed({ - sessionID: input.session.id, - reason: "manual", - error: { type: "compaction.unavailable", message: "Nothing to compact yet" }, - inputID: input.inputID, - }) - const resolved = yield* input.resolveModel(input.session).pipe( - Effect.catch((cause) => - failed({ - sessionID: input.session.id, - reason: "manual", - error: toSessionError(cause), - inputID: input.inputID, + // Install settlement before restoring work; a nested mask would not restore interruptibility. + const compacted = yield* input + .restore( + Effect.gen(function* () { + const content = planContent(yield* input.messages, state.get().tokens) + if (!content) + return yield* failed({ + sessionID: input.session.id, + reason: "manual", + error: { type: "compaction.unavailable", message: "Nothing to compact yet" }, + inputID: input.inputID, + }) + const resolved = yield* input.resolveModel(input.session).pipe( + Effect.catch((cause) => + failed({ + sessionID: input.session.id, + reason: "manual", + error: toSessionError(cause), + inputID: input.inputID, + }), + ), + ) + if ("status" in resolved) return resolved + return yield* execute({ + session: input.session, + resolved, + prepare: input.prepare, + reason: "manual", + inputID: input.inputID, + ...content, + }) }), - ), - ) - if ("status" in resolved) return resolved - return yield* execute({ - session: input.session, - resolved, - prepare: input.prepare, + ) + .pipe(Effect.exit) + if (Exit.isSuccess(compacted)) return compacted.value + yield* failed({ + sessionID: input.session.id, reason: "manual", + error: Cause.hasInterruptsOnly(compacted.cause) + ? { type: "aborted", message: "Compaction cancelled" } + : { type: "compaction.failed", message: Cause.pretty(compacted.cause) }, inputID: input.inputID, - started: input.started, - ...content, }) + return yield* Effect.failCause(compacted.cause) }) return Service.of({ transform: state.transform, diff --git a/packages/core/src/session/runner/llm.ts b/packages/core/src/session/runner/llm.ts index 15550fc6eeff..f2bdd25fffc6 100644 --- a/packages/core/src/session/runner/llm.ts +++ b/packages/core/src/session/runner/llm.ts @@ -1,7 +1,7 @@ export * as SessionRunnerLLM from "./llm.js" import { Message } from "@opencode-ai/ai" -import { Cause, Config, Effect, Exit, FiberMap, Layer, Pull, Schedule } from "effect" +import { Config, Effect, FiberMap, Layer, Pull, Schedule } from "effect" import { Database } from "../../database/database.js" import { Bus } from "../../bus.js" import { InstructionState } from "../instruction-state.js" @@ -132,29 +132,14 @@ const layer = Layer.effect( if (pending?.type === "compaction") { const session = yield* store.get(sessionID) if (!session) return yield* Effect.die(new Error(`Session not found: ${sessionID}`)) - const compacted = yield* restore( - Effect.gen(function* () { - return yield* compaction.compactManual({ - session, - resolveModel: context.resolveModel, - prepare: context.prepare, - messages: yield* store.context(sessionID), - inputID: pending.id, - started: true, - }) - }), - ).pipe(Effect.exit) - if (Exit.isFailure(compacted)) { - yield* bus.publish(SessionEvent.Compaction.Failed, { - sessionID, - reason: "manual", - error: Cause.hasInterruptsOnly(compacted.cause) - ? { type: "aborted", message: "Compaction cancelled" } - : { type: "compaction.failed", message: Cause.pretty(compacted.cause) }, - inputID: pending.id, - }) - return yield* Effect.failCause(compacted.cause) - } + yield* compaction.compactManual({ + session, + resolveModel: context.resolveModel, + prepare: context.prepare, + messages: store.context(sessionID), + inputID: pending.id, + restore, + }) force = false continue } diff --git a/packages/core/test/config/compaction.test.ts b/packages/core/test/config/compaction.test.ts index 938ee7fe1a32..92fa65e9ab0d 100644 --- a/packages/core/test/config/compaction.test.ts +++ b/packages/core/test/config/compaction.test.ts @@ -79,11 +79,17 @@ describe("ConfigCompactionPlugin.Plugin", () => { .subscribe(SessionEvent.Compaction.Started) .pipe(Stream.runHead, Effect.forkScoped({ startImmediately: true })) expect( - yield* compaction.compactManual({ + yield* compaction.compact({ session, - resolveModel: () => Effect.succeed(resolved), + resolved, prepare: modelRequests.prepare, messages: [ + { + id: SessionMessage.ID.create(), + type: "user", + text: "Oldest context ".repeat(5_000), + time: { created: DateTime.makeUnsafe(0) }, + }, { id: SessionMessage.ID.create(), type: "user", @@ -97,10 +103,11 @@ describe("ConfigCompactionPlugin.Plugin", () => { time: { created: DateTime.makeUnsafe(1) }, }, ], - inputID: SessionMessage.ID.make("msg_compaction_manual"), }), ).toEqual({ status: "completed" }) - expect(Option.getOrThrow(yield* Fiber.join(started)).data.recent).toContain("Recent context") + const recent = Option.getOrThrow(yield* Fiber.join(started)).data.recent + expect(recent).toContain("Recent context") + expect(recent).not.toContain("Older context") yield* config.setEntries([ new Document({ diff --git a/packages/core/test/session-compaction.test.ts b/packages/core/test/session-compaction.test.ts index c2167485f31b..d611277fe0ed 100644 --- a/packages/core/test/session-compaction.test.ts +++ b/packages/core/test/session-compaction.test.ts @@ -217,6 +217,34 @@ const insertSession = (id: Session.ID, overrides?: Partial (session ? Effect.succeed(session) : Effect.die(`session missing: ${id}`)))) }) +const compactManual = Effect.fnUntraced(function* ( + input: Pick & { + readonly messages: readonly SessionMessage.Info[] + }, +) { + const compaction = yield* SessionCompaction.Service + const bus = yield* Bus.Service + const modelRequests = yield* SessionModelRequest.Service + return yield* Effect.uninterruptibleMask((restore) => + Effect.gen(function* () { + // Unit fixtures begin at the runner's already-started manual-control boundary. + yield* bus.publish(SessionEvent.Compaction.Started, { + sessionID: input.session.id, + reason: "manual", + recent: "", + inputID: input.inputID, + }) + return yield* compaction.compactManual({ + ...input, + messages: Effect.succeed(input.messages), + resolveModel: () => Effect.succeed(resolved), + prepare: modelRequests.prepare, + restore, + }) + }), + ) +}) + it.effect("manual compaction summarizes short context instead of no-op", () => Effect.gen(function* () { requests = [] @@ -240,17 +268,15 @@ it.effect("manual compaction summarizes short context instead of no-op", () => time: { created: DateTime.makeUnsafe(0) }, } const session = yield* insertSession(sessionID, { parent_id: parentID }) - const modelRequests = yield* SessionModelRequest.Service + yield* compaction.transform((draft) => draft.configure({ auto: false })) const delta = yield* bus .subscribe(SessionEvent.Compaction.Delta) .pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped) yield* Effect.yieldNow expect( - yield* compaction.compactManual({ + yield* compactManual({ session, - resolveModel: () => Effect.succeed(resolved), - prepare: modelRequests.prepare, messages: [userMessage], inputID: SessionMessage.ID.make("msg_manual_compaction"), }), @@ -272,7 +298,7 @@ it.effect("manual compaction summarizes short context instead of no-op", () => 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(yield* store.context(sessionID)).toMatchObject([ - { type: "compaction", reason: "manual", summary: "manual summary", recent: "" }, + { type: "compaction", status: "completed", reason: "manual", summary: "manual summary", recent: "" }, ]) expect(yield* store.get(sessionID)).toMatchObject({ cost: 0.0000233, @@ -297,19 +323,16 @@ it.effect("manual compaction summarizes short context instead of no-op", () => it.effect("forked session compaction reuses the fork root prompt cache key", () => Effect.gen(function* () { requests = [] - const compaction = yield* SessionCompaction.Service + const store = yield* SessionStore.Service const sessionID = Session.ID.make("ses_fork_compaction") const rootID = Session.ID.make("ses_fork_compaction_root") const session = yield* insertSession(sessionID, { fork_session_id: rootID, fork_boundary: { type: "before", messageID: SessionMessage.ID.create() }, }) - const modelRequests = yield* SessionModelRequest.Service expect( - yield* compaction.compactManual({ + yield* compactManual({ session, - resolveModel: () => Effect.succeed(resolved), - prepare: modelRequests.prepare, messages: [ { id: SessionMessage.ID.create(), @@ -321,6 +344,7 @@ it.effect("forked session compaction reuses the fork root prompt cache key", () inputID: SessionMessage.ID.make("msg_fork_compaction"), }), ).toEqual({ status: "completed" }) + expect(yield* store.context(sessionID)).toMatchObject([{ type: "compaction", status: "completed" }]) expect(requests).toHaveLength(1) expect(requests[0]?.promptCacheKey).toBe(rootID) @@ -330,7 +354,7 @@ it.effect("forked session compaction reuses the fork root prompt cache key", () it.effect("keeps session context hooks away from compaction requests", () => Effect.gen(function* () { requests = [] - const compaction = yield* SessionCompaction.Service + const store = yield* SessionStore.Service // Context hooks shape the agent conversation; compaction is not part of it, // so it opts out and the transcript passes through unchanged. const hooks = yield* PluginHooks.Service @@ -340,12 +364,9 @@ it.effect("keeps session context hooks away from compaction requests", () => }), ) const session = yield* insertSession(Session.ID.make("ses_hook_compaction")) - const modelRequests = yield* SessionModelRequest.Service expect( - yield* compaction.compactManual({ + yield* compactManual({ session, - resolveModel: () => Effect.succeed(resolved), - prepare: modelRequests.prepare, messages: [ { id: SessionMessage.ID.create(), @@ -357,6 +378,7 @@ it.effect("keeps session context hooks away from compaction requests", () => inputID: SessionMessage.ID.make("msg_hook_compaction"), }), ).toEqual({ status: "completed" }) + expect(yield* store.context(session.id)).toMatchObject([{ type: "compaction", status: "completed" }]) expect(requests).toHaveLength(1) expect(requests[0]?.system).toEqual([]) diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index 32c3f398fc4d..e546a1c412e9 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -2273,6 +2273,87 @@ describe("SessionRunnerLLM", () => { }) }) + scenario("settles manual compaction interrupted after delivery before summary entry", function* (s) { + yield* s.llm.push(TestLLM.text("Earlier answer", "text-manual-entry-history")) + yield* s.runPrompt("Earlier question") + s.requests.length = 0 + s.modelResolveHook = Effect.die("summary resolution must not start") + + const delivered = yield* Deferred.make() + const release = yield* Effect.acquireRelease(Deferred.make(), (deferred) => + Deferred.succeed(deferred, undefined), + ) + yield* s.bus.listen((event) => + event.type === SessionEvent.Compaction.Started.type + ? Deferred.succeed(delivered, undefined).pipe(Effect.andThen(Deferred.await(release))) + : Effect.void, + ) + const compaction = yield* SessionInbox.admitCompaction(s.db, s.bus, { + id: SessionMessage.ID.create(), + sessionID, + delivery: "steer", + }) + const runner = yield* SessionRunner.Service + const run = yield* runner.drain({ sessionID, force: false }).pipe(Effect.forkChild) + yield* Deferred.await(delivered) + const interruption = yield* Fiber.interrupt(run).pipe(Effect.forkChild({ startImmediately: true })) + yield* Deferred.succeed(release, undefined) + yield* Fiber.join(interruption) + + const exit = yield* Fiber.await(run) + expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true) + expect(s.requests).toHaveLength(0) + expect(yield* SessionInbox.find(s.db, compaction.id)).toBeUndefined() + expect((yield* s.messages).find((message) => message.id === compaction.id)).toMatchObject({ + type: "compaction", + status: "failed", + error: { type: "aborted", message: "Compaction cancelled" }, + }) + expect((yield* recordedEventTypes(sessionID)).slice(-3)).toEqual([ + Bus.versionedType(SessionEvent.InboxDelivered.type, 1), + Bus.versionedType(SessionEvent.Compaction.Started.type, 1), + Bus.versionedType(SessionEvent.Compaction.Failed.type, 1), + ]) + }) + + scenario("allows inbox cancellation while manual compaction is streaming", function* (s) { + yield* s.llm.push(TestLLM.text("Earlier answer", "text-manual-unlocked-history")) + yield* s.runPrompt("Earlier question") + s.requests.length = 0 + + yield* s.llm.push(TestLLM.text("Manual summary", "text-manual-unlocked-summary")) + const summary = yield* s.llm.gate + const compaction = yield* SessionInbox.admitCompaction(s.db, s.bus, { + id: SessionMessage.ID.create(), + sessionID, + delivery: "steer", + }) + const queued = yield* s.session.prompt({ sessionID, text: "Cancel this", delivery: "queue", resume: false }) + const run = yield* s.resume.pipe(Effect.forkChild) + yield* summary.started + yield* Effect.gen(function* () { + const cancellation = yield* s.session + .cancelInbox({ sessionID, inboxID: queued.id }) + .pipe(Effect.forkChild({ startImmediately: true })) + const completion = yield* Fiber.await(cancellation).pipe(Effect.timeout("1 second"), Effect.forkChild) + yield* TestClock.adjust("1 second") + expect(yield* Fiber.join(completion)).toMatchObject({ _tag: "Success" }) + expect(yield* SessionInbox.find(s.db, queued.id)).toBeUndefined() + expect((yield* s.messages).find((message) => message.id === compaction.id)).toMatchObject({ + type: "compaction", + status: "running", + }) + }).pipe(Effect.ensuring(summary.release)) + yield* Fiber.join(run) + + expect(s.requests).toHaveLength(1) + expect((yield* s.messages).find((message) => message.id === compaction.id)).toMatchObject({ + type: "compaction", + status: "completed", + summary: "Manual summary", + }) + }) + scenario("settles an admitted manual compaction when pre-start resolution throws", function* (s) { yield* s.llm.push(TestLLM.text("Earlier answer", "text-manual-resolution-history")) yield* s.runPrompt("Earlier question") From c476515eb11b3e46688d1efd40d1dbefe9c8454d Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 27 Aug 2026 23:15:51 -0400 Subject: [PATCH 2/2] test(core): keep compaction coverage without extraction --- packages/core/src/session/compaction.ts | 77 ++++++++----------- packages/core/src/session/runner/llm.ts | 33 +++++--- packages/core/test/config/compaction.test.ts | 15 +--- packages/core/test/session-compaction.test.ts | 52 ++++--------- 4 files changed, 73 insertions(+), 104 deletions(-) diff --git a/packages/core/src/session/compaction.ts b/packages/core/src/session/compaction.ts index 3542653d4d99..c04e13bca4b4 100644 --- a/packages/core/src/session/compaction.ts +++ b/packages/core/src/session/compaction.ts @@ -3,13 +3,12 @@ export * as SessionCompaction from "./compaction.js" import { LLMClient, AIError, LLMEvent, Message, type LLMRequest } from "@opencode-ai/ai" import type { StreamOptions } from "@opencode-ai/ai/route" import { SessionError } from "@opencode-ai/schema/session-error" -import { Cause, Context, Effect, Exit, Layer, Stream } from "effect" +import { Context, Effect, Layer, Stream } from "effect" import { Bus } from "../bus.js" import { makeLocationNode } from "@opencode-ai/util/effect/app-node" import { llmClient } from "../effect/app-node-platform.js" import { SessionEvent } from "./event.js" import type { SessionContext } from "./context.js" -import type { MessageDecodeError } from "./error.js" import type { SessionMessage } from "./message.js" import type { SessionModelRequest } from "./model-request.js" import type { SessionRunnerModel } from "./runner/model.js" @@ -84,9 +83,9 @@ type RequiredInput = Pick export type ManualInput = { readonly session: SessionSchema.Info - readonly messages: Effect.Effect + readonly messages: readonly SessionMessage.Info[] readonly inputID: SessionMessage.ID - readonly restore: (effect: Effect.Effect) => Effect.Effect + readonly started?: boolean /** Invoked after content planning, not when the caller captures the operation. */ readonly resolveModel: SessionContext.Interface["resolveModel"] readonly prepare: SessionModelRequest.Interface["prepare"] @@ -99,6 +98,7 @@ type Plan = { readonly prompt: string readonly recent: string readonly inputID?: SessionMessage.ID + readonly started?: boolean readonly prepare: SessionModelRequest.Interface["prepare"] } @@ -110,8 +110,7 @@ export interface Interface extends State.Transformable { readonly enabled: () => boolean readonly required: (input: RequiredInput) => boolean readonly compact: (input: AutoInput) => Effect.Effect - /** Runs an already-started control under its caller's interruption mask. */ - readonly compactManual: (input: ManualInput) => Effect.Effect + readonly compactManual: (input: ManualInput) => Effect.Effect } export class Service extends Context.Service()("@opencode/SessionCompaction") {} @@ -263,7 +262,7 @@ const make = (dependencies: Dependencies) => { return { status: "failed" as const, error: input.error } }) const execute = Effect.fn("SessionCompaction.execute")(function* (plan: Plan) { - if (plan.reason === "auto") + if (!plan.started) yield* dependencies.bus.publish(SessionEvent.Compaction.Started, { sessionID: plan.session.id, reason: plan.reason, @@ -385,50 +384,34 @@ const make = (dependencies: Dependencies) => { return used >= promptCeiling } const compactManual = Effect.fn("SessionCompaction.compactManual")(function* (input: ManualInput) { - // Install settlement before restoring work; a nested mask would not restore interruptibility. - const compacted = yield* input - .restore( - Effect.gen(function* () { - const content = planContent(yield* input.messages, state.get().tokens) - if (!content) - return yield* failed({ - sessionID: input.session.id, - reason: "manual", - error: { type: "compaction.unavailable", message: "Nothing to compact yet" }, - inputID: input.inputID, - }) - const resolved = yield* input.resolveModel(input.session).pipe( - Effect.catch((cause) => - failed({ - sessionID: input.session.id, - reason: "manual", - error: toSessionError(cause), - inputID: input.inputID, - }), - ), - ) - if ("status" in resolved) return resolved - return yield* execute({ - session: input.session, - resolved, - prepare: input.prepare, - reason: "manual", - inputID: input.inputID, - ...content, - }) + const content = planContent(input.messages, state.get().tokens) + if (!content) + return yield* failed({ + sessionID: input.session.id, + reason: "manual", + error: { type: "compaction.unavailable", message: "Nothing to compact yet" }, + inputID: input.inputID, + }) + const resolved = yield* input.resolveModel(input.session).pipe( + Effect.catch((cause) => + failed({ + sessionID: input.session.id, + reason: "manual", + error: toSessionError(cause), + inputID: input.inputID, }), - ) - .pipe(Effect.exit) - if (Exit.isSuccess(compacted)) return compacted.value - yield* failed({ - sessionID: input.session.id, + ), + ) + if ("status" in resolved) return resolved + return yield* execute({ + session: input.session, + resolved, + prepare: input.prepare, reason: "manual", - error: Cause.hasInterruptsOnly(compacted.cause) - ? { type: "aborted", message: "Compaction cancelled" } - : { type: "compaction.failed", message: Cause.pretty(compacted.cause) }, inputID: input.inputID, + started: input.started, + ...content, }) - return yield* Effect.failCause(compacted.cause) }) return Service.of({ transform: state.transform, diff --git a/packages/core/src/session/runner/llm.ts b/packages/core/src/session/runner/llm.ts index f2bdd25fffc6..15550fc6eeff 100644 --- a/packages/core/src/session/runner/llm.ts +++ b/packages/core/src/session/runner/llm.ts @@ -1,7 +1,7 @@ export * as SessionRunnerLLM from "./llm.js" import { Message } from "@opencode-ai/ai" -import { Config, Effect, FiberMap, Layer, Pull, Schedule } from "effect" +import { Cause, Config, Effect, Exit, FiberMap, Layer, Pull, Schedule } from "effect" import { Database } from "../../database/database.js" import { Bus } from "../../bus.js" import { InstructionState } from "../instruction-state.js" @@ -132,14 +132,29 @@ const layer = Layer.effect( if (pending?.type === "compaction") { const session = yield* store.get(sessionID) if (!session) return yield* Effect.die(new Error(`Session not found: ${sessionID}`)) - yield* compaction.compactManual({ - session, - resolveModel: context.resolveModel, - prepare: context.prepare, - messages: store.context(sessionID), - inputID: pending.id, - restore, - }) + const compacted = yield* restore( + Effect.gen(function* () { + return yield* compaction.compactManual({ + session, + resolveModel: context.resolveModel, + prepare: context.prepare, + messages: yield* store.context(sessionID), + inputID: pending.id, + started: true, + }) + }), + ).pipe(Effect.exit) + if (Exit.isFailure(compacted)) { + yield* bus.publish(SessionEvent.Compaction.Failed, { + sessionID, + reason: "manual", + error: Cause.hasInterruptsOnly(compacted.cause) + ? { type: "aborted", message: "Compaction cancelled" } + : { type: "compaction.failed", message: Cause.pretty(compacted.cause) }, + inputID: pending.id, + }) + return yield* Effect.failCause(compacted.cause) + } force = false continue } diff --git a/packages/core/test/config/compaction.test.ts b/packages/core/test/config/compaction.test.ts index 92fa65e9ab0d..938ee7fe1a32 100644 --- a/packages/core/test/config/compaction.test.ts +++ b/packages/core/test/config/compaction.test.ts @@ -79,17 +79,11 @@ describe("ConfigCompactionPlugin.Plugin", () => { .subscribe(SessionEvent.Compaction.Started) .pipe(Stream.runHead, Effect.forkScoped({ startImmediately: true })) expect( - yield* compaction.compact({ + yield* compaction.compactManual({ session, - resolved, + resolveModel: () => Effect.succeed(resolved), prepare: modelRequests.prepare, messages: [ - { - id: SessionMessage.ID.create(), - type: "user", - text: "Oldest context ".repeat(5_000), - time: { created: DateTime.makeUnsafe(0) }, - }, { id: SessionMessage.ID.create(), type: "user", @@ -103,11 +97,10 @@ describe("ConfigCompactionPlugin.Plugin", () => { time: { created: DateTime.makeUnsafe(1) }, }, ], + inputID: SessionMessage.ID.make("msg_compaction_manual"), }), ).toEqual({ status: "completed" }) - const recent = Option.getOrThrow(yield* Fiber.join(started)).data.recent - expect(recent).toContain("Recent context") - expect(recent).not.toContain("Older context") + expect(Option.getOrThrow(yield* Fiber.join(started)).data.recent).toContain("Recent context") yield* config.setEntries([ new Document({ diff --git a/packages/core/test/session-compaction.test.ts b/packages/core/test/session-compaction.test.ts index d611277fe0ed..c2167485f31b 100644 --- a/packages/core/test/session-compaction.test.ts +++ b/packages/core/test/session-compaction.test.ts @@ -217,34 +217,6 @@ const insertSession = (id: Session.ID, overrides?: Partial (session ? Effect.succeed(session) : Effect.die(`session missing: ${id}`)))) }) -const compactManual = Effect.fnUntraced(function* ( - input: Pick & { - readonly messages: readonly SessionMessage.Info[] - }, -) { - const compaction = yield* SessionCompaction.Service - const bus = yield* Bus.Service - const modelRequests = yield* SessionModelRequest.Service - return yield* Effect.uninterruptibleMask((restore) => - Effect.gen(function* () { - // Unit fixtures begin at the runner's already-started manual-control boundary. - yield* bus.publish(SessionEvent.Compaction.Started, { - sessionID: input.session.id, - reason: "manual", - recent: "", - inputID: input.inputID, - }) - return yield* compaction.compactManual({ - ...input, - messages: Effect.succeed(input.messages), - resolveModel: () => Effect.succeed(resolved), - prepare: modelRequests.prepare, - restore, - }) - }), - ) -}) - it.effect("manual compaction summarizes short context instead of no-op", () => Effect.gen(function* () { requests = [] @@ -268,15 +240,17 @@ it.effect("manual compaction summarizes short context instead of no-op", () => time: { created: DateTime.makeUnsafe(0) }, } const session = yield* insertSession(sessionID, { parent_id: parentID }) - yield* compaction.transform((draft) => draft.configure({ auto: false })) + const modelRequests = yield* SessionModelRequest.Service const delta = yield* bus .subscribe(SessionEvent.Compaction.Delta) .pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped) yield* Effect.yieldNow expect( - yield* compactManual({ + yield* compaction.compactManual({ session, + resolveModel: () => Effect.succeed(resolved), + prepare: modelRequests.prepare, messages: [userMessage], inputID: SessionMessage.ID.make("msg_manual_compaction"), }), @@ -298,7 +272,7 @@ it.effect("manual compaction summarizes short context instead of no-op", () => 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(yield* store.context(sessionID)).toMatchObject([ - { type: "compaction", status: "completed", reason: "manual", summary: "manual summary", recent: "" }, + { type: "compaction", reason: "manual", summary: "manual summary", recent: "" }, ]) expect(yield* store.get(sessionID)).toMatchObject({ cost: 0.0000233, @@ -323,16 +297,19 @@ it.effect("manual compaction summarizes short context instead of no-op", () => it.effect("forked session compaction reuses the fork root prompt cache key", () => Effect.gen(function* () { requests = [] - const store = yield* SessionStore.Service + const compaction = yield* SessionCompaction.Service const sessionID = Session.ID.make("ses_fork_compaction") const rootID = Session.ID.make("ses_fork_compaction_root") const session = yield* insertSession(sessionID, { fork_session_id: rootID, fork_boundary: { type: "before", messageID: SessionMessage.ID.create() }, }) + const modelRequests = yield* SessionModelRequest.Service expect( - yield* compactManual({ + yield* compaction.compactManual({ session, + resolveModel: () => Effect.succeed(resolved), + prepare: modelRequests.prepare, messages: [ { id: SessionMessage.ID.create(), @@ -344,7 +321,6 @@ it.effect("forked session compaction reuses the fork root prompt cache key", () inputID: SessionMessage.ID.make("msg_fork_compaction"), }), ).toEqual({ status: "completed" }) - expect(yield* store.context(sessionID)).toMatchObject([{ type: "compaction", status: "completed" }]) expect(requests).toHaveLength(1) expect(requests[0]?.promptCacheKey).toBe(rootID) @@ -354,7 +330,7 @@ it.effect("forked session compaction reuses the fork root prompt cache key", () it.effect("keeps session context hooks away from compaction requests", () => Effect.gen(function* () { requests = [] - const store = yield* SessionStore.Service + const compaction = yield* SessionCompaction.Service // Context hooks shape the agent conversation; compaction is not part of it, // so it opts out and the transcript passes through unchanged. const hooks = yield* PluginHooks.Service @@ -364,9 +340,12 @@ it.effect("keeps session context hooks away from compaction requests", () => }), ) const session = yield* insertSession(Session.ID.make("ses_hook_compaction")) + const modelRequests = yield* SessionModelRequest.Service expect( - yield* compactManual({ + yield* compaction.compactManual({ session, + resolveModel: () => Effect.succeed(resolved), + prepare: modelRequests.prepare, messages: [ { id: SessionMessage.ID.create(), @@ -378,7 +357,6 @@ it.effect("keeps session context hooks away from compaction requests", () => inputID: SessionMessage.ID.make("msg_hook_compaction"), }), ).toEqual({ status: "completed" }) - expect(yield* store.context(session.id)).toMatchObject([{ type: "compaction", status: "completed" }]) expect(requests).toHaveLength(1) expect(requests[0]?.system).toEqual([])