From 6220133fb67bc192082e9f4cb7125188ae4c7218 Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 15 Sep 2026 01:29:49 +0200 Subject: [PATCH 1/3] fix(ci): transfer proof brokers through bounded run storage Session-Id: 01a09c40-ce3b-7f11-a7df-b6b7ccab6fd9 --- .github/workflows/relayflow-pr-proof.yml | 5 +- .github/workflows/test.yml | 3 + docs/pr-proof-broker-transfer.md | 43 +++ scripts/pr-proof/broker-transfer.mjs | 301 +++++++++++++++++++++ scripts/pr-proof/broker-transfer.test.mjs | 309 ++++++++++++++++++++++ scripts/pr-proof/run-arm.mjs | 12 +- scripts/pr-proof/run-cloud.mjs | 31 ++- 7 files changed, 695 insertions(+), 9 deletions(-) create mode 100644 docs/pr-proof-broker-transfer.md create mode 100644 scripts/pr-proof/broker-transfer.mjs create mode 100644 scripts/pr-proof/broker-transfer.test.mjs diff --git a/.github/workflows/relayflow-pr-proof.yml b/.github/workflows/relayflow-pr-proof.yml index 6dbf27530f..85c49bb191 100644 --- a/.github/workflows/relayflow-pr-proof.yml +++ b/.github/workflows/relayflow-pr-proof.yml @@ -112,9 +112,8 @@ jobs: if: steps.prepare.outputs.required == 'true' run: | git add -f -- .relayflow/pr-proof-input.json - if [[ "${{ steps.prepare.outputs.broker_artifacts_required }}" == "true" ]]; then - git add -f -- .relayflow/pr-proof-binaries - fi + # Exact binaries go through authenticated run storage, never Relayfile seed. + test -z "$(git ls-files -- .relayflow/pr-proof-binaries)" - name: Install the released Agent Relay CLI if: steps.prepare.outputs.required == 'true' diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index 6232863b5b..d2e7dd9806 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -49,6 +49,9 @@ jobs: - name: Ensure rollup optional dependencies are installed run: npm install --no-save rollup || true + - name: Check PR proof transfer guards + run: node --test scripts/pr-proof/broker-transfer.test.mjs scripts/pr-proof/run-cloud.test.mjs + - name: Check subscription proof guards run: npm run test:subscriptions:proof diff --git a/docs/pr-proof-broker-transfer.md b/docs/pr-proof-broker-transfer.md new file mode 100644 index 0000000000..c07ef7d46a --- /dev/null +++ b/docs/pr-proof-broker-transfer.md @@ -0,0 +1,43 @@ +# Broker artifacts in Cloud PR proof + +The trusted dispatcher verifies exact base/head build manifests, then stages +only the validated proof input for code sync. Raw binaries stay on the GitHub +runner and enter authenticated run storage after the prepared-run ID arrives. + +Each artifact is limited to 100 MiB, the pair to 200 MiB, and each raw part to +2 MiB (at most 3 MiB as JSON). Each part binds run ID, nonce, arm, source SHA, +final build SHA-256, total byte length, index/count, and part checksum. +Exclusive `If-None-Match: *` writes reject collisions. The readiness manifest +is published only after all parts for that arm are acknowledged. Lost upload +acknowledgments fail the proof without replaying writes or creating a new run. + +The isolated arm waits at most 60 seconds for the manifest, reads each part +once with bounded response sizes, and verifies the final raw build hash. +Requests time out after 15 seconds with a 120-second transfer lifetime. +Missing, duplicated, misbound, oversized, or corrupt parts fail closed. Only +verified bytes enter the existing private mode-0500 executable, isolation, +and post-execution attestation checks. Broker requirements and red/green +outcomes remain unchanged. + +The existing Cloud storage contract was checked at revision `2ca32360`: CI +prepared-run-grant PUT, sandbox-scoped GET, and exclusive creation are present. +No new endpoint is required. CI credentials stay in the GitHub environment; +case processes receive neither credentials nor transfer authorization headers. + +Each arm consumes and tombstones its manifest/parts immediately after bounded +assembly while its sandbox token is live. Only then may the verified bytes +reach the private executable. Private executable cleanup remains in the arm's +`finally` block. On submission failure/interruption, CI tombstones its nonce's +objects before cancellation revokes the prepared-run write grant. Successful +proof releases only local buffers because both arms already consumed storage. +Cleanup failures fail the dispatcher. Cloud has no general object-delete or +prepared-grant-revocation route; this change does not invent one or revoke +shared CI credentials. Existing cancellation revokes run API sessions. +Forced runner termination or independent remote token revocation can prevent +cleanup; remaining objects contain public binaries, remain run-scoped, and +follow platform storage lifecycle and token expiry. + +This prerequisite must be admitted to the trusted base before rerunning another +PR's `pull_request_target` proof. A change on that PR's head cannot repair the +base dispatcher. Resolve and attest both exact binaries again after admission. +Local mocked transfer checks do not establish a successful Cloud proof. diff --git a/scripts/pr-proof/broker-transfer.mjs b/scripts/pr-proof/broker-transfer.mjs new file mode 100644 index 0000000000..4f7e835d7d --- /dev/null +++ b/scripts/pr-proof/broker-transfer.mjs @@ -0,0 +1,301 @@ +/** Run-scoped broker transfer. No artifact bytes or credentials enter code sync. */ +import { createHash } from 'node:crypto'; +import { readFile } from 'node:fs/promises'; +import { setTimeout as delay } from 'node:timers/promises'; +import { validateProofInput } from './contract.mjs'; +import { loadBrokerArtifact } from './stage-broker-artifacts.mjs'; + +export const PART_BYTES = 2 * 1024 * 1024; +export const MAX_ARTIFACT_BYTES = 100 * 1024 * 1024; +export const MAX_PAIR_BYTES = 2 * MAX_ARTIFACT_BYTES; +const MAX_PART_JSON = 3 * 1024 * 1024; +const sha256 = (bytes) => createHash('sha256').update(bytes).digest('hex'); +const equal = (a, b) => JSON.stringify(a) === JSON.stringify(b); + +function binding(input, arm, runId, totalBytes) { + const artifact = input.runtimeArtifacts?.broker?.[arm]; + if ( + !/^[A-Za-z0-9][A-Za-z0-9_-]{0,127}$/.test(runId ?? '') || + !/^[0-9a-f]{32}$/.test(input.handoffNonce ?? '') || + !['base', 'head'].includes(arm) || + artifact?.sourceSha !== (arm === 'base' ? input.baseSha : input.headSha) || + !/^[0-9a-f]{40}$/.test(artifact?.sourceSha ?? '') || + !/^[0-9a-f]{64}$/.test(artifact?.sha256 ?? '') || + !Number.isSafeInteger(totalBytes) || + totalBytes < 1 || + totalBytes > MAX_ARTIFACT_BYTES + ) { + throw new Error('Invalid broker transfer binding or size'); + } + return { + version: 1, + runId, + nonce: input.handoffNonce, + arm, + sourceSha: artifact.sourceSha, + sha256: artifact.sha256, + totalBytes, + count: Math.ceil(totalBytes / PART_BYTES), + }; +} + +function urlFor(env, input, arm, runId, object) { + const artifact = input.runtimeArtifacts.broker[arm]; + // Validate identifiers independently before constructing a trusted-origin URL. + binding(input, arm, runId, 1); + const base = new URL(env.CLOUD_API_URL); + if ( + !['https:', 'http:'].includes(base.protocol) || + base.username || + base.password || + base.search || + base.hash + ) { + throw new Error('Invalid broker storage origin'); + } + const key = `pr-proof/${input.handoffNonce}/brokers/${arm}/${artifact.sourceSha}/${artifact.sha256}/${object}`; + return new URL( + `api/v1/workflows/runs/${encodeURIComponent(runId)}/storage/${key}`, + base.href.replace(/\/?$/, '/') + ); +} + +async function jsonResponse(response, limit) { + if (!response.body) throw new Error('Empty broker transfer response'); + const reader = response.body.getReader(); + const chunks = []; + let length = 0; + try { + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + length += value.length; + if (length > limit) throw new Error('Broker transfer response exceeds its size limit'); + chunks.push(value); + } + try { + return JSON.parse(Buffer.concat(chunks).toString('utf8')); + } catch { + throw new Error('Invalid broker transfer JSON'); + } + } finally { + await reader.cancel().catch(() => {}); + } +} + +function client(env, tokenName, options) { + const token = env[tokenName]?.trim(); + if (!token) throw new Error(`Broker transfer requires ${tokenName}`); + const lifetime = AbortSignal.timeout(options.timeoutMs ?? 120_000); + return async (url, body, exclusive = false) => { + try { + return await (options.fetchImpl ?? fetch)(url, { + method: body === undefined ? 'GET' : 'PUT', + redirect: 'error', + signal: AbortSignal.any([ + lifetime, + options.signal ?? new AbortController().signal, + AbortSignal.timeout(15_000), + ]), + headers: { + authorization: `Bearer ${token}`, + 'content-type': 'application/json', + ...(exclusive ? { 'if-none-match': '*' } : {}), + }, + ...(body === undefined ? {} : { body }), + }); + } catch { + throw new Error('Broker storage request failed'); + } + }; +} + +/** Prepare and verify both artifacts locally before starting the remote run. */ +export async function prepareBrokerTransfer(inputPath, options = {}) { + const text = await readFile(inputPath, 'utf8'); + if (Buffer.byteLength(text) > 128 * 1024) throw new Error('Proof input exceeds its size limit'); + const input = validateProofInput(JSON.parse(text)); + if (!input.runtimeArtifacts?.broker) return null; + const artifacts = []; + for (const arm of ['base', 'head']) { + const loaded = await loadBrokerArtifact({ arm, expectedSha: input[`${arm}Sha`], root: options.root }); + if (!equal(loaded.artifact, input.runtimeArtifacts.broker[arm])) + throw new Error('Broker upload binding changed'); + artifacts.push({ arm, bytes: loaded.contents }); + } + if (artifacts.reduce((total, item) => total + item.bytes.length, 0) > MAX_PAIR_BYTES) { + throw new Error('Broker transfer exceeds aggregate size limit'); + } + return createBrokerTransfer(input, artifacts, options); +} + +export function createBrokerTransfer(input, artifacts, options = {}) { + if ( + artifacts.length !== 2 || + artifacts[0].arm !== 'base' || + artifacts[1].arm !== 'head' || + artifacts.reduce((total, item) => total + item.bytes.length, 0) > MAX_PAIR_BYTES + ) { + throw new Error('Broker transfer requires one bounded artifact per arm'); + } + let ownedRunId; + let pending; + let cleaned = false; + const env = options.env ?? process.env; + const uploadAbort = new AbortController(); + const request = client(env, 'CLOUD_API_KEY', { + ...options, + signal: AbortSignal.any([uploadAbort.signal, options.signal ?? new AbortController().signal]), + }); + const objectNames = (item) => [ + 'manifest.json', + ...Array.from({ length: Math.ceil(item.bytes.length / PART_BYTES) }, (_, index) => `part-${index}.json`), + ]; + for (const item of artifacts) { + binding(input, item.arm, 'validation', item.bytes.length); + if (sha256(item.bytes) !== input.runtimeArtifacts.broker[item.arm].sha256) + throw new Error('Broker upload hash mismatch'); + } + return { + start(runId) { + if (ownedRunId && ownedRunId !== runId) throw new Error('Broker upload run changed'); + if (cleaned) throw new Error('Broker upload already cleaned'); + ownedRunId = runId; + pending ??= (async () => { + for (const item of artifacts) { + const identity = binding(input, item.arm, runId, item.bytes.length); + for (let index = 0; index < identity.count; index++) { + const bytes = item.bytes.subarray(index * PART_BYTES, (index + 1) * PART_BYTES); + const body = JSON.stringify({ + ...identity, + index, + partSha256: sha256(bytes), + data: bytes.toString('base64'), + }); + if (Buffer.byteLength(body) > MAX_PART_JSON) + throw new Error('Broker chunk exceeds its size limit'); + const response = await request( + urlFor(env, input, item.arm, runId, `part-${index}.json`), + body, + true + ); + if (!response.ok) throw new Error(`Broker chunk upload refused (${response.status})`); + } + // Publish the readiness marker only after every exclusive part write succeeded. + const response = await request( + urlFor(env, input, item.arm, runId, 'manifest.json'), + JSON.stringify(identity), + true + ); + if (!response.ok) throw new Error(`Broker manifest upload refused (${response.status})`); + } + })(); + return pending; + }, + release() { + cleaned = true; + for (const item of artifacts) item.bytes.fill(0); + }, + async cleanup() { + if (cleaned) return; + cleaned = true; + uploadAbort.abort(); + await pending?.catch(() => {}); + if (!ownedRunId) return; + let failed = false; + const cleanupRequest = client(env, 'CLOUD_API_KEY', { + ...options, + signal: undefined, + timeoutMs: 60_000, + }); + // Cloud has no object DELETE route. Tombstone only this nonce/run's keys, + // manifest first, before cancellation revokes the prepared-run write grant. + for (const item of artifacts) + for (const object of objectNames(item)) { + try { + const response = await cleanupRequest(urlFor(env, input, item.arm, ownedRunId, object), ''); + if (!response.ok) failed = true; + } catch { + failed = true; + } + } + for (const item of artifacts) item.bytes.fill(0); + if (failed) throw new Error('Broker transfer cleanup incomplete'); + }, + }; +} + +/** Read each bounded part once, then validate the original build hash in memory. */ +export async function downloadBrokerArtifact(input, arm, options = {}) { + const env = options.env ?? process.env; + if ( + env.RUN_ID && + env.AGENT_RELAY_CLOUD_WORKER_RUN_ID && + env.RUN_ID !== env.AGENT_RELAY_CLOUD_WORKER_RUN_ID + ) { + throw new Error('Conflicting broker storage run IDs'); + } + const runId = env.RUN_ID || env.AGENT_RELAY_CLOUD_WORKER_RUN_ID; + const request = client(env, 'CLOUD_API_ACCESS_TOKEN', options); + const deadline = Date.now() + 60_000; + let manifest; + for (;;) { + const response = await request(urlFor(env, input, arm, runId, 'manifest.json')); + if (response.status !== 404) { + if (!response.ok) throw new Error(`Broker manifest unavailable (${response.status})`); + manifest = await jsonResponse(response, 4096); + break; + } + if (Date.now() >= deadline) throw new Error('Broker manifest was not published before transfer deadline'); + await delay(250, undefined, { signal: options.signal }); + } + const identity = binding(input, arm, runId, manifest?.totalBytes); + try { + if (!equal(manifest, identity)) throw new Error('Broker manifest binding mismatch'); + const chunks = []; + for (let index = 0; index < identity.count; index++) { + const response = await request(urlFor(env, input, arm, runId, `part-${index}.json`)); + if (!response.ok) throw new Error(`Broker part unavailable (${response.status})`); + const part = await jsonResponse(response, MAX_PART_JSON); + const { data, partSha256, index: receivedIndex, ...receivedIdentity } = part ?? {}; + if (!equal(receivedIdentity, identity) || receivedIndex !== index || typeof data !== 'string') { + throw new Error('Broker part binding mismatch'); + } + const bytes = Buffer.from(data, 'base64'); + const expectedLength = Math.min(PART_BYTES, identity.totalBytes - index * PART_BYTES); + if ( + bytes.length !== expectedLength || + bytes.toString('base64') !== data || + sha256(bytes) !== partSha256 + ) { + throw new Error('Broker part length, encoding, or hash mismatch'); + } + chunks.push(bytes); + } + const contents = Buffer.concat(chunks); + if (contents.length !== identity.totalBytes || sha256(contents) !== identity.sha256) + throw new Error('Broker final build hash mismatch'); + return { artifact: input.runtimeArtifacts.broker[arm], contents }; + } finally { + // Consume the transfer before the arm returns. Run completion/cancellation + // may revoke its token, so cleanup cannot be deferred to the dispatcher. + const cleanupRequest = client(env, 'CLOUD_API_ACCESS_TOKEN', { + ...options, + signal: undefined, + timeoutMs: 60_000, + }); + let failed = false; + for (const object of [ + 'manifest.json', + ...Array.from({ length: identity.count }, (_, index) => `part-${index}.json`), + ]) { + try { + const response = await cleanupRequest(urlFor(env, input, arm, runId, object), ''); + if (!response.ok) failed = true; + } catch { + failed = true; + } + } + if (failed) throw new Error('Broker transfer consumption cleanup incomplete'); + } +} diff --git a/scripts/pr-proof/broker-transfer.test.mjs b/scripts/pr-proof/broker-transfer.test.mjs new file mode 100644 index 0000000000..bbc72c31ed --- /dev/null +++ b/scripts/pr-proof/broker-transfer.test.mjs @@ -0,0 +1,309 @@ +import assert from 'node:assert/strict'; +import { createHash } from 'node:crypto'; +import { describe, it } from 'node:test'; +import { + createBrokerTransfer, + downloadBrokerArtifact, + PART_BYTES, + MAX_PAIR_BYTES, +} from './broker-transfer.mjs'; + +function fixture() { + const bytes = Buffer.alloc(PART_BYTES + 31, 7); + const digest = createHash('sha256').update(bytes).digest('hex'); + const input = { + baseSha: 'a'.repeat(40), + headSha: 'b'.repeat(40), + handoffNonce: 'c'.repeat(32), + runtimeArtifacts: { broker: {} }, + }; + for (const arm of ['base', 'head']) + input.runtimeArtifacts.broker[arm] = { + path: `.relayflow/pr-proof-binaries/${arm}/agent-relay-broker`, + sha256: digest, + sourceSha: input[`${arm}Sha`], + }; + const objects = new Map(); + const calls = []; + const fetchImpl = async (url, init) => { + const key = new URL(url).pathname; + calls.push({ key, init }); + assert.equal(init.redirect, 'error'); + if (init.method === 'PUT') { + assert.ok(['Bearer ci-fixture', 'Bearer sandbox-fixture'].includes(init.headers.authorization)); + if (init.headers['if-none-match'] === '*' && objects.has(key)) return new Response('', { status: 412 }); + objects.set(key, init.body); + return new Response('{}'); + } + assert.equal(init.headers.authorization, 'Bearer sandbox-fixture'); + return new Response(objects.get(key) ?? '', { status: objects.has(key) ? 200 : 404 }); + }; + const env = { + CLOUD_API_URL: 'https://cloud.example', + CLOUD_API_KEY: 'ci-fixture', + CLOUD_API_ACCESS_TOKEN: 'sandbox-fixture', + RUN_ID: 'run-fixture', + }; + const artifacts = ['base', 'head'].map((arm) => ({ arm, bytes: Buffer.from(bytes) })); + const options = { env, fetchImpl }; + return { + input, + objects, + calls, + bytes, + env, + options, + artifacts, + transfer: createBrokerTransfer(input, artifacts, options), + read: (arm) => downloadBrokerArtifact(input, arm, options), + }; +} + +describe('bounded run-scoped broker transfer', () => { + it('uploads once, publishes readiness last, reconstructs each part once, and tombstones on cleanup', async () => { + const f = fixture(); + const first = f.transfer.start('run-fixture'); + assert.equal(f.transfer.start('run-fixture'), first); + await first; + assert.throws(() => f.transfer.start('different-run'), /run changed/); + for (const arm of ['base', 'head']) { + const writes = f.calls.filter((call) => call.key.includes(`/brokers/${arm}/`)); + assert.match(writes.at(-1).key, /manifest.json$/); + assert.equal(writes.length, 3); + assert.ok(writes.every((call) => Buffer.byteLength(call.init.body) < 3 * 1024 * 1024)); + const actual = await f.read(arm); + assert.deepEqual(actual.contents, f.bytes); + assert.deepEqual(actual.artifact, f.input.runtimeArtifacts.broker[arm]); + assert.equal( + f.calls.filter((call) => call.init.method === 'GET' && call.key.includes(`/brokers/${arm}/`)).length, + 3 + ); + } + assert.ok( + [...f.objects.values()].every( + (body) => !body.includes('ci-fixture') && !body.includes('sandbox-fixture') + ) + ); + await f.transfer.cleanup(); + await f.transfer.cleanup(); + assert.ok([...f.objects.values()].every((body) => body === '')); + await assert.rejects(f.read('base'), /Invalid broker transfer JSON/); + assert.throws(() => f.transfer.start('run-fixture'), /already cleaned/); + }); + for (const [field, value] of [ + ['runId', 'other-run'], + ['arm', 'head'], + ['sourceSha', 'd'.repeat(40)], + ['sha256', 'd'.repeat(64)], + ['totalBytes', 10], + ['count', 8], + ['index', 1], + ['nonce', 'd'.repeat(32)], + ]) { + it(`rejects a chunk with changed ${field}`, async () => { + const f = fixture(); + await f.transfer.start('run-fixture'); + const key = [...f.objects.keys()].find((key) => key.includes('/base/') && key.endsWith('part-0.json')); + const part = JSON.parse(f.objects.get(key)); + part[field] = value; + f.objects.set(key, JSON.stringify(part)); + await assert.rejects(f.read('base'), /binding mismatch/); + }); + } + it('rejects missing/duplicate parts, malformed base64, corrupt bytes, and oversized bodies', async () => { + for (const mutation of ['missing', 'duplicate', 'encoding', 'bytes', 'oversized']) { + const f = fixture(); + await f.transfer.start('run-fixture'); + const key = [...f.objects.keys()].find((key) => key.includes('/base/') && key.endsWith('part-1.json')); + const part = JSON.parse(f.objects.get(key)); + if (mutation === 'missing') f.objects.delete(key); + else if (mutation === 'duplicate') f.objects.set(key, f.objects.get(key.replace('part-1', 'part-0'))); + else if (mutation === 'oversized') f.objects.set(key, ' '.repeat(3 * 1024 * 1024 + 1)); + else { + part.data = mutation === 'encoding' ? '!!' + part.data : Buffer.alloc(31, 8).toString('base64'); + f.objects.set(key, JSON.stringify(part)); + } + await assert.rejects(f.read('base'), /unavailable|mismatch|size limit/); + } + }); + it('checks final build SHA even if a changed part has a self-consistent checksum', async () => { + const f = fixture(); + await f.transfer.start('run-fixture'); + const key = [...f.objects.keys()].find((key) => key.includes('/base/') && key.endsWith('part-1.json')); + const part = JSON.parse(f.objects.get(key)); + const forged = Buffer.alloc(31, 8); + part.data = forged.toString('base64'); + part.partSha256 = createHash('sha256').update(forged).digest('hex'); + f.objects.set(key, JSON.stringify(part)); + await assert.rejects(f.read('base'), /final build hash mismatch/); + }); + it('refuses size violations and hash changes before any PUT', () => { + const f = fixture(); + const tooLarge = [{ arm: 'base', bytes: { length: MAX_PAIR_BYTES } }, f.artifacts[1]]; + assert.throws(() => createBrokerTransfer(f.input, tooLarge, f.options), /bounded artifact/); + f.artifacts[0].bytes[0] = 9; + assert.throws(() => createBrokerTransfer(f.input, f.artifacts, f.options), /hash mismatch/); + assert.equal(f.calls.length, 0); + }); + it('does not publish a manifest after a partial or ambiguous upload, then cleans attempted keys', async () => { + for (const ambiguous of [false, true]) { + const f = fixture(); + let uploads = 0; + const transfer = createBrokerTransfer(f.input, f.artifacts, { + env: f.env, + fetchImpl: async (url, init) => { + if (init.headers['if-none-match'] === '*' && ++uploads === 2) { + if (ambiguous) { + await f.options.fetchImpl(url, init); + throw new Error('ci-fixture secret'); + } + return new Response('', { status: 412 }); + } + return f.options.fetchImpl(url, init); + }, + }); + await assert.rejects( + transfer.start('run-fixture'), + /Broker (chunk upload refused|storage request failed)/ + ); + assert.ok([...f.objects.keys()].every((key) => !key.endsWith('manifest.json'))); + await transfer.cleanup(); + assert.ok([...f.objects.values()].every((body) => body === '')); + } + }); + it('does not return executable bytes when consumption cleanup is refused', async () => { + const f = fixture(); + await f.transfer.start('run-fixture'); + await assert.rejects( + downloadBrokerArtifact(f.input, 'base', { + env: f.env, + fetchImpl: async (url, init) => { + if (init.method === 'PUT') return new Response('private-error-body', { status: 403 }); + return f.options.fetchImpl(url, init); + }, + }), + /consumption cleanup incomplete/ + ); + }); + it('waits for a racing manifest and refuses conflicting sandbox run identities', async () => { + const f = fixture(); + const reading = f.read('base'); + await f.transfer.start('run-fixture'); + assert.deepEqual((await reading).contents, f.bytes); + await assert.rejects( + downloadBrokerArtifact(f.input, 'base', { + ...f.options, + env: { ...f.env, AGENT_RELAY_CLOUD_WORKER_RUN_ID: 'other' }, + }), + /Conflicting/ + ); + }); +}); + +import { mkdtemp, mkdir, writeFile, readFile, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +import { main as runCloud } from './run-cloud.mjs'; + +describe('trusted dispatcher transfer integration', () => { + for (const failUpload of [false, true]) { + it(`uses prepared identity and cleans after ${failUpload ? 'upload failure and cancellation' : 'terminal success'}`, async () => { + const f = fixture(); + const directory = await mkdtemp(path.join(tmpdir(), 'broker-transfer-dispatch-')); + const originalDirectory = process.cwd(); + const originalEnvironment = { ...process.env }; + const originalFetch = globalThis.fetch; + try { + const caseId = 'fixture-broker-transfer'; + const input = { + ...f.input, + version: 1, + repository: 'AgentWorkforce/relay', + pullRequest: 7, + caseId, + kind: 'feature', + manifest: { + version: 1, + id: caseId, + kind: 'feature', + title: 'Fixture broker task', + runner: { command: ['node', `tests/relayflows/cases/${caseId}/run.mjs`] }, + requirements: ['broker-linux-x64'], + timeoutSeconds: 180, + expected: { + base: { outcome: 'absent', signature: 'fixture_absent' }, + head: { outcome: 'fixed', signature: 'fixture_fixed' }, + }, + }, + }; + for (const item of f.artifacts) { + const root = path.join(directory, '.relayflow/pr-proof-binaries', item.arm); + await mkdir(root, { recursive: true }); + await writeFile(path.join(root, 'agent-relay-broker'), item.bytes); + await writeFile( + path.join(root, 'broker-manifest.json'), + JSON.stringify({ + version: 1, + sourceSha: input[`${item.arm}Sha`], + sha256: input.runtimeArtifacts.broker[item.arm].sha256, + }) + ); + } + await writeFile(path.join(directory, '.relayflow/pr-proof-input.json'), JSON.stringify(input)); + const commands = path.join(directory, 'commands.jsonl'); + const cli = path.join(directory, 'fake-cli.mjs'); + await writeFile( + cli, + `#!/usr/bin/env node\nimport {appendFileSync} from 'node:fs';\nconst command=process.argv[3];appendFileSync(${JSON.stringify(commands)},JSON.stringify(process.argv.slice(2))+'\\n');\nif(command==='run'){console.error('AGENT_RELAY_CLOUD_PREPARED_RUN_ID=run-fixture');console.log(JSON.stringify({runId:'run-fixture'}));}\nelse if(command==='status')console.log(JSON.stringify({status:'completed'}));\nelse if(command==='logs')console.log('fixture log');\nelse if(command==='cancel')console.log('{}');\nelse process.exitCode=2;`, + { mode: 0o700 } + ); + process.chdir(directory); + Object.assign(process.env, f.env, { + PR_PROOF_AGENT_RELAY_BIN: cli, + PR_PROOF_POLL_MS: '1000', + PR_PROOF_CLOUD_LOG_PATH: path.join(directory, 'cloud.log'), + }); + for (const key of ['PR_PROOF_INPUT_PATH', 'GITHUB_OUTPUT', 'GITHUB_STEP_SUMMARY']) + delete process.env[key]; + globalThis.fetch = async (url, init) => { + if (failUpload && init.headers['if-none-match'] === '*') return new Response('', { status: 503 }); + if (init.method === 'PUT' && init.body === '' && failUpload) { + assert.doesNotMatch(await readFile(commands, 'utf8'), /"cancel"/); + } + const response = await f.options.fetchImpl(url, init); + if ( + !failUpload && + init.headers['if-none-match'] === '*' && + String(url).includes('/head/') && + String(url).endsWith('/manifest.json') + ) { + await f.read('base'); + await f.read('head'); + } + return response; + }; + if (failUpload) await assert.rejects(runCloud(), /Broker chunk upload refused/); + else await runCloud(); + const invocations = (await readFile(commands, 'utf8')) + .trim() + .split('\n') + .map((line) => JSON.parse(line)); + assert.equal(invocations.filter((args) => args[1] === 'run').length, 1); + assert.equal(invocations.filter((args) => args[1] === 'cancel').length, failUpload ? 1 : 0); + assert.ok([...f.objects.values()].every((body) => body === '')); + if (!failUpload) + assert.doesNotMatch( + await readFile(path.join(directory, 'cloud.log'), 'utf8'), + /ci-fixture|sandbox-fixture/ + ); + } finally { + globalThis.fetch = originalFetch; + process.chdir(originalDirectory); + for (const key of Object.keys(process.env)) + if (!(key in originalEnvironment)) delete process.env[key]; + Object.assign(process.env, originalEnvironment); + await rm(directory, { recursive: true, force: true }); + } + }); + } +}); diff --git a/scripts/pr-proof/run-arm.mjs b/scripts/pr-proof/run-arm.mjs index 1b092f838b..4511a29095 100644 --- a/scripts/pr-proof/run-arm.mjs +++ b/scripts/pr-proof/run-arm.mjs @@ -20,6 +20,7 @@ import { import { uploadCloudEvidence, validateCloudEvidenceEnvironment } from './cloud-storage.mjs'; import { runBoundedProcess } from './process-runner.mjs'; import { loadBrokerArtifact } from './stage-broker-artifacts.mjs'; +import { downloadBrokerArtifact } from './broker-transfer.mjs'; const MAX_OBSERVATION_FILE_BYTES = 64 * 1024; const LANDLOCK_PROBE_TIMEOUT_MS = 5_000; @@ -481,7 +482,13 @@ export async function runLandlockedProcess(command, args, { writableRoots, ...op return runProcess(invocation.command, invocation.args, { ...options, env: caseEnvironment }); } -export async function openVerifiedBrokerExecutable({ input, arm, privateRoot, root = process.cwd() }) { +export async function openVerifiedBrokerExecutable({ + input, + arm, + privateRoot, + root = process.cwd(), + loadArtifact = loadBrokerArtifact, +}) { const artifact = input.runtimeArtifacts?.broker?.[arm]; if (!artifact) return null; if (process.platform !== 'linux') { @@ -490,7 +497,7 @@ export async function openVerifiedBrokerExecutable({ input, arm, privateRoot, ro const artifactPath = path.resolve(root, artifact.path); const expectedRoot = path.join(root, '.relayflow', 'pr-proof-binaries') + path.sep; if (!artifactPath.startsWith(expectedRoot)) throw new Error('broker artifact escaped its staging root'); - const loaded = await loadBrokerArtifact({ + const loaded = await loadArtifact({ arm, expectedSha: arm === 'base' ? input.baseSha : input.headSha, root, @@ -633,6 +640,7 @@ export async function main() { input, arm, privateRoot: temporaryRoot, + loadArtifact: () => downloadBrokerArtifact(input, arm), }); const brokerPath = brokerExecutable?.path ?? null; diff --git a/scripts/pr-proof/run-cloud.mjs b/scripts/pr-proof/run-cloud.mjs index 903290a043..2915822d5c 100644 --- a/scripts/pr-proof/run-cloud.mjs +++ b/scripts/pr-proof/run-cloud.mjs @@ -5,6 +5,7 @@ import path from 'node:path'; import { setTimeout as delay } from 'node:timers/promises'; import { pathToFileURL } from 'node:url'; +import { prepareBrokerTransfer } from './broker-transfer.mjs'; import { runBoundedProcess } from './process-runner.mjs'; const TERMINAL_SUCCESS = new Set(['completed', 'succeeded', 'success']); @@ -593,8 +594,14 @@ export async function main() { label: 'PR_PROOF_CLOUD_COMMAND_TIMEOUT_MS', }); const auth = createCliApiKeyEnvironment(process.env); + const brokerTransfer = await prepareBrokerTransfer( + process.env.PR_PROOF_INPUT_PATH ?? '.relayflow/pr-proof-input.json', + { env: auth.cliEnv } + ); + let transferPromise; let runId = null; let terminal = false; + let proofSucceeded = false; let cancelPromise = null; let shuttingDown = false; let activeCommandController = null; @@ -608,6 +615,12 @@ export async function main() { throw new Error(`Cloud prepare/run ID mismatch: ${runId} != ${preparedRunId}`); } runId = preparedRunId; + transferPromise ??= brokerTransfer?.start(runId); + transferPromise?.catch((error) => { + // The existing bounded launch finishes before cancellation/cleanup. + // Keep the original upload failure even if CLI output arrives later. + launchProgressError ??= error; + }); } catch (error) { launchProgressError ??= error; } @@ -667,6 +680,7 @@ export async function main() { shuttingDown = true; activeCommandController?.abort(); void (async () => { + await brokerTransfer?.cleanup().catch(() => console.warn('Broker transfer cleanup incomplete')); await cancelRemote(signal).catch((error) => console.warn(sanitizeCloudCommandOutput(error.message, auth.diagnosticSecretValues)) ); @@ -687,6 +701,8 @@ export async function main() { onStderr: (text) => captureLaunchProgressError(() => launchProgress.write(text)), }); captureLaunchProgressError(() => launchProgress.end()); + if (brokerTransfer && !transferPromise) throw new Error('Broker transfer requires prepared run identity'); + await transferPromise; if (launchProgressError) throw launchProgressError; if (launch.aborted) throw new Error('Cloud workflow submission was interrupted'); if (launch.timedOut) { @@ -801,6 +817,7 @@ export async function main() { if (!TERMINAL_SUCCESS.has(terminalStatus)) { throw new Error(`Cloud RelayFlow finished with status ${terminalStatus}`); } + proofSucceeded = true; if (process.env.GITHUB_STEP_SUMMARY) { await appendFile( process.env.GITHUB_STEP_SUMMARY, @@ -814,10 +831,16 @@ export async function main() { activeCommandController?.abort(); process.removeListener('SIGINT', signalHandler); process.removeListener('SIGTERM', signalHandler); - if (runId && !terminal) - await cancelRemote('dispatcher exiting').catch((error) => - console.warn(sanitizeCloudCommandOutput(error.message, auth.diagnosticSecretValues)) - ); + try { + if (terminal) + brokerTransfer?.release(); // each completed arm consumed its transfer + else await brokerTransfer?.cleanup(); + } finally { + if (runId && !terminal) + await cancelRemote('dispatcher exiting').catch((error) => + console.warn(sanitizeCloudCommandOutput(error.message, auth.diagnosticSecretValues)) + ); + } } } From 51e6f399ccad84e582a1886443e7bc2e164abe8d Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 15 Sep 2026 01:37:44 +0200 Subject: [PATCH 2/3] fix(ci): require secure broker transfer and clean failed proofs Session-Id: 01a09c40-ce3b-7f11-a7df-b6b7ccab6fd9 --- docs/pr-proof-broker-transfer.md | 3 ++ scripts/pr-proof/broker-transfer.mjs | 6 +++ scripts/pr-proof/broker-transfer.test.mjs | 47 +++++++++++++++++++++-- scripts/pr-proof/run-cloud.mjs | 2 +- 4 files changed, 53 insertions(+), 5 deletions(-) diff --git a/docs/pr-proof-broker-transfer.md b/docs/pr-proof-broker-transfer.md index c07ef7d46a..ae1cf64bc0 100644 --- a/docs/pr-proof-broker-transfer.md +++ b/docs/pr-proof-broker-transfer.md @@ -41,3 +41,6 @@ This prerequisite must be admitted to the trusted base before rerunning another PR's `pull_request_target` proof. A change on that PR's head cannot repair the base dispatcher. Resolve and attest both exact binaries again after admission. Local mocked transfer checks do not establish a successful Cloud proof. + +Bearer transport requires HTTPS except literal loopback IP addresses for local +tests. Hostnames that merely resemble loopback addresses are refused. diff --git a/scripts/pr-proof/broker-transfer.mjs b/scripts/pr-proof/broker-transfer.mjs index 4f7e835d7d..ad847036f5 100644 --- a/scripts/pr-proof/broker-transfer.mjs +++ b/scripts/pr-proof/broker-transfer.mjs @@ -1,5 +1,6 @@ /** Run-scoped broker transfer. No artifact bytes or credentials enter code sync. */ import { createHash } from 'node:crypto'; +import { isIP } from 'node:net'; import { readFile } from 'node:fs/promises'; import { setTimeout as delay } from 'node:timers/promises'; import { validateProofInput } from './contract.mjs'; @@ -53,6 +54,11 @@ function urlFor(env, input, arm, runId, object) { ) { throw new Error('Invalid broker storage origin'); } + const loopback = + base.hostname === '[::1]' || (isIP(base.hostname) === 4 && base.hostname.startsWith('127.')); + if (base.protocol !== 'https:' && !loopback) { + throw new Error('Broker transfer credentials require HTTPS outside literal loopback addresses'); + } const key = `pr-proof/${input.handoffNonce}/brokers/${arm}/${artifact.sourceSha}/${artifact.sha256}/${object}`; return new URL( `api/v1/workflows/runs/${encodeURIComponent(runId)}/storage/${key}`, diff --git a/scripts/pr-proof/broker-transfer.test.mjs b/scripts/pr-proof/broker-transfer.test.mjs index bbc72c31ed..6ce70b618c 100644 --- a/scripts/pr-proof/broker-transfer.test.mjs +++ b/scripts/pr-proof/broker-transfer.test.mjs @@ -60,6 +60,31 @@ function fixture() { } describe('bounded run-scoped broker transfer', () => { + for (const origin of ['http://cloud.example', 'http://192.168.1.1', 'http://127.0.0.1.example']) { + it(`refuses bearer transfer to ${origin} before any request`, async () => { + const f = fixture(); + const transfer = createBrokerTransfer(f.input, f.artifacts, { + ...f.options, + env: { ...f.env, CLOUD_API_URL: origin }, + }); + await assert.rejects(transfer.start('run-fixture'), /require HTTPS/); + assert.equal(f.calls.length, 0); + }); + } + for (const origin of ['http://127.0.0.1:1234', 'http://127.9.8.7:1234', 'http://[::1]:1234']) { + it(`permits literal loopback fixture ${origin}`, async () => { + const f = fixture(); + const env = { ...f.env, CLOUD_API_URL: origin }; + const transfer = createBrokerTransfer(f.input, f.artifacts, { ...f.options, env }); + await transfer.start('run-fixture'); + assert.deepEqual( + (await downloadBrokerArtifact(f.input, 'base', { ...f.options, env })).contents, + f.bytes + ); + await transfer.cleanup(); + }); + } + it('uploads once, publishes readiness last, reconstructs each part once, and tombstones on cleanup', async () => { const f = fixture(); const first = f.transfer.start('run-fixture'); @@ -206,13 +231,17 @@ import path from 'node:path'; import { main as runCloud } from './run-cloud.mjs'; describe('trusted dispatcher transfer integration', () => { - for (const failUpload of [false, true]) { - it(`uses prepared identity and cleans after ${failUpload ? 'upload failure and cancellation' : 'terminal success'}`, async () => { + for (const scenario of ['completed', 'upload-failure', 'failed', 'cancelled', 'error', 'failed-revoked']) { + const failUpload = scenario === 'upload-failure'; + const revoked = scenario === 'failed-revoked'; + const terminalStatus = revoked ? 'failed' : failUpload ? 'completed' : scenario; + it(`preserves cleanup authority for ${scenario}`, async () => { const f = fixture(); const directory = await mkdtemp(path.join(tmpdir(), 'broker-transfer-dispatch-')); const originalDirectory = process.cwd(); const originalEnvironment = { ...process.env }; const originalFetch = globalThis.fetch; + let cleanupAttempts = 0; try { const caseId = 'fixture-broker-transfer'; const input = { @@ -254,7 +283,7 @@ describe('trusted dispatcher transfer integration', () => { const cli = path.join(directory, 'fake-cli.mjs'); await writeFile( cli, - `#!/usr/bin/env node\nimport {appendFileSync} from 'node:fs';\nconst command=process.argv[3];appendFileSync(${JSON.stringify(commands)},JSON.stringify(process.argv.slice(2))+'\\n');\nif(command==='run'){console.error('AGENT_RELAY_CLOUD_PREPARED_RUN_ID=run-fixture');console.log(JSON.stringify({runId:'run-fixture'}));}\nelse if(command==='status')console.log(JSON.stringify({status:'completed'}));\nelse if(command==='logs')console.log('fixture log');\nelse if(command==='cancel')console.log('{}');\nelse process.exitCode=2;`, + `#!/usr/bin/env node\nimport {appendFileSync} from 'node:fs';\nconst command=process.argv[3];appendFileSync(${JSON.stringify(commands)},JSON.stringify(process.argv.slice(2))+'\\n');\nif(command==='run'){console.error('AGENT_RELAY_CLOUD_PREPARED_RUN_ID=run-fixture');console.log(JSON.stringify({runId:'run-fixture'}));}\nelse if(command==='status')console.log(JSON.stringify({status:${JSON.stringify(terminalStatus)}}));\nelse if(command==='logs')console.log('fixture log');\nelse if(command==='cancel')console.log('{}');\nelse process.exitCode=2;`, { mode: 0o700 } ); process.chdir(directory); @@ -266,6 +295,10 @@ describe('trusted dispatcher transfer integration', () => { for (const key of ['PR_PROOF_INPUT_PATH', 'GITHUB_OUTPUT', 'GITHUB_STEP_SUMMARY']) delete process.env[key]; globalThis.fetch = async (url, init) => { + if (init.method === 'PUT' && init.body === '') { + cleanupAttempts++; + if (revoked) return new Response('', { status: 403 }); + } if (failUpload && init.headers['if-none-match'] === '*') return new Response('', { status: 503 }); if (init.method === 'PUT' && init.body === '' && failUpload) { assert.doesNotMatch(await readFile(commands, 'utf8'), /"cancel"/); @@ -273,6 +306,7 @@ describe('trusted dispatcher transfer integration', () => { const response = await f.options.fetchImpl(url, init); if ( !failUpload && + scenario === 'completed' && init.headers['if-none-match'] === '*' && String(url).includes('/head/') && String(url).endsWith('/manifest.json') @@ -283,6 +317,9 @@ describe('trusted dispatcher transfer integration', () => { return response; }; if (failUpload) await assert.rejects(runCloud(), /Broker chunk upload refused/); + else if (revoked) await assert.rejects(runCloud(), /Broker transfer cleanup incomplete/); + else if (scenario !== 'completed') + await assert.rejects(runCloud(), /Cloud RelayFlow finished with status/); else await runCloud(); const invocations = (await readFile(commands, 'utf8')) .trim() @@ -290,7 +327,9 @@ describe('trusted dispatcher transfer integration', () => { .map((line) => JSON.parse(line)); assert.equal(invocations.filter((args) => args[1] === 'run').length, 1); assert.equal(invocations.filter((args) => args[1] === 'cancel').length, failUpload ? 1 : 0); - assert.ok([...f.objects.values()].every((body) => body === '')); + assert.ok(cleanupAttempts > 0 || scenario === 'completed'); + if (revoked) assert.ok([...f.objects.values()].some((body) => body !== '')); + else assert.ok([...f.objects.values()].every((body) => body === '')); if (!failUpload) assert.doesNotMatch( await readFile(path.join(directory, 'cloud.log'), 'utf8'), diff --git a/scripts/pr-proof/run-cloud.mjs b/scripts/pr-proof/run-cloud.mjs index 2915822d5c..4775ac26df 100644 --- a/scripts/pr-proof/run-cloud.mjs +++ b/scripts/pr-proof/run-cloud.mjs @@ -832,7 +832,7 @@ export async function main() { process.removeListener('SIGINT', signalHandler); process.removeListener('SIGTERM', signalHandler); try { - if (terminal) + if (proofSucceeded) brokerTransfer?.release(); // each completed arm consumed its transfer else await brokerTransfer?.cleanup(); } finally { From 747ffa738ed30d84bb34c730bc2e6d40c5183a82 Mon Sep 17 00:00:00 2001 From: Khaliq Gant Date: Mon, 14 Sep 2026 22:51:19 -0700 Subject: [PATCH 3/3] fix(ci): start transfer lifetime on upload and tombstone before cancel Address PR #1773 review findings from Bugbot and CodeRabbit: - The 120-second aggregate upload lifetime now starts inside start(), when the prepared run ID is known, instead of at createBrokerTransfer. A slow `cloud run` prepare no longer consumes the budget before the first exclusive PUT. - downloadBrokerArtifact runs consumption cleanup after capturing the primary download/validation/integrity error and rethrows that error when cleanup also fails; the cleanup failure is thrown only when the download itself succeeded, still before any bytes are returned. - The dispatcher's submission-timeout and poll-deadline paths tombstone the transfer before cancelRemote revokes the prepared-run write grant. A cleanup failure is warned and never replaces the dispatcher's own failure. Tests: the fake fetch now honors abort signals so the lifetime case bites; new cases cover lifetime start, cleanup-error precedence, and a hung submission whose tombstones must precede `cloud cancel`. All four affected cases fail on the previous sources and pass now (54/54). Co-Authored-By: Claude Opus 5 (1M context) --- docs/pr-proof-broker-transfer.md | 22 +++--- scripts/pr-proof/broker-transfer.mjs | 61 ++++++++++------- scripts/pr-proof/broker-transfer.test.mjs | 83 +++++++++++++++++++++-- scripts/pr-proof/run-cloud.mjs | 29 +++++++- 4 files changed, 150 insertions(+), 45 deletions(-) diff --git a/docs/pr-proof-broker-transfer.md b/docs/pr-proof-broker-transfer.md index ae1cf64bc0..5474c52160 100644 --- a/docs/pr-proof-broker-transfer.md +++ b/docs/pr-proof-broker-transfer.md @@ -13,7 +13,8 @@ acknowledgments fail the proof without replaying writes or creating a new run. The isolated arm waits at most 60 seconds for the manifest, reads each part once with bounded response sizes, and verifies the final raw build hash. -Requests time out after 15 seconds with a 120-second transfer lifetime. +Requests time out after 15 seconds with a 120-second transfer lifetime that +starts with the first upload or download request, not at preparation. Missing, duplicated, misbound, oversized, or corrupt parts fail closed. Only verified bytes enter the existing private mode-0500 executable, isolation, and post-execution attestation checks. Broker requirements and red/green @@ -27,15 +28,16 @@ case processes receive neither credentials nor transfer authorization headers. Each arm consumes and tombstones its manifest/parts immediately after bounded assembly while its sandbox token is live. Only then may the verified bytes reach the private executable. Private executable cleanup remains in the arm's -`finally` block. On submission failure/interruption, CI tombstones its nonce's -objects before cancellation revokes the prepared-run write grant. Successful -proof releases only local buffers because both arms already consumed storage. -Cleanup failures fail the dispatcher. Cloud has no general object-delete or -prepared-grant-revocation route; this change does not invent one or revoke -shared CI credentials. Existing cancellation revokes run API sessions. -Forced runner termination or independent remote token revocation can prevent -cleanup; remaining objects contain public binaries, remain run-scoped, and -follow platform storage lifecycle and token expiry. +`finally` block. On submission failure, timeout, or interruption, CI +tombstones its nonce's objects before cancellation revokes the prepared-run +write grant. Successful proof releases only local buffers because both arms +already consumed storage. Cleanup failures are reported without replacing +the original download or dispatcher failure. Cloud has no general +object-delete or prepared-grant-revocation route; this change does not invent +one or revoke shared CI credentials. Existing cancellation revokes run API +sessions. Forced runner termination or independent remote token revocation +can prevent cleanup; remaining objects contain public binaries, remain +run-scoped, and follow platform storage lifecycle and token expiry. This prerequisite must be admitted to the trusted base before rerunning another PR's `pull_request_target` proof. A change on that PR's head cannot repair the diff --git a/scripts/pr-proof/broker-transfer.mjs b/scripts/pr-proof/broker-transfer.mjs index ad847036f5..01da336bcc 100644 --- a/scripts/pr-proof/broker-transfer.mjs +++ b/scripts/pr-proof/broker-transfer.mjs @@ -149,10 +149,6 @@ export function createBrokerTransfer(input, artifacts, options = {}) { let cleaned = false; const env = options.env ?? process.env; const uploadAbort = new AbortController(); - const request = client(env, 'CLOUD_API_KEY', { - ...options, - signal: AbortSignal.any([uploadAbort.signal, options.signal ?? new AbortController().signal]), - }); const objectNames = (item) => [ 'manifest.json', ...Array.from({ length: Math.ceil(item.bytes.length / PART_BYTES) }, (_, index) => `part-${index}.json`), @@ -168,6 +164,13 @@ export function createBrokerTransfer(input, artifacts, options = {}) { if (cleaned) throw new Error('Broker upload already cleaned'); ownedRunId = runId; pending ??= (async () => { + // The aggregate transfer lifetime starts with the first upload, not at + // preparation: the prepared run ID may arrive well after this transfer + // was created, and that launch delay must not consume the budget. + const request = client(env, 'CLOUD_API_KEY', { + ...options, + signal: AbortSignal.any([uploadAbort.signal, options.signal ?? new AbortController().signal]), + }); for (const item of artifacts) { const identity = binding(input, item.arm, runId, item.bytes.length); for (let index = 0; index < identity.count; index++) { @@ -256,6 +259,8 @@ export async function downloadBrokerArtifact(input, arm, options = {}) { await delay(250, undefined, { signal: options.signal }); } const identity = binding(input, arm, runId, manifest?.totalBytes); + let result; + let failure; try { if (!equal(manifest, identity)) throw new Error('Broker manifest binding mismatch'); const chunks = []; @@ -281,27 +286,33 @@ export async function downloadBrokerArtifact(input, arm, options = {}) { const contents = Buffer.concat(chunks); if (contents.length !== identity.totalBytes || sha256(contents) !== identity.sha256) throw new Error('Broker final build hash mismatch'); - return { artifact: input.runtimeArtifacts.broker[arm], contents }; - } finally { - // Consume the transfer before the arm returns. Run completion/cancellation - // may revoke its token, so cleanup cannot be deferred to the dispatcher. - const cleanupRequest = client(env, 'CLOUD_API_ACCESS_TOKEN', { - ...options, - signal: undefined, - timeoutMs: 60_000, - }); - let failed = false; - for (const object of [ - 'manifest.json', - ...Array.from({ length: identity.count }, (_, index) => `part-${index}.json`), - ]) { - try { - const response = await cleanupRequest(urlFor(env, input, arm, runId, object), ''); - if (!response.ok) failed = true; - } catch { - failed = true; - } + result = { artifact: input.runtimeArtifacts.broker[arm], contents }; + } catch (error) { + failure = error; + } + // Consume the transfer before the arm returns. Run completion/cancellation + // may revoke its token, so cleanup cannot be deferred to the dispatcher. + const cleanupRequest = client(env, 'CLOUD_API_ACCESS_TOKEN', { + ...options, + signal: undefined, + timeoutMs: 60_000, + }); + let failed = false; + for (const object of [ + 'manifest.json', + ...Array.from({ length: identity.count }, (_, index) => `part-${index}.json`), + ]) { + try { + const response = await cleanupRequest(urlFor(env, input, arm, runId, object), ''); + if (!response.ok) failed = true; + } catch { + failed = true; } - if (failed) throw new Error('Broker transfer consumption cleanup incomplete'); } + // A download or integrity failure is the diagnostic that matters; a cleanup + // failure must not replace it. Refused cleanup after a good download still + // withholds the bytes. + if (failure) throw failure; + if (failed) throw new Error('Broker transfer consumption cleanup incomplete'); + return result; } diff --git a/scripts/pr-proof/broker-transfer.test.mjs b/scripts/pr-proof/broker-transfer.test.mjs index 6ce70b618c..a0addc3434 100644 --- a/scripts/pr-proof/broker-transfer.test.mjs +++ b/scripts/pr-proof/broker-transfer.test.mjs @@ -1,6 +1,7 @@ import assert from 'node:assert/strict'; import { createHash } from 'node:crypto'; import { describe, it } from 'node:test'; +import { setTimeout as delay } from 'node:timers/promises'; import { createBrokerTransfer, downloadBrokerArtifact, @@ -26,6 +27,8 @@ function fixture() { const objects = new Map(); const calls = []; const fetchImpl = async (url, init) => { + // Real fetch rejects on an already-aborted signal; the lifetime tests rely on it. + init.signal?.throwIfAborted(); const key = new URL(url).pathname; calls.push({ key, init }); assert.equal(init.redirect, 'error'); @@ -196,6 +199,50 @@ describe('bounded run-scoped broker transfer', () => { assert.ok([...f.objects.values()].every((body) => body === '')); } }); + it('starts the upload lifetime at start(), not at preparation, and still bounds the upload', async () => { + const f = fixture(); + const delayed = createBrokerTransfer(f.input, f.artifacts, { ...f.options, timeoutMs: 30 }); + await delay(80); // a slow prepare must not consume the transfer budget + await delayed.start('run-fixture'); + assert.ok([...f.objects.keys()].some((key) => key.endsWith('manifest.json'))); + await delayed.cleanup(); + + const g = fixture(); + const bounded = createBrokerTransfer(g.input, g.artifacts, { + ...g.options, + timeoutMs: 40, + fetchImpl: async (url, init) => { + await delay(30); + return g.options.fetchImpl(url, init); + }, + }); + await assert.rejects(bounded.start('run-fixture'), /Broker storage request failed/); + assert.ok([...g.objects.keys()].every((key) => !key.endsWith('manifest.json'))); + await bounded.cleanup(); + }); + it('keeps the download failure when consumption cleanup also fails', async () => { + const f = fixture(); + await f.transfer.start('run-fixture'); + const key = [...f.objects.keys()].find((key) => key.includes('/base/') && key.endsWith('part-0.json')); + const part = JSON.parse(f.objects.get(key)); + part.runId = 'other-run'; + f.objects.set(key, JSON.stringify(part)); + let cleanupAttempts = 0; + await assert.rejects( + downloadBrokerArtifact(f.input, 'base', { + env: f.env, + fetchImpl: async (url, init) => { + if (init.method === 'PUT') { + cleanupAttempts++; + return new Response('', { status: 403 }); + } + return f.options.fetchImpl(url, init); + }, + }), + /Broker part binding mismatch/ + ); + assert.equal(cleanupAttempts, 3); // manifest + both parts were still attempted + }); it('does not return executable bytes when consumption cleanup is refused', async () => { const f = fixture(); await f.transfer.start('run-fixture'); @@ -231,18 +278,31 @@ import path from 'node:path'; import { main as runCloud } from './run-cloud.mjs'; describe('trusted dispatcher transfer integration', () => { - for (const scenario of ['completed', 'upload-failure', 'failed', 'cancelled', 'error', 'failed-revoked']) { + for (const scenario of [ + 'completed', + 'upload-failure', + 'submission-timeout', + 'failed', + 'cancelled', + 'error', + 'failed-revoked', + ]) { const failUpload = scenario === 'upload-failure'; + const hangSubmission = scenario === 'submission-timeout'; const revoked = scenario === 'failed-revoked'; - const terminalStatus = revoked ? 'failed' : failUpload ? 'completed' : scenario; + const cancels = failUpload || hangSubmission; + const terminalStatus = revoked ? 'failed' : failUpload || hangSubmission ? 'completed' : scenario; it(`preserves cleanup authority for ${scenario}`, async () => { const f = fixture(); const directory = await mkdtemp(path.join(tmpdir(), 'broker-transfer-dispatch-')); const originalDirectory = process.cwd(); const originalEnvironment = { ...process.env }; const originalFetch = globalThis.fetch; + const originalWarn = console.warn; + const warnings = []; let cleanupAttempts = 0; try { + console.warn = (...args) => warnings.push(args.join(' ')); const caseId = 'fixture-broker-transfer'; const input = { ...f.input, @@ -283,13 +343,14 @@ describe('trusted dispatcher transfer integration', () => { const cli = path.join(directory, 'fake-cli.mjs'); await writeFile( cli, - `#!/usr/bin/env node\nimport {appendFileSync} from 'node:fs';\nconst command=process.argv[3];appendFileSync(${JSON.stringify(commands)},JSON.stringify(process.argv.slice(2))+'\\n');\nif(command==='run'){console.error('AGENT_RELAY_CLOUD_PREPARED_RUN_ID=run-fixture');console.log(JSON.stringify({runId:'run-fixture'}));}\nelse if(command==='status')console.log(JSON.stringify({status:${JSON.stringify(terminalStatus)}}));\nelse if(command==='logs')console.log('fixture log');\nelse if(command==='cancel')console.log('{}');\nelse process.exitCode=2;`, + `#!/usr/bin/env node\nimport {appendFileSync} from 'node:fs';\nconst command=process.argv[3];appendFileSync(${JSON.stringify(commands)},JSON.stringify(process.argv.slice(2))+'\\n');\nif(command==='run'){console.error('AGENT_RELAY_CLOUD_PREPARED_RUN_ID=run-fixture');if(${hangSubmission})setTimeout(()=>{},60_000);else console.log(JSON.stringify({runId:'run-fixture'}));}\nelse if(command==='status')console.log(JSON.stringify({status:${JSON.stringify(terminalStatus)}}));\nelse if(command==='logs')console.log('fixture log');\nelse if(command==='cancel')console.log('{}');\nelse process.exitCode=2;`, { mode: 0o700 } ); process.chdir(directory); Object.assign(process.env, f.env, { PR_PROOF_AGENT_RELAY_BIN: cli, PR_PROOF_POLL_MS: '1000', + PR_PROOF_CLOUD_COMMAND_TIMEOUT_MS: '1000', PR_PROOF_CLOUD_LOG_PATH: path.join(directory, 'cloud.log'), }); for (const key of ['PR_PROOF_INPUT_PATH', 'GITHUB_OUTPUT', 'GITHUB_STEP_SUMMARY']) @@ -300,7 +361,8 @@ describe('trusted dispatcher transfer integration', () => { if (revoked) return new Response('', { status: 403 }); } if (failUpload && init.headers['if-none-match'] === '*') return new Response('', { status: 503 }); - if (init.method === 'PUT' && init.body === '' && failUpload) { + if (init.method === 'PUT' && init.body === '' && cancels) { + // Tombstones must land before cancellation revokes the write grant. assert.doesNotMatch(await readFile(commands, 'utf8'), /"cancel"/); } const response = await f.options.fetchImpl(url, init); @@ -317,25 +379,32 @@ describe('trusted dispatcher transfer integration', () => { return response; }; if (failUpload) await assert.rejects(runCloud(), /Broker chunk upload refused/); - else if (revoked) await assert.rejects(runCloud(), /Broker transfer cleanup incomplete/); + else if (hangSubmission) + await assert.rejects(runCloud(), /submission command timed out and its prepared run was cancelled/); else if (scenario !== 'completed') await assert.rejects(runCloud(), /Cloud RelayFlow finished with status/); else await runCloud(); + // A revoked grant is reported, but never hides the dispatcher's own failure. + assert.equal( + warnings.some((line) => /Broker transfer cleanup incomplete/.test(line)), + revoked + ); const invocations = (await readFile(commands, 'utf8')) .trim() .split('\n') .map((line) => JSON.parse(line)); assert.equal(invocations.filter((args) => args[1] === 'run').length, 1); - assert.equal(invocations.filter((args) => args[1] === 'cancel').length, failUpload ? 1 : 0); + assert.equal(invocations.filter((args) => args[1] === 'cancel').length, cancels ? 1 : 0); assert.ok(cleanupAttempts > 0 || scenario === 'completed'); if (revoked) assert.ok([...f.objects.values()].some((body) => body !== '')); else assert.ok([...f.objects.values()].every((body) => body === '')); - if (!failUpload) + if (!cancels) assert.doesNotMatch( await readFile(path.join(directory, 'cloud.log'), 'utf8'), /ci-fixture|sandbox-fixture/ ); } finally { + console.warn = originalWarn; globalThis.fetch = originalFetch; process.chdir(originalDirectory); for (const key of Object.keys(process.env)) diff --git a/scripts/pr-proof/run-cloud.mjs b/scripts/pr-proof/run-cloud.mjs index 4775ac26df..e951ce8808 100644 --- a/scripts/pr-proof/run-cloud.mjs +++ b/scripts/pr-proof/run-cloud.mjs @@ -606,6 +606,7 @@ export async function main() { let shuttingDown = false; let activeCommandController = null; let launchProgressError = null; + let dispatchFailure = null; let lastStatusOutput = ''; let statusPollFailures = 0; @@ -675,12 +676,26 @@ export async function main() { await cancelPromise; }; + // Tombstone this nonce's objects while the prepared-run write grant is live. + // Cancellation may revoke that grant, so every path that cancels the remote + // run must clean up first. The dispatcher's original failure is the useful + // diagnostic; a cleanup failure is reported and returned, never thrown here. + const cleanupTransfer = async () => { + try { + await brokerTransfer?.cleanup(); + return null; + } catch (error) { + console.warn(sanitizeCloudCommandOutput(error.message, auth.diagnosticSecretValues)); + return error; + } + }; + const signalHandler = (signal) => { if (shuttingDown) return; shuttingDown = true; activeCommandController?.abort(); void (async () => { - await brokerTransfer?.cleanup().catch(() => console.warn('Broker transfer cleanup incomplete')); + await cleanupTransfer(); await cancelRemote(signal).catch((error) => console.warn(sanitizeCloudCommandOutput(error.message, auth.diagnosticSecretValues)) ); @@ -706,6 +721,7 @@ export async function main() { if (launchProgressError) throw launchProgressError; if (launch.aborted) throw new Error('Cloud workflow submission was interrupted'); if (launch.timedOut) { + await cleanupTransfer(); await cancelRemote('submission command timed out'); throw new Error( 'Cloud workflow submission command timed out and its prepared run was cancelled; it is not retried' @@ -782,6 +798,7 @@ export async function main() { statusPollFailures, diagnosticSecretValues: auth.diagnosticSecretValues, }); + await cleanupTransfer(); await cancelRemote('deadline exceeded'); terminal = true; throw new Error(`Cloud RelayFlow exceeded ${timeoutMs}ms`); @@ -827,14 +844,20 @@ export async function main() { )}\`\n- Cloud status: **${terminalStatus}**\n` ); } + } catch (error) { + dispatchFailure = error; + throw error; } finally { activeCommandController?.abort(); process.removeListener('SIGINT', signalHandler); process.removeListener('SIGTERM', signalHandler); try { - if (proofSucceeded) + if (proofSucceeded) { brokerTransfer?.release(); // each completed arm consumed its transfer - else await brokerTransfer?.cleanup(); + } else { + const cleanupError = await cleanupTransfer(); + if (cleanupError && !dispatchFailure) throw cleanupError; + } } finally { if (runId && !terminal) await cancelRemote('dispatcher exiting').catch((error) =>