diff --git a/.github/workflows/upstream-watch.yml b/.github/workflows/upstream-watch.yml index 907a0cda..38f361c3 100644 --- a/.github/workflows/upstream-watch.yml +++ b/.github/workflows/upstream-watch.yml @@ -58,8 +58,10 @@ jobs: set -e { echo "## Upstream watch preview (exit $code)" - jq -r 'if .blind then "blind \(.blind): \(.error // "unknown error"), could not check \(.fetchErrors | length)" else "since \(.since) (\(.sinceSource)), new records \(.records | length), could not check \(.fetchErrors | length), blind \(.blind), notice \(.notice.post), would fire \(.wouldFire | length)" end' watch.json + jq -r 'if .blind then "blind \(.blind): \(.error // "unknown error"), could not check \(.fetchErrors | length), dispatch errors \(.dispatchErrors | length), deferred \(.deferred | length)" else "since \(.since) (\(.sinceSource)), new records \(.records | length), could not check \(.fetchErrors | length), dispatch errors \(.dispatchErrors | length), blind \(.blind), notice \(.notice.post), would fire \(.wouldFire | length), deferred \(.deferred | length)" end' watch.json jq -r '(.wouldFire // [])[] | "- would fire \(.id) \(.version) \(.branch)"' watch.json + jq -r '(.deferred // [])[] | "- deferred \(.id) \(.version) \(.branch)"' watch.json + jq -r '(.fired // [])[] | "- observed session before ledger failure: \(.id) \(.fields.session)"' watch.json echo; echo '```text'; cat errors.txt; echo '```' } >> "$GITHUB_STEP_SUMMARY" exit "$code" @@ -102,8 +104,10 @@ jobs: set -e { echo "## Upstream watch (exit $code)" - jq -r 'if .blind then "blind \(.blind): \(.error // "unknown error"), could not check \(.fetchErrors | length), dispatch errors \(.dispatchErrors | length)" else "since \(.since) (\(.sinceSource)), new records \(.records | length), could not check \(.fetchErrors | length), dispatch errors \(.dispatchErrors | length), blind \(.blind), commit \(.commit // "none"), would fire \(.wouldFire | length)" end' watch.json + jq -r 'if .blind then "blind \(.blind): \(.error // "unknown error"), could not check \(.fetchErrors | length), dispatch errors \(.dispatchErrors | length), deferred \(.deferred | length)" else "since \(.since) (\(.sinceSource)), new records \(.records | length), could not check \(.fetchErrors | length), dispatch errors \(.dispatchErrors | length), blind \(.blind), commit \(.commit // "none"), would fire \(.wouldFire | length), deferred \(.deferred | length)" end' watch.json jq -r '(.wouldFire // [])[] | "- would fire \(.id) \(.version) \(.branch)"' watch.json + jq -r '(.deferred // [])[] | "- deferred \(.id) \(.version) \(.branch)"' watch.json + jq -r '(.fired // [])[] | "- observed session before ledger failure: \(.id) \(.fields.session)"' watch.json echo; echo '```text'; cat errors.txt; echo '```' jq -r '.notice.body' watch.json } >> "$GITHUB_STEP_SUMMARY" diff --git a/docs/archive/2026-09-29-plan-main-watch-reconciliation.md b/docs/archive/2026-09-29-plan-main-watch-reconciliation.md new file mode 100644 index 00000000..5aaea1fa --- /dev/null +++ b/docs/archive/2026-09-29-plan-main-watch-reconciliation.md @@ -0,0 +1,47 @@ +# Main watcher reconciliation + +## Status + +Implemented and independently reviewed through `7c2f3054`. The controller integrated +validated V6 develop at `3bbee599`; the reviewed watcher files were unchanged by +that merge. All eight local gates passed on that exact clean commit: 6,440 unit +tests passed, seven skipped, legacy suites passed, and browser checks passed. +Coverage was 94.18% lines, 83.80% branches and 93.50% functions. Feature PR CI, +squash integration and final human main review remain separate gates at archival. +No real routine trigger was performed. The limits below describe the worker scope. + +## Scope and acceptance + +Combine develop PR observation and blind reporting with main #280 dispatch pacing. +Preserve pending eligibility, the strict seven day observation window, two recorded +firings per thread, three day cooldown, three fixes per run, and 15 second spacing. +Then restrict trigger retries to documented HTTP 500/503, sanitize token-bearing +error metadata, preserve deferred work in bounded previews and blind results, and +align the living guide and workflow summaries. Use focused synthetic regression +tests before behavior changes and scoped static checks. + +## Ownership and policy receipt + +The controller assigned sole writing ownership of watcher source, workflow, tests, +and guide in `fix/main-watch-reconciliation`, and authorized the prepared two-parent +merge commit followed by three conventional unit commits. Shared manifests, +lockfiles, V6 fixtures, develop/main integration, archive index changes, publication, +and final full gates remain controller-owned. This plan is ready for the controller +to archive with its index update in the completing pull request. + +## Worker limits + +No real trigger, provider turn, external mutation, full suite, UI tests, pnpm, +installed CLI changes, user data writes, or delegated workers. Inject fetch, +execution, and sleeps. Retry evidence does not establish exactly-once execution or +absence of server-side sessions. Deferred work remains eligible for later checks; +repeated earlier failures can starve it. + +## Units + +1. Reconcile and commit both parent contracts; establish focused baseline. +2. Test and fix bounded retry and sanitized error diagnostics. +3. Test and fix preview cap, deferred blind result, and later-run discoverability. +4. Clarify guide, workflow preview/record summaries, and evidence limits. + +Record command evidence and limitations in the ignored worker implementation report. diff --git a/docs/archive/README.md b/docs/archive/README.md index 3929324e..5394fe9c 100644 --- a/docs/archive/README.md +++ b/docs/archive/README.md @@ -69,6 +69,7 @@ reconfirmed by this metadata audit. The per-file inventory and limitations are r | File | Original location | What it was | Why it's historical | |---|---|---|---| | [2026-09-28-plan-follow-ups-v2.md](2026-09-28-plan-follow-ups-v2.md) | `docs/plans/2026-09-28-follow-ups-v2.md` | V4 product, CLI, memory, process-lifecycle and upstream integration follow-ups. | All eight local gates and independent whole-branch review passed at `29152654`; final-head PR CI and squash integration were pending at archival. Conditional Ruflo #3419 guidance remains deferred; native Windows AQE was not tested. | +| [2026-09-29-plan-main-watch-reconciliation.md](2026-09-29-plan-main-watch-reconciliation.md) | `docs/plans/2026-09-29-main-watch-reconciliation.md` | Completed main #280 and develop watcher reconciliation | Both contracts preserved; bounded documented retry, sanitized diagnostics and faithful deferred preview. Independent review and all eight local gates passed through `3bbee599`; PR CI and squash integration remained pending at archival. Current contract: [Upstream watch](../upstream-watch.md). | | [2026-09-28-plan-usage-accuracy.md](2026-09-28-plan-usage-accuracy.md) | `docs/plans/2026-09-28-usage-accuracy.md` | Completed V6 usage and session evidence plan | All 24 units independently accepted; whole-branch review and all eight local gates passed through `1c02db91`. Feature PR CI and develop integration remain separate gates at archival. Current contracts: [Usage metrics](../usage-scorecard-metrics.md) and [ADR-0060](../adr/0060-session-surface-initiator-and-product-names.md). | | [2026-09-29-plan-v6-surfaces-ui.md](2026-09-29-plan-v6-surfaces-ui.md) | `docs/plans/2026-09-29-v6-surfaces-ui.md` | Completed V6 session presentation plan | Shared vocabulary, legacy filter compatibility, independent evidence dimensions and source-coverage disclosures; included in the V6 gates at `1c02db91`. | | [2026-09-29-plan-v6-opencode-cost.md](2026-09-29-plan-v6-opencode-cost.md) | `docs/plans/2026-09-29-v6-opencode-cost.md` | Completed OpenCode reported-zero trust plan | Scoped cost work and mandatory core cache handoff accepted; versioned observation semantics are documented in the living usage guide. | diff --git a/docs/upstream-watch.md b/docs/upstream-watch.md index 32be1c39..39a49d13 100644 --- a/docs/upstream-watch.md +++ b/docs/upstream-watch.md @@ -101,8 +101,10 @@ Every command also takes `--concurrency <1-16>` (default 4) and `--registry event.event === 'released' && event.fields.branch); const eligibleIds = new Set(registry.watch.filter((entry) => PENDING.has(entry.status)).map((entry) => entry.id)); - const fired = await dispatch({ released, records: ledger.records, dispatcher, repo, sentinel, now, recordedAt: runAt, dryRun: options.dryRun, eligibleIds }); + const fired = await dispatch({ released, records: ledger.records, dispatcher, repo, sentinel, now, recordedAt: runAt, dryRun: options.dryRun, eligibleIds, pause: sleep }); const records = [...withoutRecorded(all, recorded).map((event) => toRecord(event, runAt)), ...fired.records]; const checkedAt = fetchErrors.length ? (ledger.checkedAt ?? since) : runAt; let commit = null; @@ -298,20 +298,21 @@ async function record(registry, fetcher, options, { stdout, stderr, now, ledgerS sentences: records.map((item) => commitSafe(sentence(item))), }); } catch (error) { - // The routine already ran for these; without the commit the next run fires again. + // These sessions were observed; without the commit, later runs may fire again. const sessions = fired.records.filter((item) => item.event === 'fired'); for (const item of sessions) stderr.write(`Fired ${item.id} before the ledger commit failed: session ${item.fields.session}\n`); - return blindRecord(`Could not build the ledger commit: ${error.message}`, io, { fetchErrors, dispatchErrors: fired.errors, fired: sessions }); + return blindRecord(`Could not build the ledger commit: ${error.message}`, io, { fetchErrors, dispatchErrors: fired.errors, fired: sessions, deferred: fired.deferred }); } } const body = renderNotice({ records, mention, date: runAt.slice(0, 10), recordedAt: runAt }); const result = { - since, sinceSource, checkedAt, blind: false, records, fetchErrors, dispatchErrors: fired.errors, wouldFire: fired.wouldFire, + since, sinceSource, checkedAt, blind: false, records, fetchErrors, dispatchErrors: fired.errors, wouldFire: fired.wouldFire, deferred: fired.deferred, parent: ledger.commit, commit, notice: { post: Boolean(body), body }, }; for (const item of fetchErrors) stderr.write(`Could not check ${item.id}: ${item.error}\n`); for (const item of fired.errors) stderr.write(`Dispatch ${item.id}: ${item.error}\n`); - const lines = [...records.map((item) => item.line), ...fired.wouldFire.map((item) => `Would fire ${item.id} ${item.version} ${item.branch}`)]; + const lines = [...records.map((item) => item.line), ...fired.wouldFire.map((item) => `Would fire ${item.id} ${item.version} ${item.branch}`), + ...fired.deferred.map((item) => `Deferred for a later check: ${item.id} ${item.version} ${item.branch}`)]; stdout.write(options.json ? `${JSON.stringify(result, null, 2)}\n` : lines.length ? `${lines.join('\n')}\n` : 'No new records.\n'); return 0; } diff --git a/scripts/upstream-watch/dispatch.mjs b/scripts/upstream-watch/dispatch.mjs index f3bebe1e..d4b1e76b 100644 --- a/scripts/upstream-watch/dispatch.mjs +++ b/scripts/upstream-watch/dispatch.mjs @@ -1,8 +1,10 @@ // Dispatch of released upstream fixes (spec 2026-09-28): fire the claude.ai // routine's API trigger for each released fix whose branch does not exist yet, // at most MAX_FIRES times and not again within REFIRE_AFTER_DAYS, and record -// the draft pull request once the branch has one. A dry run lists what would -// fire instead of firing. Injectable exec and fetch. +// the draft pull request once the branch has one. A run fires at most +// MAX_FIRES_PER_RUN fixes, spaced apart, and stops at the first failed call; +// the rest stay eligible for later checks. A dry run lists what would fire +// instead of firing. Injectable exec, fetch and pause. import { eventLine } from './classify.mjs'; import { run } from './fetch.mjs'; import { toRecord } from './ledger-branch.mjs'; @@ -13,6 +15,13 @@ export const REFIRE_AFTER_DAYS = 3; export const PR_OBSERVE_DAYS = 7; export const MAX_FIRES = 2; export const FIRE_TIMEOUT_MS = 30_000; +// Sessions start gradually: a few per run, spaced apart; later checks revisit the rest. +export const MAX_FIRES_PER_RUN = 3; +export const FIRE_SPACING_MS = 15_000; +// The endpoint documents retries for HTTP 500/503, but has no idempotency key. +// These bounded retries do not guarantee that only one server-side session exists. +export const FIRE_ATTEMPTS = 3; +export const FIRE_BACKOFF_MS = 2_000; // `gh pr list --head` matches the branch name in any fork; only a pull request // from this repository is the routine's. export const SAME_REPO_PR = '[.[] | select(.isCrossRepository | not)][0].number // empty'; @@ -20,12 +29,18 @@ const ROUTINE = /^trig_[A-Za-z0-9]+$/; const DISPATCH_BRANCH = /^upstream\/[a-z0-9._-]+$/; const DAY = 86_400_000; +function diagnostic(value, token) { + if (typeof value !== 'string') return ''; + // Redact before bounding, so truncation cannot expose the start of a token. + return value.split(token).join('[REDACTED]').replace(/[\p{Cc}\p{Cf}]/gu, ' ').slice(0, 200); +} + export function sessionList(sessions) { if (sessions.length < 3) return sessions.join(' and '); return `${sessions.slice(0, -1).join(', ')}, and ${sessions.at(-1)}`; } -export function createDispatcher({ exec = run, fetchImpl = globalThis.fetch, env = process.env } = {}) { +export function createDispatcher({ exec = run, fetchImpl = globalThis.fetch, env = process.env, sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)) } = {}) { return { async branchExists(branch) { if (!DISPATCH_BRANCH.test(branch ?? '')) throw new Error(`not a dispatch branch: ${branch}`); @@ -47,25 +62,40 @@ export function createDispatcher({ exec = run, fetchImpl = globalThis.fetch, env const token = env.UPSTREAM_DISPATCH_TOKEN ?? ''; if (!ROUTINE.test(routine)) throw new Error('UPSTREAM_DISPATCH_ROUTINE is not a routine id'); if (!token) throw new Error('UPSTREAM_DISPATCH_TOKEN is not set'); - const response = await fetchImpl(FIRE_URL(routine), { - method: 'POST', headers: { ...FIRE_HEADERS, authorization: `Bearer ${token}` }, body: JSON.stringify({ text }), - signal: AbortSignal.timeout(FIRE_TIMEOUT_MS), - }); - const body = await response.json().catch(() => null); - const session = body?.claude_code_session_url; - if (response.status !== 200 || typeof session !== 'string') { - throw new Error(`the routine trigger answered HTTP ${response.status}${typeof session === 'string' ? '' : ' without a session'}`); + let failure = ''; + for (let attempt = 1; attempt <= FIRE_ATTEMPTS; attempt++) { + if (attempt > 1) await sleep(FIRE_BACKOFF_MS * (attempt - 1)); + const response = await fetchImpl(FIRE_URL(routine), { + method: 'POST', headers: { ...FIRE_HEADERS, authorization: `Bearer ${token}` }, body: JSON.stringify({ text }), + signal: AbortSignal.timeout(FIRE_TIMEOUT_MS), + }); + const body = await response.json().catch((error) => { + // Only a complete non-JSON response uses the status-only fallback. + // A body-stream failure is ambiguous and must not cause another POST. + if (error instanceof SyntaxError) return null; + throw error; + }); + const session = body?.claude_code_session_url; + if (response.status === 200 && typeof session === 'string') return session; + const requestId = diagnostic(response.headers?.get?.('request-id'), token); + const detail = diagnostic(body?.error?.message ?? (typeof body?.error === 'string' ? body.error : ''), token); + failure = `the routine trigger answered HTTP ${response.status}${typeof session === 'string' ? '' : ' without a session'}` + + `${requestId ? ` (request-id ${requestId})` : ''}${detail ? `: ${detail}` : ''}`; + if (typeof session === 'string' || ![500, 503].includes(response.status)) break; } - return session; + throw new Error(failure); }, }; } /** Fire (or, in a dry run, list in `wouldFire`) each released fix; the branch and pull request lookups only read. */ -export async function dispatch({ released, records, dispatcher, repo, sentinel, now, recordedAt, dryRun = false, eligibleIds = null }) { +export async function dispatch({ released, records, dispatcher, repo, sentinel, now, recordedAt, dryRun = false, eligibleIds = null, pause = (ms) => new Promise((resolve) => setTimeout(resolve, ms)) }) { const out = []; const errors = []; const wouldFire = []; + const deferred = []; + let attempted = 0; + let triggerFailed = false; const today = recordedAt.slice(0, 10); const recordsOf = (id, event) => records.filter((item) => item.id === id && item.event === event); for (const event of released) { @@ -79,12 +109,23 @@ export async function dispatch({ released, records, dispatcher, repo, sentinel, } const newest = Math.max(...firings.map((item) => Date.parse(item.recordedAt)), 0); if (newest && now.getTime() - newest < REFIRE_AFTER_DAYS * DAY) continue; + if (triggerFailed || attempted >= MAX_FIRES_PER_RUN) { + deferred.push({ id: event.id, version, branch }); + continue; + } if (dryRun) { + attempted++; wouldFire.push({ id: event.id, version, branch }); continue; } - const session = await dispatcher.fire(`${event.id} ${version} ${branch}`); - out.push(toRecord(eventLine(sentinel, event.id, 'fired', today, { branch, session }), recordedAt)); + if (attempted++) await pause(FIRE_SPACING_MS); + try { + const session = await dispatcher.fire(`${event.id} ${version} ${branch}`); + out.push(toRecord(eventLine(sentinel, event.id, 'fired', today, { branch, session }), recordedAt)); + } catch (error) { + triggerFailed = true; + throw error; + } } catch (error) { errors.push({ id: event.id, error: error.message }); } @@ -102,5 +143,5 @@ export async function dispatch({ released, records, dispatcher, repo, sentinel, errors.push({ id, error: error.message }); } } - return { records: out, errors, wouldFire }; + return { records: out, errors, wouldFire, deferred }; } diff --git a/tests/kit/upstream-watch-dispatch.test.mjs b/tests/kit/upstream-watch-dispatch.test.mjs index 11ee674c..03976d76 100644 --- a/tests/kit/upstream-watch-dispatch.test.mjs +++ b/tests/kit/upstream-watch-dispatch.test.mjs @@ -27,7 +27,7 @@ function fakeDispatcher({ exists = false, pr = null, session = 'https://claude.a }; } const run = (dispatcher, records = [], { dryRun, list = [released], eligibleIds = new Set([ID]), now = NOW } = {}) => dispatch({ released: list, records, dispatcher, repo: 'pacphi/agentic-kit', sentinel: 'UPSTREAM-WATCH', now, recordedAt: RECORDED_AT, dryRun, eligibleIds }); -const NOTHING = { records: [], errors: [], wouldFire: [] }; +const NOTHING = { records: [], errors: [], wouldFire: [], deferred: [] }; test('exhausted firing links read naturally at any count', () => { assert.equal(sessionList(['one']), 'one'); @@ -83,7 +83,7 @@ test('a dry run lists what would fire, calls no trigger and records no firing', const result = await run(dispatcher, [], { dryRun: true }); assert.deepEqual(dispatcher.calls.fire, []); assert.deepEqual(dispatcher.calls.exists, [BRANCH], 'the branch check still runs'); - assert.deepEqual(result, { records: [], errors: [], wouldFire: [{ id: ID, version: '3.13.10', branch: BRANCH }] }); + assert.deepEqual(result, { records: [], errors: [], wouldFire: [{ id: ID, version: '3.13.10', branch: BRANCH }], deferred: [] }); const spent = await run(fakeDispatcher(), [fired('2026-09-20T14:17:00Z', 'https://claude.ai/code/session_a'), fired('2026-09-25T14:17:00Z', 'https://claude.ai/code/session_b')], { dryRun: true }); assert.deepEqual([spent.wouldFire, spent.errors.length], [[], 1], 'the firing limit is still an error'); const lookup = fakeDispatcher({ exists: true, pr: 261 }); @@ -130,6 +130,28 @@ test('PR observation stops on ineligible status or seven days after latest firin assert.deepEqual(boundary.calls.pr, []); }); +const fix = (n) => eventLine('UPSTREAM-WATCH', `proffesor-for-testing/agentic-qe#${n}`, 'released', '2026-08-06', { version: '3.14.5', branch: `upstream/proffesor-for-testing-agentic-qe-${n}` }); +const gentle = (dispatcher, list, pauses = []) => dispatch({ released: list, records: [], dispatcher, repo: 'pacphi/agentic-kit', sentinel: 'UPSTREAM-WATCH', now: NOW, recordedAt: RECORDED_AT, pause: async (ms) => { pauses.push(ms); } }); + +test('a run fires at most three fixes, spaced apart, and defers the rest without an error', async () => { + const dispatcher = fakeDispatcher(); + const pauses = []; + const result = await gentle(dispatcher, [1, 2, 3, 4, 5].map(fix), pauses); + assert.equal(dispatcher.calls.fire.length, 3); + assert.equal(result.records.length, 3); + assert.deepEqual(result.errors, []); + assert.deepEqual(result.deferred.map((item) => item.id), ['proffesor-for-testing/agentic-qe#4', 'proffesor-for-testing/agentic-qe#5']); + assert.equal(pauses.length, 2, 'a pause between fires, none before the first'); +}); + +test('after one failed firing the run stops firing and defers the rest', async () => { + const dispatcher = fakeDispatcher({ fireError: 'the routine trigger answered HTTP 503 without a session' }); + const result = await gentle(dispatcher, [1, 2, 3].map(fix)); + assert.equal(dispatcher.calls.fire.length, 1, 'no further call once the trigger is failing'); + assert.equal(result.errors.length, 1); + assert.equal(result.deferred.length, 2); +}); + test('the trigger call sends the payload with the documented headers and never prints the token', async () => { const requests = []; const fetchImpl = async (url, init) => { requests.push({ url, init }); return { status: 200, json: async () => ({ claude_code_session_url: 'https://claude.ai/code/session_x' }) }; }; @@ -157,6 +179,31 @@ test('the trigger call sends the payload with the documented headers and never p await assert.rejects(createDispatcher({ fetchImpl: empty, env }).fire('x'), /without a session/); }); +test('the trigger call retries HTTP 503, then reports the request id and body', async () => { + const env = { UPSTREAM_DISPATCH_ROUTINE: 'trig_01LmNVKJ4K86joHPvvPtc7yx', UPSTREAM_DISPATCH_TOKEN: 'sk-secret-token' }; + const ok = { status: 200, json: async () => ({ claude_code_session_url: 'https://claude.ai/code/session_r' }) }; + const busy = { status: 503, headers: new Headers({ 'request-id': 'req_1' }), json: async () => ({ error: { message: 'overloaded' } }) }; + const sleeps = []; + const sleep = async (ms) => { sleeps.push(ms); }; + let calls = 0; + const flaky = async () => (++calls < 3 ? busy : ok); + assert.equal(await createDispatcher({ fetchImpl: flaky, env, sleep }).fire('x'), 'https://claude.ai/code/session_r'); + assert.equal(calls, 3); + assert.equal(sleeps.length, 2); + calls = 0; + const down = async () => { calls++; return busy; }; + const error = await createDispatcher({ fetchImpl: down, env, sleep }).fire('x').catch((failure) => failure); + assert.equal(calls, 3, 'gives up after three attempts'); + assert.match(error.message, /HTTP 503/); + assert.match(error.message, /request-id req_1/); + assert.match(error.message, /overloaded/); + assert.doesNotMatch(error.message, /sk-secret-token/); + calls = 0; + const refused = async () => { calls++; return { status: 401, json: async () => ({}) }; }; + await assert.rejects(createDispatcher({ fetchImpl: refused, env, sleep }).fire('x'), /HTTP 401/); + assert.equal(calls, 1, 'a 4xx answer is not retried'); +}); + test('branch and pull request lookups use git and gh without a shell', async () => { const calls = []; const exec = async (command, args) => { @@ -186,3 +233,139 @@ test('a pull request from a fork is never taken for the dispatch pull request', assert.equal(await createDispatcher({ exec: forkOnly }).openPullRequest('pacphi/agentic-kit', BRANCH), null); assert.deepEqual(calls[0].slice(-4), ['--json', 'number,isCrossRepository', '--jq', SAME_REPO_PR]); }); + +// These cases catch retrying an ambiguous or undocumented outcome. +for (const status of [200, 400, 401, 403, 404, 429, 502, 504, 529]) { + test(`HTTP ${status} without a session is not retried`, async () => { + let calls = 0; + const waits = []; + const dispatcher = createDispatcher({ + env: { UPSTREAM_DISPATCH_ROUTINE: 'trig_test', UPSTREAM_DISPATCH_TOKEN: 'test-secret' }, + fetchImpl: async () => { calls++; return { status, json: async () => ({}) }; }, + sleep: async (ms) => { waits.push(ms); }, + }); + await assert.rejects(dispatcher.fire('x'), new RegExp(`HTTP ${status}`)); + assert.equal(calls, 1); + assert.deepEqual(waits, []); + }); +} + +for (const status of [500, 503]) { + test(`HTTP ${status} retries exactly twice with bounded waits`, async () => { + let calls = 0; + const waits = []; + const dispatcher = createDispatcher({ + env: { UPSTREAM_DISPATCH_ROUTINE: 'trig_test', UPSTREAM_DISPATCH_TOKEN: 'test-secret' }, + fetchImpl: async () => { calls++; return { status, json: async () => ({}) }; }, + sleep: async (ms) => { waits.push(ms); }, + }); + await assert.rejects(dispatcher.fire('x'), new RegExp(`HTTP ${status}`)); + assert.equal(calls, 3); + assert.deepEqual(waits, [2000, 4000]); + }); + test(`HTTP ${status} with a session URL never repeats the request`, async () => { + let calls = 0; + const waits = []; + const dispatcher = createDispatcher({ + env: { UPSTREAM_DISPATCH_ROUTINE: 'trig_test', UPSTREAM_DISPATCH_TOKEN: 'test-secret' }, + fetchImpl: async () => { calls++; return { status, json: async () => ({ claude_code_session_url: 'https://claude.ai/code/session_seen' }) }; }, + sleep: async (ms) => { waits.push(ms); }, + }); + await assert.rejects(dispatcher.fire('x'), new RegExp(`HTTP ${status}`)); + assert.equal(calls, 1); + assert.deepEqual(waits, []); + }); +} + +for (const failure of [new Error('connection reset'), new DOMException('timed out', 'TimeoutError')]) { + test(`a thrown ${failure.name} is not retried`, async () => { + let calls = 0; + const waits = []; + const dispatcher = createDispatcher({ + env: { UPSTREAM_DISPATCH_ROUTINE: 'trig_test', UPSTREAM_DISPATCH_TOKEN: 'test-secret' }, + fetchImpl: async () => { calls++; throw failure; }, + sleep: async (ms) => { waits.push(ms); }, + }); + await assert.rejects(dispatcher.fire('x'), { message: failure.message }); + assert.equal(calls, 1); + assert.deepEqual(waits, []); + }); +} + +test('error metadata redacts the configured token before truncation and removes control characters', async () => { + const token = 'synthetic-secret-longer-than-the-remaining-space'; + const dispatcher = createDispatcher({ + env: { UPSTREAM_DISPATCH_ROUTINE: 'trig_test', UPSTREAM_DISPATCH_TOKEN: token }, + fetchImpl: async () => ({ + status: 401, + headers: { get: () => `req\n\x1b${token} ${'r'.repeat(190)}${token}` }, + json: async () => ({ error: { message: `message\r\n${token} ${'m'.repeat(170)}${token}` } }), + }), + }); + const error = await dispatcher.fire('x').catch((failure) => failure); + assert.doesNotMatch(error.message, /synthetic|secret|\p{Cc}/u); + assert.match(error.message, /request-id req.*\[REDACTED\]/); + assert.match(error.message, /message.*\[REDACTED\]/); + assert.ok(error.message.length < 500, 'both metadata fields are bounded'); +}); + +test('non-string error metadata is ignored without coercing objects', async () => { + const dispatcher = createDispatcher({ + env: { UPSTREAM_DISPATCH_ROUTINE: 'trig_test', UPSTREAM_DISPATCH_TOKEN: 'test-secret' }, + fetchImpl: async () => ({ status: 401, headers: { get: () => ({ invalid: true }) }, json: async () => ({ error: { message: { invalid: true } } }) }), + }); + await assert.rejects(dispatcher.fire('x'), { message: 'the routine trigger answered HTTP 401 without a session' }); +}); + +test('a dry run previews three eligible fixes, defers the rest and never sleeps', async () => { + const dispatcher = fakeDispatcher(); + const result = await dispatch({ released: [1, 2, 3, 4, 5].map(fix), records: [], dispatcher, repo: 'pacphi/agentic-kit', sentinel: 'UPSTREAM-WATCH', now: NOW, recordedAt: RECORDED_AT, dryRun: true, pause: async () => { assert.fail('dry run slept'); } }); + assert.deepEqual(result.wouldFire.map((item) => item.id), [1, 2, 3].map((n) => `proffesor-for-testing/agentic-qe#${n}`)); + assert.deepEqual(result.deferred.map((item) => item.id), [4, 5].map((n) => `proffesor-for-testing/agentic-qe#${n}`)); + assert.deepEqual(dispatcher.calls.fire, []); + assert.deepEqual(result.records, []); +}); + +for (const fail of [false, true]) { + test(`PR observation remains bounded after ${fail ? 'a fire failure' : 'the fire cap'}`, async () => { + const dispatcher = fakeDispatcher({ pr: 261, fireError: fail ? 'unavailable' : null }); + const waits = []; + const observed = (n, at) => ({ ...fired(at), id: `proffesor-for-testing/agentic-qe#${n}`, fields: { branch: fix(n).fields.branch, session: `https://claude.ai/code/session_${n}` } }); + const records = [observed(10, '2026-10-01T14:17:00Z'), observed(11, '2026-09-25T14:17:00Z'), observed(12, '2026-10-01T14:17:00Z'), observed(13, '2026-10-01T14:17:00Z'), { ...observed(13, '2026-10-01T14:17:00Z'), event: 'dispatch-pr' }]; + const result = await dispatch({ released: [10, 1, 2, 3, 4, 5].map(fix), records, dispatcher, repo: 'pacphi/agentic-kit', sentinel: 'UPSTREAM-WATCH', now: NOW, recordedAt: RECORDED_AT, eligibleIds: new Set([10, 11, 13].map((n) => `proffesor-for-testing/agentic-qe#${n}`)), pause: async (ms) => { waits.push(ms); } }); + assert.equal(dispatcher.calls.fire.length, fail ? 1 : 3); + assert.deepEqual(waits, fail ? [] : [15000, 15000]); + assert.equal(result.deferred.length, fail ? 4 : 2); + assert.deepEqual(dispatcher.calls.pr, [['pacphi/agentic-kit', 'upstream/proffesor-for-testing-agentic-qe-10']]); + assert.deepEqual(result.records.filter((item) => item.event === 'dispatch-pr').map((item) => item.id), ['proffesor-for-testing/agentic-qe#10']); + }); +} + +for (const status of [500, 503]) { + for (const failure of [new DOMException('body timed out', 'TimeoutError'), new DOMException('body aborted', 'AbortError'), new TypeError('body stream network failure')]) { + test(`HTTP ${status} body ${failure.name} propagates without another POST`, async () => { + let calls = 0; + const waits = []; + const dispatcher = createDispatcher({ + env: { UPSTREAM_DISPATCH_ROUTINE: 'trig_test', UPSTREAM_DISPATCH_TOKEN: 'test-secret' }, + fetchImpl: async () => { calls++; return { status, json: async () => { throw failure; } }; }, + sleep: async (ms) => { waits.push(ms); }, + }); + await assert.rejects(dispatcher.fire('x'), (error) => error === failure); + assert.equal(calls, 1); + assert.deepEqual(waits, []); + }); + } + test(`HTTP ${status} complete malformed JSON retains the bounded status retry`, async () => { + let calls = 0; + const waits = []; + const dispatcher = createDispatcher({ + env: { UPSTREAM_DISPATCH_ROUTINE: 'trig_test', UPSTREAM_DISPATCH_TOKEN: 'test-secret' }, + fetchImpl: async () => { calls++; return new Response('not JSON', { status }); }, + sleep: async (ms) => { waits.push(ms); }, + }); + await assert.rejects(dispatcher.fire('x'), new RegExp(`HTTP ${status} without a session`)); + assert.equal(calls, 3); + assert.deepEqual(waits, [2000, 4000]); + }); +} diff --git a/tests/kit/upstream-watch-record.test.mjs b/tests/kit/upstream-watch-record.test.mjs index 6332946d..30c1ec27 100644 --- a/tests/kit/upstream-watch-record.test.mjs +++ b/tests/kit/upstream-watch-record.test.mjs @@ -75,6 +75,7 @@ test('a released fix fires the routine with its id, version and branch, and the assert.match(result.notice.body, /\nThe full record: `node scripts\/upstream-watch\.mjs ledger --recorded-since 2026-09-26T23:00:00Z`\n$/); assert.deepEqual(result.dispatchErrors, []); assert.deepEqual(result.wouldFire, [], 'a real run lists nothing it would fire'); + assert.deepEqual(result.deferred, [], 'nothing deferred within the per-run cap'); }); }); @@ -153,7 +154,7 @@ test('a ledger commit that cannot be built after a firing keeps the session link test('record is blind (exit 3) when gh, the ledger, the registry or every upstream thread fails', async () => { await withRegistryFile([entry('ruvnet/ruflo#3153', { relation: 'commented' })], async (file) => { const offline = await record(file, [], { fetcher: fixtureFetcher({ authenticated: false }) }); - assert.deepEqual([offline.code, offline.result.blind, offline.ledgerStore.built, offline.result.wouldFire, offline.result.fired], [3, true, [], [], []]); + assert.deepEqual([offline.code, offline.result.blind, offline.ledgerStore.built, offline.result.wouldFire, offline.result.deferred, offline.result.fired], [3, true, [], [], [], []]); const broken = { read: async () => { throw new Error('events.ndjson line 2 is not JSON'); }, build: async () => { throw new Error('not called'); } }; const unreadable = await record(file, [], { ledgerStore: broken }); assert.equal(unreadable.code, 3); @@ -199,3 +200,55 @@ test('invalid registry takes precedence over a future --since', async () => { assert.doesNotMatch(err, /--since is in the future/); }); }); + +// Reuse the recorded release fact for distinct synthetic issue ids; all I/O is injected. +const backlog = () => [1, 2, 3, 4, 5].map((n) => entry(`proffesor-for-testing/agentic-qe#${n}`, { doneWhen: { state: 'closed-completed', release: { channel: 'npm', name: 'agentic-qe', minVersion: '3.13.10' } } })); +const backlogFetcher = () => { + const fetcher = fixtureFetcher(); + return { ...fetcher, thread: async () => fetcher.thread('proffesor-for-testing/agentic-qe#617') }; +}; + +test('advanced Checked-At preserves deferred releases for a later run', async () => { + await withRegistryFile(backlog(), async (file) => { + const first = await record(file, [], { fetcher: backlogFetcher() }); + assert.equal(first.code, 0); + assert.equal(first.dispatcher.fired.length, 3); + assert.deepEqual(first.result.deferred.map((item) => item.id), ['proffesor-for-testing/agentic-qe#4', 'proffesor-for-testing/agentic-qe#5']); + const second = await record(file, [], { fetcher: backlogFetcher(), ledgerStore: memoryLedger({ commit: first.result.commit, records: first.result.records, checkedAt: first.result.checkedAt }), now: new Date('2026-09-27T23:00:00Z') }); + assert.equal(second.result.since, '2026-09-26T23:00:00Z'); + assert.equal(second.dispatcher.fired.length, 2); + assert.deepEqual(second.result.records.map((item) => [item.id, item.event]), [['proffesor-for-testing/agentic-qe#4', 'fired'], ['proffesor-for-testing/agentic-qe#5', 'fired']]); + assert.deepEqual(second.result.deferred, []); + }); +}); + +test('ledger build failure retains successful firings, dispatch errors and the deferred backlog', async () => { + await withRegistryFile(backlog(), async (file) => { + const dispatcher = fakeDispatcher(); + const fire = dispatcher.fire; + dispatcher.fire = async (text) => { if (dispatcher.fired.length === 1) throw new Error('HTTP 503'); return fire(text); }; + const ledgerStore = { read: async () => ({ commit: null, records: [], checkedAt: null }), build: async () => { throw new Error('disk full'); } }; + const { code, result } = await record(file, [], { fetcher: backlogFetcher(), dispatcher, ledgerStore }); + assert.equal(code, 3); + assert.equal(result.blind, true); + assert.equal(result.fired.length, 1); + assert.equal(result.fired[0].fields.session, 'https://claude.ai/code/session_new'); + assert.deepEqual(result.dispatchErrors, [{ id: 'proffesor-for-testing/agentic-qe#2', error: 'HTTP 503' }]); + assert.deepEqual(result.deferred.map((item) => item.id), ['proffesor-for-testing/agentic-qe#3', 'proffesor-for-testing/agentic-qe#4', 'proffesor-for-testing/agentic-qe#5']); + }); +}); + +test('record dry run reports its bounded preview and deferred backlog in JSON and console', async () => { + await withRegistryFile(backlog(), async (file) => { + const preview = await record(file, ['--dry-run'], { fetcher: backlogFetcher() }); + assert.equal(preview.result.wouldFire.length, 3); + assert.equal(preview.result.deferred.length, 2); + assert.deepEqual(preview.dispatcher.fired, []); + assert.deepEqual(preview.ledgerStore.built, []); + const out = capture(); + const code = await main(['record', '--dry-run', '--registry', file], { fetcher: backlogFetcher(), ledgerStore: memoryLedger(), dispatcher: fakeDispatcher(), sleep: async () => { assert.fail('preview slept'); }, stdout: out.stream, stderr: capture().stream, now: NOW }); + assert.equal(code, 0); + assert.equal(out.text().match(/^Would fire /gm)?.length, 3); + assert.equal(out.text().match(/^Deferred /gm)?.length, 2); + }); +}); diff --git a/tests/kit/upstream-watch-workflow.test.mjs b/tests/kit/upstream-watch-workflow.test.mjs index 6a659153..e14dd237 100644 --- a/tests/kit/upstream-watch-workflow.test.mjs +++ b/tests/kit/upstream-watch-workflow.test.mjs @@ -59,6 +59,8 @@ test('both summaries say how many routine sessions a run would start, and which' assert.match(step, /would fire \\\(\.wouldFire \| length\)/, name); assert.match(step, /jq -r '\(\.wouldFire \/\/ \[\]\)\[\] \| "- would fire \\\(\.id\) \\\(\.version\) \\\(\.branch\)"' watch\.json/, name); } + assert.match(record, /deferred \\\(\.deferred \| length\)/, 'the summary counts fixes deferred to the next run'); + assert.match(record, /jq -r '\(\.deferred \/\/ \[\]\)\[\] \| "- deferred \\\(\.id\) \\\(\.version\) \\\(\.branch\)"' watch\.json/); assert.match(record, /# 3 = blind: [^\n]*\n\s*# or a ledger commit that could not be built\./); });