diff --git a/.ai/contexts/trigger-watcher.md b/.ai/contexts/trigger-watcher.md index ff685e57..1cf5f8b2 100644 --- a/.ai/contexts/trigger-watcher.md +++ b/.ai/contexts/trigger-watcher.md @@ -812,6 +812,79 @@ still applies once a rise is observed (unchanged by the rise-wait bound)` is true). The two pre-existing settle-window spec tests above — including the documented deadline-vs-settle trade-off — remain green unmodified. +### Readiness and edge proof from the CLI descriptor (issue #407) + +Field case (2026-10-01): after a `/compact` chain step, the next step's text +landed in the composer but its Enter became a line break, while the log said +`Chain step 1 sent`. Two defects stacked. The level probe +(`pollForBusyObserved`) accepts any `_cliBusy = true` within the window, and +`_cliBusy` comes from OSC titles / OSC 9;4, so the spinner that ends a +compaction satisfied it and the recovery Enter was never armed. And nothing +waited for the CLI to be back at its prompt: #185 (settle), #190 (rise wait) +and the `midBusy` gate all read the same terminal-derived signal, which is +not the CLI's own account of its state (the 0.0.64 incident above: even a +lone retry `\r` seconds later was absorbed). + +The CLI's own descriptor (`~/.claude/sessions/.json`, read by +`cli-session-state.js`, see `cli-session-state.md`) is the better source: +`status` is `idle` at the prompt, `waiting` when a dialog is open, `busy` +while working, and `statusUpdatedAt` is written on change. The watcher reaches +it through the optional `ctx.getCliStatus(sessionId)` (`trigger-context.js`, +wired to `cliSessionState.getStatus` in `main.js`; `undefined` for a remote +session, which has no local descriptor). No new watcher: it reuses the cache +`cli-session-state.js` already keeps. + +- **Readiness.** A chain step that follows a `/compact` step first waits + (`waitForCliIdleAfter`) for `status: "idle"` with a `statusUpdatedAt` later + than the compact step's send time. Bounded by + `SWITCHBOARD_CLI_READY_WAIT_MS` (default 60 000 ms) and by the step's own + deadline; on expiry the step is written anyway, with the warning `CLI not + idle after /compact within N ms, writing chain step N anyway`. Its time is + counted in the step's and the chain's `waited_ms`. +- **Proof of submission by edge.** When a descriptor with an integer + `statusUpdatedAt` is available, a submission counts when the descriptor + shows ANY status write (`busy`, `idle` or `waiting`) with a + `statusUpdatedAt` at or after the moment of our Enter (`cliReactedSince`): + the CLI reacted. `idle` alone covers a turn too fast for a poll to see + `busy`; `waiting` is a permission dialog our Enter opened. A spinner on the + level probe, or a status that began earlier, proves nothing. Otherwise the + existing recovery applies (one bare ` `, only into a free composer, same + window), and the reaction is looked for again. Still nothing: + `confirmed: false`. +- **The recovery Enter is never written while the descriptor reads `waiting` + or `busy`** (`cliForbidsRecoveryEnter`, in edge and fallback modes, whenever + the descriptor has a status, even without a usable timestamp): a bare Enter + would answer the dialog with its default, or land in a running turn. The + step reports `recoverySkipped` and `confirmed: false`. +- **A descriptor whose `statusUpdatedAt` is not an integer** is treated as no + descriptor for readiness and proof (old behaviour), not as one that never + matches. +- **Result and log.** `submitWithVerify` returns `confirmed`: `true` (edge + seen), `false` (descriptor available, no edge even after the recovery Enter) + or `null` (no descriptor: the level probe decides, as before). `true` logs + `Chain step N submitted to ...` (single command: `Submitted command`); + `false` logs the warning `Chain step N not confirmed submitted to ...` + (`Command not confirmed submitted`) and never `sent`; `null` keeps `sent`. + A confirmed step reads `submitted: "confirmed"`; an unconfirmed one reads + `"assumed"`, so the chain fold drops to `assumed` too. The step carries + `submit_confirmed` and the chain result lists `unconfirmed_steps` (indexes) + when there are any. Both fields are absent without a descriptor. +- **What it does not cover.** A command that never makes the CLI busy (a + local slash command that answers at once) cannot show a busy edge: it takes + the recovery `\r` (a no-op on an empty composer) and is reported + unconfirmed. The descriptor is written by the CLI process; the edge is as + fresh as `cli-session-state.js`'s watch of that file. + +Unverified against a real CLI: whether waiting for the descriptor's idle +actually makes the Enter after `/compact` submit. The deciding measurement +(isolated instance, real CLI) is `/compact` then text with +`SWITCHBOARD_SUBMIT_ENTER_DELAY_MS` 50 / 500 / 3000, with and without the +readiness wait. Until then the fix guarantees the failure is recovered once or +named, not that the first Enter always lands. + +Tests: `test/trigger-descriptor-proof.test.js` (fake timers and a fake +descriptor for the helpers; the real watcher for the chain wiring). + ### Why `composerEmptyAfterWrite` cannot be made to prove submission, even by feeding it our own writes A proposal, considered and rejected 2026-09-04: since `submitToPty` writes diff --git a/CHANGELOG.md b/CHANGELOG.md index d2c90046..2d5686a6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ What changes for you in each release of Switchboard. How to write an entry: [doc ## Unreleased ### Fixed +- A step of a trigger chain that follows `/compact` now waits for the CLI to be back at its prompt before it is written, and a step whose Enter did not start a turn is retried once and then reported as "not confirmed submitted" in the log and the result instead of "sent". (#407) - 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) - A session that has exited no longer keeps a busy dot in the sidebar, and the status bar's running count drops as soon as the session ends instead of waiting for the next refresh. (#375) diff --git a/docs/automation.md b/docs/automation.md index b5601d5c..e0c0cc9b 100644 --- a/docs/automation.md +++ b/docs/automation.md @@ -263,6 +263,7 @@ spends the same budget. | `SWITCHBOARD_TRIGGER_MAX_AGE_MS` | The staleness limit | 300 000 | | `SWITCHBOARD_SUBMIT_ENTER_DELAY_MS` | Delay between the text and its Enter | 50 | | `SWITCHBOARD_SUBMIT_VERIFY_MS` | How long a submission is watched for a turn | 2 000 | +| `SWITCHBOARD_CLI_READY_WAIT_MS` | How long a chain step after `/compact` waits for the CLI to report idle | 60 000 | | `SWITCHBOARD_BUSY_FALL_SETTLE_MS` | How long "not busy" must hold between chain steps | 300 | The triggers directory does not move with `SWITCHBOARD_DATA_DIR`: an instance diff --git a/main.js b/main.js index ea0ddd24..2935eb04 100644 --- a/main.js +++ b/main.js @@ -3183,7 +3183,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 })); + require('./trigger-watcher').start(createTriggerContext({ activeSessions, log, getCliStatus: (id) => cliSessionState.getStatus(id) })); } catch (err) { log.error('[trigger-watcher] Failed to start trigger watcher:', err.message); } diff --git a/test/trigger-context.test.js b/test/trigger-context.test.js index 3b8eaf9f..574ff68e 100644 --- a/test/trigger-context.test.js +++ b/test/trigger-context.test.js @@ -137,3 +137,21 @@ test('log is forwarded, and isPtyAlive is only present when supplied', () => { }); assert.equal(withProbe.isPtyAlive, probe); }); + +test('getCliStatus is only present when supplied, and answers for local live sessions only', () => { + assert.equal('getCliStatus' in createTriggerContext({ activeSessions: new Map(), log: silentLog }), false); + + const sessions = new Map([ + ['local', { pty: {}, host: null }], + ['remote', { pty: {}, host: 'box', handle: {} }], + ]); + const seen = []; + const ctx = createTriggerContext({ + activeSessions: sessions, log: silentLog, + getCliStatus: (id) => { seen.push(id); return { status: 'idle', statusUpdatedAt: 5 }; }, + }); + assert.deepEqual(ctx.getCliStatus('local'), { status: 'idle', statusUpdatedAt: 5 }); + assert.equal(ctx.getCliStatus('remote'), undefined, 'a remote session has no local descriptor'); + assert.equal(ctx.getCliStatus('unknown'), undefined); + assert.deepEqual(seen, ['local']); +}); diff --git a/test/trigger-descriptor-proof.test.js b/test/trigger-descriptor-proof.test.js new file mode 100644 index 00000000..61c085ae --- /dev/null +++ b/test/trigger-descriptor-proof.test.js @@ -0,0 +1,380 @@ +// test/trigger-descriptor-proof.test.js +// +// Readiness after /compact and proof of submission by edge, both read from the +// CLI's own descriptor (ctx.getCliStatus). See +// .ai/contexts/trigger-watcher.md, "Readiness and edge proof from the CLI descriptor". +'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'; + +const test = require('node:test'); +const assert = require('node:assert/strict'); +const fs = require('fs'); +const os = require('os'); +const path = require('path'); + +const { + submitWithVerify, waitForCliIdleAfter, isCompactCommand, start, +} = require('../trigger-watcher'); + +const T0 = 1_000_000; + +function fakeSession({ levelBusy = false, withDescriptor = true, onEnter } = {}) { + const desc = { status: 'idle', statusUpdatedAt: T0 - 10_000 }; + const writes = []; + let enters = 0; + const handle = { + write(data) { + writes.push(data); + if (data === '\r') { + enters += 1; + if (onEnter) onEnter(enters, desc); + } + }, + isAlive() { return true; }, + }; + const ctx = { + getPtyForSession: () => ({ ptyProcess: {}, handle }), + isSessionBusy: () => levelBusy, + getComposerState: () => ({ pending: 0, lastInputAt: 0 }), + }; + if (withDescriptor) ctx.getCliStatus = () => ({ ...desc }); + return { desc, writes, handle, ctx }; +} + +async function settle(t, promise, maxMs = 5000) { + let done = false; + let value; + let error; + promise.then((v) => { done = true; value = v; }, (e) => { done = true; error = e; }); + for (let elapsed = 0; elapsed < maxMs && !done; elapsed += 5) { + t.mock.timers.tick(5); + await new Promise((r) => setImmediate(r)); + } + if (error) throw error; + assert.ok(done, `promise still pending after ${maxMs} mocked ms`); + return value; +} + +function enableClock(t) { + t.mock.timers.enable({ apis: ['Date', 'setTimeout'], now: T0 }); +} + +function busyEdgeAfter(ms) { + return (_n, desc) => { + setTimeout(() => { desc.status = 'busy'; desc.statusUpdatedAt = Date.now(); }, ms); + }; +} + +test('edge: descriptor goes busy after our Enter -> confirmed, no recovery Enter', async (t) => { + enableClock(t); + const s = fakeSession({ onEnter: busyEdgeAfter(120) }); + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.deepEqual(s.writes, ['hello', '\r']); + assert.equal(v.confirmed, true); + assert.equal(v.composerConfirmed, true); + assert.equal(v.submit_retries, 0); +}); + +test('edge: a spinner on the level probe and a stale idle descriptor prove nothing -> one recovery Enter, then unconfirmed', async (t) => { + enableClock(t); + const s = fakeSession({ levelBusy: true }); + s.desc.statusUpdatedAt = T0 - 500; + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.deepEqual(s.writes, ['hello', '\r', '\r']); + assert.equal(v.confirmed, false); + assert.equal(v.composerConfirmed, false); + assert.equal(v.sawBusy, false); + assert.equal(v.submit_retries, 1); +}); + +test('edge: first Enter absorbed, the recovery Enter starts the turn -> confirmed after one retry', async (t) => { + enableClock(t); + const s = fakeSession({ onEnter: (n, desc) => { if (n === 2) busyEdgeAfter(120)(n, desc); } }); + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.deepEqual(s.writes, ['hello', '\r', '\r']); + assert.equal(v.confirmed, true); + assert.equal(v.submit_retries, 1); +}); + +test('edge: a fast turn seen only as idle with a newer timestamp still proves the CLI reacted', async (t) => { + enableClock(t); + const s = fakeSession({ onEnter: (_n, desc) => { + setTimeout(() => { desc.status = 'idle'; desc.statusUpdatedAt = Date.now(); }, 120); + } }); + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.deepEqual(s.writes, ['hello', '\r']); + assert.equal(v.confirmed, true); + assert.equal(v.submit_retries, 0); +}); + +test('edge: a dialog opened by our Enter (waiting, newer timestamp) is a submission and gets no recovery Enter', async (t) => { + enableClock(t); + const s = fakeSession({ onEnter: (_n, desc) => { + setTimeout(() => { desc.status = 'waiting'; desc.statusUpdatedAt = Date.now(); }, 120); + } }); + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.deepEqual(s.writes, ['hello', '\r']); + assert.equal(v.confirmed, true); +}); + +for (const status of ['waiting', 'busy']) { + test(`recovery: the descriptor reading "${status}" without a reaction to our Enter never gets a recovery Enter`, async (t) => { + enableClock(t); + const s = fakeSession(); + s.desc.status = status; + s.desc.statusUpdatedAt = T0 - 500; + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.deepEqual(s.writes, ['hello', '\r']); + assert.equal(v.recoverySkipped, true); + assert.equal(v.confirmed, false); + }); +} + +test('recovery: a descriptor without a usable timestamp still forbids the recovery Enter while a dialog is open', async (t) => { + enableClock(t); + const s = fakeSession(); + s.ctx.getCliStatus = () => ({ status: 'waiting', statusUpdatedAt: null }); + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.deepEqual(s.writes, ['hello', '\r']); + assert.equal(v.recoverySkipped, true); + assert.equal(v.confirmed, null); +}); + +test('fallback: a descriptor whose statusUpdatedAt is not an integer is treated as no descriptor', async (t) => { + enableClock(t); + const s = fakeSession({ levelBusy: true }); + s.ctx.getCliStatus = () => ({ status: 'idle', statusUpdatedAt: null }); + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.equal(v.confirmed, null); + assert.equal(v.sawBusy, true); + const r = await settle(t, waitForCliIdleAfter('sid', s.ctx, T0, T0 + 60_000)); + assert.equal(r.available, false); +}); + +test('fallback: no ctx.getCliStatus -> the level probe still decides and confirmed stays null', async (t) => { + enableClock(t); + const s = fakeSession({ levelBusy: true, withDescriptor: false }); + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.deepEqual(s.writes, ['hello', '\r']); + assert.equal(v.sawBusy, true); + assert.equal(v.confirmed, null); + assert.equal(v.submit_retries, 0); +}); + +test('fallback: a descriptor unknown for this session behaves like no descriptor', async (t) => { + enableClock(t); + const s = fakeSession({ levelBusy: true }); + s.ctx.getCliStatus = () => undefined; + const v = await settle(t, submitWithVerify(s.handle, 'sid', 'hello', s.ctx)); + assert.equal(v.confirmed, null); + assert.equal(v.sawBusy, true); +}); + +test('readiness: idle with a statusUpdatedAt after the compact send -> ready at once', async (t) => { + enableClock(t); + const s = fakeSession(); + s.desc.statusUpdatedAt = T0 + 1; + const r = await settle(t, waitForCliIdleAfter('sid', s.ctx, T0, T0 + 60_000)); + assert.equal(r.ready, true); + assert.equal(r.available, true); +}); + +test('readiness: idle older than the compact send, then a later idle -> waits for the later one', async (t) => { + enableClock(t); + const s = fakeSession(); + s.desc.statusUpdatedAt = T0 - 1; + setTimeout(() => { s.desc.status = 'busy'; s.desc.statusUpdatedAt = Date.now(); }, 200); + setTimeout(() => { s.desc.status = 'idle'; s.desc.statusUpdatedAt = Date.now(); }, 1000); + const r = await settle(t, waitForCliIdleAfter('sid', s.ctx, T0, T0 + 60_000)); + assert.equal(r.ready, true); + assert.ok(r.waited_ms >= 1000, `waited ${r.waited_ms} ms, expected to hold until the later idle`); +}); + +test('readiness: a dialog ("waiting") is not idle', async (t) => { + enableClock(t); + const s = fakeSession(); + s.desc.status = 'waiting'; + s.desc.statusUpdatedAt = T0 + 5; + const r = await settle(t, waitForCliIdleAfter('sid', s.ctx, T0, T0 + 1000)); + assert.equal(r.ready, false); + assert.equal(r.timedOut, true); +}); + +test('readiness: never idle -> bounded, reports timedOut', async (t) => { + enableClock(t); + const s = fakeSession(); + s.desc.status = 'busy'; + const r = await settle(t, waitForCliIdleAfter('sid', s.ctx, T0, T0 + 2000)); + assert.equal(r.ready, false); + assert.equal(r.timedOut, true); + assert.ok(r.waited_ms >= 2000); +}); + +test('readiness: no descriptor -> available:false so the caller keeps today\'s behaviour', async (t) => { + enableClock(t); + const s = fakeSession({ withDescriptor: false }); + const r = await settle(t, waitForCliIdleAfter('sid', s.ctx, T0, T0 + 60_000)); + assert.equal(r.available, false); + assert.equal(r.ready, false); + assert.equal(r.timedOut, false); +}); + +test('isCompactCommand: /compact with or without arguments, nothing else', () => { + assert.equal(isCompactCommand('/compact'), true); + assert.equal(isCompactCommand(' /compact keep the plan'), true); + assert.equal(isCompactCommand('/compactify'), false); + assert.equal(isCompactCommand('run /compact'), false); +}); + +// ── chain, through the real watcher ───────────────────────────────────────── + +function mkTmp() { + return fs.realpathSync.native(fs.mkdtempSync(path.join(os.tmpdir(), 'sw-trigger-desc-'))); +} + +function recordingLog() { + const lines = []; + const mk = (level) => (...args) => { lines.push({ level, text: args.join(' ') }); }; + return { lines, info: mk('info'), warn: mk('warn'), error: mk('error'), debug: () => {} }; +} + +function chainSession(sessionId, { log, onEnter }) { + const written = []; + const desc = { status: 'idle', statusUpdatedAt: Date.now() - 10_000 }; + const ptyProcess = { + pid: process.pid, + write(data) { + written.push({ data, at: Date.now() }); + if (data === '\r') onEnter(written.filter((w) => w.data === '\r').length, desc); + }, + }; + let busy = false; + const ctx = { + log, + getPtyForSession: (id) => (id === sessionId ? { ptyProcess } : null), + isSessionBusy: () => busy, + isPtyAlive: () => true, + getComposerState: () => ({ pending: 0, lastInputAt: 0 }), + getCliStatus: (id) => (id === sessionId ? { ...desc } : undefined), + }; + return { ctx, written, desc, setBusy(v) { busy = v; } }; +} + +async function runChain(chain, session, uuid) { + const tmp = mkTmp(); + process.env.SWITCHBOARD_TRIGGERS_DIR = tmp; + process.env.SWITCHBOARD_TRIGGER_IDLE_TIMEOUT_MS = '2000'; + const watcher = start(session.ctx); + try { + fs.writeFileSync(path.join(tmp, uuid + '.json'), + JSON.stringify({ sessionId: uuid, wait: 'idle', chain, timeout_ms: 20000 }), 'utf8'); + const resultPath = path.join(tmp, 'processed', uuid + '.result.json'); + const deadline = Date.now() + 15000; + 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; + delete process.env.SWITCHBOARD_TRIGGER_IDLE_TIMEOUT_MS; + fs.rmSync(tmp, { recursive: true, force: true }); + } +} + +test('chain: the step after /compact is held until the descriptor is idle after the compact, then confirmed by edge', async () => { + const uuid = 'sess-desc-ready-' + Date.now(); + const log = recordingLog(); + let compactIdleAt = null; + const session = chainSession(uuid, { + log, + onEnter(n, desc) { + session.setBusy(true); + if (n === 1) { + desc.status = 'busy'; desc.statusUpdatedAt = Date.now(); + setTimeout(() => { session.setBusy(false); }, 80); + setTimeout(() => { + desc.status = 'idle'; desc.statusUpdatedAt = Date.now(); compactIdleAt = Date.now(); + }, 900); + } else { + desc.status = 'busy'; desc.statusUpdatedAt = Date.now(); + setTimeout(() => { session.setBusy(false); desc.status = 'idle'; desc.statusUpdatedAt = Date.now(); }, 100); + } + }, + }); + + const result = await runChain([{ command: '/compact' }, { command: 'resume the work' }], session, uuid); + + const nextText = session.written.find((w) => w.data === 'resume the work'); + assert.ok(compactIdleAt, 'the compact never reached idle'); + assert.ok(nextText.at >= compactIdleAt, `step 1 written at +${nextText.at - compactIdleAt} ms relative to the compact idle`); + assert.equal(result.ok, true); + assert.equal(result.steps[1].submit_confirmed, true); + assert.equal(result.steps[1].submitted, 'confirmed'); + assert.equal(result.unconfirmed_steps, undefined); + assert.ok(log.lines.some((l) => l.level === 'info' && /Chain step 1 submitted to/.test(l.text))); + assert.ok(!log.lines.some((l) => /Chain step 1 sent/.test(l.text))); +}); + +test('chain: a CLI that never goes idle after /compact -> bounded wait, warning, step still written', async () => { + process.env.SWITCHBOARD_CLI_READY_WAIT_MS = '300'; + try { + const uuid = 'sess-desc-timeout-' + Date.now(); + const log = recordingLog(); + const session = chainSession(uuid, { + log, + onEnter(n, desc) { + if (n === 1) { + desc.status = 'busy'; desc.statusUpdatedAt = Date.now(); + setTimeout(() => session.setBusy(false), 60); + } + }, + }); + session.setBusy(false); + + const started = Date.now(); + const result = await runChain([{ command: '/compact' }, { command: 'resume the work' }], session, uuid); + + const nextText = session.written.find((w) => w.data === 'resume the work'); + assert.ok(nextText, 'the step must still be written after the bounded wait'); + assert.ok(log.lines.some((l) => l.level === 'warn' && /CLI not idle after \/compact/.test(l.text))); + assert.ok(nextText.at - started >= 300, 'the readiness wait must have been honoured up to its bound'); + assert.equal(result.ok, true); + } finally { + delete process.env.SWITCHBOARD_CLI_READY_WAIT_MS; + } +}); + +test('chain: an Enter that never starts a turn is reported "not confirmed submitted", never "sent"', async () => { + process.env.SWITCHBOARD_CLI_READY_WAIT_MS = '200'; + try { + const uuid = 'sess-desc-unconfirmed-' + Date.now(); + const log = recordingLog(); + const session = chainSession(uuid, { + log, + onEnter(n, desc) { + if (n === 1) { + desc.status = 'busy'; desc.statusUpdatedAt = Date.now(); + setTimeout(() => { desc.status = 'idle'; desc.statusUpdatedAt = Date.now(); }, 100); + } + }, + }); + + const result = await runChain([{ command: '/compact' }, { command: 'resume the work' }], session, uuid); + + assert.deepEqual(session.written.map((w) => w.data), ['/compact', '\r', 'resume the work', '\r', '\r']); + assert.ok(log.lines.some((l) => l.level === 'warn' && /Chain step 1 not confirmed submitted/.test(l.text))); + assert.ok(!log.lines.some((l) => /Chain step 1 (sent|submitted)/.test(l.text))); + assert.equal(result.steps[1].submit_confirmed, false); + assert.equal(result.steps[1].submitted, 'assumed'); + assert.deepEqual(result.unconfirmed_steps, [1]); + assert.equal(result.submitted, 'assumed'); + } finally { + delete process.env.SWITCHBOARD_CLI_READY_WAIT_MS; + } +}); diff --git a/trigger-context.js b/trigger-context.js index d083c8ed..d2ad7e6a 100644 --- a/trigger-context.js +++ b/trigger-context.js @@ -24,9 +24,10 @@ function createLocalSessionHandle(ptyProcess) { * @param {Map} deps.activeSessions * @param {object} deps.log electron-log compatible logger * @param {function} [deps.isPtyAlive] (ptyProcess) => boolean + * @param {function} [deps.getCliStatus] (sessionId) => { status, statusUpdatedAt } | undefined * @returns {object} ctx */ -function createTriggerContext({ activeSessions, log, isPtyAlive }) { +function createTriggerContext({ activeSessions, log, isPtyAlive, getCliStatus }) { const ctx = { log, getPtyForSession(sessionId) { @@ -50,6 +51,13 @@ function createTriggerContext({ activeSessions, log, isPtyAlive }) { }, }; if (isPtyAlive) ctx.isPtyAlive = isPtyAlive; + if (getCliStatus) { + ctx.getCliStatus = (sessionId) => { + const session = activeSessions.get(sessionId); + if (!session || session.host != null) return undefined; + return getCliStatus(sessionId); + }; + } return ctx; } diff --git a/trigger-watcher.js b/trigger-watcher.js index 306523e8..8e06be54 100644 --- a/trigger-watcher.js +++ b/trigger-watcher.js @@ -143,7 +143,9 @@ function delayWithBusyPoll(ms, sessionId, ctx) { // then send Enter as a SEPARATE write so it is read as a discrete "submit" // keypress rather than a trailing newline. See DEFAULT_SUBMIT_ENTER_DELAY_MS. // -// Returns `midBusy`: busy polled continuously between the text write and the +// Returns `{ midBusy, enterAt }` (enterAt: the clock just before the Enter +// write, the lower bound the descriptor edge proof is compared against). +// `midBusy`: busy polled continuously between the text write and the // Enter write, true if observed at any point in that window. A single sample // -- at the start, or at the end -- can miss a busy that rises and falls // entirely inside the window, which is a real, reproduced case (see @@ -156,8 +158,9 @@ async function submitToPty(handle, command, sessionId, ctx) { const envMs = envNumber('SWITCHBOARD_SUBMIT_ENTER_DELAY_MS'); const delayMs = envMs !== undefined ? envMs : DEFAULT_SUBMIT_ENTER_DELAY_MS; const midBusy = await delayWithBusyPoll(delayMs, sessionId, ctx); + const enterAt = Date.now(); handle.write('\r'); - return midBusy; + return { midBusy, enterAt }; } // A variable set to the empty string is a launcher artefact, not a value: @@ -313,7 +316,8 @@ function getBusyRiseWaitMs() { * - timedOut: global deadline fired before any observation * - sessionExited: PTY vanished during the poll */ -function pollForBusyObserved(sessionId, ctx, windowMs, deadlineMs) { +function pollForBusyObserved(sessionId, ctx, windowMs, deadlineMs, probe) { + const isBusy = probe || (() => ctx.isSessionBusy(sessionId)); const start = Date.now(); const windowEnd = start + windowMs; @@ -326,7 +330,7 @@ function pollForBusyObserved(sessionId, ctx, windowMs, deadlineMs) { if (!ctx.getPtyForSession(sessionId)) { return resolve({ sawBusy: false, timedOut: false, sessionExited: true, waited_ms: now - start }); } - if (ctx.isSessionBusy(sessionId)) { + if (isBusy()) { return resolve({ sawBusy: true, timedOut: false, sessionExited: false, waited_ms: now - start }); } if (now >= windowEnd) { @@ -337,6 +341,74 @@ function pollForBusyObserved(sessionId, ctx, windowMs, deadlineMs) { }); } +// see .ai/contexts/trigger-watcher.md, "Readiness and edge proof from the CLI descriptor" +const DEFAULT_CLI_READY_WAIT_MS = 60_000; // ms +function getCliReadyWaitMs() { + const v = envNumber('SWITCHBOARD_CLI_READY_WAIT_MS'); + return v !== undefined ? v : DEFAULT_CLI_READY_WAIT_MS; +} + +function readCliStatusRaw(ctx, sessionId) { + if (typeof ctx.getCliStatus !== 'function') return null; + try { + const s = ctx.getCliStatus(sessionId); + return s && typeof s.status === 'string' ? s : null; + } catch (_) { + return null; + } +} + +function readCliStatus(ctx, sessionId) { + const s = readCliStatusRaw(ctx, sessionId); + return s && Number.isInteger(s.statusUpdatedAt) ? s : null; +} + +const CLI_REACTION_STATUSES = ['busy', 'idle', 'waiting']; + +function isCompactCommand(command) { + return /^\/compact(\s|$)/.test(String(command).trim()); +} + +function cliReactedSince(ctx, sessionId, sinceMs) { + const s = readCliStatus(ctx, sessionId); + return !!s && CLI_REACTION_STATUSES.includes(s.status) && s.statusUpdatedAt >= sinceMs; +} + +function cliForbidsRecoveryEnter(ctx, sessionId) { + const s = readCliStatusRaw(ctx, sessionId); + return !!s && (s.status === 'waiting' || s.status === 'busy'); +} + +/** + * Wait until the CLI's own descriptor reports "idle" with a statusUpdatedAt + * later than `afterMs`, bounded by `deadlineMs`. + * + * Returns { ready, available, timedOut, sessionExited, waited_ms }. `available: + * false` means no descriptor could be read (at the start or later): the caller + * keeps its pre-descriptor behaviour. + */ +function waitForCliIdleAfter(sessionId, ctx, afterMs, deadlineMs) { + const start = Date.now(); + return pollLoop((resolve, scheduleNext) => { + const now = Date.now(); + const waited_ms = now - start; + if (!ctx.getPtyForSession(sessionId)) { + return resolve({ ready: false, available: true, timedOut: false, sessionExited: true, waited_ms }); + } + const s = readCliStatus(ctx, sessionId); + if (!s) { + return resolve({ ready: false, available: false, timedOut: false, sessionExited: false, waited_ms }); + } + if (s.status === 'idle' && Number.isInteger(s.statusUpdatedAt) && s.statusUpdatedAt > afterMs) { + return resolve({ ready: true, available: true, timedOut: false, sessionExited: false, waited_ms }); + } + if (now >= deadlineMs) { + return resolve({ ready: false, available: true, timedOut: true, sessionExited: false, waited_ms }); + } + scheduleNext(); + }); +} + /** * Submit a command, then verify it. * @@ -375,7 +447,9 @@ async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs) { // Sampled before the write — see .ai/contexts/trigger-watcher.md ("submitted"). const preBusy = ctx.isSessionBusy(sessionId); - const midBusy = await submitToPty(handle, command, sessionId, ctx); + const edgeMode = readCliStatus(ctx, sessionId) !== null; + + const { midBusy, enterAt } = await submitToPty(handle, command, sessionId, ctx); // Composer read-back: unconditional, immediate, never gated on activity. const postWriteState = (typeof ctx.getComposerState === 'function') @@ -388,12 +462,16 @@ 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 first = await pollForBusyObserved(sessionId, ctx, windowMs, effectiveDeadline); + const probe = edgeMode ? () => cliReactedSince(ctx, sessionId, enterAt) : undefined; + const first = await pollForBusyObserved(sessionId, ctx, windowMs, effectiveDeadline, probe); if (first.sawBusy || first.sessionExited || first.timedOut) { return { submit_retries: 0, sawBusy: first.sawBusy, - composerConfirmed: first.sawBusy && preBusy === false && midBusy === false && composerEmptyAfterWrite, + confirmed: edgeMode ? first.sawBusy : null, + composerConfirmed: edgeMode + ? first.sawBusy + : first.sawBusy && preBusy === false && midBusy === false && composerEmptyAfterWrite, sessionExited: first.sessionExited, timedOut: first.timedOut, waited_ms: first.waited_ms, @@ -403,11 +481,25 @@ async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs) { // Nothing observed in the window — retry the Enter ONCE (bare '\r', never the text), // and only into a free composer. See docs/automation.md. const recoveryDeadline = Math.min(effectiveDeadline, Date.now() + windowMs); + if (cliForbidsRecoveryEnter(ctx, sessionId)) { + return { + submit_retries: 0, + sawBusy: false, + confirmed: edgeMode ? false : null, + composerConfirmed: false, + sessionExited: false, + timedOut: false, + recoverySkipped: true, + recoveryReason: 'the CLI reports a dialog open or a turn running, a bare Enter could answer it', + waited_ms: first.waited_ms, + }; + } const polite = await waitForComposerFree(sessionId, ctx, recoveryDeadline); if (!polite.free) { return { submit_retries: 0, sawBusy: false, + confirmed: edgeMode ? false : null, composerConfirmed: false, sessionExited: false, timedOut: false, @@ -424,6 +516,7 @@ async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs) { return { submit_retries: 1, sawBusy: false, + confirmed: edgeMode ? false : null, composerConfirmed: false, sessionExited: false, timedOut: false, @@ -432,11 +525,12 @@ async function submitWithVerify(handle, sessionId, command, ctx, deadlineMs) { }; } - const second = await pollForBusyObserved(sessionId, ctx, windowMs, effectiveDeadline); + const second = await pollForBusyObserved(sessionId, ctx, windowMs, effectiveDeadline, probe); return { submit_retries: 1, sawBusy: second.sawBusy, - composerConfirmed: false, + confirmed: edgeMode ? second.sawBusy : null, + composerConfirmed: edgeMode && second.sawBusy, sessionExited: second.sessionExited, timedOut: second.timedOut, waited_ms: first.waited_ms + second.waited_ms, @@ -980,9 +1074,11 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR let composerConfirmed = false; let recoverySkipped = false; let recoveryReason = null; + let submitConfirmed = null; try { const v = await submitWithVerify(handle, sessionId, command, ctx); submitRetries = v.submit_retries; + submitConfirmed = v.confirmed; sawBusy = v.sawBusy; composerConfirmed = !!v.composerConfirmed; recoverySkipped = !!v.recoverySkipped; @@ -1000,8 +1096,13 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR if (recoverySkipped) { ctx.log.warn('[trigger-watcher] Recovery Enter withheld for ' + sessionId + ': ' + recoveryReason); } - ctx.log.info(`[trigger-watcher] Sent command to ${sessionId}: ${command}` + - (submitRetries ? ` (submit retried ${submitRetries}x)` : '')); + if (submitConfirmed === false) { + ctx.log.warn(`[trigger-watcher] Command not confirmed submitted to ${sessionId}: ${command}` + + (submitRetries ? ` (submit retried ${submitRetries}x)` : '')); + } else { + ctx.log.info(`[trigger-watcher] ${submitConfirmed ? 'Submitted' : 'Sent'} command to ${sessionId}: ${command}` + + (submitRetries ? ` (submit retried ${submitRetries}x)` : '')); + } await writeResult({ ok: true, @@ -1011,6 +1112,7 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR sent_at: new Date().toISOString(), waited_ms, submit_retries: submitRetries, + ...(submitConfirmed === null ? {} : { submit_confirmed: submitConfirmed }), ...(recoverySkipped ? { reason: recoveryReason } : {}), }); return; @@ -1063,8 +1165,13 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR } } + let compactSentAtMs = null; + const unconfirmedSteps = []; + for (let i = 0; i < chain.length; i++) { const step = chain[i]; + const readyAfterMs = compactSentAtMs; + compactSentAtMs = null; // Check deadline before each step if (Date.now() >= globalDeadline) { @@ -1085,6 +1192,7 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR // Inject the step command const stepSentAt = new Date().toISOString(); if (i === 0) step0SentAt = stepSentAt; + if (isCompactCommand(step.command)) compactSentAtMs = Date.parse(stepSentAt); // Per-step timeout_ms (if set) bounds THIS whole step (verify + retry + the // busy-fall wait for non-final steps), capped by the remaining global @@ -1132,6 +1240,28 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR return; } + let readyWaitedMs = 0; + if (readyAfterMs !== null && !readCliStatus(ctx, sessionId)) { + ctx.log.info(`[trigger-watcher] No usable CLI descriptor for ${sessionId}, readiness wait skipped before chain step ${i}`); + } + if (readyAfterMs !== null && readCliStatus(ctx, sessionId)) { + const readyDeadline = Math.min(stepDeadline, Date.now() + getCliReadyWaitMs()); + const ready = await waitForCliIdleAfter(sessionId, ctx, readyAfterMs, readyDeadline); + readyWaitedMs = ready.waited_ms; + totalWaitedMs += readyWaitedMs; + if (ready.sessionExited) { + ctx.log.warn(`[trigger-watcher] Session exited waiting for the CLI to be ready at chain step ${i}:`, sessionId); + await writeResult({ ok: false, submitted: (i > 0) ? chainSubmitted : SUBMITTED_NO, error: 'session exited during wait', partial: i > 0, steps_completed: i, sessionId, sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs }); + return; + } + if (!ready.available) { + ctx.log.info(`[trigger-watcher] CLI descriptor vanished for ${sessionId}, readiness wait ended early before chain step ${i}`); + } + if (ready.timedOut) { + ctx.log.warn(`[trigger-watcher] CLI not idle after /compact within ${readyWaitedMs} ms, writing chain step ${i} anyway:`, sessionId); + } + } + // Submit the step, then look for activity on the session. // The verify poll IS this step's Phase 1 — for non-final steps we proceed // straight to the busy-FALL wait, never re-observing busy. @@ -1141,7 +1271,7 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR // Includes this step's own politeness wait (see "total_waited_ms" in // docs/automation.md) so steps[i].waited_ms accounts for everything this // step spent, not just the submit-verify portion. - let stepWaitedMs = polite.waited_ms; + let stepWaitedMs = polite.waited_ms + readyWaitedMs; let verify; try { verify = await submitWithVerify(entryHandle, sessionId, step.command, ctx, stepDeadline); @@ -1168,19 +1298,26 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR ctx.log.warn(`[trigger-watcher] Recovery Enter withheld at chain step ${i} for ` + `${sessionId}: ${verify.recoveryReason}`); } - ctx.log.info(`[trigger-watcher] Chain step ${i} sent to ${sessionId}: ${step.command}` + - (submitRetries ? ` (submit retried ${submitRetries}x)` : '')); + const stepConfirmed = verify.confirmed === undefined ? null : verify.confirmed; + if (stepConfirmed === false) { + unconfirmedSteps.push(i); + ctx.log.warn(`[trigger-watcher] Chain step ${i} not confirmed submitted to ${sessionId}: ${step.command}` + + (submitRetries ? ` (submit retried ${submitRetries}x)` : '')); + } else { + ctx.log.info(`[trigger-watcher] Chain step ${i} ${stepConfirmed ? 'submitted' : 'sent'} to ${sessionId}: ${step.command}` + + (submitRetries ? ` (submit retried ${submitRetries}x)` : '')); + } // 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 }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); 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 }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); await writeResult({ ok: false, submitted: chainSubmitted, error: 'chain timeout', partial: true, steps_completed: i, sessionId, sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs }); return; } @@ -1200,20 +1337,20 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR 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 }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); 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 }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); await writeResult({ ok: false, submitted: chainSubmitted, error: 'chain timeout', 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 }); + steps.push({ idx: i, command: step.command, sent_at: stepSentAt, waited_ms: stepWaitedMs, submit_retries: submitRetries, submitted: stepSubmitted, ...(stepConfirmed === null ? {} : { submit_confirmed: stepConfirmed }) }); } await writeResult({ @@ -1223,6 +1360,7 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR sent_at: step0SentAt, steps, total_waited_ms: totalWaitedMs, + ...(unconfirmedSteps.length ? { unconfirmed_steps: unconfirmedSteps } : {}), }); } catch (err) { @@ -1249,6 +1387,8 @@ async function processTriggerFile(name, ctx, triggersDir, processedDir, onEntryR * @param {function} ctx.isSessionBusy (sessionId: string) => boolean * @param {function} [ctx.getComposerState] (sessionId) => { pending, lastInputAt } | null; * absent or null means busy, never free + * @param {function} [ctx.getCliStatus] (sessionId) => { status, statusUpdatedAt } | undefined; + * the CLI's own descriptor; absent keeps the level probe * @param {function} [ctx.isPtyAlive] (ptyProcess) => boolean (default: handle.isAlive()) * @param {object} ctx.log electron-log compatible logger * @returns {{ close(): void }} @@ -1381,4 +1521,4 @@ function start(ctx) { }; } -module.exports = { start, weakestSubmitted, SUBMITTED_RANK, normalizeCwd }; +module.exports = { start, weakestSubmitted, SUBMITTED_RANK, normalizeCwd, submitWithVerify, waitForCliIdleAfter, isCompactCommand };