Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion .github/workflows/upstream-watch.yml
Original file line number Diff line number Diff line change
Expand Up @@ -102,8 +102,9 @@ jobs:
set -e
{
echo "## Upstream watch (exit $code)"
jq -r '"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)"' watch.json
jq -r '"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)"' watch.json
jq -r '(.wouldFire // [])[] | "- would fire \(.id) \(.version) \(.branch)"' watch.json
jq -r '(.deferred // [])[] | "- deferred \(.id) \(.version) \(.branch)"' watch.json
echo; echo '```text'; cat errors.txt; echo '```'
jq -r '.notice.body' watch.json
} >> "$GITHUB_STEP_SUMMARY"
Expand Down
9 changes: 5 additions & 4 deletions scripts/upstream-watch.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -251,7 +251,7 @@ const WEEK = 7 * 86_400_000;
function blindRecord(error, { stdout, stderr, json }, extra = {}) {
stderr.write(`${error}\n`);
const result = {
blind: true, error, records: [], fetchErrors: [], dispatchErrors: [], wouldFire: [], fired: [], parent: null, commit: null, notice: { post: false, body: '' }, ...extra,
blind: true, error, records: [], fetchErrors: [], dispatchErrors: [], wouldFire: [], deferred: [], fired: [], parent: null, commit: null, notice: { post: false, body: '' }, ...extra,
};
if (json) stdout.write(`${JSON.stringify(result, null, 2)}\n`);
return BLIND;
Expand Down Expand Up @@ -286,7 +286,7 @@ async function record(registry, fetcher, options, { stdout, stderr, now, ledgerS
const recorded = ledger.records.map((item) => item.line).join('\n');
const all = ledgerEvents(report, registry, { since });
const released = all.filter((event) => event.event === 'released' && event.fields.branch);
const fired = await dispatch({ released, records: ledger.records, dispatcher, repo, sentinel, now, recordedAt: runAt, dryRun: options.dryRun });
const fired = await dispatch({ released, records: ledger.records, dispatcher, repo, sentinel, now, recordedAt: runAt, dryRun: options.dryRun, 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;
Expand All @@ -306,12 +306,13 @@ async function record(registry, fetcher, options, { stdout, stderr, now, ledgerS
}
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 to the next run: ${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;
}
Expand Down
60 changes: 44 additions & 16 deletions scripts/upstream-watch/dispatch.mjs
Original file line number Diff line number Diff line change
@@ -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 are deferred to the next run. 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';
Expand All @@ -12,14 +14,20 @@ export const FIRE_HEADERS = { 'anthropic-beta': 'experimental-cc-routine-2026-04
export const REFIRE_AFTER_DAYS = 3;
export const MAX_FIRES = 2;
export const FIRE_TIMEOUT_MS = 30_000;
// Sessions start gradually: a few per run, spaced apart; the rest wait for the next run.
export const MAX_FIRES_PER_RUN = 3;
export const FIRE_SPACING_MS = 15_000;
// A 5xx answer means no session was started, so the call is safe to repeat.
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';
const ROUTINE = /^trig_[A-Za-z0-9]+$/;
const DISPATCH_BRANCH = /^upstream\/[a-z0-9._-]+$/;
const DAY = 86_400_000;

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}`);
Expand All @@ -41,25 +49,35 @@ 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(() => null);
const session = body?.claude_code_session_url;
if (response.status === 200 && typeof session === 'string') return session;
const requestId = response.headers?.get?.('request-id');
const detail = body?.error?.message ?? (typeof body?.error === 'string' ? body.error : '');
failure = `the routine trigger answered HTTP ${response.status}${typeof session === 'string' ? '' : ' without a session'}`
+ `${requestId ? ` (request-id ${requestId})` : ''}${detail ? `: ${String(detail).slice(0, 200)}` : ''}`;
if (response.status < 500) 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 }) {
export async function dispatch({ released, records, dispatcher, repo, sentinel, now, recordedAt, dryRun = false, 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) {
Expand All @@ -77,8 +95,18 @@ export async function dispatch({ released, records, dispatcher, repo, sentinel,
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 (triggerFailed || attempted >= MAX_FIRES_PER_RUN) {
deferred.push({ id: event.id, version, branch });
continue;
}
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 });
}
Expand All @@ -93,5 +121,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 };
}
51 changes: 49 additions & 2 deletions tests/kit/upstream-watch-dispatch.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ function fakeDispatcher({ exists = false, pr = null, session = 'https://claude.a
};
}
const run = (dispatcher, records = [], { dryRun, list = [released] } = {}) => dispatch({ released: list, records, dispatcher, repo: 'pacphi/agentic-kit', sentinel: 'UPSTREAM-WATCH', now: NOW, recordedAt: RECORDED_AT, dryRun });
const NOTHING = { records: [], errors: [], wouldFire: [] };
const NOTHING = { records: [], errors: [], wouldFire: [], deferred: [] };

test('a released fix without a branch fires once and is recorded with its session', async () => {
const dispatcher = fakeDispatcher();
Expand Down Expand Up @@ -77,7 +77,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 });
Expand Down Expand Up @@ -108,6 +108,28 @@ test('a fired branch with an open pull request records dispatch-pr once', async
assert.deepEqual(again.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' }) }; };
Expand Down Expand Up @@ -135,6 +157,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 a 5xx answer, 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) => {
Expand Down
3 changes: 2 additions & 1 deletion tests/kit/upstream-watch-record.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,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');
});
});

Expand Down Expand Up @@ -146,7 +147,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);
Expand Down
2 changes: 2 additions & 0 deletions tests/kit/upstream-watch-workflow.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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\./);
});

Expand Down
Loading