/** * Attention triggers (§11.6) and their dedupe (R-SLEEP-17). * * R-IMPL-1: pure functions over plain data. Everything here takes a `RunStatus`, a * clock and a set of already-fired triggers, and returns a decision. No I/O, no * timers, no `Date.now()` default that a test cannot control — this is the logic most * likely to be subtly wrong (an idle check that fires while a tool is running turns * every `npm install` into a false alarm), so it is the primary unit-test target. */ import type { AttentionConfig } from "../worker/config.ts"; import { type RunEvent, type RunStatus, isEndedState, isRunActiveTool } from "../worker/status.ts"; export type AttentionTrigger = | "idle" | "repeated_failure" | "long_running" | "stalled_tool" | "context_pressure" | "compaction_churn"; const ATTENTION_TRIGGERS = new Set([ "idle", "repeated_failure", "long_running", "stalled_tool", "context_pressure", "compaction_churn", ]); export function isAttentionTrigger(value: string): value is AttentionTrigger { return ATTENTION_TRIGGERS.has(value); } export interface AttentionFinding { trigger: AttentionTrigger; /** One line of evidence, quoted into the payload (R-SLEEP-18). */ detail: string; /** * R-SLEEP-16: `long_running` alone goes to the widget and the event log and * **never** to chat. A worker legitimately taking 20 minutes is normal, and waking * the orchestrator to say "still working" is pure cost. */ chatWorthy: boolean; } /** Order matters: the first finding is what the payload leads with. */ const TRIGGER_PRIORITY: AttentionTrigger[] = [ "stalled_tool", "repeated_failure", "context_pressure", "compaction_churn", "idle", "long_running", ]; function seconds(ms: number): string { return `${Math.round(ms / 1000)}s`; } function minutes(ms: number): string { const total = Math.round(ms / 60000); if (total < 60) return `${total}m`; return `${Math.floor(total / 60)}h${total % 60}m`; } /** * Count tool failures inside the rolling window, and separately the worst per-path * count. §11.6: "3 consecutive failures, **or** 3 failures on the same path, within a * 5-minute window". Both halves matter — a worker alternating a broken edit with a * successful read is not consecutive-failing, but it is looping. */ export function analyzeFailures( events: RunEvent[], now: number, windowMs: number, ): { consecutive: number; samePath: number; path: string | undefined; lastError: string | undefined } { let consecutive = 0; let bestConsecutive = 0; const perPath = new Map(); let lastError: string | undefined; for (const event of events) { if (event.kind !== "tool_end") continue; const at = Date.parse(String(event.ts)); if (Number.isNaN(at) || now - at > windowMs) { // Outside the window: it cannot contribute, and it must not carry a // consecutive streak across the boundary either. consecutive = 0; continue; } if (event.isError === true) { consecutive += 1; bestConsecutive = Math.max(bestConsecutive, consecutive); const target = typeof event.target === "string" ? event.target : undefined; if (target !== undefined) perPath.set(target, (perPath.get(target) ?? 0) + 1); if (typeof event.tool === "string") lastError = event.tool; } else { consecutive = 0; } } let samePath = 0; let path: string | undefined; for (const [candidate, count] of perPath) { if (count > samePath) { samePath = count; path = candidate; } } return { consecutive: bestConsecutive, samePath, path, lastError }; } export interface EvaluateInput { status: RunStatus; /** Recent `tool_start`/`tool_end` events, for the failure window. */ events: RunEvent[]; config: AttentionConfig; now: number; } interface ActiveToolEvidence { id: string; tool: string; target: string | null; startedAt: number; lastProgressAt: number; } function timestamp(value: string | null | undefined): number | undefined { if (value === null || value === undefined) return undefined; const parsed = Date.parse(value); return Number.isNaN(parsed) ? undefined : parsed; } function evidenceOf(invocation: unknown): ActiveToolEvidence | undefined { if (!isRunActiveTool(invocation)) return undefined; const startedAt = timestamp(invocation.startedAt); if (startedAt === undefined) return undefined; return { id: invocation.id, tool: invocation.tool ?? "unknown", target: invocation.target, startedAt, lastProgressAt: Math.max(startedAt, timestamp(invocation.lastProgressAt) ?? startedAt), }; } function activeToolState(status: RunStatus): { present: boolean; evidence: ActiveToolEvidence[] } { const authoritative = status.activity.activeTools; if (Array.isArray(authoritative)) { return { present: authoritative.length > 0, evidence: authoritative.map(evidenceOf).filter((value): value is ActiveToolEvidence => value !== undefined), }; } // Compatibility with schema-1 status files written before per-invocation // activity existed. Their worker-wide lastEventAt remains the best available // approximation of progress for the one legacy projection. if (status.activity.currentTool === null) return { present: false, evidence: [] }; const startedAt = timestamp(status.activity.currentToolStartedAt); if (startedAt === undefined) return { present: true, evidence: [] }; return { present: true, evidence: [{ id: "legacy", tool: status.activity.currentTool, target: status.activity.currentPath, startedAt, lastProgressAt: Math.max(startedAt, timestamp(status.activity.lastEventAt) ?? startedAt), }], }; } /** * Every trigger this run currently satisfies, in priority order. Dedupe is applied by * the caller (`selectAttention`) so this stays a pure predicate over current state. */ export function evaluateAttention(input: EvaluateInput): AttentionFinding[] { const { status, config, now } = input; if (!config.enabled) return []; if (isEndedState(status.state)) return []; // A queued run has no process and no activity; every timer below would measure the // time it spent waiting for a slot, which is not a worker problem. if (status.state === "queued") return []; const findings = new Map(); const activeTools = activeToolState(status); // Stalled tool — per-invocation progress means healthy sibling activity cannot // hide a quiet call, and a progressing old call is not falsely blamed. const stalled = activeTools.evidence .map((tool) => ({ tool, runtime: now - tool.startedAt, quietFor: now - tool.lastProgressAt })) .filter(({ runtime, quietFor }) => runtime >= config.stalledToolMs && quietFor >= config.stalledToolMs) .sort((a, b) => b.quietFor - a.quietFor || b.runtime - a.runtime || a.tool.id.localeCompare(b.tool.id))[0]; if (stalled !== undefined) { findings.set("stalled_tool", { trigger: "stalled_tool", detail: `${stalled.tool.tool}${stalled.tool.target === null ? "" : ` "${stalled.tool.target}"`} ` + `has run ${seconds(stalled.runtime)} with no progress for ${seconds(stalled.quietFor)}.`, chatWorthy: true, }); } // Idle. The "no tool currently running" condition is essential: a worker inside a // five-minute `npm install` is busy, not idle, and without this every slow build // produces a false alarm. const lastEventAt = status.activity.lastEventAt === null ? undefined : Date.parse(status.activity.lastEventAt); if (!activeTools.present && lastEventAt !== undefined && !Number.isNaN(lastEventAt) && now - lastEventAt > config.idleMs) { findings.set("idle", { trigger: "idle", detail: `No worker event for ${seconds(now - lastEventAt)} and no tool is running.`, chatWorthy: true, }); } // Repeated tool failure. const failures = analyzeFailures(input.events, now, config.failureWindowMs); if (failures.consecutive >= config.failedToolAttempts || failures.samePath >= config.failedToolAttempts) { const detail = failures.samePath >= config.failedToolAttempts && failures.path !== undefined ? `${failures.samePath} tool failures on "${failures.path}" within ${minutes(config.failureWindowMs)}.` : `${failures.consecutive} consecutive tool failures within ${minutes(config.failureWindowMs)}.`; findings.set("repeated_failure", { trigger: "repeated_failure", detail, chatWorthy: true }); } // Long-running. Widget and event log only (R-SLEEP-16). const startedAt = Date.parse(status.startedAt ?? status.createdAt); if (!Number.isNaN(startedAt) && now - startedAt > config.longRunningMs) { findings.set("long_running", { trigger: "long_running", detail: `Running for ${minutes(now - startedAt)}.`, chatWorthy: false, }); } // Context pressure: the worker is about to compact repeatedly. if (status.context !== null && status.context.percent > config.contextPercent) { findings.set("context_pressure", { trigger: "context_pressure", detail: `Worker context is at ${status.context.percent}% (threshold ${config.contextPercent}%).`, chatWorthy: true, }); } // Compaction churn: the task is too large (E36). if (status.counters.compactions >= config.compactionChurn) { findings.set("compaction_churn", { trigger: "compaction_churn", detail: `Worker has compacted ${status.counters.compactions} times; the task is probably too large.`, chatWorthy: true, }); } const ordered: AttentionFinding[] = []; for (const trigger of TRIGGER_PRIORITY) { const finding = findings.get(trigger); if (finding !== undefined) ordered.push(finding); } return ordered; } export type FiredSet = ReadonlySet; export function attentionKey(runId: string, trigger: AttentionTrigger): string { return `${runId}\u0000${trigger}`; } export interface AttentionSelection { /** Findings that have not fired for this run before, in priority order. */ fresh: AttentionFinding[]; /** The leading fresh finding, or undefined. */ primary: AttentionFinding | undefined; /** Keys to record as fired. */ keys: string[]; } /** * R-SLEEP-17: each `(runId, trigger)` pair fires at most once per run. Without it an * idle worker generates an attention wake on every poll interval — the same wake, * forever, for as long as it stays idle. */ export function selectAttention(runId: string, findings: AttentionFinding[], fired: FiredSet): AttentionSelection { const fresh = findings.filter((finding) => !fired.has(attentionKey(runId, finding.trigger))); return { fresh, primary: fresh[0], keys: fresh.map((finding) => attentionKey(runId, finding.trigger)), }; } function sanitizeAttentionLabel(value: string | null | undefined, max = 120): string | null { if (value === null || value === undefined) return null; const oneLine = value .replace(/[\u0000-\u001f\u007f-\u009f]+/g, " ") .replace(/\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(?:\.\d+)?Z/gi, "[time]") .replace(/\$\s*\d+(?:\.\d+)?/g, "[cost]") .replace(/\b\d+(?:\.\d+)?%/g, "[percentage]") .replace(/\b\d+(?:\.\d+)?\s*(?:turns?|tool calls?|tool errors?|compactions?|tokens?|cycles?)\b/gi, "[metric]") .replace(/\d+/g, "#") .replace(/\s+/g, " ") .trim(); return oneLine.length === 0 ? null : oneLine.slice(0, max); } export function modelFacingAttentionEvidence( finding: { trigger: string; detail: string }, status?: Pick, ): string { switch (finding.trigger) { case "stalled_tool": { // The legacy currentTool projection follows the most recently progressing // invocation, which may be a healthy sibling. The selected stalled finding // already carries the exact invocation in its durable detail, so derive the // bounded model-facing label from that evidence. This also survives a restart, // when only the durable trigger/detail pair is rearmed. const targeted = /^(.+?) "([^"]*)" has run /.exec(finding.detail); const untargeted = targeted === null ? /^(.+?) has run /.exec(finding.detail) : null; const tool = sanitizeAttentionLabel(targeted?.[1] ?? untargeted?.[1] ?? status?.activity.currentTool, 60); const target = sanitizeAttentionLabel(targeted?.[2] ?? (targeted === null && untargeted === null ? status?.activity.currentPath : null), 120); return `The worker${tool === null ? "" : `'s ${tool} tool`}${target === null ? "" : ` on "${target}"`} appears stalled.`; } case "repeated_failure": { const target = /on "([^"]+)"/.exec(finding.detail)?.[1]; const safeTarget = sanitizeAttentionLabel(target, 120); return safeTarget === null ? "The worker is repeatedly encountering tool failures." : `The worker is repeatedly encountering tool failures on "${safeTarget}".`; } case "context_pressure": return "The worker is under context pressure."; case "compaction_churn": return "The worker is compacting repeatedly; the task may be too large."; case "idle": return "No tool is running and the worker has stopped producing events."; case "long_running": return "The worker has been running unusually long."; default: return "The worker needs review."; } } /** R-SLEEP-18: inject only the trigger and bounded evidence; standing actions live in the system prompt. */ export function renderAttentionPayload(status: RunStatus, finding: AttentionFinding, options: { controlAvailable: boolean }): string { void options; return [ `Agent ${status.name} needs attention: ${finding.trigger}.`, modelFacingAttentionEvidence(finding, status), ].join("\n"); }