feat(tasks): wake an idle owner when a background task completes
Completion notices went through agent.inject(), which never reserves a driver, so a task settling after its turn closed left the notice parked until unrelated input woke the agent — while the same prompt told the model not to poll for it. An unreported completion now picks its lane from the owner's state: a busy owner is injected as before, an idle owner is woken with followup(). This adopts the delivery rule the subagent continuation manager already ships. maxConsecutiveWakes bounds the self-exciting chain and is reset by user-authored input; completionDelivery: quiet restores the old lane for deterministic transcripts. Teardown cancellation now claims the terminal report the way kill() already does, so an owner being destroyed is never woken, and settle() announces completion last so a reporter that opens a turn synchronously sees a committed record.
This commit is contained in:
@@ -224,11 +224,12 @@ export class LocalTaskService extends TaskService {
|
||||
}
|
||||
const onAbort = (): void => {
|
||||
task.waitResolvers.delete(onSettled)
|
||||
// A settled task cannot reach here: settlement releases every waiter
|
||||
// before it announces completion, and each released waiter detaches
|
||||
// this listener in the same synchronous span, so nothing that reacts
|
||||
// to a settlement can abort a wait the settlement already owed.
|
||||
if (timeoutOf(d.signal, TASK_WAIT_TIMEOUT) !== undefined) {
|
||||
resolve()
|
||||
} else if (isTerminal(task.status)) {
|
||||
// Settlement suppressed the notice for this waiter; deliver it.
|
||||
resolve()
|
||||
} else {
|
||||
uncount()
|
||||
reject(new Error('wait aborted'))
|
||||
@@ -364,9 +365,12 @@ export class LocalTaskService extends TaskService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Record the first terminal outcome, notify contained listeners, and release
|
||||
* waiters. First-wins preserves a teardown force-failure against late producer
|
||||
* settlement. Pending waits mark the task reported before listeners run.
|
||||
* Record the first terminal outcome, release waiters, then announce
|
||||
* completion. First-wins preserves a teardown force-failure against late
|
||||
* producer settlement. Pending waits mark the task reported before listeners
|
||||
* run. Completion is announced last because a reporter may open a model turn
|
||||
* synchronously: every other observer of this settlement must already have
|
||||
* seen the committed record.
|
||||
*/
|
||||
private settle(task: TrackedTask, outcome: TaskOutcome): void {
|
||||
if (isTerminal(task.status)) return
|
||||
@@ -375,24 +379,23 @@ export class LocalTaskService extends TaskService {
|
||||
task.output = outcome.output
|
||||
task.finishedAt = Date.now()
|
||||
if (task.waiters > 0) task.reported = true
|
||||
if (!this.listenersClosed) {
|
||||
const snapshot = this.snapshot(task)
|
||||
for (const listener of this.listenersFor(task.owner)) {
|
||||
try {
|
||||
const returned = listener(snapshot, task.owner)
|
||||
void Promise.resolve(returned).catch((error: unknown) => {
|
||||
this.selfCtx.logger.warn(`tasks: onTaskDone listener rejected for ${task.id}: ${String(error)}`)
|
||||
})
|
||||
} catch (error: unknown) {
|
||||
this.selfCtx.logger.warn(`tasks: onTaskDone listener threw for ${task.id}: ${String(error)}`)
|
||||
}
|
||||
}
|
||||
}
|
||||
const snapshot = this.snapshot(task)
|
||||
const waitResolvers = [...task.waitResolvers]
|
||||
task.waitResolvers.clear()
|
||||
for (const resolveWait of waitResolvers) resolveWait()
|
||||
task.markSettled()
|
||||
this.notifyChanged(task.owner)
|
||||
if (this.listenersClosed) return
|
||||
for (const listener of this.listenersFor(task.owner)) {
|
||||
try {
|
||||
const returned = listener(snapshot, task.owner)
|
||||
void Promise.resolve(returned).catch((error: unknown) => {
|
||||
this.selfCtx.logger.warn(`tasks: onTaskDone listener rejected for ${task.id}: ${String(error)}`)
|
||||
})
|
||||
} catch (error: unknown) {
|
||||
this.selfCtx.logger.warn(`tasks: onTaskDone listener threw for ${task.id}: ${String(error)}`)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -466,6 +469,11 @@ export class LocalTaskService extends TaskService {
|
||||
try {
|
||||
task.cancel(reason)
|
||||
task.status = 'stopping'
|
||||
// Teardown cancellation is a kill without a caller, so it claims the
|
||||
// terminal report the same way `kill()` does. Nothing will read a
|
||||
// notice for a task whose owner or service is being destroyed, and a
|
||||
// waking reporter would spend a model request per teardown layer.
|
||||
task.reported = true
|
||||
// Teardown reaches settlement only after the producer releases, which a
|
||||
// slow stop can defer; announcing the transition here is what keeps an
|
||||
// observer from showing `running` for that whole window.
|
||||
|
||||
Reference in New Issue
Block a user