From fa5ccac7073fe6c589ad5d08c8f1a24f08bf13bd Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Fri, 28 Aug 2026 21:20:26 -0400 Subject: [PATCH] refactor(core): share subagent completion delivery (#46054) --- .../core/src/session/execution/restart.ts | 30 ++--- .../core/src/session/subagent-completion.ts | 32 ++++++ packages/core/src/tool/plugin/subagent.ts | 52 +++------ packages/core/test/session-execution.test.ts | 45 +++++++- packages/core/test/tool-subagent.test.ts | 106 +++++++++++++++++- 5 files changed, 204 insertions(+), 61 deletions(-) create mode 100644 packages/core/src/session/subagent-completion.ts diff --git a/packages/core/src/session/execution/restart.ts b/packages/core/src/session/execution/restart.ts index d0c0effe42f..1e28ca4f8a4 100644 --- a/packages/core/src/session/execution/restart.ts +++ b/packages/core/src/session/execution/restart.ts @@ -9,6 +9,7 @@ import { SessionEvent } from "../event.js" import { SessionExecution } from "../execution.js" import { SessionSchema } from "../schema.js" import { SessionStore } from "../store.js" +import { SubagentCompletion } from "../subagent-completion.js" const CONTINUE_AFTER_SERVER_RESTART = "The server restarted while you were working. Continue from where you left off without repeating completed work." @@ -143,29 +144,12 @@ export const layer = (options?: Options) => } const notify = Effect.fnUntraced(function* (result: Pick) { - if (result.status === "running") return - const text = - result.status === "completed" - ? (result.output ?? "Subagent completed without a text response.") - : result.status === "error" - ? (result.error ?? "Subagent failed") - : "Subagent cancelled" - yield* sessions - .synthetic({ - id: background.notificationID, - sessionID: recovery.parentSessionID, - ...(suspended.has(recovery.parentSessionID) ? { resume: false } : {}), - description: recovery.description, - text: `\n${text}\n`, - metadata: { - source: "subagent", - childID: recovery.childSessionID, - agent: recovery.agent, - state: result.status, - }, - }) - .pipe(Effect.orDie) - yield* jobs.completeBackground(background.notificationID) + yield* SubagentCompletion.deliver(sessions, jobs, { + ...result, + recovery, + notificationID: background.notificationID, + resume: suspended.has(recovery.parentSessionID) ? false : undefined, + }).pipe(Effect.orDie) }) if (background.status !== "running") { diff --git a/packages/core/src/session/subagent-completion.ts b/packages/core/src/session/subagent-completion.ts new file mode 100644 index 00000000000..9ccf24363f4 --- /dev/null +++ b/packages/core/src/session/subagent-completion.ts @@ -0,0 +1,32 @@ +export * as SubagentCompletion from "./subagent-completion.js" + +import { Effect } from "effect" +import type { Job } from "../job.js" +import type { Session } from "../session.js" + +export const deliver = Effect.fnUntraced(function* ( + sessions: Pick, + jobs: Pick, + input: Pick & { + recovery: Extract + resume?: boolean + }, +) { + if (input.status === "running") return + const recovery = input.recovery + const text = + input.status === "completed" + ? (input.output ?? "Subagent completed without a text response.") + : input.status === "error" + ? (input.error ?? "Subagent failed") + : "Subagent cancelled" + yield* sessions.synthetic({ + ...(input.notificationID ? { id: input.notificationID } : {}), + sessionID: recovery.parentSessionID, + ...(input.resume === false ? { resume: false } : {}), + description: recovery.description, + text: `\n${text}\n`, + metadata: { source: "subagent", childID: recovery.childSessionID, agent: recovery.agent, state: input.status }, + }) + if (input.notificationID) yield* jobs.completeBackground(input.notificationID) +}) diff --git a/packages/core/src/tool/plugin/subagent.ts b/packages/core/src/tool/plugin/subagent.ts index 4a43f022710..dd0c949a98a 100644 --- a/packages/core/src/tool/plugin/subagent.ts +++ b/packages/core/src/tool/plugin/subagent.ts @@ -5,9 +5,11 @@ import type { Context } from "@opencode-ai/plugin/effect/plugin" import { Effect, Schema, Scope } from "effect" import { Agent } from "../../agent.js" import { Config } from "../../config.js" +import type { Job } from "../../job.js" import { PluginRuntime } from "../../plugin/runtime.js" import { Permission } from "../../permission.js" import { SessionSchema } from "../../session/schema.js" +import { SubagentCompletion } from "../../session/subagent-completion.js" export const name = "subagent" @@ -79,32 +81,15 @@ export const Plugin = { }) const notifyWhenDone = Effect.fn("SubagentTool.notifyWhenDone")(function* ( - parentID: SessionSchema.ID, - childID: SessionSchema.ID, - agent: string, - description: string, + recovery: Extract, startedAt: number, ) { - const key = `${childID}:${startedAt}` + const key = `${recovery.childSessionID}:${startedAt}` if (notifications.has(key)) return notifications.add(key) yield* Effect.gen(function* () { - const info = (yield* runtime.job.wait({ id: childID })).info - if (!info || info.status === "running") return - const text = - info.status === "completed" - ? (info.output ?? NO_TEXT) - : info.status === "error" - ? (info.error ?? "Subagent failed") - : "Subagent cancelled" - yield* runtime.session.synthetic({ - ...(info.notificationID ? { id: info.notificationID } : {}), - sessionID: parentID, - text: `\n${text}\n`, - description, - metadata: { source: "subagent", childID, agent, state: info.status }, - }) - if (info.notificationID) yield* runtime.job.completeBackground(info.notificationID) + const info = (yield* runtime.job.wait({ id: recovery.childSessionID })).info + if (info) yield* SubagentCompletion.deliver(runtime.session, runtime.job, { ...info, recovery }) }).pipe( Effect.ensuring(Effect.sync(() => notifications.delete(key))), Effect.forkIn(scope, { startImmediately: true }), @@ -232,24 +217,25 @@ export const Plugin = { ), ) + const recovery = { + kind: "subagent" as const, + parentSessionID: context.sessionID, + childSessionID: child.id, + agent: agent.name, + description: input.description, + } const info = yield* runtime.job.start({ id: child.id, type: name, title: input.description, metadata: {}, - recovery: { - kind: "subagent", - parentSessionID: context.sessionID, - childSessionID: child.id, - agent: agent.name, - description: input.description, - }, + recovery, run: runtime.session.resume(child.id).pipe(Effect.andThen(latestAssistantText(child.id))), }) if (background) { yield* runtime.job.background(info.id) - yield* notifyWhenDone(context.sessionID, child.id, agent.name, input.description, info.started_at) + yield* notifyWhenDone(recovery, info.started_at) return backgroundResult(child.id) } @@ -261,13 +247,7 @@ export const Plugin = { ), ) if (result?.type === "backgrounded") { - yield* notifyWhenDone( - context.sessionID, - child.id, - agent.name, - input.description, - result.info.started_at, - ) + yield* notifyWhenDone(recovery, result.info.started_at) return backgroundResult(child.id) } // Failure surfaces keep the sessionID visible so the model can continue the child. diff --git a/packages/core/test/session-execution.test.ts b/packages/core/test/session-execution.test.ts index 9dabbde747b..6735291e83f 100644 --- a/packages/core/test/session-execution.test.ts +++ b/packages/core/test/session-execution.test.ts @@ -21,7 +21,7 @@ import { SessionMessage } from "@opencode-ai/core/session/message" import { SessionRunner } from "@opencode-ai/core/session/runner/index" import { SessionInboxTable, SessionTable } from "@opencode-ai/core/session/sql" import { SessionStore } from "@opencode-ai/core/session/store" -import { Context, Deferred, Effect, Exit, Fiber, Layer, LayerMap, Scope } from "effect" +import { Cause, Context, Deferred, Effect, Exit, Fiber, Layer, LayerMap, Scope } from "effect" import { eq } from "drizzle-orm" import { testEffect } from "./lib/effect" @@ -888,6 +888,49 @@ describe("SessionRestart background recovery", () => { }), ) + it.effect("retains a subagent completion marker when synthetic admission conflicts", () => + Effect.gen(function* () { + const database = yield* Database.Service + const bus = yield* Bus.Service + const jobs = yield* Job.Service + const sessions = yield* Session.Service + const parent = Session.ID.make("ses_completion_conflict_parent") + const child = Session.ID.make("ses_completion_conflict_child") + yield* seedSessions(database, [parent]) + yield* seedSessions(database, [child], { parent_id: parent }) + yield* jobs.start({ + id: child, + type: "subagent", + recovery: { + kind: "subagent", + parentSessionID: parent, + childSessionID: child, + agent: "explore", + description: "Completed inspection", + }, + run: Effect.succeed("Recovered result"), + }) + yield* jobs.wait({ id: child }) + yield* jobs.background(child) + const marker = (yield* jobs.pendingBackground)[0] + if (!marker) return yield* Effect.die("background record missing") + yield* SessionInbox.admit(database.db, bus, { + id: marker.notificationID, + sessionID: parent, + item: { type: "user", payload: { text: "User input" }, delivery: "steer" }, + }) + + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const context = yield* buildExecution(scope, () => Effect.die("Admission must not wake the parent")) + const exit = yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions.pipe(Effect.exit) + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Session.SyntheticConflictError) + expect(yield* jobs.pendingBackground).toEqual([marker]) + expect(yield* sessions.inbox(parent)).toMatchObject([{ type: "user", payload: { text: "User input" } }]) + }), + ) + for (const resumeAttempts of [1, 2]) { it.effect(`honors a suspended parent's restart budget after ${resumeAttempts} attempts before notifying it`, () => Effect.gen(function* () { diff --git a/packages/core/test/tool-subagent.test.ts b/packages/core/test/tool-subagent.test.ts index 89810c10228..9aed4a76565 100644 --- a/packages/core/test/tool-subagent.test.ts +++ b/packages/core/test/tool-subagent.test.ts @@ -1,8 +1,12 @@ import { describe, expect } from "bun:test" -import { Effect, Fiber, Layer, Schema, Stream } from "effect" +import { Deferred, Effect, Fiber, Layer, Schema, Stream } from "effect" +import { LanguageModel } from "@opencode-ai/ai" +import { OpenAIChat } from "@opencode-ai/ai/protocols" +import { TestLLM } from "@opencode-ai/ai/testing" import path from "path" import { Money } from "@opencode-ai/schema/money" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" +import { LayerNodePlatform } from "@opencode-ai/core/effect/app-node-platform" import { LayerNode } from "@opencode-ai/util/effect/layer-node" import { Global } from "@opencode-ai/util/global" import { makeGlobalNode, makeLocationNode } from "@opencode-ai/util/effect/app-node" @@ -15,10 +19,12 @@ import { Provider } from "@opencode-ai/core/provider" import { AbsolutePath } from "@opencode-ai/core/schema" import { Agent } from "@opencode-ai/core/agent" import { Job } from "@opencode-ai/core/job" +import { KV } from "@opencode-ai/core/kv" import { LocationServiceMap } from "@opencode-ai/core/location-service-map" import { Session } from "@opencode-ai/core/session" import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionExecution } from "@opencode-ai/core/session/execution" +import { SessionRestart } from "@opencode-ai/core/session/execution/restart" import { SessionInbox } from "@opencode-ai/core/session/inbox" import { SessionMessage } from "@opencode-ai/core/session/message" import { SessionRunnerModel } from "@opencode-ai/core/session/runner/model" @@ -120,6 +126,29 @@ const replacements = [ ] satisfies LayerNode.Replacements const productionIt = testEffect(AppNodeBuilder.build(nodes, replacements)) const it = testEffect(AppNodeBuilder.build(nodes, [...replacements, [PluginSupervisor.node, subagentPluginSupervisor]])) +const completionIt = testEffect( + AppNodeBuilder.build(LayerNode.group([nodes, SessionRestart.node, KV.node]), [ + [Global.node, tempGlobalLayer], + [PluginSupervisor.node, subagentPluginSupervisor], + [LayerNodePlatform.llmClient, TestLLM.testLayer({ fallback: TestLLM.text(childText, "completion") })], + [ + SessionRunnerModel.node, + Layer.succeed(SessionRunnerModel.Service, { + resolve: () => + Effect.succeed( + SessionRunnerModel.resolved( + LanguageModel.make({ id: "child", provider: "test", route: OpenAIChat.route }), + { + capabilities: { tools: true, input: ["text"], output: ["text"] }, + cost: [], + limit: { context: 200_000, output: 32_000 }, + }, + ), + ), + }), + ], + ]), +) const withSubagent = (location: Location.Ref) => Effect.gen(function* () { @@ -147,6 +176,81 @@ const withSubagent = (location: Location.Ref) => }) describe("SubagentTool", () => { + completionIt.live("admits one durable completion across live delivery and restart replay", () => + Effect.acquireRelease( + Effect.promise(() => tmpdir()), + (dir) => Effect.promise(() => dir[Symbol.asyncDispose]()), + ).pipe( + Effect.flatMap((dir) => + Effect.gen(function* () { + const sessions = yield* Session.Service + const parent = yield* sessions.create({ + location: Location.Ref.make({ directory: AbsolutePath.make(dir.path) }), + model: parentModel, + title: "Completion recipient", + }) + yield* withSubagent(parent.location) + const locations = yield* LocationServiceMap.Service + const registry = yield* Tool.Service.pipe(Effect.provide(locations.get(parent.location))) + const jobs = yield* Job.Service + const bus = yield* Bus.Service + const admitted = yield* Deferred.make() + const notifications: SessionMessage.ID[] = [] + yield* bus.project(SessionEvent.InboxEnqueued, (event) => + Effect.gen(function* () { + if (event.data.sessionID !== parent.id || event.data.item.type !== "synthetic") return + notifications.push(event.data.inboxID) + const marker = (yield* jobs.pendingBackground).find((job) => job.notificationID === event.data.inboxID) + // The marker must survive until admission commits, not merely until delivery starts. + expect(marker?.status).toBe("completed") + if (marker) yield* Deferred.succeed(admitted, marker) + }), + ) + + const result = yield* executeTool(registry, { + sessionID: parent.id, + ...toolIdentity, + call: { + type: "tool-call", + id: "call-completion-replay", + name: SubagentTool.name, + input: { agent: "reviewer", description: "background review", prompt: "review", background: true }, + }, + }) + const marker = yield* Deferred.await(admitted) + yield* jobs.pendingBackground.pipe(Effect.repeat({ until: (pending) => pending.length === 0 })) + yield* sessions.wait(parent.id) + const messages = (yield* sessions.context(parent.id)).filter((message) => message.type === "synthetic") + expect(messages).toEqual([ + expect.objectContaining({ + id: marker.notificationID, + description: "background review", + text: `\n${childText}\n`, + metadata: { + source: "subagent", + childID: outputSessionID(result.metadata), + agent: "reviewer", + state: "completed", + }, + }), + ]) + + // Reproduce a crash after admission but before acknowledgment using the real persisted marker. + const kv = yield* KV.Service + yield* kv.set(`job.background/${marker.notificationID}`, marker) + const restart = yield* SessionRestart.Service + yield* restart.resumeSuspendedSessions + yield* sessions.wait(parent.id) + expect(notifications).toEqual([marker.notificationID]) + expect((yield* sessions.context(parent.id)).filter((message) => message.type === "synthetic")).toEqual( + messages, + ) + expect(yield* jobs.pendingBackground).toEqual([]) + }), + ), + ), + ) + productionIt.live("registers globally while resolving agents from the caller location", () => Effect.acquireRelease( Effect.promise(() => tmpdir()),