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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 2 additions & 3 deletions .github/workflows/relayflow-pr-proof.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
3 changes: 3 additions & 0 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
48 changes: 48 additions & 0 deletions docs/pr-proof-broker-transfer.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
# 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 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
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, 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
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.
318 changes: 318 additions & 0 deletions scripts/pr-proof/broker-transfer.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,318 @@
/** 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';
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 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}`,
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),
]),
Comment thread
cursor[bot] marked this conversation as resolved.
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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 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 () => {
// 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++) {
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);
let result;
let failure;
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');
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;
}
}
// 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;
}
Loading
Loading