diff --git a/scripts/pm/fleet-write/dispatch.mjs b/scripts/pm/fleet-write/dispatch.mjs index 0db3add5674..18e88ff3e5a 100644 --- a/scripts/pm/fleet-write/dispatch.mjs +++ b/scripts/pm/fleet-write/dispatch.mjs @@ -112,16 +112,80 @@ * `write-pace.mjs` counts the dispatch as THE write — one POST, paced and * leased like any other; the runner's executor paces its own writes there. * + * ## The read-back — the run's `success` is the executor's, not the write's + * + * A run concluding `success` proves the executor's requests were ANSWERED 2xx. + * It does not prove the board holds the bytes this seat sent: measured, a + * 41,699-byte `issue_patch` body went through a `success` run and was stored + * with two multi-byte characters replaced by U+FFFD, deterministically on a + * re-send — and this file said exit 0 both times. The executor cannot catch + * that from where it stands: whatever reached it is what it wrote and what the + * platform echoed, so only the SEAT, which holds the bytes it meant, can. + * + * So after a `success` run `sendFleetWrite` reads back every action that + * carried a `body` — `BODY_OPS`, derived from `ops.mjs`, so a new body op is + * read back the day it lands (`READ_BACK_LOCATORS` must name it, pinned) — and + * judges the stored bytes against the bytes sent with post-stamped's own + * `classifyReadBack` / `sentBodyLanded`: the exact-bytes classifier whose + * declared normalisations (the platform's footer block, a stripped trailing + * newline) are the ones this fleet measured, ⛔ never a second copy here. It is + * loaded with a dynamic import because post-stamped imports THIS file + * statically; a static import back would be a cycle whose evaluation order + * decides whether either file's constants exist yet. One cell is added, for + * the one surface post-stamped never writes: a `pr_create` body the platform + * stored as the bytes sent followed by exactly the session-URL footer block + * (`PR_CREATE_FOOTER_TAIL`) — `platform-readings.md`'s fourth shape, where the + * sent body is a strict prefix of the stored one. It forgives an append after + * EVERY byte sent and nothing inside them. + * + * Each action's verdict is one of three, and the locator decides which a + * missing match may be: + * + * landed every byte sent is on the platform (identical, or a declared + * normalisation). + * not-stored the object the action wrote was READ and a byte sent is not + * the byte stored at that offset — the result is `failure` + * with `notStored: true`, exit `EXIT_NOT_STORED` (4), naming the + * op, the object and the first differing byte. + * unverified the stored body could not be read, or — for a `comment`, + * which only its content can find — no comment created since + * the dispatch holds the bytes sent. Result `unverified`, + * exit `EXIT_UNCONFIRMED` (6): an object that was not found was + * not measured, so ⛔ it is never called "not stored", and ⛔ never + * success either. + * + * ⛔ No retry on a mismatch: a split that happened once happens again on an + * identical re-send (measured), and a body edit retried writes the damage over + * and over. `4` means what it means in every other fleet tool — the write + * HAPPENED and the board disagrees: go READ it. + * + * How the callers read the new outcomes — one exit vocabulary, no caller edited: + * - `failure` + `notStored` keeps the state every caller already reads as "the + * platform did not keep it whole — go READ, never fall back": post-stamped's + * `relayExitFor` maps `failure` to its own `EXIT_NOT_STORED` (4); the tools + * that call `exitForResult` (issue-create, label-write, issue-transfer, + * close-cards) get 4 from it — issue-create's `EXIT_READ_BACK_MISMATCH`. + * - `unverified` is a state no caller names, so each falls to its non-success + * branch: `exitForResult` answers 6, post-stamped's `relayExitFor` answers + * its default 6. Neither is ever a fall-back to direct (that is `no-run`'s + * alone) and neither is 0. + * - a stroke carrying no body (labels, assignees, state, a transfer, the + * GraphQL ops) reads nothing back and its outcome is unchanged. + * * ## Exit codes — capture them BEFORE any pipe * - * 0 the run completed with conclusion `success`. + * 0 the run completed with conclusion `success`, and every body it wrote + * reads back as the bytes sent. * 2 usage, or the packed payload was refused by the validator (nothing sent). * 3 PREREQUISITE NOT MET — no token, no session in dispatch mode, the * dispatch route dead (401/407/0), or the run list unreadable. + * 4 NOT STORED — the run succeeded and the read-back shows a body the board + * does not hold as sent (op, object and first differing byte printed). + * Go READ it; ⛔ do not retry. * 5 the platform REFUSED the dispatch (403/404/422), or the run FAILED. - * 6 UNCONFIRMED — no run in the start window, or no completion within - * the ceiling. The run URL (when any) is printed. Go READ it; ⛔ do not - * retry blind. + * 6 UNCONFIRMED — no run in the start window, no completion within the + * ceiling, or a body the read-back could not find or read. The run URL + * (when any) is printed. Go READ it; ⛔ do not retry blind. * 10 the write throttle refused the dispatch. */ @@ -135,8 +199,8 @@ import { isEntrypoint } from '../../invoked-as.mjs'; import { EXIT_PREREQUISITE_NOT_MET, PROXY_FLAG, proxyRearmPlan, proxyRoute } from '../check-half-states.mjs'; import { classifyHttp } from '../label-write.mjs'; import { EXIT_WRITE_PACE_REFUSED, isWriteMethod, noteResponse, paceWrite, releaseWriteLease } from '../write-pace.mjs'; -import { CONTAINER_SESSION_ENV, CONTAINER_SESSION_SHAPE, MAX_REQUEST_ID_CHARS, RELAY_EVENT_TYPE, RELAY_REPO, SESSION_ENV, SESSION_SHAPE, TRANSPORTS, TRANSPORT_ENV } from './ops.mjs'; -import { refusalText, validatePayload } from './validate.mjs'; +import { CONTAINER_SESSION_ENV, CONTAINER_SESSION_SHAPE, MAX_REQUEST_ID_CHARS, OPS, OP_NAMES, RELAY_EVENT_TYPE, RELAY_REPO, SESSION_ENV, SESSION_SHAPE, TRANSPORTS, TRANSPORT_ENV } from './ops.mjs'; +import { PAYLOAD_ENV, refusalText, validatePayload } from './validate.mjs'; const SELF_PATH = fileURLToPath(import.meta.url); const DEFAULT_API = 'https://api.github.com'; @@ -145,6 +209,12 @@ const PROXY_REARM_GUARD = 'OS_FLEET_DISPATCH_PROXY_REARMED'; export const EXIT_OK = 0; export const EXIT_USAGE = 2; export const EXIT_PREREQUISITE = EXIT_PREREQUISITE_NOT_MET; +/** + * The run succeeded and the read-back READ a body the board does not hold as + * sent — the value every fleet tool gives "written, and the board disagrees": + * post-stamped's `EXIT_NOT_STORED`, issue-create's `EXIT_READ_BACK_MISMATCH`. + */ +export const EXIT_NOT_STORED = 4; export const EXIT_PLATFORM_REFUSAL = 5; /** The one exit that means "dispatched, outcome not confirmed" — shared by every tool that takes the relay. */ export const EXIT_UNCONFIRMED = 6; @@ -395,6 +465,135 @@ export function fallbackText(result, tool = 'fleet-write') { ); } +// --------------------------------------------------------------------------- +// The read-back — pure halves (the header's read-back section is the authority) +// --------------------------------------------------------------------------- + +/** The clock slack a created object's timestamp may carry before the dispatch — post-stamped's relay read-back allows the same minute. */ +export const READ_BACK_SLACK_MS = 60_000; + +/** Every op the table lets carry a `body` — derived from `ops.mjs`, never listed, so a new body op is read back the day it lands. */ +export const BODY_OPS = Object.freeze(OP_NAMES.filter((op) => [...OPS[op].required, ...OPS[op].optional].includes('body'))); + +/** + * Where each body op's stored body is read from, and how it is FOUND — which + * decides what a missing match may mean. `address`: the action named the + * object. `key`: the run created it and a key the action sent (the title, the + * head) finds it. Found either way, the object is identified independently of + * its body, so a difference is a measurement of THAT object: `not-stored`. + * `content`: only the body can find it (a new comment), so a body that matches + * nothing is `unverified` — an object that was not found was not measured. + * The self-test pins that this table names exactly `BODY_OPS`. + */ +export const READ_BACK_LOCATORS = Object.freeze({ + issue_patch: Object.freeze({ found: 'address', where: 'GET /repos/{repo}/issues/{issue}' }), + comment_edit: Object.freeze({ found: 'address', where: 'GET /repos/{repo}/issues/comments/{comment_id}' }), + issue_create: Object.freeze({ found: 'key', where: 'the newest issue created since the dispatch whose title is the title sent' }), + pr_create: Object.freeze({ found: 'key', where: 'the newest pull request on the head sent, created since the dispatch' }), + comment: Object.freeze({ found: 'content', where: 'the newest comment on the issue, created since the dispatch, whose stored body holds the bytes sent' }), +}); + +/** + * The tail the platform appends to a `pr_create` body that does not already + * end in the block — `.claude/skills/pm-dispatch/references/platform-readings.md`, + * the fourth shape: a rule plus ONE session-URL footer, the sent body a strict + * prefix of the stored one, 90/91 bytes (the session id's length). A pattern + * only because the session id varies; everything around it is literal. + */ +export const PR_CREATE_FOOTER_TAIL = /^\n+---\n_Generated by \[Claude Code\]\(https:\/\/claude\.ai\/code\/session_[A-Za-z0-9]+\)_$/u; + +/** The read-back's judge — post-stamped's exact-bytes classifier, loaded when a body is read back. ⛔ Never restated here (header). */ +export async function loadReadBackJudge() { + const judge = await import('../post-stamped.mjs'); + return { classifyReadBack: judge.classifyReadBack, sentBodyLanded: judge.sentBodyLanded }; +} + +/** Every action in a judged payload that carried a body, with the keys its locator needs. Pure. */ +export function readBackTargets(payload) { + const out = []; + const actions = Array.isArray(payload?.actions) ? payload.actions : []; + actions.forEach((a, i) => { + if (!BODY_OPS.includes(a?.op) || typeof a.body !== 'string') return; + const target = { action: i + 1, op: a.op, sent: a.body }; + for (const key of ['issue', 'comment_id', 'title', 'head']) if (key in a) target[key] = a[key]; + out.push(target); + }); + return out; +} + +const replacementCount = (s) => (String(s ?? '').match(/\uFFFD/gu) ?? []).length; + +/** + * One stored body against the bytes sent. Pure: `judge` is post-stamped's pair + * (`loadReadBackJudge`). Returns `{ verdict, cls, offset, sentBytes, + * storedBytes, sentContext, storedContext, replacements }` — `verdict` is + * `landed` or `not-stored` for a body that was read, `unverified` for one that + * was not a string. `replacements` is how many MORE U+FFFD the stored body + * carries than the sent one: the mark of a byte-level split upstream, printed + * because the context window drops a replacement glyph at its edges. + */ +export function judgeStoredBody({ op, sent, stored }, judge) { + const sentText = String(sent ?? ''); + const sentBytes = Buffer.byteLength(sentText, 'utf8'); + if (typeof stored !== 'string') return { verdict: 'unverified', cls: 'unreadable', offset: null, sentBytes, storedBytes: null, replacements: 0 }; + const storedBytes = Buffer.byteLength(stored, 'utf8'); + const rb = judge.classifyReadBack({ sent: sentText, stored }); + if (judge.sentBodyLanded(rb)) return { verdict: 'landed', cls: rb.class, offset: null, sentBytes, storedBytes, replacements: 0 }; + const content = sentText.replace(/\n+$/u, ''); + if (op === 'pr_create' && stored.startsWith(content) && PR_CREATE_FOOTER_TAIL.test(stored.slice(content.length))) { + return { verdict: 'landed', cls: 'pr-create-footer-appended', offset: null, sentBytes, storedBytes, replacements: 0 }; + } + return { + verdict: 'not-stored', + cls: rb.class, + offset: rb.offset, + sentBytes, + storedBytes, + sentContext: rb.sentContext, + storedContext: rb.storedContext, + replacements: Math.max(0, replacementCount(stored) - replacementCount(sentText)), + }; +} + +/** The stroke's verdict from its rows: any `not-stored` wins, then any `unverified`; `none` when nothing carried a body. Pure. */ +export function strokeReadBackState(rows) { + if (!rows.length) return 'none'; + if (rows.some((r) => r.verdict === 'not-stored')) return 'not-stored'; + if (rows.some((r) => r.verdict === 'unverified')) return 'unverified'; + return 'landed'; +} + +/** The transcript line for one read-back row. Pure. */ +export function readBackLine(row) { + const head = `fleet-write: read-back action ${row.action} ${row.op}${row.where ? ` ${row.where}` : ''}`; + if (row.verdict === 'landed') return `${head}: ${row.sentBytes} byte(s) sent, ${row.storedBytes} stored — ${row.cls}; every byte sent is on the platform.`; + if (row.verdict === 'unverified') return `${head}: ⚠️ UNVERIFIED — ${row.why ?? `the stored body could not be read (${row.cls})`}.`; + const fffd = row.replacements > 0 ? ` The stored body carries ${row.replacements} U+FFFD replacement character(s) the sent one does not — the mark of a byte-level split upstream.` : ''; + return `${head}: ✗ NOT STORED — sent ${row.sentBytes} byte(s), stored ${row.storedBytes}; first difference at byte ${row.offset}: sent ${row.sentContext} | stored ${row.storedContext}.${fffd}`; +} + +/** What the CLI prints on a `failure` the read-back measured. */ +export function notStoredText(result, tool = 'fleet-write') { + const bad = (result.readBack?.rows ?? []).filter((r) => r.verdict === 'not-stored'); + const which = bad.map((r) => `action ${r.action} (${r.op}${r.where ? ` ${r.where}` : ''}) first differs at byte ${r.offset}`).join('; '); + return ( + `${tool}: NOT STORED — run ${result.run?.id ?? '?'} concluded success, and the read-back shows the board does NOT hold the bytes sent: ${which}.` + + `${result.run?.url ? ` Run: ${result.run.url}` : ''}\n` + + ' The write HAPPENED. Go READ what is stored and send a body that can land. ⛔ Do not re-run blind: an identical re-send\n' + + ` reproduces a deterministic split, and a body edit retried writes the damage again. Exit ${EXIT_NOT_STORED}.` + ); +} + +/** What the CLI prints when the run succeeded and a body could not be read back. */ +export function unverifiedText(result, tool = 'fleet-write') { + const open = (result.readBack?.rows ?? []).filter((r) => r.verdict === 'unverified'); + return ( + `${tool}: UNCONFIRMED — run ${result.run?.id ?? '?'} concluded success, but ${open.length} body/bodies could not be read back: ` + + `${open.map((r) => `action ${r.action} (${r.op}) — ${r.why ?? r.cls}`).join('; ')}.${result.run?.url ? ` Run: ${result.run.url}` : ''}\n` + + ` Go READ the target; ⛔ do not re-run blind — a second dispatch is a second write. Exit ${EXIT_UNCONFIRMED}.` + ); +} + // --------------------------------------------------------------------------- // Transport — the one POST, paced; the reads around it are not. // --------------------------------------------------------------------------- @@ -439,12 +638,97 @@ async function rest(api, path, { method = 'GET', body = null } = {}, deps = {}) return { status: res.status, rateRemaining: rateRemaining === null ? null : Number(rateRemaining), json, detail: typeof json?.message === 'string' ? json.message : '', call: `${method} ${path}` }; } +/** Every page of a list endpoint (at most `maxPages`), or the first failing answer. Reads only — never paced. */ +async function listAll(api, path, t, maxPages = 10) { + const rows = []; + for (let page = 1; page <= maxPages; page++) { + const r = await rest(api, `${path}${path.includes('?') ? '&' : '?'}per_page=100&page=${page}`, {}, t); + if (r.status !== 200 || !Array.isArray(r.json)) return { ok: false, rows, failed: r }; + rows.push(...r.json); + if (r.json.length < 100) break; + } + return { ok: true, rows }; +} + +/** + * Read back every body a completed stroke wrote and judge it against the bytes + * sent. Returns `{ state, rows }` — `state` from `strokeReadBackState`. Never + * throws on a status: a read that fails makes its row `unverified`, with the + * call and status in `why`. `taken` keeps two actions of one stroke from being + * judged against the same created object. + */ +export async function readBackStroke(payload, { dispatchedAt }, deps = {}) { + const targets = readBackTargets(payload); + if (!targets.length) return { state: 'none', rows: [] }; + const api = deps.api ?? DEFAULT_API; + const t = { fetch: deps.fetch, token: deps.token }; + const judge = deps.judge ?? (await loadReadBackJudge()); + const since = dispatchedAt - READ_BACK_SLACK_MS; + const sinceIso = encodeURIComponent(new Date(since).toISOString()); + const created = (row) => Date.parse(row?.created_at ?? '') >= since; + const repo = payload.repo; + const taken = new Set(); + const rows = []; + const unread = (target, where, r) => ({ action: target.action, op: target.op, where, verdict: 'unverified', cls: 'unread', why: `${r.call} -> HTTP ${r.status}${r.detail ? ` (${r.detail})` : ''}` }); + for (const target of targets) { + const judged = (where, stored) => ({ action: target.action, op: target.op, where, ...judgeStoredBody({ op: target.op, sent: target.sent, stored }, judge) }); + if (target.op === 'issue_patch' || target.op === 'comment_edit') { + const path = target.op === 'issue_patch' ? `/repos/${repo}/issues/${target.issue}` : `/repos/${repo}/issues/comments/${target.comment_id}`; + const where = target.op === 'issue_patch' ? `${repo}#${target.issue}` : `${repo} comment ${target.comment_id}`; + const r = await rest(api, path, {}, t); + rows.push(r.status === 200 && r.json ? judged(where, r.json.body) : unread(target, where, r)); + continue; + } + if (target.op === 'issue_create' || target.op === 'pr_create') { + const isPr = target.op === 'pr_create'; + const head = isPr ? (String(target.head).includes(':') ? String(target.head) : `${repo.split('/')[0]}:${target.head}`) : null; + const path = isPr + ? `/repos/${repo}/pulls?state=all&head=${encodeURIComponent(head)}&sort=created&direction=desc` + : `/repos/${repo}/issues?state=all&sort=created&direction=desc&since=${sinceIso}`; + const listed = await listAll(api, path, t, isPr ? 1 : 3); + const hit = listed.rows.find((o) => created(o) && !taken.has(o.id) && (isPr ? true : !o.pull_request && String(o.title ?? '').trim() === String(target.title ?? '').trim())); + if (!listed.ok && !hit) { + rows.push(unread(target, isPr ? `${repo} head ${head}` : `${repo} (new issue)`, listed.failed)); + continue; + } + if (!hit) { + rows.push({ action: target.action, op: target.op, where: isPr ? `${repo} head ${head}` : `${repo} (new issue)`, verdict: 'unverified', cls: 'unfound', why: isPr ? `no pull request on ${head} was created since the dispatch` : `no issue created on ${repo} since the dispatch carries the title sent` }); + continue; + } + taken.add(hit.id); + rows.push(judged(`${repo}#${hit.number}`, hit.body)); + continue; + } + // `comment`: only the body finds a new comment, so a body that matches nothing is unverified, never not-stored. + const where = `${repo}#${target.issue} (new comment)`; + const listed = await listAll(api, `/repos/${repo}/issues/${target.issue}/comments?since=${sinceIso}`, t); + const candidates = listed.rows.filter((c) => created(c) && !taken.has(c.id)); + const hit = [...candidates].reverse().find((c) => judgeStoredBody({ op: target.op, sent: target.sent, stored: c.body }, judge).verdict === 'landed'); + if (hit) { + taken.add(hit.id); + rows.push(judged(`${repo}#${target.issue} comment ${hit.id}`, hit.body)); + continue; + } + if (!listed.ok) { + rows.push(unread(target, where, listed.failed)); + continue; + } + // One candidate is named with its first difference as a LEAD, never as a verdict: it may be someone else's comment. + const lone = candidates.length === 1 ? judgeStoredBody({ op: target.op, sent: target.sent, stored: candidates[0].body }, judge) : null; + const lead = lone && lone.verdict === 'not-stored' ? `; the one comment created since then (${candidates[0].id}) first differs at byte ${lone.offset}${lone.replacements > 0 ? ` and carries ${lone.replacements} extra U+FFFD` : ''}` : ''; + rows.push({ action: target.action, op: target.op, where, verdict: 'unverified', cls: 'unfound', why: `${candidates.length} comment(s) created on #${target.issue} since the dispatch, none holding the bytes sent${lead}` }); + } + return { state: strokeReadBackState(rows), rows }; +} + /** * Send one packed payload and wait for its run. * - * Returns `{ state, ok, status, verdict, run, requestId, dispatchedAt, startMs, ceilingMs, detail }` where `state` is one of - * `success` · `failure` (the run completed otherwise) · `no-run` · `timeout` · `refused` (the dispatch itself). - * Never throws on an HTTP status. + * Returns `{ state, ok, status, verdict, run, requestId, dispatchedAt, startMs, ceilingMs, detail, readBack }` where `state` is one of + * `success` (the run succeeded and every body reads back as sent) · `failure` (the run completed otherwise, or — `notStored: + * true` — it succeeded and the read-back measured a body the board does not hold as sent) · `unverified` (it succeeded and a + * body could not be read back) · `no-run` · `timeout` · `refused` (the dispatch itself). `readBack` is `readBackStroke`'s + * answer, present once the run succeeded. Never throws on an HTTP status. */ export async function sendFleetWrite(payload, deps = {}) { const api = deps.api ?? DEFAULT_API; @@ -500,14 +784,28 @@ export async function sendFleetWrite(payload, deps = {}) { } const ok = run.conclusion === 'success'; log(`fleet-write: run ${run.id} completed — conclusion ${run.conclusion}${run.url ? ` ${run.url}` : ''}.`); - return { ...base, run, state: ok ? 'success' : 'failure', ok, status: 204, verdict: 'ok', dispatchedAt, detail: `conclusion ${run.conclusion}` }; + const done = { ...base, run, status: 204, verdict: 'ok', dispatchedAt }; + if (!ok) return { ...done, state: 'failure', ok: false, detail: `conclusion ${run.conclusion}` }; + + // ── read it back — the run's success is the executor's, not the write's ───── + const readBack = await readBackStroke(payload, { dispatchedAt }, { api, fetch: deps.fetch, token: deps.token, judge: deps.judge }); + for (const row of readBack.rows) log(readBackLine(row)); + if (readBack.state === 'not-stored') { + const first = readBack.rows.find((r) => r.verdict === 'not-stored'); + return { ...done, state: 'failure', ok: false, notStored: true, readBack, detail: `conclusion success, but NOT STORED — action ${first.action} (${first.op} ${first.where}) first differs from the bytes sent at byte ${first.offset}` }; + } + if (readBack.state === 'unverified') { + const first = readBack.rows.find((r) => r.verdict === 'unverified'); + return { ...done, state: 'unverified', ok: false, readBack, detail: `conclusion success, but UNVERIFIED — action ${first.action} (${first.op}): ${first.why ?? first.cls}` }; + } + return { ...done, state: 'success', ok: true, readBack, detail: `conclusion ${run.conclusion}` }; } /** The exit a tool takes from a result that is not `success`. */ export function exitForResult(result) { if (result.state === 'refused') return result.verdict === 'prerequisite' || result.verdict === 'ratelimit' ? EXIT_PREREQUISITE : EXIT_PLATFORM_REFUSAL; - if (result.state === 'failure') return EXIT_PLATFORM_REFUSAL; - if (result.state === 'no-run' || result.state === 'timeout') return EXIT_UNCONFIRMED; + if (result.state === 'failure') return result.notStored ? EXIT_NOT_STORED : EXIT_PLATFORM_REFUSAL; + if (result.state === 'no-run' || result.state === 'timeout' || result.state === 'unverified') return EXIT_UNCONFIRMED; return EXIT_OK; } @@ -568,10 +866,14 @@ const SELF_TEST_BATTERIES = Object.freeze({ 'the dispatch: one paced POST to the board repo with event_type and client_payload; 204 is acceptance, anything else a refusal': 7, 'the run-poller: the run named after the request id, its completion, its conclusion': 6, 'the ceilings: no run in the start window, no completion in the ceiling — UNCONFIRMED, never retried': 7, - 'the wiring: the POST is paced and on the roster, the reads are not, the token never reaches the log': 4, + 'the read-back targets: every op the table lets carry a body — derived, each with a locator — and no other op is read back': 6, + "the verdict: post-stamped's own classifier, imported — declared normalisations land, a split, lost or truncated byte is NOT STORED at its first differing byte, an unreadable body is unverified, the PR-create footer forgiven on pr_create alone": 13, + "the round trip: the card's 41,699 bytes with a multi-byte character across every 16 KiB boundary of every stream, byte for byte through pack, the wire, the runner's env text, the validator and the executor; a per-chunk decode is NOT STORED": 8, + 'the read-back end to end: after a success run each body at its locator — a corrupted read-back exits 4 through the CLI, an unreadable or unfound one 6, a body-less stroke reads nothing, never a retry': 17, + 'the wiring: the POST is paced and on the roster, the reads are not, the token never reaches the log': 5, 'the CLI: a dry run sends nothing, usage, the exit ladder, the session derived from the container, a route read behind a dead proxy refuses': 11, }); -const SELF_TEST_BATTERY_FLOOR = 11; +const SELF_TEST_BATTERY_FLOOR = 15; const UNATTRIBUTED_BATTERY = '(unattributed)'; const batteryCases = new Map(); @@ -596,6 +898,10 @@ export async function selfTest() { const TOKEN = 'ghs_FixtureTokenNotRealAtAll0000000000000'; const SESSION = 'session_01ABCDEFGHJKMNPQRSTVWXYZ'; const ACTIONS = [{ op: 'comment', issue: 19701, body: 'Hello from the relay.' }]; + // The read-back's judge, loaded the way the live path loads it — so this self-test, run from this file's own entry, + // is also the proof that the entry settles before the import back into it (the header's no-top-level-await note). + const JUDGE = await loadReadBackJudge(); + const { PLATFORM_COMMENT_FOOTER, classifyReadBack: POST_STAMPED_CLASSIFIER } = await import('../post-stamped.mjs'); // ── the selector ──────────────────────────────────────────────────────── battery('the selector: OS_FLEET_TRANSPORT wins; auto is dispatch only under all three conditions and direct otherwise'); @@ -811,7 +1117,16 @@ export async function selfTest() { * function of the elapsed ms since the dispatch, so appearance latency and * completion are modelled as time, not as call counts. */ - const platform = ({ dispatch = { status: 204 }, runs = () => [], one = () => null }, seen, clock) => async (url, init) => { + /** + * The BOARD the read-back reads: `store` maps a read call (`GET `) to + * `{ status, json }`, or to a function of the parsed query. The default + * holds the one comment the default stroke writes, created just after the + * dispatch and stored as sent — so a stroke that succeeds reads back clean. + */ + const T0 = Date.UTC(2026, 8, 22, 9, 4, 0); + const at = (ms) => new Date(T0 + ms).toISOString(); + const DEFAULT_STORE = { [`GET /repos/objectstack-ai/objectstack/issues/19701/comments`]: { status: 200, json: [{ id: 501, created_at: at(8_000), body: ACTIONS[0].body }] } }; + const platform = ({ dispatch = { status: 204 }, runs = () => [], one = () => null, store = DEFAULT_STORE }, seen, clock) => async (url, init) => { const u = new URL(url); const call = `${init?.method ?? 'GET'} ${u.pathname}`; seen.push({ call, query: u.search, body: init?.body ? JSON.parse(init.body) : null, auth: init?.headers?.authorization ?? '' }); @@ -830,17 +1145,24 @@ export async function selfTest() { const r = one(Number(m[1]), clock.elapsed()); return r ? { status: 200, headers, json: async () => r } : { status: 404, headers, json: async () => ({ message: 'Not Found' }) }; } + if (call in store) { + const s = typeof store[call] === 'function' ? store[call](u.searchParams) : store[call]; + // A list answers its rows on page 1 and nothing after, as the platform's pagination does. + const page = Number(u.searchParams.get('page') ?? '1'); + const json = Array.isArray(s.json) && page > 1 ? [] : s.json; + return { status: s.status, headers, json: async () => json }; + } return { status: 404, headers, json: async () => ({ message: 'Not Found' }) }; }; - const drive = async (script, { file, ceilings } = {}) => { + const drive = async (script, { file, ceilings, stroke = payload } = {}) => { const seen = []; const logs = []; - let nowMs = Date.UTC(2026, 8, 22, 9, 4, 0); + let nowMs = T0; const t0 = nowMs; const clock = { now: () => nowMs, elapsed: () => nowMs - t0 }; const pace = paceFor(file ?? join(dir, `pace-${paceCase++}.jsonl`)); try { - const r = await sendFleetWrite(payload, { + const r = await sendFleetWrite(stroke, { fetch: platform(script, seen, clock), token: TOKEN, pace, @@ -908,6 +1230,203 @@ export async function selfTest() { t('the ceilings read the environment with malformed overrides ignored and said', [ceilingsFrom({ OS_FLEET_DISPATCH_START_MS: '10', OS_FLEET_DISPATCH_TIMEOUT_MS: 'soon' }).startMs, ceilingsFrom({ OS_FLEET_DISPATCH_TIMEOUT_MS: 'soon' }).ceilingMs, ceilingsFrom({ OS_FLEET_DISPATCH_TIMEOUT_MS: 'soon' }).notes.length], [10, DEFAULT_CEILING_MS, 1]); } + // ── the read-back: which actions are read, and where ───────────────────── + battery('the read-back targets: every op the table lets carry a body — derived, each with a locator — and no other op is read back'); + { + t('BODY_OPS is derived from the op table: exactly the five ops that can carry a body', [...BODY_OPS].sort(), ['comment', 'comment_edit', 'issue_create', 'issue_patch', 'pr_create']); + t('⛔ READ_BACK_LOCATORS names exactly BODY_OPS — a body op without a locator reds here the day it lands in ops.mjs', Object.keys(READ_BACK_LOCATORS).sort(), [...BODY_OPS].sort()); + t('an addressed object and one found by a key are MEASURED; a new comment only its content can find', Object.fromEntries(Object.entries(READ_BACK_LOCATORS).map(([op, l]) => [op, l.found])), { issue_patch: 'address', comment_edit: 'address', issue_create: 'key', pr_create: 'key', comment: 'content' }); + const mixed = { repo: 'objectstack-ai/objectstack', actions: [{ op: 'labels_add', issue: 1, labels: ['a'] }, { op: 'issue_patch', issue: 2, state: 'closed' }, { op: 'issue_patch', issue: 3, body: 'B' }, { op: 'pr_create', title: 'T', head: 'h', base: 'main' }, { op: 'pr_create', title: 'T', head: 'h2', base: 'main', body: 'P' }] }; + t('only an action that SENT a body is a target, numbered as the stroke numbers it', readBackTargets(mixed).map((x) => [x.action, x.op]), [[3, 'issue_patch'], [5, 'pr_create']]); + t('…carrying the keys its locator needs', [readBackTargets(mixed)[0].issue, readBackTargets(mixed)[1].head, readBackTargets(mixed)[1].sent], [3, 'h2', 'P']); + let fetched = 0; + const none = await readBackStroke({ repo: 'objectstack-ai/objectstack', actions: [{ op: 'labels_add', issue: 1, labels: ['a'] }] }, { dispatchedAt: T0 }, { + judge: JUDGE, + fetch: async () => { + fetched += 1; + throw new Error('no read expected'); + }, + }); + t('a stroke with no body reads NOTHING back — state none, zero requests', [none.state, none.rows.length, fetched], ['none', 0, 0]); + } + + // ── the verdict ───────────────────────────────────────────────────────── + battery("the verdict: post-stamped's own classifier, imported — declared normalisations land, a split, lost or truncated byte is NOT STORED at its first differing byte, an unreadable body is unverified, the PR-create footer forgiven on pr_create alone"); + { + const J = (op, sent, stored) => judgeStoredBody({ op, sent, stored }, JUDGE); + const SENT = 'A census table:\n| card | title |\n|---|---|\n| #1 | no… 全部 |\n'; + const TRIMMED = SENT.replace(/\n+$/u, ''); + t('identical bytes land', [J('issue_patch', SENT, SENT).verdict, J('issue_patch', SENT, SENT).cls], ['landed', 'identical']); + t("the platform's footer block appended lands (post-stamped's footer-appended)", [J('issue_patch', SENT, `${TRIMMED}${PLATFORM_COMMENT_FOOTER}`).verdict, J('issue_patch', SENT, `${TRIMMED}${PLATFORM_COMMENT_FOOTER}`).cls], ['landed', 'footer-appended']); + t('a stripped trailing newline lands', J('comment', SENT, TRIMMED).verdict, 'landed'); + const ell = Buffer.byteLength(SENT.slice(0, SENT.indexOf('…')), 'utf8'); + const SPLIT_ELL = SENT.replace('…', '\uFFFD\uFFFD\uFFFD'); + const e1 = J('issue_patch', SENT, SPLIT_ELL); + t("⭐ the card's first shape — '…' stored as three U+FFFD — is NOT STORED at the byte where '…' began, three extra replacement characters", [e1.verdict, e1.offset, e1.replacements], ['not-stored', ell, 3]); + const quan = Buffer.byteLength(SENT.slice(0, SENT.indexOf('全')), 'utf8'); + const e2 = J('issue_patch', SENT, SENT.replace('全', '\uFFFD\uFFFD')); + t("⭐ …and its second — '全' stored as two U+FFFD — at the byte where '全' began", [e2.verdict, e2.offset, e2.replacements], ['not-stored', quan, 2]); + t('a truncation is NOT STORED at the first byte the stored body lacks', [J('comment_edit', SENT, SENT.slice(0, 10)).verdict, J('comment_edit', SENT, SENT.slice(0, 10)).offset], ['not-stored', 10]); + t('a stored body that is not a string is unverified — never landed, never not-stored', [J('issue_patch', SENT, null).verdict, J('issue_patch', SENT, undefined).verdict], ['unverified', 'unverified']); + const PR_TAIL = '\n\n---\n_Generated by [Claude Code](https://claude.ai/code/session_01KTZmMfzVzjNvyaLyQ8mHvg)_'; + t("pr_create: the platform's rule + ONE session-URL footer after every byte sent lands as pr-create-footer-appended", [J('pr_create', SENT, `${TRIMMED}${PR_TAIL}`).verdict, J('pr_create', SENT, `${TRIMMED}${PR_TAIL}`).cls], ['landed', 'pr-create-footer-appended']); + t("⛔ …that cell is pr_create's alone: the same tail on an issue body is NOT STORED", J('issue_patch', SENT, `${TRIMMED}${PR_TAIL}`).verdict, 'not-stored'); + t('⛔ …it forgives nothing INSIDE the bytes sent: a split character before that footer is still NOT STORED', J('pr_create', SENT, `${SPLIT_ELL.replace(/\n+$/u, '')}${PR_TAIL}`).verdict, 'not-stored'); + t('⛔ …nor an append that is not exactly that block', J('pr_create', SENT, `${SENT}\nextra`).verdict, 'not-stored'); + const line = readBackLine({ action: 1, op: 'issue_patch', where: 'objectstack-ai/objectui#11041', ...e1 }); + t('the NOT STORED line names the op, the object, the first differing byte and the replacement characters', line.includes('NOT STORED') && line.includes('issue_patch') && line.includes('#11041') && line.includes(`first difference at byte ${ell}`) && line.includes('3 U+FFFD')); + const own = readReal(SELF_PATH, 'utf8'); + t("⛔ the judge is post-stamped's, imported — this file declares no classifier of its own", [JUDGE.classifyReadBack === POST_STAMPED_CLASSIFIER, /import\('\.\.\/post-stamped\.mjs'\)/u.test(own), new RegExp('function\\s+(classifyReadBack|footerReAnchoring|firstDifferingByte|sentBodyLanded)\\b').test(own)], [true, true, false]); + } + + // ── the round trip ────────────────────────────────────────────────────── + // The platform legs are MODELLED here, not measured: GitHub parses the wire, and the runner's `toJSON` is taken as + // the 2-space JSON the workflow's expression engine renders. Everything this repository executes is the real code. + battery("the round trip: the card's 41,699 bytes with a multi-byte character across every 16 KiB boundary of every stream, byte for byte through pack, the wire, the runner's env text, the validator and the executor; a per-chunk decode is NOT STORED"); + { + const { executeFleetWrite } = await import('./execute.mjs'); + const KIB16 = 16 * 1024; + const CARD_BYTES = 41_699; + const UNIT = '𠮷😀全😀…😀é\n"'; // 4-, 3- and 2-byte characters, and the newline and quote the JSON streams escape + const UNIT_BYTES = Buffer.byteLength(UNIT, 'utf8'); + const build = (pad) => { + let s = 'x'.repeat(pad); + while (Buffer.byteLength(s, 'utf8') + UNIT_BYTES <= CARD_BYTES) s += UNIT; + return s + 'x'.repeat(CARD_BYTES - Buffer.byteLength(s, 'utf8')); + }; + const streamsOf = (body) => { + const p = packRequest({ repo: 'objectstack-ai/objectui', session: SESSION, actions: [{ op: 'issue_patch', issue: 11041, body }], requestId: 'fw-20260929T061500Z-abc123' }).payload; + const wire = JSON.stringify({ event_type: RELAY_EVENT_TYPE, client_payload: p }); + return { payload: p, wire, envText: JSON.stringify(JSON.parse(wire).client_payload, null, 2), exec: JSON.stringify({ body }) }; + }; + /** Every 16 KiB boundary of the stream falls on a UTF-8 continuation byte — inside a character, never between two. */ + const straddles = (text) => { + const buf = Buffer.from(text, 'utf8'); + for (let k = KIB16; k < buf.length; k += KIB16) if ((buf[k] & 0xc0) !== 0x80) return false; + return true; + }; + let pad = -1; + for (let p = 0; p < 2000 && pad < 0; p++) { + const body = build(p); + const s = streamsOf(body); + if ([body, s.wire, s.envText, s.exec].every(straddles)) pad = p; + } + const BODY = build(Math.max(pad, 0)); + const SENT_BYTES = Buffer.from(BODY, 'utf8'); + const S = streamsOf(BODY); + t("a body of exactly the card's 41,699 bytes exists whose EVERY 16 KiB boundary — in the body, the dispatch wire, the runner's env text and the executor's request — falls inside a multi-byte character", [pad >= 0, SENT_BYTES.length, [BODY, S.wire, S.envText, S.exec].every(straddles)], [true, CARD_BYTES, true]); + t('…and every one of those streams crosses at least two such boundaries', [BODY, S.wire, S.envText, S.exec].map((x) => Buffer.byteLength(x, 'utf8') > 2 * KIB16), [true, true, true, true]); + t('① pack → the dispatch wire → the platform\'s parse: byte for byte', Buffer.compare(Buffer.from(JSON.parse(S.wire).client_payload.actions[0].body, 'utf8'), SENT_BYTES), 0); + const child = spawnSync(process.execPath, [fileURLToPath(new URL('./validate.mjs', import.meta.url)), '--from-env', '--json'], { encoding: 'utf8', env: { ...process.env, [PAYLOAD_ENV]: S.envText, GITHUB_STEP_SUMMARY: '' }, maxBuffer: 1 << 24 }); + const judged = child.status === 0 ? JSON.parse(child.stdout.trim().split('\n')[1]) : null; + t("② the runner's env text through FLEET_WRITE_PAYLOAD into validate.mjs --from-env, in its own process: byte for byte", [child.status, judged ? Buffer.compare(Buffer.from(judged.actions[0].body, 'utf8'), SENT_BYTES) : null], [0, 0], child.stderr.slice(-300)); + let captured = null; + const execRun = await executeFleetWrite({ payload: JSON.parse(S.envText), sender: 'os-support-ai', token: TOKEN }, { + fetch: async (url, init) => { + if (init?.method === 'PATCH') captured = init.body; + const json = new URL(url).pathname.endsWith('/permission') ? { permission: 'write', role_name: 'write' } : { number: 11041 }; + return { status: 200, headers: new Headers({ 'x-ratelimit-remaining': '4999' }), json: async () => json }; + }, + pace: paceFor(join(dir, 'pace-roundtrip.jsonl')), + log: () => {}, + }); + t("③ the executor's own PATCH is exactly the request the stream check straddled, carrying the bytes sent", [execRun.exit, captured === S.exec, captured ? Buffer.compare(Buffer.from(JSON.parse(captured).body, 'utf8'), SENT_BYTES) : null], [0, true, 0]); + const answer = (body) => async () => ({ status: 200, headers: new Headers(), json: async () => ({ number: 11041, body }) }); + const back = await readBackStroke(S.payload, { dispatchedAt: T0 }, { judge: JUDGE, token: TOKEN, fetch: answer(`${BODY}${PLATFORM_COMMENT_FOOTER}`) }); + t("④ what that request stored, read back with the platform's footer, lands", [back.state, back.rows[0]?.cls], ['landed', 'footer-appended']); + // The defect the card suspects: a reader that decodes each 16 KiB chunk of the wire on its own. + const wireBuf = Buffer.from(S.wire, 'utf8'); + const chunked = Array.from({ length: Math.ceil(wireBuf.length / KIB16) }, (_, i) => wireBuf.subarray(i * KIB16, (i + 1) * KIB16).toString('utf8')).join(''); + const damaged = JSON.parse(chunked).client_payload.actions[0].body; + // Where the first wire boundary falls in BODY coordinates: walk the body a character at a time, each one as wide + // on the wire as its JSON escape, until the character that spans wire byte 16384. + let predicted = null; + { + let wirePos = Buffer.byteLength(S.wire.slice(0, S.wire.indexOf('"body":"') + '"body":"'.length), 'utf8'); + let bodyPos = 0; + for (const ch of BODY) { + const w = Buffer.byteLength(JSON.stringify(ch).slice(1, -1), 'utf8'); + if (wirePos < KIB16 && KIB16 < wirePos + w) { + predicted = bodyPos; + break; + } + wirePos += w; + bodyPos += Buffer.byteLength(ch, 'utf8'); + } + } + const broken = await readBackStroke(S.payload, { dispatchedAt: T0 }, { judge: JUDGE, token: TOKEN, fetch: answer(damaged) }); + t('⑤ that body decoded chunk by chunk — every boundary splitting a character — is NOT STORED, with the replacement characters counted', [broken.state, damaged === BODY, broken.rows[0]?.replacements > 0], ['not-stored', false, true]); + t('…at exactly the body byte where the first 16 KiB wire boundary split a character', [predicted !== null, broken.rows[0]?.offset], [true, predicted]); + } + + // ── the read-back, end to end ─────────────────────────────────────────── + battery('the read-back end to end: after a success run each body at its locator — a corrupted read-back exits 4 through the CLI, an unreadable or unfound one 6, a body-less stroke reads nothing, never a retry'); + { + const UI = 'objectstack-ai/objectui'; + const strokeOf = (actions, repo = 'objectstack-ai/objectstack') => packRequest({ repo, session: SESSION, actions, requestId: 'fw-test-1' }).payload; + const OK_RUN = { runs: () => [RUN('completed', 'success')] }; + const B = 'Body with a multi-byte tail: no… 全部\n'; + const SPLIT = B.replace('全', '\uFFFD\uFFFD'); + const AT = Buffer.byteLength(B.slice(0, B.indexOf('全')), 'utf8'); + const patch = strokeOf([{ op: 'issue_patch', issue: 11041, body: B }], UI); + const ISSUE = `GET /repos/${UI}/issues/11041`; + const good = await drive({ ...OK_RUN, store: { [ISSUE]: { status: 200, json: { number: 11041, body: `${B.replace(/\n+$/u, '')}${PLATFORM_COMMENT_FOOTER}` } } } }, { stroke: patch }); + t('a success run whose body reads back as sent (plus the footer) is success, its row landed, exit 0', [good.state, good.ok, good.readBack?.state, exitForResult(good)], ['success', true, 'landed', EXIT_OK]); + t('…read at the address the action named, once, AFTER the run completed', [good.seen.filter((s) => s.call === ISSUE).length, good.seen.findIndex((s) => s.call === ISSUE) > good.seen.findIndex((s) => s.call === RUNS)], [1, true]); + const bad = await drive({ ...OK_RUN, store: { [ISSUE]: { status: 200, json: { number: 11041, body: SPLIT } } } }, { stroke: patch }); + t('⭐ a deliberately corrupted read-back: failure + notStored, exit 4, naming the op, the object and the first differing byte', [bad.state, bad.ok, bad.notStored, exitForResult(bad), bad.detail.includes('issue_patch') && bad.detail.includes('#11041') && bad.detail.includes(`byte ${AT}`)], ['failure', false, true, EXIT_NOT_STORED, true], bad.detail); + t('…the transcript says NOT STORED and counts the U+FFFD, and ⛔ there was exactly ONE dispatch — never a retry', [bad.logs.some((l) => l.includes('NOT STORED') && l.includes('U+FFFD')), bad.seen.filter((s) => s.call === DISPATCH).length], [true, 1]); + t('…and the CLI text says the write HAPPENED, go READ, never re-run blind, exit 4', notStoredText(bad).includes('HAPPENED') && notStoredText(bad).includes('Do not re-run') && notStoredText(bad).includes(`Exit ${EXIT_NOT_STORED}`)); + const unread = await drive({ ...OK_RUN, store: { [ISSUE]: { status: 503, json: { message: 'down' } } } }, { stroke: patch }); + t('an object the read-back cannot read is unverified, exit 6 — ⛔ never success, ⛔ never not-stored — and says which call failed', [unread.state, unread.ok, unread.notStored ?? false, exitForResult(unread), unverifiedText(unread).includes('HTTP 503')], ['unverified', false, false, EXIT_UNCONFIRMED, true]); + const EDIT = 'GET /repos/objectstack-ai/objectstack/issues/comments/77'; + const edit = await drive({ ...OK_RUN, store: { [EDIT]: { status: 200, json: { id: 77, body: SPLIT } } } }, { stroke: strokeOf([{ op: 'comment_edit', comment_id: 77, body: B }]) }); + t('comment_edit: read at its address; a split character is NOT STORED, exit 4', [edit.state, edit.notStored, exitForResult(edit)], ['failure', true, EXIT_NOT_STORED]); + const CMTS = 'GET /repos/objectstack-ai/objectstack/issues/19701/comments'; + const two = await drive({ ...OK_RUN, store: { [CMTS]: { status: 200, json: [{ id: 601, created_at: at(5_000), body: B }, { id: 602, created_at: at(6_000), body: B }] } } }, { stroke: strokeOf([{ op: 'comment', issue: 19701, body: B }, { op: 'comment', issue: 19701, body: B }]) }); + t('comment: found by its content among those created since the dispatch; two identical comments match two DIFFERENT stored ones', [two.state, two.readBack?.rows?.map((r) => r.where)], ['success', ['objectstack-ai/objectstack#19701 comment 602', 'objectstack-ai/objectstack#19701 comment 601']]); + const lost = await drive({ ...OK_RUN, store: { [CMTS]: { status: 200, json: [{ id: 603, created_at: at(5_000), body: SPLIT }] } } }, { stroke: strokeOf([{ op: 'comment', issue: 19701, body: B }]) }); + t("⛔ a comment only its content can find is UNVERIFIED when nothing matches — exit 6, never not-stored — the one candidate's first difference named as a LEAD", [lost.state, exitForResult(lost), lost.readBack?.rows?.[0]?.why?.includes(`first differs at byte ${AT}`) ?? false, lost.readBack?.rows?.[0]?.why?.includes('U+FFFD') ?? false], ['unverified', EXIT_UNCONFIRMED, true, true]); + const stale = await drive({ ...OK_RUN, store: { [CMTS]: { status: 200, json: [{ id: 604, created_at: at(-120_000), body: B }] } } }, { stroke: strokeOf([{ op: 'comment', issue: 19701, body: B }]) }); + t('…and a comment created BEFORE the dispatch (outside the slack) is never taken for this write', [stale.state, exitForResult(stale)], ['unverified', EXIT_UNCONFIRMED]); + const ISSUES = 'GET /repos/objectstack-ai/objectstack/issues'; + const created = await drive({ ...OK_RUN, store: { [ISSUES]: { status: 200, json: [{ id: 9, number: 20999, title: 'Other', created_at: at(4_000), body: B }, { id: 8, number: 20998, title: 'Card', created_at: at(3_000), body: SPLIT }] } } }, { stroke: strokeOf([{ op: 'issue_create', title: 'Card', body: B }]) }); + t('issue_create: the new issue found by the title sent and its body judged — a split character is NOT STORED on THAT issue, exit 4', [created.state, created.notStored, created.readBack?.rows?.[0]?.where, exitForResult(created)], ['failure', true, 'objectstack-ai/objectstack#20998', EXIT_NOT_STORED]); + const PULLS = 'GET /repos/objectstack-ai/objectstack/pulls'; + const PR_TAIL = '\n\n---\n_Generated by [Claude Code](https://claude.ai/code/session_01KTZmMfzVzjNvyaLyQ8mHvg)_'; + const prStroke = strokeOf([{ op: 'pr_create', title: 'T', head: 'claude/issue-1-x', base: 'main', body: B }]); + const pr = await drive({ ...OK_RUN, store: { [PULLS]: (q) => ({ status: 200, json: q.get('head') === 'objectstack-ai:claude/issue-1-x' ? [{ id: 31, number: 20600, created_at: at(4_000), body: `${B.replace(/\n+$/u, '')}${PR_TAIL}` }] : [] }) } }, { stroke: prStroke }); + t("pr_create: found by owner:head; the platform's session-URL footer after every byte sent is success", [pr.state, pr.readBack?.rows?.[0]?.cls, pr.readBack?.rows?.[0]?.where], ['success', 'pr-create-footer-appended', 'objectstack-ai/objectstack#20600']); + const noPr = await drive({ ...OK_RUN, store: { [PULLS]: { status: 200, json: [] } } }, { stroke: prStroke }); + t('…and a pull the head does not find is unverified, exit 6', [noPr.state, exitForResult(noPr)], ['unverified', EXIT_UNCONFIRMED]); + const labels = await drive(OK_RUN, { stroke: strokeOf([{ op: 'labels_add', issue: 1, labels: ['a'] }]) }); + t('a stroke carrying no body reads nothing back: success, read-back none, no request beyond the dispatch and the run', [labels.state, labels.readBack?.state, labels.seen.filter((s) => !s.call.includes('/actions/runs') && s.call !== DISPATCH).length], ['success', 'none', 0]); + t("the callers' one vocabulary: failure + notStored → 4, failure → 5, unverified → 6", [exitForResult({ state: 'failure', notStored: true }), exitForResult({ state: 'failure' }), exitForResult({ state: 'unverified' })], [EXIT_NOT_STORED, EXIT_PLATFORM_REFUSAL, EXIT_UNCONFIRMED]); + // The CLI, end to end: the exit code a seat actually reads, from the actions file the card's seat sent. + writeFileSync(join(dir, 'patch.json'), JSON.stringify([{ op: 'issue_patch', issue: 11041, body: B }]), 'utf8'); + const cli = async (stored) => { + let nowMs = T0; + const clock = { now: () => nowMs, elapsed: () => nowMs - T0 }; + const out = []; + const [log, err] = [console.log, console.error]; + console.log = (l) => out.push(String(l)); + console.error = () => {}; + try { + const code = await main(['--repo', UI, '--actions-file', join(dir, 'patch.json'), '--request-id', 'fw-test-1', '--json'], { + env: { GITHUB_TOKEN: TOKEN, [SESSION_ENV]: SESSION }, + send: { fetch: platform({ ...OK_RUN, store: { [ISSUE]: { status: 200, json: { number: 11041, body: stored } } } }, [], clock), pace: paceFor(join(dir, `pace-cli-${paceCase++}.jsonl`)), now: clock.now, sleep: async (ms) => { nowMs += ms; }, log: () => {}, ceilings: { startMs: 90_000, ceilingMs: 300_000, pollMs: 5_000, notes: [] } }, + }); + return { code, json: out.length ? JSON.parse(out[out.length - 1]) : null }; + } finally { + [console.log, console.error] = [log, err]; + } + }; + const cliBad = await cli(SPLIT); + t("⭐ the CLI: the card's stroke with a deliberately corrupted read-back EXITS 4, and --json says not_stored with the first differing byte", [cliBad.code, cliBad.json?.not_stored, cliBad.json?.read_back?.[0]?.first_difference_byte], [EXIT_NOT_STORED, true, AT]); + const cliGood = await cli(`${B.replace(/\n+$/u, '')}${PLATFORM_COMMENT_FOOTER}`); + t('…and the same stroke read back intact exits 0, its row landed', [cliGood.code, cliGood.json?.read_back?.[0]?.verdict], [EXIT_OK, 'landed']); + } + // ── the wiring ────────────────────────────────────────────────────────── battery('the wiring: the POST is paced and on the roster, the reads are not, the token never reaches the log'); { @@ -916,8 +1435,10 @@ export async function selfTest() { t('write-pace lists this file among the wired write tools', WIRED_WRITE_TOOLS.includes('scripts/pm/fleet-write/dispatch.mjs')); t('this file calls both halves, guarded by the write-verb predicate', own.includes('paceWrite(') && own.includes('noteResponse(') && own.includes('isWriteMethod(')); const records = readReal(paceFile, 'utf8'); - t('the dispatches above were paced; the run-list and run reads were not', [records.includes('fleet-write dispatch POST'), records.includes('fleet-write dispatch GET')], [true, false]); + t('the dispatches above were paced; the run-list, run and read-back reads were not', [records.includes('fleet-write dispatch POST'), records.includes('fleet-write dispatch GET')], [true, false]); t('⛔ and the token never reached the throttle\'s log', records.includes(TOKEN), false); + const entry = own.slice(own.lastIndexOf('if (isEntrypoint(import.meta.url))')); + t('⛔ the entry settles before main runs — no top-level await on main — so the read-back\'s import of post-stamped (which imports this file) cannot wait on it', [/\bawait\s+main\(/u.test(entry), /main\(process\.argv\.slice\(2\)\)\.then\(/u.test(entry)], [false, true]); } // ── the CLI ───────────────────────────────────────────────────────────── @@ -985,7 +1506,8 @@ export async function selfTest() { } console.log( `✓ fleet-write/dispatch self-test: ${cases.length} cases pass across ${declared.length} batteries — the transport selector that never guesses, ` + - 'the session on the envelope, one paced dispatch per stroke, a run found by its request id and waited to its conclusion, both ceilings answered UNCONFIRMED and never retried.', + 'the session on the envelope, one paced dispatch per stroke, a run found by its request id and waited to its conclusion, both ceilings answered UNCONFIRMED and never retried, ' + + 'and every body read back after a success run — the card\'s 41,699 bytes byte-exact across every 16 KiB boundary, a corrupted read-back NOT STORED (exit 4), an unfound one UNCONFIRMED (6).', ); selfTestReachedVerdict = true; return 0; @@ -1040,8 +1562,9 @@ const USAGE = [ ` ${CONTAINER_SESSION_ENV} (cse_ → session_); ${SESSION_ENV} overrides the container (a local checkout, a test).`, ` ⛔ Never as a prefix on the command line: the seats' allow rules are literal command prefixes. The dispatch always`, ` goes to ${RELAY_REPO}; --repo names the target.`, - ` Exits: ${EXIT_OK} run succeeded · ${EXIT_USAGE} usage / payload refused · ${EXIT_PREREQUISITE} prerequisite · ${EXIT_PLATFORM_REFUSAL} dispatch refused or run failed ·`, - ` ${EXIT_UNCONFIRMED} UNCONFIRMED (no run, or no completion within the ceiling) · ${EXIT_WRITE_PACE_REFUSED} throttle refused`, + ` Exits: ${EXIT_OK} run succeeded and every body reads back as sent · ${EXIT_USAGE} usage / payload refused · ${EXIT_PREREQUISITE} prerequisite ·`, + ` ${EXIT_NOT_STORED} NOT STORED (the run succeeded, a body reads back otherwise — go READ, never retry) · ${EXIT_PLATFORM_REFUSAL} dispatch refused or run failed ·`, + ` ${EXIT_UNCONFIRMED} UNCONFIRMED (no run, no completion within the ceiling, or a body not read back) · ${EXIT_WRITE_PACE_REFUSED} throttle refused`, ].join('\n'); function rearmThroughProxy(args) { @@ -1067,7 +1590,14 @@ function rearmThroughProxy(args) { return null; } -export async function main(argv) { +/** + * The CLI. `deps` is the self-test's door and nothing else's: `env` stands in + * for `process.env` and `send` is handed to `sendFleetWrite` (a fake platform, + * a throttle file, a clock), so a corrupted read-back can be driven end to end + * through the exit code a seat actually reads. + */ +export async function main(argv, deps = {}) { + const env = deps.env ?? process.env; const opts = parseArgs(argv); if (opts.help) { console.log(USAGE); @@ -1094,10 +1624,10 @@ export async function main(argv) { console.error(`fleet-write/dispatch: --actions-file ${opts.actionsFile} cannot be read as JSON (${e?.message ?? 'unreadable'}). Nothing was sent.`); return EXIT_USAGE; } - const session = opts.session ?? sessionFrom(process.env); + const session = opts.session ?? sessionFrom(env); if (!session || !SESSION_SHAPE.test(session)) { console.error( - `fleet-write/dispatch: PREREQUISITE NOT MET — no session id: ${sessionSource(process.env).reason}. Pass --session session_…, or set ${SESSION_ENV} (a cloud seat container supplies it as ${CONTAINER_SESSION_ENV}=cse_; the explicit variable overrides it). The envelope carries the dispatching seat's identity; it is never invented. Nothing was sent.`, + `fleet-write/dispatch: PREREQUISITE NOT MET — no session id: ${sessionSource(env).reason}. Pass --session session_…, or set ${SESSION_ENV} (a cloud seat container supplies it as ${CONTAINER_SESSION_ENV}=cse_; the explicit variable overrides it). The envelope carries the dispatching seat's identity; it is never invented. Nothing was sent.`, ); return EXIT_PREREQUISITE; } @@ -1112,27 +1642,49 @@ export async function main(argv) { console.log(JSON.stringify({ event_type: RELAY_EVENT_TYPE, client_payload: packed.payload }, null, 2)); return EXIT_OK; } - const token = process.env.GITHUB_TOKEN ?? process.env.GH_TOKEN ?? ''; + const token = env.GITHUB_TOKEN ?? env.GH_TOKEN ?? ''; if (!token) { console.error('fleet-write/dispatch: PREREQUISITE NOT MET — no GITHUB_TOKEN / GH_TOKEN in the environment. Nothing was sent.'); return EXIT_PREREQUISITE; } - const rearmed = rearmThroughProxy(argv); - if (rearmed !== null) return rearmed; - const result = await sendFleetWrite(packed.payload, { token }); - if (!result.ok && (result.state === 'no-run' || result.state === 'timeout')) console.error(unconfirmedText(result, 'scripts/pm/fleet-write/dispatch.mjs')); - if (opts.json) console.log(JSON.stringify({ request_id: packed.payload.request_id, state: result.state, ok: result.ok, run: result.run, dispatched_at: new Date(result.dispatchedAt).toISOString() })); + if (!deps.send) { + const rearmed = rearmThroughProxy(argv); + if (rearmed !== null) return rearmed; + } + const tool = 'scripts/pm/fleet-write/dispatch.mjs'; + const result = await sendFleetWrite(packed.payload, { token, ...(deps.send ?? {}) }); + if (!result.ok && (result.state === 'no-run' || result.state === 'timeout')) console.error(unconfirmedText(result, tool)); + if (result.notStored) console.error(notStoredText(result, tool)); + if (result.state === 'unverified') console.error(unverifiedText(result, tool)); + if (opts.json) { + console.log( + JSON.stringify({ + request_id: packed.payload.request_id, + state: result.state, + ok: result.ok, + not_stored: result.notStored === true, + run: result.run, + dispatched_at: new Date(result.dispatchedAt).toISOString(), + read_back: (result.readBack?.rows ?? []).map((r) => ({ action: r.action, op: r.op, where: r.where, verdict: r.verdict, class: r.cls, first_difference_byte: r.offset ?? null, sent_bytes: r.sentBytes ?? null, stored_bytes: r.storedBytes ?? null })), + }), + ); + } return result.ok ? EXIT_OK : exitForResult(result); } +// ⛔ No top-level await here. The read-back imports post-stamped, which imports THIS file: were this module's evaluation +// still pending on `await main(…)`, that import would wait on it while it waits on the import — a deadlock node reports +// as exit 13, "unsettled top-level await" (measured, the first time the read-back ran from this entry). Settling the +// evaluation first and running `main` after it is what makes the dynamic import in `loadReadBackJudge` safe. if (isEntrypoint(import.meta.url)) { - const code = await main(process.argv.slice(2)); - if (process.argv.includes('--self-test') && !selfTestReachedVerdict) { - console.error( - '\n✗ fleet-write/dispatch self-test: selfTest() returned without reaching its verdict, so no success line was\n' + - 'printed. Exiting 0 here would report a self-test that never finished as one that passed.\n', - ); - process.exit(1); - } - process.exit(code); + main(process.argv.slice(2)).then((code) => { + if (process.argv.includes('--self-test') && !selfTestReachedVerdict) { + console.error( + '\n✗ fleet-write/dispatch self-test: selfTest() returned without reaching its verdict, so no success line was\n' + + 'printed. Exiting 0 here would report a self-test that never finished as one that passed.\n', + ); + process.exit(1); + } + process.exit(code); + }); }