diff --git a/.ai/contexts/trigger-watcher.md b/.ai/contexts/trigger-watcher.md index ed034166..58fc8675 100644 --- a/.ai/contexts/trigger-watcher.md +++ b/.ai/contexts/trigger-watcher.md @@ -8,6 +8,7 @@ |---|---|---| | `trigger-watcher.js` | ~1050 | The entire module: directory setup, `fs.watch` listener, idle-wait logic, single + chained trigger processing, submit-with-verify busy-rise/fall polling, input validation, PTY write, result file. | | `trigger-context.js` | ~35 | `createTriggerContext({ activeSessions, log })` — builds the whole `ctx` object out of `main.js`'s session map. | +| `transcript-turn.js` | ~140 | Whether a session transcript's main turn is closed, read from its tail; backs `ctx.getTranscriptTurn` (see "Transcript fallback while the descriptor stays busy"). | | `terminal-input.js` | ~20 | `handleTerminalInput(activeSessions, sessionId, data, now)` — the body of the `terminal-input` IPC handler; feeds `session.composerState`. | | `main.js` (wiring) | 3 | `require('./trigger-watcher').start(createTriggerContext({ activeSessions, log }))` in the `app.whenReady` block, right after `startScheduler`, plus the one-line `terminal-input` registration. | @@ -948,13 +949,197 @@ state was the cause. settle window, ends the wait even when `_cliBusy` is stuck true. An idle older than the Enter proves nothing (the Enter may have been absorbed) and leaves the `_cliBusy` logic in charge, as it does when no usable descriptor - exists. Not measured as fixed for sessions with background agents: the - descriptor stays `busy` until the last agent ends, so the busy-fall still - waits for it. + exists. For sessions with background agents the descriptor stays `busy` + until the last agent ends; the transcript fallback below covers that case. Tests: `test/trigger-every-step-readiness.test.js` (the real watcher with a fake descriptor, plus the wait helpers under mocked timers). +### Transcript fallback while the descriptor stays busy (issue #360) + +Measured: the descriptor keeps `status: "busy"` while background agents +(`run_in_background`) run, even when the prompt is free and the user can type, +and `shell` while background shell jobs run. A chain then waited out its +whole deadline after step 0 (issue comment of 2026-10-02: `steps_completed` +0, `waited_ms` 599959). Maintainer decision (2026-10-03): when the descriptor +says `busy` or `shell` but the transcript shows the turn is over, the prompt +is treated as free. + +**Where it applies.** Chains only: the readiness wait before every step +(`waitForCliIdleAfter`, 7th argument `{ transcriptAfterMs }`) and the +busy-fall wait after a non-final step (`waitForBusyFall`). Single triggers do +not pass the option and are unchanged. An `idle` descriptor never reaches the +fallback: its own path decides, including an idle older than the `/compact` +anchor, and the transcript is not even read. `waiting` (a dialog) never +reaches it either. + +**All of these must hold** (`transcriptShowsTurnOver`): + +- The descriptor reads `busy` or `shell`. +- No dialog: `waiting` was not sampled within the settle window + (`createDialogProbe`, the same detection as everywhere else). +- The last main-thread message entry of the transcript closes a turn and + has a parseable `timestamp` (`transcript-turn.js`, `classifyTranscriptTail`). + Three entries close one: an `assistant` entry with `stop_reason` `end_turn` + or `stop_sequence`, followed by a main-thread `system` `turn_duration` + entry; a `user` entry whose string content is one `` + block, start to end (the output of `/compact`, `/clear`, `/model`, …); a + `user` entry with `isCompactSummary: true` whose previous main-thread entry + is a `compact_boundary` with `compactMetadata.trigger` `manual` (an + automatic compaction happens inside a turn that goes on). Any other `user` + entry (a prompt, a meta prompt + such as a scheduled task or the ``, a + `` entry, a `tool_result`) means a turn is in progress, and so + does a `tool_use` stop. A slash command that expands into a prompt writes + its `` entry and then a model turn, so it is closed only by + that turn's `end_turn`. + Entries with `isSidechain` are skipped. Bookkeeping entries (`system`, + `attachment`, `last-prompt`, `file-history-*`, …) are skipped. A + `queue-operation` after the closed turn counts: any `dequeue`, or more + `enqueue` than `remove`, means a queued prompt is about to run. +- The closed turn is stamped at or after the anchor: the Enter of the step + just written (busy-fall), or the Enter of the previous step (readiness; no + anchor before step 0). A turn that closed before our Enter says nothing + about our step. +- The transcript file has not changed for `SWITCHBOARD_TRANSCRIPT_QUIET_MS` + (default 3000 ms, `DEFAULT_TRANSCRIPT_QUIET_MS`). With `turn_duration` + required, the quiet window no longer decides when a turn ends; it lets the + entries written in the same burst (a queued prompt's `enqueue`/`dequeue`) + land before the tail is trusted. + +**Why `turn_duration`, measured.** An `end_turn` alone is not the end of a +turn. Measured on the 19 real transcripts of this machine (2026-10-03, read +only), over 4344 main-thread `end_turn` messages: + +- 4120 (94.8 %) are followed by a `turn_duration` before the next message + entry. The gap is p50 0.27 s, p90 2.9 s, p99 42.8 s, max 124 s; 404 gaps + (9.8 %) are 3 s or more, so a quiet window alone would have read those + turns as over while their Stop hook still ran. +- The 224 others are all continued, never ended: 173 by a meta `user` entry, + 24 by a `user` prompt, 17 by a ``, 5 by a synthetic + `stop_sequence` message 9 to 76 s later, 2 by a local command, and 3 end + the file (a session still running, or killed). +- Where both are present (3695 times), `stop_hook_summary` always comes + before `turn_duration` (0 exceptions; `turn_duration` lands p50 16 ms, max + 5.5 s after it): `turn_duration` is written after the Stop hook. 39 + `end_turn` have a `stop_hook_summary` and no `turn_duration`; all 39 were + continued, so `stop_hook_summary` alone does not close. +- `stop_sequence` is only the synthetic message (`model` ``, 39 + entries): 19 are followed by a `turn_duration`, 20 are not. It closes on the + same condition as `end_turn`. +- No `turn_duration` follows a `` entry (0 of 95) or a + compaction summary (0 of 82), so those two keep the rule above, without + `turn_duration`. All 82 compactions of the corpus are `manual`. + +**What `/compact` leaves.** Measured on a real transcript (CLI of +2026-10-03, two compactions), in file order once compaction ends: a `system` +`compact_boundary`; the summary, a `user` entry with `isCompactSummary` and +`isVisibleInTranscriptOnly`, stamped about 0.8 s before the boundary; the +`` (`isMeta`) and the `/compact` entries, +both stamped at the command's Enter; the `` entry, +stamped last; then `attachment` entries and bookkeeping. The stdout entry is +the last message entry, so a `compact-now.sh` chain's step after `/compact` +is released from the transcript while background agents hold the descriptor +busy. The summary alone (read before the stdout lands) also closes, and the +quiet window covers the gap. The test fixture copies this shape with +synthetic text. + +**Proof of submission.** In edge mode a busy descriptor that does not change +status writes no new `statusUpdatedAt`, so our Enter would never count as +seen and the chain would stop on `step not confirmed`. For a chain step only +(`submitWithVerify(…, { transcriptReaction: true })`; single triggers do not +pass it), while the descriptor reads `busy` or `shell`, the step's own entry +stamped at or after the Enter also counts as the CLI's reaction +(`transcriptReactedSince`): a non-meta main-thread `user` entry with string +content, or an `enqueue` `queue-operation`, whose text equals the step's +command once trimmed (`promptMatches`). Two shapes are matched besides: + +- a slash command: content made only of ``, + ``, `` and `` elements, whose + `` element, wherever it sits, equals the command's first word. + Measured (2026-10-03, read only): of 116 main-thread entries holding a + ``, the 111 command entries are made only of those elements, + 99 with `` first and 12 with `` first (skills + and custom commands); the 5 others are compaction summaries quoting them; +- a `!cmd` step: content that is one `` element whose text, + trimmed, equals the command without its `!` (23 of 23 such entries are one + whole element). + +Attachments, `system` entries, meta entries and other texts (another agent's +notice) never confirm. The reader keeps the last 50 such entries of the tail +(`prompts`). Each step records `confirm_source` (`descriptor` or +`transcript`) when its submission was confirmed in edge mode. + +**A chain step written while the CLI stays busy.** The CLI writes the +`/compact` entry only when compaction ends, one to three +minutes after the Enter although it is stamped at the Enter, and a descriptor +held `busy` writes no new stamp. In 79 of 82 measured compactions a plain +`user` entry `/compact` is also written at the Enter, which the 2 s verify +window does confirm; in the other 3 it sees nothing, so +an unconfirmed chain step whose recovery Enter was withheld is not a failure +when the descriptor reads `busy` or `shell` and the session's transcript is +readable: it is pending (`waitForPendingConfirmation`), up to the step's own +deadline, and no Enter is written meanwhile. The step is confirmed when the +descriptor reacts after the Enter (`confirm_source` `descriptor`), or when +the step's own entry, stamped at or after the Enter, appears and the turn is +closed after that entry under the rules above, quiet window included +(`confirm_source` `transcript`). A turn closed before the step's own entry, +or another turn with no own entry, does not count. The dialog check is not +applied here: a `waiting` written after the Enter carries a new stamp, so the +descriptor branch, checked first, confirms the step on it. + +Two limits, `error` `step not confirmed` and nothing more typed at either: + +- No own entry within `SWITCHBOARD_PENDING_OWN_ENTRY_MS` of the Enter + (default 30 s, `DEFAULT_PENDING_OWN_ENTRY_MS`): `reasonNoOwnEntry`. A + prompt or slash command shows at once, as a `user` entry or an `enqueue`, + so a swallowed Enter fails then instead of spending the chain's budget. +- `/compact` is exempt: 3 of the 82 manual compactions measured wrote no + entry naming `/compact` before compaction ended (Enter to boundary: 0.5 s + to 332 s, median 129 s). It waits to the step deadline, as does any step + whose own entry has appeared and whose turn is not closed yet: + `REASON_UNCONFIRMED_BEFORE_DEADLINE`. A chain starting with `/compact` + under busy should set `timeout_ms` to 600 000. A dialog, a remote +session or a missing transcript keep the immediate `step not confirmed`. Once +confirmed, the step goes on as any other: for a non-final step the busy-fall +wait reads the same closed turn and ends at once. + +**The result records the signal.** Each chain step carries `ready_source` +(`descriptor` or `transcript`: what released the readiness wait before it; +absent when no descriptor was available) and, for a non-final step, +`idle_source` (`descriptor`, `transcript`, `busy_flag` for the `_cliBusy` +level probe, `no_rise` when no turn was ever observed). A step ended from the +transcript also logs one info line. + +**Reading the transcript.** `trigger-context.js` builds +`ctx.getTranscriptTurn(sessionId)` when `main.js` passes `projectsDir` +(`PROJECTS_DIR`): `//.jsonl`, local sessions only (a remote session returns `null`). The reader +stats the file and re-reads only when its mtime or size changed, the last +256 KB, so a poll costs a `stat`. The cache holds one entry per path, so +two chains on two sessions do not re-read each other's files. A chain +evicts its session's entry when it ends, whatever the outcome +(`ctx.forgetTranscriptTurn`, which remembers the path read for the session, +so it works after the session left `activeSessions`). A missing file, a read error, or a tail +with no message entry is never "closed". + +**Known limits.** + +- An `end_turn` with no `turn_duration` never reads closed: in the corpus, + 3 of 4344 ended the file that way; a chain there waits for the descriptor. +- An interrupted turn ends on a `user` entry and never reads closed. +- A step whose text is pasted and stored differently by the CLI (a + `[Pasted text]` placeholder, for example) is not matched: its Enter is + confirmed by the descriptor or not at all. Not measured. +- `popAll` queue operations (26 in the corpus) are not counted. +- The anchors compare the CLI's transcript timestamps with this process's + clock; both are the same machine's clock. + +Tests: `test/trigger-transcript-fallback.test.js` (the classification on +JSONL text, the reader on real files, the wait helpers under mocked timers, +and the chain through the real watcher and the real trigger context over a +real transcript file). + ### A blocked session tells its driver (issue #379) Every wait that can end a trigger on its deadline reports a dialog, not only @@ -1623,4 +1808,5 @@ observation layer. - If you rename `session.composerState` or stop feeding it from `terminal-input.js`, `getComposerState` returns `null` and **every trigger renounces with `not sent`** — the safe direction, but the channel goes silent. Tests for the model live in `test/composer-state.test.js`; the handler and the ctx are exercised in `test/terminal-input-handler.test.js` and `test/trigger-context.test.js`, and `test/main-wiring-source-check.test.js` reads `main.js` as text to check the remaining glue is still written down (source only — it proves nothing at runtime). - If you rename `activeSessions` or change the structure (`session.pty` → `session.ptyProcess`), update `getPtyForSession` and `isSessionBusy` in `trigger-context.js`, and `handleTerminalInput` in `terminal-input.js`. - If you rename `session.cwd` in `main.js`, update `getPtyForSession` in `trigger-context.js` — the target guard silently falls back to "indeterminate" (refuses every guarded trigger) rather than throwing, so this one fails quiet, not loud. +- If you rename `session.projectFolder` or `session.realSessionId` in `main.js`, update `getTranscriptTurn` in `trigger-context.js` — it returns `null` and the transcript fallback silently never fires (chains wait for the descriptor as before). - Tests live in `test/trigger-watcher.test.js`. They use `SWITCHBOARD_TRIGGERS_DIR` env override — do not hardcode paths there. diff --git a/CHANGELOG.md b/CHANGELOG.md index f0e85c99..0a14add5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,7 @@ What changes for you in each release of Switchboard. How to write an entry: [doc - Switchboard now checks once per host, at the first successful refresh and then every six hours (every 30 minutes while one is missing), whether `tmux` and `inotifywait` are installed. A host with `tmux` and no session running no longer shows attach as missing; a host without `tmux` no longer offers to attach to a session and opens its transcript, saying why in the tooltip; and the host's tooltip says when live updates are off because `inotifywait` is missing. On a remote host, the new-session button's tooltip now gives the reason, and Send a prompt… is disabled, with the reason, while no live session on the host reports a messaging socket. Stop is never disabled. (#218) - A trigger that gave up waiting for a session now says, in its result file's `reason`, when the session was blocked on a dialog such as a permission prompt or a question: for a single trigger, a chain's first wait, and a chain step whose turn never finished. Without a dialog the result is as before. (#379) ### Fixed +- A trigger chain no longer stops after its first step while background agents are running: when the CLI keeps reporting itself busy but the session transcript shows the turn finished, with no new write for 3 seconds and no dialog open, the next step is sent, and the step's `ready_source` / `idle_source` in the result file says `transcript`. A step typed while the CLI reports busy, such as `/compact`, is confirmed once the transcript shows it and its turn finished (`confirm_source` says `transcript`), up to the step's deadline; past it, or after 30 seconds for a step other than `/compact` that never shows, the chain stops with `step not confirmed` and nothing more is typed. Give a chain that starts with `/compact` a `timeout_ms` of 600000. (#360) - The Changes panel no longer shows `fatal: .git/index: index file open failed: Permission denied` now and then on Windows: a session's local changes are read one git call at a time instead of three at once. (#421) ## v0.0.87 — 2026-10-02 @@ -26,7 +27,7 @@ What changes for you in each release of Switchboard. How to write an entry: [doc ### Changed - After three failed refreshes of a remote host in a row, a row that would have attached opens its transcript and says why in its tooltip, instead of failing when clicked. Stop is never disabled: it runs its own ssh. (#218) ### Fixed -- A step of a trigger chain, the first one included, is no longer typed while the CLI reads busy or waiting on a dialog: it waits for the CLI to be at its prompt, up to the step's deadline, then fails cleanly with a reason instead of being written; a step whose Enter did not start a turn is retried once, or stops the chain when that retry is withheld because the CLI is busy or waiting on a dialog, or you have typed input pending, and is reported as "not confirmed submitted" instead of "sent". Without a readable CLI descriptor a step is still written as before, but no longer once its own deadline has passed. Single triggers are not covered. (#407, #360) +- A step of a trigger chain, the first one included, is no longer typed while the CLI reads busy or waiting on a dialog: it waits for the CLI to be at its prompt, up to the step's deadline, then fails cleanly with a reason instead of being written; a step whose Enter did not start a turn is retried once, or stops the chain when that retry is withheld because the CLI is busy (for a local session, after waiting for the step to show in the transcript) or waiting on a dialog, or you have typed input pending, and is reported as "not confirmed submitted" instead of "sent". Without a readable CLI descriptor a step is still written as before, but no longer once its own deadline has passed. Single triggers are not covered. (#407, #360) - After an upgrade, schedules keep running in a project that has settings of its own and in a git checkout that already holds a schedule; any other project, opened before the upgrade or not, runs no schedule until you open a session in it or add it. (#385) - Stopping a terminal twice in quick succession, or resizing it while it is being stopped, no longer closes the Windows pseudo console twice, which could kill the whole app with no error. (#405) - A sandboxed session, or a sandboxed schedule, whose Additional Directories include a `.claude` or `.git` directory, or a path inside one, is now refused instead of binding it read-write over its read-only protection; add the project directory instead. A session started in a `.claude` or `.git` directory is refused too, except below `.claude/worktrees`, and Additional Directories naming your home directory or a parent of it are refused however the path is written. A relative `add-dirs` entry in a schedule is taken from the schedule's directory. (#385) diff --git a/docs/automation.md b/docs/automation.md index 20d409bd..3ba8dea7 100644 --- a/docs/automation.md +++ b/docs/automation.md @@ -257,7 +257,11 @@ Size the budget by what the steps do, not by how many there are. A chain that compacts and then resumes spends most of it waiting for the session to go idle after `/compact`: `{"chain": [{"command": "/compact"}, {"command": "…"}], "timeout_ms": 600000}` is the shape that fits. A session busy for another reason -spends the same budget. +spends the same budget. Keep `timeout_ms` at 600 000 for a chain that starts +with `/compact` while background agents keep the CLI busy: the step is +confirmed only when compaction ends, and the 82 manual compactions measured +on one machine took from under a second to 332 s from the Enter, median +129 s. ### Environment overrides @@ -268,6 +272,7 @@ spends the same budget. | `SWITCHBOARD_TRIGGER_QUIET_MS` | The politeness quiet window | 3 000 | | `SWITCHBOARD_TRIGGER_MAX_AGE_MS` | The staleness limit | 300 000 | | `SWITCHBOARD_SUBMIT_ENTER_DELAY_MS` | Delay between the text and its Enter | 50 | +| `SWITCHBOARD_PENDING_OWN_ENTRY_MS` | How long a chain step written while the CLI reads `busy` may take to show in the transcript (not `/compact`) | 30 000 | | `SWITCHBOARD_SUBMIT_VERIFY_MS` | How long a submission is watched for a turn | 2 000 | | `SWITCHBOARD_BUSY_FALL_SETTLE_MS` | How long "not busy" must hold between chain steps | 300 | @@ -440,7 +445,7 @@ never in `error`: `not sent: input pending` is not `not sent`. |---|---|---| | `not sent` | **not one byte reached the session**: no idle came, politeness never allowed a write, or the trigger was refused before any write (stale, bad `wait`, bad `expectedCwd`, target guard) | nothing happened; it is safe to send again | | `chain timeout` | at least one step **was written**, and the expected effect was not observed before the deadline | assume the written steps landed | -| `step not confirmed` | a chain step **was written**, its submission was not confirmed by the CLI's descriptor, and the recovery Enter was withheld (the descriptor reads `busy` or `waiting`, or input of your own is pending in the composer); the chain stopped there and nothing more was typed | the step may sit unsubmitted in the composer: look before sending again | +| `step not confirmed` | a chain step **was written**, its submission was not confirmed by the CLI's descriptor, and the recovery Enter was withheld (the descriptor reads `busy` or `waiting`, or input of your own is pending in the composer); the chain stopped there and nothing more was typed. For a local session whose descriptor reads `busy`, the chain first waits for the step to show in the session transcript with its turn finished: up to 30 s for the step to show at all (up to the step's deadline for `/compact`, which shows only when compaction ends), then up to the step's deadline for its turn to finish; `reason` says which wait ran out | after a dialog or input of your own, the step may sit unsubmitted in the composer: look before sending again. After a wait under `busy`, the step may still be running or queued in the CLI: look at the session before sending again | | anything else | free text: `session not found`, `target process not running`, `missing required field`, `invalid timeout_ms`, `command and chain are mutually exclusive`, `trigger too large (max 64 KB)`, `command too long (max 4 KB)`, `trigger must be a regular file`, `pty write failed: …` | read `submitted` to know whether anything landed | The two reserved values mean opposite things: diff --git a/main.js b/main.js index 8ca046e5..3abd74b6 100644 --- a/main.js +++ b/main.js @@ -3215,7 +3215,7 @@ if (!gotSingleInstanceLock) { // I3: wrapped in try/catch so a boot failure here doesn't abort // app.whenReady (auto-updater, etc. would otherwise be silently lost). try { - require('./trigger-watcher').start(createTriggerContext({ activeSessions, log, getCliStatus: (id) => cliSessionState.getStatus(id) })); + require('./trigger-watcher').start(createTriggerContext({ activeSessions, log, getCliStatus: (id) => cliSessionState.getStatus(id), projectsDir: PROJECTS_DIR })); } catch (err) { log.error('[trigger-watcher] Failed to start trigger watcher:', err.message); } diff --git a/test/trigger-busy-chain-stall.test.js b/test/trigger-busy-chain-stall.test.js new file mode 100644 index 00000000..192d8884 --- /dev/null +++ b/test/trigger-busy-chain-stall.test.js @@ -0,0 +1,148 @@ +// test/trigger-busy-chain-stall.test.js +// +// Issue #360: a chain stalls after step 0 while the CLI descriptor stays +// "busy" because background agents run, although the turn is over and the +// session transcript says so. End to end through the real watcher and the +// real trigger context, with the transcript as a file on disk. +'use strict'; + +process.env.SWITCHBOARD_SUBMIT_ENTER_DELAY_MS = '1'; +process.env.SWITCHBOARD_SUBMIT_VERIFY_MS = '400'; +process.env.SWITCHBOARD_BUSY_FALL_SETTLE_MS = '50'; +process.env.SWITCHBOARD_BUSY_RISE_WAIT_MS = '100'; +process.env.SWITCHBOARD_TRANSCRIPT_QUIET_MS = '300'; + +const test = require('node:test'); +const assert = require('node:assert/strict'); +const fs = require('fs'); +const os = require('os'); +const path = require('path'); + +const { start } = require('../trigger-watcher'); +const { createTriggerContext } = require('../trigger-context'); + +function mkTmp(prefix) { + return fs.realpathSync.native(fs.mkdtempSync(path.join(os.tmpdir(), prefix))); +} + +const iso = (ms) => new Date(ms).toISOString(); +const jsonl = (entries) => entries.map((e) => JSON.stringify(e)).join('\n') + '\n'; +const prompt = (at, content) => ({ type: 'user', timestamp: iso(at), message: { role: 'user', content } }); +const endTurn = (at) => ({ type: 'assistant', timestamp: iso(at), message: { role: 'assistant', stop_reason: 'end_turn', content: [{ type: 'text', text: 'x' }] } }); +const system = (at, subtype, extra = {}) => ({ type: 'system', subtype, timestamp: iso(at), ...extra }); +const enqueue = (at, content) => ({ type: 'queue-operation', operation: 'enqueue', content, timestamp: iso(at) }); +const dequeue = (at) => ({ type: 'queue-operation', operation: 'dequeue', timestamp: iso(at) }); + +// What a turn answered while background agents keep the descriptor busy +// leaves: the queued prompt, the answer, then the Stop hook and turn_duration. +function closedTurn(command, delayMs) { + return (append) => { + append([enqueue(Date.now(), command), dequeue(Date.now()), prompt(Date.now(), command)]); + setTimeout(() => append([endTurn(Date.now()), system(Date.now(), 'stop_hook_summary'), system(Date.now(), 'turn_duration')]), delayMs); + }; +} + +// What /compact leaves, written only when compaction ends (synthetic text). +function compactionOutput(enterAt) { + const done = Date.now(); + return [ + system(done + 800, 'compact_boundary', { compactMetadata: { trigger: 'manual', preTokens: 1 } }), + { type: 'user', isCompactSummary: true, isVisibleInTranscriptOnly: true, timestamp: iso(done), message: { role: 'user', content: 'This session is being continued. Synthetic summary.' } }, + { type: 'user', isMeta: true, timestamp: iso(enterAt), message: { role: 'user', content: 'Caveat: synthetic.' } }, + prompt(enterAt, '/compact\n compact\n '), + prompt(done + 900, 'Compacted'), + { type: 'attachment', timestamp: iso(done + 100), attachment: { type: 'file' } }, + ]; +} + +function busySession(sessionId, onEnter) { + const projectsDir = mkTmp('sw-busy-stall-projects-'); + fs.mkdirSync(path.join(projectsDir, 'C--proj')); + const file = path.join(projectsDir, 'C--proj', sessionId + '.jsonl'); + const old = Date.now() - 60_000; + fs.writeFileSync(file, jsonl([prompt(old - 1000, 'earlier'), endTurn(old), system(old + 100, 'turn_duration')])); + const past = new Date(old + 200); + fs.utimesSync(file, past, past); + + const written = []; + let lastText = null; + const desc = { status: 'busy', statusUpdatedAt: Date.now() - 30_000 }; + const append = (entries) => { + if (fs.existsSync(projectsDir)) fs.appendFileSync(file, jsonl(entries)); + }; + const session = { + pty: { + pid: process.pid, + write(data) { + written.push({ data, at: Date.now() }); + if (data !== '\r') lastText = data; + if (data === '\r') onEnter({ append, command: lastText, n: written.filter((w) => w.data === '\r').length }); + }, + }, + projectFolder: 'C--proj', + _cliBusy: true, + composerState: { pending: 0, lastInputAt: 0 }, + }; + const ctx = createTriggerContext({ + activeSessions: new Map([[sessionId, session]]), + log: { info() {}, warn() {}, error() {}, debug() {} }, + isPtyAlive: () => true, + getCliStatus: () => ({ ...desc }), + projectsDir, + }); + return { ctx, written, cleanup: () => fs.rmSync(projectsDir, { recursive: true, force: true }) }; +} + +async function runChain(chain, s, uuid, timeoutMs) { + const tmp = mkTmp('sw-busy-stall-triggers-'); + process.env.SWITCHBOARD_TRIGGERS_DIR = tmp; + const watcher = start(s.ctx); + try { + fs.writeFileSync(path.join(tmp, uuid + '.json'), + JSON.stringify({ sessionId: uuid, wait: 'none', chain, timeout_ms: timeoutMs }), 'utf8'); + const resultPath = path.join(tmp, 'processed', uuid + '.result.json'); + const deadline = Date.now() + timeoutMs + 5000; + while (!fs.existsSync(resultPath)) { + if (Date.now() > deadline) throw new Error('no result file'); + await new Promise((r) => setTimeout(r, 20)); + } + await new Promise((r) => setTimeout(r, 20)); + return JSON.parse(fs.readFileSync(resultPath, 'utf8')); + } finally { + watcher.close(); + delete process.env.SWITCHBOARD_TRIGGERS_DIR; + fs.rmSync(tmp, { recursive: true, force: true }); + } +} + +test('#360: a descriptor held busy by background agents, the turn closed in the transcript -> step 1 is written', async () => { + const uuid = 'sess-busy-stall-' + Date.now(); + const s = busySession(uuid, ({ append, command }) => closedTurn(command, 100)(append)); + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 6000); + + assert.ok(s.written.some((w) => w.data === 'second step'), 'step 1 was never written: ' + JSON.stringify(result)); + assert.equal(result.ok, true); + assert.equal(result.steps[0].submit_confirmed, true); + } finally { + s.cleanup(); + } +}); + +test('#360: /compact under a descriptor held busy, its output written when compaction ends -> step 1 is written', async () => { + const uuid = 'sess-busy-stall-compact-' + Date.now(); + const s = busySession(uuid, ({ append, command, n }) => { + const enterAt = Date.now(); + if (n === 1) setTimeout(() => append(compactionOutput(enterAt)), 3000); + else closedTurn(command, 100)(append); + }); + try { + const result = await runChain([{ command: '/compact' }, { command: 'second step' }], s, uuid, 10000); + + assert.ok(s.written.some((w) => w.data === 'second step'), 'step 1 was never written after /compact: ' + JSON.stringify(result)); + assert.equal(result.ok, true); + assert.equal(result.steps[0].submit_confirmed, true); + } finally { + s.cleanup(); + } +}); diff --git a/test/trigger-transcript-fallback.test.js b/test/trigger-transcript-fallback.test.js new file mode 100644 index 00000000..d61233d5 --- /dev/null +++ b/test/trigger-transcript-fallback.test.js @@ -0,0 +1,991 @@ +// test/trigger-transcript-fallback.test.js +// +// The transcript fallback for a descriptor held "busy" (or "shell") by +// background work while the prompt is free (issue #360). See +// .ai/contexts/trigger-watcher.md, "Transcript fallback while the descriptor +// stays busy". +'use strict'; + +process.env.SWITCHBOARD_SUBMIT_ENTER_DELAY_MS = '1'; +process.env.SWITCHBOARD_SUBMIT_VERIFY_MS = '400'; +process.env.SWITCHBOARD_BUSY_FALL_SETTLE_MS = '50'; +process.env.SWITCHBOARD_BUSY_RISE_WAIT_MS = '100'; +process.env.SWITCHBOARD_TRANSCRIPT_QUIET_MS = '300'; +process.env.SWITCHBOARD_PENDING_OWN_ENTRY_MS = '1000'; + +const test = require('node:test'); +const assert = require('node:assert/strict'); +const fs = require('fs'); +const os = require('os'); +const path = require('path'); + +const { start, waitForBusyFall, waitForCliIdleAfter } = require('../trigger-watcher'); +const { createTriggerContext } = require('../trigger-context'); +const { classifyTranscriptTail, createTranscriptTurnReader, promptMatches } = require('../transcript-turn'); + +function mkTmp(prefix) { + return fs.realpathSync.native(fs.mkdtempSync(path.join(os.tmpdir(), prefix))); +} + +const iso = (ms) => new Date(ms).toISOString(); +const userPrompt = (at, extra = {}) => ({ type: 'user', timestamp: iso(at), message: { role: 'user', content: 'do it' }, ...extra }); +const toolResult = (at) => ({ type: 'user', timestamp: iso(at), message: { role: 'user', content: [{ type: 'tool_result', content: 'ok' }] } }); +const assistant = (at, stopReason, extra = {}) => ({ type: 'assistant', timestamp: iso(at), message: { role: 'assistant', stop_reason: stopReason, content: [{ type: 'text', text: 'x' }] }, ...extra }); +const systemEntry = (at, subtype) => ({ type: 'system', subtype, timestamp: iso(at) }); +const queueOp = (at, operation) => ({ type: 'queue-operation', operation, timestamp: iso(at) }); +const jsonl = (entries) => entries.map((e) => JSON.stringify(e)).join('\n') + '\n'; + +// ── The transcript classification ─────────────────────────────────────────── + +test('classify: a closed assistant turn followed only by bookkeeping entries is closed, stamped at that entry', () => { + const t = classifyTranscriptTail(jsonl([ + userPrompt(1000), assistant(2000, 'tool_use'), toolResult(2500), assistant(3000, 'end_turn'), + systemEntry(3100, 'stop_hook_summary'), systemEntry(3200, 'turn_duration'), { type: 'last-prompt' }, + ])); + assert.equal(t.closed, true); + assert.equal(t.closedAt, 3000); +}); + +test('classify: a last assistant entry that is a tool_use is not closed', () => { + const t = classifyTranscriptTail(jsonl([userPrompt(1000), assistant(2000, 'tool_use')])); + assert.equal(t.closed, false); +}); + +for (const [name, entry] of [['a user prompt', userPrompt(4000)], ['a tool_result', toolResult(4000)], ['a meta user prompt', userPrompt(4000, { isMeta: true })]]) { + test(`classify: ${name} after a closed turn means a turn is in progress`, () => { + const t = classifyTranscriptTail(jsonl([assistant(3000, 'end_turn'), entry])); + assert.equal(t.closed, false); + }); +} + +test('classify: sidechain entries after the closed turn do not count as the main turn', () => { + const t = classifyTranscriptTail(jsonl([ + assistant(3000, 'end_turn'), turnDuration(3100), userPrompt(3500, { isSidechain: true }), assistant(3600, 'tool_use', { isSidechain: true }), + ])); + assert.equal(t.closed, true); + assert.equal(t.closedAt, 3000); +}); + +test('classify: a prompt enqueued after the closed turn and not removed means a turn is coming', () => { + assert.equal(classifyTranscriptTail(jsonl([assistant(3000, 'end_turn'), turnDuration(3100), queueOp(3500, 'enqueue')])).closed, false); + assert.equal(classifyTranscriptTail(jsonl([assistant(3000, 'end_turn'), turnDuration(3100), queueOp(3500, 'enqueue'), queueOp(3600, 'dequeue')])).closed, false); + assert.equal(classifyTranscriptTail(jsonl([assistant(3000, 'end_turn'), turnDuration(3100), queueOp(3500, 'enqueue'), queueOp(3600, 'remove')])).closed, true); +}); + +// The shape /compact leaves in a real transcript (CLI measured 2026-10-03), +// with synthetic text: boundary, summary, caveat, the command, its stdout, +// then attachments. +const localStdout = (at, text = 'Compacted') => ({ type: 'user', timestamp: iso(at), message: { role: 'user', content: `${text}` } }); +const slashCommand = (at, name) => ({ type: 'user', timestamp: iso(at), message: { role: 'user', content: `/${name}\n ${name}\n ` } }); +const caveat = (at) => ({ type: 'user', isMeta: true, timestamp: iso(at), message: { role: 'user', content: 'Caveat: synthetic.' } }); +const compactSummary = (at) => ({ type: 'user', isCompactSummary: true, isVisibleInTranscriptOnly: true, timestamp: iso(at), message: { role: 'user', content: 'This session is being continued. Synthetic summary.' } }); +const attachment = (at) => ({ type: 'attachment', timestamp: iso(at), attachment: { type: 'file' } }); +const boundary = (at, trigger) => ({ ...systemEntry(at, 'compact_boundary'), ...(trigger ? { compactMetadata: { trigger, preTokens: 1 } } : {}) }); +const compactWrites = (enterAt, doneAt) => [ + boundary(doneAt + 800, 'manual'), compactSummary(doneAt), caveat(enterAt), slashCommand(enterAt, 'compact'), + localStdout(doneAt + 900), attachment(doneAt + 100), attachment(doneAt + 110), +]; + +test('classify: what /compact leaves (stdout of the local command last) is a closed turn, stamped at the stdout', () => { + const t = classifyTranscriptTail(jsonl([userPrompt(1000), assistant(2000, 'end_turn'), ...compactWrites(3000, 50_000)])); + assert.equal(t.closed, true); + assert.equal(t.closedAt, 50_900); +}); + +test('classify: the compaction summary as the last main-thread entry is a closed turn, stamped at the summary', () => { + const t = classifyTranscriptTail(jsonl([assistant(2000, 'end_turn'), boundary(50_800, 'manual'), compactSummary(50_000)])); + assert.equal(t.closed, true); + assert.equal(t.closedAt, 50_000); +}); + +test('classify: the output of another local command (/model) last is a closed turn', () => { + const t = classifyTranscriptTail(jsonl([assistant(2000, 'end_turn'), caveat(3000), slashCommand(3000, 'model'), localStdout(3100, 'Set model to x')])); + assert.equal(t.closed, true); + assert.equal(t.closedAt, 3100); +}); + +for (const [name, entries] of [ + ['a slash command expanding into a prompt', [slashCommand(3000, 'tdd')]], + ['the local command caveat alone', [caveat(3000)]], + ['user text quoting a stdout block', [{ ...localStdout(3000), message: { role: 'user', content: 'see this: x' } }]], + ['a stdout block followed by more user text', [{ ...localStdout(3000), message: { role: 'user', content: 'x and go on' } }]], + ['a tool_result carrying a stdout block', [{ type: 'user', timestamp: iso(3000), message: { role: 'user', content: [{ type: 'tool_result', content: 'x' }] } }]], + ['a user prompt followed by a sidechain stdout', [userPrompt(2900), { ...localStdout(3000, 'x'), isSidechain: true }]], +]) { + test(`classify: ${name} last is not a closed turn`, () => { + const t = classifyTranscriptTail(jsonl([assistant(2000, 'end_turn'), ...entries])); + assert.equal(t.closed, false); + }); +} + +test('classify: a local command output followed by a queued prompt, or unstamped, is not closed', () => { + assert.equal(classifyTranscriptTail(jsonl([localStdout(3000), queueOp(3500, 'enqueue')])).closed, false); + const unstamped = localStdout(3000); + delete unstamped.timestamp; + assert.equal(classifyTranscriptTail(jsonl([unstamped])).closed, false); +}); + +// ── Review of PR #433: the turn end, the summary, the queue, the cache ────── + +const turnDuration = (at) => systemEntry(at, 'turn_duration'); + +test('classify: an end_turn with no turn_duration after it is not closed (the turn can still resume)', () => { + assert.equal(classifyTranscriptTail(jsonl([userPrompt(1000), assistant(2000, 'end_turn')])).closed, false); +}); + +test('classify: an end_turn followed only by stop_hook_summary is not closed, the turn_duration closes it', () => { + assert.equal(classifyTranscriptTail(jsonl([assistant(2000, 'end_turn'), systemEntry(2300, 'stop_hook_summary')])).closed, false); + const t = classifyTranscriptTail(jsonl([assistant(2000, 'end_turn'), systemEntry(2300, 'stop_hook_summary'), turnDuration(2320)])); + assert.equal(t.closed, true); + assert.equal(t.closedAt, 2000); +}); + +test('classify: a turn_duration of an earlier turn, or of a sidechain, does not close the last end_turn', () => { + assert.equal(classifyTranscriptTail(jsonl([assistant(1000, 'end_turn'), turnDuration(1100), userPrompt(1500), assistant(2000, 'end_turn')])).closed, false); + assert.equal(classifyTranscriptTail(jsonl([assistant(2000, 'end_turn'), { ...turnDuration(2100), isSidechain: true }])).closed, false); +}); + +test('classify: a tool_use followed by a turn_duration is still not closed', () => { + assert.equal(classifyTranscriptTail(jsonl([assistant(2000, 'tool_use'), turnDuration(2100)])).closed, false); +}); + +test('classify: the synthetic stop_sequence message followed by its turn_duration is closed', () => { + assert.equal(classifyTranscriptTail(jsonl([assistant(2000, 'stop_sequence', { message: { role: 'assistant', model: '', stop_reason: 'stop_sequence', content: [] } }), turnDuration(2100)])).closed, true); +}); + +test('classify: a dequeue after the closed turn means a queued prompt is running', () => { + assert.equal(classifyTranscriptTail(jsonl([assistant(2000, 'end_turn'), turnDuration(2100), queueOp(2500, 'dequeue')])).closed, false); +}); + +test('classify: the compaction summary closes only after a manual compact_boundary', () => { + assert.equal(classifyTranscriptTail(jsonl([boundary(50_800, 'auto'), compactSummary(50_000)])).closed, false); + assert.equal(classifyTranscriptTail(jsonl([boundary(50_800), compactSummary(50_000)])).closed, false); + assert.equal(classifyTranscriptTail(jsonl([{ ...systemEntry(50_800, 'informational'), compactMetadata: { trigger: 'manual' } }, compactSummary(50_000)])).closed, false); + assert.equal(classifyTranscriptTail(jsonl([assistant(2000, 'end_turn'), compactSummary(50_000)])).closed, false); + assert.equal(classifyTranscriptTail(jsonl([boundary(50_800, 'manual'), compactSummary(50_000)])).closed, true); +}); + +test('classify: prompts lists the main-thread user prompts and enqueued contents with their stamps, newest last', () => { + const t = classifyTranscriptTail(jsonl([ + userPrompt(1000, { message: { role: 'user', content: 'first step' } }), + userPrompt(1100, { isMeta: true, message: { role: 'user', content: 'meta text' } }), + userPrompt(1200, { isSidechain: true, message: { role: 'user', content: 'side text' } }), + toolResult(1300), + { ...queueOp(1400, 'enqueue'), content: 'second step' }, + { ...queueOp(1450, 'remove'), content: 'removed step' }, + attachment(1500), + { type: 'user', message: { role: 'user', content: 'unstamped' } }, + assistant(1600, 'end_turn'), + ])); + assert.deepEqual(t.prompts, [{ at: 1000, text: 'first step' }, { at: 1400, text: 'second step' }]); +}); + +test('promptMatches: the step\'s own text, trimmed, or the of a slash command; nothing else', () => { + assert.equal(promptMatches('first step', 'first step'), true); + assert.equal(promptMatches(' first step\n', 'first step '), true); + assert.equal(promptMatches('/compact\ncompact\n', '/compact'), true); + assert.equal(promptMatches('/compact\nkeep x', '/compact keep x'), true); + assert.equal(promptMatches('/compactor', '/compact'), false); + assert.equal(promptMatches('/compact', 'compact'), false); + assert.equal(promptMatches('please run first step now', 'first step'), false); + assert.equal(promptMatches(' ', ''), false); + assert.equal(promptMatches('see /compact', '/compact'), false); + assert.equal(promptMatches(undefined, 'first step'), false); +}); + +test('promptMatches: a skill or custom command written first still matches its element', () => { + const skill = 'update-config\n/update-config\nx'; + assert.equal(promptMatches(skill, '/update-config x'), true); + assert.equal(promptMatches('loop\n/loop', '/loop'), true); + assert.equal(promptMatches(skill, '/update'), false); + assert.equal(promptMatches('x then /loop', '/loop'), false); + assert.equal(promptMatches('/compact and more text', '/compact'), false); +}); + +test('promptMatches: a !cmd step matches its element, and nothing else', () => { + assert.equal(promptMatches('ls -la', '!ls -la'), true); + assert.equal(promptMatches(' ls -la', '! ls -la'), true); + assert.equal(promptMatches('ls', '!rm'), false); + assert.equal(promptMatches('ls', 'ls'), false); + assert.equal(promptMatches('ls more', '!ls'), false); + assert.equal(promptMatches('', '!'), false); +}); + +test('reader: the tail cache is kept per path, so alternating sessions do not re-read', () => { + const dir = mkTmp('sw-transcript-cache-'); + const realOpen = fs.openSync; + const opened = []; + fs.openSync = (p, ...rest) => { opened.push(String(p)); return realOpen(p, ...rest); }; + try { + const a = path.join(dir, 'a.jsonl'); + const b = path.join(dir, 'b.jsonl'); + fs.writeFileSync(a, jsonl([assistant(1000, 'end_turn'), turnDuration(1100)])); + fs.writeFileSync(b, jsonl([userPrompt(1000)])); + const reader = createTranscriptTurnReader(); + assert.equal(reader.read(a).closed, true); + assert.equal(reader.read(b).closed, false); + assert.equal(reader.read(a).closed, true); + assert.equal(reader.read(b).closed, false); + assert.equal(opened.filter((p) => p === a).length, 1); + assert.equal(opened.filter((p) => p === b).length, 1); + } finally { + fs.openSync = realOpen; + fs.rmSync(dir, { recursive: true, force: true }); + } +}); + +test('classify: no message entry at all, or an unstamped closed turn, is never closed', () => { + assert.equal(classifyTranscriptTail(jsonl([systemEntry(1, 'x')])).closed, false); + const unstamped = assistant(3000, 'end_turn'); + delete unstamped.timestamp; + assert.equal(classifyTranscriptTail(jsonl([unstamped, turnDuration(3100)])).closed, false); +}); + +test('classify: the latest stamp of any main-thread entry is reported, sidechain stamps are not', () => { + const t = classifyTranscriptTail(jsonl([assistant(3000, 'end_turn'), queueOp(5000, 'enqueue'), userPrompt(9000, { isSidechain: true })])); + assert.equal(t.lastEntryAt, 5000); +}); + +test('classify: a tail cut mid-line skips the unparseable partial line', () => { + const text = '"stop_reason":"end_turn"}}\n' + jsonl([toolResult(4000)]); + assert.equal(classifyTranscriptTail(text).closed, false); +}); + +test('reader: reads the file tail, reports its mtime, and returns null for a missing file', () => { + const dir = mkTmp('sw-transcript-reader-'); + try { + const file = path.join(dir, 's.jsonl'); + fs.writeFileSync(file, 'x'.repeat(5000) + '\n' + jsonl([assistant(3000, 'end_turn'), turnDuration(3100)])); + const reader = createTranscriptTurnReader({ tailBytes: 1024 }); + const t = reader.read(file); + assert.equal(t.closed, true); + assert.equal(t.mtimeMs, fs.statSync(file).mtimeMs); + fs.appendFileSync(file, jsonl([toolResult(4000)])); + assert.equal(reader.read(file).closed, false); + assert.equal(reader.read(path.join(dir, 'missing.jsonl')), null); + } finally { + fs.rmSync(dir, { recursive: true, force: true }); + } +}); + +test('trigger context: getTranscriptTurn reads //.jsonl of a local session only', () => { + const dir = mkTmp('sw-transcript-ctx-'); + try { + fs.mkdirSync(path.join(dir, 'C--proj')); + fs.writeFileSync(path.join(dir, 'C--proj', 'real-id.jsonl'), jsonl([assistant(3000, 'end_turn'), turnDuration(3100)])); + const base = { pty: { pid: process.pid, write() {} }, projectFolder: 'C--proj' }; + const ctx = createTriggerContext({ + activeSessions: new Map([ + ['tmp-id', { ...base, realSessionId: 'real-id' }], + ['remote', { ...base, host: 'h', realSessionId: 'real-id' }], + ['nofolder', { pty: base.pty }], + ]), + log: { info() {}, warn() {}, error() {} }, + projectsDir: dir, + }); + assert.equal(ctx.getTranscriptTurn('tmp-id').closed, true); + assert.equal(ctx.getTranscriptTurn('remote'), null); + assert.equal(ctx.getTranscriptTurn('nofolder'), null); + assert.equal(ctx.getTranscriptTurn('unknown'), null); + assert.equal(createTriggerContext({ activeSessions: new Map(), log: {} }).getTranscriptTurn, undefined); + } finally { + fs.rmSync(dir, { recursive: true, force: true }); + } +}); + +// ── The wait helpers, under mocked timers ─────────────────────────────────── + +function fakeClock(t) { + t.mock.timers.enable({ apis: ['Date', 'setTimeout'], now: 1_000_000 }); +} + +async function settleRun(t, promise, maxMs = 5000) { + let done = false; + let value; + promise.then((v) => { done = true; value = v; }); + for (let i = 0; i < maxMs && !done; i += 5) { + t.mock.timers.tick(5); + await new Promise((r) => setImmediate(r)); + } + assert.ok(done, 'still pending'); + return value; +} + +function fallbackCtx(desc, turn) { + return { + getPtyForSession: () => ({}), + isSessionBusy: () => true, + getCliStatus: () => ({ ...desc }), + getTranscriptTurn: () => (turn ? { ...turn } : null), + }; +} + +const T = { transcriptAfterMs: -Infinity }; + +test('readiness: busy descriptor, closed turn quiet for the window -> ready from the transcript', async (t) => { + fakeClock(t); + const ctx = fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, { closed: true, closedAt: 999_000, lastEntryAt: 999_000, mtimeMs: 999_900 }); + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 3000, 50, false, T)); + assert.equal(r.ready, true); + assert.equal(r.source, 'transcript'); + assert.ok(r.waited_ms >= 200, 'ready after ' + r.waited_ms + ' ms, before the file had been quiet for 300 ms'); +}); + +test('readiness: the fallback also applies to a "shell" descriptor', async (t) => { + fakeClock(t); + const ctx = fallbackCtx({ status: 'shell', statusUpdatedAt: 900_000 }, { closed: true, closedAt: 999_000, lastEntryAt: 999_000, mtimeMs: 990_000 }); + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 3000, 50, false, T)); + assert.equal(r.ready, true); + assert.equal(r.source, 'transcript'); +}); + +test('readiness: a transcript still in a turn keeps the wait going to the deadline', async (t) => { + fakeClock(t); + const ctx = fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, { closed: false, closedAt: 980_000, lastEntryAt: 990_000, mtimeMs: 990_000 }); + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 1000, 50, false, T)); + assert.equal(r.ready, false); + assert.equal(r.timedOut, true); +}); + +test('readiness: a closed turn whose file changed inside the window waits until it is quiet', async (t) => { + fakeClock(t); + const turn = { closed: true, closedAt: 999_000, lastEntryAt: 999_000, mtimeMs: 1_000_000 }; + const ctx = { ...fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, null), getTranscriptTurn: () => ({ ...turn }) }; + setTimeout(() => { turn.mtimeMs = Date.now(); }, 200); + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 3000, 50, false, T)); + assert.equal(r.ready, true); + assert.ok(r.waited_ms >= 500, 'ready after ' + r.waited_ms + ' ms, the write at 200 ms must restart the quiet window'); +}); + +test('readiness: a closed turn older than the anchor (the previous step\'s Enter) is not readiness', async (t) => { + fakeClock(t); + const ctx = fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, { closed: true, closedAt: 999_000, lastEntryAt: 999_000, mtimeMs: 990_000 }); + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 1000, 50, false, { transcriptAfterMs: 999_500 })); + assert.equal(r.ready, false); + assert.equal(r.timedOut, true); +}); + +test('readiness: a dialog seen inside the window blocks the fallback', async (t) => { + fakeClock(t); + const desc = { status: 'waiting', statusUpdatedAt: 1_000_000 }; + const ctx = { ...fallbackCtx(desc, { closed: true, closedAt: 999_000, lastEntryAt: 999_000, mtimeMs: 990_000 }), getCliStatus: () => ({ ...desc }) }; + setTimeout(() => { desc.status = 'busy'; desc.statusUpdatedAt = Date.now(); }, 50); + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 2000, 300, false, T)); + assert.equal(r.ready, true); + assert.equal(r.source, 'transcript'); + assert.ok(r.waited_ms >= 350, 'ready after ' + r.waited_ms + ' ms, while the dialog was still inside the window'); +}); + +test('readiness: an idle descriptor older than the compact anchor stays in charge, the transcript is not read', async (t) => { + fakeClock(t); + let transcriptReads = 0; + const ctx = { + ...fallbackCtx({ status: 'idle', statusUpdatedAt: 999_000 }, null), + getTranscriptTurn: () => { transcriptReads += 1; return { closed: true, closedAt: 999_900, lastEntryAt: 999_900, mtimeMs: 990_000 }; }, + }; + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, 999_500, 1_000_000 + 1000, 50, false, T)); + assert.equal(r.ready, false); + assert.equal(transcriptReads, 0); +}); + +test('readiness: without the transcript option the busy descriptor keeps the wait going (single triggers unchanged)', async (t) => { + fakeClock(t); + const ctx = fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, { closed: true, closedAt: 999_000, lastEntryAt: 999_000, mtimeMs: 990_000 }); + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 1000, 50)); + assert.equal(r.ready, false); +}); + +test('readiness: an idle descriptor keeps its own path and reports it', async (t) => { + fakeClock(t); + let transcriptReads = 0; + const ctx = { ...fallbackCtx({ status: 'idle', statusUpdatedAt: 900_000 }, null), getTranscriptTurn: () => { transcriptReads += 1; return null; } }; + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 1000, 50, false, T)); + assert.equal(r.ready, true); + assert.equal(r.source, 'descriptor'); + assert.equal(transcriptReads, 0); +}); + +test('readiness: a transcript reader that throws is not readiness', async (t) => { + fakeClock(t); + const ctx = { ...fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, null), getTranscriptTurn: () => { throw new Error('EBUSY'); } }; + const r = await settleRun(t, waitForCliIdleAfter('sid', ctx, -Infinity, 1_000_000 + 600, 50, false, T)); + assert.equal(r.ready, false); + assert.equal(r.timedOut, true); +}); + +test('busy-fall: busy descriptor, _cliBusy stuck, turn closed after the Enter and quiet -> ends from the transcript', async (t) => { + fakeClock(t); + const ctx = fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, { closed: true, closedAt: 1_000_100, lastEntryAt: 1_000_100, mtimeMs: 1_000_100 }); + const r = await settleRun(t, waitForBusyFall('sid', ctx, 1_000_000 + 3000, 1_000_000)); + assert.equal(r.timedOut, false); + assert.equal(r.source, 'transcript'); + assert.ok(r.waited_ms >= 350, 'ended after ' + r.waited_ms + ' ms, before the file had been quiet for 300 ms'); +}); + +test('busy-fall: a closed turn that predates the Enter does not end the wait', async (t) => { + fakeClock(t); + const ctx = fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, { closed: true, closedAt: 999_000, lastEntryAt: 999_000, mtimeMs: 999_000 }); + const r = await settleRun(t, waitForBusyFall('sid', ctx, 1_000_000 + 1000, 1_000_000)); + assert.equal(r.timedOut, true); +}); + +test('busy-fall: a long tool call (transcript not closed) keeps the wait going', async (t) => { + fakeClock(t); + const ctx = fallbackCtx({ status: 'busy', statusUpdatedAt: 900_000 }, { closed: false, closedAt: null, lastEntryAt: 1_000_050, mtimeMs: 1_000_050 }); + const r = await settleRun(t, waitForBusyFall('sid', ctx, 1_000_000 + 1000, 1_000_000)); + assert.equal(r.timedOut, true); +}); + +test('busy-fall: the descriptor idle path reports its source', async (t) => { + fakeClock(t); + const desc = { status: 'busy', statusUpdatedAt: 900_000 }; + setTimeout(() => { desc.status = 'idle'; desc.statusUpdatedAt = Date.now(); }, 100); + const ctx = { ...fallbackCtx(desc, null), getCliStatus: () => ({ ...desc }) }; + const r = await settleRun(t, waitForBusyFall('sid', ctx, 1_000_000 + 3000, 1_000_000)); + assert.equal(r.source, 'descriptor'); +}); + +test('busy-fall: the level probe and the never-rose bound report their sources', async (t) => { + fakeClock(t); + let busy = true; + setTimeout(() => { busy = false; }, 50); + const level = { getPtyForSession: () => ({}), isSessionBusy: () => busy }; + assert.equal((await settleRun(t, waitForBusyFall('sid', level, 1_000_000 + 3000))).source, 'busy_flag'); + const never = { getPtyForSession: () => ({}), isSessionBusy: () => false }; + assert.equal((await settleRun(t, waitForBusyFall('sid', never, Date.now() + 3000))).source, 'no_rise'); +}); + +// ── The chain, through the real watcher and the real trigger context ─────── + +function transcriptSession(sessionId, { onEnter }) { + const projectsDir = mkTmp('sw-transcript-projects-'); + fs.mkdirSync(path.join(projectsDir, 'C--proj')); + const file = path.join(projectsDir, 'C--proj', sessionId + '.jsonl'); + const old = Date.now() - 60_000; + fs.writeFileSync(file, jsonl([userPrompt(old - 1000), assistant(old, 'end_turn'), systemEntry(old + 100, 'turn_duration')])); + const past = new Date(old + 200); + fs.utimesSync(file, past, past); + + const written = []; + let lastText = null; + const desc = { status: 'busy', statusUpdatedAt: Date.now() - 30_000 }; + const append = (entries) => { + if (fs.existsSync(projectsDir)) fs.appendFileSync(file, jsonl(entries)); + }; + const sessions = new Map(); + const session = { + pty: { + pid: process.pid, + write(data) { + written.push({ data, at: Date.now() }); + if (data !== '\r') lastText = data; + if (data === '\r') onEnter({ append, desc, n: written.filter((w) => w.data === '\r').length, command: lastText }); + }, + }, + projectFolder: 'C--proj', + _cliBusy: true, + composerState: { pending: 0, lastInputAt: 0 }, + }; + sessions.set(sessionId, session); + const ctx = createTriggerContext({ + activeSessions: sessions, + log: { info() {}, warn() {}, error() {}, debug() {} }, + isPtyAlive: () => true, + getCliStatus: () => ({ ...desc }), + projectsDir, + }); + return { ctx, written, desc, file, session, sessions, cleanup: () => fs.rmSync(projectsDir, { recursive: true, force: true }) }; +} + +async function runChain(chain, session, uuid, timeoutMs) { + const tmp = mkTmp('sw-transcript-triggers-'); + process.env.SWITCHBOARD_TRIGGERS_DIR = tmp; + const watcher = start(session.ctx); + try { + fs.writeFileSync(path.join(tmp, uuid + '.json'), + JSON.stringify({ sessionId: uuid, wait: 'none', chain, timeout_ms: timeoutMs }), 'utf8'); + const resultPath = path.join(tmp, 'processed', uuid + '.result.json'); + const deadline = Date.now() + timeoutMs + 5000; + while (!fs.existsSync(resultPath)) { + if (Date.now() > deadline) throw new Error('no result file'); + await new Promise((r) => setTimeout(r, 20)); + } + await new Promise((r) => setTimeout(r, 20)); + return JSON.parse(fs.readFileSync(resultPath, 'utf8')); + } finally { + watcher.close(); + delete process.env.SWITCHBOARD_TRIGGERS_DIR; + fs.rmSync(tmp, { recursive: true, force: true }); + } +} + +const ownPrompt = (at, command) => userPrompt(at, { message: { role: 'user', content: command } }); + +function closedTurnAfter(delayMs) { + return ({ append, command }) => { + append([{ ...queueOp(Date.now(), 'enqueue'), content: command }, queueOp(Date.now(), 'dequeue'), ownPrompt(Date.now(), command)]); + setTimeout(() => append([assistant(Date.now(), 'end_turn'), systemEntry(Date.now(), 'turn_duration')]), delayMs); + }; +} + +test('chain: a descriptor held busy by background agents, the turn closed and quiet -> step 1 is written', async () => { + const uuid = 'sess-tx-proceeds-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: closedTurnAfter(100) }); + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 6000); + + assert.ok(s.written.some((w) => w.data === 'second step'), 'step 1 was never written: ' + JSON.stringify(result)); + assert.equal(result.ok, true); + assert.equal(result.steps[0].idle_source, 'transcript'); + assert.equal(result.steps[0].ready_source, 'transcript'); + assert.equal(result.steps[1].ready_source, 'transcript'); + assert.equal(result.steps[0].submit_confirmed, true); + assert.equal(result.steps[0].confirm_source, 'transcript'); + } finally { + s.cleanup(); + } +}); + +test('chain: a closed turn older than the previous step\'s Enter does not release the next step', async () => { + const uuid = 'sess-tx-anchor-' + Date.now(); + let s; + s = transcriptSession(uuid, { + onEnter: ({ desc }) => { + desc.statusUpdatedAt = Date.now(); + s.session._cliBusy = true; + setTimeout(() => { s.session._cliBusy = false; }, 100); + }, + }); + s.session._cliBusy = false; + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 2500); + + assert.ok(!s.written.some((w) => w.data === 'second step'), 'the turn closed before step 0 released step 1'); + assert.equal(result.ok, false); + assert.equal(result.steps_completed, 1); + assert.equal(result.steps[0].idle_source, 'busy_flag'); + } finally { + s.cleanup(); + } +}); + +test('chain: busy descriptor and a long tool call (last entry tool_use) -> step 1 is never written', async () => { + const uuid = 'sess-tx-tooluse-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ append, command }) => { + append([ownPrompt(Date.now(), command)]); + setTimeout(() => append([assistant(Date.now(), 'tool_use')]), 100); + }, + }); + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 2500); + + assert.ok(!s.written.some((w) => w.data === 'second step'), 'false idle: step 1 was written during a tool call'); + assert.equal(result.ok, false); + assert.equal(result.error, 'chain timeout'); + } finally { + s.cleanup(); + } +}); + +test('chain: busy descriptor and the last entry a tool_result -> step 1 is never written', async () => { + const uuid = 'sess-tx-toolresult-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ append, command }) => { + append([ownPrompt(Date.now(), command)]); + setTimeout(() => append([assistant(Date.now(), 'tool_use'), toolResult(Date.now())]), 100); + }, + }); + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 2500); + + assert.ok(!s.written.some((w) => w.data === 'second step')); + assert.equal(result.ok, false); + } finally { + s.cleanup(); + } +}); + +test('chain: a closed turn followed by sidechain entries still releases step 1', async () => { + const uuid = 'sess-tx-sidechain-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ append, command }) => { + append([ownPrompt(Date.now(), command)]); + setTimeout(() => append([assistant(Date.now(), 'end_turn'), turnDuration(Date.now()), userPrompt(Date.now(), { isSidechain: true }), assistant(Date.now(), 'tool_use', { isSidechain: true })]), 100); + }, + }); + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 6000); + + assert.ok(s.written.some((w) => w.data === 'second step'), JSON.stringify(result)); + assert.equal(result.ok, true); + } finally { + s.cleanup(); + } +}); + +// The CLI writes everything /compact leaves, its entry +// included, only when compaction ends: here 3 s after the Enter, far beyond +// the verify window (400 ms in this file, 2 s in production). +const COMPACTION_MS = 3000; + +test('chain: /compact as step 0 with the descriptor held busy -> confirmed from the transcript when compaction ends, then step 1', async () => { + const uuid = 'sess-tx-compact-' + Date.now(); + let outputAt = null; + const s = transcriptSession(uuid, { + onEnter: ({ append, n, command }) => { + const enterAt = Date.now(); + if (n === 1) setTimeout(() => { outputAt = Date.now(); append(compactWrites(enterAt, Date.now())); }, COMPACTION_MS); + else closedTurnAfter(100)({ append, command }); + }, + }); + try { + const result = await runChain([{ command: '/compact' }, { command: 'second step' }], s, uuid, 10000); + + const second = s.written.find((w) => w.data === 'second step'); + assert.ok(second, 'step 1 was never written after /compact: ' + JSON.stringify(result)); + assert.ok(second.at > outputAt, 'step 1 was typed before the compaction output appeared'); + assert.equal(s.written.filter((w) => w.data === '\r').length, 2, 'a recovery Enter was typed into the busy CLI'); + assert.equal(result.ok, true); + assert.equal(result.steps[0].submit_confirmed, true); + assert.equal(result.steps[0].confirm_source, 'transcript'); + assert.equal(result.steps[0].submitted, 'confirmed'); + assert.equal(result.steps[0].idle_source, 'transcript'); + assert.equal(result.steps[1].ready_source, 'transcript'); + assert.equal(result.unconfirmed_steps, undefined); + } finally { + s.cleanup(); + } +}); + +test('chain: /compact as the only step with the descriptor held busy -> confirmed when compaction ends', async () => { + const uuid = 'sess-tx-compact-last-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ append }) => { + const enterAt = Date.now(); + setTimeout(() => append(compactWrites(enterAt, Date.now())), COMPACTION_MS); + }, + }); + try { + const result = await runChain([{ command: '/compact' }], s, uuid, 10000); + + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.steps[0].submit_confirmed, true); + assert.equal(result.steps[0].confirm_source, 'transcript'); + } finally { + s.cleanup(); + } +}); + +test('chain: a swallowed Enter under a busy descriptor fails once the own-entry wait is over, not at the step deadline, the next step never typed', async () => { + const uuid = 'sess-tx-swallowed-busy-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: () => {} }); + const forgotten = []; + const forget = s.ctx.forgetTranscriptTurn; + s.ctx.forgetTranscriptTurn = (id) => { forgotten.push(id); return forget(id); }; + const startedAt = Date.now(); + try { + const result = await runChain([{ command: 'first step', timeout_ms: 5000 }, { command: 'second step' }], s, uuid, 8000); + const elapsed = Date.now() - startedAt; + + assert.equal(result.ok, false); + assert.equal(result.error, 'step not confirmed'); + assert.match(result.reason, /did not show the step within 1 s of its Enter/); + assert.equal(result.steps[0].submit_confirmed, false); + assert.equal(result.steps_completed, 0); + assert.ok(elapsed >= 900, 'failed before the own-entry wait was over: ' + elapsed); + assert.ok(elapsed < 4000, 'waited for the step deadline: ' + elapsed); + assert.ok(!s.written.some((w) => w.data === 'second step')); + assert.equal(s.written.filter((w) => w.data === '\r').length, 1, 'a recovery Enter was typed into the busy CLI'); + assert.deepEqual(forgotten, [uuid]); + } finally { + s.cleanup(); + } +}); + +test('chain: a swallowed /compact under a busy descriptor waits for the step deadline, its entry only comes when compaction ends', async () => { + const uuid = 'sess-tx-swallowed-compact-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: () => {} }); + const startedAt = Date.now(); + try { + const result = await runChain([{ command: '/compact', timeout_ms: 2500 }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.error, 'step not confirmed', JSON.stringify(result)); + assert.match(result.reason, /before the step deadline/); + assert.ok(Date.now() - startedAt >= 2400, 'gave up on /compact before the step deadline'); + assert.ok(!s.written.some((w) => w.data === 'second step')); + } finally { + s.cleanup(); + } +}); + +test('chain: the step\'s own entry under a busy descriptor with the turn still running fails at the deadline, the next step never typed', async () => { + const uuid = 'sess-tx-pending-open-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ append, command }) => setTimeout(() => append([ownPrompt(Date.now(), command), assistant(Date.now(), 'tool_use')]), 800), + }); + const startedAt = Date.now(); + try { + const result = await runChain([{ command: 'first step', timeout_ms: 2500 }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.error, 'step not confirmed', JSON.stringify(result)); + assert.ok(Date.now() - startedAt >= 2400, 'stopped waiting for the closed turn once the own entry was there'); + assert.match(result.reason, /before the step deadline/); + assert.ok(!s.written.some((w) => w.data === 'second step')); + } finally { + s.cleanup(); + } +}); + +test('chain: under a busy descriptor, another turn closing after the Enter without the step\'s own entry never confirms it', async () => { + const uuid = 'sess-tx-pending-foreign-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ append }) => setTimeout(() => append([ + userPrompt(Date.now(), { message: { role: 'user', content: 'done' } }), + assistant(Date.now(), 'end_turn'), turnDuration(Date.now()), + ]), 300), + }); + try { + const result = await runChain([{ command: 'first step', timeout_ms: 2500 }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.error, 'step not confirmed', JSON.stringify(result)); + assert.match(result.reason, /did not show the step within/); + assert.ok(!s.written.some((w) => w.data === 'second step')); + } finally { + s.cleanup(); + } +}); + +test('chain: under a busy descriptor, a turn that closed before the step\'s own entry does not confirm it', async () => { + const uuid = 'sess-tx-pending-before-own-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ append, command }) => { + setTimeout(() => append([assistant(Date.now(), 'end_turn'), turnDuration(Date.now())]), 200); + setTimeout(() => append([{ ...queueOp(Date.now(), 'enqueue'), content: command }, { ...queueOp(Date.now(), 'remove'), content: command }]), 700); + }, + }); + try { + const result = await runChain([{ command: 'first step', timeout_ms: 2500 }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.error, 'step not confirmed', JSON.stringify(result)); + assert.match(result.reason, /before the step deadline/); + } finally { + s.cleanup(); + } +}); + +test('chain: a session that exits while its step is pending ends on session exited', async () => { + const uuid = 'sess-tx-pending-exit-' + Date.now(); + let s; + s = transcriptSession(uuid, { onEnter: () => setTimeout(() => s.sessions.delete(uuid), 1000) }); + try { + const result = await runChain([{ command: 'first step', timeout_ms: 4000 }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.error, 'session exited during wait', JSON.stringify(result)); + assert.ok(!s.written.some((w) => w.data === 'second step')); + } finally { + s.cleanup(); + } +}); + +test('chain: a pending step confirmed by the descriptor reacting later reports the descriptor as its source', async () => { + const uuid = 'sess-tx-pending-desc-' + Date.now(); + let s; + s = transcriptSession(uuid, { + onEnter: ({ desc, n, command, append }) => { + if (n === 1) { + setTimeout(() => { desc.status = 'idle'; desc.statusUpdatedAt = Date.now(); }, 800); + } else { + append([ownPrompt(Date.now(), command)]); + } + }, + }); + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.steps[0].submit_confirmed, true, JSON.stringify(result)); + assert.equal(result.steps[0].confirm_source, 'descriptor'); + assert.ok(s.written.some((w) => w.data === 'second step')); + } finally { + s.cleanup(); + } +}); + +test('chain: without a readable transcript a busy unconfirmed step stops at once, as before', async () => { + const uuid = 'sess-tx-no-transcript-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: ({ desc }) => { desc.status = 'busy'; } }); + s.desc.status = 'idle'; + s.ctx.getTranscriptTurn = () => null; + const startedAt = Date.now(); + try { + const result = await runChain([{ command: 'first step', timeout_ms: 4000 }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.error, 'step not confirmed', JSON.stringify(result)); + assert.match(result.reason, /recovery Enter was withheld/); + assert.ok(Date.now() - startedAt < 3000, 'waited for the deadline without a transcript'); + } finally { + s.cleanup(); + } +}); + +test('chain: a dialog right after the Enter stops the step at once, it is never pending', async () => { + const uuid = 'sess-tx-pending-dialog-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: ({ desc }) => { desc.status = 'waiting'; } }); + const startedAt = Date.now(); + try { + const result = await runChain([{ command: 'first step', timeout_ms: 4000 }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.error, 'step not confirmed', JSON.stringify(result)); + assert.match(result.reason, /recovery Enter was withheld/); + assert.ok(Date.now() - startedAt < 3000, 'waited for the deadline on a dialog'); + } finally { + s.cleanup(); + } +}); + +test('trigger context: forgetTranscriptTurn drops the cached tail of a session, even once it left activeSessions', () => { + const dir = mkTmp('sw-transcript-forget-'); + const realOpen = fs.openSync; + const opened = []; + fs.openSync = (p, ...rest) => { opened.push(String(p)); return realOpen(p, ...rest); }; + try { + fs.mkdirSync(path.join(dir, 'C--proj')); + const file = path.join(dir, 'C--proj', 'sid.jsonl'); + fs.writeFileSync(file, jsonl([userPrompt(1000)])); + const sessions = new Map([['sid', { pty: { pid: process.pid, write() {} }, projectFolder: 'C--proj' }]]); + const ctx = createTriggerContext({ activeSessions: sessions, log: { info() {}, warn() {}, error() {} }, projectsDir: dir }); + ctx.getTranscriptTurn('sid'); + ctx.getTranscriptTurn('sid'); + assert.equal(opened.filter((p) => p === file).length, 1); + ctx.forgetTranscriptTurn('sid'); + ctx.getTranscriptTurn('sid'); + assert.equal(opened.filter((p) => p === file).length, 2); + sessions.delete('sid'); + ctx.forgetTranscriptTurn('sid'); + sessions.set('sid', { pty: { pid: process.pid, write() {} }, projectFolder: 'C--proj' }); + ctx.getTranscriptTurn('sid'); + assert.equal(opened.filter((p) => p === file).length, 3); + ctx.forgetTranscriptTurn('unknown'); + } finally { + fs.openSync = realOpen; + fs.rmSync(dir, { recursive: true, force: true }); + } +}); + +test('chain: a slash command expanding into a prompt, model turn still running -> step 1 is never written', async () => { + const uuid = 'sess-tx-skill-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ append }) => { + setTimeout(() => append([slashCommand(Date.now(), 'tdd'), userPrompt(Date.now(), { isMeta: true })]), 100); + }, + }); + try { + const result = await runChain([{ command: '/tdd' }, { command: 'second step' }], s, uuid, 2500); + + assert.ok(!s.written.some((w) => w.data === 'second step'), 'false idle: step 1 was written while the expanded prompt ran'); + assert.equal(result.ok, false); + } finally { + s.cleanup(); + } +}); + +test('chain: a swallowed Enter with only a pre-Enter entry of the same text ends on step not confirmed', async () => { + const uuid = 'sess-tx-swallowed-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: () => {} }); + fs.appendFileSync(s.file, jsonl([ownPrompt(Date.now() - 5000, 'first step'), assistant(Date.now() - 4900, 'end_turn'), turnDuration(Date.now() - 4800)])); + const past = new Date(Date.now() - 4000); + fs.utimesSync(s.file, past, past); + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 4000); + + assert.equal(result.ok, false); + assert.equal(result.error, 'step not confirmed'); + assert.equal(result.steps[0].submit_confirmed, false); + assert.ok(!s.written.some((w) => w.data === 'second step')); + } finally { + s.cleanup(); + } +}); + +test('chain: an idle descriptor that never moves is not confirmed by a transcript entry', async () => { + const uuid = 'sess-tx-idle-noreact-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: ({ append, command }) => append([ownPrompt(Date.now(), command)]) }); + s.desc.status = 'idle'; + s.session._cliBusy = false; + try { + const result = await runChain([{ command: 'first step' }], s, uuid, 4000); + + assert.equal(result.steps[0].submit_confirmed, false, JSON.stringify(result)); + } finally { + s.cleanup(); + } +}); + +for (const [name, entries] of [ + ['an attachment and a system entry', (at) => [attachment(at), systemEntry(at, 'informational')]], + ['another agent\'s notice', (at) => [userPrompt(at, { message: { role: 'user', content: 'done' } })]], + ['an enqueue of another text', (at) => [{ ...queueOp(at, 'enqueue'), content: 'something else' }]], +]) { + test(`chain: ${name} written after the Enter does not confirm it`, async () => { + const uuid = 'sess-tx-foreign-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: ({ append }) => append(entries(Date.now())) }); + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 4000); + + assert.equal(result.error, 'step not confirmed', JSON.stringify(result)); + assert.equal(result.steps[0].submit_confirmed, false); + } finally { + s.cleanup(); + } + }); +} + +test('single trigger: the transcript reaction does not confirm it (chains only)', async () => { + const uuid = 'sess-tx-single-' + Date.now(); + const s = transcriptSession(uuid, { onEnter: ({ append, command }) => append([{ ...queueOp(Date.now(), 'enqueue'), content: command }, ownPrompt(Date.now(), command)]) }); + const tmp = mkTmp('sw-transcript-triggers-'); + process.env.SWITCHBOARD_TRIGGERS_DIR = tmp; + const watcher = start(s.ctx); + try { + fs.writeFileSync(path.join(tmp, uuid + '.json'), JSON.stringify({ sessionId: uuid, command: 'first step', wait: 'none' }), 'utf8'); + const resultPath = path.join(tmp, 'processed', uuid + '.result.json'); + const deadline = Date.now() + 8000; + while (!fs.existsSync(resultPath)) { + if (Date.now() > deadline) throw new Error('no result file'); + await new Promise((r) => setTimeout(r, 20)); + } + await new Promise((r) => setTimeout(r, 20)); + const result = JSON.parse(fs.readFileSync(resultPath, 'utf8')); + + assert.ok(s.written.some((w) => w.data === 'first step'), JSON.stringify(result)); + assert.equal(result.submit_confirmed, false, JSON.stringify(result)); + } finally { + watcher.close(); + delete process.env.SWITCHBOARD_TRIGGERS_DIR; + fs.rmSync(tmp, { recursive: true, force: true }); + s.cleanup(); + } +}); + +test('chain: an idle descriptor keeps the descriptor path, the transcript is not the source', async () => { + const uuid = 'sess-tx-idle-' + Date.now(); + const s = transcriptSession(uuid, { + onEnter: ({ desc, append }) => { + desc.status = 'busy'; desc.statusUpdatedAt = Date.now(); + append([userPrompt(Date.now())]); + setTimeout(() => { desc.status = 'idle'; desc.statusUpdatedAt = Date.now(); }, 100); + }, + }); + s.desc.status = 'idle'; + try { + const result = await runChain([{ command: 'first step' }, { command: 'second step' }], s, uuid, 6000); + + assert.equal(result.ok, true); + assert.equal(result.steps[0].ready_source, 'descriptor'); + assert.equal(result.steps[0].idle_source, 'descriptor'); + assert.equal(result.steps[0].confirm_source, 'descriptor'); + assert.equal(result.steps[1].ready_source, 'descriptor'); + } finally { + s.cleanup(); + } +}); diff --git a/transcript-turn.js b/transcript-turn.js new file mode 100644 index 00000000..db431438 --- /dev/null +++ b/transcript-turn.js @@ -0,0 +1,151 @@ +// transcript-turn.js — whether a session transcript shows its main turn closed. +// see .ai/contexts/trigger-watcher.md, "Transcript fallback while the descriptor stays busy" +'use strict'; + +const fs = require('fs'); + +const CLOSED_STOP_REASONS = new Set(['end_turn', 'stop_sequence']); +const DEFAULT_TAIL_BYTES = 256 * 1024; + +function stampOf(entry) { + const t = typeof entry.timestamp === 'string' ? Date.parse(entry.timestamp) : NaN; + return Number.isFinite(t) ? t : null; +} + +const MAX_PROMPTS = 50; +const COMMAND_ELEMENTS_ONLY = /^(?:\s*<(command-message|command-name|command-args|command-contents)>[\s\S]*?<\/\1>)+\s*$/; + +function isLocalCommandOutput(entry) { + const content = entry.message && entry.message.content; + if (typeof content !== 'string') return false; + const t = content.trim(); + return t.startsWith('') && t.endsWith(''); +} + +function isManualCompactBoundary(entry) { + return !!entry && entry.type === 'system' && entry.subtype === 'compact_boundary' + && !!entry.compactMetadata && entry.compactMetadata.trigger === 'manual'; +} + +function endsTurn(entry, turnDurationAfter, previous) { + if (entry.type === 'user') { + if (entry.isCompactSummary === true) return isManualCompactBoundary(previous); + return isLocalCommandOutput(entry); + } + return turnDurationAfter && CLOSED_STOP_REASONS.has(entry.message && entry.message.stop_reason); +} + +function parseMainThread(text) { + const entries = []; + for (const raw of String(text).split('\n')) { + const line = raw.trim(); + if (!line) continue; + let entry; + try { entry = JSON.parse(line); } catch (_) { continue; } + if (!entry || typeof entry !== 'object' || entry.isSidechain) continue; + entries.push(entry); + } + return entries; +} + +function promptOf(entry) { + if (entry.type === 'queue-operation') { + return entry.operation === 'enqueue' && typeof entry.content === 'string' ? entry.content : null; + } + if (entry.type !== 'user' || entry.isMeta) return null; + const content = entry.message && entry.message.content; + return typeof content === 'string' ? content : null; +} + +function collectPrompts(entries) { + const prompts = []; + for (const entry of entries) { + const text = promptOf(entry); + const at = stampOf(entry); + if (text !== null && at !== null) prompts.push({ at, text }); + } + return prompts.slice(-MAX_PROMPTS); +} + +function promptMatches(text, command) { + if (typeof text !== 'string' || typeof command !== 'string') return false; + const want = command.trim(); + if (!want) return false; + const got = text.trim(); + if (got === want) return true; + if (want.startsWith('!')) { + const bash = got.match(/^([\s\S]*)<\/bash-input>$/); + const typed = want.slice(1).trim(); + return !!bash && !!typed && bash[1].trim() === typed; + } + if (!COMMAND_ELEMENTS_ONLY.test(got)) return false; + const name = got.match(/([^<]*)<\/command-name>/); + return !!name && name[1].trim() === want.split(/\s/)[0]; +} + +function classifyTranscriptTail(text) { + const entries = parseMainThread(text); + const prompts = collectPrompts(entries); + let lastEntryAt = null; + let enqueued = 0; + let removed = 0; + let dequeued = 0; + let turnDurationAfter = false; + for (let i = entries.length - 1; i >= 0; i -= 1) { + const entry = entries[i]; + const at = stampOf(entry); + if (at !== null && (lastEntryAt === null || at > lastEntryAt)) lastEntryAt = at; + if (entry.type === 'queue-operation') { + if (entry.operation === 'enqueue') enqueued += 1; + else if (entry.operation === 'remove') removed += 1; + else if (entry.operation === 'dequeue') dequeued += 1; + continue; + } + if (entry.type === 'system' && entry.subtype === 'turn_duration') turnDurationAfter = true; + if (entry.type !== 'user' && entry.type !== 'assistant') continue; + const queueIdle = enqueued <= removed && dequeued === 0; + const closed = endsTurn(entry, turnDurationAfter, entries[i - 1]) && queueIdle && at !== null; + return { closed, closedAt: closed ? at : null, lastEntryAt, prompts }; + } + return { closed: false, closedAt: null, lastEntryAt, prompts }; +} + +function readTail(filePath, size, tailBytes) { + const start = Math.max(0, size - tailBytes); + const length = size - start; + const buf = Buffer.alloc(length); + const fd = fs.openSync(filePath, 'r'); + try { + let off = 0; + while (off < length) { + const n = fs.readSync(fd, buf, off, length - off, start + off); + if (n <= 0) break; + off += n; + } + return buf.toString('utf8', 0, off); + } finally { + fs.closeSync(fd); + } +} + +function createTranscriptTurnReader({ tailBytes = DEFAULT_TAIL_BYTES } = {}) { + const cache = new Map(); + return { + read(filePath) { + let stat; + try { stat = fs.statSync(filePath); } catch (_) { return null; } + const hit = cache.get(filePath); + if (hit && hit.mtimeMs === stat.mtimeMs && hit.size === stat.size) return { ...hit.turn }; + let tail; + try { tail = readTail(filePath, stat.size, tailBytes); } catch (_) { return null; } + const turn = { ...classifyTranscriptTail(tail), mtimeMs: stat.mtimeMs }; + cache.set(filePath, { mtimeMs: stat.mtimeMs, size: stat.size, turn }); + return { ...turn }; + }, + forget(filePath) { + cache.delete(filePath); + }, + }; +} + +module.exports = { classifyTranscriptTail, createTranscriptTurnReader, promptMatches, CLOSED_STOP_REASONS }; diff --git a/trigger-context.js b/trigger-context.js index d2ad7e6a..b9c7eb33 100644 --- a/trigger-context.js +++ b/trigger-context.js @@ -1,6 +1,10 @@ // trigger-context.js — see .ai/contexts/trigger-watcher.md 'use strict'; +const path = require('path'); + +const { createTranscriptTurnReader } = require('./transcript-turn'); + // Local session handle — see .ai/contexts/trigger-watcher.md, "Session handle". function createLocalSessionHandle(ptyProcess) { return { @@ -25,9 +29,10 @@ function createLocalSessionHandle(ptyProcess) { * @param {object} deps.log electron-log compatible logger * @param {function} [deps.isPtyAlive] (ptyProcess) => boolean * @param {function} [deps.getCliStatus] (sessionId) => { status, statusUpdatedAt } | undefined + * @param {string} [deps.projectsDir] root of the CLI's transcript folders (~/.claude/projects) * @returns {object} ctx */ -function createTriggerContext({ activeSessions, log, isPtyAlive, getCliStatus }) { +function createTriggerContext({ activeSessions, log, isPtyAlive, getCliStatus, projectsDir }) { const ctx = { log, getPtyForSession(sessionId) { @@ -51,6 +56,25 @@ function createTriggerContext({ activeSessions, log, isPtyAlive, getCliStatus }) }, }; if (isPtyAlive) ctx.isPtyAlive = isPtyAlive; + // see .ai/contexts/trigger-watcher.md, "Transcript fallback while the descriptor stays busy" + if (projectsDir) { + const reader = createTranscriptTurnReader(); + const readPaths = new Map(); + ctx.getTranscriptTurn = (sessionId) => { + const session = activeSessions.get(sessionId); + if (!session || session.exited || session.host != null || !session.projectFolder) return null; + const id = session.realSessionId || sessionId; + const filePath = path.join(projectsDir, session.projectFolder, id + '.jsonl'); + readPaths.set(sessionId, filePath); + return reader.read(filePath); + }; + ctx.forgetTranscriptTurn = (sessionId) => { + const filePath = readPaths.get(sessionId); + if (filePath === undefined) return; + readPaths.delete(sessionId); + reader.forget(filePath); + }; + } if (getCliStatus) { ctx.getCliStatus = (sessionId) => { const session = activeSessions.get(sessionId); diff --git a/trigger-watcher.js b/trigger-watcher.js index af9941be..e5ba28b2 100644 --- a/trigger-watcher.js +++ b/trigger-watcher.js @@ -37,6 +37,7 @@ const path = require('path'); const os = require('os'); const { createLocalSessionHandle } = require('./trigger-context'); +const { promptMatches } = require('./transcript-turn'); const DEFAULT_TRIGGERS_DIR = path.join(os.homedir(), '.switchboard', 'triggers'); // Default idle-wait timeout: 5 minutes. @@ -94,6 +95,10 @@ const REASON_DIALOG_OPEN_AFTER_WRITE = 'the CLI reports a dialog open (waiting) const REASON_DEADLINE_BEFORE_WRITE = 'the step deadline passed before it could be written; nothing was written'; const REASON_CLI_BUSY = 'the CLI still reported a turn running (busy) at the deadline; nothing was written'; const REASON_CLI_NOT_IDLE = 'the CLI never reported idle before the deadline; nothing was written'; +const REASON_UNCONFIRMED_BEFORE_DEADLINE = 'the step was written while the CLI reported busy; neither the CLI nor the transcript (the step\'s own entry, then a closed turn) confirmed it before the step deadline; nothing more was typed'; +function reasonNoOwnEntry(ms) { + return `the step was written while the CLI reported busy; the CLI did not react and the transcript did not show the step within ${Math.round(ms / 1000)} s of its Enter; nothing more was typed`; +} const REASON_IDLE_UNSETTLED = 'the CLI was idle only briefly before the deadline; it never held long enough to settle; nothing was written'; const ACCEPTED_WAITS = ['idle', 'none']; @@ -304,6 +309,85 @@ function getBusyFallSettleMs() { return v !== undefined ? v : DEFAULT_BUSY_FALL_SETTLE_MS; } +// see .ai/contexts/trigger-watcher.md, "Transcript fallback while the descriptor stays busy" +const DEFAULT_TRANSCRIPT_QUIET_MS = 3000; // ms +function getTranscriptQuietMs() { + const v = envNumber('SWITCHBOARD_TRANSCRIPT_QUIET_MS'); + return v !== undefined ? v : DEFAULT_TRANSCRIPT_QUIET_MS; +} + +const TRANSCRIPT_FALLBACK_STATUSES = ['busy', 'shell']; + +// see .ai/contexts/trigger-watcher.md, "A chain step written while the CLI stays busy" +const DEFAULT_PENDING_OWN_ENTRY_MS = 30000; // ms +function getPendingOwnEntryMs() { + const v = envNumber('SWITCHBOARD_PENDING_OWN_ENTRY_MS'); + return v !== undefined ? v : DEFAULT_PENDING_OWN_ENTRY_MS; +} + +function readTranscriptTurn(ctx, sessionId) { + if (typeof ctx.getTranscriptTurn !== 'function') return null; + try { + return ctx.getTranscriptTurn(sessionId) || null; + } catch (_) { + return null; + } +} + +// see .ai/contexts/trigger-watcher.md, "Transcript fallback while the descriptor stays busy" +function transcriptShowsTurnOver(ctx, sessionId, desc, afterMs, now, dialogSeen) { + if (!desc || !TRANSCRIPT_FALLBACK_STATUSES.includes(desc.status) || dialogSeen) return false; + const turn = readTranscriptTurn(ctx, sessionId); + if (!turn || turn.closed !== true || !Number.isFinite(turn.closedAt) || !Number.isFinite(turn.mtimeMs)) return false; + if (!(turn.closedAt >= afterMs)) return false; + return now - turn.mtimeMs >= getTranscriptQuietMs(); +} + +function transcriptOwnEntryAt(ctx, sessionId, sinceMs, command) { + const turn = readTranscriptTurn(ctx, sessionId); + if (!turn || !Array.isArray(turn.prompts)) return null; + const own = turn.prompts.find((p) => Number.isFinite(p.at) && p.at >= sinceMs && promptMatches(p.text, command)); + return own ? own.at : null; +} + +function cliReadsBusyOrShell(ctx, sessionId) { + const desc = readCliStatusRaw(ctx, sessionId); + return !!desc && TRANSCRIPT_FALLBACK_STATUSES.includes(desc.status); +} + +function transcriptReactedSince(ctx, sessionId, sinceMs, command) { + if (!cliReadsBusyOrShell(ctx, sessionId)) return false; + return transcriptOwnEntryAt(ctx, sessionId, sinceMs, command) !== null; +} + +// see .ai/contexts/trigger-watcher.md, "A chain step written while the CLI stays busy" +function waitForPendingConfirmation(sessionId, ctx, enterAtMs, command, deadlineMs) { + const start = Date.now(); + const ownEntryDeadline = isCompactCommand(command) ? deadlineMs : Math.min(deadlineMs, enterAtMs + getPendingOwnEntryMs()); + return pollLoop((resolve, scheduleNext) => { + const now = Date.now(); + const waited_ms = now - start; + if (!ctx.getPtyForSession(sessionId)) { + return resolve({ confirmed: false, sessionExited: true, timedOut: false, waited_ms }); + } + if (cliReactedSince(ctx, sessionId, enterAtMs)) { + return resolve({ confirmed: true, source: 'descriptor', sessionExited: false, timedOut: false, waited_ms }); + } + const ownAt = transcriptOwnEntryAt(ctx, sessionId, enterAtMs, command); + if (ownAt !== null + && transcriptShowsTurnOver(ctx, sessionId, readCliStatusRaw(ctx, sessionId), ownAt, now, false)) { + return resolve({ confirmed: true, source: 'transcript', sessionExited: false, timedOut: false, waited_ms }); + } + if (ownAt === null && now >= ownEntryDeadline && ownEntryDeadline < deadlineMs) { + return resolve({ confirmed: false, sessionExited: false, timedOut: true, noOwnEntry: true, waited_ms }); + } + if (now >= deadlineMs) { + return resolve({ confirmed: false, sessionExited: false, timedOut: true, waited_ms }); + } + scheduleNext(); + }); +} + // see .ai/contexts/trigger-watcher.md, "waitForBusyFall waits for the rise too" function getBusyRiseWaitMs() { const v = envNumber('SWITCHBOARD_BUSY_RISE_WAIT_MS'); @@ -414,7 +498,7 @@ function cliForbidsRecoveryEnter(ctx, sessionId) { } // see .ai/contexts/trigger-watcher.md, "Readiness before every step" (return shape, descriptor loss, deadline) -function waitForCliIdleAfter(sessionId, ctx, afterMs, deadlineMs, settleMs = 0, trustIdleStamp = false) { +function waitForCliIdleAfter(sessionId, ctx, afterMs, deadlineMs, settleMs = 0, trustIdleStamp = false, transcript = null) { const start = Date.now(); let idleSince = null; let idleStamp = null; @@ -434,6 +518,7 @@ function waitForCliIdleAfter(sessionId, ctx, afterMs, deadlineMs, settleMs = 0, return resolve({ ready: false, available: false, timedOut: false, sessionExited: false, waited_ms, lastStatus, waitingSeen }); } let idleHeld = false; + let source = 'descriptor'; if (!s) { idleSince = null; } else { @@ -448,10 +533,14 @@ function waitForCliIdleAfter(sessionId, ctx, afterMs, deadlineMs, settleMs = 0, idleHeld = now - idleSince >= settleMs; } else { idleSince = null; + if (transcript && transcriptShowsTurnOver(ctx, sessionId, s, transcript.transcriptAfterMs, now, dialog.seen(now))) { + idleHeld = true; + source = 'transcript'; + } } } if (idleHeld && now < deadlineMs) { - return resolve({ ready: true, available: true, timedOut: false, sessionExited: false, waited_ms, lastStatus, waitingSeen }); + return resolve({ ready: true, available: true, timedOut: false, sessionExited: false, waited_ms, lastStatus, waitingSeen, source }); } if (now >= deadlineMs) { return resolve({ ready: false, available: true, timedOut: true, sessionExited: false, waited_ms, lastStatus, waitingSeen: dialog.seen(now) }); @@ -494,7 +583,7 @@ function waitForCliIdleAfter(sessionId, ctx, afterMs, deadlineMs, settleMs = 0, * the caller keeps the legacy instant-reply semantics — submit_retries traces * that the verification could not confirm a turn started. */ -async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs) { +async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs, { transcriptReaction = false } = {}) { // Sampled before the write — see .ai/contexts/trigger-watcher.md ("submitted"). const preBusy = ctx.isSessionBusy(sessionId); @@ -513,7 +602,14 @@ async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs) { // must fire on window expiry, so the deadline must NOT coincide with it. const effectiveDeadline = (deadlineMs !== undefined) ? deadlineMs : Infinity; - const probe = edgeMode ? () => cliReactedSince(ctx, sessionId, enterAt) : undefined; + let confirmSource = null; + const probe = edgeMode + ? () => { + if (cliReactedSince(ctx, sessionId, enterAt)) { confirmSource = 'descriptor'; return true; } + if (transcriptReaction && transcriptReactedSince(ctx, sessionId, enterAt, command)) { confirmSource = 'transcript'; return true; } + return false; + } + : undefined; const first = await pollForBusyObserved(sessionId, ctx, windowMs, effectiveDeadline, probe); if (first.sawBusy || first.sessionExited || first.timedOut) { return { @@ -527,6 +623,7 @@ async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs) { sessionExited: first.sessionExited, timedOut: first.timedOut, waited_ms: first.waited_ms, + confirmSource: first.sawBusy ? confirmSource : null, }; } @@ -590,6 +687,7 @@ async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs) { sessionExited: second.sessionExited, timedOut: second.timedOut, waited_ms: first.waited_ms + second.waited_ms, + confirmSource: second.sawBusy ? confirmSource : null, }; } @@ -643,10 +741,13 @@ function waitForBusyFall(sessionId, ctx, deadlineMs, enterAtMs) { descIdleStamp = desc.statusUpdatedAt; } if (now - descIdleSince >= settleMs) { - return resolve({ timedOut: false, sessionExited: false, waited_ms: now - start }); + return resolve({ timedOut: false, sessionExited: false, waited_ms: now - start, source: 'descriptor' }); } } else { descIdleSince = null; + if (desc && transcriptShowsTurnOver(ctx, sessionId, desc, enterAtMs, now, dialog.seen(now))) { + return resolve({ timedOut: false, sessionExited: false, waited_ms: now - start, source: 'transcript' }); + } } if (ctx.isSessionBusy(sessionId)) { hasRisen = true; @@ -654,11 +755,11 @@ function waitForBusyFall(sessionId, ctx, deadlineMs, enterAtMs) { } else if (hasRisen) { if (idleSince === null) idleSince = now; if (now - idleSince >= settleMs) { - return resolve({ timedOut: false, sessionExited: false, waited_ms: now - start }); + return resolve({ timedOut: false, sessionExited: false, waited_ms: now - start, source: 'busy_flag' }); } } else if (now >= riseDeadline) { // Never rose within the bound: turn never observed, not an error. - return resolve({ timedOut: false, sessionExited: false, waited_ms: now - start }); + return resolve({ timedOut: false, sessionExited: false, waited_ms: now - start, source: 'no_rise' }); } scheduleNext(); }); @@ -819,6 +920,7 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR } let releaseSessionLock; + let forgetTranscript = null; let stepsTotal = 0; try { @@ -1243,6 +1345,7 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR } // ── 6. Chain path ───────────────────────────────────────────────────────── + if (typeof ctx.forgetTranscriptTurn === 'function') forgetTranscript = () => ctx.forgetTranscriptTurn(sessionId); // Global deadline for the whole chain const globalTimeout = (resolvedTimeoutMs !== undefined) ? resolvedTimeoutMs : getIdleTimeout(); const globalDeadline = Date.now() + globalTimeout; @@ -1290,6 +1393,7 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR } let compactSentAtMs = null; + let previousEnterAtMs = -Infinity; const unconfirmedSteps = []; for (let i = 0; i < chain.length; i++) { @@ -1364,9 +1468,12 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR } let readyWaitedMs = 0; + let readySource = null; { - const ready = await waitForCliIdleAfter(sessionId, ctx, readyAfterMs === null ? -Infinity : readyAfterMs, stepDeadline, getBusyFallSettleMs()); + const ready = await waitForCliIdleAfter(sessionId, ctx, readyAfterMs === null ? -Infinity : readyAfterMs, stepDeadline, + getBusyFallSettleMs(), false, { transcriptAfterMs: previousEnterAtMs }); readyWaitedMs = ready.waited_ms; + if (ready.ready) readySource = ready.source; totalWaitedMs += readyWaitedMs; if (ready.sessionExited) { ctx.log.warn(`[trigger-watcher] Session exited waiting for the CLI to be ready at chain step ${i}:`, sessionId); @@ -1425,7 +1532,7 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR let stepWaitedMs = polite.waited_ms + readyWaitedMs; let verify; try { - verify = await submitWithVerify(entryHandle, sessionId, step.command, ctx, stepDeadline); + verify = await submitWithVerify(entryHandle, sessionId, step.command, ctx, stepDeadline, { transcriptReaction: true }); } catch (err) { ctx.log.error(`[trigger-watcher] PTY write failed at chain step ${i}:`, err.message); await writeResult({ ok: false, error: 'pty write failed: ' + err.message, partial: true, steps_completed: i, sessionId, sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs }); @@ -1439,9 +1546,27 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR if (isCompactCommand(step.command)) { compactSentAtMs = Number.isFinite(verify.enterAt) ? verify.enterAt : Date.parse(stepSentAt); } + previousEnterAtMs = Number.isFinite(verify.enterAt) ? verify.enterAt : Date.parse(stepSentAt); + const sources = readySource ? { ready_source: readySource } : {}; submitRetries = verify.submit_retries; stepWaitedMs += verify.waited_ms; totalWaitedMs += verify.waited_ms; + + if (verify.confirmed === false && verify.recoverySkipped && !verify.timedOut && !verify.sessionExited + && cliReadsBusyOrShell(ctx, sessionId) && readTranscriptTurn(ctx, sessionId) !== null) { + ctx.log.info(`[trigger-watcher] Chain step ${i} written while the CLI reads busy, waiting for its confirmation:`, sessionId); + const pending = await waitForPendingConfirmation(sessionId, ctx, verify.enterAt, step.command, stepDeadline); + stepWaitedMs += pending.waited_ms; + totalWaitedMs += pending.waited_ms; + if (pending.confirmed) { + verify = { ...verify, confirmed: true, composerConfirmed: true, sawBusy: true, recoverySkipped: false, confirmSource: pending.source }; + } else if (pending.sessionExited) { + verify = { ...verify, sessionExited: true }; + } else { + verify = { ...verify, pendingTimedOut: true, pendingNoOwnEntry: !!pending.noOwnEntry }; + } + } + if (verify.confirmSource) sources.confirm_source = verify.confirmSource; // This step's own submitted -- same classification the chain fold below // uses, attached to the step itself so a consumer can ask "was THIS step // (e.g. the last one) confirmed?" instead of only the chain's weakest. @@ -1465,24 +1590,28 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR // Session exited / global timeout observed during verify. if (verify.sessionExited) { ctx.log.warn(`[trigger-watcher] Session exited during chain step ${i} submit verify:`, sessionId); - steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }), ...sources }); await writeResult({ ok: false, submitted: chainSubmitted, error: 'session exited during wait', partial: true, steps_completed: i, sessionId, sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs }); return; } if (verify.timedOut) { ctx.log.warn(`[trigger-watcher] Chain timeout during step ${i} submit verify:`, sessionId); - steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }), ...sources }); await writeResult({ ok: false, submitted: chainSubmitted, error: 'chain timeout', partial: true, steps_completed: i, sessionId, sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs }); return; } if (stepConfirmed === false && verify.recoverySkipped) { - steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, submit_confirmed: false }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, submit_confirmed: false, ...sources }); await writeResult({ ok: false, submitted: chainSubmitted, error: ERROR_UNCONFIRMED, - reason: `chain step ${i} was typed but its submission was not confirmed and the recovery Enter was withheld (${verify.recoveryReason}); nothing more was typed`, + reason: verify.pendingNoOwnEntry + ? `chain step ${i}: ${reasonNoOwnEntry(getPendingOwnEntryMs())}` + : verify.pendingTimedOut + ? `chain step ${i}: ${REASON_UNCONFIRMED_BEFORE_DEADLINE}` + : `chain step ${i} was typed but its submission was not confirmed and the recovery Enter was withheld (${verify.recoveryReason}); nothing more was typed`, partial: true, steps_completed: i, sessionId, sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs, }); @@ -1501,23 +1630,29 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR const result = await waitForBusyFall(sessionId, ctx, stepDeadline, verify.enterAt); stepWaitedMs += result.waited_ms; totalWaitedMs += result.waited_ms; + if (result.source) { + sources.idle_source = result.source; + if (result.source === 'transcript') { + ctx.log.info(`[trigger-watcher] Chain step ${i} turn end read from the transcript, the descriptor still reads busy:`, sessionId); + } + } if (result.sessionExited) { ctx.log.warn(`[trigger-watcher] Session exited during chain step ${i} turn wait:`, sessionId); - steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }), ...sources }); await writeResult({ ok: false, submitted: chainSubmitted, error: 'session exited during wait', partial: true, steps_completed: i, sessionId, sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs }); return; } if (result.timedOut) { ctx.log.warn(`[trigger-watcher] Chain timeout at step ${i}:`, sessionId); - steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }), ...sources }); await writeResult({ ok: false, submitted: chainSubmitted, error: 'chain timeout', ...(result.waitingSeen ? { reason: REASON_DIALOG_OPEN_AFTER_WRITE } : {}), partial: true, steps_completed: i, sessionId, sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs }); return; } } - steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }), ...sources }); } await writeResult({ @@ -1536,6 +1671,9 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR name, err && err.message); await writeResult({ ok: false, error: 'internal error: ' + (err && err.message), internal: true }); } finally { + if (forgetTranscript) { + try { forgetTranscript(); } catch (_) { /* swallow */ } + } if (releaseSessionLock) releaseSessionLock(); } }