fix(background-agent): defer busy parent wake
This commit is contained in:
@@ -88,9 +88,24 @@ import {
|
|||||||
resolveSubagentSpawnContext,
|
resolveSubagentSpawnContext,
|
||||||
type SubagentSpawnContext,
|
type SubagentSpawnContext,
|
||||||
} from "./subagent-spawn-limits"
|
} from "./subagent-spawn-limits"
|
||||||
|
import { settleAfterSessionIdle } from "../../hooks/shared/session-idle-settle"
|
||||||
|
|
||||||
type OpencodeClient = PluginInput["client"]
|
type OpencodeClient = PluginInput["client"]
|
||||||
|
|
||||||
|
type ParentWakePromptContext = {
|
||||||
|
agent?: string
|
||||||
|
model?: { providerID: string; modelID: string }
|
||||||
|
variant?: string
|
||||||
|
tools?: Record<string, boolean>
|
||||||
|
}
|
||||||
|
|
||||||
|
type SessionStatusInfo = { type?: string }
|
||||||
|
|
||||||
|
const BACKGROUND_PARENT_WAKE_PROMPT = `<system-reminder>
|
||||||
|
[BACKGROUND TASK NOTIFICATION READY]
|
||||||
|
A background task notification was already added to this session. Continue from that notification.
|
||||||
|
</system-reminder>`
|
||||||
|
|
||||||
interface MessagePartInfo {
|
interface MessagePartInfo {
|
||||||
id?: string
|
id?: string
|
||||||
sessionID?: string
|
sessionID?: string
|
||||||
@@ -210,6 +225,7 @@ export class BackgroundManager {
|
|||||||
private completedTaskSummaries: Map<string, BackgroundTaskNotificationTask[]> = new Map()
|
private completedTaskSummaries: Map<string, BackgroundTaskNotificationTask[]> = new Map()
|
||||||
private idleDeferralTimers: Map<string, ReturnType<typeof setTimeout>> = new Map()
|
private idleDeferralTimers: Map<string, ReturnType<typeof setTimeout>> = new Map()
|
||||||
private notificationQueueByParent: Map<string, Promise<void>> = new Map()
|
private notificationQueueByParent: Map<string, Promise<void>> = new Map()
|
||||||
|
private pendingParentWakes: Map<string, ParentWakePromptContext> = new Map()
|
||||||
private observedOutputSessions: Set<string> = new Set()
|
private observedOutputSessions: Set<string> = new Set()
|
||||||
private observedIncompleteTodosBySession: Map<string, boolean> = new Map()
|
private observedIncompleteTodosBySession: Map<string, boolean> = new Map()
|
||||||
private rootDescendantCounts: Map<string, number>
|
private rootDescendantCounts: Map<string, number>
|
||||||
@@ -1360,6 +1376,12 @@ The fallback retry session is now created and can be inspected directly.
|
|||||||
|
|
||||||
if (event.type === "session.idle") {
|
if (event.type === "session.idle") {
|
||||||
if (!props || typeof props !== "object") return
|
if (!props || typeof props !== "object") return
|
||||||
|
const sessionID = typeof props.sessionID === "string" ? props.sessionID : undefined
|
||||||
|
if (sessionID) {
|
||||||
|
void this.enqueueNotificationForParent(sessionID, () => this.flushPendingParentWake(sessionID)).catch((error) => {
|
||||||
|
log("[background-agent] Failed to flush pending parent wake:", { sessionID, error })
|
||||||
|
})
|
||||||
|
}
|
||||||
handleSessionIdleBackgroundEvent({
|
handleSessionIdleBackgroundEvent({
|
||||||
properties: props as Record<string, unknown>,
|
properties: props as Record<string, unknown>,
|
||||||
findBySession: (id) => {
|
findBySession: (id) => {
|
||||||
@@ -2152,24 +2174,32 @@ The task was re-queued on a fallback model after a retryable failure.
|
|||||||
const shouldReply = allComplete || isTaskFailure
|
const shouldReply = allComplete || isTaskFailure
|
||||||
|
|
||||||
const variant = promptContext?.model?.variant
|
const variant = promptContext?.model?.variant
|
||||||
|
const parentPromptContext: ParentWakePromptContext = {
|
||||||
|
...(agent !== undefined ? { agent } : {}),
|
||||||
|
...(model !== undefined ? { model } : {}),
|
||||||
|
...(variant !== undefined ? { variant } : {}),
|
||||||
|
...(resolvedTools ? { tools: resolvedTools } : {}),
|
||||||
|
}
|
||||||
|
const shouldDeferReply = shouldReply && await this.isSessionActive(task.parentSessionId)
|
||||||
|
|
||||||
try {
|
try {
|
||||||
await this.client.session.promptAsync({
|
await this.client.session.promptAsync({
|
||||||
path: { id: task.parentSessionId },
|
path: { id: task.parentSessionId },
|
||||||
body: {
|
body: {
|
||||||
noReply: !shouldReply,
|
noReply: shouldDeferReply || !shouldReply,
|
||||||
...(agent !== undefined ? { agent } : {}),
|
...parentPromptContext,
|
||||||
...(model !== undefined ? { model } : {}),
|
|
||||||
...(variant !== undefined ? { variant } : {}),
|
|
||||||
...(resolvedTools ? { tools: resolvedTools } : {}),
|
|
||||||
parts: [createInternalAgentTextPart(notification)],
|
parts: [createInternalAgentTextPart(notification)],
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
if (shouldDeferReply) {
|
||||||
|
this.pendingParentWakes.set(task.parentSessionId, parentPromptContext)
|
||||||
|
}
|
||||||
log("[background-agent] Sent notification to parent session:", {
|
log("[background-agent] Sent notification to parent session:", {
|
||||||
taskId: task.id,
|
taskId: task.id,
|
||||||
allComplete,
|
allComplete,
|
||||||
isTaskFailure,
|
isTaskFailure,
|
||||||
noReply: !shouldReply,
|
noReply: shouldDeferReply || !shouldReply,
|
||||||
|
deferredReply: shouldDeferReply,
|
||||||
})
|
})
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
if (isAbortedSessionError(error)) {
|
if (isAbortedSessionError(error)) {
|
||||||
@@ -2201,6 +2231,60 @@ The task was re-queued on a fallback model after a retryable failure.
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private async isSessionActive(sessionID: string): Promise<boolean> {
|
||||||
|
const sessionStatusMethod = this.client?.session?.status
|
||||||
|
if (typeof sessionStatusMethod !== "function") {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
|
const statusResult = await this.client.session.status()
|
||||||
|
const statuses = normalizeSDKResponse(
|
||||||
|
statusResult,
|
||||||
|
{} as Record<string, SessionStatusInfo>,
|
||||||
|
)
|
||||||
|
const status = statuses[sessionID]
|
||||||
|
return typeof status?.type === "string" && isActiveSessionStatus(status.type)
|
||||||
|
} catch (error) {
|
||||||
|
log("[background-agent] Unable to check parent session status before wake:", {
|
||||||
|
sessionID,
|
||||||
|
error,
|
||||||
|
})
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private async flushPendingParentWake(sessionID: string): Promise<void> {
|
||||||
|
const wakeContext = this.pendingParentWakes.get(sessionID)
|
||||||
|
if (!wakeContext) return
|
||||||
|
|
||||||
|
if (await this.isSessionActive(sessionID)) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
this.pendingParentWakes.delete(sessionID)
|
||||||
|
await settleAfterSessionIdle()
|
||||||
|
|
||||||
|
if (await this.isSessionActive(sessionID)) {
|
||||||
|
this.pendingParentWakes.set(sessionID, wakeContext)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
|
await this.client.session.promptAsync({
|
||||||
|
path: { id: sessionID },
|
||||||
|
body: {
|
||||||
|
noReply: false,
|
||||||
|
...wakeContext,
|
||||||
|
parts: [createInternalAgentTextPart(BACKGROUND_PARENT_WAKE_PROMPT)],
|
||||||
|
},
|
||||||
|
})
|
||||||
|
log("[background-agent] Sent deferred parent wake:", { sessionID })
|
||||||
|
} catch (error) {
|
||||||
|
log("[background-agent] Failed to send deferred parent wake:", { sessionID, error })
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private pruneStaleTasksAndNotifications(): void {
|
private pruneStaleTasksAndNotifications(): void {
|
||||||
pruneStaleTasksAndNotifications({
|
pruneStaleTasksAndNotifications({
|
||||||
tasks: this.tasks,
|
tasks: this.tasks,
|
||||||
@@ -2525,6 +2609,7 @@ The task was re-queued on a fallback model after a retryable failure.
|
|||||||
this.pendingNotifications.clear()
|
this.pendingNotifications.clear()
|
||||||
this.pendingByParent.clear()
|
this.pendingByParent.clear()
|
||||||
this.notificationQueueByParent.clear()
|
this.notificationQueueByParent.clear()
|
||||||
|
this.pendingParentWakes.clear()
|
||||||
this.rootDescendantCounts.clear()
|
this.rootDescendantCounts.clear()
|
||||||
this.queuesByKey.clear()
|
this.queuesByKey.clear()
|
||||||
this.processingKeys.clear()
|
this.processingKeys.clear()
|
||||||
|
|||||||
@@ -50,11 +50,19 @@ function createTask(overrides: Partial<BackgroundTask> & { id: string; parentSes
|
|||||||
function createManager(enableParentSessionNotifications: boolean): {
|
function createManager(enableParentSessionNotifications: boolean): {
|
||||||
manager: BackgroundManager
|
manager: BackgroundManager
|
||||||
promptAsyncCalls: PromptAsyncCall[]
|
promptAsyncCalls: PromptAsyncCall[]
|
||||||
|
}
|
||||||
|
function createManager(
|
||||||
|
enableParentSessionNotifications: boolean,
|
||||||
|
sessionStatuses?: Record<string, { type: string }>,
|
||||||
|
): {
|
||||||
|
manager: BackgroundManager
|
||||||
|
promptAsyncCalls: PromptAsyncCall[]
|
||||||
} {
|
} {
|
||||||
const promptAsyncCalls: PromptAsyncCall[] = []
|
const promptAsyncCalls: PromptAsyncCall[] = []
|
||||||
const client = {
|
const client = {
|
||||||
session: {
|
session: {
|
||||||
messages: async () => [],
|
messages: async () => [],
|
||||||
|
status: async () => ({ data: sessionStatuses ?? {} }),
|
||||||
prompt: async () => ({}),
|
prompt: async () => ({}),
|
||||||
promptAsync: async (call: PromptAsyncCall) => {
|
promptAsync: async (call: PromptAsyncCall) => {
|
||||||
promptAsyncCalls.push(call)
|
promptAsyncCalls.push(call)
|
||||||
@@ -143,6 +151,10 @@ async function notifyParentSessionForTest(manager: BackgroundManager, task: Back
|
|||||||
return notifyParentSession.call(manager, task)
|
return notifyParentSession.call(manager, task)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function waitForDeferredWake(): Promise<void> {
|
||||||
|
return new Promise((resolve) => setTimeout(resolve, 180))
|
||||||
|
}
|
||||||
|
|
||||||
function getRequiredTimer(manager: BackgroundManager, taskID: string): ReturnType<typeof setTimeout> {
|
function getRequiredTimer(manager: BackgroundManager, taskID: string): ReturnType<typeof setTimeout> {
|
||||||
const timer = getCompletionTimers(manager).get(taskID)
|
const timer = getCompletionTimers(manager).get(taskID)
|
||||||
expect(timer).toBeDefined()
|
expect(timer).toBeDefined()
|
||||||
@@ -232,6 +244,52 @@ describe("BackgroundManager.notifyParentSession cleanup scheduling", () => {
|
|||||||
expect(allCompletePayload).toContain(taskA.description)
|
expect(allCompletePayload).toContain(taskA.description)
|
||||||
expect(allCompletePayload).toContain(taskB.description)
|
expect(allCompletePayload).toContain(taskB.description)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
test("#when parent session is busy #then all-complete notification does not start an overlapping parent reply", async () => {
|
||||||
|
// given
|
||||||
|
const sessionStatuses: Record<string, { type: string }> = {
|
||||||
|
"parent-1": { type: "busy" },
|
||||||
|
}
|
||||||
|
const { manager, promptAsyncCalls } = createManager(true, sessionStatuses)
|
||||||
|
managerUnderTest = manager
|
||||||
|
const task = createTask({ id: "task-a", parentSessionId: "parent-1", description: "task A", status: "completed", completedAt: new Date("2026-03-11T00:01:00.000Z") })
|
||||||
|
getTasks(manager).set(task.id, task)
|
||||||
|
getPendingByParent(manager).set(task.parentSessionId, new Set([task.id]))
|
||||||
|
|
||||||
|
// when
|
||||||
|
await notifyParentSessionForTest(manager, task)
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(promptAsyncCalls).toHaveLength(1)
|
||||||
|
expect(promptAsyncCalls[0]?.body.noReply).toBe(true)
|
||||||
|
expect(JSON.stringify(promptAsyncCalls[0]?.body.parts)).toContain("ALL BACKGROUND TASKS COMPLETE")
|
||||||
|
})
|
||||||
|
|
||||||
|
test("#when deferred parent session becomes idle #then wake prompt is sent once without duplicating the notification", async () => {
|
||||||
|
// given
|
||||||
|
const sessionStatuses: Record<string, { type: string }> = {
|
||||||
|
"parent-1": { type: "busy" },
|
||||||
|
}
|
||||||
|
const { manager, promptAsyncCalls } = createManager(true, sessionStatuses)
|
||||||
|
managerUnderTest = manager
|
||||||
|
const task = createTask({ id: "task-a", parentSessionId: "parent-1", description: "task A", status: "completed", completedAt: new Date("2026-03-11T00:01:00.000Z") })
|
||||||
|
getTasks(manager).set(task.id, task)
|
||||||
|
getPendingByParent(manager).set(task.parentSessionId, new Set([task.id]))
|
||||||
|
await notifyParentSessionForTest(manager, task)
|
||||||
|
|
||||||
|
// when
|
||||||
|
sessionStatuses["parent-1"] = { type: "idle" }
|
||||||
|
manager.handleEvent({ type: "session.idle", properties: { sessionID: "parent-1" } })
|
||||||
|
await waitForDeferredWake()
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(promptAsyncCalls).toHaveLength(2)
|
||||||
|
expect(promptAsyncCalls[0]?.body.noReply).toBe(true)
|
||||||
|
expect(promptAsyncCalls[1]?.body.noReply).toBe(false)
|
||||||
|
const wakePayload = JSON.stringify(promptAsyncCalls[1]?.body.parts)
|
||||||
|
expect(wakePayload).toContain("BACKGROUND TASK NOTIFICATION READY")
|
||||||
|
expect(wakePayload).not.toContain("ALL BACKGROUND TASKS COMPLETE")
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
describe("#given a completed task with cleanup timer scheduled", () => {
|
describe("#given a completed task with cleanup timer scheduled", () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user