mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-28 06:20:15 -04:00
fix(core): wake sessions for recovered shell outcomes (#45781)
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@opencode-ai/core": patch
|
||||
---
|
||||
|
||||
Wake idle sessions when recovering background shell outcomes after a server restart. Admit shell outcomes before resuming background child sessions, while preserving restart retry budgets and completion-driven parent notifications.
|
||||
@@ -98,6 +98,7 @@ export const layer = (options?: Options) =>
|
||||
const recoverShell = Effect.fnUntraced(function* (
|
||||
background: Job.Background,
|
||||
recovery: Extract<Job.Recovery, { kind: "shell" }>,
|
||||
suspended: ReadonlySet<SessionSchema.ID>,
|
||||
) {
|
||||
const state = background.status === "running" ? "cancelled" : background.status
|
||||
const text =
|
||||
@@ -121,7 +122,7 @@ export const layer = (options?: Options) =>
|
||||
shellID: recovery.shellID,
|
||||
state,
|
||||
},
|
||||
resume: false,
|
||||
...(suspended.has(recovery.sessionID) ? { resume: false } : {}),
|
||||
})
|
||||
.pipe(
|
||||
Effect.catchTag("Session.NotFoundError", () => Effect.void),
|
||||
@@ -211,23 +212,25 @@ export const layer = (options?: Options) =>
|
||||
return Service.of({
|
||||
resumeSuspendedSessions: Effect.gen(function* () {
|
||||
const active = yield* execution.active
|
||||
// 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]
|
||||
: [],
|
||||
),
|
||||
const children = pending.flatMap((background) =>
|
||||
background.status === "running" && background.recovery.kind === "subagent"
|
||||
? [background.recovery.childSessionID]
|
||||
: [],
|
||||
)
|
||||
// Early notices wait for recovery's accounting, including Sessions that exhaust their budget.
|
||||
const suspended = new Set(
|
||||
[...(yield* store.listSuspended()), ...children].filter((sessionID) => !active.has(sessionID)),
|
||||
)
|
||||
yield* store.releaseChildClaims(children)
|
||||
yield* Effect.forEach(
|
||||
pending,
|
||||
// Admit shell outcomes before a recovered child can start its first model request.
|
||||
pending.toSorted((a, b) => Number(a.recovery.kind === "subagent") - Number(b.recovery.kind === "subagent")),
|
||||
Effect.fnUntraced(function* (background) {
|
||||
if ((yield* jobs.get(background.id))?.status === "running") return
|
||||
const recovery = background.recovery
|
||||
yield* recovery.kind === "shell"
|
||||
? recoverShell(background, recovery)
|
||||
? recoverShell(background, recovery, suspended)
|
||||
: recoverSubagent(background, recovery, suspended)
|
||||
}),
|
||||
{ discard: true },
|
||||
|
||||
@@ -368,7 +368,7 @@ describe("SessionExecution lifecycle", () => {
|
||||
})
|
||||
|
||||
describe("SessionRestart background recovery", () => {
|
||||
it.effect("admits orphaned shell notices without waking and delivers them once on the next run", () =>
|
||||
it.effect("wakes idle shell owners and delivers recovered notices exactly once", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const store = yield* SessionStore.Service
|
||||
@@ -401,50 +401,44 @@ describe("SessionRestart background recovery", () => {
|
||||
restarted,
|
||||
)
|
||||
const restart = Context.get(context, SessionRestart.Service)
|
||||
const execution = Context.get(context, SessionExecution.Service)
|
||||
yield* restart.resumeSuspendedSessions
|
||||
yield* Effect.forEach([parent, child], execution.awaitIdle, { discard: true })
|
||||
|
||||
expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, parent)).toMatchObject([
|
||||
expect(drained.toSorted()).toEqual([parent, child].toSorted())
|
||||
expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).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",
|
||||
},
|
||||
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([
|
||||
expect((yield* store.context(child)).filter((message) => message.type === "synthetic")).toMatchObject([
|
||||
{
|
||||
type: "synthetic",
|
||||
payload: {
|
||||
metadata: {
|
||||
source: "shell",
|
||||
jobID: "call-child-shell",
|
||||
shellID: "sh_child_orphan",
|
||||
state: "cancelled",
|
||||
},
|
||||
metadata: {
|
||||
source: "shell",
|
||||
jobID: "call-child-shell",
|
||||
shellID: "sh_child_orphan",
|
||||
state: "cancelled",
|
||||
},
|
||||
},
|
||||
])
|
||||
expect(drained).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, parent)).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, child)).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(drained).toHaveLength(2)
|
||||
expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toHaveLength(1)
|
||||
expect((yield* store.context(child)).filter((message) => message.type === "synthetic")).toHaveLength(1)
|
||||
}),
|
||||
)
|
||||
|
||||
@@ -469,7 +463,7 @@ describe("SessionRestart background recovery", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("preserves a silent shell failure persisted before its completion notification", () =>
|
||||
it.effect("wakes the owner for a silent shell failure persisted before its completion notification", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const jobs = yield* Job.Service
|
||||
@@ -494,9 +488,17 @@ describe("SessionRestart background recovery", () => {
|
||||
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)
|
||||
const drained: Session.ID[] = []
|
||||
const context = yield* buildExecution(
|
||||
scope,
|
||||
({ sessionID }) => Effect.sync(() => void drained.push(sessionID)),
|
||||
undefined,
|
||||
restarted,
|
||||
)
|
||||
yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions
|
||||
yield* Context.get(context, SessionExecution.Service).awaitIdle(sessionID)
|
||||
|
||||
expect(drained).toEqual([sessionID])
|
||||
expect(yield* SessionInbox.list(database.db, sessionID)).toMatchObject([
|
||||
{
|
||||
type: "synthetic",
|
||||
@@ -617,6 +619,145 @@ describe("SessionRestart background recovery", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("does not bypass a claimed owner's restart budget for a shell notice", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const jobs = yield* Job.Service
|
||||
const sessionID = Session.ID.make("ses_shell_recovery_exhausted")
|
||||
yield* seedSessions(database, [sessionID], { time_suspended: Date.now(), resume_attempts: 2 })
|
||||
yield* seedBackground(jobs, sessionID, [
|
||||
{ id: "call-exhausted-shell", shellID: "sh_exhausted", command: "sleep 60" },
|
||||
])
|
||||
|
||||
const drained: Session.ID[] = []
|
||||
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.sync(() => void drained.push(sessionID)),
|
||||
{ maxAttempts: 2 },
|
||||
restarted,
|
||||
)
|
||||
const restart = Context.get(context, SessionRestart.Service)
|
||||
const execution = Context.get(context, SessionExecution.Service)
|
||||
yield* restart.resumeSuspendedSessions
|
||||
yield* execution.awaitIdle(sessionID)
|
||||
|
||||
expect(drained).toEqual([])
|
||||
expect(yield* claims(database)).toEqual({ [sessionID]: false })
|
||||
expect(yield* attempts(database, sessionID)).toBe(0)
|
||||
expect(yield* SessionInbox.list(database.db, sessionID)).toMatchObject([
|
||||
{ payload: { metadata: { source: "shell", state: "cancelled" } } },
|
||||
])
|
||||
expect(yield* restarted.pendingBackground).toEqual([])
|
||||
yield* restart.resumeSuspendedSessions
|
||||
expect(drained).toEqual([])
|
||||
}),
|
||||
)
|
||||
|
||||
for (const shellFirst of [false, true]) {
|
||||
for (const resumeAttempts of [1, 2]) {
|
||||
it.effect(
|
||||
`recovers a child's shell ${shellFirst ? "before" : "after"} its job record without bypassing attempt ${resumeAttempts + 1}`,
|
||||
() =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const jobs = yield* Job.Service
|
||||
const bus = yield* Bus.Service
|
||||
const store = yield* SessionStore.Service
|
||||
const parent = Session.ID.make("ses_shell_child_parent")
|
||||
const child = Session.ID.make("ses_shell_child")
|
||||
yield* seedSessions(database, [parent])
|
||||
yield* seedSessions(database, [child], {
|
||||
parent_id: parent,
|
||||
time_suspended: Date.now(),
|
||||
resume_attempts: resumeAttempts,
|
||||
})
|
||||
const shell = seedBackground(jobs, child, [
|
||||
{ id: "call-child-shell", shellID: "sh_child", command: "sleep 60" },
|
||||
])
|
||||
if (shellFirst) yield* shell
|
||||
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)
|
||||
if (!shellFirst) yield* shell
|
||||
expect((yield* jobs.pendingBackground).map((job) => job.recovery.kind)).toEqual(
|
||||
shellFirst ? ["shell", "subagent"] : ["subagent", "shell"],
|
||||
)
|
||||
|
||||
const observed = yield* Deferred.make<{ attempts: number | undefined; notices: string[] }>()
|
||||
const release = yield* Deferred.make<void>()
|
||||
const parentWoken = yield* Deferred.make<void>()
|
||||
const drained: Session.ID[] = []
|
||||
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* SessionInbox.promote(database.db, bus, sessionID, "steer")
|
||||
if (sessionID === parent) {
|
||||
yield* Deferred.succeed(parentWoken, undefined)
|
||||
return
|
||||
}
|
||||
yield* Deferred.succeed(observed, {
|
||||
attempts: yield* attempts(database, sessionID),
|
||||
notices: (yield* store.context(sessionID))
|
||||
.filter((message) => message.type === "synthetic")
|
||||
.map((message) => message.text),
|
||||
})
|
||||
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(observed)).toEqual({
|
||||
attempts: 2,
|
||||
notices: [
|
||||
expect.stringContaining("The server restarted while you were working"),
|
||||
expect.stringContaining("Command cancelled because the server restarted"),
|
||||
],
|
||||
})
|
||||
expect(drained).toEqual([child])
|
||||
expect(yield* execution.active).not.toContain(parent)
|
||||
expect(yield* SessionInbox.list(database.db, parent)).toEqual([])
|
||||
}
|
||||
yield* Deferred.succeed(release, undefined)
|
||||
yield* Deferred.await(parentWoken)
|
||||
yield* Effect.forEach([parent, child], execution.awaitIdle, { discard: true })
|
||||
expect(drained).toEqual(resumeAttempts < 2 ? [child, parent] : [parent])
|
||||
expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toMatchObject([
|
||||
{ metadata: { source: "subagent", childID: child, state: resumeAttempts < 2 ? "completed" : "error" } },
|
||||
])
|
||||
expect(yield* claims(database)).toEqual({ [parent]: false, [child]: false })
|
||||
expect(yield* attempts(database, child)).toBe(0)
|
||||
expect(yield* restarted.pendingBackground).toEqual([])
|
||||
|
||||
yield* restart.resumeSuspendedSessions
|
||||
expect(drained).toHaveLength(resumeAttempts < 2 ? 2 : 1)
|
||||
}),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
it.effect("resumes a background subagent and notifies its parent exactly once", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
|
||||
Reference in New Issue
Block a user