fix(tasks): close background task review findings
This commit is contained in:
@@ -8,7 +8,7 @@ The background task registry (`ctx.tasks`): a runtime-global, CONCRETE service (
|
||||
- `get(id, caller?)` / `list(caller?)` — non-consuming snapshots; `list` returns only caller-owned plus unowned tasks (a global listing would leak foreign labels).
|
||||
- `read(id, caller?): TaskRead` — stream kinds consume the per-task cursor (v1's single intended reader is the owning model — a non-consuming multi-reader surface would be a cursor/snapshot API extension, not a `read` change); final kinds read the terminal output idempotently.
|
||||
- `kill(id, caller?, reason?)` — `'requested'` (live task: producer `cancel` runs first — a throw fails the kill loud and leaves the task untouched — then `stopping`) or `'already-terminal'`. Every successful kill marks the task `reported` (the killer saw the end → completion notice suppressed).
|
||||
- `wait(id, timeoutMs, caller?, signal?)` — resolves with the terminal snapshot (marked `reported`), or the live snapshot at timeout; an aborted signal rejects the WAIT only — unless the task already settled, in which case the wait still resolves and delivers the terminal snapshot (settlement suppressed the completion notice on this waiter's behalf, so rejecting would leave the finish both unreported and un-noticed). Timing is a [`dsh-timeout`](../../util/timeout/README.md) `deadline()` scoped to the `TASK_WAIT_TIMEOUT` code, so a nested foreign deadline never misreads as a wait timeout.
|
||||
- `wait(id, timeoutMs, caller?, signal?)` — resolves with the terminal snapshot (marked `reported`), or the live snapshot at timeout; an aborted signal rejects the WAIT only — unless the task already settled, in which case the wait still resolves and delivers the terminal snapshot (settlement suppressed the completion notice on this waiter's behalf, so rejecting would leave the finish both unreported and un-noticed). Timing is a [`dsh-timeout`](../../util/timeout/README.md) `deadline()` scoped to the `TASK_WAIT_TIMEOUT` code, so a nested foreign deadline never misreads as a wait timeout; timeout and abort detach their settlement resolver immediately, keeping retention bounded while the task remains live.
|
||||
- `onTaskDone(listener)` — exactly once per terminal task record; effect-scoped, per-listener containment, silent after service disposal.
|
||||
- `attachSurface(name)` — declares a control surface exists (the model tools, or a deployment's custom surface); effect-scoped.
|
||||
|
||||
|
||||
@@ -83,6 +83,8 @@ interface TrackedTask {
|
||||
markSettled: () => void
|
||||
/** Live {@link TaskService.wait} calls — a settlement with waiters marks the task reported. */
|
||||
waiters: number
|
||||
/** Removable resolvers for live waits; timeout/abort unregister before the task settles. */
|
||||
waitResolvers: Set<() => void>
|
||||
}
|
||||
|
||||
/** True for the three terminal {@link TaskStatus} values. */
|
||||
@@ -169,6 +171,7 @@ export class TaskService extends Service {
|
||||
settled,
|
||||
markSettled,
|
||||
waiters: 0,
|
||||
waitResolvers: new Set(),
|
||||
}
|
||||
this.store.set(id, task)
|
||||
|
||||
@@ -276,8 +279,10 @@ export class TaskService extends Service {
|
||||
* has already settled: settlement saw this live waiter and suppressed the
|
||||
* completion notice on its behalf, so the wait still resolves and delivers
|
||||
* the terminal snapshot it owes (an abort must never leave a finished task
|
||||
* both unreported and notice-suppressed). Throws for an unknown id, a task
|
||||
* owned by another session, or a non-positive timeout.
|
||||
* both unreported and notice-suppressed). Each live wait uses a removable
|
||||
* resolver that timeout/abort detaches, so a long-running task does not
|
||||
* retain expired waits. Throws for an unknown id, a task owned by another
|
||||
* session, or a non-positive timeout.
|
||||
* @param id - the task to wait for.
|
||||
* @param timeoutMs - max wait in milliseconds (positive, finite; the surface caps it).
|
||||
* @param caller - the waiting agent, checked against the task's owner.
|
||||
@@ -313,7 +318,13 @@ export class TaskService extends Service {
|
||||
// the wait. `using` clears the timer on every exit path.
|
||||
using d = deadline(signal, timeoutMs, TASK_WAIT_TIMEOUT)
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
const onSettled = (): void => {
|
||||
task.waitResolvers.delete(onSettled)
|
||||
d.signal.removeEventListener('abort', onAbort)
|
||||
resolve()
|
||||
}
|
||||
const onAbort = (): void => {
|
||||
task.waitResolvers.delete(onSettled)
|
||||
if (timeoutOf(d.signal, TASK_WAIT_TIMEOUT) !== undefined) {
|
||||
resolve()
|
||||
} else if (isTerminal(task.status)) {
|
||||
@@ -325,11 +336,8 @@ export class TaskService extends Service {
|
||||
reject(new Error('wait aborted'))
|
||||
}
|
||||
}
|
||||
task.waitResolvers.add(onSettled)
|
||||
d.signal.addEventListener('abort', onAbort, { once: true })
|
||||
void task.settled.then(() => {
|
||||
d.signal.removeEventListener('abort', onAbort)
|
||||
resolve()
|
||||
})
|
||||
})
|
||||
} finally {
|
||||
uncount()
|
||||
@@ -437,6 +445,9 @@ export class TaskService extends Service {
|
||||
}
|
||||
}
|
||||
}
|
||||
const waitResolvers = [...task.waitResolvers]
|
||||
task.waitResolvers.clear()
|
||||
for (const resolveWait of waitResolvers) resolveWait()
|
||||
task.markSettled()
|
||||
}
|
||||
|
||||
|
||||
@@ -71,7 +71,7 @@ export interface TaskStart {
|
||||
label: string
|
||||
/**
|
||||
* The spawning agent. Its `session.header.id` becomes the task's owner
|
||||
* token (read/kill/wait/list are fenced to that session), and its `ctx` scope
|
||||
* identity (read/kill/wait/list are fenced to that session), and its `ctx` scope
|
||||
* owns an async cleanup that cancels and awaits the task during disposal. It
|
||||
* must be the exact live instance currently registered under its agent id;
|
||||
* a stale object whose id has been reused is rejected before work starts.
|
||||
|
||||
@@ -59,6 +59,14 @@ async function harness() {
|
||||
/** Let the settlement continuation (a `done.then`) run. */
|
||||
const tick = () => new Promise<void>(r => setTimeout(r, 0))
|
||||
|
||||
/** Inspect the internal resolver registry to pin bounded retention while a task stays live. */
|
||||
function waitResolverCount(ctx: Context, id: TaskId): number {
|
||||
const service = ctx.tasks as unknown as { store: Map<TaskId, { waitResolvers: Set<() => void> }> }
|
||||
const task = service.store.get(id)
|
||||
if (task === undefined) throw new Error(`missing test task ${id}`)
|
||||
return task.waitResolvers.size
|
||||
}
|
||||
|
||||
describe('TaskService.start', () => {
|
||||
it('preserves the SessionId brand on public owner snapshots', () => {
|
||||
expectTypeOf<TaskSnapshot['ownerSession']>().toEqualTypeOf<SessionId | undefined>()
|
||||
@@ -254,6 +262,26 @@ describe('TaskService.wait', () => {
|
||||
expect(await ctx.tasks.wait(id, 5)).toMatchObject({ status: 'running', reported: false })
|
||||
})
|
||||
|
||||
it('unregisters timed-out and aborted wait resolvers while the task remains live', async () => {
|
||||
const ctx = await harness()
|
||||
const id = ctx.tasks.start(producer().spec)
|
||||
|
||||
for (let index = 0; index < 3; index += 1) {
|
||||
const wait = ctx.tasks.wait(id, 5)
|
||||
expect(waitResolverCount(ctx, id)).toBe(1)
|
||||
await expect(wait).resolves.toMatchObject({ status: 'running' })
|
||||
expect(waitResolverCount(ctx, id)).toBe(0)
|
||||
}
|
||||
|
||||
const controller = new AbortController()
|
||||
const wait = ctx.tasks.wait(id, 5_000, undefined, controller.signal)
|
||||
expect(waitResolverCount(ctx, id)).toBe(1)
|
||||
controller.abort()
|
||||
await expect(wait).rejects.toThrow('wait aborted')
|
||||
expect(waitResolverCount(ctx, id)).toBe(0)
|
||||
expect(ctx.tasks.get(id).status).toBe('running')
|
||||
})
|
||||
|
||||
it('returns immediately for an already-terminal task', async () => {
|
||||
const ctx = await harness()
|
||||
const p = producer()
|
||||
|
||||
Reference in New Issue
Block a user