/** * The Scheduler: arming, ticking, backoff, attention polling, user interjection and * the headless drain (§11). * * The model is **idle with an armed wake** (§11.1), not a blocking wait tool. The * orchestrator ends its turn, pi goes idle, no tokens are consumed, and a wake later * delivers a `custom_message` with `triggerTurn: true`. A blocking wait would keep the * orchestrator inside one enormous turn whose context grows with total delegated work, * which inverts the state layer's whole purpose (D6/D10). * * Ordering constraints that are not obvious from the rules: * * - Wakes are armed on `agent_settled`, never `agent_end` (R-SLEEP-2 / F2). * - A wake is **never** sent synchronously from inside an `agent_settled` handler. * `agent-session.ts:_emitAgentSettled` clears `_isAgentRunActive` *before* awaiting * extension handlers, so `sendMessage({triggerTurn:true})` from inside one is seen * as idle and re-enters `_runAgentPrompt` from within the previous run's `finally`. * Every send is therefore deferred through a timer. * - Delivery is `deliverAs: "followUp"`, never `"steer"` (R-SLEEP-4): a completion * should land after the orchestrator finishes its current reasoning, not inside it. */ import type { ExtensionContext } from "@earendil-works/pi-coding-agent"; import type { WorkerConfig } from "../worker/config.ts"; import { type RegistryEntry, elapsed, loadAllRuns } from "../worker/registry.ts"; import { type RunStatus, isCompletionOutcomeState, isEndedState, isTerminalState, patchStatus, readEvents, readFileIfExists, } from "../worker/status.ts"; import { humanizeActivity } from "../worker/widget.ts"; import { type AttentionFinding, type AttentionTrigger, attentionKey, evaluateAttention, isAttentionTrigger, renderAttentionPayload, selectAttention, } from "./attention.ts"; import { type BatcherTimers, type CompletionItem, NotificationBatcher, completionItem, renderCompletionPayload, renderSleepPayload, renderTickPayload, } from "./notify.ts"; import { type WakeRecord, armWake, clearWakes, loadWakes, removeWake, settleWake } from "./wake.ts"; import { isReviewEligibleState, renderWorkerReview } from "../worker/trajectory.ts"; export const TICK_SUPPRESS_MS = 30_000; /** How often the scheduler re-evaluates attention and refreshes the indicator. */ export const SCHEDULER_POLL_MS = 500; export const WORKER_WAIT_MS = 300_000; export const WAKE_CUSTOM_TYPE = "agi-wake"; export type WakeReason = "worker_complete" | "worker_attention" | "worker_check" | "tick" | "sleep" | "user" | "startup"; /** R-UI-19: wake payloads are `custom_message` because they must reach the model. */ export interface WakePayload { reason: WakeReason; text: string; runIds?: string[]; } export interface SchedulerHost { /** Live ctx, or undefined when AGI mode is off / the instance was replaced (F4). */ ctx(): ExtensionContext | undefined; config(ctx: ExtensionContext): WorkerConfig; /** Delivers the wake. Returns false when the send could not be attempted. */ send(payload: WakePayload): boolean; /** R-SLEEP-19/20: the reason surfaced in the next context block. */ onWakeReason(reason: WakeReason, detail: string): void; /** Emitted for the fleet widget and the sleep indicator. */ onRefresh?: () => void; /** Whether concrete steer/interrupt actions may be offered in attention text. */ controlAvailable?: () => boolean; } export interface SchedulerTimers extends BatcherTimers { setInterval(fn: () => void, ms: number): unknown; clearInterval(handle: unknown): void; } export const realSchedulerTimers: SchedulerTimers = { setTimeout: (fn, ms) => { const timer = setTimeout(fn, ms); timer.unref?.(); return timer; }, clearTimeout: (handle) => clearTimeout(handle as ReturnType), setInterval: (fn, ms) => { const timer = setInterval(fn, ms); timer.unref?.(); return timer; }, clearInterval: (handle) => clearInterval(handle as ReturnType), now: () => Date.now(), }; export type SleepRequest = | { source: "timer"; durationMs: number; note?: string } | { source: "worker" | "tick" | "user"; note?: string }; export type SchedulerExit = "blocked"; export interface SchedulerSnapshot { armed: boolean; tickIntervalMs: number; quietTicks: number; sleepStreak: number; sleepRemainingMs: number | undefined; sleepNote: string | undefined; tickInMs: number | undefined; blocked: boolean; exit: SchedulerExit | undefined; } interface TimedWake { token: string; kind: "tick" | "sleep" | "worker_check"; firesAt: number; handle: unknown; note?: string; streak?: number; durationMs?: number; } interface DispatchClaim { payload?: WakePayload; failed?: boolean; } interface PendingDispatch { handle: unknown; epoch: number; sessionId: string; reason: WakeReason; tokens: string[]; /** Mutable scheduler-owned items until an idle worker-completion send succeeds. */ completionItems?: CompletionItem[]; } export class Scheduler { private host: SchedulerHost; private timers: SchedulerTimers; private batcher: NotificationBatcher; private cwd: string | undefined; private sessionId = ""; private enabled = false; private idle = false; private blockedThisTurn = false; private pendingSleep: SleepRequest | undefined; private timed: TimedWake | undefined; /** R-SLEEP-5: durable completion records, one per run, even during nudge hold. */ private workerWakeTokens = new Map(); /** Runs whose completion wake was consumed, headless-delivered, or durably claimed. */ private handledWorkerWakes = new Set(); /** Evidence signature last durably claimed for each handled run in this session. */ private deliveredWorkerEvidence = new Map(); private pollHandle: unknown; /** R-SLEEP-12: suppression begins only after a worker wake was actually sent. */ private lastWorkerWakeAt = 0; private tickIntervalMs = 600_000; private quietTicks = 0; /** R-TOOL-21d: internal continuity for the same named external wait. */ private sleepStreak = 0; private sleepNote: string | undefined; private actedThisTurn = false; private sleptThisTurn = false; /** R-SLEEP-19/Issue 5: unread completions remain scheduler-owned until idle. */ private deferred = new Map(); private draining = false; private dispatchEpoch = 0; private dispatches = new Set(); private exitReason: SchedulerExit | undefined; constructor(host: SchedulerHost, timers: SchedulerTimers = realSchedulerTimers) { this.host = host; this.timers = timers; this.batcher = new NotificationBatcher({ deliver: (items) => this.deliverCompletions(items), timers, }); } snapshot(): SchedulerSnapshot { const now = this.timers.now(); return { armed: this.timed !== undefined, tickIntervalMs: this.tickIntervalMs, quietTicks: this.quietTicks, sleepStreak: this.sleepStreak, sleepRemainingMs: this.timed?.kind === "sleep" || this.timed?.kind === "worker_check" ? Math.max(0, this.timed.firesAt - now) : undefined, sleepNote: this.sleepNote, tickInMs: this.timed?.kind === "tick" ? Math.max(0, this.timed.firesAt - now) : undefined, blocked: this.blockedThisTurn, exit: this.exitReason, }; } /** Toggle ON, and on `session_start` for a restored session. */ activate(ctx: ExtensionContext): void { this.dispatchEpoch += 1; this.idle = false; this.cwd = ctx.cwd; const sessionId = ctx.sessionManager.getSessionId(); if (this.sessionId !== sessionId) { this.handledWorkerWakes.clear(); this.deliveredWorkerEvidence.clear(); } this.sessionId = sessionId; this.enabled = true; this.exitReason = undefined; this.tickIntervalMs = this.host.config(ctx).tickMs; this.quietTicks = 0; this.rearmFromDisk(); this.rearmUnconsumedCompletions(); if (this.pollHandle === undefined) { this.pollHandle = this.timers.setInterval(() => this.poll(), SCHEDULER_POLL_MS); } } /** * A worker may settle while Pi is down, before any scheduler wake record can be * created. Project those unread terminal results back through the ordinary durable * completion path instead of duplicating them in a broad startup recovery block. */ private rearmUnconsumedCompletions(): void { if (this.cwd === undefined) return; for (const entry of loadAllRuns(this.cwd).entries) { if (!isCompletionOutcomeState(entry.status.state)) continue; if (entry.status.owner !== undefined && entry.status.owner.sessionId !== this.sessionId) continue; if (entry.status.resultConsumed || entry.status.headlessDeliveredAt !== undefined) continue; if (this.handledWorkerWakes.has(entry.runId) || this.workerWakeTokens.has(entry.runId)) continue; const record = armWake(this.cwd, { kind: "worker", sessionId: this.sessionId, runId: entry.runId, now: this.timers.now(), }); if (record === undefined) continue; this.workerWakeTokens.set(entry.runId, record.token); const item = completionItem(entry.runId, entry.status, this.readReport(entry), this.timers.now()); this.queueCompletionItem({ ...item, wakeToken: record.token }); } } /** Shared run storage may contain another orchestrator's work. */ private belongsToSession(status: RunStatus): boolean { return status.owner === undefined || status.owner.sessionId === this.sessionId; } private reviewEligibleEntries(cwd: string): RegistryEntry[] { return loadAllRuns(cwd).entries.filter( (entry) => this.belongsToSession(entry.status) && isReviewEligibleState(entry.status.state), ); } deactivate(): void { this.enabled = false; this.idle = false; this.dispatchEpoch += 1; this.disarmAll(true); if (this.pollHandle !== undefined) { this.timers.clearInterval(this.pollHandle); this.pollHandle = undefined; } // E64g: a wake that outlives its purpose would wake an orchestrator with nothing // to do, so the durable records go too. Worker records are dropped as well: this // session is no longer supervising, and the next activation re-derives them from // the run directory. } /** * R-SLEEP-5 / E64f. Records armed before a restart are honoured, and one whose * `firesAt` is already past fires immediately rather than being silently extended. */ private rearmFromDisk(): void { if (this.cwd === undefined) return; this.sleepStreak = 0; this.sleepNote = undefined; const { own } = loadWakes(this.cwd, this.sessionId); let best: WakeRecord | undefined; for (const record of own) { if (record.kind === "worker_check" && this.reviewEligibleEntries(this.cwd).length === 0) { removeWake(this.cwd, record.token); continue; } if (record.kind === "attention") { if ( record.runId === undefined || record.attentionTrigger === undefined || record.attentionDetail === undefined || !isAttentionTrigger(record.attentionTrigger) ) { removeWake(this.cwd, record.token); continue; } const entry = loadAllRuns(this.cwd).entries.find((candidate) => candidate.runId === record.runId); if (entry === undefined || !this.belongsToSession(entry.status) || isEndedState(entry.status.state)) { removeWake(this.cwd, record.token); continue; } const finding: AttentionFinding = { trigger: record.attentionTrigger, detail: record.attentionDetail, chatWorthy: true, }; const payload: WakePayload = { reason: "worker_attention", text: renderAttentionPayload(entry.status, finding, { controlAvailable: this.host.controlAvailable?.() === true }), runIds: [entry.runId], }; this.dispatch(payload, { tokens: [record.token], claim: () => { const settled = settleWake(this.cwd ?? "", record.token); return settled === undefined ? { failed: true } : { payload }; }, }); continue; } if (record.kind === "worker") { // R-SLEEP-5: a completion armed before a crash is still a wake, not merely // recovery metadata. Re-submit it through the normal hold/batch path so it // reaches the model once; `resultConsumed` can still cancel it. if (record.runId === undefined) { removeWake(this.cwd, record.token); continue; } const entry = loadAllRuns(this.cwd).entries.find((candidate) => candidate.runId === record.runId); if (entry === undefined || !this.belongsToSession(entry.status) || !isTerminalState(entry.status.state)) { removeWake(this.cwd, record.token); continue; } if (entry.status.resultConsumed || entry.status.headlessDeliveredAt !== undefined) { this.handledWorkerWakes.add(record.runId); removeWake(this.cwd, record.token); continue; } const existingToken = this.workerWakeTokens.get(record.runId); if (existingToken !== undefined) { // Older versions could persist duplicates for one run. Keep one durable // ownership token and discard the redundant record during rearm. if (existingToken !== record.token) removeWake(this.cwd, record.token); continue; } const token = record.token; this.workerWakeTokens.set(record.runId, token); const item = completionItem(entry.runId, entry.status, this.readReport(entry), this.timers.now()); this.queueCompletionItem({ ...item, wakeToken: token }); continue; } if (record.firesAt === undefined) { removeWake(this.cwd, record.token); continue; } // Only one timed wake can be pending at a time; keep the soonest and drop the // rest, which is what an interrupted arm/disarm sequence leaves behind. if (best === undefined || Date.parse(record.firesAt) < Date.parse(best.firesAt ?? "")) { if (best !== undefined) removeWake(this.cwd, best.token); best = record; } else { removeWake(this.cwd, record.token); } } if (best === undefined) return; this.sleepStreak = best.streak ?? 0; this.sleepNote = best.note; this.scheduleExisting(best); } private scheduleExisting(record: WakeRecord): void { const firesAt = Date.parse(record.firesAt ?? ""); if (Number.isNaN(firesAt)) return; const kind = record.kind === "sleep" || record.kind === "worker_check" ? record.kind : "tick"; this.clearTimed(); const delay = Math.max(0, firesAt - this.timers.now()); const timed: TimedWake = { token: record.token, kind, firesAt, handle: this.timers.setTimeout(() => this.fireTimed(), delay), ...(record.note === undefined ? {} : { note: record.note }), ...(record.streak === undefined ? {} : { streak: record.streak }), ...(record.durationMs === undefined ? {} : { durationMs: record.durationMs }), }; this.timed = timed; } /** R-UI-10: the working row may only be taken over between settle and next start. */ onAgentStart(): void { if (this.exitReason === "blocked") this.exitReason = undefined; this.idle = false; // A zero-delay completion dispatch may already be queued. Pi has not accepted // it yet, so reclaim it before this active turn can make it uncancellable. this.deferPendingCompletionDispatches(); this.blockedThisTurn = false; this.actedThisTurn = false; this.sleptThisTurn = false; } /** Record real activity; settle uses it for tick backoff and wait-series state. */ noteAction(toolName: string): void { if (toolName === "agi_sleep") return; this.actedThisTurn = true; } /** R-TOOL-21: a `blocked` note suppresses sleep for the rest of the turn (E62). */ noteBlocked(): void { this.blockedThisTurn = true; } isBlocked(): boolean { return this.blockedThisTurn; } /** `agi_sleep` hands the selected wake source here; armed at settle. */ requestSleep(request: SleepRequest): void { this.pendingSleep = request; this.sleptThisTurn = true; } sleepStreakCount(): number { return this.sleepStreak; } /** * R-SLEEP-2, and the §11.2 arming table: * * mode OFF → nothing * last note this turn was blocked → do not arm, hand to the user * any run active → arm five-minute WORKER CHECK * explicit agi_sleep request → arm its selected source/timer * otherwise → do not arm * * A pending external `agi_sleep({ms})` overrides the periodic tick. Worker * completion and attention remain independent event wakes. */ onAgentSettled(ctx: ExtensionContext): void { this.idle = true; this.cwd = ctx.cwd; this.sessionId = ctx.sessionManager.getSessionId(); if (!this.enabled) return; // Activity always resets quiet-tick backoff. It does not immediately erase the // internal wait continuity: sleep -> one cheap check -> the same sleep is still // one diagnostic series. No warning or count is exposed to the model or TUI. if (this.actedThisTurn) { this.quietTicks = 0; this.tickIntervalMs = this.host.config(ctx).tickMs; } if (this.blockedThisTurn) { // E62: the badge shows blocked and the user must respond. Arming here would // wake an orchestrator that is still blocked, burning a turn — forever. this.exitReason = "blocked"; this.disarmAll(true); this.host.onRefresh?.(); return; } this.exitReason = undefined; // R-SLEEP-19: a worker wake deferred behind the user's turn is delivered now // that the user's turn is over. this.flushDeferred(); const active = this.reviewEligibleEntries(ctx.cwd); if (this.pendingSleep !== undefined) { const request = this.pendingSleep; this.pendingSleep = undefined; if (request.source === "user") { this.clearTimed(); this.sleepStreak = 0; this.sleepNote = request.note; return; } const previousNote = this.sleepNote; const nextNote = request.note ?? previousNote; const sameWait = this.sleepStreak > 0 && nextNote === previousNote; this.sleepStreak = sameWait ? this.sleepStreak + 1 : 1; this.sleepNote = nextNote; if (request.source === "timer") { this.armTimed("sleep", request.durationMs, ctx, { ...(this.sleepNote === undefined ? {} : { note: this.sleepNote }), streak: this.sleepStreak, durationMs: request.durationMs, }); } else if (request.source === "tick") { this.armTimed("tick", this.tickIntervalMs, ctx, { ...(this.sleepNote === undefined ? {} : { note: this.sleepNote }), streak: this.sleepStreak, }); } else { if (active.length > 0) { this.armTimed("worker_check", WORKER_WAIT_MS, ctx, { ...(this.sleepNote === undefined ? {} : { note: this.sleepNote }), streak: this.sleepStreak, durationMs: WORKER_WAIT_MS, }); } else this.clearTimed(); } return; } if (this.sleptThisTurn) return; // A turn that does not return to the same external wait ends the wait series. // User input also resets it eagerly in onUserInput(). if (this.actedThisTurn) { this.sleepStreak = 0; this.sleepNote = undefined; } if (active.length > 0) { this.armTimed("worker_check", WORKER_WAIT_MS, ctx, { durationMs: WORKER_WAIT_MS }); return; } // No implicit plan-driven tick: without an explicit sleep or a worker event, // the orchestrator remains idle until the user speaks. this.clearTimed(); } private armTimed(kind: "tick" | "sleep" | "worker_check", delayMs: number, ctx: ExtensionContext, extra: { note?: string; streak?: number; durationMs?: number }): void { this.clearTimed(); const now = this.timers.now(); const firesAt = now + delayMs; const record = armWake(ctx.cwd, { kind, sessionId: this.sessionId, firesAt, now, ...extra, }); if (record === undefined) return; this.timed = { token: record.token, kind, firesAt, handle: this.timers.setTimeout(() => this.fireTimed(), delayMs), ...(extra.note === undefined ? {} : { note: extra.note }), ...(extra.streak === undefined ? {} : { streak: extra.streak }), ...(extra.durationMs === undefined ? {} : { durationMs: extra.durationMs }), }; } private clearTimed(): void { this.dropTimed(true); } private clearDispatches(predicate: (dispatch: PendingDispatch) => boolean, removeRecords: boolean): void { for (const dispatch of [...this.dispatches]) { if (!predicate(dispatch)) continue; this.timers.clearTimeout(dispatch.handle); this.dispatches.delete(dispatch); if (removeRecords && this.cwd !== undefined) { for (const token of dispatch.tokens) removeWake(this.cwd, token); } } } private disarmAll(removeDurable: boolean): void { if (removeDurable) this.clearTimed(); else this.dropTimed(false); this.batcher.clear(); this.deferred.clear(); this.pendingSleep = undefined; this.workerWakeTokens.clear(); this.clearDispatches(() => true, removeDurable); if (removeDurable && this.cwd !== undefined) { clearWakes(this.cwd, (record) => record.sessionId === this.sessionId); } } /** * `remove: false` drops the in-process timer but leaves the durable record alone. * * This distinction is R-SLEEP-5 / E70 and it is easy to lose: `session_shutdown` * must keep the record so the next `session_start` re-arms it (a `/reload` or a * crash mid-sleep is exactly the case durable wakes exist for), while toggle OFF * and goal-met must delete it (E64g — a wake outliving its purpose would wake an * orchestrator with nothing to do). Deleting on shutdown makes every restart lose * the pending sleep silently, which looks like nothing at all going wrong. */ private dropTimed(remove: boolean): void { const timed = this.timed; this.timed = undefined; if (timed === undefined) return; this.timers.clearTimeout(timed.handle); if (remove && this.cwd !== undefined) removeWake(this.cwd, timed.token); } private fireTimed(): void { const timed = this.timed; if (timed === undefined) return; const ctx = this.host.ctx(); if (ctx === undefined || !this.enabled) return; this.timed = undefined; this.timers.clearTimeout(timed.handle); if (timed.kind === "sleep") { this.deliverSleep(ctx, timed); return; } if (timed.kind === "worker_check") { this.deliverWorkerCheck(ctx, timed); return; } this.deliverTick(ctx, timed); } private deliverWorkerCheck(ctx: ExtensionContext, timed: TimedWake): void { const active = this.reviewEligibleEntries(ctx.cwd); if (active.length === 0) { settleWake(ctx.cwd, timed.token); return; } const review = renderWorkerReview(active); const payload: WakePayload = { reason: "worker_check", text: review.text, runIds: review.runIds }; this.dispatch(payload, { tokens: [timed.token], claim: () => settleWake(ctx.cwd, timed.token) === undefined ? { failed: true } : { payload }, }); } private deliverSleep(ctx: ExtensionContext, timed: TimedWake): void { const text = renderSleepPayload({ ...(timed.note === undefined ? {} : { note: timed.note }), }); const payload = { reason: "sleep" as const, text }; this.dispatch(payload, { tokens: [timed.token], claim: () => { const settled = settleWake(ctx.cwd, timed.token); return settled === undefined ? { failed: true } : { payload }; }, }); } /** * R-SLEEP-12: a worker wake within TICK_SUPPRESS_MS already gave the orchestrator a * turn, so the tick is skipped and simply re-armed (E59). */ private deliverTick(ctx: ExtensionContext, timed?: TimedWake): void { const now = this.timers.now(); if (now - this.lastWorkerWakeAt < TICK_SUPPRESS_MS) { if (timed !== undefined && settleWake(ctx.cwd, timed.token) === undefined) { this.armRecoveryTick(ctx); return; } this.armTimed("tick", this.tickIntervalMs, ctx, { ...(timed?.note === undefined ? {} : { note: timed.note }), ...(timed?.streak === undefined ? {} : { streak: timed.streak }), }); return; } const config = this.host.config(ctx); const text = renderTickPayload({ ...(timed?.note === undefined ? {} : { note: timed.note }) }); // R-SLEEP-11: after tickQuietThreshold consecutive no-op ticks the interval // doubles up to tickMaxMs. `actedThisTurn` is what makes a tick "no-op": the // counter is reset in onAgentSettled by any orchestrator action, and by any // worker event or user message below. this.quietTicks += 1; if (this.quietTicks >= config.tickQuietThreshold) { this.tickIntervalMs = Math.min(config.tickMaxMs, this.tickIntervalMs * 2); this.quietTicks = 0; } const payload = { reason: "tick" as const, text }; this.dispatch(payload, timed === undefined ? undefined : { tokens: [timed.token], claim: () => { const settled = settleWake(ctx.cwd, timed.token); return settled === undefined ? { failed: true } : { payload }; }, }); } /** * Called from the Supervisor's terminal-state callback. Arms the completion wake * through the nudge hold and the batch window. */ onRunTerminal(runId: string, status: RunStatus): void { const ctx = this.host.ctx(); if (!this.enabled || ctx === undefined) return; const entry = loadAllRuns(ctx.cwd).entries.find((candidate) => candidate.runId === runId); if (entry === undefined) return; this.host.onRefresh?.(); if (this.blockedThisTurn || this.exitReason !== undefined) return; if (!this.host.config(ctx).notifyOnComplete) return; // A result the orchestrator already read needs no wake at all (R-SLEEP-13). if (status.resultConsumed || status.headlessDeliveredAt !== undefined) { this.onResultConsumed(entry.runId); return; } const report = this.readReport(entry); const item = completionItem(entry.runId, status, report, this.timers.now()); if (this.handledWorkerWakes.has(entry.runId)) { const delivered = this.deliveredWorkerEvidence.get(entry.runId); if (delivered === undefined || delivered === item.evidenceKey) return; // New semantic evidence after a claimed wake gets one bounded refresh. An // exact replay keeps the same key and remains suppressed. this.handledWorkerWakes.delete(entry.runId); } if (this.timed?.kind === "worker_check") this.clearTimed(); this.clearDispatches((dispatch) => dispatch.reason === "worker_check", true); // Only a genuinely new worker event resets backoff and supersedes the current // fleet trajectory check. Duplicate/no-op callbacks must not cancel a sibling's. this.quietTicks = 0; // R-SLEEP-5: persist before the 200ms hold. A process crash inside that tiny // window must not lose the completion wake — "durable" includes inconvenient // timing, not just tick and sleep records. let token = this.workerWakeTokens.get(entry.runId); if (token === undefined) { const record = armWake(ctx.cwd, { kind: "worker", sessionId: this.sessionId, runId: entry.runId, now: this.timers.now(), }); if (record === undefined) return; token = record.token; this.workerWakeTokens.set(entry.runId, token); } this.queueCompletionItem({ ...item, wakeToken: token }); } private readReport(entry: RegistryEntry): string | undefined { try { return readFileIfExists(entry.paths.result); } catch { return undefined; } } /** Upsert one run into its single current scheduler-owned completion layer. */ private queueCompletionItem(item: CompletionItem): void { let pending = false; for (const dispatch of this.dispatches) { if (dispatch.reason !== "worker_complete" || dispatch.completionItems === undefined) continue; const index = dispatch.completionItems.findIndex((candidate) => candidate.runId === item.runId); if (index < 0) continue; dispatch.completionItems[index] = item; if (item.wakeToken !== undefined && !dispatch.tokens.includes(item.wakeToken)) dispatch.tokens.push(item.wakeToken); pending = true; } if (pending) { this.deferred.delete(item.runId); this.batcher.cancel(item.runId); return; } if (this.deferred.has(item.runId)) { this.deferred.set(item.runId, item); this.batcher.cancel(item.runId); return; } this.batcher.submit(item); } private deferCompletions(items: CompletionItem[]): void { for (const item of items) this.deferred.set(item.runId, item); } /** Reclaim zero-delay completion sends before an active parent turn begins. */ private deferPendingCompletionDispatches(): void { for (const dispatch of [...this.dispatches]) { if (dispatch.reason !== "worker_complete" || dispatch.completionItems === undefined) continue; this.timers.clearTimeout(dispatch.handle); this.dispatches.delete(dispatch); this.deferCompletions(dispatch.completionItems); } } /** R-SLEEP-13: cancel this run across every scheduler-owned completion layer. */ onResultConsumed(runId: string): void { this.batcher.cancel(runId); const tokens = new Set(); const mappedToken = this.workerWakeTokens.get(runId); if (mappedToken !== undefined) tokens.add(mappedToken); const deferred = this.deferred.get(runId); if (deferred?.wakeToken !== undefined) tokens.add(deferred.wakeToken); this.deferred.delete(runId); for (const dispatch of [...this.dispatches]) { if (dispatch.reason !== "worker_complete" || dispatch.completionItems === undefined) continue; const removed = dispatch.completionItems.filter((item) => item.runId === runId); for (const item of removed) if (item.wakeToken !== undefined) tokens.add(item.wakeToken); dispatch.completionItems.splice(0, dispatch.completionItems.length, ...dispatch.completionItems.filter((item) => item.runId !== runId)); dispatch.tokens = dispatch.tokens.filter((token) => !tokens.has(token)); if (dispatch.completionItems.length > 0) continue; this.timers.clearTimeout(dispatch.handle); this.dispatches.delete(dispatch); } this.workerWakeTokens.delete(runId); this.handledWorkerWakes.add(runId); this.deliveredWorkerEvidence.delete(runId); if (this.cwd !== undefined) { for (const token of tokens) removeWake(this.cwd, token); clearWakes(this.cwd, (record) => record.sessionId === this.sessionId && record.kind === "worker" && record.runId === runId); } } private deliverCompletions(items: CompletionItem[]): void { const ctx = this.host.ctx(); if (ctx === undefined || !this.enabled || this.blockedThisTurn || this.exitReason !== undefined) { // Not deliverable here. Leave the records on disk so the next session_start // re-arms them rather than silently dropping the completions. return; } const unique = new Map(); for (const item of items) unique.set(item.runId, item); const dispatchItems = [...unique.values()]; if (!this.idle) { this.deferCompletions(dispatchItems); return; } const tokens = dispatchItems.flatMap((item) => item.wakeToken === undefined ? [] : [item.wakeToken]); this.dispatch({ reason: "worker_complete", text: "", runIds: dispatchItems.map((item) => item.runId), }, { tokens, completionItems: dispatchItems, claim: (liveCtx) => { const claimed: CompletionItem[] = []; let failed = false; const currentByRun = new Map(loadAllRuns(liveCtx.cwd).entries.map((entry) => [entry.runId, entry])); for (const item of dispatchItems) { const token = item.wakeToken; if (token === undefined) continue; this.workerWakeTokens.delete(item.runId); const current = currentByRun.get(item.runId); if (current === undefined || !isTerminalState(current.status.state)) { removeWake(liveCtx.cwd, token); continue; } if (current.status.resultConsumed || current.status.headlessDeliveredAt !== undefined) { this.onResultConsumed(item.runId); continue; } const currentItem = completionItem(current.runId, current.status, this.readReport(current), this.timers.now()); if (settleWake(liveCtx.cwd, token) === undefined) { failed = true; // R-SLEEP-6 recovers through one ordinary tick. A repeated terminal // callback must not arm a second worker record beside that recovery. this.handledWorkerWakes.add(item.runId); this.deliveredWorkerEvidence.set(item.runId, currentItem.evidenceKey); continue; } this.handledWorkerWakes.add(item.runId); this.deliveredWorkerEvidence.set(item.runId, currentItem.evidenceKey); clearWakes(liveCtx.cwd, (record) => record.sessionId === this.sessionId && record.kind === "worker" && record.runId === item.runId); claimed.push({ ...currentItem, wakeToken: token }); } if (claimed.length === 0) return { failed }; return { failed, payload: { reason: "worker_complete", text: renderCompletionPayload(claimed), runIds: claimed.map((item) => item.runId), }, }; }, }); } /** * R-SLEEP-19/20/21. A user message supersedes pending **tick** and **sleep** wakes: * the user's turn is the wake. Worker completion wakes are not cancelled — they * report an independent event — and are re-queued as `followUp` behind the user's * turn (E61). */ onUserInput(): void { this.idle = false; this.quietTicks = 0; this.sleepStreak = 0; this.sleepNote = undefined; // User input can arrive after agi_sleep requested a wait but before the turn's // settle callback. Drop that stale request while retaining sleptThisTurn, so // settlement cannot fall through and manufacture a fresh worker check. if (this.pendingSleep !== undefined) this.pendingSleep = undefined; this.host.onWakeReason("user", "the user typed"); // E64c: a timed sleep has no independent event to report, so cancelling it loses // nothing. A tick is likewise superseded. this.clearTimed(); this.clearDispatches( (dispatch) => dispatch.reason === "tick" || dispatch.reason === "sleep" || dispatch.reason === "worker_check", true, ); this.deferPendingCompletionDispatches(); const pending = this.batcher.pendingRunIds(); if (pending.length === 0) return; // Deliberately not cancelled: the completion still has to reach the model. It is // held until the user's turn settles, which is what "re-queued behind the user's // turn" means when the wake would otherwise land as a steer. this.batcher.clear(); const ctx = this.host.ctx(); if (ctx === undefined) return; const { entries } = loadAllRuns(ctx.cwd); const items: CompletionItem[] = []; for (const runId of pending) { const entry = entries.find((candidate) => candidate.runId === runId); if (entry === undefined || entry.status.resultConsumed || entry.status.headlessDeliveredAt !== undefined) { this.onResultConsumed(runId); continue; } const item = completionItem(entry.runId, entry.status, this.readReport(entry), this.timers.now()); const token = this.workerWakeTokens.get(entry.runId); items.push(token === undefined ? item : { ...item, wakeToken: token }); } if (items.length === 0) return; // Carried as items rather than rendered text, so `deliverCompletions` still // performs the R-SLEEP-6 claim and the resultConsumed re-check when it finally // runs. Rendering here would leave the durable records armed *and* deliver them, // which is a double wake after a restart. this.deferCompletions(items); } private flushDeferred(): void { if (this.deferred.size === 0) return; const items = [...this.deferred.values()]; this.deferred.clear(); this.deliverCompletions(items); } /** * Every send goes through here, and every send is deferred by a timer. * * `_emitAgentSettled` sets `_isAgentRunActive = false` before awaiting extension * handlers, so a `sendMessage({triggerTurn:true})` issued synchronously from an * `agent_settled` handler is treated as "idle" and re-enters `_runAgentPrompt` * from inside the previous run's `finally` block. Deferring by a tick puts the * send back on the macrotask queue where pi expects it. */ private dispatch( payload: WakePayload, options: { tokens?: string[]; completionItems?: CompletionItem[]; claim?: (ctx: ExtensionContext) => DispatchClaim } = {}, ): void { const epoch = this.dispatchEpoch; const sessionId = this.sessionId; const pending: PendingDispatch = { handle: undefined, epoch, sessionId, reason: payload.reason, tokens: options.tokens ?? [], ...(options.completionItems === undefined ? {} : { completionItems: options.completionItems }), }; pending.handle = this.timers.setTimeout(() => { this.dispatches.delete(pending); const ctx = this.host.ctx(); if ( ctx === undefined || !this.enabled || epoch !== this.dispatchEpoch || sessionId !== this.sessionId || ctx.sessionManager.getSessionId() !== sessionId ) return; if (payload.reason === "worker_complete" && !this.idle) { this.deferCompletions(pending.completionItems ?? []); return; } if (this.blockedThisTurn || this.exitReason !== undefined) { this.exitReason = "blocked"; this.disarmAll(true); this.host.onRefresh?.(); return; } const claim = options.claim?.(ctx) ?? { payload }; if (claim.payload === undefined) { if (claim.failed) this.armRecoveryTick(ctx); return; } let sent = false; try { sent = this.host.send(claim.payload) === true; } catch { sent = false; } if (!sent) { this.armRecoveryTick(ctx); return; } if (claim.payload.reason === "worker_complete") this.lastWorkerWakeAt = this.timers.now(); this.host.onWakeReason(claim.payload.reason, wakeDetail(claim.payload)); if (claim.failed) this.armRecoveryTick(ctx); }, 0); this.dispatches.add(pending); } private armRecoveryTick(ctx: ExtensionContext): void { if (!this.enabled || this.blockedThisTurn || this.exitReason !== undefined) return; this.armTimed("tick", this.tickIntervalMs, ctx, {}); } /** * Attention polling (§11.6). Runs on the scheduler interval rather than on worker * events, because the two most valuable triggers — idle and stalled_tool — are * defined by the *absence* of events. */ private poll(): void { const ctx = this.host.ctx(); if (ctx === undefined || !this.enabled) return; this.host.onRefresh?.(); if (this.blockedThisTurn || this.exitReason !== undefined) return; let config: WorkerConfig; try { config = this.host.config(ctx); } catch { return; } if (!config.attention.enabled) return; const now = this.timers.now(); let entries: RegistryEntry[]; try { entries = loadAllRuns(ctx.cwd).entries; } catch { return; } for (const entry of entries) { if (!this.belongsToSession(entry.status)) continue; if (isEndedState(entry.status.state)) continue; let findings: AttentionFinding[]; try { findings = evaluateAttention({ status: entry.status, events: readEvents(entry.paths, 200), config: config.attention, now, }); } catch { continue; } const fired = new Set( (entry.status.attention?.fired ?? []).filter(isAttentionTrigger).map((trigger) => attentionKey(entry.runId, trigger)), ); const selection = selectAttention(entry.runId, findings, fired); if (selection.primary === undefined) continue; // R-SLEEP-16: long_running goes to the widget and the event log only, never // to chat. `chatWorthy` is false for exactly that trigger. const chatWorthy = selection.fresh.filter((finding) => finding.chatWorthy); if (chatWorthy.length === 0 || !config.notifyOnAttention) { this.recordAttention(entry, selection.primary, selection.fresh.map((finding) => finding.trigger)); continue; } const primary = chatWorthy[0]; if (primary === undefined) continue; const record = armWake(ctx.cwd, { kind: "attention", sessionId: this.sessionId, runId: entry.runId, attentionTrigger: primary.trigger, attentionDetail: primary.detail, now, }); if (record === undefined) continue; if (!this.recordAttention(entry, selection.primary, selection.fresh.map((finding) => finding.trigger))) { removeWake(ctx.cwd, record.token); continue; } if (this.timed?.kind === "worker_check") this.clearTimed(); this.clearDispatches((dispatch) => dispatch.reason === "worker_check", true); const payload: WakePayload = { reason: "worker_attention", text: renderAttentionPayload(entry.status, primary, { controlAvailable: this.host.controlAvailable?.() === true }), runIds: [entry.runId], }; this.dispatch(payload, { tokens: [record.token], claim: () => { const settled = settleWake(ctx.cwd, record.token); return settled === undefined ? { failed: true } : { payload }; }, }); } } /** * `status.attention` is an externally-owned field: the pump merges rather than * overwrites it, so this read-modify-write is the only writer and must preserve * everything else in the file it re-reads. */ private recordAttention(entry: RegistryEntry, finding: AttentionFinding, fired: AttentionTrigger[]): boolean { try { const next = patchStatus(entry.paths, (current) => { const previous = current.attention?.fired ?? []; return { ...current, attention: { reason: finding.trigger, since: new Date(this.timers.now()).toISOString(), detail: finding.detail, fired: [...new Set([...previous, ...fired])], }, }; }); if (next === undefined) return false; entry.status = next; return true; } catch { // E37: a status write that fails must never crash the orchestrator. return false; } } /** State for the sleep indicator (§12.3). */ indicatorState(ctx: ExtensionContext): { activeRuns: number; dominant?: { name: string; activity: string; elapsedMs: number }; sleepRemainingMs?: number; sleepNote?: string; sleepStreak?: number; tickInMs?: number; } { const now = this.timers.now(); let entries: RegistryEntry[] = []; try { entries = loadAllRuns(ctx.cwd).entries; } catch { entries = []; } const active = entries.filter((entry) => !isTerminalState(entry.status.state) && entry.status.state !== "queued"); // R-UI-7: the dominant worker is the longest-running one, which is the one an // observer is most likely to be wondering about. const dominant = active .slice() .sort((a, b) => Date.parse(a.status.startedAt ?? a.status.createdAt) - Date.parse(b.status.startedAt ?? b.status.createdAt))[0]; const snapshot = this.snapshot(); return { activeRuns: active.length, ...(dominant === undefined ? {} : { dominant: { name: dominant.status.name, activity: humanizeActivity(dominant, now) || "thinking…", elapsedMs: Math.max(0, now - Date.parse(dominant.status.startedAt ?? dominant.status.createdAt)), }, }), ...(snapshot.sleepRemainingMs === undefined ? {} : { sleepRemainingMs: snapshot.sleepRemainingMs }), ...(snapshot.sleepNote === undefined ? {} : { sleepNote: snapshot.sleepNote }), ...(this.sleepStreak === 0 ? {} : { sleepStreak: this.sleepStreak }), ...(snapshot.tickInMs === undefined ? {} : { tickInMs: snapshot.tickInMs }), }; } reviewEligibleRunCount(ctx: ExtensionContext): number { try { return this.reviewEligibleEntries(ctx.cwd).length; } catch { return 0; } } isIdle(): boolean { return this.idle; } exitState(): SchedulerExit | undefined { return this.exitReason; } /** * R-SLEEP-22 (headless drain). In `print` and `json` modes there is nobody to hand * control back to and the process exits when prompts are done, so detached workers * would be orphaned. * * The `while` loop, not a single wait, is deliberate: the orchestrator may delegate * *more* work in response to a completion, and that work must also be drained * (E69). On the deadline with runs still active they are stopped and a clear error * is emitted — never a silent exit with orphans (E68). * * R-SLEEP-23: RPC behaves like interactive and does **not** auto-drain; its client * owns the lifecycle. The caller enforces that by only invoking this for `!hasUI`, * and `hasUI` is true in RPC. */ async drainHeadless( ctx: ExtensionContext, options: { deadlineMs: number; stopAll: (entries: RegistryEntry[]) => void; deliver: (payload: WakePayload) => Promise | void; pollMs?: number; now?: () => number; sleep?: (ms: number) => Promise; }, ): Promise<{ drained: number; timedOut: boolean }> { if (this.draining) return { drained: 0, timedOut: false }; this.draining = true; const now = options.now ?? (() => this.timers.now()); const pollMs = options.pollMs ?? 500; const sleep = options.sleep ?? ((ms: number) => new Promise((resolve) => setTimeout(resolve, ms))); const deadline = now() + options.deadlineMs; const reported = new Set(); let drained = 0; try { while (now() < deadline) { const { entries } = loadAllRuns(ctx.cwd); const active = entries.filter( (entry) => this.belongsToSession(entry.status) && !isTerminalState(entry.status.state), ); const finished = entries.filter( (entry) => this.belongsToSession(entry.status) && isTerminalState(entry.status.state) && !entry.status.resultConsumed && entry.status.headlessDeliveredAt === undefined && !reported.has(entry.runId), ); for (const entry of finished) { reported.add(entry.runId); const item = completionItem(entry.runId, entry.status, this.readReport(entry), now()); // Delivered as a followUp so the orchestrator can react — including by // delegating more work, which the loop then picks up. await options.deliver({ reason: "worker_complete", text: renderCompletionPayload([item]), runIds: [entry.runId] }); const deliveredAt = new Date(now()).toISOString(); const marked = patchStatus(entry.paths, (current) => ({ ...current, headlessDeliveredAt: deliveredAt })); if (marked !== undefined) { drained += 1; this.onResultConsumed(entry.runId); } } if (active.length === 0) { // Re-check after delivery: the orchestrator may have delegated more work // while reacting to the completion above. const recheck = loadAllRuns(ctx.cwd).entries.filter( (entry) => this.belongsToSession(entry.status) && !isTerminalState(entry.status.state), ); if (recheck.length === 0 && finished.length === 0) return { drained, timedOut: false }; } await sleep(pollMs); } const { entries } = loadAllRuns(ctx.cwd); const stillActive = entries.filter( (entry) => this.belongsToSession(entry.status) && !isTerminalState(entry.status.state), ); if (stillActive.length === 0) return { drained, timedOut: false }; options.stopAll(stillActive); return { drained, timedOut: true }; } finally { this.draining = false; } } /** R-UI-9 helper: everything the scheduler installed, gone. */ shutdown(): void { // R-SLEEP-5 / E70: the durable record deliberately survives, so the next // session_start re-arms it. Only the in-process timer is dropped here. this.dispatchEpoch += 1; this.idle = false; this.disarmAll(false); if (this.pollHandle !== undefined) { this.timers.clearInterval(this.pollHandle); this.pollHandle = undefined; } this.enabled = false; } } function wakeDetail(payload: WakePayload): string { switch (payload.reason) { case "worker_complete": return `${payload.runIds?.join(", ") ?? ""} finished`; case "worker_attention": return `${payload.runIds?.join(", ") ?? ""} needs attention`; case "tick": return "scheduled check"; case "sleep": return "timed sleep elapsed"; case "startup": return "session start with unconsumed runs"; default: return ""; } } export type { AttentionTrigger }; export { attentionKey, elapsed };