diff --git a/packages/core/src/job.ts b/packages/core/src/job.ts index 75100f333f4..1137ec0c026 100644 --- a/packages/core/src/job.ts +++ b/packages/core/src/job.ts @@ -1,11 +1,41 @@ export * as Job from "./job.js" -import { Cause, Clock, Context, Deferred, Effect, Exit, Layer, Scope, SynchronizedRef } from "effect" +import { Array, Cause, Clock, Context, Deferred, Effect, Exit, Layer, Schema, Scope, SynchronizedRef } from "effect" import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" import { Identifier } from "./id/id.js" +import { KV } from "./kv.js" +import { SessionMessage } from "./session/message.js" import { SessionSchema } from "./session/schema.js" -export type Status = "running" | "completed" | "error" | "cancelled" +const Background = Schema.Struct({ + id: Schema.String, + notificationID: SessionMessage.ID, + recovery: Schema.Union([ + Schema.Struct({ + kind: Schema.Literal("shell"), + sessionID: SessionSchema.ID, + shellID: Schema.String, + command: Schema.String, + }), + Schema.Struct({ + kind: Schema.Literal("subagent"), + parentSessionID: SessionSchema.ID, + childSessionID: SessionSchema.ID, + agent: Schema.String, + description: Schema.String, + }), + ]), + status: Schema.Literals(["running", "completed", "error", "cancelled"]), + output: Schema.optionalKey(Schema.String), + error: Schema.optionalKey(Schema.String), +}) + +export type Background = typeof Background.Type +export type Recovery = Background["recovery"] +export type Status = Background["status"] + +const decodeBackground = Schema.decodeUnknownResult(Background) +const backgroundPrefix = "job.background/" export type Info = { id: string @@ -17,6 +47,7 @@ export type Info = { output?: string error?: string metadata?: Record + notificationID?: SessionMessage.ID } type Active = { @@ -27,6 +58,7 @@ type Active = { token: object blockingSessions: Map isBackgrounded: boolean + recovery?: Recovery } type State = { @@ -63,6 +95,8 @@ export type StartInput = { type: string title?: string metadata?: Record + recovery?: Recovery + notificationID?: SessionMessage.ID run: Effect.Effect } @@ -96,6 +130,8 @@ export interface Interface { readonly background: (id: string) => Effect.Effect readonly backgroundAll: (input: BackgroundAllInput) => Effect.Effect readonly cancel: (id: string) => Effect.Effect + readonly pendingBackground: Effect.Effect + readonly completeBackground: (notificationID: SessionMessage.ID) => Effect.Effect } export class Service extends Context.Service()("@opencode/Job") {} @@ -126,43 +162,57 @@ function decrementSession(input: Map, sessionID: Sessi } /** - * Makes one scoped, process-local registry. Entries are intentionally not - * durable: process restart or owner-scope closure loses status and interrupts - * live work. Persisted observation, restart recovery, and remote workers need a - * separate durable ownership slice rather than pretending this registry has - * those semantics. + * Makes one scoped, process-local registry. Explicitly recoverable background + * work also owns a durable notification marker until its notification is admitted. */ export const make = Effect.gen(function* () { + const kv = yield* KV.Service const state: State = { jobs: yield* SynchronizedRef.make(new Map()), scope: yield* Scope.Scope, } + const persistBackground = Effect.fnUntraced(function* (job: Active) { + if (!job.recovery || !job.info.notificationID) return + yield* kv.set(`${backgroundPrefix}${job.info.notificationID}`, { + id: job.info.id, + notificationID: job.info.notificationID, + recovery: job.recovery, + status: job.info.status, + ...(job.info.output !== undefined ? { output: job.info.output } : {}), + ...(job.info.error !== undefined ? { error: job.info.error } : {}), + }) + }) + const settle = Effect.fnUntraced(function* (id: string, token: object, exit: Exit.Exit) { const completed_at = yield* Clock.currentTimeMillis - const result = yield* SynchronizedRef.modify(state.jobs, (jobs): readonly [FinishResult, Map] => { - const job = jobs.get(id) - if (!job) return [{}, jobs] - if (job.token !== token) return [{}, jobs] - if (job.info.status !== "running") return [{ info: snapshot(job) }, jobs] - const status: Exclude = Exit.isSuccess(exit) - ? "completed" - : Cause.hasInterruptsOnly(exit.cause) - ? "cancelled" - : "error" - const next = { - ...job, - blockingSessions: new Map(), - info: { - ...job.info, - status, - completed_at, - ...(Exit.isSuccess(exit) ? { output: exit.value } : {}), - ...(Exit.isFailure(exit) ? { error: errorText(Cause.squash(exit.cause)) } : {}), - }, - } - return [{ info: snapshot(next), done: job.done, scope: job.scope }, new Map(jobs).set(id, next)] - }) + const result = yield* SynchronizedRef.modifyEffect( + state.jobs, + Effect.fnUntraced(function* (jobs): Effect.fn.Return]> { + const job = jobs.get(id) + if (!job) return [{}, jobs] + if (job.token !== token) return [{}, jobs] + if (job.info.status !== "running") return [{ info: snapshot(job) }, jobs] + const status: Exclude = Exit.isSuccess(exit) + ? "completed" + : Cause.hasInterruptsOnly(exit.cause) + ? "cancelled" + : "error" + const next = { + ...job, + blockingSessions: new Map(), + info: { + ...job.info, + status, + completed_at, + ...(Exit.isSuccess(exit) ? { output: exit.value } : {}), + ...(Exit.isFailure(exit) ? { error: errorText(Cause.squash(exit.cause)) } : {}), + }, + } + if (status !== "cancelled") yield* persistBackground(next) + return [{ info: snapshot(next), done: job.done, scope: job.scope }, new Map(jobs).set(id, next)] + }), + ) if (result.info && result.done) yield* Deferred.succeed(result.done, result.info).pipe(Effect.ignore) if (result.scope) { yield* Scope.close(result.scope, Exit.void).pipe(Effect.forkIn(state.scope, { startImmediately: true })) @@ -170,22 +220,6 @@ export const make = Effect.gen(function* () { return result.info }) - const fork = Effect.fnUntraced(function* ( - scope: Scope.Scope, - id: string, - token: object, - run: Effect.Effect, - ) { - return yield* run.pipe( - Effect.matchCauseEffect({ - onSuccess: (output) => settle(id, token, Exit.succeed(output)), - onFailure: (cause) => settle(id, token, Exit.failCause(cause)), - }), - Effect.asVoid, - Effect.forkIn(scope, { startImmediately: true }), - ) - }) - const get: Interface["get"] = Effect.fn("Job.get")(function* (id) { const job = (yield* SynchronizedRef.get(state.jobs)).get(id) if (!job) return undefined @@ -201,10 +235,10 @@ export const make = Effect.gen(function* () { const backgrounded = yield* Deferred.make() const result = yield* SynchronizedRef.modifyEffect( state.jobs, - Effect.fnUntraced(function* (jobs) { + Effect.fnUntraced(function* (jobs): Effect.fn.Return]> { const existing = jobs.get(id) if (existing?.info.status === "running") { - return [{ info: snapshot(existing) }, jobs] as readonly [StartResult, Map] + return [{ info: snapshot(existing) }, jobs] } const scope = yield* Scope.fork(state.scope, "parallel") const token = {} @@ -216,6 +250,7 @@ export const make = Effect.gen(function* () { status: "running" as const, started_at, metadata: input.metadata, + ...(input.notificationID ? { notificationID: input.notificationID } : {}), }, done, backgrounded, @@ -223,14 +258,18 @@ export const make = Effect.gen(function* () { token, blockingSessions: new Map(), isBackgrounded: false, + recovery: input.recovery, } - return [{ info: snapshot(job), scope, token }, new Map(jobs).set(id, job)] as readonly [ - StartResult, - Map, - ] + return [{ info: snapshot(job), scope, token }, new Map(jobs).set(id, job)] }), ) - if ("scope" in result) yield* fork(result.scope, id, result.token, restore(input.run)) + if ("scope" in result) + yield* restore(input.run).pipe( + Effect.exit, + Effect.flatMap((exit) => settle(id, result.token, exit)), + Effect.asVoid, + Effect.forkIn(result.scope, { startImmediately: true }), + ) return result.info }), ) @@ -281,20 +320,31 @@ export const make = Effect.gen(function* () { ).pipe(Effect.ensuring(removeBlock(input))) }) - const background: Interface["background"] = Effect.fn("Job.background")(function* (id) { - const result = yield* SynchronizedRef.modify( - state.jobs, - (jobs): readonly [BackgroundResult, Map] => { - const job = jobs.get(id) - if (!job || job.info.status !== "running") return [{}, jobs] - if (job.isBackgrounded) return [{ info: snapshot(job) }, jobs] - const next = { - ...job, - isBackgrounded: true, - blockingSessions: new Map(), - } - return [{ info: snapshot(next), backgrounded: job.backgrounded }, new Map(jobs).set(id, next)] + const markBackground = Effect.fnUntraced(function* (job: Active) { + const next = { + ...job, + isBackgrounded: true, + blockingSessions: new Map(), + info: { + ...job.info, + ...(job.recovery ? { notificationID: job.info.notificationID ?? SessionMessage.ID.create() } : {}), }, + } + yield* persistBackground(next) + return next + }) + + const background: Interface["background"] = Effect.fn("Job.background")(function* (id) { + const result = yield* SynchronizedRef.modifyEffect( + state.jobs, + Effect.fnUntraced(function* (jobs): Effect.fn.Return]> { + const job = jobs.get(id) + // Recoverable work may finish before the caller backgrounds it. + if (!job || (job.info.status !== "running" && !job.recovery)) return [{}, jobs] + if (job.isBackgrounded) return [{ info: snapshot(job) }, jobs] + const next = yield* markBackground(job) + return [{ info: snapshot(next), backgrounded: job.backgrounded }, new Map(jobs).set(id, next)] + }), ) if (result.info && result.backgrounded) yield* Deferred.succeed(result.backgrounded, result.info).pipe(Effect.ignore) @@ -302,60 +352,83 @@ export const make = Effect.gen(function* () { }) const backgroundAll: Interface["backgroundAll"] = Effect.fn("Job.backgroundAll")(function* (input) { - const result = yield* SynchronizedRef.modify( + const result = yield* SynchronizedRef.modifyEffect( state.jobs, - (jobs): readonly [BackgroundResult[], Map] => { - const results: BackgroundResult[] = [] + Effect.fnUntraced(function* (jobs): Effect.fn.Return< + readonly [Required[], Map] + > { + const results: Required[] = [] const next = new Map(jobs) for (const [id, job] of jobs) { if (job.info.status !== "running") continue if (job.isBackgrounded) continue if (input.type !== undefined && job.info.type !== input.type) continue if (!job.blockingSessions.has(input.sessionID)) continue - const updated = { - ...job, - isBackgrounded: true, - blockingSessions: new Map(), - } + const updated = yield* markBackground(job) results.push({ info: snapshot(updated), backgrounded: job.backgrounded }) next.set(id, updated) } return [results, next] - }, + }), ) - yield* Effect.forEach( - result, - (item) => (item.info && item.backgrounded ? Deferred.succeed(item.backgrounded, item.info) : Effect.void), - { discard: true }, - ) - return result.flatMap((item) => (item.info ? [item.info] : [])) + yield* Effect.forEach(result, (item) => Deferred.succeed(item.backgrounded, item.info), { discard: true }) + return result.map((item) => item.info) }) const cancel: Interface["cancel"] = Effect.fn("Job.cancel")(function* (id) { const completed_at = yield* Clock.currentTimeMillis - const result = yield* SynchronizedRef.modify(state.jobs, (jobs): readonly [FinishResult, Map] => { - const job = jobs.get(id) - if (!job) return [{}, jobs] - if (job.info.status !== "running") return [{ info: snapshot(job) }, jobs] - const next = { - ...job, - blockingSessions: new Map(), - info: { - ...job.info, - status: "cancelled" as const, - completed_at, - }, - } - return [{ info: snapshot(next), done: job.done, scope: job.scope }, new Map(jobs).set(id, next)] - }) + const result = yield* SynchronizedRef.modifyEffect( + state.jobs, + Effect.fnUntraced(function* (jobs): Effect.fn.Return]> { + const job = jobs.get(id) + if (!job) return [{}, jobs] + if (job.info.status !== "running") return [{ info: snapshot(job) }, jobs] + const next = { + ...job, + blockingSessions: new Map(), + info: { + ...job.info, + status: "cancelled" as const, + completed_at, + }, + } + yield* persistBackground(next) + return [{ info: snapshot(next), done: job.done, scope: job.scope }, new Map(jobs).set(id, next)] + }), + ) if (result.info && result.done) yield* Deferred.succeed(result.done, result.info).pipe(Effect.ignore) if (result.scope) yield* Scope.close(result.scope, Exit.void) return result.info }) - return Service.of({ get, start, wait, block, background, backgroundAll, cancel }) + const pendingBackground: Interface["pendingBackground"] = Effect.gen(function* () { + const recovered: Background[] = [] + let after: string | undefined + do { + const page = yield* kv.scan({ prefix: backgroundPrefix, after }) + recovered.push(...Array.filterMap(page.entries, (entry) => decodeBackground(entry.value))) + after = page.next + } while (after) + return recovered + }).pipe(Effect.withSpan("Job.pendingBackground")) + + const completeBackground: Interface["completeBackground"] = Effect.fn("Job.completeBackground")((notificationID) => + kv.remove(`${backgroundPrefix}${notificationID}`), + ) + + return Service.of({ + get, + start, + wait, + block, + background, + backgroundAll, + cancel, + pendingBackground, + completeBackground, + }) }) const layer = Layer.effect(Service, make) -export const node = makeGlobalNode({ service: Service, layer, deps: [] }) +export const node = makeGlobalNode({ service: Service, layer, deps: [KV.node] }) diff --git a/packages/core/src/plugin/runtime.ts b/packages/core/src/plugin/runtime.ts index 36751a74f7e..3dfb9c0af12 100644 --- a/packages/core/src/plugin/runtime.ts +++ b/packages/core/src/plugin/runtime.ts @@ -29,7 +29,7 @@ export interface Interface { | "wait" | "context" > - readonly job: Pick + readonly job: Pick readonly location: { readonly agent: { readonly list: ( @@ -92,6 +92,8 @@ export const layerWithCell = (cell: Cell) => block: (input) => require(cell, (runtime) => runtime.job.block(input)), background: (id) => require(cell, (runtime) => runtime.job.background(id)), cancel: (id) => require(cell, (runtime) => runtime.job.cancel(id)), + completeBackground: (notificationID) => + require(cell, (runtime) => runtime.job.completeBackground(notificationID)), }, location: { agent: { diff --git a/packages/core/src/session/execution.ts b/packages/core/src/session/execution.ts index dd256b5cf1d..c49526973e0 100644 --- a/packages/core/src/session/execution.ts +++ b/packages/core/src/session/execution.ts @@ -3,6 +3,7 @@ export * as SessionExecution from "./execution.js" import { Cause, Context, Effect, Exit, Layer } from "effect" import { Bus } from "../bus.js" import { Database } from "../database/database.js" +import { Job } from "../job.js" import { LocationServiceMap } from "../location-service-map.js" import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" import { SessionEvent } from "./event.js" @@ -52,6 +53,7 @@ export const layer = Layer.effect( const store = yield* SessionStore.Service const locations = yield* LocationServiceMap.Service const bus = yield* Bus.Service + const jobs = yield* Job.Service const db = (yield* Database.Service).db const reportLifecycle = (sessionID: SessionSchema.ID, effect: Effect.Effect) => effect.pipe( @@ -118,6 +120,7 @@ export const layer = Layer.effect( if (outcome.type === "interrupted") { // A user cancel releases the claim: the turn must not resurrect at the next // boot. Shutdown interruption keeps it for restart continuity. + if (outcome.reason === "user") yield* jobs.cancel(sessionID) yield* bus.publish( SessionEvent.Execution.Interrupted, { sessionID, reason: outcome.reason }, @@ -167,7 +170,7 @@ export const layer = Layer.effect( export const node = makeGlobalNode({ service: Service, layer, - deps: [SessionStore.node, LocationServiceMap.node, Bus.node, Database.node], + deps: [SessionStore.node, LocationServiceMap.node, Bus.node, Database.node, Job.node], }) /** Low-level compatibility layer for callers that only need durable Session recording. */ diff --git a/packages/core/src/session/execution/restart.ts b/packages/core/src/session/execution/restart.ts index 7413d690385..303305bac37 100644 --- a/packages/core/src/session/execution/restart.ts +++ b/packages/core/src/session/execution/restart.ts @@ -3,6 +3,8 @@ export * as SessionRestart from "./restart.js" import { Context, Effect, Layer } from "effect" import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" import { Bus } from "../../bus.js" +import { Job } from "../../job.js" +import { Session } from "../../session.js" import { SessionEvent } from "../event.js" import { SessionExecution } from "../execution.js" import { SessionSchema } from "../schema.js" @@ -45,6 +47,9 @@ export interface Interface { * process: crash, SIGKILL, isolate eviction, and graceful restart all leave * the same durable signature. * + * Recovery is at-least-once: local coordination prevents concurrent drains, + * not repeated external side effects after a crash. + * * The sweep assumes every orphaned claim's owner is dead. The managed-server * protocol guarantees this: a successor is only spawned after the previous * process is confirmed dead (client service `kill`/`evict` poll the PID), the @@ -62,14 +67,16 @@ export const layer = (options?: Options) => const store = yield* SessionStore.Service const execution = yield* SessionExecution.Service const bus = yield* Bus.Service + const jobs = yield* Job.Service + const sessions = yield* Session.Service const scope = yield* Effect.scope const maxAttempts = options?.maxAttempts ?? DEFAULT_MAX_ATTEMPTS - const resumeOne = Effect.fnUntraced(function* (sessionID: SessionSchema.ID) { + const prepareResume = Effect.fnUntraced(function* (sessionID: SessionSchema.ID) { // Durable before the resume runs, so a crash inside the resumed turn is // counted by the next sweep and the budget cannot be dodged. const attempts = yield* store.countResume(sessionID) - if (attempts === undefined) return // the Session was deleted since listing + if (attempts === undefined) return false if (attempts > maxAttempts) { // Terminalize instead: the release hook clears the claim and resets the // counter atomically with the terminal event. @@ -78,31 +85,166 @@ export const layer = (options?: Options) => { sessionID, error: RESUME_EXHAUSTED }, { commit: () => store.release(sessionID) }, ) - return + return false } yield* bus.publish(SessionEvent.Synthetic, { sessionID, text: CONTINUE_AFTER_SERVER_RESTART, description: "Continuing after restart", }) - // Forked into the service scope so boot never waits on resumed turns; - // resuming an already-live Session joins its execution. Drain failures - // are logged and durably recorded by the execution layer. - yield* execution.resume(sessionID).pipe(Effect.ignore, Effect.forkIn(scope)) + return true + }) + + const recoverShell = Effect.fnUntraced(function* ( + background: Job.Background, + recovery: Extract, + ) { + const state = background.status === "running" ? "cancelled" : background.status + const text = + background.status === "running" + ? "Command cancelled because the server restarted" + : state === "completed" + ? (background.output ?? "Command completed") + : state === "error" + ? (background.error ?? "Command failed") + : "Command cancelled" + + yield* sessions + .synthetic({ + id: background.notificationID, + sessionID: recovery.sessionID, + description: recovery.command, + text: `\n${text}\n`, + metadata: { + source: "shell", + jobID: background.id, + shellID: recovery.shellID, + state, + }, + resume: false, + }) + .pipe( + Effect.catchTag("Session.NotFoundError", () => Effect.void), + Effect.orDie, + ) + yield* jobs.completeBackground(background.notificationID) + }) + + const recoverSubagent = Effect.fnUntraced(function* ( + background: Job.Background, + recovery: Extract, + suspended: ReadonlySet, + ) { + const child = yield* store.get(recovery.childSessionID) + if (!child || child.parentID !== recovery.parentSessionID || !(yield* store.get(recovery.parentSessionID))) { + yield* jobs.completeBackground(background.notificationID) + return + } + + 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) + }) + + if (background.status !== "running") { + yield* notify(background) + return + } + if ((yield* execution.active).has(recovery.childSessionID)) return + if (!(yield* prepareResume(recovery.childSessionID))) { + yield* notify({ status: "error", error: RESUME_EXHAUSTED.message }) + return + } + + yield* jobs.start({ + id: background.id, + type: "subagent", + title: recovery.description, + notificationID: background.notificationID, + recovery, + run: execution.resume(recovery.childSessionID).pipe( + Effect.andThen(store.context(recovery.childSessionID)), + Effect.map((messages) => { + const assistant = messages.findLast( + (message) => + message.type === "assistant" && message.time.completed !== undefined && message.error === undefined, + ) + if (assistant?.type !== "assistant") return "Subagent completed without a text response." + return ( + assistant.content + .filter((part) => part.type === "text") + .map((part) => part.text) + .join("") || "Subagent completed without a text response." + ) + }), + ), + }) + yield* jobs.background(background.id) + yield* jobs.wait({ id: background.id }).pipe( + Effect.flatMap((result) => (result.info ? notify(result.info) : Effect.void)), + Effect.ignore, + Effect.forkIn(scope), + ) }) return Service.of({ resumeSuspendedSessions: Effect.gen(function* () { - // Child claims never drive recovery (children are not resumed), so a - // dead child's claim is noise no terminal will ever release. Clearing - // is safe even against a live child: claims are recovery markers, not - // locks, and children are excluded from that recovery. - yield* store.releaseChildClaims const active = yield* execution.active - // Sessions already draining in this process keep their claim; resuming - // them would only inject a stray continuation into a live turn. - const orphaned = (yield* store.listSuspended()).filter((sessionID) => !active.has(sessionID)) - yield* Effect.forEach(orphaned, resumeOne, { concurrency: "unbounded", discard: true }) + // Early notices wait for root recovery's accounting, including roots that exhaust their budget. + const suspended = new Set((yield* store.listSuspended()).filter((sessionID) => !active.has(sessionID))) + const pending = yield* jobs.pendingBackground + yield* store.releaseChildClaims( + pending.flatMap((background) => + background.status === "running" && background.recovery.kind === "subagent" + ? [background.recovery.childSessionID] + : [], + ), + ) + yield* Effect.forEach( + pending, + Effect.fnUntraced(function* (background) { + if ((yield* jobs.get(background.id))?.status === "running") return + const recovery = background.recovery + yield* recovery.kind === "shell" + ? recoverShell(background, recovery) + : recoverSubagent(background, recovery, suspended) + }), + { discard: true }, + ) + + // Background completion can wake a parent, so inspect local ownership only after recovery. + const resumed = yield* execution.active + yield* Effect.forEach( + (yield* store.listSuspended()).filter((sessionID) => !resumed.has(sessionID)), + (sessionID) => + execution + .resume(sessionID) + .pipe(Effect.ignore, Effect.forkIn(scope), Effect.when(prepareResume(sessionID))), + { concurrency: "unbounded", discard: true }, + ) + // Async observers consult this set at delivery; later completions wake parents normally. + suspended.clear() }), }) }), @@ -111,5 +253,5 @@ export const layer = (options?: Options) => export const node = makeGlobalNode({ service: Service, layer: layer(), - deps: [SessionStore.node, SessionExecution.node, Bus.node], + deps: [SessionStore.node, SessionExecution.node, Bus.node, Job.node, Session.node], }) diff --git a/packages/core/src/session/runner/llm.ts b/packages/core/src/session/runner/llm.ts index 205f021d537..9656228cf8d 100644 --- a/packages/core/src/session/runner/llm.ts +++ b/packages/core/src/session/runner/llm.ts @@ -301,11 +301,18 @@ const layer = Layer.effect( if (message.type !== "assistant") continue for (const tool of message.content) { if (tool.type !== "tool" || (tool.state.status !== "streaming" && tool.state.status !== "running")) continue + const metadata = tool.state.status === "running" ? tool.state.metadata : undefined + const childID = + tool.name === "subagent" && typeof metadata?.sessionID === "string" ? metadata.sessionID : undefined yield* bus.publish(SessionEvent.Tool.Failed, { sessionID, assistantMessageID: message.id, id: tool.id, - error: { type: "aborted", message: `Tool execution interrupted: ${tool.name}` }, + error: { + type: "aborted", + message: `Tool execution interrupted: ${tool.name}${childID ? ` (sessionID: ${childID})` : ""}`, + }, + ...(metadata && Object.keys(metadata).length > 0 ? { metadata } : {}), executed: tool.executed === true, }) } diff --git a/packages/core/src/session/runner/publish-llm-event.ts b/packages/core/src/session/runner/publish-llm-event.ts index d50a9578fa9..a03da9734b8 100644 --- a/packages/core/src/session/runner/publish-llm-event.ts +++ b/packages/core/src/session/runner/publish-llm-event.ts @@ -327,7 +327,10 @@ export const createLLMEventPublisher = (bus: Pick, inp sessionID: input.sessionID, assistantMessageID, id, - error, + error: + tool.name === "subagent" && error.type === "aborted" && typeof tool.progress?.sessionID === "string" + ? { ...error, message: `${error.message} (sessionID: ${tool.progress.sessionID})` } + : error, ...failureSnapshot(tool, metadata), executed: tool.providerExecuted, }) diff --git a/packages/core/src/session/store.ts b/packages/core/src/session/store.ts index 908170c1ab8..f657c3a26c6 100644 --- a/packages/core/src/session/store.ts +++ b/packages/core/src/session/store.ts @@ -1,6 +1,6 @@ export * as SessionStore from "./store.js" -import { and, eq, isNotNull, isNull, sql } from "drizzle-orm" +import { and, eq, isNotNull, isNull, notInArray, sql } from "drizzle-orm" import { Context, Effect, Layer } from "effect" import { Database } from "../database/database.js" import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" @@ -18,9 +18,8 @@ export interface Interface { messageID: SessionMessage.ID, ) => Effect.Effect<{ readonly sessionID: Session.ID; readonly message: SessionMessage.Info } | undefined> /** - * Top-level Sessions holding an execution claim. Child (subagent) Sessions - * are excluded: a resumed parent re-runs its tool call and spawns fresh - * children, so resuming orphaned children would duplicate their work. + * Top-level Sessions holding an execution claim. Recoverable background + * children are resumed separately through their durable Job records. */ readonly listSuspended: () => Effect.Effect> /** @@ -33,11 +32,10 @@ export interface Interface { /** Releases the claim and resets resume accounting. Terminal events call this on commit. */ readonly release: (sessionID: Session.ID) => Effect.Effect /** - * Clears orphaned child (subagent) claims. Children are never resumed - * independently, so a dead child's claim is noise no terminal will ever - * release. + * Clears orphaned child claims except children owned by recoverable + * background subagent jobs. */ - readonly releaseChildClaims: Effect.Effect + readonly releaseChildClaims: (recoverable: ReadonlyArray) => Effect.Effect /** * Durably counts one more resume of an orphaned claim, returning the new * total — or undefined when the Session no longer exists. @@ -103,12 +101,20 @@ const layer = Layer.effect( .run() .pipe(Effect.orDie) }), - releaseChildClaims: db - .update(SessionTable) - .set({ time_suspended: null, resume_attempts: 0, time_updated: sql`${SessionTable.time_updated}` }) - .where(and(isNotNull(SessionTable.time_suspended), isNotNull(SessionTable.parent_id))) - .run() - .pipe(Effect.orDie, Effect.asVoid, Effect.withSpan("SessionStore.releaseChildClaims")), + releaseChildClaims: Effect.fn("SessionStore.releaseChildClaims")((recoverable) => + db + .update(SessionTable) + .set({ time_suspended: null, resume_attempts: 0, time_updated: sql`${SessionTable.time_updated}` }) + .where( + and( + isNotNull(SessionTable.time_suspended), + isNotNull(SessionTable.parent_id), + recoverable.length > 0 ? notInArray(SessionTable.id, Array.from(recoverable)) : undefined, + ), + ) + .run() + .pipe(Effect.orDie, Effect.asVoid), + ), countResume: Effect.fn("SessionStore.countResume")(function* (sessionID) { const row = yield* db .update(SessionTable) diff --git a/packages/core/src/shell.ts b/packages/core/src/shell.ts index aad62a84823..e1d8467dd49 100644 --- a/packages/core/src/shell.ts +++ b/packages/core/src/shell.ts @@ -134,8 +134,8 @@ const layer = () => Effect.gen(function* () { for (const session of sessions.values()) { if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) - // Unblock waiters still pending at teardown; succeed is a no-op once already resolved. - yield* Deferred.fail(session.done, new NotFoundError({ id: Shell.ID.make(session.info.id) })) + // Teardown interrupts pending commands; it is not a terminal command failure. + yield* Deferred.interrupt(session.done) } sessions.clear() exitOrder.length = 0 diff --git a/packages/core/src/tool/plugin/shell.ts b/packages/core/src/tool/plugin/shell.ts index 75d25d46767..beb2111a7ec 100644 --- a/packages/core/src/tool/plugin/shell.ts +++ b/packages/core/src/tool/plugin/shell.ts @@ -115,56 +115,45 @@ export const Plugin = { const permission = yield* Permission.Service const config = yield* Config.Service - const notifyWhenDone = Effect.fn("ShellTool.notifyWhenDone")(function* ( - sessionID: SessionSchema.ID, - id: string, - shellID: string, - command: string, - settled: Deferred.Deferred, - ) { - yield* runtime.job.wait({ id: id }).pipe( - Effect.flatMap((result) => - Effect.gen(function* () { - const info = result.info - if (!info) return - const state = - info.status === "completed" - ? "completed" - : info.status === "error" - ? "error" - : info.status === "cancelled" - ? "cancelled" - : undefined - if (state === undefined) return - const output = state === "completed" ? yield* Deferred.await(settled) : undefined - const text = output - ? resultMessages(output).join("\n\n") - : state === "error" - ? (info.error ?? "Command failed") - : "Command cancelled" - yield* runtime.session.synthetic({ - sessionID, - text: `\n${text}\n`, - description: command, - metadata: { - source: "shell", - jobID: id, - shellID, - state, - ...(output - ? { - truncated: output.truncated, - ...(output.exit !== undefined ? { exit: output.exit } : {}), - ...(output.timeout !== undefined ? { timeout: output.timeout } : {}), - } - : {}), - }, - }) - }), - ), - Effect.forkIn(scope, { startImmediately: true }), - ) - }) + const notifyWhenDone = Effect.fn("ShellTool.notifyWhenDone")( + function* ( + sessionID: SessionSchema.ID, + id: string, + shellID: string, + command: string, + settled: Deferred.Deferred, + ) { + const info = (yield* runtime.job.wait({ id })).info + if (!info || info.status === "running") return + const output = info.status === "completed" ? yield* Deferred.await(settled) : undefined + const text = output + ? resultMessages(output).join("\n\n") + : info.status === "error" + ? (info.error ?? "Command failed") + : "Command cancelled" + yield* runtime.session.synthetic({ + ...(info.notificationID ? { id: info.notificationID } : {}), + sessionID, + text: `\n${text}\n`, + description: command, + metadata: { + source: "shell", + jobID: id, + shellID, + state: info.status, + ...(output + ? { + truncated: output.truncated, + ...(output.exit !== undefined ? { exit: output.exit } : {}), + ...(output.timeout !== undefined ? { timeout: output.timeout } : {}), + } + : {}), + }, + }) + if (info.notificationID) yield* runtime.job.completeBackground(info.notificationID) + }, + Effect.forkIn(scope, { startImmediately: true }), + ) yield* ctx.tool .transform((draft) => @@ -286,7 +275,7 @@ export const Plugin = { const settled = yield* Deferred.make() const run = settleShell().pipe( Effect.tap((output) => Deferred.succeed(settled, output)), - Effect.map((output) => output.output), + Effect.map((output) => resultMessages(output).join("\n\n")), Effect.onInterrupt(() => shell.remove(info.id).pipe(Effect.ignore)), ) const job = yield* runtime.job.start({ @@ -294,6 +283,12 @@ export const Plugin = { type: name, title: info.command, metadata: { sessionID: context.sessionID, shellID: info.id }, + recovery: { + kind: "shell", + sessionID: context.sessionID, + shellID: info.id, + command: info.command, + }, run, }) diff --git a/packages/core/src/tool/plugin/subagent.ts b/packages/core/src/tool/plugin/subagent.ts index aa620856094..163a7a5435d 100644 --- a/packages/core/src/tool/plugin/subagent.ts +++ b/packages/core/src/tool/plugin/subagent.ts @@ -78,22 +78,6 @@ export const Plugin = { return text.length > 0 ? text : NO_TEXT }) - const injectCompletion = Effect.fn("SubagentTool.injectCompletion")(function* ( - parentID: SessionSchema.ID, - childID: SessionSchema.ID, - agent: string, - description: string, - state: "completed" | "error" | "cancelled", - text: string, - ) { - yield* runtime.session.synthetic({ - sessionID: parentID, - text: `\n${text}\n`, - description, - metadata: { source: "subagent", childID, agent, state }, - }) - }) - const notifyWhenDone = Effect.fn("SubagentTool.notifyWhenDone")(function* ( parentID: SessionSchema.ID, childID: SessionSchema.ID, @@ -104,23 +88,24 @@ export const Plugin = { const key = `${childID}:${startedAt}` if (notifications.has(key)) return notifications.add(key) - yield* runtime.job.wait({ id: childID }).pipe( - Effect.flatMap((result) => { - if (result.info?.status === "completed") - return injectCompletion(parentID, childID, agent, description, "completed", result.info.output ?? NO_TEXT) - if (result.info?.status === "error") - return injectCompletion( - parentID, - childID, - agent, - description, - "error", - result.info.error ?? "Subagent failed", - ) - if (result.info?.status === "cancelled") - return injectCompletion(parentID, childID, agent, description, "cancelled", "Subagent cancelled") - return Effect.void - }), + 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) + }).pipe( Effect.ensuring(Effect.sync(() => notifications.delete(key))), Effect.forkIn(scope, { startImmediately: true }), ) @@ -239,6 +224,7 @@ export const Plugin = { existing === undefined ? ["You are a subagent spawned by another session.", input.prompt].join("\n") : input.prompt, + ...(background && existing === undefined ? { resume: false } : {}), }) .pipe( Effect.mapError( @@ -246,17 +232,19 @@ export const Plugin = { ), ) - const run = Effect.gen(function* () { - yield* runtime.session.resume(child.id) - return yield* latestAssistantText(child.id) - }).pipe(Effect.onInterrupt(() => runtime.session.interrupt(child.id))) - const info = yield* runtime.job.start({ id: child.id, type: name, title: input.description, metadata: {}, - run, + recovery: { + kind: "subagent", + parentSessionID: context.sessionID, + childSessionID: child.id, + agent: agent.name, + description: input.description, + }, + run: runtime.session.resume(child.id).pipe(Effect.andThen(latestAssistantText(child.id))), }) if (background) { diff --git a/packages/core/test/job.test.ts b/packages/core/test/job.test.ts index a725b349af5..e08159dfc5a 100644 --- a/packages/core/test/job.test.ts +++ b/packages/core/test/job.test.ts @@ -1,11 +1,13 @@ import { describe, expect } from "bun:test" import { Job } from "@opencode-ai/core/job" +import { KV } from "@opencode-ai/core/kv" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" +import { LayerNode } from "@opencode-ai/util/effect/layer-node" import { Deferred, Effect, Exit, Fiber, Scope } from "effect" import { SessionSchema } from "@opencode-ai/core/session/schema" import { testEffect } from "./lib/effect" -const it = testEffect(AppNodeBuilder.build(Job.node)) +const it = testEffect(AppNodeBuilder.build(LayerNode.group([Job.node, KV.node]))) describe("Job", () => { it.live("tracks process-local work through explicit observation", () => @@ -145,6 +147,177 @@ describe("Job", () => { }), ) + it.live("retains background ownership and terminal output until notification acknowledgment", () => + Effect.gen(function* () { + const jobs = yield* Job.Service + const latch = yield* Deferred.make() + const recovery = { + kind: "shell" as const, + sessionID: SessionSchema.ID.make("ses_background_shell"), + shellID: "shell_background", + command: "echo done", + } + const job = yield* jobs.start({ type: "shell", recovery, run: Deferred.await(latch).pipe(Effect.as("done")) }) + + expect((yield* jobs.pendingBackground).find((item) => item.id === job.id)).toBeUndefined() + const background = yield* jobs.background(job.id) + + const running = (yield* jobs.pendingBackground).find((item) => item.id === job.id) + expect(running).toMatchObject({ id: job.id, recovery, status: "running" }) + expect(running?.notificationID).toStartWith("msg_") + expect(background?.notificationID).toBe(running?.notificationID) + + yield* Deferred.succeed(latch, undefined) + yield* jobs.wait({ id: job.id }) + + const completed = (yield* jobs.pendingBackground).find((item) => item.id === job.id) + expect(completed).toMatchObject({ + id: job.id, + notificationID: running?.notificationID, + recovery, + status: "completed", + output: "done", + }) + if (!completed) return yield* Effect.die("background marker missing") + + yield* jobs.completeBackground(completed.notificationID) + expect((yield* jobs.pendingBackground).find((item) => item.id === job.id)).toBeUndefined() + }), + ) + + it.live("persists backgroundAll ownership before releasing a blocked subagent", () => + Effect.gen(function* () { + const jobs = yield* Job.Service + const parentSessionID = SessionSchema.ID.make("ses_background_parent") + const latch = yield* Deferred.make() + const recovery = { + kind: "subagent" as const, + parentSessionID, + childSessionID: SessionSchema.ID.make("ses_background_child"), + agent: "explore", + description: "Explore background recovery", + } + const job = yield* jobs.start({ type: "subagent", recovery, run: Deferred.await(latch).pipe(Effect.as("done")) }) + const waiting = yield* jobs + .block({ id: job.id, sessionID: parentSessionID }) + .pipe(Effect.forkIn(yield* Scope.Scope, { startImmediately: true })) + + yield* jobs.backgroundAll({ sessionID: parentSessionID }) + expect(yield* Fiber.join(waiting)).toMatchObject({ type: "backgrounded", info: { id: job.id } }) + + const marker = (yield* jobs.pendingBackground).find((item) => item.id === job.id) + expect(marker).toMatchObject({ id: job.id, recovery, status: "running" }) + if (!marker) return yield* Effect.die("background marker missing") + + yield* jobs.cancel(job.id) + expect((yield* jobs.pendingBackground).find((item) => item.id === job.id)).toMatchObject({ + notificationID: marker.notificationID, + status: "cancelled", + }) + yield* jobs.completeBackground(marker.notificationID) + }), + ) + + it.live("retains terminal errors for recovery until notification acknowledgment", () => + Effect.gen(function* () { + const jobs = yield* Job.Service + const latch = yield* Deferred.make() + const job = yield* jobs.start({ + type: "shell", + recovery: { + kind: "shell", + sessionID: SessionSchema.ID.make("ses_background_error"), + shellID: "shell_error", + command: "exit 1", + }, + run: Deferred.await(latch).pipe(Effect.andThen(Effect.fail(new Error("shell failed")))), + }) + + yield* jobs.background(job.id) + yield* Deferred.succeed(latch, undefined) + yield* jobs.wait({ id: job.id }) + + const marker = (yield* jobs.pendingBackground).find((item) => item.id === job.id) + expect(marker).toMatchObject({ id: job.id, status: "error", error: "shell failed" }) + if (!marker) return yield* Effect.die("background marker missing") + yield* jobs.completeBackground(marker.notificationID) + }), + ) + + it.live("durably backgrounds recoverable work that has already failed", () => + Effect.gen(function* () { + const jobs = yield* Job.Service + const job = yield* jobs.start({ + type: "shell", + recovery: { + kind: "shell", + sessionID: SessionSchema.ID.make("ses_immediate_error"), + shellID: "shell_immediate_error", + command: "exit 1", + }, + run: Effect.fail(new Error("shell failed")), + }) + expect((yield* jobs.wait({ id: job.id })).info?.status).toBe("error") + + const background = yield* jobs.background(job.id) + expect(background?.notificationID).toStartWith("msg_") + expect(yield* jobs.pendingBackground).toMatchObject([ + { id: job.id, notificationID: background?.notificationID, status: "error", error: "shell failed" }, + ]) + }), + ) + + it.live("recovers a background marker after its process-local registry closes", () => + Effect.gen(function* () { + const scope = yield* Scope.make() + const previous = yield* Job.make.pipe(Scope.provide(scope)) + const job = yield* previous.start({ + type: "shell", + recovery: { + kind: "shell", + sessionID: SessionSchema.ID.make("ses_background_restart"), + shellID: "shell_restart", + command: "sleep 60", + }, + run: Effect.never, + }) + yield* previous.background(job.id) + yield* Scope.close(scope, Exit.void) + + const current = yield* Job.make + const marker = (yield* current.pendingBackground).find((item) => item.id === job.id) + expect(marker).toMatchObject({ id: job.id, status: "running" }) + if (!marker) return yield* Effect.die("background marker missing") + yield* current.completeBackground(marker.notificationID) + }), + ) + + it.live("preserves running background ownership when its work is interrupted", () => + Effect.gen(function* () { + const jobs = yield* Job.Service + const interrupted = yield* Deferred.make() + const job = yield* jobs.start({ + type: "subagent", + recovery: { + kind: "subagent", + parentSessionID: SessionSchema.ID.make("ses_interrupted_parent"), + childSessionID: SessionSchema.ID.make("ses_interrupted_child"), + agent: "explore", + description: "Continue after shutdown", + }, + run: Deferred.await(interrupted).pipe(Effect.andThen(Effect.interrupt)), + }) + yield* jobs.background(job.id) + yield* Deferred.succeed(interrupted, undefined) + yield* jobs.wait({ id: job.id }) + + const marker = (yield* jobs.pendingBackground).find((item) => item.id === job.id) + expect(marker).toMatchObject({ id: job.id, status: "running" }) + if (!marker) return yield* Effect.die("background marker missing") + yield* jobs.completeBackground(marker.notificationID) + }), + ) + it.live("interrupts live work without promising settlement after the owning process-local scope closes", () => Effect.gen(function* () { const scope = yield* Scope.make() diff --git a/packages/core/test/npm.test.ts b/packages/core/test/npm.test.ts index 23f33455a58..4c542dabd86 100644 --- a/packages/core/test/npm.test.ts +++ b/packages/core/test/npm.test.ts @@ -223,6 +223,7 @@ describe("Npm.add", () => { ).toBeTruthy() }) + // Several real Git installs and refreshes exceed Bun's default timeout on Windows. test("refreshes mutable Git packages once per service lifetime and preserves pinned or cached installs", async () => { await using tmp = await tmpdir() const fixture = await createGitFixture(tmp.path) @@ -262,7 +263,7 @@ describe("Npm.add", () => { return yield* npm.add(mutable, { refresh: true }) }).pipe(Effect.scoped, Effect.provide(npmLayer(cache)), Effect.runPromise) expect(await Bun.file(path.join(offline.directory, "index.js")).text()).toContain('root: "second"') - }) + }, 30_000) }) describe("Npm.resolve", () => { diff --git a/packages/core/test/session-execution.test.ts b/packages/core/test/session-execution.test.ts index bdf2109cb2b..cb72952d737 100644 --- a/packages/core/test/session-execution.test.ts +++ b/packages/core/test/session-execution.test.ts @@ -4,6 +4,8 @@ import { Database } from "@opencode-ai/core/database/database" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" import { LayerNode } from "@opencode-ai/util/effect/layer-node" import { Bus } from "@opencode-ai/core/bus" +import { Job } from "@opencode-ai/core/job" +import { KV } from "@opencode-ai/core/kv" import { LocationServiceMap } from "@opencode-ai/core/location-service-map" import type { LocationServices } from "@opencode-ai/core/location-services" import { Project } from "@opencode-ai/core/project" @@ -23,7 +25,9 @@ import { Context, Deferred, Effect, Exit, Fiber, Layer, LayerMap, Scope } from " import { eq } from "drizzle-orm" import { testEffect } from "./lib/effect" -const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionStore.node]))) +const it = testEffect( + AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionStore.node, Job.node, KV.node, Session.node])), +) describe("SessionExecution lifecycle", () => { test("classifies success and typed failure terminals", () => { @@ -60,14 +64,13 @@ describe("SessionExecution lifecycle", () => { const idle = Session.ID.make("ses_recover_idle") yield* seedSessions(database, [parent], { time_suspended: Date.now() }) yield* seedSessions(database, [idle]) - // An orphaned child is never resumed: the resumed parent re-runs its - // tool call and spawns a fresh child instead. + // Children recover through background Job records, never through the root claim sweep. yield* seedSessions(database, [child], { time_suspended: Date.now(), parent_id: parent }) expect(yield* store.listSuspended()).toEqual([parent]) // The sweep clears orphaned child claims outright; parents keep theirs. - yield* store.releaseChildClaims + yield* store.releaseChildClaims([]) expect(yield* claims(database)).toEqual({ [parent]: true, [child]: false, [idle]: false }) }), ) @@ -147,6 +150,66 @@ describe("SessionExecution lifecycle", () => { }), ) + it.effect("does not resume a user-cancelled background child whose notification was not admitted", () => + Effect.gen(function* () { + const database = yield* Database.Service + const parent = Session.ID.make("ses_cancelled_background_parent") + const child = Session.ID.make("ses_cancelled_background_child") + yield* seedSessions(database, [parent]) + yield* seedSessions(database, [child], { parent_id: parent }) + + const running = yield* Deferred.make() + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const jobs = yield* Job.make.pipe(Scope.provide(scope)) + const context = yield* buildExecution( + scope, + () => Deferred.succeed(running, undefined).pipe(Effect.andThen(Effect.never)), + undefined, + jobs, + ) + const execution = Context.get(context, SessionExecution.Service) + yield* jobs.start({ + id: child, + type: "subagent", + recovery: { + kind: "subagent", + parentSessionID: parent, + childSessionID: child, + agent: "general", + description: "Cancelled inspection", + }, + run: execution.resume(child).pipe(Effect.as("unused")), + }) + yield* jobs.background(child) + yield* Deferred.await(running) + expect(yield* execution.interrupt(child)).toBeTrue() + yield* execution.awaitIdle(child) + expect((yield* jobs.wait({ id: child })).info?.status).toBe("cancelled") + expect(yield* jobs.pendingBackground).toMatchObject([{ id: child, status: "cancelled" }]) + expect((yield* claims(database))[child]).toBe(false) + yield* Scope.close(scope, Exit.void) + + const restartedScope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(restartedScope, Exit.void)) + const restartedJobs = yield* Job.make.pipe(Scope.provide(restartedScope)) + const drained: Session.ID[] = [] + const restarted = yield* buildExecution( + restartedScope, + ({ sessionID }) => Effect.sync(() => void drained.push(sessionID)), + undefined, + restartedJobs, + ) + yield* Context.get(restarted, SessionRestart.Service).resumeSuspendedSessions + yield* Context.get(restarted, SessionExecution.Service).awaitIdle(parent) + expect(drained).toEqual([parent]) + expect(yield* SessionInbox.list(database.db, parent)).toMatchObject([ + { payload: { text: expect.stringContaining("Subagent cancelled"), metadata: { state: "cancelled" } } }, + ]) + expect(yield* restartedJobs.pendingBackground).toEqual([]) + }), + ) + it.effect("starts every claimed execution without waiting for earlier drains to finish", () => Effect.gen(function* () { const database = yield* Database.Service @@ -304,6 +367,518 @@ describe("SessionExecution lifecycle", () => { ) }) +describe("SessionRestart background recovery", () => { + it.effect("admits orphaned shell notices without waking and delivers them once on the next run", () => + Effect.gen(function* () { + const database = yield* Database.Service + const store = yield* SessionStore.Service + const jobs = yield* Job.Service + const bus = yield* Bus.Service + const parent = Session.ID.make("ses_background_recovery_parent") + const child = Session.ID.make("ses_background_recovery_child") + yield* seedSessions(database, [parent]) + yield* seedSessions(database, [child], { parent_id: parent, time_suspended: Date.now() }) + yield* seedBackground(jobs, parent, [ + { id: "call-background-shell", shellID: "sh_background_orphan", command: "sleep 60" }, + ]) + yield* seedBackground(jobs, child, [{ id: "call-child-shell", shellID: "sh_child_orphan", command: "sleep 30" }]) + + expect(yield* store.listSuspended()).toEqual([]) + expect(yield* jobs.pendingBackground).toHaveLength(2) + + const drained: Session.ID[] = [] + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Effect.provideService(Scope.Scope, scope)) + const context = yield* buildExecution( + scope, + ({ sessionID }) => + Effect.sync(() => void drained.push(sessionID)).pipe( + Effect.andThen(SessionInbox.promote(database.db, bus, sessionID, "steer")), + Effect.asVoid, + ), + undefined, + restarted, + ) + const restart = Context.get(context, SessionRestart.Service) + yield* restart.resumeSuspendedSessions + + expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toEqual([]) + expect(yield* SessionInbox.list(database.db, parent)).toMatchObject([ + { + type: "synthetic", + payload: { + description: "sleep 60", + text: expect.stringContaining("server restarted"), + metadata: { + source: "shell", + jobID: "call-background-shell", + shellID: "sh_background_orphan", + state: "cancelled", + }, + }, + }, + ]) + expect(yield* SessionInbox.list(database.db, child)).toMatchObject([ + { + type: "synthetic", + payload: { + metadata: { + source: "shell", + jobID: "call-child-shell", + shellID: "sh_child_orphan", + state: "cancelled", + }, + }, + }, + ]) + expect(drained).toEqual([]) + expect(yield* claims(database)).toEqual({ [parent]: false, [child]: false }) + expect(yield* restarted.pendingBackground).toEqual([]) + + yield* restart.resumeSuspendedSessions + expect(yield* SessionInbox.list(database.db, parent)).toHaveLength(1) + expect(drained).toEqual([]) + const execution = Context.get(context, SessionExecution.Service) + yield* execution.resume(parent) + expect(yield* SessionInbox.list(database.db, parent)).toEqual([]) + expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toHaveLength(1) + yield* execution.resume(parent) + expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toHaveLength(1) + }), + ) + + it.effect("preserves locally running background work", () => + Effect.gen(function* () { + const database = yield* Database.Service + const store = yield* SessionStore.Service + const jobs = yield* Job.Service + const parent = Session.ID.make("ses_background_existing_parent") + yield* seedSessions(database, [parent]) + yield* seedBackground(jobs, parent, [{ id: "call-running-shell", shellID: "sh_running", command: "sleep 60" }]) + + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const context = yield* buildExecution(scope, () => Effect.void) + const restart = Context.get(context, SessionRestart.Service) + yield* restart.resumeSuspendedSessions + + expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toEqual([]) + expect(yield* jobs.get("call-running-shell")).toMatchObject({ status: "running" }) + expect(yield* jobs.pendingBackground).toHaveLength(1) + }), + ) + + it.effect("preserves a silent shell failure persisted before its completion notification", () => + Effect.gen(function* () { + const database = yield* Database.Service + const jobs = yield* Job.Service + const sessionID = Session.ID.make("ses_background_completed_shell") + yield* seedSessions(database, [sessionID]) + const complete = yield* Deferred.make() + yield* jobs.start({ + id: "call-completed-shell", + type: "shell", + recovery: { + kind: "shell", + sessionID, + shellID: "sh_completed", + command: "exit 7", + }, + run: Deferred.await(complete), + }) + yield* jobs.background("call-completed-shell") + yield* Deferred.succeed(complete, "(no output)\n\nCommand exited with code 7.") + yield* jobs.wait({ id: "call-completed-shell" }) + + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Effect.provideService(Scope.Scope, scope)) + const context = yield* buildExecution(scope, () => Effect.void, undefined, restarted) + yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions + + expect(yield* SessionInbox.list(database.db, sessionID)).toMatchObject([ + { + type: "synthetic", + payload: { + text: expect.stringContaining("(no output)\n\nCommand exited with code 7."), + metadata: { source: "shell", shellID: "sh_completed", state: "completed" }, + }, + }, + ]) + expect(yield* restarted.pendingBackground).toEqual([]) + }), + ) + + for (const delivered of [false, true]) { + it.effect(`does not duplicate a shell notification already ${delivered ? "delivered" : "admitted"}`, () => + Effect.gen(function* () { + const database = yield* Database.Service + const bus = yield* Bus.Service + const jobs = yield* Job.Service + const sessions = yield* Session.Service + const sessionID = Session.ID.make("ses_shell_notification_retry") + yield* seedSessions(database, [sessionID]) + yield* seedBackground(jobs, sessionID, [ + { id: "call-shell-notified", shellID: "sh_notified", command: "echo done" }, + ]) + const background = (yield* jobs.pendingBackground)[0] + if (!background) return yield* Effect.die("background record missing") + yield* sessions.synthetic({ + id: background.notificationID, + sessionID, + text: "Command already completed", + metadata: { source: "shell", shellID: "sh_notified", state: "completed" }, + resume: false, + }) + if (delivered) yield* SessionInbox.promote(database.db, bus, sessionID, "steer") + + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Effect.provideService(Scope.Scope, scope)) + const context = yield* buildExecution(scope, () => Effect.void, undefined, restarted) + yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions + + expect(yield* restarted.pendingBackground).toEqual([]) + expect(yield* SessionInbox.list(database.db, sessionID)).toHaveLength(delivered ? 0 : 1) + yield* SessionInbox.promote(database.db, bus, sessionID, "steer") + expect(yield* sessions.messages({ sessionID })).toMatchObject([ + { + id: background.notificationID, + type: "synthetic", + text: "Command already completed", + metadata: { state: "completed" }, + }, + ]) + expect(yield* sessions.messages({ sessionID })).toHaveLength(1) + }), + ) + } + + it.effect("acknowledges recovery markers when their owning session is deleted", () => + Effect.gen(function* () { + const database = yield* Database.Service + const jobs = yield* Job.Service + const sessionID = Session.ID.make("ses_background_deleted") + yield* seedSessions(database, [sessionID]) + yield* seedBackground(jobs, sessionID, [{ id: "call-deleted-shell", shellID: "sh_deleted", command: "sleep 60" }]) + yield* database.db.delete(SessionTable).where(eq(SessionTable.id, sessionID)).run().pipe(Effect.orDie) + + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Effect.provideService(Scope.Scope, scope)) + const context = yield* buildExecution(scope, () => Effect.void, undefined, restarted) + yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions + + expect(yield* restarted.pendingBackground).toEqual([]) + }), + ) + + it.effect("delivers cancellation at the resumed parent's next step", () => + Effect.gen(function* () { + const database = yield* Database.Service + const jobs = yield* Job.Service + const store = yield* SessionStore.Service + const bus = yield* Bus.Service + const parent = Session.ID.make("ses_background_claimed_parent") + yield* seedSessions(database, [parent], { time_suspended: Date.now() }) + yield* seedBackground(jobs, parent, [{ id: "call-claimed-shell", shellID: "sh_claimed", command: "sleep 60" }]) + + const observed = yield* Deferred.make() + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Effect.provideService(Scope.Scope, scope)) + const context = yield* buildExecution( + scope, + ({ sessionID }) => + SessionInbox.promote(database.db, bus, sessionID, "steer").pipe( + Effect.andThen(store.context(sessionID)), + Effect.orDie, + Effect.flatMap((messages) => + Deferred.succeed( + observed, + messages.filter((message) => message.type === "synthetic").map((message) => message.text), + ), + ), + Effect.asVoid, + ), + undefined, + restarted, + ) + const execution = Context.get(context, SessionExecution.Service) + yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions + expect(yield* Deferred.await(observed)).toEqual([ + "The server restarted while you were working. Continue from where you left off without repeating completed work.", + expect.stringContaining("Command cancelled because the server restarted"), + ]) + yield* execution.awaitIdle(parent) + expect(yield* SessionInbox.list(database.db, parent)).toEqual([]) + expect((yield* claims(database))[parent]).toBe(false) + }), + ) + + it.effect("resumes a background subagent and notifies its parent exactly once", () => + Effect.gen(function* () { + const database = yield* Database.Service + const jobs = yield* Job.Service + const parent = Session.ID.make("ses_subagent_recovery_parent") + const child = Session.ID.make("ses_subagent_recovery_child") + const unrelated = Session.ID.make("ses_subagent_unrelated_child") + yield* seedSessions(database, [parent], { time_suspended: Date.now(), resume_attempts: 1 }) + yield* seedSessions(database, [child, unrelated], { parent_id: parent, time_suspended: Date.now() }) + yield* jobs.start({ + id: child, + type: "subagent", + recovery: { + kind: "subagent", + parentSessionID: parent, + childSessionID: child, + agent: "explore", + description: "Inspect recovery", + }, + run: Effect.never, + }) + yield* jobs.background(child) + + const resumed = yield* Deferred.make() + const release = yield* Deferred.make() + const parentResumed = yield* Deferred.make() + const parentWoken = yield* Deferred.make() + const drained: Session.ID[] = [] + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Effect.provideService(Scope.Scope, scope)) + const context = yield* buildExecution( + scope, + ({ sessionID }) => + Effect.gen(function* () { + drained.push(sessionID) + if (sessionID === child) { + yield* Deferred.succeed(resumed, undefined) + yield* Deferred.await(release) + return + } + yield* Deferred.succeed( + drained.filter((id) => id === parent).length === 1 ? parentResumed : parentWoken, + undefined, + ) + }), + undefined, + restarted, + ) + const restart = Context.get(context, SessionRestart.Service) + const execution = Context.get(context, SessionExecution.Service) + yield* restart.resumeSuspendedSessions + yield* Deferred.await(resumed) + yield* Deferred.await(parentResumed) + yield* execution.awaitIdle(parent) + + yield* restart.resumeSuspendedSessions + expect(drained.toSorted()).toEqual([child, parent].toSorted()) + expect(yield* claims(database)).toEqual({ [parent]: false, [child]: true, [unrelated]: false }) + expect(yield* attempts(database, child)).toBe(1) + expect(yield* restarted.get(child)).toMatchObject({ status: "running" }) + + yield* Deferred.succeed(release, undefined) + yield* Deferred.await(parentWoken) + expect(drained.filter((id) => id === child)).toHaveLength(1) + expect(drained.filter((id) => id === parent)).toHaveLength(2) + expect(yield* SessionInbox.list(database.db, parent)).toMatchObject([ + { + payload: { + description: "Inspect recovery", + metadata: { source: "subagent", childID: child, agent: "explore", state: "completed" }, + }, + }, + ]) + expect(yield* restarted.pendingBackground).toEqual([]) + yield* restart.resumeSuspendedSessions + expect(yield* SessionInbox.list(database.db, parent)).toHaveLength(1) + }), + ) + + it.effect("delivers a subagent result persisted before restart without rerunning the child", () => + Effect.gen(function* () { + const database = yield* Database.Service + const jobs = yield* Job.Service + const parent = Session.ID.make("ses_subagent_completed_parent") + const child = Session.ID.make("ses_subagent_completed_child") + yield* seedSessions(database, [parent]) + yield* seedSessions(database, [child], { parent_id: parent }) + const complete = yield* Deferred.make() + yield* jobs.start({ + id: child, + type: "subagent", + recovery: { + kind: "subagent", + parentSessionID: parent, + childSessionID: child, + agent: "explore", + description: "Completed inspection", + }, + run: Deferred.await(complete), + }) + yield* jobs.background(child) + yield* Deferred.succeed(complete, "Recovered result") + yield* jobs.wait({ id: child }) + + const parentWoken = yield* Deferred.make() + const drained: Session.ID[] = [] + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Effect.provideService(Scope.Scope, scope)) + const context = yield* buildExecution( + scope, + ({ sessionID }) => + Effect.sync(() => void drained.push(sessionID)).pipe( + Effect.andThen(Deferred.succeed(parentWoken, undefined)), + ), + undefined, + restarted, + ) + yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions + yield* Deferred.await(parentWoken) + + expect(drained).toEqual([parent]) + expect(yield* SessionInbox.list(database.db, parent)).toMatchObject([ + { payload: { text: expect.stringContaining("Recovered result"), metadata: { state: "completed" } } }, + ]) + expect(yield* restarted.pendingBackground).toEqual([]) + }), + ) + + for (const resumeAttempts of [1, 2]) { + it.effect(`honors a suspended parent's restart budget after ${resumeAttempts} attempts before notifying it`, () => + Effect.gen(function* () { + const database = yield* Database.Service + const bus = yield* Bus.Service + const jobs = yield* Job.Service + const parent = Session.ID.make("ses_subagent_budget_parent") + const children = [ + Session.ID.make("ses_subagent_budget_child_1"), + Session.ID.make("ses_subagent_budget_child_2"), + ] + yield* seedSessions(database, [parent], { time_suspended: Date.now(), resume_attempts: resumeAttempts }) + yield* seedSessions(database, children, { parent_id: parent }) + const complete = yield* Deferred.make() + for (const child of children) { + yield* jobs.start({ + id: child, + type: "subagent", + recovery: { + kind: "subagent", + parentSessionID: parent, + childSessionID: child, + agent: "explore", + description: "Completed inspection", + }, + run: Deferred.await(complete), + }) + yield* jobs.background(child) + } + yield* Deferred.succeed(complete, "Recovered result") + yield* Effect.forEach(children, (id) => jobs.wait({ id }), { discard: true }) + + const draining = yield* Deferred.make() + const release = yield* Deferred.make() + const drained: Session.ID[] = [] + const continued: Session.ID[] = [] + yield* bus.project(SessionEvent.Synthetic, (event) => + Effect.sync(() => void continued.push(event.data.sessionID)), + ) + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Scope.provide(scope)) + const context = yield* buildExecution( + scope, + ({ sessionID }) => + Effect.gen(function* () { + drained.push(sessionID) + yield* Deferred.succeed(draining, yield* attempts(database, sessionID)) + yield* Deferred.await(release) + }), + { maxAttempts: 2 }, + restarted, + ) + const restart = Context.get(context, SessionRestart.Service) + const execution = Context.get(context, SessionExecution.Service) + yield* restart.resumeSuspendedSessions + + if (resumeAttempts < 2) { + expect(yield* Deferred.await(draining)).toBe(2) + expect(drained).toEqual([parent]) + expect(continued).toEqual([parent]) + yield* Deferred.succeed(release, undefined) + yield* execution.awaitIdle(parent) + } + if (resumeAttempts === 2) { + expect(drained).toEqual([]) + expect(continued).toEqual([]) + } + expect((yield* claims(database))[parent]).toBe(false) + expect(yield* attempts(database, parent)).toBe(0) + expect(yield* SessionInbox.list(database.db, parent)).toHaveLength(2) + expect(yield* restarted.pendingBackground).toEqual([]) + yield* restart.resumeSuspendedSessions + expect(drained).toHaveLength(resumeAttempts < 2 ? 1 : 0) + }), + ) + } + + it.effect("terminalizes a recovered subagent that exhausts its resume budget", () => + Effect.gen(function* () { + const database = yield* Database.Service + const jobs = yield* Job.Service + const parent = Session.ID.make("ses_subagent_exhausted_parent") + const child = Session.ID.make("ses_subagent_exhausted_child") + yield* seedSessions(database, [parent]) + yield* seedSessions(database, [child], { parent_id: parent, time_suspended: Date.now(), resume_attempts: 2 }) + yield* jobs.start({ + id: child, + type: "subagent", + recovery: { + kind: "subagent", + parentSessionID: parent, + childSessionID: child, + agent: "explore", + description: "Exhausted inspection", + }, + run: Effect.never, + }) + yield* jobs.background(child) + + const parentWoken = yield* Deferred.make() + const drained: Session.ID[] = [] + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)) + const restarted = yield* Job.make.pipe(Effect.provideService(Scope.Scope, scope)) + const context = yield* buildExecution( + scope, + ({ sessionID }) => + Effect.sync(() => void drained.push(sessionID)).pipe( + Effect.andThen(Deferred.succeed(parentWoken, undefined)), + ), + { maxAttempts: 2 }, + restarted, + ) + yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions + yield* Deferred.await(parentWoken) + + expect(drained).toEqual([parent]) + expect((yield* claims(database))[child]).toBe(false) + expect(yield* SessionInbox.list(database.db, parent)).toMatchObject([ + { + payload: { + text: expect.stringContaining("will not be resumed automatically"), + metadata: { source: "subagent", childID: child, state: "error" }, + }, + }, + ]) + expect(yield* restarted.pendingBackground).toEqual([]) + }), + ) +}) + describe("SessionExecution interrupt continuation", () => { it.effect("resumes only steering input after an interrupt with continue", () => Effect.gen(function* () { @@ -446,6 +1021,27 @@ describe("SessionExecution interrupt continuation", () => { ) }) +function seedBackground( + jobs: Job.Interface, + sessionID: Session.ID, + background: ReadonlyArray<{ readonly id: string; readonly shellID: string; readonly command: string }>, +) { + return Effect.forEach( + background, + (job) => + Effect.gen(function* () { + yield* jobs.start({ + id: job.id, + type: "shell", + recovery: { kind: "shell", sessionID, shellID: job.shellID, command: job.command }, + run: Effect.never, + }) + yield* jobs.background(job.id) + }), + { discard: true }, + ) +} + /** Plain deliveries seed user prompts; objects seed control items. */ function seedInbox( database: Database.Service["Service"], @@ -531,11 +1127,27 @@ function buildExecution( scope: Scope.Closeable, drain: (input: Parameters[0]) => Effect.Effect, options?: SessionRestart.Options, + overrideJobs?: Job.Interface, ) { return Effect.gen(function* () { const database = yield* Database.Service const bus = yield* Bus.Service const store = yield* SessionStore.Service + const jobs = overrideJobs ?? (yield* Job.Service) + const sessions = yield* Session.Service + const sessionLayer = Layer.effect( + Session.Service, + Effect.gen(function* () { + const execution = yield* SessionExecution.Service + return Session.Service.of({ + ...sessions, + synthetic: (input) => + sessions + .synthetic({ ...input, resume: false }) + .pipe(Effect.tap(() => (input.resume === false ? Effect.void : execution.wake(input.sessionID)))), + }) + }), + ) const runner = Layer.succeed( SessionRunner.Service, SessionRunner.Service.of({ @@ -553,10 +1165,12 @@ function buildExecution( ) return yield* Layer.buildWithScope( SessionRestart.layer(options).pipe( - Layer.provideMerge(SessionExecution.layer), + Layer.provideMerge(sessionLayer), + Layer.provideMerge(Layer.fresh(SessionExecution.layer)), Layer.provide(Layer.succeed(Database.Service, database)), Layer.provide(Layer.succeed(Bus.Service, bus)), Layer.provide(Layer.succeed(SessionStore.Service, store)), + Layer.provide(Layer.succeed(Job.Service, jobs)), Layer.provide(locations), ), scope, diff --git a/packages/core/test/session-runner-tool-events.test.ts b/packages/core/test/session-runner-tool-events.test.ts index e6f80c30933..ea1b83eb5b7 100644 --- a/packages/core/test/session-runner-tool-events.test.ts +++ b/packages/core/test/session-runner-tool-events.test.ts @@ -128,6 +128,35 @@ test("interrupted progress metadata remains in the terminal failure snapshot", a }) }) +test("interrupted subagent failures expose their existing child session to the model", async () => { + const { published, publisher } = capture("anthropic", { interruptProgress: true }) + const subagent = LLMEvent.toolCall({ + id: "call-subagent", + name: "subagent", + input: { agent: "general", description: "Recover child", prompt: "Continue working" }, + }) + await Effect.runPromise(publisher.publish(subagent)) + await Effect.runPromiseExit(publisher.progress(subagent.id, { sessionID: "ses_existing_child", status: "running" })) + await Effect.runPromise(publisher.failUnsettledTools({ type: "aborted", message: "Tool execution interrupted" })) + + expect(published.find((event) => event.type === "session.tool.failed.2")?.data).toMatchObject({ + error: { type: "aborted", message: "Tool execution interrupted (sessionID: ses_existing_child)" }, + metadata: { sessionID: "ses_existing_child", status: "running" }, + }) +}) + +test("interrupted non-subagent failures do not expose their progress session IDs", async () => { + const { published, publisher } = capture() + await Effect.runPromise(publisher.publish(call)) + await Effect.runPromise(publisher.progress(call.id, { sessionID: "ses_private", status: "running" })) + await Effect.runPromise(publisher.failUnsettledTools({ type: "aborted", message: "Tool execution interrupted" })) + + expect(published.find((event) => event.type === "session.tool.failed.2")?.data).toMatchObject({ + error: { type: "aborted", message: "Tool execution interrupted" }, + metadata: { sessionID: "ses_private", status: "running" }, + }) +}) + test("local failure metadata completes the progress snapshot", async () => { const { published, publisher } = capture() await Effect.runPromise(publisher.publish(call)) diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index 22d339ac67c..f8257111002 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -79,7 +79,7 @@ import { Provider } from "@opencode-ai/core/provider" import { Cause, Deferred, Effect, Exit, Fiber, Layer, Queue, Schema, Scope, Stream } from "effect" import { TestClock } from "effect/testing" import { HttpClientRequest, HttpClientResponse } from "effect/unstable/http" -import { asc, desc, eq } from "drizzle-orm" +import { asc, desc, eq, sql } from "drizzle-orm" import { testEffect } from "./lib/effect" import { permissionLayer } from "./lib/permission" import { agentHost, catalogHost, host } from "./plugin/host" @@ -3640,6 +3640,81 @@ describe("SessionRunnerLLM", () => { }), ) + it.effect("preserves a stale subagent child session in its model-visible failure", () => + Effect.gen(function* () { + const session = yield* setup + const bus = yield* Bus.Service + const database = yield* Database.Service + yield* admit(session, "Recover interrupted subagent") + yield* SessionInbox.promote(database.db, bus, sessionID, "steer") + const assistantMessageID = SessionMessage.ID.create() + yield* bus.publish(SessionEvent.Step.Started, { + sessionID, + assistantMessageID, + agent: Agent.ID.make("build"), + model: { id: ID.make("fake-model"), providerID: Provider.ID.make("fake") }, + }) + yield* bus.publish(SessionEvent.Tool.Input.Started, { + sessionID, + assistantMessageID, + id: "call-interrupted-subagent", + name: "subagent", + }) + yield* bus.publish(SessionEvent.Tool.Input.Ended, { + sessionID, + assistantMessageID, + id: "call-interrupted-subagent", + text: '{"agent":"general"}', + }) + yield* bus.publish(SessionEvent.Tool.Called, { + sessionID, + assistantMessageID, + id: "call-interrupted-subagent", + input: { agent: "general" }, + executed: false, + }) + yield* database.db + .update(SessionMessageTable) + .set({ + data: sql`json_set( + ${SessionMessageTable.data}, + '$.content[0].state.metadata', + json('{"sessionID":"ses_existing_child","status":"running","internal":"private"}') + )`, + }) + .where(eq(SessionMessageTable.id, assistantMessageID)) + .run() + .pipe(Effect.orDie) + requests.length = 0 + yield* TestLLM.push([]) + yield* session.resume(sessionID) + + expect(yield* session.context(sessionID)).toMatchObject([ + { type: "user", text: "Recover interrupted subagent" }, + { + type: "assistant", + content: [ + { + type: "tool", + id: "call-interrupted-subagent", + state: { + status: "error", + error: { + type: "aborted", + message: "Tool execution interrupted: subagent (sessionID: ses_existing_child)", + }, + metadata: { sessionID: "ses_existing_child", status: "running", internal: "private" }, + }, + }, + ], + }, + ]) + const modelResult = JSON.stringify(requests[0]?.messages.at(-1)) + expect(modelResult).toContain("ses_existing_child") + expect(modelResult).not.toContain("private") + }), + ) + it.effect("durably fails hosted tools left running by a prior process before continuing inline", () => Effect.gen(function* () { const session = yield* setup diff --git a/packages/core/test/tool-shell.test.ts b/packages/core/test/tool-shell.test.ts index f22de91bd8d..a924a5ec5a0 100644 --- a/packages/core/test/tool-shell.test.ts +++ b/packages/core/test/tool-shell.test.ts @@ -790,6 +790,48 @@ describe("ShellTool", () => { ), ) + it.live("persists a silent command that finishes before backgrounding", () => + Effect.acquireUseRelease( + Effect.promise(() => tmpdir()), + (tmp) => { + reset() + return withSession(tmp.path, (registry) => + Effect.gen(function* () { + const bus = yield* Bus.Service + const jobs = yield* Job.Service + const shell = yield* Shell.Service + const persisted = yield* Deferred.make() + yield* bus.project(SessionEvent.InboxEnqueued, (event) => + event.data.sessionID === sessionID && event.data.item.type === "synthetic" + ? jobs.pendingBackground.pipe( + Effect.flatMap((background) => Deferred.succeed(persisted, background)), + Effect.asVoid, + ) + : Effect.void, + ) + yield* executeTool(registry, { + ...call({ command: "exit 7", background: true }, "call-background-silent-nonzero"), + // The command can finish while its initial progress update is being published. + progress: (update) => + typeof update.shellID === "string" + ? shell.wait(ShellSchema.ID.make(update.shellID)).pipe(Effect.orDie, Effect.asVoid) + : Effect.void, + }) + + expect(yield* Deferred.await(persisted)).toMatchObject([ + { + id: "call-background-silent-nonzero", + status: "completed", + output: "(no output)\n\nCommand exited with code 7.", + }, + ]) + }), + ) + }, + (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]().then(() => undefined)), + ), + ) + it.live( "updates and clears a running shell timeout", () => diff --git a/packages/tui/src/routes/session/index.tsx b/packages/tui/src/routes/session/index.tsx index 6c7a1310f9e..e9da715c10d 100644 --- a/packages/tui/src/routes/session/index.tsx +++ b/packages/tui/src/routes/session/index.tsx @@ -360,7 +360,8 @@ export function Session(props: { createEffect(() => { if (restored || !synced() || !rowsSynced() || !scroll || scroll.isDestroyed) return restored = true - restoreScrollPosition() + // Initial synchronization can finish after the reader has already navigated. + if (!isAwayFromBottom()) restoreScrollPosition() }) let awayTimer: ReturnType | undefined onCleanup(() => { diff --git a/packages/www/openapi.json b/packages/www/openapi.json index 238ca367b22..3539338410c 100644 --- a/packages/www/openapi.json +++ b/packages/www/openapi.json @@ -3950,6 +3950,143 @@ }, "description": "Retrieve one projected message owned by the Session.", "summary": "Get session message" + }, + "patch": { + "tags": ["session"], + "operationId": "v2.session.messageUpdate", + "parameters": [ + { + "name": "sessionID", + "in": "path", + "schema": { + "type": "string", + "pattern": "^ses" + }, + "required": true + }, + { + "name": "messageID", + "in": "path", + "schema": { + "type": "string", + "pattern": "^msg_" + }, + "required": true + } + ], + "security": [], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "data": { + "$ref": "#/components/schemas/Session.Message.Assistant" + } + }, + "required": ["data"], + "additionalProperties": false + } + } + } + }, + "400": { + "description": "InvalidRequestError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + }, + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + } + ] + } + } + } + }, + "401": { + "description": "UnauthorizedError", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UnauthorizedErrorEncoded" + } + } + } + }, + "404": { + "description": "SessionNotFoundError | MessageNotFoundError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" + }, + { + "$ref": "#/components/schemas/MessageNotFoundErrorEncoded" + } + ] + } + } + } + }, + "409": { + "description": "SessionBusyError | ConflictError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/SessionBusyErrorEncoded" + }, + { + "$ref": "#/components/schemas/ConflictErrorEncoded" + } + ] + } + } + } + } + }, + "description": "Replace the content of a completed assistant message in an idle session.", + "summary": "Update assistant message content", + "requestBody": { + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "content": { + "type": "array", + "items": { + "anyOf": [ + { + "$ref": "#/components/schemas/Session.Message.Assistant.Text" + }, + { + "$ref": "#/components/schemas/Session.Message.Assistant.Reasoning" + }, + { + "$ref": "#/components/schemas/Session.Message.Assistant.Tool" + } + ] + } + } + }, + "required": ["content"], + "additionalProperties": false + } + } + }, + "required": true + } } }, "/api/session/{sessionID}/environment": { diff --git a/packages/www/public/openapi.json b/packages/www/public/openapi.json index 238ca367b22..3539338410c 100644 --- a/packages/www/public/openapi.json +++ b/packages/www/public/openapi.json @@ -3950,6 +3950,143 @@ }, "description": "Retrieve one projected message owned by the Session.", "summary": "Get session message" + }, + "patch": { + "tags": ["session"], + "operationId": "v2.session.messageUpdate", + "parameters": [ + { + "name": "sessionID", + "in": "path", + "schema": { + "type": "string", + "pattern": "^ses" + }, + "required": true + }, + { + "name": "messageID", + "in": "path", + "schema": { + "type": "string", + "pattern": "^msg_" + }, + "required": true + } + ], + "security": [], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "data": { + "$ref": "#/components/schemas/Session.Message.Assistant" + } + }, + "required": ["data"], + "additionalProperties": false + } + } + } + }, + "400": { + "description": "InvalidRequestError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + }, + { + "$ref": "#/components/schemas/InvalidRequestErrorEncoded" + } + ] + } + } + } + }, + "401": { + "description": "UnauthorizedError", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UnauthorizedErrorEncoded" + } + } + } + }, + "404": { + "description": "SessionNotFoundError | MessageNotFoundError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/SessionNotFoundErrorEncoded" + }, + { + "$ref": "#/components/schemas/MessageNotFoundErrorEncoded" + } + ] + } + } + } + }, + "409": { + "description": "SessionBusyError | ConflictError", + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/SessionBusyErrorEncoded" + }, + { + "$ref": "#/components/schemas/ConflictErrorEncoded" + } + ] + } + } + } + } + }, + "description": "Replace the content of a completed assistant message in an idle session.", + "summary": "Update assistant message content", + "requestBody": { + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "content": { + "type": "array", + "items": { + "anyOf": [ + { + "$ref": "#/components/schemas/Session.Message.Assistant.Text" + }, + { + "$ref": "#/components/schemas/Session.Message.Assistant.Reasoning" + }, + { + "$ref": "#/components/schemas/Session.Message.Assistant.Tool" + } + ] + } + } + }, + "required": ["content"], + "additionalProperties": false + } + } + }, + "required": true + } } }, "/api/session/{sessionID}/environment": {