diff --git a/.gitignore b/.gitignore index 3c45938..6fd523b 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,4 @@ node_modules/ *.log .DS_Store +*.swp diff --git a/README.md b/README.md index dd00aaf..3f8cd4c 100644 --- a/README.md +++ b/README.md @@ -95,6 +95,7 @@ defaults: "options": { "maxTurns": 20, "stallLimit": 2, + "pollLimit": 3, "judgeModel": { "providerID": "openrouter", "id": "google/gemini-3-flash-preview" } } } @@ -106,6 +107,7 @@ defaults: | --- | --- | --- | | `maxTurns` | `20` | Automatic continuation turns before the loop auto-pauses. A whole number, or `"unlimited"` / `null` for no ceiling. The initial `/goal` turn is not counted, so `20` allows 21 agent executions in total. See below. | | `stallLimit` | `2` | Consecutive turns that ran **no tools** before the loop is declared stalled. | +| `pollLimit` | `3` | Consecutive turns that re-read an unchanged result before the loop is called polling. Distinct from `stallLimit`: a polling turn *does* use tools, it just learns nothing. | | `judgeModel` | session model | Model used for the `done` / `continue` / `blocked` verdict. | | `quiet` | `false` | Stop posting the loop's turn banner and completion notices into the transcript while the panel is open. | @@ -143,6 +145,7 @@ until the session ends — it never rewrites your config. | --- | --- | --- | | `/goal budget ` | `maxTurns` | See [Removing the turn limit](#removing-the-turn-limit) | | `/goal stall ` | `stallLimit` | Turns with **no tool calls** before giving up | +| `/goal poll ` | `pollLimit` | Turns reading the **same unchanged result** before calling it polling | | `/goal quiet ` | `quiet` | When on, the panel replaces the loop's transcript notices | | `/goal judge ` | `judgeModel` | Cheaper and sharper models judge better and cost less | | `/goal settings` | — | Shows all four, and whether each is a session override or the config default | @@ -206,9 +209,11 @@ It applies to the running goal immediately and to every goal set afterwards in t | Still stops it | | | --- | --- | | The judge says `done` | the goal is met, with evidence | -| The judge says `blocked` | the goal is unreachable, or the agent is going in circles | +| The judge says `blocked` | the goal is unreachable, on goal-level evidence | | Stall guard | `stallLimit` consecutive turns with no tool calls | +| Polling guard | `pollLimit` consecutive turns with an unchanged result | | Repetition guard | the same reply twice running | +| Polling guard | 3 turns that re-read an unchanged result | | You | `/goal pause`, `/goal clear`, or esc | So an agent that keeps making small, genuine-looking progress can now run indefinitely. Nothing @@ -299,6 +304,23 @@ non-trivial. Comma or space separated, with an optional trailing `on`/`off`: /goal display footer,sidebar off turn both off, leave the rest alone ``` +A change appears immediately, without reopening anything. + +That took a fix worth recording, because the cause was not obvious from the symptom. +`context.storage.store` is durable but **not reactive**, and a slot's `render` runs once, so +`` read a snapshot of the placements and never re-read it. The +visible effect was that a placement change only appeared after something forced the slot to +re-render — in practice, toggling the panel with `/goal panel`, which is exactly the manual +workaround that should not be needed. + +So the stored value is mirrored into a signal, which is the mechanism the goal state already +used and which demonstrably re-renders in this host. The store stays the single durable +record; the signal is only what the UI reads, and it leads the write so a change shows on +the frame it happens rather than after the storage round trip. Because the UI shows four +placements but only ever changes the named one, `applyMutation` returns the **whole** value +rather than the keys that changed — otherwise an untouched placement would silently revert on +the next write. Both properties are pinned in `test/display.test.ts`. + With an explicit `on`/`off` the named placements are set and the rest untouched; without one, each named placement toggles. @@ -366,26 +388,65 @@ Only these exact prefixes are recognised, so an ordinary goal containing a colon ## How the loop stops -The point of this plugin is that it terminates. Five conditions, checked in order: +The point of this plugin is that it terminates. Six conditions, checked in order: | # | Condition | Kind | | --- | --- | --- | | 1 | Judge returns `done` — the reply carries concrete evidence, such as a passing command and its output | model | -| 2 | Judge returns `blocked` — impossible, out of scope, needs credentials or hardware you do not have, or the agent is going in circles | model | +| 2 | Judge returns `blocked` — impossible, out of scope, needs credentials or hardware you do not have | model | | 3 | **Stall** — `stallLimit` consecutive turns ran no tools at all, so nothing changed however confident the prose | deterministic | | 4 | **Repetition** — the agent produced the same reply twice running | deterministic | -| 5 | **Budget** — `maxTurns` continuation turns spent (21 executions by default) | deterministic | +| 5 | **Polling** — `pollLimit` consecutive turns ran tools and read back the same unchanged result | deterministic | +| 6 | **Budget** — `maxTurns` continuation turns spent (21 executions by default) | deterministic | -Conditions 3–5 do not consult the model. This is deliberate: in testing, a weak judge model +Conditions 3–6 do not consult the model. This is deliberate: in testing, a weak judge model answered `continue` to twenty byte-identical replies and happily spent the entire budget. The deterministic guards are what actually stopped it, in two turns and about $0.004 instead of twenty turns and roughly $0.02. Treat the judge's `blocked` verdict as a useful fourth opinion, not as the safety net. +### `blocked` means the goal, not the turn + +The one rule worth stating on its own, because getting it wrong stops work that would have +finished. An earlier version told the judge to answer `blocked` when "the loop is going in +circles". That mapped a **turn-level** observation onto a **goal-level** verdict, and the +failure mode was concrete: an agent waiting on a five-minute build spends each turn reading +an unchanged log, and the judge — seeing a turn that changed nothing — declared the whole +goal unachievable and paused. The work was not unachievable; the agent was waiting. + +So repetition is now explicitly **not** a reason to answer `blocked`, however many turns it +has happened, because repetition is recoverable: the next turn can do something different. A +turn that checks on work already in flight is **waiting**, not circling, and the judge is told +so and told not to assume work did *not* happen off-screen. + +### The polling guard + +The stall and repetition guards miss a specific shape. An agent waiting on a long job says +something different every turn (so the reply digest moves) and calls a tool every turn (so +the tool count is non-zero), while learning nothing. Only the model noticed, and it reached +for the terminal verdict. + +So the loop also digests **what the tool calls reported**, ignoring the commands that +produced them, and pauses after `pollLimit` turns whose observation is unchanged. Its message +says the work is not moving, says to wait for a running command rather than re-read it, and +says plainly that this is not a statement about whether the goal is reachable. + +`pollLimit` is configurable exactly like `stallLimit` — `pollLimit` in `opencode.json`, or +`/goal poll ` for the session — because the right threshold is a property of the +work: a build that takes five minutes wants a higher limit than a test that takes five +seconds, and a hardcoded 3 was the one number here a user could not argue with. + +It fails open: an unrecognised tool-state shape yields an empty digest and the guard stays +quiet, because a heuristic that pauses a loop on a guess is worse than one that misses. + ## Caveats - **The judge is only as good as its model.** A weak judge is permissive. Set `judgeModel` to something small and sharp, and expect to use `/goal status` to sanity-check its verdicts. +- **A `blocked` verdict is worth reading twice.** It means the goal looks unreachable, which + is a strong claim resting on one model call. If the agent was mid-way through a long + command, a large batch, or a build, the more likely story is that it was waiting and the + judge read a quiet turn as a dead end. `/goal resume` costs nothing but the turn. - **Your agent model must actually use tools.** If the session model narrates intentions without calling tools, the stall guard fires after two turns. That is the guard working, but it means the goal will not get done. diff --git a/src/display.ts b/src/display.ts index a99d852..253727f 100644 --- a/src/display.ts +++ b/src/display.ts @@ -99,3 +99,17 @@ export function applyAll(enabled: boolean): Display { return { panel: enabled, footer: enabled, composer: enabled, sidebar: enabled } } +/** + * Apply a mutation to a display and return the complete next value. + * + * The UI shows all four placements but only ever changes the named ones, so persisting + * has to write the *whole* value rather than the keys the caller happened to touch. + * Returning a fresh object rather than editing in place also gives the UI something new + * to compare, which is what makes a change visible on the frame it happens. + */ +export function applyMutation(current: Display, mutate: (draft: Display) => void): Display { + const next: Display = { ...current } + mutate(next) + return { panel: next.panel, footer: next.footer, composer: next.composer, sidebar: next.sidebar } +} + diff --git a/src/listen.ts b/src/listen.ts new file mode 100644 index 0000000..4b8cc56 --- /dev/null +++ b/src/listen.ts @@ -0,0 +1,116 @@ +/** + * A subscription that survives. + * + * The client's own `rpc.events.on` is one `for await` loop with a catch at the + * end, and nothing that restarts it. Measured against the real library: + * + * - a handler that throws once ends the loop, so every later event is dropped + * and the error goes to `console.error`; + * - a stream that simply ends ends the loop with no error at all, and nothing + * opens a second connection. + * + * Either way the listener stays registered and looks alive while being + * permanently deaf, and the UI it feeds stops updating until something else + * refreshes it - here, reopening the panel, which reads state directly. + * + * `listen` fixes both. A handler error is reported and the next event is still + * delivered; when the stream ends or fails it is reopened after a growing delay, + * and `onResume` runs first so the caller can re-read whatever it missed while + * nothing was listening. + * + * No dependency on the client, Solid or JSX, so it can be tested with plain + * async iterables. + */ + +export type Stop = () => void + +/** + * Where a failure came from. A dropped stream is routine and heals itself, while + * a throwing handler is a bug, so callers usually want to treat them differently. + */ +export type Phase = "handler" | "stream" | "resume" + +export interface ListenOptions { + /** + * Runs after a gap in delivery, before the stream is reopened. Events sent + * during the gap are gone, so this is where the caller re-reads state. It is + * not called before the first connection. + */ + onResume?: () => void | Promise + /** Told about a failure and where it came from. May throw; that is contained. */ + onError?: (error: unknown, phase: Phase) => void + /** Injectable so tests need not wait. */ + sleep?: (ms: number, signal: AbortSignal) => Promise + /** First reconnect delay. Default 500ms. */ + minDelay?: number + /** Ceiling for the doubling delay. Default 10s. */ + maxDelay?: number +} + +const sleepUnlessAborted = (ms: number, signal: AbortSignal) => + new Promise((resolve) => { + if (signal.aborted) return resolve() + const finish = () => { + clearTimeout(timer) + signal.removeEventListener("abort", finish) + resolve() + } + const timer = setTimeout(finish, ms) + signal.addEventListener("abort", finish, { once: true }) + }) + +export function listen( + subscribe: (signal: AbortSignal) => AsyncIterable, + handler: (event: E) => void | Promise, + options: ListenOptions = {}, +): Stop { + const controller = new AbortController() + const { signal } = controller + const sleep = options.sleep ?? sleepUnlessAborted + const minDelay = options.minDelay ?? 500 + const maxDelay = options.maxDelay ?? 10_000 + const report = (error: unknown, phase: Phase) => { + try { + options.onError?.(error, phase) + } catch { + // A broken error sink must not take the listener down with it. + } + } + + void (async () => { + let delay = minDelay + let reconnecting = false + + while (!signal.aborted) { + if (reconnecting) { + await sleep(delay, signal) + if (signal.aborted) return + delay = Math.min(delay * 2, maxDelay) + try { + await options.onResume?.() + } catch (error) { + report(error, "resume") + } + if (signal.aborted) return + } + reconnecting = true + + try { + for await (const event of subscribe(signal)) { + if (signal.aborted) return + // Something arrived, so the connection is healthy: forget the backoff. + delay = minDelay + try { + await handler(event) + } catch (error) { + report(error, "handler") + } + } + } catch (error) { + if (!signal.aborted) report(error, "stream") + } + } + })() + + return () => controller.abort() +} diff --git a/src/observation.ts b/src/observation.ts new file mode 100644 index 0000000..da4b370 --- /dev/null +++ b/src/observation.ts @@ -0,0 +1,55 @@ +/** + * Turning a turn's tool calls into one stable observation string. + * + * The point of this is narrow and specific: an agent waiting on a long build produces a + * turn that calls tools, says something new, and learns nothing. That defeats both of the + * loop's existing no-progress checks — the reply digest moves, so the reply-repetition + * guard never fires, and the tool count is non-zero, so the stall guard never fires. Only + * the judge noticed, and its answer was to declare the *goal* unachievable from + * turn-level evidence. + * + * So this digests what the tool calls *reported*, not what the agent said or asked for. + * Two deliberate choices: + * + * - Tool inputs are skipped. An agent that reworded the same command each turn was still + * reading the same unchanged result, and that is precisely the case worth catching. + * - It fails open. An unrecognised state shape yields an empty digest, and the guard + * stays quiet, because a heuristic that pauses a loop on a guess is worse than one + * that misses. + */ + +/** Keys that describe *how* the tool was called rather than what it found. */ +const CALL_SHAPE = new Set(["input", "args", "description", "title", "callID", "id"]) + +/** Everything a turn's tool calls reported, as one stable string. */ +export function observationOf(messages: readonly unknown[]): string { + const parts: string[] = [] + const walk = (value: unknown, depth = 0): void => { + if (depth > 6 || value === null || value === undefined) return + if (typeof value === "string") { + parts.push(value) + return + } + if (typeof value === "number" || typeof value === "boolean") { + parts.push(String(value)) + return + } + if (Array.isArray(value)) { + for (const item of value) walk(item, depth + 1) + return + } + if (typeof value === "object") + for (const [key, nested] of Object.entries(value as Record)) + if (!CALL_SHAPE.has(key)) walk(nested, depth + 1) + } + + for (const message of messages) { + const record = message as { type?: unknown; content?: unknown } | null + if (record?.type !== "assistant") continue + for (const part of (record.content ?? []) as readonly unknown[]) { + const entry = part as { type?: unknown; state?: unknown } | null + if (entry?.type === "tool") walk(entry.state) + } + } + return parts.join("\n") +} diff --git a/src/rpc.ts b/src/rpc.ts index 1bd80ab..e5bff5e 100644 --- a/src/rpc.ts +++ b/src/rpc.ts @@ -35,6 +35,7 @@ export const HELP_TEXT = `/goal set the goal and start working /goal budget inf no turn limit (also: unlim, unlimited, none, ∞) /goal budget default back to the configured maxTurns /goal stall turns with no tools before the loop gives up +/goal poll turns re-reading an unchanged result before it counts as polling /goal quiet whether the panel replaces the transcript notices /goal judge provider/model[#variant] used to judge each turn /goal settings show every setting, and where it came from @@ -57,6 +58,9 @@ export type GoalView = { maxTurns: number | null stalled: number repeats: number + /** Turns that re-read an unchanged result. Optional: added after v1 shipped, so it is + * deliberately not in `required` and older views simply do not carry it. */ + observing?: number /** The judge's most recent one-sentence reason, or "". */ reason: string /** The goal's own proof condition, or "". */ @@ -74,6 +78,7 @@ const state = { maxTurns: { anyOf: [{ type: "number" }, { type: "null" }] }, stalled: { type: "number" }, repeats: { type: "number" }, + observing: { type: "number" }, reason: { type: "string" }, verification: { type: "string" }, updatedAt: { type: "number" }, diff --git a/src/server.ts b/src/server.ts index 3f12bae..da3fe0d 100644 --- a/src/server.ts +++ b/src/server.ts @@ -1,6 +1,8 @@ import { realpathSync } from "node:fs" import { Plugin } from "@opencode/plugin" import { Goal, HELP_TEXT, type GoalView } from "./rpc.js" +import { observationOf } from "./observation.js" +import { backoffSeconds, MAX_WAIT_MS, outputPathOf, pendingBackground, waitForBackground } from "./waiting.js" import { budgetLabel, formatTurns, @@ -12,6 +14,7 @@ import { type Budget, } from "./budget.js" import { + DEFAULT_POLL, DEFAULT_STALL, modelLabel, NO_OVERRIDES, @@ -83,6 +86,15 @@ type GoalState = { /** Digest of the previous turn's reply, to spot a loop the judge may miss. */ lastDigest?: string repeats: number + /** Digest of the previous turn's tool results, to spot polling that changes nothing. + * Distinct from `lastDigest`: an agent can vary its prose every turn and still be + * reading the same unchanged output, which is how waiting gets mistaken for spinning. */ + lastObservation?: string + observing: number + /** Consecutive turns spent waiting on background work, and the time that cost. While + * a background command is unfinished the no-progress guards stand down, within a bound. */ + waits?: number + waitedMs?: number reason?: string } @@ -169,6 +181,7 @@ function toView(state: GoalState): Omit { maxTurns: state.maxTurns, stalled: state.stalled ?? 0, repeats: state.repeats ?? 0, + observing: state.observing ?? 0, reason: state.reason ?? "", verification: state.contract.verification ?? "", updatedAt: Date.now(), @@ -185,6 +198,7 @@ const JUDGE_PROMPT = `You are a strict completion judge for an autonomous coding Turn {{TURN}} of at most {{MAX_TURNS}} have been spent on this goal. Tool calls made in the turn you are judging: {{TOOLCALLS}} Consecutive turns that changed nothing at all: {{STALLED}} +Consecutive turns that repeated one observation without a result changing: {{OBSERVING}} Your verdict on the previous turn was: "continue", because: {{PREVIOUS}} @@ -198,10 +212,12 @@ Reply with exactly one line of strict JSON and nothing else: Rules: - "done" ONLY when the response carries concrete evidence the whole goal is satisfied: a command that passed along with its output, files that were actually created or changed, a test suite that is green. A claim, a plan, an intention, or "I will now..." is never done. If the goal's own proof condition is named above, that specific proof must be present. -- "blocked" when the goal cannot be concluded as written. That covers a goal that is impossible, self-contradictory, outside the repository's scope, or dependent on credentials, hardware, or decisions the agent does not have. It also covers a goal asking for something no amount of the agent's work can produce, even when the agent has not admitted that yet. -- "blocked" when the loop is going in circles: this turn repeats the previous turn's promise, edit, or failing command without new information, or the turn changed nothing while earlier turns did not either. Judge the loop state above, not just the prose. -- "continue" when real, non-repeating progress is still possible. Slow is fine. Stalled is not. -- Judge only what is in front of you. Do not assume work happened off-screen.` +- "blocked" ONLY when the GOAL cannot be concluded as written: impossible, self-contradictory, outside the repository's scope, or dependent on credentials, hardware, or decisions the agent does not have. It also covers a goal asking for something no amount of the agent's work can produce, even when the agent has not admitted that yet. +- "blocked" requires goal-level evidence. A turn that was unproductive, slow, repeated, or made no change is evidence about THIS TURN, never about whether the goal is reachable. Repetition is recoverable — the next turn can do something different — so repetition alone must never produce "blocked", however many turns it has happened. +- "continue" when real progress is still possible, AND when this turn was unproductive. Slow is fine. Waiting is fine. Stalled is not. +- A turn that checked on work already in flight — a background build, a long test run, a command started earlier — is WAITING, not going in circles, unless the underlying work shows no progress at all across several turns. Do not read a repeated read-only check as a dead goal. +- If the loop state above suggests repetition, the correct verdict is still "continue": say in the reason what the agent should change (do different work, or stop polling and wait for the result), and let the next turn act on it. +- Judge only what is in front of you. Do not assume work happened off-screen. Equally, do not assume work did NOT happen off-screen: a long-running command may still be in flight.` function continuationPrompt(state: GoalState, reason: string, quiet: boolean): string { // With a TUI watching, the panel already shows which turn this is, so the @@ -215,7 +231,7 @@ function continuationPrompt(state: GoalState, reason: string, quiet: boolean): s Judge's note: ${reason} -The goal is not met yet. Take the next concrete step. Do not repeat a command that just failed or an edit you just made unless something has changed. If the goal cannot be completed as written, say so plainly and name the blocker — that pauses the loop instead of burning the remaining turns.` +The goal is not met yet. Take the next concrete step. Do not repeat a command that just failed or an edit you just made unless something has changed. If a background command you started is still running and there is nothing else to do, call goal_wait instead of replying with a status update. If the goal cannot be completed as written, say so plainly and name the blocker — that pauses the loop instead of burning the remaining turns.` } function statusReport(state: GoalState | undefined): string { @@ -235,7 +251,7 @@ function statusReport(state: GoalState | undefined): string { async function lastAssistantTurn( ctx: any, sessionID: string, -): Promise<{ text: string; model?: ModelRef; toolCalls: number }> { +): Promise<{ text: string; model?: ModelRef; toolCalls: number; observation: string; pending: number }> { const messages = (await ctx.session.context({ sessionID })) as readonly any[] // A turn is every assistant message after the last user or synthetic input. @@ -251,10 +267,11 @@ async function lastAssistantTurn( } } + const turn = messages.slice(start) let model: ModelRef | undefined let toolCalls = 0 let text = "" - for (const message of messages.slice(start)) { + for (const message of turn) { if (message?.type !== "assistant") continue model ??= message.model as ModelRef | undefined const content = (message.content ?? []) as readonly any[] @@ -266,7 +283,13 @@ async function lastAssistantTurn( .trim() if (said) text = said } - return { text, model, toolCalls } + return { + text, + model, + toolCalls, + observation: digest(observationOf(turn)), + pending: pendingBackground(messages).length, + } } function parseVerdict(raw: string): { verdict: Verdict; reason: string } { @@ -283,6 +306,20 @@ function parseVerdict(raw: string): { verdict: Verdict; reason: string } { return { verdict: parsed.verdict, reason: parsed.reason ?? "" } } +/** True while some process has the file open. Unknown (no lsof) counts as still running. */ +async function holdsOpen(path: string): Promise { + try { + const { execFile } = await import("node:child_process") + return await new Promise((resolve) => + execFile("lsof", ["-t", path], (error, stdout) => + resolve(error && (error as { code?: unknown }).code === "ENOENT" ? true : stdout.trim().length > 0), + ), + ) + } catch { + return true + } +} + export default Plugin.define({ id: "goal", async setup(ctx) { @@ -295,6 +332,11 @@ export default Plugin.define({ typeof ctx.options.stallLimit === "number" && ctx.options.stallLimit >= 1 ? ctx.options.stallLimit : DEFAULT_STALL + // Same shape as the stall limit: a whole number of at least 1, else the default. + const configuredPoll = + typeof ctx.options.pollLimit === "number" && ctx.options.pollLimit >= 1 + ? ctx.options.pollLimit + : DEFAULT_POLL const configuredJudge = (ctx.options.judgeModel ?? null) as ModelRef | null /** * Whether to drop the loop's reporting from the transcript. Off by default: @@ -323,7 +365,7 @@ export default Plugin.define({ /** * All four settings in one record, so the summary, the setters and the loop * itself cannot disagree about what is in force. `settings:` holds - * `{ maxTurns?, stall?, quiet?, judge? }` with only the overridden keys + * `{ maxTurns?, stall?, poll?, quiet?, judge? }` with only the overridden keys * present; a key that is `null` means the plugin option applies. */ const settingsKey = (sessionID: string) => `settings:${sessionID}` @@ -331,6 +373,7 @@ export default Plugin.define({ const configured = { maxTurns: configuredBudget, stall: configuredStall, + poll: configuredPoll, quiet: configuredQuiet, judge: configuredJudge, } @@ -349,6 +392,7 @@ export default Plugin.define({ // "unlimited", not an absent value. Absent keys stay undefined. if ("maxTurns" in stored) overrides.maxTurns = (stored.maxTurns as Budget | null) ?? null if (typeof stored.stall === "number") overrides.stall = stored.stall + if (typeof stored.poll === "number") overrides.poll = stored.poll if (typeof stored.quiet === "boolean") overrides.quiet = stored.quiet // The judge has no third value the way maxTurns does: not set is // exactly "use the session's model", so absent is the only unset form. @@ -368,7 +412,9 @@ export default Plugin.define({ const writeOverride = async ( sessionID: string, - patch: Partial>, + patch: Partial< + Record<"maxTurns" | "stall" | "poll" | "quiet" | "judge", Budget | number | boolean | ModelRef | null> + >, ) => { const stored = (await ctx.storage.get(settingsKey(sessionID))) as Record | undefined const next: Record = { ...(stored ?? {}) } @@ -615,12 +661,71 @@ export default Plugin.define({ .replace("{{MAX_TURNS}}", budgetLabel(state.maxTurns)) .replace("{{TOOLCALLS}}", String(toolCalls)) .replace("{{STALLED}}", String(state.stalled ?? 0)) + .replace("{{OBSERVING}}", String(state.observing ?? 0)) .replace("{{PREVIOUS}}", state.reason || "(this is the first turn)") .replace("{{RESPONSE}}", text.slice(-4000) || "(the agent produced no text this turn)") const result = await ctx.generate.text({ model: chosen, prompt }) return parseVerdict(result.text) } + // Lets the agent say "I am waiting" instead of ending the turn with a status reply. + // The call blocks, so the turn stays open and no repeat or stall can form. + await ctx.tool.transform((editor) => { + editor.add({ + name: "goal_wait", + description: + "While a standing goal is active and a background command you started is still running, " + + "call this instead of replying with a status update. It returns when the command " + + "finishes, or after `seconds`, whichever comes first.", + input: { + type: "object", + properties: { + seconds: { type: "number", description: "Longest to wait, 1 to 3600. Default 600." }, + }, + }, + execute: async (raw: unknown, context) => { + const input = (raw ?? {}) as { seconds?: number } + const sessionID = context.sessionID as string + const state = await read(sessionID) + if (state?.status !== "active") { + return { content: "No goal is active, so there is nothing to wait for." } + } + const budget = Math.max(0, MAX_WAIT_MS - (state.waitedMs ?? 0)) / 1000 + const seconds = Math.min(Math.max(1, Number(input?.seconds) || 600), 3600, budget) + if (seconds <= 0) return { content: "The waiting allowance for this goal is used up." } + const { end, waitedMs } = await waitForBackground({ + seconds, + intervalMs: 5000, + signal: context.signal, + // The completion notice is only delivered once the turn ends, so inside a + // turn it never shows in the context. Whether the process still holds its + // output file open is the signal that is visible now. + pending: async () => { + const messages = (await ctx.session.context({ sessionID })) as readonly unknown[] + let alive = 0 + for (const id of pendingBackground(messages)) { + const out = outputPathOf(messages, id) + if (!out || (await holdsOpen(out))) alive++ + } + return alive + }, + active: async () => (await read(sessionID))?.status === "active", + }) + const current = await read(sessionID) + if (current) await write(sessionID, { ...current, waitedMs: (current.waitedMs ?? 0) + waitedMs }) + const minutes = Math.round(waitedMs / 6000) / 10 + return { + content: + end === "finished" + ? `Background work finished after ${minutes} min. Continue with its result.` + : end === "timeout" + ? `Still running after ${minutes} min. Call goal_wait again, or do other work.` + : `Wait ended early (${end}).`, + } + }, + }) + }) + const command = await ctx.command.transform((editor) => { editor.add({ name: "goal", @@ -710,7 +815,9 @@ export default Plugin.define({ turns: 0, stalled: 0, repeats: 0, + observing: 0, lastDigest: undefined, + lastObservation: undefined, reason: undefined, } await write(sessionID, resumed) @@ -770,7 +877,7 @@ export default Plugin.define({ return } - if (sub === "stall" || sub === "quiet" || sub === "judge" || sub === "settings") { + if (sub === "stall" || sub === "poll" || sub === "quiet" || sub === "judge" || sub === "settings") { if (sub === "settings") { const summary = renderSettings(await settingsFor(sessionID)) if (await tuiHere(sessionID)) { @@ -805,6 +912,28 @@ export default Plugin.define({ ) return } + if (sub === "poll") { + const parsed = parseCount(argument, "poll") + if (parsed.kind === "invalid") { + await announce( + sessionID, + "Poll limit", + "Give me a whole number of observations, or default. For example: /goal poll 5.", + "error", + ) + return + } + const next = parsed.kind === "clear" ? configuredPoll : parsed.value + await writeOverride(sessionID, { poll: parsed.kind === "clear" ? undefined : next }) + await announce( + sessionID, + "Poll limit", + parsed.kind === "clear" + ? `reset to the configured default (${next})` + : `${next} turns with an unchanged result`, + ) + return + } if (sub === "quiet") { const parsed = parseFlag(argument) if (parsed.kind === "invalid") { @@ -898,7 +1027,9 @@ export default Plugin.define({ maxTurns: budget, stalled: 0, repeats: 0, + observing: 0, lastDigest: undefined, + lastObservation: undefined, reason: undefined, } await write(sessionID, state) @@ -976,23 +1107,47 @@ export default Plugin.define({ // Two deterministic no-progress checks. These, not the judge, are what // actually stop a runaway: a weak judge model will happily answer // "continue" to twenty identical replies. - const { toolCalls, text: replyText } = await lastAssistantTurn(ctx, sessionID) + const { toolCalls, text: replyText, observation, pending } = await lastAssistantTurn(ctx, sessionID) const reply = normalize(replyText) - const repeated = reply.length > 0 && state.lastDigest === digest(reply) + // A background command still running makes a quiet turn legitimate: there is + // nothing to do but wait. The guards stand down for it, bounded by MAX_WAIT_MS so + // a hung job cannot hold the loop open forever. + const waiting = pending > 0 && (state.waitedMs ?? 0) < MAX_WAIT_MS + const waitSeconds = waiting ? backoffSeconds(state.waits ?? 0) : 0 + const repeated = !waiting && reply.length > 0 && state.lastDigest === digest(reply) const repeats = repeated ? (state.repeats ?? 0) + 1 : 0 - const stalled = toolCalls === 0 ? (state.stalled ?? 0) + 1 : 0 + const stalled = waiting ? 0 : toolCalls === 0 ? (state.stalled ?? 0) + 1 : 0 + // The third no-progress shape: the agent used tools and said something new each + // turn, but read back the same unchanged result — polling. Neither guard above can + // see it, because the reply digest moves and the tool count is non-zero. It is + // counted separately and, on the first turn, only *reported* to the judge, so the + // agent is told to stop polling before anything is stopped for it. + const sameObservation = + observation.length > 0 && state.lastObservation === observation + const observing = !waiting && sameObservation ? (state.observing ?? 0) + 1 : 0 const looping: GoalState = { ...progressed, lastDigest: reply.length ? digest(reply) : state.lastDigest, + lastObservation: observation.length ? observation : state.lastObservation, repeats, stalled, + observing, + waits: waiting ? (state.waits ?? 0) + 1 : 0, + waitedMs: (state.waitedMs ?? 0) + waitSeconds * 1000, } if (repeats >= 2) { await write(sessionID, { ...looping, status: "paused" }) return terminal( sessionID, - `⏸ Goal paused — looping: the agent has now repeated the same reply ${repeats} turns running without acting. Judge said: ${verdict.reason}\nThe goal is not reachable in this session. Re-scope it with /goal .`, + `⏸ Goal paused — looping: the agent has now repeated the same reply ${repeats} turns running without acting. Judge said: ${verdict.reason}\nThis loop is not making progress. Either resume with /goal resume, or re-scope it with /goal .`, + ) + } + if (observing >= (await settingsFor(sessionID)).poll) { + await write(sessionID, { ...looping, status: "paused" }) + return terminal( + sessionID, + `⏸ Goal paused — polling: ${observing} turns ran tools and read back the same unchanged result, so the work is not moving.\nIf a command is still running, wait for its completion instead of re-reading it; otherwise do different work, then /goal resume. This says nothing about whether the goal is reachable.`, ) } if (stalled >= (await settingsFor(sessionID)).stall) { @@ -1015,6 +1170,11 @@ export default Plugin.define({ const next: GoalState = { ...looping, turns: state.turns + 1 } await write(sessionID, next) + if (waitSeconds > 0) { + await new Promise((resolve) => setTimeout(resolve, waitSeconds * 1000)) + // Someone may have paused, cleared or replaced the goal while we waited. + if ((await read(sessionID))?.status !== "active") continue + } // No status message here on purpose: a synthetic inbox message is a // real prompt, so it would start another execution, re-enter this // loop, and burn the budget twice as fast. The continuation prompt diff --git a/src/settings.ts b/src/settings.ts index 7c480cc..fde2711 100644 --- a/src/settings.ts +++ b/src/settings.ts @@ -1,6 +1,6 @@ /** - * Session settings: the turn budget, the stall limit, the quiet flag and the - * judge model. + * Session settings: the turn budget, the stall limit, the polling limit, the + * quiet flag and the judge model. * * Each has a value in the plugin options, which is the default, and may be * overridden for one session with a `/goal …` subcommand. Pure functions only, @@ -29,6 +29,7 @@ export type ModelRef = { providerID: string; id: string; variant?: string } export type Overrides = { maxTurns: Budget | undefined stall: number | undefined + poll: number | undefined quiet: boolean | undefined judge: ModelRef | undefined } @@ -36,6 +37,7 @@ export type Overrides = { export const NO_OVERRIDES: Overrides = { maxTurns: undefined, stall: undefined, + poll: undefined, quiet: undefined, judge: undefined, } @@ -44,15 +46,18 @@ export const NO_OVERRIDES: Overrides = { export type Effective = { maxTurns: Budget stall: number + poll: number quiet: boolean judge: ModelRef | null - overridden: { maxTurns: boolean; stall: boolean; quiet: boolean; judge: boolean } + overridden: { maxTurns: boolean; stall: boolean; poll: boolean; quiet: boolean; judge: boolean } } export const DEFAULT_STALL = 2 +/** The polling guard waits for three unchanged observations before it pauses. */ +export const DEFAULT_POLL = 3 export function resolve( - options: { maxTurns: Budget; stall: number; quiet: boolean; judge: ModelRef | null }, + options: { maxTurns: Budget; stall: number; poll: number; quiet: boolean; judge: ModelRef | null }, overrides: Overrides, ): Effective { const pick = (override: T | undefined, fallback: T): T => @@ -60,11 +65,13 @@ export function resolve( return { maxTurns: pick(overrides.maxTurns, options.maxTurns), stall: pick(overrides.stall, options.stall), + poll: pick(overrides.poll, options.poll), quiet: pick(overrides.quiet, options.quiet), judge: pick(overrides.judge, options.judge), overridden: { maxTurns: overrides.maxTurns !== undefined, stall: overrides.stall !== undefined, + poll: overrides.poll !== undefined, quiet: overrides.quiet !== undefined, judge: overrides.judge !== undefined, }, @@ -144,17 +151,19 @@ export function renderSettings(state: Effective): string { return [ `Turn budget ${budgetText} (${mark(state.overridden.maxTurns)})`, `Stall limit ${state.stall} turns with no tools (${mark(state.overridden.stall)})`, + `Poll limit ${state.poll} unchanged observations (${mark(state.overridden.poll)})`, `Quiet mode ${state.quiet ? "on" : "off"} (${mark(state.overridden.quiet)})`, `Judge model ${modelLabel(state.judge)} (${mark(state.overridden.judge)})`, ``, `Change any of these for this session:`, ` /goal budget `, ` /goal stall `, + ` /goal poll `, ` /goal quiet `, ` /goal judge `, ``, `Session overrides last until this session ends. The config default is`, - `whatever maxTurns, stallLimit, quiet and judgeModel are set to in`, + `whatever maxTurns, stallLimit, pollLimit, quiet and judgeModel are set to in`, `opencode.json. Panel and display are TUI-only; the rest work anywhere.`, ].join("\n") } diff --git a/src/tui.tsx b/src/tui.tsx index dc1a3e0..8eb3507 100644 --- a/src/tui.tsx +++ b/src/tui.tsx @@ -2,8 +2,10 @@ import { Plugin } from "@opencode/plugin/tui" import { createEffect, createSignal, onCleanup, Show } from "solid-js" import { Goal, HELP_TEXT as HELP, type GoalView } from "./rpc.js" import { formatTurns } from "./budget.js" +import { listen } from "./listen.js" import { applyAll, + applyMutation, applyParsed, applySet, describe, @@ -54,9 +56,33 @@ export default Plugin.define({ setup(context) { const rpc = context.client.rpc(Goal) const [states, setStates] = createSignal>({}) - const [display, setDisplay] = context.storage.store("display", { + /** + * Where the goal is shown. + * + * `storage.store` is the durable record, but it is not reactive: a slot's `render` + * runs once, so `` read a snapshot and never re-read it. + * The result was that a placement change only appeared after something forced the + * slot to re-render — in practice, toggling the panel with `/goal panel`, which is + * exactly the manual workaround this is meant to remove. + * + * So the stored value is mirrored into a signal, the same mechanism the goal state + * already uses and which is known to re-render here. The store stays the single + * durable record; the signal is only what the UI reads. + */ + const [storedDisplay, setStoredDisplay] = context.storage.store("display", { initial: seedDisplay(context.options as Record), }) + const asDisplay = (value: unknown): Display => seedDisplay(value) + const [display, setDisplayNow] = createSignal(asDisplay(storedDisplay)) + /** Update the signal immediately, then persist. The signal leads so the change is + * visible on this frame rather than after the storage round trip. */ + const setDisplay = async (mutate: (draft: Display) => void): Promise => { + const next = applyMutation(asDisplay(storedDisplay), mutate) + setDisplayNow(next) + await setStoredDisplay((draft: Display) => { + for (const key of PLACEMENTS) draft[key] = next[key] + }) + } /** Last terminal state toasted per session, so a repaint cannot repeat it. */ const announced = new Map() @@ -151,17 +177,21 @@ export default Plugin.define({ const fingerprint = `${state.status}:${state.turns}:${state.reason}` if (announced.get(state.sessionID) === fingerprint) return announced.set(state.sessionID, fingerprint) - context.ui.toast.show({ - title: - state.status === "done" - ? "Goal achieved" - : state.status === "blocked" - ? "Goal judged unachievable" - : "Goal paused", - message: state.reason || state.goal, - variant: state.status === "done" ? "success" : "warning", - duration: 8000, - }) + try { + context.ui.toast.show({ + title: + state.status === "done" + ? "Goal achieved" + : state.status === "blocked" + ? "Goal judged unachievable" + : "Goal paused", + message: state.reason || state.goal, + variant: state.status === "done" ? "success" : "warning", + duration: 8000, + }) + } catch { + // The state is already tracked; losing the toast is not worth the listener. + } } /** Events arrive for every location, so ignore other projects entirely. */ @@ -172,6 +202,105 @@ export default Plugin.define({ return trim(directory) === trim(mine) } + /** + * Events are the fast path, but they are only pushed, so anything that missed + * them has to ask. Two things did: a session opened after its goal started + * (the footer, composer and sidebar never fetched anything, only the panel + * did), and a listener that had silently stopped hearing. + */ + const fingerprint = (state: GoalView) => `${state.status}:${state.turns}:${state.reason}` + const sameView = (a: GoalView, b: GoalView) => + JSON.stringify({ ...a, updatedAt: 0 }) === JSON.stringify({ ...b, updatedAt: 0 }) + + /** Read a session's state straight from the server. */ + const refresh = async (sessionID: string) => { + try { + const result = asGet(await rpc.get({ sessionID })) + const before = states()[sessionID] + if (!result.state) { + if (before) { + announced.delete(sessionID) + track(sessionID, null) + } + return + } + // updatedAt is stamped when the view is built, so it always differs; without + // this every poll would replace the state and repaint for nothing. + if (before && sameView(before, result.state)) return + track(sessionID, result.state) + if (result.state.status === "active") return + if (before?.status === "active") { + // It ended while nothing was listening. Say so, once. + announce(result.state) + } else { + // Already over when first seen. Record that, so a later event carrying + // the same result does not announce a goal that finished long ago. + announced.set(sessionID, fingerprint(result.state)) + } + } catch { + // Server unreachable; the next attempt tries again. + } + } + + /** Sessions something is showing, so a resync knows what to re-read. */ + const fetched = new Set() + const ensureFetched = (sessionID: string) => { + if (fetched.has(sessionID)) return + fetched.add(sessionID) + void refresh(sessionID) + } + const resync = async () => { + const known = new Set([...fetched, ...Object.keys(states())]) + await Promise.all([...known].map(refresh)) + } + + /** Backstop for a half-open stream, which reports no error to reconnect on. */ + const pollTimer = setInterval(() => { + for (const [sessionID, state] of Object.entries(states())) { + if (state.status === "active") void refresh(sessionID) + } + }, 5000) + + type EventName = "changed" | "panel" | "display" | "notice" | "settings" | "help" + type RpcEvent = { readonly data: unknown; readonly location?: { readonly directory?: string } } + + /** + * A handler that throws is a bug worth seeing, so it gets a toast, but never + * more than one per half minute. A dropped stream is routine and stays quiet. + */ + let lastHandlerErrorAt = 0 + const reportListenError = (error: unknown, phase: "handler" | "stream" | "resume") => { + if (phase !== "handler") return + if (Date.now() - lastHandlerErrorAt < 30_000) return + lastHandlerErrorAt = Date.now() + try { + context.ui.toast.show({ + title: "Goal plugin", + message: (error instanceof Error ? error.message : String(error)).slice(0, 120), + variant: "error", + duration: 6000, + }) + } catch { + // Nowhere left to report to. + } + } + + /** + * Subscribe to a plugin event. Not `rpc.events.on`: that ends for good the + * first time a handler throws or the stream closes, and says nothing, which + * is how the areas stopped updating until the panel was reopened. + */ + const on = ( + name: EventName, + handler: (event: RpcEvent) => void | Promise, + onResume?: () => void | Promise, + ) => + listen( + (signal) => rpc.events.subscribe(name, { signal }) as AsyncIterable, + handler, + { onResume, onError: reportListenError }, + ) + // Each registration is guarded on its own. A throw anywhere in setup // discards the entire plugin, so one bad call would otherwise cost the // panel, the toasts and the live updates together. @@ -185,7 +314,7 @@ export default Plugin.define({ } guard("rpc events", () => - rpc.events.on("changed", (event) => { + on("changed", (event) => { if (!isLocal(event.location?.directory)) { return } @@ -197,11 +326,18 @@ export default Plugin.define({ } track(sessionID, state) announce(state) - }), + }, resync), ) /** One line, for the slots with no room for the full view. */ const Compact = (props: { sessionID?: string }) => { + // Without this the footer, composer and sidebar showed nothing for a goal + // already running when the session was opened: they only ever waited for an + // event, and no event was coming. + createEffect(() => { + const sessionID = props.sessionID + if (sessionID) ensureFetched(sessionID) + }) const view = () => (props.sessionID ? states()[props.sessionID] : undefined) const word = (status: GoalView["status"]) => status === "active" ? "running" : status === "done" ? "achieved" : status === "blocked" ? "unachievable" : "paused" @@ -227,7 +363,7 @@ export default Plugin.define({ context.ui.slot({ append: "session.panel", render: (panel) => ( - + ), @@ -240,7 +376,7 @@ export default Plugin.define({ context.ui.slot({ append: "prompt.footer.status", render: (input) => ( - + ), @@ -251,7 +387,7 @@ export default Plugin.define({ context.ui.slot({ append: "session.composer.top", render: (input) => ( - + ), @@ -262,7 +398,7 @@ export default Plugin.define({ context.ui.slot({ append: "sidebar.footer", render: (input) => ( - + ), @@ -271,19 +407,9 @@ export default Plugin.define({ /** Pull state for a session the panel opened before any event arrived. */ const prime = async (sessionID: string) => { - try { - const result = asGet(await rpc.get({ sessionID })) - if (result.state) { - track(sessionID, result.state) - announced.set( - sessionID, - `${result.state.status}:${result.state.turns}:${result.state.reason}`, - ) - void attach(sessionID) - } - } catch { - // Nothing to show yet. - } + fetched.add(sessionID) + await refresh(sessionID) + if (states()[sessionID]) void attach(sessionID) } const Body = (props: { view: GoalView; width?: number }) => { @@ -382,7 +508,7 @@ export default Plugin.define({ } guard("panel request", () => - rpc.events.on("panel", (event) => { + on("panel", (event) => { if (!isLocal(event.location?.directory)) return try { if (context.ui.panel.current()?.name === PANEL) context.ui.panel.close() @@ -412,7 +538,7 @@ export default Plugin.define({ title: "Where should the goal be shown?", options: [ ...PLACEMENTS.map((key) => ({ - title: `${display[key] ? "on " : "off"} ${key}`, + title: `${display()[key] ? "on " : "off"} ${key}`, value: key, description: PLACEMENT_HELP[key], })), @@ -425,7 +551,7 @@ export default Plugin.define({ }) if (choice === undefined) return const next = - choice === "__off" ? applyAll(false) : applySet(display, choice, !display[choice]) + choice === "__off" ? applyAll(false) : applySet(display(), choice, !display()[choice]) await setDisplay((draft) => { for (const key of PLACEMENTS) draft[key] = next[key] }) @@ -437,7 +563,7 @@ export default Plugin.define({ } guard("notice request", () => - rpc.events.on("notice", (event) => { + on("notice", (event) => { if (!isLocal(event.location?.directory)) return const { title, message, variant } = event.data as { title?: string @@ -450,7 +576,7 @@ export default Plugin.define({ ) guard("settings request", () => - rpc.events.on("settings", async (event) => { + on("settings", async (event) => { if (!isLocal(event.location?.directory)) return try { // The server owns the values; the TUI just renders them. `dialog.alert` @@ -467,7 +593,7 @@ export default Plugin.define({ ) guard("help request", () => - rpc.events.on("help", async (event) => { + on("help", async (event) => { if (!isLocal(event.location?.directory)) return try { context.ui.dialog.set({ size: "large", centered: true }) @@ -479,7 +605,7 @@ export default Plugin.define({ ) guard("display request", () => - rpc.events.on("display", async (event) => { + on("display", async (event) => { if (!isLocal(event.location?.directory)) return const { placements, enabled } = event.data as { placements?: string[] @@ -494,16 +620,16 @@ export default Plugin.define({ // No placement named: open the multi-select dialog, so several can be // changed in one go. The command line still works for scripting. // - // Deliberately not wrapped in a try/catch. A swallowed failure here left - // an open dialog that accepted no keys at all and reported nothing, so - // if this ever throws again let it surface. + // Not wrapped in a try/catch, on purpose: a failure here should be seen. + // `on` reports it as a toast and keeps listening, where the stock + // rpc.events.on used to end the subscription for the rest of the session. if (named.length === 0) { await openDisplayPicker() return } if (named.length > 0) { - const next = applyParsed(display, { placements: named, enabled, unknown: [] }) + const next = applyParsed(display(), { placements: named, enabled, unknown: [] }) await setDisplay((draft) => { for (const key of named) draft[key] = next[key] }) @@ -526,6 +652,7 @@ export default Plugin.define({ return () => { clearInterval(presenceTimer) + clearInterval(pollTimer) for (const dispose of cleanups.reverse()) { try { dispose() diff --git a/src/waiting.ts b/src/waiting.ts new file mode 100644 index 0000000..75b9996 --- /dev/null +++ b/src/waiting.ts @@ -0,0 +1,101 @@ +/** + * Waiting on background work is not looping. + * + * An agent that started a twenty-minute build and has nothing else to do can only say + * "still running" each turn. The reply-repeat and stall guards read that as a runaway and + * pause the goal, even though the work is moving. This module answers one question — is a + * background command this session started still unfinished? — so the loop can back off + * instead of stopping. + * + * It fails open in both directions: a shape it does not recognise reports nothing pending, + * so the ordinary guards stay in force, and the wait is bounded by the caller. + */ + +const BACKGROUND_ID = /\bsh_[A-Za-z0-9]+\b/ + +/** Ids of background commands launched in these messages. */ +function launched(messages: readonly unknown[]): string[] { + const ids: string[] = [] + for (const message of messages) { + const record = message as { type?: unknown; content?: unknown } | null + if (record?.type !== "assistant") continue + for (const part of (record.content ?? []) as readonly unknown[]) { + const entry = part as { type?: unknown; state?: { input?: { background?: unknown }; output?: unknown; content?: unknown } } | null + if (entry?.type !== "tool" || entry.state?.input?.background !== true) continue + const found = BACKGROUND_ID.exec(JSON.stringify([entry.state.content ?? "", entry.state.output ?? ""])) + if (found) ids.push(found[0]) + } + } + return ids +} + +/** True when a completion notice for this id appears anywhere after it started. */ +function finished(id: string, messages: readonly unknown[]): boolean { + for (const message of messages) { + const record = message as { type?: unknown } | null + if (record?.type === "assistant") continue + const text = JSON.stringify(message) + if (text.includes(id) && /state=\\?"(completed|failed|killed|timed_out)/.test(text)) return true + } + return false +} + +/** Background commands started in this session that have not reported completion. */ +export function pendingBackground(messages: readonly unknown[]): string[] { + return launched(messages).filter((id) => !finished(id, messages)) +} + +/** Seconds to wait before the next prompt: 15s, doubling each wait, capped at two minutes. */ +export function backoffSeconds(waits: number): number { + return Math.min(120, 15 * 2 ** Math.max(0, waits)) +} + +/** Total time a goal may spend waiting on background work before the guards take over. */ +export const MAX_WAIT_MS = 60 * 60 * 1000 + +export type WaitEnd = "finished" | "timeout" | "aborted" | "inactive" + +/** + * Hold a turn open until background work finishes, the time is up, the caller aborts, or + * the goal stops being active. Polling is injected so this stays a pure, testable loop. + * Holding the turn open is the point: no reply is produced, so nothing repeats. + */ +export async function waitForBackground(options: { + seconds: number + intervalMs: number + pending: () => Promise + active: () => Promise + signal?: AbortSignal + sleep?: (ms: number) => Promise + now?: () => number +}): Promise<{ end: WaitEnd; waitedMs: number }> { + const sleep = options.sleep ?? ((ms) => new Promise((resolve) => setTimeout(resolve, ms))) + const now = options.now ?? Date.now + const started = now() + const deadline = started + options.seconds * 1000 + const result = (end: WaitEnd) => ({ end, waitedMs: now() - started }) + for (;;) { + if (options.signal?.aborted) return result("aborted") + if (!(await options.active())) return result("inactive") + if ((await options.pending()) === 0) return result("finished") + const left = deadline - now() + if (left <= 0) return result("timeout") + await sleep(Math.min(options.intervalMs, left)) + } +} + +/** Where a background command streams its output, read from the launch notice. */ +export function outputPathOf(messages: readonly unknown[], id: string): string | undefined { + for (const message of messages) { + const record = message as { type?: unknown; content?: unknown } | null + if (record?.type !== "assistant") continue + for (const part of (record.content ?? []) as readonly unknown[]) { + const state = (part as { type?: unknown; state?: { content?: unknown } } | null)?.state + const text = JSON.stringify(state?.content ?? "") + if (!text.includes(id)) continue + const found = new RegExp(`streaming to: (\\S*${id}\\.out)`).exec(text.replace(/\\n/g, " ")) + if (found) return found[1] + } + } + return undefined +} diff --git a/test/display.test.ts b/test/display.test.ts index 2873335..76c3812 100644 --- a/test/display.test.ts +++ b/test/display.test.ts @@ -10,6 +10,7 @@ import { describe, test } from "node:test" import assert from "node:assert/strict" import { applyAll, + applyMutation, applyParsed, applySet, DEFAULTS, @@ -270,3 +271,37 @@ describe("summarise", () => { assert.equal(summarise(applyAll(false)), "hidden everywhere") }) }) + +describe("applyMutation", () => { + test("returns every placement, not only the ones changed", () => { + // The UI shows all four but only ever changes the named one, so persistence has to + // write the whole value or an untouched placement silently reverts. + const next = applyMutation(DEFAULTS, (draft) => { + draft.panel = false + }) + assert.equal(next.panel, false) + assert.equal(next.footer, DEFAULTS.footer) + assert.equal(next.composer, DEFAULTS.composer) + assert.equal(next.sidebar, DEFAULTS.sidebar) + }) + + test("does not edit the value it was given", () => { + const before = { ...DEFAULTS } + applyMutation(before, (draft) => { + draft.sidebar = true + }) + assert.deepEqual(before, DEFAULTS) + }) + + test("returns a new object, so a change is visible when compared", () => { + // A mutation in place would leave the UI looking at the same reference, which is + // how a change fails to appear on the frame it happens. + const current = { ...DEFAULTS } + const next = applyMutation(current, (draft) => { + draft.composer = true + }) + assert.notEqual(next, current) + assert.equal(next.composer, true) + assert.equal(current.composer, false) + }) +}) diff --git a/test/listen.test.ts b/test/listen.test.ts new file mode 100644 index 0000000..5188b46 --- /dev/null +++ b/test/listen.test.ts @@ -0,0 +1,356 @@ +/** + * The resilient listener. + * + * The two cases that matter are the ones measured against the client's stock + * `rpc.events.on`: a handler that throws once, and a stream that ends. Both left + * a listener that was registered, looked alive, and never delivered again. + */ +import { describe, test } from "node:test" +import assert from "node:assert/strict" +import { listen } from "../src/listen.ts" + +/** Resolve on the next macrotask, letting every queued microtask run first. */ +const settle = () => new Promise((resolve) => setImmediate(resolve)) + +/** An async iterable that yields the given events then ends. */ +async function* finite(...events: E[]): AsyncGenerator { + for (const event of events) yield event +} + +/** Sleep that resolves at once and records what it was asked for. */ +const instantSleep = (delays: number[]) => async (ms: number) => { + delays.push(ms) +} + +/** + * Serves a queue of streams, one per connection, and stops the listener once + * they run out so a test never loops forever. + */ +function streams(plan: Array<() => AsyncIterable>, stop: () => void) { + let connections = 0 + const subscribe = (signal: AbortSignal) => { + const make = plan[connections++] + if (!make) { + stop() + return finite() + } + void signal + return make() + } + return { subscribe, connections: () => connections } +} + +describe("listen", () => { + test("delivers events in order", async () => { + const seen: number[] = [] + let stop: () => void = () => {} + const source = streams([() => finite(1, 2, 3)], () => stop()) + stop = listen(source.subscribe, (event: number) => void seen.push(event), { + sleep: instantSleep([]), + }) + await settle() + assert.deepEqual(seen, [1, 2, 3]) + stop() + }) + + // Measured on the stock client: one throw ended the loop, so events 2 and 3 + // below never arrived. + test("a handler that throws does not end the subscription", async () => { + const seen: number[] = [] + const errors: unknown[] = [] + let stop: () => void = () => {} + const source = streams([() => finite(1, 2, 3)], () => stop()) + stop = listen( + source.subscribe, + (event: number) => { + seen.push(event) + if (event === 1) throw new Error("boom") + }, + { onError: (error) => errors.push(error), sleep: instantSleep([]) }, + ) + await settle() + assert.deepEqual(seen, [1, 2, 3]) + assert.equal(errors.length, 1) + stop() + }) + + test("a handler that rejects does not end the subscription either", async () => { + const seen: number[] = [] + let stop: () => void = () => {} + const source = streams([() => finite(1, 2)], () => stop()) + stop = listen( + source.subscribe, + async (event: number) => { + seen.push(event) + throw new Error("async boom") + }, + { onError: () => {}, sleep: instantSleep([]) }, + ) + await settle() + assert.deepEqual(seen, [1, 2]) + stop() + }) + + // Measured on the stock client: the stream ended after one event and nothing + // ever opened a second connection, with no error to say so. + test("reopens the stream when it ends", async () => { + const seen: number[] = [] + let stop: () => void = () => {} + const source = streams([() => finite(1), () => finite(2)], () => stop()) + stop = listen(source.subscribe, (event: number) => void seen.push(event), { + sleep: instantSleep([]), + }) + await settle() + assert.deepEqual(seen, [1, 2]) + assert.ok(source.connections() >= 2) + stop() + }) + + test("reopens the stream when it throws, and reports why", async () => { + const seen: number[] = [] + const errors: string[] = [] + let stop: () => void = () => {} + const source = streams( + [ + async function* () { + yield 1 + throw new Error("socket closed") + }, + () => finite(2), + ], + () => stop(), + ) + stop = listen(source.subscribe, (event: number) => void seen.push(event), { + onError: (error) => errors.push((error as Error).message), + sleep: instantSleep([]), + }) + await settle() + assert.deepEqual(seen, [1, 2]) + assert.deepEqual(errors, ["socket closed"]) + stop() + }) + + test("runs onResume between a gap and the reconnect, never before the first connection", async () => { + const order: string[] = [] + let stop: () => void = () => {} + const source = streams( + [ + async function* () { + order.push("stream 1") + yield 1 + }, + async function* () { + order.push("stream 2") + yield 2 + }, + ], + () => stop(), + ) + stop = listen(source.subscribe, (event: number) => void order.push(`event ${event}`), { + onResume: () => void order.push("resume"), + sleep: instantSleep([]), + }) + await settle() + assert.deepEqual(order.slice(0, 5), ["stream 1", "event 1", "resume", "stream 2", "event 2"]) + stop() + }) + + test("says where each kind of failure came from", async () => { + const phases: string[] = [] + let stop: () => void = () => {} + const source = streams( + [ + async function* () { + yield 1 + throw new Error("stream") + }, + () => finite(2), + ], + () => stop(), + ) + stop = listen( + source.subscribe, + (event: number) => { + if (event === 1) throw new Error("handler") + }, + { + onResume: () => { + throw new Error("resume") + }, + onError: (error, phase) => phases.push(`${phase}:${(error as Error).message}`), + sleep: instantSleep([]), + }, + ) + await settle() + assert.deepEqual(phases.slice(0, 3), ["handler:handler", "stream:stream", "resume:resume"]) + stop() + }) + + test("a failing onResume does not stop the reconnect", async () => { + const seen: number[] = [] + const errors: unknown[] = [] + let stop: () => void = () => {} + const source = streams([() => finite(1), () => finite(2)], () => stop()) + stop = listen(source.subscribe, (event: number) => void seen.push(event), { + onResume: () => { + throw new Error("resync failed") + }, + onError: (error) => errors.push(error), + sleep: instantSleep([]), + }) + await settle() + assert.deepEqual(seen, [1, 2]) + assert.ok(errors.length >= 1) + stop() + }) + + test("a throwing error sink cannot take the listener down", async () => { + const seen: number[] = [] + let stop: () => void = () => {} + const source = streams([() => finite(1, 2)], () => stop()) + stop = listen( + source.subscribe, + (event: number) => { + seen.push(event) + throw new Error("handler") + }, + { + onError: () => { + throw new Error("sink") + }, + sleep: instantSleep([]), + }, + ) + await settle() + assert.deepEqual(seen, [1, 2]) + stop() + }) + + test("backs off while the stream keeps ending empty, up to the ceiling", async () => { + const delays: number[] = [] + let stop: () => void = () => {} + const source = streams( + [() => finite(), () => finite(), () => finite(), () => finite(), () => finite()], + () => stop(), + ) + stop = listen(source.subscribe, () => {}, { + sleep: instantSleep(delays), + minDelay: 100, + maxDelay: 500, + }) + await settle() + assert.deepEqual(delays.slice(0, 5), [100, 200, 400, 500, 500]) + stop() + }) + + test("forgets the backoff once an event proves the connection is healthy", async () => { + const delays: number[] = [] + let stop: () => void = () => {} + const source = streams( + [() => finite(), () => finite(), () => finite(1), () => finite()], + () => stop(), + ) + stop = listen(source.subscribe, () => {}, { + sleep: instantSleep(delays), + minDelay: 100, + maxDelay: 10_000, + }) + await settle() + // Two empty gaps double the delay (100, 200). The third connection delivers + // an event, which resets it, so the gap after that sleeps 100 again rather + // than the 400 it would have reached, and then doubles from there. + assert.deepEqual(delays.slice(0, 4), [100, 200, 100, 200]) + stop() + }) + + test("stopping aborts the signal handed to the stream", async () => { + let received: AbortSignal | undefined + const stop = listen( + (signal) => { + received = signal + return (async function* () { + await new Promise(() => {}) + })() + }, + () => {}, + { sleep: instantSleep([]) }, + ) + await settle() + assert.equal(received?.aborted, false) + stop() + assert.equal(received?.aborted, true) + }) + + test("does not reconnect once stopped", async () => { + let connections = 0 + let stop: () => void = () => {} + stop = listen( + () => { + connections++ + // subscribe runs synchronously inside listen(), before `stop` has been + // assigned, so stopping has to wait a microtask to reach the real one. + queueMicrotask(() => stop()) + return finite() + }, + () => {}, + { sleep: instantSleep([]) }, + ) + await settle() + assert.equal(connections, 1) + }) + + test("stopping during the backoff sleep prevents the reconnect", async () => { + let connections = 0 + let release: () => void = () => {} + const stop = listen( + () => { + connections++ + return finite() + }, + () => {}, + { + sleep: () => new Promise((resolve) => (release = resolve)), + }, + ) + await settle() + assert.equal(connections, 1) + stop() + release() + await settle() + assert.equal(connections, 1) + }) + + test("the default sleep wakes as soon as the listener is stopped", async () => { + let connections = 0 + const stop = listen( + () => { + connections++ + return finite() + }, + () => {}, + { minDelay: 60_000 }, + ) + await settle() + const started = Date.now() + stop() + await settle() + assert.ok(Date.now() - started < 1000) + assert.equal(connections, 1) + }) + + test("waits the real delay before reconnecting when not stopped", async () => { + let connections = 0 + const stop = listen( + () => { + connections++ + return finite() + }, + () => {}, + { minDelay: 40, maxDelay: 40 }, + ) + await new Promise((resolve) => setTimeout(resolve, 130)) + stop() + assert.ok(connections >= 2, `expected a reconnect, saw ${connections} connection(s)`) + assert.ok(connections <= 5, `expected backoff, saw ${connections} connections`) + }) +}) diff --git a/test/observation.test.ts b/test/observation.test.ts new file mode 100644 index 0000000..1387425 --- /dev/null +++ b/test/observation.test.ts @@ -0,0 +1,87 @@ +/** + * The polling guard. + * + * An agent waiting on a five-minute gate produces a turn that uses tools, says something + * new, and learns nothing. That shape defeats both existing guards: the reply digest + * moves, so `repeats` never fires, and the tool count is non-zero, so `stalled` never + * fires. Only the judge sees it — and before this change the judge's answer was to + * declare the *goal* unachievable, from turn-level evidence. + * + * `observationOf` is therefore a pure function worth pinning on its own: it has to ignore + * the tool inputs (so the same read with a reworded command still counts as the same + * observation) and it has to fail open on a shape it does not recognise. + */ +import { describe, test } from "node:test" +import assert from "node:assert/strict" +import { observationOf } from "../src/observation.ts" + +describe("observationOf", () => { + test("is empty for a turn with no tool calls", () => { + assert.equal( + observationOf([{ type: "assistant", content: [{ type: "text", text: "thinking about it" }] }]), + "", + ) + }) + + test("ignores the turn's prose entirely", () => { + const state = { status: "completed", output: "identical output" } + const first = observationOf([ + { type: "assistant", content: [{ type: "tool", state }] }, + ]) + const second = observationOf([ + { type: "assistant", content: [{ type: "tool", state }] }, + ]) + assert.equal(first, second) + assert.notEqual(first, "") + }) + + test("is the same when only the tool input is reworded", () => { + // This is the shape that slipped through: the agent varies the command it runs to + // read the same unchanged result. + const a = observationOf([ + { + type: "assistant", + content: [{ type: "tool", input: { command: "tail -5 log" }, state: { output: "same" } }], + }, + ]) + const b = observationOf([ + { + type: "assistant", + content: [{ type: "tool", input: { command: "tail -20 log" }, state: { output: "same" } }], + }, + ]) + assert.equal(a, b) + }) + + test("differs when the result actually differs", () => { + const a = observationOf([ + { type: "assistant", content: [{ type: "tool", state: { output: "exit 0" } }] }, + ]) + const b = observationOf([ + { type: "assistant", content: [{ type: "tool", state: { output: "exit 1" } }] }, + ]) + assert.notEqual(a, b) + }) + + test("fails open on a shape it does not recognise", () => { + // The right way for a heuristic to fail: no observation, so the guard stays quiet + // rather than pausing a loop on a guess. + assert.equal( + observationOf([{ type: "assistant", content: [{ type: "tool", state: null }] }]), + "", + ) + assert.equal(observationOf([{ type: "assistant", content: [] }]), "") + assert.equal(observationOf([]), "") + }) + + test("tolerates a cycle in the message values", () => { + const cyclic: Record = { output: "x" } + cyclic.self = cyclic + const messages = [{ type: "assistant", content: [{ type: "tool", state: cyclic }] }] + assert.doesNotThrow(() => observationOf(messages as never)) + }) + + test("ignores messages that are not the turn's assistant output", () => { + assert.equal(observationOf([{ type: "user", content: [{ type: "text", text: "hi" }] }]), "") + }) +}) diff --git a/test/settings.test.ts b/test/settings.test.ts index adf071e..d6abb33 100644 --- a/test/settings.test.ts +++ b/test/settings.test.ts @@ -8,6 +8,7 @@ import { describe, test } from "node:test" import assert from "node:assert/strict" import { + DEFAULT_POLL, DEFAULT_STALL, modelLabel, NO_OVERRIDES, @@ -19,7 +20,13 @@ import { type Overrides, } from "../src/settings.ts" -const OPTIONS = { maxTurns: 20 as const, stall: DEFAULT_STALL, quiet: false, judge: null } +const OPTIONS = { + maxTurns: 20 as const, + stall: DEFAULT_STALL, + poll: DEFAULT_POLL, + quiet: false, + judge: null, +} describe("parseCount", () => { test("reads a whole number", () => { @@ -119,9 +126,10 @@ describe("resolve", () => { assert.deepEqual(resolve(OPTIONS, NO_OVERRIDES), { maxTurns: 20, stall: 2, + poll: DEFAULT_POLL, quiet: false, judge: null, - overridden: { maxTurns: false, stall: false, quiet: false, judge: false }, + overridden: { maxTurns: false, stall: false, poll: false, quiet: false, judge: false }, }) }) @@ -132,9 +140,10 @@ describe("resolve", () => { assert.deepEqual(resolve(OPTIONS, { ...NO_OVERRIDES, maxTurns: null }), { maxTurns: null, stall: 2, + poll: DEFAULT_POLL, quiet: false, judge: null, - overridden: { maxTurns: true, stall: false, quiet: false, judge: false }, + overridden: { maxTurns: true, stall: false, poll: false, quiet: false, judge: false }, }) }) @@ -143,16 +152,29 @@ describe("resolve", () => { }) test("only the overridden settings change", () => { - const mixed: Overrides = { maxTurns: undefined, stall: 5, quiet: true, judge: undefined } + const mixed: Overrides = { maxTurns: undefined, stall: 5, poll: undefined, quiet: true, judge: undefined } assert.deepEqual(resolve(OPTIONS, mixed), { maxTurns: 20, stall: 5, + poll: DEFAULT_POLL, quiet: true, judge: null, - overridden: { maxTurns: false, stall: true, quiet: true, judge: false }, + overridden: { maxTurns: false, stall: true, poll: false, quiet: true, judge: false }, }) }) + test("a poll override beats the configured one", () => { + const effective = resolve({ ...OPTIONS, poll: 3 }, { ...NO_OVERRIDES, poll: 6 }) + assert.equal(effective.poll, 6) + assert.equal(effective.overridden.poll, true) + }) + + test("an absent poll override leaves the configured default alone", () => { + const effective = resolve({ ...OPTIONS, poll: 3 }, NO_OVERRIDES) + assert.equal(effective.poll, 3) + assert.equal(effective.overridden.poll, false) + }) + test("a judge override is reported as overridden", () => { const judge = { providerID: "a", id: "b" } const resolved = resolve(OPTIONS, { ...NO_OVERRIDES, judge }) @@ -200,6 +222,10 @@ describe("renderSettings", () => { assert.ok(rendered.includes("Stall limit 4 turns with no tools"), rendered) }) + test("shows the poll limit", () => { + assert.ok(rendered.includes("Poll limit"), rendered) + }) + test("shows quiet as on", () => { assert.ok(rendered.includes("Quiet mode on"), rendered) }) @@ -233,8 +259,8 @@ describe("renderSettings", () => { }) // Count the parenthesised row markers, not the bare words: the closing // paragraph legitimately mentions both "config default" and "this session". - test("marks all four rows as the config default when nothing is set", () => { - assert.equal(untouched.match(/\(config default\)/g)?.length, 4) + test("marks all five rows as the config default when nothing is set", () => { + assert.equal(untouched.match(/\(config default\)/g)?.length, 5) assert.equal(untouched.match(/\(this session\)/g), null) }) }) diff --git a/test/waiting.test.ts b/test/waiting.test.ts new file mode 100644 index 0000000..2c3a1c1 --- /dev/null +++ b/test/waiting.test.ts @@ -0,0 +1,112 @@ +import { describe, test } from "node:test" +import assert from "node:assert/strict" +import { backoffSeconds, outputPathOf, pendingBackground } from "../src/waiting.ts" + +const launch = (id: string) => ({ + type: "assistant", + content: [{ type: "tool", state: { input: { background: true }, output: `started ${id}` } }], +}) +const notice = (id: string) => ({ + type: "synthetic", + text: ``, +}) + +describe("pendingBackground", () => { + test("reads the id from the tool result content, as the host reports it", () => { + const message = { + type: "assistant", + content: [ + { + type: "tool", + state: { input: { background: true }, content: [{ type: "text", text: "moved (shell ID: sh_real9)." }] }, + }, + ], + } + assert.deepEqual(pendingBackground([message]), ["sh_real9"]) + }) + + test("reports a background command with no completion notice", () => { + assert.deepEqual(pendingBackground([launch("sh_abc123")]), ["sh_abc123"]) + }) + + test("clears it once the completion notice arrives", () => { + assert.deepEqual(pendingBackground([launch("sh_abc123"), notice("sh_abc123")]), []) + }) + + test("keeps an unfinished command when another one finished", () => { + assert.deepEqual( + pendingBackground([launch("sh_one"), launch("sh_two"), notice("sh_one")]), + ["sh_two"], + ) + }) + + test("ignores foreground commands", () => { + const message = { type: "assistant", content: [{ type: "tool", state: { input: {}, output: "sh_zzz" } }] } + assert.deepEqual(pendingBackground([message]), []) + }) + + test("fails open on shapes it does not recognise", () => { + assert.deepEqual(pendingBackground([null, {}, { type: "assistant", content: [{ type: "tool" }] }]), []) + }) +}) + +describe("backoffSeconds", () => { + test("doubles from 15s and caps at two minutes", () => { + assert.deepEqual([0, 1, 2, 3, 4, 9].map(backoffSeconds), [15, 30, 60, 120, 120, 120]) + }) +}) + +import { waitForBackground } from "../src/waiting.ts" + +describe("waitForBackground", () => { + const clock = () => { + let t = 0 + return { now: () => t, sleep: async (ms: number) => void (t += ms) } + } + const base = { seconds: 100, intervalMs: 10_000, active: async () => true } + + test("returns as soon as nothing is pending", async () => { + let calls = 0 + const c = clock() + const out = await waitForBackground({ ...base, ...c, pending: async () => (++calls < 3 ? 1 : 0) }) + assert.equal(out.end, "finished") + assert.equal(out.waitedMs, 20_000) + }) + + test("stops at the time limit while work is still pending", async () => { + const out = await waitForBackground({ ...base, ...clock(), pending: async () => 1 }) + assert.equal(out.end, "timeout") + assert.equal(out.waitedMs, 100_000) + }) + + test("stops when the caller aborts", async () => { + const controller = new AbortController() + controller.abort() + const out = await waitForBackground({ ...base, ...clock(), signal: controller.signal, pending: async () => 1 }) + assert.equal(out.end, "aborted") + }) + + test("stops when the goal is paused or cleared", async () => { + const out = await waitForBackground({ ...base, ...clock(), active: async () => false, pending: async () => 1 }) + assert.equal(out.end, "inactive") + }) +}) + +describe("outputPathOf", () => { + test("reads the streaming path from the launch notice", () => { + const message = { + type: "assistant", + content: [ + { + type: "tool", + state: { + input: { background: true }, + content: [{ type: "text", text: "moved (shell ID: sh_a1).\nOutput is streaming to: /tmp/x/sh_a1.out" }], + }, + }, + ], + } + assert.equal(outputPathOf([message], "sh_a1"), "/tmp/x/sh_a1.out") + assert.equal(outputPathOf([message], "sh_other"), undefined) + }) +})