From 718127b58684d5228ad58c513b796d608f78a049 Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 17 Sep 2026 11:11:08 -0700 Subject: [PATCH 1/4] feat(sdk): flows run --cloud --sync-code and flows sync, Cloud-API transport only MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit v1's `agent-relay cloud run --sync-code` uploaded the working tree so the hosted run executed inside it, and `cloud sync ` pulled the sandbox's diff back. Cloud already honours both for v2 runs — the v2 executor's cwd is the code mount and patch generation runs after either engine — but `flows run --cloud` never sent code, so a hosted authored flow ran in an empty directory unless a deployment supplied a repository grant. `--sync-code` now performs the same prepare → upload → submit sequence, through the Cloud API alone: `prepare` must answer with a `cloud-api` workflow-storage backend (R2), the gzip'd ustar goes to `/workflows/runs//storage/` with the run-scoped credential the receipt carries, and the run is submitted against that prepared ID. Any other backend is refused as `unsupported_storage_backend` before a byte is uploaded; the SDK carries no AWS client. Sync precedes submission, so a refused backend or failed upload never leaves a launched run pointing at a tree Cloud does not hold. Inside a Git checkout the archive is `git ls-files --cached --others --exclude-standard`; outside Git, everything but `.git`/`node_modules`. The tar writer is in-process (no `tar` dependency), deterministic, and refuses paths the ustar prefix split cannot carry rather than truncating. `flows sync ` fetches `/workflows/runs//patch` and applies it with `git apply` behind a `--check`, so a conflict leaves the tree untouched. Multi-path runs are refused as `sync_unsupported`. `flows run --cloud` also accepts `--input` for an authored `.flow.ts`, parsed exactly as a local direct run parses it; without it the CLI could not submit an authored flow at all, since Cloud requires the input. Co-Authored-By: Claude Opus 5 (1M context) --- docs/CLOUD.md | 52 ++++- packages/sdk/src/cli.ts | 58 ++++-- packages/sdk/src/cli/cloud-run.ts | 29 ++- packages/sdk/src/cli/cloud-sync.ts | 34 ++++ packages/sdk/src/cloud-http.ts | 40 +++- packages/sdk/src/cloud-run.ts | 33 +++- packages/sdk/src/cloud-sync.ts | 210 ++++++++++++++++++++ packages/sdk/src/index.ts | 4 + packages/sdk/tests/cloud-sync.test.ts | 274 ++++++++++++++++++++++++++ 9 files changed, 710 insertions(+), 24 deletions(-) create mode 100644 packages/sdk/src/cli/cloud-sync.ts create mode 100644 packages/sdk/src/cloud-sync.ts create mode 100644 packages/sdk/tests/cloud-sync.test.ts diff --git a/docs/CLOUD.md b/docs/CLOUD.md index a1de1de7d..293a7945a 100644 --- a/docs/CLOUD.md +++ b/docs/CLOUD.md @@ -36,8 +36,52 @@ pinned runtime must be enabled by that deployment's operator. ```sh flows run --cloud examples/cloud-gates/cloud-gates.flow.yaml flows run --cloud --wait --json examples/cloud-gates/cloud-gates.flow.yaml +flows run --cloud --input '{}' review.flow.ts ``` +An authored `.flow.ts` takes `--input` exactly as a local direct run does (an +existing JSON file, otherwise inline JSON), and travels as one self-contained +source with its pinned Surface authority. + +## Code sync + +```sh +cd +flows run --cloud --sync-code --wait review.flow.ts --input '{"pr": 7}' +flows sync # apply the run's changes to this checkout +``` + +`--sync-code` uploads the invoking directory before submission, so the hosted +run — every `f.run` and every `f.agent` — executes inside that tree, the way +v1's `agent-relay cloud run --sync-code` did. Inside a Git checkout the upload +is `git ls-files --cached --others --exclude-standard`: `.gitignore` governs, +untracked files ride along, `.git` and `node_modules` never do. Outside Git, +every file except those two directories. The limit is 256 MiB uncompressed. +The flow source itself is still sent in the request body, so it must stay +self-contained; sibling imports inside the tree are not resolved by the hosted +runner. + +The transport is the Cloud API only. `POST /api/v1/workflows/prepare` must +answer with a `cloud-api` workflow-storage backend (Cloud's R2), the archive +is `PUT` to `/api/v1/workflows/runs//storage/` with the run-scoped +credential that receipt carries, and the run is submitted against that +prepared run ID. A `prepare` that offers any other backend is refused as +`unsupported_storage_backend` before a byte is uploaded; this SDK carries no +AWS client and never uploads to a bucket directly. Sync happens before +submission, so a refused backend or failed upload never leaves a launched run +pointing at a tree Cloud does not hold. + +`flows sync ` fetches the sandbox's post-run diff from +`/api/v1/workflows/runs//patch` and applies it with `git apply` after a +`--check` pass, so a conflicting patch leaves the tree untouched +(`patch_conflict`, exit 2). Runs that declared several mounted paths carry one +patch per path and are refused here (`sync_unsupported`). `--dir ` +targets a checkout other than the current directory. + +A synced run and a Cloud repository grant are mutually exclusive on the +server: `--sync-code` is the local-driven development loop, and +webhook-triggered deployments keep cloning through the grant. + Without `--wait`, exit 0 means the server accepted the run. With `--wait`, it means Cloud reported `completed` with a validated `success` completion reason. Failed/cancelled runs and observation/transport failures return 1. Local input, @@ -72,10 +116,10 @@ by this endpoint and is not synthesized by the SDK. ## Current limits and scope -- Cloud's v2 bootstrap explicitly rejects authored `.flow.ts` files. The SDK - refuses those before HTTP rather than uploading code that cannot run. Inputs, - local imports, CLI configuration files, and workspace files are not bundled - or uploaded by this path. +- An authored `.flow.ts` is submitted as one self-contained source. Local + imports and `use:` dependencies are refused before HTTP. With `--sync-code` + the working tree is uploaded for the run to execute in, but the runner still + loads the flow from the request body, not from the tree. - SDK observation has no fixed execution deadline. The merged authored-agent executor follows worker leases rather than the former 30-second limit. Cloud's separate `relayflow-v2-executor.ts` still has a one-hour execution diff --git a/packages/sdk/src/cli.ts b/packages/sdk/src/cli.ts index e4dc4fe87..bcc8b93a7 100644 --- a/packages/sdk/src/cli.ts +++ b/packages/sdk/src/cli.ts @@ -24,6 +24,7 @@ import { runDirectFlow } from './cli/direct-run.js'; import { parseReplayArgs, replayJournal, type ReplayArgs } from './cli/replay.js'; import { checkTypeScriptFlow } from './cli/check-typescript.js'; import { runCloudCli } from './cli/cloud-run.js'; +import { runCloudSyncCli } from './cli/cloud-sync.js'; import { isAuthoredFlowPath } from './direct-input.js'; import { parseDeployArgs, runDeploy, type DeployArgs } from './cli/deploy.js'; import { parseDigestReference } from './bundle-transport.js'; @@ -51,7 +52,8 @@ type ParsedArgs = | BuildArgs | DeployArgs | { command: 'serve-webhook'; dataDir: string; port: number; admitted?: readonly string[] } - | { command: 'cloud-run'; value: string; json: boolean; wait: boolean } + | { command: 'cloud-run'; value: string; json: boolean; wait: boolean; input: string | undefined; syncCode: boolean } + | { command: 'sync'; runId: string; json: boolean; root: string } | { command: 'check'; json: boolean; watch: boolean; value: string } | { command: 'run'; bucket: string | undefined; reuseFromRunId: string | undefined; localAgent: boolean; dataDir: string; input: string | undefined; json: boolean; spawn: boolean; noObserverLink: boolean; allowHumanInfluenced: boolean; value: string } | { command: 'resume'; localAgent: boolean; dataDir: string; json: boolean; spawn: boolean; noObserverLink: boolean; allowHumanInfluenced: boolean; value: string } @@ -71,7 +73,9 @@ const USAGE = [ 'flows check [--watch] [--json] ', 'flows serve-webhook --data-dir --port

[--allow [,]]', 'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir

] [--local-agent] [--reuse-from ] ', - 'flows run --cloud [--json] [--wait] ', + 'flows run --cloud [--json] [--wait] [--sync-code] ', + 'flows run --cloud [--json] [--wait] [--sync-code] --input ', + 'flows sync [--json] [--dir ] ', 'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir ] [--local-agent] --input ', 'flows tick start --schedule-id --interval-ms [--epoch-ms ] [--max-catch-up ] [--poll-interval-ms ] [--data-dir ] ', 'flows resume [--allow-human-influenced] [--json] [--no-spawn] [--no-observer-link] [--data-dir ] [--local-agent] ', @@ -116,6 +120,7 @@ export async function runCli( if (parsed.command === 'serve-webhook') return runServeWebhook(parsed, io); if (parsed.command === 'cloud-run') return runCloudCli(parsed, io); + if (parsed.command === 'sync') return runCloudSyncCli(parsed, io); if (parsed.command === 'replay') return replayJournal(parsed, io); if (parsed.command === 'build') return runBuild(parsed, io); if (parsed.command === 'deploy') return runDeploy(parsed, io); @@ -424,12 +429,14 @@ function parseArgs(args: readonly string[]): ParsedArgs | undefined { if (command === 'hn-monitor') return parseHnMonitorArgs(args.slice(1)); if (command === 'tick') return parseTickArgs(args.slice(1)); if (command === 'observer') return parseObserverArgs(args.slice(1)); + if (command === 'sync') return parseSyncArgs(args.slice(1)); if (command !== 'check' && command !== 'run' && command !== 'resume') return undefined; let json = false; let watch = false; let cloud = false; let wait = false; + let syncCode = false; let localAgent = false; let allowHumanInfluenced = false; let dataDir = DEFAULT_DATA_DIR; @@ -443,10 +450,11 @@ function parseArgs(args: readonly string[]): ParsedArgs | undefined { const positionals: string[] = []; for (let index = 1; index < args.length; index += 1) { const argument = args[index]!; - if (argument === '--cloud' || argument === '--wait') { - if (command !== 'run' || (argument === '--cloud' ? cloud : wait)) return undefined; + if (argument === '--cloud' || argument === '--wait' || argument === '--sync-code') { + if (command !== 'run' || (argument === '--cloud' ? cloud : argument === '--wait' ? wait : syncCode)) return undefined; if (argument === '--cloud') cloud = true; - else wait = true; + else if (argument === '--wait') wait = true; + else syncCode = true; continue; } if (argument === '--allow-human-influenced') { @@ -520,13 +528,15 @@ function parseArgs(args: readonly string[]): ParsedArgs | undefined { if (bucket !== undefined && (cloud || !parseDigestReference(positionals[0]!))) return undefined; if (cloud) { // `--cloud` submits the spec to Cloud, so every flag that only describes a - // local run -- an inline input, a data dir, a suppressed daemon, a local - // agent, a local observer-link opt-out -- describes nothing there and is - // refused rather than ignored. - if (allowHumanInfluenced || sawInput || sawDataDir || !spawn || localAgent || noObserverLink || reuseFromRunId !== undefined) return undefined; - return { command: 'cloud-run', value: positionals[0]!, json, wait }; + // local run -- a data dir, a suppressed daemon, a local agent, a local + // observer-link opt-out -- describes nothing there and is refused rather + // than ignored. `--input` is the authored body's argument and travels with + // the source, so it is accepted exactly where a local run accepts it. + if (allowHumanInfluenced || sawDataDir || !spawn || localAgent || noObserverLink || reuseFromRunId !== undefined) return undefined; + if (sawInput && !isAuthoredFlowPath(positionals[0]!)) return undefined; + return { command: 'cloud-run', value: positionals[0]!, json, wait, input, syncCode }; } - if (wait) return undefined; + if (wait || syncCode) return undefined; if (reuseFromRunId !== undefined && isAuthoredFlowPath(positionals[0]!)) return undefined; if (command === 'run' && input !== undefined && !isAuthoredFlowPath(positionals[0]!)) return undefined; @@ -579,6 +589,32 @@ function parseHnMonitorArgs(rest: readonly string[]): ParsedArgs | undefined { * directory at all -- the mint is a pure Relaycast API round-trip. No * positional argument, no other flags. */ +/** `flows sync [--json] [--dir ] `: apply a hosted run's patch to a local tree. */ +function parseSyncArgs(args: readonly string[]): ParsedArgs | undefined { + let json = false; + let root: string | undefined; + const positionals: string[] = []; + for (let index = 0; index < args.length; index += 1) { + const argument = args[index]!; + if (argument === '--json') { + if (json) return undefined; + json = true; + continue; + } + if (argument === '--dir') { + const value = args[index + 1]; + if (root !== undefined || value === undefined || value.startsWith('-')) return undefined; + root = value; + index += 1; + continue; + } + if (argument.startsWith('-')) return undefined; + positionals.push(argument); + } + if (positionals.length !== 1) return undefined; + return { command: 'sync', runId: positionals[0]!, json, root: root ?? '.' }; +} + function parseObserverArgs(rest: readonly string[]): ParsedArgs | undefined { let dataDir = DEFAULT_DATA_DIR; let sawDataDir = false; diff --git a/packages/sdk/src/cli/cloud-run.ts b/packages/sdk/src/cli/cloud-run.ts index 49212b337..4cbeb0ba1 100644 --- a/packages/sdk/src/cli/cloud-run.ts +++ b/packages/sdk/src/cli/cloud-run.ts @@ -1,10 +1,14 @@ import { CloudFlowError } from '../cloud-http.js'; -import { runInCloud, waitForCloudFlowRun } from '../cloud-run.js'; +import { runInCloud, waitForCloudFlowRun, type RunInCloudOptions } from '../cloud-run.js'; +import { DirectInputError, isAuthoredFlowPath, parseDirectInput } from '../direct-input.js'; +import { snapshotJsonValue } from '../json-value.js'; import type { CliIo } from '../cli.js'; /** Presentation only: the central CLI parser owns argv; the SDK owns the lifecycle. */ export async function runCloudCli( - { value: path, json, wait }: { value: string; json: boolean; wait: boolean }, + { value: path, json, wait, input, syncCode }: { + value: string; json: boolean; wait: boolean; input: string | undefined; syncCode: boolean; + }, io: CliIo, ): Promise<0 | 1 | 2> { const controller = new AbortController(); @@ -13,11 +17,28 @@ export async function runCloudCli( process.once('SIGTERM', abort); let runId: string | undefined; try { - const receipt = await runInCloud({ path }, { signal: controller.signal }); + const options: RunInCloudOptions = { signal: controller.signal }; + if (isAuthoredFlowPath(path)) { + // Same parse as a local direct run, so a file-or-inline argument means + // the same thing on both sides of `--cloud`. + try { + options.input = snapshotJsonValue(parseDirectInput(input), 'Cloud authored input'); + } catch (error) { + if (error instanceof DirectInputError) throw new CloudFlowError('invalid_input', error.message); + throw error; + } + } + // The tree is the invoking directory, as with v1: the flow path is where + // the body lives, not the boundary of what the run may read. + if (syncCode) options.syncCode = { root: process.cwd() }; + const receipt = await runInCloud({ path }, options); runId = receipt.runId; if (!json) { io.stdout(`ACCEPTED ${receipt.runId} (${receipt.status})`); io.stdout(receipt.apiUrl); + if (receipt.synced) { + io.stdout(`SYNCED ${receipt.synced.files} files (${receipt.synced.bytes} bytes); pull changes with: flows sync ${receipt.runId}`); + } } if (!wait) { if (json) io.stdout(JSON.stringify({ ok: true, ...receipt })); @@ -38,7 +59,7 @@ export async function runCloudCli( if (json) io.stdout(JSON.stringify({ ok: false, code, message, ...(runId ? { runId } : {}) })); else io.stderr(`${code}: ${message}${runId ? ` (run ${runId})` : ''}`); return error instanceof CloudFlowError - && (['configuration', 'unsupported_source', 'invalid_input'].includes(error.code) + && (['configuration', 'unsupported_source', 'invalid_input', 'unsupported_storage_backend', 'sync_too_large', 'sync_unsupported'].includes(error.code) || (runId === undefined && error.code === 'http_error' && [401, 403].includes(error.status ?? 0))) ? 2 : 1; } finally { process.off('SIGINT', abort); diff --git a/packages/sdk/src/cli/cloud-sync.ts b/packages/sdk/src/cli/cloud-sync.ts new file mode 100644 index 000000000..267a79fe0 --- /dev/null +++ b/packages/sdk/src/cli/cloud-sync.ts @@ -0,0 +1,34 @@ +import { CloudFlowError } from '../cloud-http.js'; +import { applyCloudPatch, downloadCloudPatch } from '../cloud-sync.js'; +import type { CliIo } from '../cli.js'; + +/** + * `flows sync `: fetch the diff a hosted run left in its synced tree + * and apply it here. The patch is the sandbox's own `git diff` against the + * uploaded baseline, so applying it reproduces exactly what the run's steps + * wrote — no re-execution, no re-upload. + */ +export async function runCloudSyncCli( + { runId, json, root }: { runId: string; json: boolean; root: string }, + io: CliIo, +): Promise<0 | 1 | 2> { + try { + const { patch, hasChanges } = await downloadCloudPatch(runId, {}); + if (!hasChanges || !patch.trim()) { + if (json) io.stdout(JSON.stringify({ ok: true, runId, hasChanges: false, applied: false })); + else io.stdout(`NO CHANGES ${runId}`); + return 0; + } + applyCloudPatch(root, patch); + const files = [...patch.matchAll(/^\+\+\+ b\/(.+)$/gmu)].map(match => match[1]!); + if (json) io.stdout(JSON.stringify({ ok: true, runId, hasChanges: true, applied: true, files })); + else io.stdout(`APPLIED ${runId}: ${files.length} file${files.length === 1 ? '' : 's'}${files.length ? `\n ${files.join('\n ')}` : ''}`); + return 0; + } catch (error) { + const code = error instanceof CloudFlowError ? error.code : 'cloud_sync_failed'; + const message = error instanceof Error ? error.message : 'Cloud sync failed.'; + if (json) io.stdout(JSON.stringify({ ok: false, code, message, runId })); + else io.stderr(`${code}: ${message} (run ${runId})`); + return error instanceof CloudFlowError && ['configuration', 'sync_unsupported', 'patch_conflict'].includes(error.code) ? 2 : 1; + } +} diff --git a/packages/sdk/src/cloud-http.ts b/packages/sdk/src/cloud-http.ts index cb19ba0b2..9bc25523c 100644 --- a/packages/sdk/src/cloud-http.ts +++ b/packages/sdk/src/cloud-http.ts @@ -10,7 +10,9 @@ export interface CloudConnectionOptions { export class CloudFlowError extends Error { constructor( - readonly code: 'configuration' | 'unsupported_source' | 'invalid_input' | 'invalid_response' | 'http_error' | 'transport_error' | 'transient_error', + readonly code: 'configuration' | 'unsupported_source' | 'invalid_input' | 'invalid_response' | 'http_error' + | 'transport_error' | 'transient_error' | 'unsupported_storage_backend' | 'sync_too_large' | 'sync_unsupported' + | 'patch_conflict', message: string, readonly status?: number, ) { @@ -43,19 +45,49 @@ export async function cloudRequest( path: string, options: CloudConnectionOptions, body?: unknown, +): Promise { + return cloudFetch(path, options, body === undefined + ? { method: 'GET' } + : { method: 'POST', body: JSON.stringify(body), contentType: 'application/json' }); +} + +export interface CloudFetchInit { + method: 'GET' | 'POST' | 'PUT'; + body?: string | Uint8Array; + contentType?: string; + /** + * A run-scoped token issued by Cloud for one upload (the `prepare` receipt's + * storage credential). Used instead of the configured token for that request + * only; it is never persisted or logged. + */ + bearerToken?: string; +} + +/** One authenticated Cloud request. Every transport error is typed; no retries. */ +export async function cloudFetch( + path: string, + options: CloudConnectionOptions, + init: CloudFetchInit, ): Promise { const { baseUrl, token } = cloudConnection(options); const timeout = options.requestTimeoutMs ?? 30_000; if (!Number.isSafeInteger(timeout) || timeout < 1 || timeout > 2_147_483_647) { throw new CloudFlowError('configuration', 'requestTimeoutMs must be a positive 32-bit integer.'); } + const bearer = init.bearerToken ?? token; + if (/[\r\n]/u.test(bearer) || !bearer.trim()) { + throw new CloudFlowError('invalid_response', 'Cloud issued an unusable storage credential.'); + } const deadline = AbortSignal.timeout(timeout); let response: Response; try { response = await fetch(`${baseUrl}${path}`, { - method: body === undefined ? 'GET' : 'POST', - headers: { authorization: `Bearer ${token}`, 'content-type': 'application/json' }, - ...(body === undefined ? {} : { body: JSON.stringify(body) }), + method: init.method, + headers: { + authorization: `Bearer ${bearer}`, + 'content-type': init.contentType ?? 'application/json', + }, + ...(init.body === undefined ? {} : { body: typeof init.body === 'string' ? init.body : new Blob([init.body]) }), signal: options.signal ? AbortSignal.any([options.signal, deadline]) : deadline, // A redirect must never carry the credential to a different origin. redirect: 'error', diff --git a/packages/sdk/src/cloud-run.ts b/packages/sdk/src/cloud-run.ts index 04dbab24f..279f9b116 100644 --- a/packages/sdk/src/cloud-run.ts +++ b/packages/sdk/src/cloud-run.ts @@ -13,6 +13,7 @@ import { CloudFlowError, cloudConnection, cloudRequest, cloudRunId, isCloudRecord, type CloudConnectionOptions, } from './cloud-http.js'; +import { packWorkingTree, prepareCloudSync, uploadCloudCode } from './cloud-sync.js'; export type CloudFlowSource = FlowSpec | { path: string }; export interface RunInCloudOptions extends CloudConnectionOptions { @@ -20,6 +21,14 @@ export interface RunInCloudOptions extends CloudConnectionOptions { workspaceId?: string; /** Exact JSON input. Required at runtime for authored .flow.ts; refused for declarative source. */ input?: JsonValue; + /** + * Upload this working tree before submission so the hosted run executes + * inside it (v1's `--sync-code`). Cloud-API storage only; see cloud-sync.ts. + * The submitted flow still travels as source in the request body, so it + * must remain self-contained: sibling imports inside the tree are not + * resolved by the hosted runner. + */ + syncCode?: { root: string }; } export interface CloudRunReceipt { runId: string; @@ -28,6 +37,8 @@ export interface CloudRunReceipt { specHash: string; /** Authenticated run API resource; this is not a public sharing URL. */ apiUrl: string; + /** Present when a working tree was synced: what was uploaded, by count and size. */ + synced?: { files: number; bytes: number }; } export interface CloudAuthoredAuthority { readonly schemaVersion: 1; @@ -118,6 +129,18 @@ export async function runInCloud( authority: authored.authority, input: authoredInput, })).digest('hex'); + // Sync before submission: `prepare` reserves the run ID and the upload lands + // under it, so the run request below names code Cloud already holds. A + // refused backend or failed upload therefore never leaves a launched run + // pointing at a tree that is not there. + let synced: { runId: string; codeKey: string; files: number; bytes: number } | undefined; + if (options.syncCode !== undefined) { + const prepared = await prepareCloudSync(options); + options.signal?.throwIfAborted(); + const packed = packWorkingTree(options.syncCode.root); + await uploadCloudCode(prepared, packed.tarball, options); + synced = { runId: prepared.runId, codeKey: prepared.codeKey, files: packed.files.length, bytes: packed.bytes }; + } const result = await cloudRequest('/api/v1/workflows/run', options, { // JSON is a YAML subset. Sending canonical data preserves the exact spec // while using the server's existing YAML-to-config admission path. @@ -127,12 +150,20 @@ export async function runInCloud( ...(authored === undefined ? {} : { authoredAuthority: authored.authority }), ...(authored === undefined ? {} : { inputs: authoredInput }), ...(options.workspaceId === undefined ? {} : { workspaceId: options.workspaceId }), + ...(synced === undefined ? {} : { runId: synced.runId, s3CodeKey: synced.codeKey }), }); if (!isCloudRecord(result) || (result.status !== 'pending' && result.status !== 'running')) { throw new CloudFlowError('invalid_response', 'Cloud did not return an accepted run.'); } const runId = cloudRunId(result.runId); - return { runId, status: result.status, specHash: hash, apiUrl: `${baseUrl}/api/v1/workflows/runs/${runId}` }; + if (synced !== undefined && runId !== synced.runId) { + throw new CloudFlowError('invalid_response', + `Cloud accepted run ${runId} but the synced code was uploaded for ${synced.runId}.`); + } + return { + runId, status: result.status, specHash: hash, apiUrl: `${baseUrl}/api/v1/workflows/runs/${runId}`, + ...(synced === undefined ? {} : { synced: { files: synced.files, bytes: synced.bytes } }), + }; } export async function getCloudFlowRun( diff --git a/packages/sdk/src/cloud-sync.ts b/packages/sdk/src/cloud-sync.ts new file mode 100644 index 000000000..ba596adcc --- /dev/null +++ b/packages/sdk/src/cloud-sync.ts @@ -0,0 +1,210 @@ +import { spawnSync } from 'node:child_process'; +import { lstatSync, readFileSync, readdirSync, readlinkSync } from 'node:fs'; +import { join, relative, resolve, sep } from 'node:path'; +import { gzipSync } from 'node:zlib'; +import { + CloudFlowError, cloudFetch, cloudRequest, cloudRunId, isCloudRecord, type CloudConnectionOptions, +} from './cloud-http.js'; + +/** + * Code sync for hosted v2 runs, Cloud-API transport only. + * + * v1's `--sync-code` had two transports: an S3 client with STS credentials and + * the Cloud API's own storage route. v2 runs on Cloudflare, so this module + * speaks only the second: `prepare` must answer with a `cloud-api` storage + * backend, the tarball is `PUT` to `/workflows/runs//storage/`, and + * the sandbox's post-run patch is read from `/workflows/runs//patch`. A + * `prepare` that hands back anything else is refused before any upload — the + * SDK never carries an AWS dependency and never sends code to a bucket it did + * not get through the Cloud API. + */ + +/** Uncompressed bytes. Matches the sandbox's extraction budget with headroom. */ +export const MAX_SYNC_BYTES = 256 * 1024 * 1024; +const ALWAYS_SKIPPED = new Set(['.git', 'node_modules']); +const CODE_KEY = /^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$/u; + +export interface PreparedCloudSync { + runId: string; + codeKey: string; + /** Run-scoped upload credential from `prepare`; falls back to the session token. */ + uploadToken?: string; +} + +/** `POST /api/v1/workflows/prepare`, refusing every non-Cloud-API storage backend. */ +export async function prepareCloudSync(options: CloudConnectionOptions): Promise { + const receipt = await cloudFetch('/api/v1/workflows/prepare', options, { method: 'POST' }); + if (!isCloudRecord(receipt)) throw new CloudFlowError('invalid_response', 'Cloud prepare returned a non-object.'); + const storage = isCloudRecord(receipt.workflowStorage) ? receipt.workflowStorage : undefined; + const credentials = isCloudRecord(receipt.s3Credentials) ? receipt.s3Credentials : undefined; + const backend = storage?.backend ?? credentials?.backend; + if (backend !== 'cloud-api') { + throw new CloudFlowError('unsupported_storage_backend', + 'Cloud offered a workflow storage backend other than the Cloud API (R2). ' + + 'flows v2 syncs code only through the Cloud API; it never uploads to AWS directly.'); + } + const codeKey = typeof receipt.s3CodeKey === 'string' ? receipt.s3CodeKey : ''; + if (!CODE_KEY.test(codeKey)) throw new CloudFlowError('invalid_response', 'Cloud prepare returned an unusable code key.'); + const uploadToken = credentials?.cloudApiAccessToken; + if (uploadToken !== undefined && typeof uploadToken !== 'string') { + throw new CloudFlowError('invalid_response', 'Cloud prepare returned an unusable storage credential.'); + } + return { + runId: cloudRunId(receipt.runId), codeKey, + ...(uploadToken === undefined ? {} : { uploadToken }), + }; +} + +export interface PackedTree { + tarball: Buffer; + /** Archive members, `root`-relative POSIX paths, sorted. */ + files: string[]; + /** Uncompressed byte total across regular files. */ + bytes: number; +} + +/** + * gzip'd ustar of the working tree at `root`. Inside a Git checkout the member + * list is `git ls-files --cached --others --exclude-standard`, so `.gitignore` + * governs and untracked-but-not-ignored files ride along, as they do in v1. + * Outside Git, every regular file and symlink except `.git`/`node_modules`. + * Ordering and mtimes are deterministic, so the same tree packs to the same + * bytes; what Cloud extracts is exactly this list, nothing implicit. + */ +export function packWorkingTree(root: string): PackedTree { + const absoluteRoot = resolve(root); + const files = listTreeFiles(absoluteRoot); + const chunks: Buffer[] = []; + let bytes = 0; + for (const path of files) { + const absolute = join(absoluteRoot, ...path.split('/')); + const stat = lstatSync(absolute); + if (stat.isSymbolicLink()) { + chunks.push(ustarHeader(path, 0, '2', readlinkSync(absolute))); + continue; + } + if (!stat.isFile()) continue; + bytes += stat.size; + if (bytes > MAX_SYNC_BYTES) { + throw new CloudFlowError('sync_too_large', + `The working tree exceeds the ${MAX_SYNC_BYTES}-byte code sync limit; narrow it with .gitignore.`); + } + const content = readFileSync(absolute); + chunks.push(ustarHeader(path, content.length, '0'), content); + const padding = (512 - (content.length % 512)) % 512; + if (padding) chunks.push(Buffer.alloc(padding)); + } + chunks.push(Buffer.alloc(1024)); + return { tarball: gzipSync(Buffer.concat(chunks), { level: 6 }), files, bytes }; +} + +function listTreeFiles(root: string): string[] { + const git = spawnSync('git', ['-C', root, 'ls-files', '-z', '--cached', '--others', '--exclude-standard'], + { encoding: 'utf8', maxBuffer: 64 * 1024 * 1024 }); + let candidates: string[]; + if (git.status === 0) { + candidates = git.stdout.split('\0').filter(Boolean); + } else { + candidates = []; + walk(root, root, candidates); + } + const files = new Set(); + for (const candidate of candidates) { + const path = candidate.split(sep).join('/'); + const segments = path.split('/'); + if (segments.some(segment => ALWAYS_SKIPPED.has(segment) || segment === '..' || segment === '')) continue; + let stat; + try { stat = lstatSync(join(root, ...segments)); } catch { continue; } + if (stat.isFile() || stat.isSymbolicLink()) files.add(path); + } + return [...files].sort(); +} + +function walk(root: string, current: string, out: string[]): void { + for (const entry of readdirSync(current, { withFileTypes: true })) { + if (ALWAYS_SKIPPED.has(entry.name)) continue; + const path = join(current, entry.name); + if (entry.isDirectory()) walk(root, path, out); + else if (entry.isFile() || entry.isSymbolicLink()) out.push(relative(root, path)); + } +} + +/** POSIX ustar header. Names beyond the 100+155 split are refused, not truncated. */ +function ustarHeader(path: string, size: number, type: '0' | '2', link = ''): Buffer { + let name = path; + let prefix = ''; + if (Buffer.byteLength(name) > 100) { + const cut = name.lastIndexOf('/', 154); + if (cut <= 0 || Buffer.byteLength(name.slice(cut + 1)) > 100 || Buffer.byteLength(name.slice(0, cut)) > 155) { + throw new CloudFlowError('sync_unsupported', `Path "${path}" is too long for the code sync archive.`); + } + prefix = name.slice(0, cut); + name = name.slice(cut + 1); + } + if (Buffer.byteLength(link) > 100) { + throw new CloudFlowError('sync_unsupported', `Symlink target of "${path}" is too long for the code sync archive.`); + } + const header = Buffer.alloc(512); + header.write(name, 0, 100); + header.write(type === '2' ? '0000777' : '0000644', 100, 8); + header.write('0000000', 108, 8); + header.write('0000000', 116, 8); + header.write(size.toString(8).padStart(11, '0'), 124, 12); + header.write('00000000000', 136, 12); + header.write(' ', 148, 8); + header.write(type, 156, 1); + header.write(link, 157, 100); + header.write('ustar\0', 257, 6); + header.write('00', 263, 2); + header.write(prefix, 345, 155); + let checksum = 0; + for (const byte of header) checksum += byte; + header.write(checksum.toString(8).padStart(6, '0') + '\0 ', 148, 8); + return header; +} + +/** `PUT` the archive through the Cloud API; the run-scoped credential wins when present. */ +export async function uploadCloudCode( + prepared: PreparedCloudSync, tarball: Buffer, options: CloudConnectionOptions, +): Promise { + await cloudFetch( + `/api/v1/workflows/runs/${encodeURIComponent(prepared.runId)}/storage/${encodeURIComponent(prepared.codeKey)}`, + { ...options, requestTimeoutMs: options.requestTimeoutMs ?? 300_000 }, + { method: 'PUT', body: tarball, contentType: 'application/gzip', + ...(prepared.uploadToken === undefined ? {} : { bearerToken: prepared.uploadToken }) }, + ); +} + +export interface CloudPatch { + patch: string; + hasChanges: boolean; +} + +/** The sandbox's post-run diff. Multi-path runs carry several patches and are refused here. */ +export async function downloadCloudPatch(runId: string, options: CloudConnectionOptions): Promise { + const payload = await cloudRequest(`/api/v1/workflows/runs/${encodeURIComponent(cloudRunId(runId))}/patch`, options); + if (!isCloudRecord(payload)) throw new CloudFlowError('invalid_response', 'Cloud patch response was not an object.'); + if (isCloudRecord(payload.patches)) { + const names = Object.keys(payload.patches); + throw new CloudFlowError('sync_unsupported', + `Run ${runId} produced ${names.length} path-scoped patches (${names.join(', ')}); flows sync applies single-tree runs only.`); + } + if (typeof payload.patch !== 'string' || typeof payload.hasChanges !== 'boolean') { + throw new CloudFlowError('invalid_response', 'Cloud patch response is missing patch or hasChanges.'); + } + return { patch: payload.patch, hasChanges: payload.hasChanges }; +} + +/** `git apply --check` then `git apply`; a conflict leaves the tree untouched. */ +export function applyCloudPatch(root: string, patch: string): void { + const args = ['-C', resolve(root), 'apply', '--whitespace=nowarn']; + const check = spawnSync('git', [...args, '--check'], { input: patch, encoding: 'utf8' }); + if (check.status !== 0) { + throw new CloudFlowError('patch_conflict', + `The patch does not apply cleanly to ${resolve(root)}:\n${check.stderr.trim()}`); + } + const apply = spawnSync('git', args, { input: patch, encoding: 'utf8' }); + if (apply.status !== 0) { + throw new CloudFlowError('patch_conflict', `git apply failed:\n${apply.stderr.trim()}`); + } +} diff --git a/packages/sdk/src/index.ts b/packages/sdk/src/index.ts index b3ead2c60..06c96829b 100644 --- a/packages/sdk/src/index.ts +++ b/packages/sdk/src/index.ts @@ -63,6 +63,10 @@ export { runInCloud, getCloudFlowRun, waitForCloudFlowRun, type CloudFlowSource, type RunInCloudOptions, type CloudRunReceipt, type CloudRunState, } from './cloud-run.js'; +export { + downloadCloudPatch, applyCloudPatch, packWorkingTree, MAX_SYNC_BYTES, + type CloudPatch, type PackedTree, +} from './cloud-sync.js'; export { canonicalize, specHash } from './canonical.js'; export { diff --git a/packages/sdk/tests/cloud-sync.test.ts b/packages/sdk/tests/cloud-sync.test.ts new file mode 100644 index 000000000..1a39c824d --- /dev/null +++ b/packages/sdk/tests/cloud-sync.test.ts @@ -0,0 +1,274 @@ +import { execFileSync } from 'node:child_process'; +import { mkdir, mkdtemp, readFile, readdir, rm, symlink, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { gunzipSync } from 'node:zlib'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { runCli } from '../src/cli.js'; +import { runInCloud } from '../src/cloud-run.js'; +import { applyCloudPatch, packWorkingTree, prepareCloudSync } from '../src/cloud-sync.js'; + +const dirs: string[] = []; +afterEach(async () => { + vi.unstubAllEnvs(); + vi.restoreAllMocks(); + for (const dir of dirs.splice(0)) await rm(dir, { recursive: true, force: true }); +}); + +async function tempDir(prefix: string): Promise { + const dir = await mkdtemp(join(tmpdir(), prefix)); + dirs.push(dir); + return dir; +} + +function git(cwd: string, ...args: string[]): string { + return execFileSync('git', ['-C', cwd, ...args], { + encoding: 'utf8', + env: { ...process.env, GIT_AUTHOR_NAME: 't', GIT_AUTHOR_EMAIL: 't@example.com', + GIT_COMMITTER_NAME: 't', GIT_COMMITTER_EMAIL: 't@example.com', GIT_CONFIG_GLOBAL: '/dev/null' }, + }); +} + +interface Call { method: string; path: string; auth: string | undefined; contentType: string | null; body: unknown } + +/** Records every request; `routes` answers by path prefix. */ +function cloud(routes: Record unknown>) { + const calls: Call[] = []; + vi.spyOn(globalThis, 'fetch').mockImplementation(async (input, init) => { + init?.signal?.throwIfAborted(); + const headers = new Headers(init?.headers); + const path = new URL(String(input)).pathname; + const raw = init?.body; + const body = raw === undefined ? undefined + : typeof raw === 'string' ? JSON.parse(raw) + : Buffer.from(await new Response(raw as BodyInit).arrayBuffer()); + const call: Call = { method: init?.method ?? 'GET', path, auth: headers.get('authorization') ?? undefined, + contentType: headers.get('content-type'), body }; + calls.push(call); + const route = Object.keys(routes).find(prefix => path.startsWith(prefix)); + if (!route) return new Response('{"error":"not found"}', { status: 404 }); + return new Response(JSON.stringify(routes[route]!(call)), { status: 200, headers: { 'content-type': 'application/json' } }); + }); + vi.stubEnv('FLOWS_CLOUD_URL', 'https://cloud-contract.example'); + vi.stubEnv('FLOWS_CLOUD_TOKEN', 'test-scoped-cloud-token'); + return { calls, options: { apiUrl: 'https://cloud-contract.example', token: 'test-scoped-cloud-token' } }; +} + +const PREPARED = { + runId: 'prepared-run', s3CodeKey: 'code.tar.gz', + workflowStorage: { backend: 'cloud-api' }, + s3Credentials: { backend: 'cloud-api', cloudApiAccessToken: 'run-scoped-upload-token' }, +}; + +async function untar(tarball: Buffer): Promise<{ dir: string; members: string[] }> { + const dir = await tempDir('cloud-sync-extract-'); + await writeFile(join(dir, 'code.tar.gz'), tarball); + const members = execFileSync('tar', ['-tzf', join(dir, 'code.tar.gz')], { encoding: 'utf8' }).trim().split('\n').sort(); + await mkdir(join(dir, 'out')); + execFileSync('tar', ['-xzf', join(dir, 'code.tar.gz'), '-C', join(dir, 'out')]); + return { dir: join(dir, 'out'), members }; +} + +describe('packWorkingTree', () => { + it('packs a Git checkout by ls-files semantics: tracked plus untracked, never ignored, .git or node_modules', async () => { + const root = await tempDir('cloud-sync-git-'); + git(root, 'init', '-q'); + await mkdir(join(root, 'src/deep'), { recursive: true }); + await mkdir(join(root, 'node_modules/dep'), { recursive: true }); + await writeFile(join(root, 'src/deep/a.ts'), 'export const a = 1;\n'); + await writeFile(join(root, 'node_modules/dep/index.js'), 'ignored by rule\n'); + await writeFile(join(root, '.gitignore'), 'dist/\n*.log\n'); + await mkdir(join(root, 'dist')); + await writeFile(join(root, 'dist/build.js'), 'ignored\n'); + await writeFile(join(root, 'debug.log'), 'ignored\n'); + await symlink('src/deep/a.ts', join(root, 'link.ts')); + git(root, 'add', '.gitignore', 'src'); + git(root, 'commit', '-q', '-m', 'baseline'); + await writeFile(join(root, 'untracked.md'), 'rides along\n'); + + const packed = packWorkingTree(root); + expect(packed.files).toEqual(['.gitignore', 'link.ts', 'src/deep/a.ts', 'untracked.md']); + expect(packed.bytes).toBe(Buffer.byteLength('dist/\n*.log\n') + Buffer.byteLength('export const a = 1;\n') + Buffer.byteLength('rides along\n')); + + const { dir, members } = await untar(packed.tarball); + expect(members).toEqual(['.gitignore', 'link.ts', 'src/deep/a.ts', 'untracked.md']); + expect(await readFile(join(dir, 'src/deep/a.ts'), 'utf8')).toBe('export const a = 1;\n'); + expect(await readFile(join(dir, 'link.ts'), 'utf8')).toBe('export const a = 1;\n'); + expect(await readdir(dir)).not.toContain('node_modules'); + }); + + it('packs a plain directory by walking it, skipping only .git and node_modules', async () => { + const root = await tempDir('cloud-sync-plain-'); + await mkdir(join(root, 'lib')); + await mkdir(join(root, 'node_modules')); + await mkdir(join(root, '.git')); + await writeFile(join(root, 'lib/x.txt'), 'x'); + await writeFile(join(root, 'node_modules/y.txt'), 'y'); + await writeFile(join(root, '.git/HEAD'), 'ref'); + await writeFile(join(root, 'top.txt'), 'top'); + const packed = packWorkingTree(root); + expect(packed.files).toEqual(['lib/x.txt', 'top.txt']); + expect((await untar(packed.tarball)).members).toEqual(['lib/x.txt', 'top.txt']); + }); + + it('is deterministic for the same tree', async () => { + const root = await tempDir('cloud-sync-det-'); + await writeFile(join(root, 'a'), 'a'); + expect(packWorkingTree(root).tarball.equals(packWorkingTree(root).tarball)).toBe(true); + }); +}); + +describe('prepareCloudSync', () => { + it('refuses a non-Cloud-API storage backend before anything is uploaded', async () => { + const { calls, options } = cloud({ + '/api/v1/workflows/prepare': () => ({ + runId: 'aws-run', s3CodeKey: 'code.tar.gz', + s3Credentials: { backend: 's3', accessKeyId: 'AKIA', secretAccessKey: 'x', bucket: 'b', prefix: 'p' }, + }), + }); + await expect(prepareCloudSync(options)).rejects.toMatchObject({ code: 'unsupported_storage_backend' }); + expect(calls.map(call => call.method)).toEqual(['POST']); + }); + + it('accepts the Cloud-API backend and surfaces the run-scoped upload token', async () => { + const { options } = cloud({ '/api/v1/workflows/prepare': () => PREPARED }); + expect(await prepareCloudSync(options)).toEqual({ + runId: 'prepared-run', codeKey: 'code.tar.gz', uploadToken: 'run-scoped-upload-token', + }); + }); +}); + +describe('runInCloud with syncCode', () => { + it('prepares, uploads through the Cloud API with the run-scoped token, then submits against the prepared run', async () => { + const root = await tempDir('cloud-sync-run-'); + await writeFile(join(root, 'flow.yaml'), JSON.stringify({ + version: '0.1.0', name: 'synced', steps: [{ id: 'ls', type: 'deterministic', command: 'ls' }], + })); + await writeFile(join(root, 'data.txt'), 'payload\n'); + const { calls, options } = cloud({ + '/api/v1/workflows/prepare': () => PREPARED, + '/api/v1/workflows/runs/prepared-run/storage/code.tar.gz': () => ({ ok: true }), + '/api/v1/workflows/run': () => ({ runId: 'prepared-run', status: 'pending' }), + }); + + const receipt = await runInCloud({ path: join(root, 'flow.yaml') }, { ...options, syncCode: { root } }); + + expect(calls.map(call => [call.method, call.path])).toEqual([ + ['POST', '/api/v1/workflows/prepare'], + ['PUT', '/api/v1/workflows/runs/prepared-run/storage/code.tar.gz'], + ['POST', '/api/v1/workflows/run'], + ]); + const upload = calls[1]!; + expect(upload.auth).toBe('Bearer run-scoped-upload-token'); + expect(upload.contentType).toBe('application/gzip'); + const tar = gunzipSync(upload.body as Buffer).toString('latin1'); + expect(tar).toContain('data.txt'); + expect(tar).toContain('flow.yaml'); + expect(calls[2]!.auth).toBe('Bearer test-scoped-cloud-token'); + expect(calls[2]!.body).toMatchObject({ runId: 'prepared-run', s3CodeKey: 'code.tar.gz', relayflowVersion: 'v2' }); + expect(receipt).toMatchObject({ runId: 'prepared-run', synced: { files: 2 } }); + }); + + it('refuses an accepted run whose ID differs from the prepared upload', async () => { + const root = await tempDir('cloud-sync-mismatch-'); + await writeFile(join(root, 'flow.yaml'), JSON.stringify({ + version: '0.1.0', name: 'synced', steps: [{ id: 'ls', type: 'deterministic', command: 'ls' }], + })); + const { options } = cloud({ + '/api/v1/workflows/prepare': () => PREPARED, + '/api/v1/workflows/runs/prepared-run/storage/code.tar.gz': () => ({ ok: true }), + '/api/v1/workflows/run': () => ({ runId: 'someone-else', status: 'pending' }), + }); + await expect(runInCloud({ path: join(root, 'flow.yaml') }, { ...options, syncCode: { root } })) + .rejects.toMatchObject({ code: 'invalid_response' }); + }); + + it('never submits when the backend is refused', async () => { + const root = await tempDir('cloud-sync-refused-'); + await writeFile(join(root, 'flow.yaml'), JSON.stringify({ + version: '0.1.0', name: 'synced', steps: [{ id: 'ls', type: 'deterministic', command: 'ls' }], + })); + const { calls, options } = cloud({ + '/api/v1/workflows/prepare': () => ({ runId: 'r', s3CodeKey: 'code.tar.gz', s3Credentials: { backend: 's3' } }), + }); + await expect(runInCloud({ path: join(root, 'flow.yaml') }, { ...options, syncCode: { root } })) + .rejects.toMatchObject({ code: 'unsupported_storage_backend' }); + expect(calls.map(call => call.path)).toEqual(['/api/v1/workflows/prepare']); + }); +}); + +describe('flows run --cloud --sync-code / flows sync', () => { + it('runs an authored flow with --input and syncs the invoking directory', async () => { + const root = await tempDir('cloud-sync-cli-'); + await symlink(join(process.cwd(), 'node_modules'), join(root, 'node_modules'), 'dir'); + await writeFile(join(root, 'review.flow.ts'), "import { flow } from '@relayflows/surface';\n" + + "export default flow('review', async f => f.done('success'));\n"); + await writeFile(join(root, 'notes.md'), 'context\n'); + const { calls } = cloud({ + '/api/v1/workflows/prepare': () => PREPARED, + '/api/v1/workflows/runs/prepared-run/storage/code.tar.gz': () => ({ ok: true }), + '/api/v1/workflows/run': () => ({ runId: 'prepared-run', status: 'pending' }), + }); + const cwd = vi.spyOn(process, 'cwd').mockReturnValue(root); + const output: string[] = []; + const code = await runCli(['run', '--cloud', '--sync-code', '--json', '--input', '{"pr":7}', join(root, 'review.flow.ts')], + { stdout: line => output.push(line), stderr: line => output.push(line) }); + cwd.mockRestore(); + expect(code, output.join('\n')).toBe(0); + expect(JSON.parse(output[0]!)).toMatchObject({ ok: true, runId: 'prepared-run', synced: { files: 2 } }); + expect(calls[2]!.body).toMatchObject({ fileType: 'ts', inputs: { pr: 7 }, runId: 'prepared-run', s3CodeKey: 'code.tar.gz' }); + const tar = gunzipSync(calls[1]!.body as Buffer).toString('latin1'); + expect(tar).toContain('notes.md'); + expect(tar).not.toContain('node_modules'); + }); + + it.each([ + ['run', '--sync-code', 'flow.yaml'], + ['run', '--cloud', '--sync-code', '--sync-code', 'flow.yaml'], + ['run', '--cloud', '--input', '{}', 'flow.yaml'], + ['sync'], ['sync', 'a', 'b'], ['sync', '--dir'], ['sync', '--json', '--json', 'r'], + ])('refuses argv %j before any request', async (...args) => { + const fetch = vi.spyOn(globalThis, 'fetch'); + expect(await runCli(args, { stdout: () => {}, stderr: () => {} })).toBe(2); + expect(fetch).not.toHaveBeenCalled(); + }); + + it('applies a single-tree patch and reports the files', async () => { + const root = await tempDir('cloud-sync-apply-'); + git(root, 'init', '-q'); + await writeFile(join(root, 'a.txt'), 'one\n'); + git(root, 'add', 'a.txt'); + git(root, 'commit', '-q', '-m', 'baseline'); + const patch = 'diff --git a/a.txt b/a.txt\n--- a/a.txt\n+++ b/a.txt\n@@ -1 +1,2 @@\n one\n+two\n' + + 'diff --git a/b.txt b/b.txt\nnew file mode 100644\n--- /dev/null\n+++ b/b.txt\n@@ -0,0 +1 @@\n+fresh\n'; + cloud({ '/api/v1/workflows/runs/done-run/patch': () => ({ patch, hasChanges: true }) }); + const output: string[] = []; + expect(await runCli(['sync', '--json', '--dir', root, 'done-run'], { stdout: line => output.push(line), stderr: line => output.push(line) })).toBe(0); + expect(JSON.parse(output[0]!)).toEqual({ ok: true, runId: 'done-run', hasChanges: true, applied: true, files: ['a.txt', 'b.txt'] }); + expect(await readFile(join(root, 'a.txt'), 'utf8')).toBe('one\ntwo\n'); + expect(await readFile(join(root, 'b.txt'), 'utf8')).toBe('fresh\n'); + }); + + it('reports no changes without touching the tree', async () => { + const root = await tempDir('cloud-sync-none-'); + cloud({ '/api/v1/workflows/runs/quiet/patch': () => ({ patch: '', hasChanges: false }) }); + const output: string[] = []; + expect(await runCli(['sync', '--dir', root, 'quiet'], { stdout: line => output.push(line), stderr: () => {} })).toBe(0); + expect(output).toEqual(['NO CHANGES quiet']); + }); + + it('refuses multi-path patches and conflicting patches without partial application', async () => { + const root = await tempDir('cloud-sync-conflict-'); + git(root, 'init', '-q'); + await writeFile(join(root, 'a.txt'), 'one\n'); + git(root, 'add', 'a.txt'); + git(root, 'commit', '-q', '-m', 'baseline'); + cloud({ '/api/v1/workflows/runs/multi/patch': () => ({ patches: { api: { patch: 'x', hasChanges: true }, web: { patch: 'y', hasChanges: true } } }) }); + const errors: string[] = []; + expect(await runCli(['sync', '--dir', root, 'multi'], { stdout: () => {}, stderr: line => errors.push(line) })).toBe(2); + expect(errors[0]).toContain('sync_unsupported'); + expect(() => applyCloudPatch(root, 'diff --git a/a.txt b/a.txt\n--- a/a.txt\n+++ b/a.txt\n@@ -1 +1 @@\n-nope\n+two\n')) + .toThrow(expect.objectContaining({ code: 'patch_conflict' })); + expect(await readFile(join(root, 'a.txt'), 'utf8')).toBe('one\n'); + }); +}); From 951b8777ff8d92c4c219083be206dc5cec699a8f Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 17 Sep 2026 11:58:10 -0700 Subject: [PATCH 2/4] =?UTF-8?q?feat(sdk):=20flows=20deploy=20,=20?= =?UTF-8?q?deployments,=20undeploy=20=E2=80=94=20hosted=20listeners=20from?= =?UTF-8?q?=20the=20CLI?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A flow that responds to GitHub issues could only be deployed through the agentrelay.com onboarding handoff into Cloud's deploy wizard. The wizard's one API call, `POST /api/v1/flows/deploy`, already admits a `cli:auth` token, so this gives it a CLI: flows deploy issue-triage.flow.ts --repo owner/name \ --on github:labels=agent --approver [--agents claude,codex] [--draft] flows deployments flows undeploy `deploy` means Cloud. The digest form, `flows deploy @sha256:… --to file://…`, keeps working; the positional decides which form is meant, so no `--cloud` flag is needed. The CLI resolves the workspace from `/api/v1/auth/whoami`, sends the exact source with `mode: activate` (or `draft`), the approver, the agent harnesses Cloud checks credentials for, and the trigger sources; a GitHub source without an explicit `repository` is scoped to `--repo`. Trigger settings are validated client-side against the same per-provider vocabulary the launcher prefilter reads. Credentials: every hosted verb now falls back to the `agent-relay cloud login` store when neither `token` nor `FLOWS_CLOUD_TOKEN` is set, taking its API URL as the default base so a login never sends its token elsewhere; an expired login is refused with the re-login remedy. Deploy routes answer refusals as `{ code, error }`; `cloudFetch` gains `detail: true` to surface exactly those two fields, so `flow_model_not_connected` or `flow_name_taken` is named instead of a bare 409. Proven on production: deployed a label-scoped issues flow against AgentWorkforce/flows (`DEPLOYED … listening`), listed it, removed it with `undeploy`. Co-Authored-By: Claude Opus 5 (1M context) --- docs/CLOUD.md | 54 +++++ packages/sdk/src/cli.ts | 28 ++- packages/sdk/src/cli/cloud-deploy.ts | 157 ++++++++++++++ packages/sdk/src/cloud-deploy.ts | 221 +++++++++++++++++++ packages/sdk/src/cloud-http.ts | 82 ++++++- packages/sdk/src/index.ts | 4 + packages/sdk/tests/cloud-deploy.test.ts | 273 ++++++++++++++++++++++++ 7 files changed, 813 insertions(+), 6 deletions(-) create mode 100644 packages/sdk/src/cli/cloud-deploy.ts create mode 100644 packages/sdk/src/cloud-deploy.ts create mode 100644 packages/sdk/tests/cloud-deploy.test.ts diff --git a/docs/CLOUD.md b/docs/CLOUD.md index 293a7945a..e21ebff53 100644 --- a/docs/CLOUD.md +++ b/docs/CLOUD.md @@ -82,6 +82,60 @@ A synced run and a Cloud repository grant are mutually exclusive on the server: `--sync-code` is the local-driven development loop, and webhook-triggered deployments keep cloning through the grant. +## Credentials + +Every hosted verb resolves its credential the same way: the `token` option, +then `FLOWS_CLOUD_TOKEN`, then the `agent-relay cloud login` store +(`~/.agentworkforce/relay/cloud-auth.json`, or `AGENT_RELAY_HOME`). The login +store also supplies the base URL unless `FLOWS_CLOUD_URL` overrides it, so a +login against one deployment never sends its token to another. An expired +login is refused with the re-login remedy rather than sent. + +Running and syncing work with either kind of token. Deploying, listing and +removing listeners need the interactive `cli:auth` credential the login +produces; a deployment (CI) token gets `session_required` and the CLI says so. + +## Listener deployments + +```sh +flows deploy issue-triage.flow.ts \ + --repo AgentWorkforce/flows \ + --on github:labels=agent \ + --approver khaliqgant [--agents claude,codex] [--name "Issue triage"] [--draft] +flows deployments +flows undeploy +``` + +`flows deploy ` is the CLI form of the agentrelay.com onboarding's +deploy wizard: `POST /api/v1/flows/deploy` stores one self-contained authored +source and creates a proactive listener whose watch rules match the chosen +ticket sources. There is no webhook to register. The workspace's GitHub App +installation (or Slack, Linear, Jira or Shortcut connection) is the ingress; +Cloud ingests events into the workspace's relayfile projection and the +listener's rules match them there. The digest form, +`flows deploy @sha256:… --to file://…`, is unchanged; the positional +decides which form is meant. + +`--on [:key=value,…]` takes `github` (`repository`, `labels`, +`contains`), `slack` (`channel`, `contains`), `linear` (`team`, `contains`), +`jira` (`project`, `contains`) or `shortcut` (`workspace`, `contains`), each +at most once. A GitHub source without `repository` is scoped to `--repo`. +Today a GitHub listener wakes on `issues.opened` and `issues.labeled` only; +pull-request and comment events are filtered out before launch. + +Each matching ticket launches one run of the stored source. Cloud clones +`--repo` at its default branch onto a fresh `relayflow/-` branch, +runs `flows run --local-agent` there, and passes the flow body +`{ approver, issue: { source, title, body, labels, repository, url, … }, event }` +as its input. The flow must therefore be the default body, +`flow(name, header, async (f, input) => …)`; `.on(github.issues(…))` +handlers are checked but are not what Cloud dispatches. `--agents` names the +coding-agent harnesses the flow uses (default `claude`); activation checks +their credentials are connected and refuses with `flow_model_not_connected` +otherwise. `--draft` saves the flow without activating it and skips those +checks. The deploy routes answer refusals as `{ code, error }`, and the CLI +names them (`flow_repository_not_connected`, `flow_name_taken`, …). + Without `--wait`, exit 0 means the server accepted the run. With `--wait`, it means Cloud reported `completed` with a validated `success` completion reason. Failed/cancelled runs and observation/transport failures return 1. Local input, diff --git a/packages/sdk/src/cli.ts b/packages/sdk/src/cli.ts index bcc8b93a7..d0e805872 100644 --- a/packages/sdk/src/cli.ts +++ b/packages/sdk/src/cli.ts @@ -25,6 +25,7 @@ import { parseReplayArgs, replayJournal, type ReplayArgs } from './cli/replay.js import { checkTypeScriptFlow } from './cli/check-typescript.js'; import { runCloudCli } from './cli/cloud-run.js'; import { runCloudSyncCli } from './cli/cloud-sync.js'; +import { parseCloudDeployArgs, runCloudDeployCli, runCloudDeploymentsCli, runCloudUndeployCli, type CloudDeployArgs } from './cli/cloud-deploy.js'; import { isAuthoredFlowPath } from './direct-input.js'; import { parseDeployArgs, runDeploy, type DeployArgs } from './cli/deploy.js'; import { parseDigestReference } from './bundle-transport.js'; @@ -54,6 +55,9 @@ type ParsedArgs = | { command: 'serve-webhook'; dataDir: string; port: number; admitted?: readonly string[] } | { command: 'cloud-run'; value: string; json: boolean; wait: boolean; input: string | undefined; syncCode: boolean } | { command: 'sync'; runId: string; json: boolean; root: string } + | CloudDeployArgs + | { command: 'deployments'; json: boolean } + | { command: 'undeploy'; agentId: string; json: boolean } | { command: 'check'; json: boolean; watch: boolean; value: string } | { command: 'run'; bucket: string | undefined; reuseFromRunId: string | undefined; localAgent: boolean; dataDir: string; input: string | undefined; json: boolean; spawn: boolean; noObserverLink: boolean; allowHumanInfluenced: boolean; value: string } | { command: 'resume'; localAgent: boolean; dataDir: string; json: boolean; spawn: boolean; noObserverLink: boolean; allowHumanInfluenced: boolean; value: string } @@ -68,6 +72,9 @@ const USAGE = [ 'flows add ', 'flows build [--out ] ', 'flows build --verify ', + 'flows deploy --repo --on [:key=value,...] [--on ...] --approver [--agents claude[,codex]] [--name ] [--draft] [--json]', + 'flows deployments [--json]', + 'flows undeploy [--json] ', 'flows deploy @sha256: --to ', 'flows run @sha256: [--bucket ] [--data-dir ] [--json]', 'flows check [--watch] [--json] ', @@ -121,6 +128,9 @@ export async function runCli( if (parsed.command === 'cloud-run') return runCloudCli(parsed, io); if (parsed.command === 'sync') return runCloudSyncCli(parsed, io); + if (parsed.command === 'cloud-deploy') return runCloudDeployCli(parsed, io); + if (parsed.command === 'deployments') return runCloudDeploymentsCli(parsed, io); + if (parsed.command === 'undeploy') return runCloudUndeployCli(parsed, io); if (parsed.command === 'replay') return replayJournal(parsed, io); if (parsed.command === 'build') return runBuild(parsed, io); if (parsed.command === 'deploy') return runDeploy(parsed, io); @@ -424,7 +434,23 @@ function parseArgs(args: readonly string[]): ParsedArgs | undefined { if (command === 'add') return args.length === 2 ? { command: 'add', value: args[1]! } : undefined; if (command === 'replay') return parseReplayArgs(args.slice(1)); if (command === 'build') return parseBuildArgs(args.slice(1)); - if (command === 'deploy') return parseDeployArgs(args.slice(1)); + if (command === 'deploy') { + // The positional decides the form: an authored source deploys a hosted + // listener; a digest reference copies a sealed bundle into a file bucket. + const source = args.slice(1).find(a => !a.startsWith('-') && isAuthoredFlowPath(a)); + return source !== undefined ? parseCloudDeployArgs(args.slice(1)) : parseDeployArgs(args.slice(1)); + } + if (command === 'undeploy') { + const rest = args.slice(1).filter(a => a !== '--json'); + const json = args.length - 1 - rest.length; + if (json > 1 || rest.length !== 1 || rest[0]!.startsWith('-')) return undefined; + return { command: 'undeploy', agentId: rest[0]!, json: json === 1 }; + } + if (command === 'deployments') { + const rest = args.slice(1); + if (rest.length > 1 || (rest.length === 1 && rest[0] !== '--json')) return undefined; + return { command: 'deployments', json: rest.length === 1 }; + } if (command === 'serve-webhook') return parseWebhookArgs(args.slice(1)); if (command === 'hn-monitor') return parseHnMonitorArgs(args.slice(1)); if (command === 'tick') return parseTickArgs(args.slice(1)); diff --git a/packages/sdk/src/cli/cloud-deploy.ts b/packages/sdk/src/cli/cloud-deploy.ts new file mode 100644 index 000000000..583e13a2f --- /dev/null +++ b/packages/sdk/src/cli/cloud-deploy.ts @@ -0,0 +1,157 @@ +import { CloudFlowError } from '../cloud-http.js'; +import { + deployToCloud, listCloudDeployments, parseAgentHarnesses, parseRepository, parseTriggerSource, undeployFromCloud, + type FlowTriggerSource, +} from '../cloud-deploy.js'; +import type { CliIo } from '../cli.js'; + +export interface CloudDeployArgs { + command: 'cloud-deploy'; + value: string; + repo: string; + on: string[]; + approver: string | undefined; + name: string | undefined; + agents: string | undefined; + draft: boolean; + json: boolean; +} + +/** + * `flows deploy --repo --on [:k=v,…] [--on …] + * --approver [--name ] [--json]` + * + * Parsed here rather than in `parseDeployArgs` because the two `deploy` forms + * share nothing but the word: the digest form copies a sealed bundle into a + * file bucket, this one creates a hosted listener. The positional decides. + */ +export function parseCloudDeployArgs(args: readonly string[]): CloudDeployArgs | undefined { + let value: string | undefined; + let repo: string | undefined; + let approver: string | undefined; + let name: string | undefined; + let agents: string | undefined; + let draft = false; + let json = false; + const on: string[] = []; + for (let i = 0; i < args.length; i++) { + const arg = args[i]!; + if (arg === '--json') { + if (json) return undefined; + json = true; + continue; + } + if (arg === '--draft') { + if (draft) return undefined; + draft = true; + continue; + } + if (arg === '--agents') { + const next = args[i + 1]; + if (agents !== undefined || next === undefined || next.startsWith('-')) return undefined; + agents = next; + i += 1; + continue; + } + if (arg === '--repo' || arg === '--approver' || arg === '--name' || arg === '--on') { + const next = args[i + 1]; + if (next === undefined || next.startsWith('-')) return undefined; + i += 1; + if (arg === '--on') { on.push(next); continue; } + if (arg === '--repo') { if (repo !== undefined) return undefined; repo = next; continue; } + if (arg === '--approver') { if (approver !== undefined) return undefined; approver = next; continue; } + if (name !== undefined) return undefined; + name = next; + continue; + } + if (arg.startsWith('-') || value !== undefined) return undefined; + value = arg; + } + if (value === undefined || repo === undefined || on.length === 0) return undefined; + return { command: 'cloud-deploy', value, repo, on, approver, name, agents, draft, json }; +} + +function describeSource(source: FlowTriggerSource): string { + const settings = Object.entries(source.settings).map(([k, v]) => `${k}=${v}`).join(' '); + return settings ? `${source.provider} ${settings}` : source.provider; +} + +export async function runCloudDeployCli(args: CloudDeployArgs, io: CliIo): Promise<0 | 1 | 2> { + try { + if (args.approver === undefined) { + throw new CloudFlowError('invalid_input', + '--approver is required: every launched run receives it as input.approver for f.human.'); + } + const deployment = await deployToCloud({ + path: args.value, + repository: parseRepository(args.repo), + sources: args.on.map(parseTriggerSource), + approver: args.approver, + draft: args.draft, + ...(args.name === undefined ? {} : { name: args.name }), + ...(args.agents === undefined ? {} : { agents: parseAgentHarnesses(args.agents) }), + }); + if (args.json) { + io.stdout(JSON.stringify({ ok: true, ...deployment })); + return 0; + } + io.stdout(`${deployment.status === 'draft' ? 'SAVED' : 'DEPLOYED'} ${deployment.agentId} ${deployment.status}`); + io.stdout(` flow: ${deployment.name} (${args.value}, sha256 ${deployment.sourceSha256.slice(0, 12)})`); + io.stdout(` repository: ${deployment.repository.owner}/${deployment.repository.name}`); + for (const source of deployment.sources) io.stdout(` on: ${describeSource(source)}`); + io.stdout(deployment.status === 'draft' + ? 'Saved without activating; activate it from the Cloud dashboard, or redeploy without --draft.' + : 'Each matching ticket launches a run of this source in a fresh branch; list with: flows deployments'); + return 0; + } catch (error) { + return reportCloudFailure(error, args.json, io); + } +} + +export async function runCloudDeploymentsCli({ json }: { json: boolean }, io: CliIo): Promise<0 | 1 | 2> { + try { + const deployments = await listCloudDeployments(); + if (json) { + io.stdout(JSON.stringify({ ok: true, deployments })); + return 0; + } + if (deployments.length === 0) { + io.stdout('No flow deployments in this workspace.'); + return 0; + } + for (const d of deployments) { + const repo = d.repository ? ` ${d.repository.owner}/${d.repository.name}` : ''; + io.stdout(`${d.agentId} ${d.status} ${JSON.stringify(d.name)}${repo}`); + for (const source of d.sources) io.stdout(` on: ${describeSource(source)}`); + } + return 0; + } catch (error) { + return reportCloudFailure(error, json, io); + } +} + +export async function runCloudUndeployCli({ agentId, json }: { agentId: string; json: boolean }, io: CliIo): Promise<0 | 1 | 2> { + try { + await undeployFromCloud(agentId); + io.stdout(json ? JSON.stringify({ ok: true, agentId, status: 'deleted' }) : `UNDEPLOYED ${agentId}`); + return 0; + } catch (error) { + return reportCloudFailure(error, json, io); + } +} + +function reportCloudFailure(error: unknown, json: boolean, io: CliIo): 1 | 2 { + const code = error instanceof CloudFlowError ? error.code : 'cloud_deploy_failed'; + let message = error instanceof Error ? error.message : 'Cloud deploy failed.'; + // The deploy routes take a browser session or a `cli:auth` token. A + // deployment (CI) token gets 403 `session_required`; say what fixes it. + if (error instanceof CloudFlowError && error.status === 403) { + message += ' Deploying needs an interactive `cli:auth` credential: run `agent-relay cloud login` ' + + '(a deployment token can run flows but not deploy them).'; + } + if (json) io.stdout(JSON.stringify({ ok: false, code, message })); + else io.stderr(`${code}: ${message}`); + return error instanceof CloudFlowError + && (['configuration', 'unsupported_source', 'invalid_input'].includes(error.code) || error.status === 403 || error.status === 401) + ? 2 : 1; +} diff --git a/packages/sdk/src/cloud-deploy.ts b/packages/sdk/src/cloud-deploy.ts new file mode 100644 index 000000000..39988f833 --- /dev/null +++ b/packages/sdk/src/cloud-deploy.ts @@ -0,0 +1,221 @@ +import { randomBytes, createHash } from 'node:crypto'; +import { readFile } from 'node:fs/promises'; +import { loadAuthoredFlow } from './authored-flow-loader.js'; +import { + CloudFlowError, cloudFetch, cloudRequest, isCloudRecord, type CloudConnectionOptions, +} from './cloud-http.js'; + +/** + * Hosted listener deployment: the CLI form of the agentrelay.com onboarding's + * deploy wizard. `POST /api/v1/flows/deploy` stores one self-contained + * authored source and creates a proactive listener whose watch rules match + * the chosen ticket sources on the workspace's relayfile projections. There + * is no webhook to register: the GitHub App installation (or Slack/Linear/ + * Jira/Shortcut connection) is the ingress, and each matching ticket launches + * a run of the stored source with `{ approver, issue, event }` as its input, + * inside a fresh branch of the deployment's repository. + */ + +export const FLOW_TRIGGER_PROVIDERS = ['github', 'linear', 'jira', 'shortcut', 'slack'] as const; +export type FlowTriggerProvider = (typeof FLOW_TRIGGER_PROVIDERS)[number]; + +/** Settings Cloud's launcher prefilter reads per provider (`flow-trigger-sources.ts`). */ +const PROVIDER_SETTINGS: Record = { + github: ['repository', 'labels', 'contains'], + slack: ['channel', 'contains'], + linear: ['team', 'contains'], + jira: ['project', 'contains'], + shortcut: ['workspace', 'contains'], +}; +const MAX_SOURCE_BYTES = 256_000; +const MAX_SETTING_LENGTH = 500; +const REPO_OWNER = /^[A-Za-z0-9-]{1,39}$/u; +const REPO_NAME = /^[A-Za-z0-9_.-]{1,100}$/u; + +export interface FlowTriggerSource { + provider: FlowTriggerProvider; + settings: Record; +} + +export interface DeployToCloudInput { + path: string; + repository: { owner: string; name: string }; + sources: FlowTriggerSource[]; + /** The `f.human` approver handle every launched run receives as `input.approver`. */ + approver: string; + /** Defaults to the flow's declared name. */ + name?: string; + /** Coding-agent harnesses the flow uses; Cloud checks their credentials are connected. Default `["claude"]`. */ + agents?: FlowAgentHarness[]; + /** Save without activating: no connection checks, no listener until activated. */ + draft?: boolean; +} + +export const FLOW_AGENT_HARNESSES = ['claude', 'codex'] as const; +export type FlowAgentHarness = (typeof FLOW_AGENT_HARNESSES)[number]; + +export function parseAgentHarnesses(value: string): FlowAgentHarness[] { + const agents = value.split(',').map(a => a.trim()).filter(Boolean); + if (agents.length === 0 || agents.length > 2 || new Set(agents).size !== agents.length + || !agents.every(a => (FLOW_AGENT_HARNESSES as readonly string[]).includes(a))) { + throw new CloudFlowError('invalid_input', `--agents takes one or two of ${FLOW_AGENT_HARNESSES.join(', ')}, got "${value}".`); + } + return agents as FlowAgentHarness[]; +} + +export interface CloudDeployment { + agentId: string; + name: string; + status: string; + repository: { owner: string; name: string }; + sources: FlowTriggerSource[]; + sourceSha256: string; +} + +export function parseRepository(value: string): { owner: string; name: string } { + const [owner, name, extra] = value.replace(/^https?:\/\/github\.com\//iu, '').replace(/\.git$/iu, '').split('/'); + if (!owner || !name || extra !== undefined || !REPO_OWNER.test(owner) || !REPO_NAME.test(name)) { + throw new CloudFlowError('invalid_input', `Expected --repo /, got "${value}".`); + } + return { owner, name }; +} + +/** `github`, `github:labels=agent,contains=urgent`, `slack:channel=#eng`. */ +export function parseTriggerSource(value: string): FlowTriggerSource { + const colon = value.indexOf(':'); + const provider = (colon === -1 ? value : value.slice(0, colon)).trim(); + if (!(FLOW_TRIGGER_PROVIDERS as readonly string[]).includes(provider)) { + throw new CloudFlowError('invalid_input', + `Unknown trigger provider "${provider}"; expected one of ${FLOW_TRIGGER_PROVIDERS.join(', ')}.`); + } + const allowed = PROVIDER_SETTINGS[provider as FlowTriggerProvider]; + const settings: Record = {}; + if (colon !== -1) { + for (const pair of value.slice(colon + 1).split(',')) { + const eq = pair.indexOf('='); + const key = (eq === -1 ? pair : pair.slice(0, eq)).trim(); + const setting = eq === -1 ? '' : pair.slice(eq + 1).trim(); + if (!allowed.includes(key)) { + throw new CloudFlowError('invalid_input', + `"${key}" is not a ${provider} trigger setting; ${provider} accepts ${allowed.join(', ')}.`); + } + if (!setting || setting.length > MAX_SETTING_LENGTH || key in settings) { + throw new CloudFlowError('invalid_input', `Trigger setting "${key}" must be given once with a non-empty value.`); + } + settings[key] = setting; + } + } + return { provider: provider as FlowTriggerProvider, settings }; +} + +export async function deployToCloud( + input: DeployToCloudInput, options: CloudConnectionOptions = {}, +): Promise { + if (!/\.flow\.ts$/iu.test(input.path)) { + throw new CloudFlowError('unsupported_source', 'flows deploy takes one authored .flow.ts source.'); + } + const bytes = await readFile(input.path); + const source = bytes.toString('utf8'); + if (!bytes.length || Buffer.from(source, 'utf8').compare(bytes) !== 0) { + throw new CloudFlowError('invalid_input', 'Authored source must be nonempty, lossless UTF-8.'); + } + if (bytes.length > MAX_SOURCE_BYTES) { + throw new CloudFlowError('invalid_input', `Authored source exceeds Cloud's ${MAX_SOURCE_BYTES}-byte deploy limit.`); + } + const loaded = await loadAuthoredFlow(input.path); + const definition = loaded.getDefinition(loaded.handle); + if (loaded.graph.length !== 1) { + throw new CloudFlowError('unsupported_source', + 'Cloud deploys one self-contained .flow.ts source without use dependencies.'); + } + if (input.sources.length === 0 || input.sources.length > 10) { + throw new CloudFlowError('invalid_input', 'Give between one and ten --on trigger sources.'); + } + if (new Set(input.sources.map(s => s.provider)).size !== input.sources.length) { + throw new CloudFlowError('invalid_input', 'Each trigger provider can be given once.'); + } + const approver = input.approver.trim(); + if (!approver) throw new CloudFlowError('invalid_input', '--approver must name who approves f.human questions.'); + const name = (input.name ?? definition.name).trim(); + if (!name || name.length > 100) throw new CloudFlowError('invalid_input', 'Deployment name must be 1-100 characters.'); + // A GitHub source scoped to nothing would wake on every repository the + // installation covers; default it to the deployment's own repository. + const sources = input.sources.map(s => s.provider === 'github' && s.settings['repository'] === undefined + ? { ...s, settings: { ...s.settings, repository: `${input.repository.owner}/${input.repository.name}` } } + : s); + options.signal?.throwIfAborted(); + + const whoami = await cloudRequest('/api/v1/auth/whoami', options); + const workspace = isCloudRecord(whoami) && isCloudRecord(whoami.currentWorkspace) ? whoami.currentWorkspace : undefined; + if (workspace === undefined || typeof workspace.id !== 'string' || !workspace.id) { + throw new CloudFlowError('invalid_response', 'Cloud did not report a current workspace for this credential.'); + } + const agents = input.agents ?? ['claude']; + const result = await cloudFetch('/api/v1/flows/deploy', options, { method: 'POST', detail: true, body: JSON.stringify({ + workspaceId: workspace.id, + mode: input.draft ? 'draft' : 'activate', + name, + // The onboarding's workflow-shape label; the CLI deploys authored source as-is. + workflow: 'flows-cli', + source, + handoffId: `flows-cli-${randomBytes(8).toString('hex')}`, + inputs: { approver, agents }, + repository: input.repository, + sources, + }) }); + if (!isCloudRecord(result) || typeof result.agentId !== 'string' || typeof result.status !== 'string') { + throw new CloudFlowError('invalid_response', 'Cloud did not return a deployment.'); + } + return { + agentId: result.agentId, name, status: result.status, + repository: input.repository, sources, + sourceSha256: createHash('sha256').update(bytes).digest('hex'), + }; +} + +export interface CloudDeploymentSummary { + agentId: string; + name: string; + status: string; + repository?: { owner: string; name: string }; + sources: FlowTriggerSource[]; + createdAt?: string; + updatedAt?: string; +} + +export async function listCloudDeployments(options: CloudConnectionOptions = {}): Promise { + const payload = await cloudRequest('/api/v1/agents/flow-deployments', options); + if (!isCloudRecord(payload) || !Array.isArray(payload.deployments)) { + throw new CloudFlowError('invalid_response', 'Cloud did not return a deployments list.'); + } + return payload.deployments.map((row): CloudDeploymentSummary => { + if (!isCloudRecord(row) || typeof row.agentId !== 'string' || typeof row.name !== 'string') { + throw new CloudFlowError('invalid_response', 'Cloud returned a malformed deployment row.'); + } + const repository = isCloudRecord(row.repository) && typeof row.repository.owner === 'string' + && typeof row.repository.name === 'string' ? { owner: row.repository.owner, name: row.repository.name } : undefined; + const sources = Array.isArray(row.sources) ? row.sources.flatMap((s): FlowTriggerSource[] => + isCloudRecord(s) && typeof s.provider === 'string' && (FLOW_TRIGGER_PROVIDERS as readonly string[]).includes(s.provider) + ? [{ provider: s.provider as FlowTriggerProvider, + settings: Object.fromEntries(Object.entries(isCloudRecord(s.settings) ? s.settings : {}) + .flatMap(([k, v]) => typeof v === 'string' ? [[k, v]] : typeof v === 'boolean' ? [[k, String(v)]] : [])) }] + : []) : []; + return { + agentId: row.agentId, name: row.name, status: typeof row.status === 'string' ? row.status : 'unknown', + ...(repository === undefined ? {} : { repository }), sources, + ...(typeof row.createdAt === 'string' ? { createdAt: row.createdAt } : {}), + ...(typeof row.updatedAt === 'string' ? { updatedAt: row.updatedAt } : {}), + }; + }); +} + +const LISTENER_ID = /^[A-Za-z0-9_-]{1,128}$/u; + +/** `DELETE /api/v1/flows/listeners/`: the listener stops; past runs and their journals stay. */ +export async function undeployFromCloud(agentId: string, options: CloudConnectionOptions = {}): Promise { + if (!LISTENER_ID.test(agentId)) throw new CloudFlowError('invalid_input', `"${agentId}" is not a deployment id.`); + const result = await cloudFetch(`/api/v1/flows/listeners/${encodeURIComponent(agentId)}`, options, { method: 'DELETE', detail: true }); + if (!isCloudRecord(result) || result.status !== 'deleted') { + throw new CloudFlowError('invalid_response', 'Cloud did not confirm the deletion.'); + } +} diff --git a/packages/sdk/src/cloud-http.ts b/packages/sdk/src/cloud-http.ts index 9bc25523c..59b61d979 100644 --- a/packages/sdk/src/cloud-http.ts +++ b/packages/sdk/src/cloud-http.ts @@ -1,3 +1,7 @@ +import { readFileSync } from 'node:fs'; +import { homedir } from 'node:os'; +import { join } from 'node:path'; + export interface CloudConnectionOptions { /** Cloud application base URL; defaults to https://agentrelay.com/cloud. */ apiUrl?: string; @@ -21,16 +25,66 @@ export class CloudFlowError extends Error { } } +/** + * The `agent-relay cloud login` credential store. Read only when neither the + * `token` option nor `FLOWS_CLOUD_TOKEN` is set, so an explicit credential + * always wins and this file can change shape without breaking a configured + * caller. Its `apiUrl` becomes the default base URL for the same reason: a + * login against one deployment must not send its token to another. + */ +export function agentRelayCloudAuthPath(env: NodeJS.ProcessEnv = process.env): string { + return join(env['AGENT_RELAY_HOME'] ?? join(homedir(), '.agentworkforce/relay'), 'cloud-auth.json'); +} + +export interface AgentRelayCloudLogin { + apiUrl: string; + accessToken: string; + /** ISO-8601; the store carries it, so an expired login refuses with a real reason. */ + accessTokenExpiresAt?: string; +} + +export function readAgentRelayCloudLogin( + env: NodeJS.ProcessEnv = process.env, + read: (path: string) => string = path => readFileSync(path, 'utf8'), +): AgentRelayCloudLogin | undefined { + let parsed: unknown; + try { + parsed = JSON.parse(read(agentRelayCloudAuthPath(env))); + } catch { + return undefined; + } + if (!isCloudRecord(parsed) || typeof parsed.apiUrl !== 'string' || typeof parsed.accessToken !== 'string' + || !parsed.accessToken.trim()) return undefined; + return { + apiUrl: parsed.apiUrl, accessToken: parsed.accessToken, + ...(typeof parsed.accessTokenExpiresAt === 'string' ? { accessTokenExpiresAt: parsed.accessTokenExpiresAt } : {}), + }; +} + export function cloudConnection(options: CloudConnectionOptions): { baseUrl: string; token: string } { - const rawToken = options.token ?? process.env['FLOWS_CLOUD_TOKEN']; + let rawToken = options.token ?? process.env['FLOWS_CLOUD_TOKEN']; + let loginApiUrl: string | undefined; + if (rawToken === undefined) { + const login = readAgentRelayCloudLogin(); + if (login !== undefined) { + const expiresAt = login.accessTokenExpiresAt === undefined ? Number.NaN : Date.parse(login.accessTokenExpiresAt); + if (Number.isFinite(expiresAt) && expiresAt <= Date.now()) { + throw new CloudFlowError('configuration', + 'The agent-relay cloud login has expired. Run `agent-relay cloud login` again, or set FLOWS_CLOUD_TOKEN.'); + } + rawToken = login.accessToken; + loginApiUrl = login.apiUrl; + } + } const token = rawToken?.trim(); if (!token || /[\r\n]/u.test(rawToken!) || /^(?:rk|ot)_live_/u.test(token)) { throw new CloudFlowError('configuration', - 'Set FLOWS_CLOUD_TOKEN to a scoped Cloud API token (workflow:invoke:write and workflow:runs:read).'); + 'Set FLOWS_CLOUD_TOKEN to a scoped Cloud API token (workflow:invoke:write and workflow:runs:read), ' + + 'or sign in with `agent-relay cloud login`.'); } let url: URL; try { - url = new URL(options.apiUrl ?? process.env['FLOWS_CLOUD_URL'] ?? 'https://agentrelay.com/cloud'); + url = new URL(options.apiUrl ?? process.env['FLOWS_CLOUD_URL'] ?? loginApiUrl ?? 'https://agentrelay.com/cloud'); } catch { throw new CloudFlowError('configuration', 'FLOWS_CLOUD_URL must be an absolute Cloud application base URL.'); } @@ -52,7 +106,7 @@ export async function cloudRequest( } export interface CloudFetchInit { - method: 'GET' | 'POST' | 'PUT'; + method: 'GET' | 'POST' | 'PUT' | 'DELETE'; body?: string | Uint8Array; contentType?: string; /** @@ -61,6 +115,13 @@ export interface CloudFetchInit { * only; it is never persisted or logged. */ bearerToken?: string; + /** + * On a non-2xx response, read a `{ code, error }` refusal from the body and + * name it. Only those two string fields are ever surfaced, so a route that + * answers with structured refusals (the deploy routes) can explain itself + * without this client echoing arbitrary response bodies. + */ + detail?: boolean; } /** One authenticated Cloud request. Every transport error is typed; no retries. */ @@ -98,7 +159,10 @@ export async function cloudFetch( } if (!response.ok) { // Do not echo server response bodies: they may contain credentials or source. - throw new CloudFlowError('http_error', `Cloud request failed with HTTP ${response.status}.`, response.status); + const refusal = init.detail ? await structuredRefusal(response) : undefined; + throw new CloudFlowError('http_error', refusal === undefined + ? `Cloud request failed with HTTP ${response.status}.` + : `Cloud refused (${refusal.code}): ${refusal.error}`, response.status); } try { return await response.json(); @@ -109,6 +173,14 @@ export async function cloudFetch( } } +async function structuredRefusal(response: Response): Promise<{ code: string; error: string } | undefined> { + let body: unknown; + try { body = await response.json(); } catch { return undefined; } + if (!isCloudRecord(body) || typeof body.code !== 'string' || typeof body.error !== 'string') return undefined; + if (!/^[a-z0-9_]{1,64}$/u.test(body.code) || body.error.length > 500) return undefined; + return { code: body.code, error: body.error }; +} + // Defensive path-segment constraint; accepting a new server ID format needs an SDK change. export function cloudRunId(value: unknown): string { if (typeof value !== 'string' || !/^[A-Za-z0-9_-]{1,128}$/u.test(value)) { diff --git a/packages/sdk/src/index.ts b/packages/sdk/src/index.ts index 06c96829b..497ba773b 100644 --- a/packages/sdk/src/index.ts +++ b/packages/sdk/src/index.ts @@ -67,6 +67,10 @@ export { downloadCloudPatch, applyCloudPatch, packWorkingTree, MAX_SYNC_BYTES, type CloudPatch, type PackedTree, } from './cloud-sync.js'; +export { + deployToCloud, listCloudDeployments, undeployFromCloud, parseRepository, parseTriggerSource, FLOW_TRIGGER_PROVIDERS, + type DeployToCloudInput, type CloudDeployment, type CloudDeploymentSummary, type FlowTriggerSource, type FlowTriggerProvider, +} from './cloud-deploy.js'; export { canonicalize, specHash } from './canonical.js'; export { diff --git a/packages/sdk/tests/cloud-deploy.test.ts b/packages/sdk/tests/cloud-deploy.test.ts new file mode 100644 index 000000000..7ba94402e --- /dev/null +++ b/packages/sdk/tests/cloud-deploy.test.ts @@ -0,0 +1,273 @@ +import { mkdir, mkdtemp, rm, symlink, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { runCli } from '../src/cli.js'; +import { deployToCloud, parseAgentHarnesses, parseRepository, parseTriggerSource } from '../src/cloud-deploy.js'; +import { cloudConnection } from '../src/cloud-http.js'; + +const dirs: string[] = []; +afterEach(async () => { + vi.unstubAllEnvs(); + vi.restoreAllMocks(); + for (const dir of dirs.splice(0)) await rm(dir, { recursive: true, force: true }); +}); + +async function tempDir(prefix: string): Promise { + const dir = await mkdtemp(join(tmpdir(), prefix)); + dirs.push(dir); + return dir; +} + +interface Call { method: string; path: string; auth: string | undefined; body: unknown } +function cloud(routes: Record { status?: number; body: unknown } | unknown>) { + const calls: Call[] = []; + vi.spyOn(globalThis, 'fetch').mockImplementation(async (input, init) => { + init?.signal?.throwIfAborted(); + const path = new URL(String(input)).pathname; + const call: Call = { method: init?.method ?? 'GET', path, + auth: new Headers(init?.headers).get('authorization') ?? undefined, + body: typeof init?.body === 'string' ? JSON.parse(init.body) : undefined }; + calls.push(call); + const route = routes[path]; + if (!route) return new Response('{"error":"not found"}', { status: 404 }); + const answer = route(call); + const { status, body } = answer !== null && typeof answer === 'object' && 'status' in answer && 'body' in answer + ? answer as { status?: number; body: unknown } : { status: 200, body: answer }; + return new Response(JSON.stringify(body), { status: status ?? 200, headers: { 'content-type': 'application/json' } }); + }); + vi.stubEnv('FLOWS_CLOUD_URL', 'https://cloud-contract.example'); + vi.stubEnv('FLOWS_CLOUD_TOKEN', 'test-cli-auth-token'); + return calls; +} + +const WHOAMI = { authenticated: true, currentWorkspace: { id: 'ws-1', slug: 'default' } }; + +async function authoredFlow(name = 'issue-triage'): Promise { + const dir = await tempDir('cloud-deploy-'); + await symlink(join(process.cwd(), 'node_modules'), join(dir, 'node_modules'), 'dir'); + const path = join(dir, `${name}.flow.ts`); + await writeFile(path, "import { flow } from '@relayflows/surface';\n" + + `export default flow<{ issue: { title: string }; approver: string }>('${name}', { budget: '$5/run' }, async (f, input) => {\n` + + " await f.run(`echo ${JSON.stringify(input.issue.title)}`);\n f.done('success');\n});\n"); + return path; +} + +describe('trigger source and repository parsing', () => { + it.each([ + ['github', { provider: 'github', settings: {} }], + ['github:labels=agent,contains=urgent', { provider: 'github', settings: { labels: 'agent', contains: 'urgent' } }], + ['slack:channel=#eng', { provider: 'slack', settings: { channel: '#eng' } }], + ['linear:team=ENG', { provider: 'linear', settings: { team: 'ENG' } }], + ])('parses %s', (value, expected) => { + expect(parseTriggerSource(value)).toEqual(expected); + }); + + it.each(['gitlab', 'github:channel=x', 'github:labels=', 'github:labels=a,labels=b', 'slack:labels=x']) + ('refuses %s', (value) => { + expect(() => parseTriggerSource(value)).toThrow(expect.objectContaining({ code: 'invalid_input' })); + }); + + it('accepts owner/name and GitHub URLs, refusing anything else', () => { + expect(parseRepository('AgentWorkforce/flows')).toEqual({ owner: 'AgentWorkforce', name: 'flows' }); + expect(parseRepository('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/AgentWorkforce/flows.git')).toEqual({ owner: 'AgentWorkforce', name: 'flows' }); + for (const bad of ['flows', 'a/b/c', 'bad owner/x', '']) { + expect(() => parseRepository(bad)).toThrow(expect.objectContaining({ code: 'invalid_input' })); + } + }); +}); + +describe('deployToCloud', () => { + it('resolves the workspace, then posts the exact source with defaulted GitHub scope and the flow name', async () => { + const path = await authoredFlow(); + const calls = cloud({ + '/api/v1/auth/whoami': () => WHOAMI, + '/api/v1/flows/deploy': () => ({ status: 201, body: { agentId: 'agent-1', status: 'listening', sources: [], repository: {} } }), + }); + const deployment = await deployToCloud({ + path, repository: { owner: 'AgentWorkforce', name: 'flows' }, + sources: [parseTriggerSource('github:labels=agent'), parseTriggerSource('slack:channel=C123')], + approver: 'khaliqgant', + }); + expect(calls.map(c => [c.method, c.path])).toEqual([['GET', '/api/v1/auth/whoami'], ['POST', '/api/v1/flows/deploy']]); + const body = calls[1]!.body as Record; + expect(body).toMatchObject({ + workspaceId: 'ws-1', mode: 'activate', name: 'issue-triage', workflow: 'flows-cli', + inputs: { approver: 'khaliqgant', agents: ['claude'] }, repository: { owner: 'AgentWorkforce', name: 'flows' }, + sources: [ + { provider: 'github', settings: { labels: 'agent', repository: 'AgentWorkforce/flows' } }, + { provider: 'slack', settings: { channel: 'C123' } }, + ], + }); + expect(body.handoffId).toMatch(/^flows-cli-[a-f0-9]{16}$/u); + expect(body.source).toContain("flow<{ issue: { title: string }; approver: string }>('issue-triage'"); + expect(deployment).toMatchObject({ agentId: 'agent-1', status: 'listening', name: 'issue-triage' }); + expect(deployment.sourceSha256).toMatch(/^[a-f0-9]{64}$/u); + }); + + it('refuses non-authored sources, empty sources, duplicate providers and a blank approver before HTTP', async () => { + const path = await authoredFlow(); + const calls = cloud({}); + const base = { path, repository: { owner: 'o', name: 'r' }, sources: [parseTriggerSource('github')], approver: 'k' }; + await expect(deployToCloud({ ...base, path: 'flow.yaml' })).rejects.toMatchObject({ code: 'unsupported_source' }); + await expect(deployToCloud({ ...base, sources: [] })).rejects.toMatchObject({ code: 'invalid_input' }); + await expect(deployToCloud({ ...base, sources: [parseTriggerSource('github'), parseTriggerSource('github')] })) + .rejects.toMatchObject({ code: 'invalid_input' }); + await expect(deployToCloud({ ...base, approver: ' ' })).rejects.toMatchObject({ code: 'invalid_input' }); + expect(calls).toHaveLength(0); + }); +}); + +describe('flows deploy / flows deployments', () => { + it('deploys an authored flow as a listener and prints the sources', async () => { + const path = await authoredFlow('triage'); + cloud({ + '/api/v1/auth/whoami': () => WHOAMI, + '/api/v1/flows/deploy': () => ({ status: 201, body: { agentId: 'agent-9', status: 'listening' } }), + }); + const out: string[] = []; + const code = await runCli(['deploy', path, '--repo', 'AgentWorkforce/flows', '--on', 'github:labels=agent', '--approver', 'khaliqgant'], + { stdout: line => out.push(line), stderr: line => out.push(`ERR ${line}`) }); + expect(code, out.join('\n')).toBe(0); + expect(out[0]).toBe('DEPLOYED agent-9 listening'); + expect(out).toContainEqual(' repository: AgentWorkforce/flows'); + expect(out).toContainEqual(' on: github labels=agent repository=AgentWorkforce/flows'); + }); + + it('passes --agents and --draft through, and names a structured refusal', async () => { + const path = await authoredFlow('drafted'); + const calls = cloud({ + '/api/v1/auth/whoami': () => WHOAMI, + '/api/v1/flows/deploy': (call) => (call.body as { mode: string }).mode === 'draft' + ? { status: 201, body: { agentId: 'agent-d', status: 'draft' } } + : { status: 409, body: { error: 'Connect an active codex subscription before activating this flow.', code: 'flow_model_not_connected' } }, + }); + const out: string[] = []; + expect(await runCli(['deploy', path, '--repo', 'o/r', '--on', 'github', '--approver', 'k', '--agents', 'claude,codex', '--draft'], + { stdout: line => out.push(line), stderr: line => out.push(`ERR ${line}`) })).toBe(0); + expect(out[0]).toBe('SAVED agent-d draft'); + expect(calls.at(-1)!.body).toMatchObject({ mode: 'draft', inputs: { approver: 'k', agents: ['claude', 'codex'] } }); + const errors: string[] = []; + expect(await runCli(['deploy', path, '--repo', 'o/r', '--on', 'github', '--approver', 'k', '--agents', 'codex'], + { stdout: () => {}, stderr: line => errors.push(line) })).toBe(1); + expect(errors[0]).toContain('flow_model_not_connected'); + expect(errors[0]).toContain('codex subscription'); + expect(() => parseAgentHarnesses('gemini')).toThrow(expect.objectContaining({ code: 'invalid_input' })); + expect(() => parseAgentHarnesses('claude,claude')).toThrow(expect.objectContaining({ code: 'invalid_input' })); + }); + + it('explains a 403 as a missing cli:auth login and exits 2', async () => { + const path = await authoredFlow(); + cloud({ + '/api/v1/auth/whoami': () => WHOAMI, + '/api/v1/flows/deploy': () => ({ status: 403, body: { error: 'Forbidden', code: 'session_required' } }), + }); + const errors: string[] = []; + expect(await runCli(['deploy', path, '--repo', 'o/r', '--on', 'github', '--approver', 'k', '--json'], + { stdout: line => errors.push(line), stderr: () => {} })).toBe(2); + expect(JSON.parse(errors[0]!)).toMatchObject({ ok: false, code: 'http_error' }); + expect(errors[0]).toContain('agent-relay cloud login'); + }); + + it.each([ + ['deploy', 'x.flow.ts'], + ['deploy', 'x.flow.ts', '--repo', 'o/r'], + ['deploy', 'x.flow.ts', '--on', 'github'], + ['deploy', 'x.flow.ts', '--repo', 'o/r', '--repo', 'p/q', '--on', 'github'], + ['deploy', 'x.flow.ts', '--repo', 'o/r', '--on'], + ['deployments', 'extra'], + ['deployments', '--data-dir', 'x'], + ])('refuses argv %j before any request', async (...args) => { + const fetch = vi.spyOn(globalThis, 'fetch'); + expect(await runCli(args, { stdout: () => {}, stderr: () => {} })).toBe(2); + expect(fetch).not.toHaveBeenCalled(); + }); + + it('requires --approver with a reason, before any request', async () => { + const path = await authoredFlow(); + const calls = cloud({}); + const errors: string[] = []; + expect(await runCli(['deploy', path, '--repo', 'o/r', '--on', 'github'], { stdout: () => {}, stderr: line => errors.push(line) })).toBe(2); + expect(errors[0]).toContain('--approver'); + expect(calls).toHaveLength(0); + }); + + it('still routes the digest form to the bucket deploy', async () => { + const fetch = vi.spyOn(globalThis, 'fetch'); + const errors: string[] = []; + const code = await runCli(['deploy', 'flow@sha256:' + 'a'.repeat(64), '--to', 'file:///nowhere'], + { stdout: () => {}, stderr: line => errors.push(line) }); + expect(code).toBe(2); + expect(errors[0]).toContain('bundle_missing_locally'); + expect(fetch).not.toHaveBeenCalled(); + }); + + it('lists deployments with their sources', async () => { + cloud({ + '/api/v1/agents/flow-deployments': () => ({ deployments: [ + { agentId: 'a-1', name: 'Garden', status: 'listening', repository: { owner: 'AgentWorkforce', name: 'flows' }, + sources: [{ provider: 'github', settings: { labels: 'garden-ready', repository: 'AgentWorkforce/flows' } }] }, + { agentId: 'a-2', name: 'Jira', status: 'paused', sources: [{ provider: 'jira', settings: { project: 'OPS' } }] }, + ] }), + }); + const out: string[] = []; + expect(await runCli(['deployments'], { stdout: line => out.push(line), stderr: () => {} })).toBe(0); + expect(out).toEqual([ + 'a-1 listening "Garden" AgentWorkforce/flows', + ' on: github labels=garden-ready repository=AgentWorkforce/flows', + 'a-2 paused "Jira"', + ' on: jira project=OPS', + ]); + }); +}); + +describe('agent-relay cloud login fallback', () => { + async function loginStore(record: Record): Promise { + const home = await tempDir('relay-home-'); + await mkdir(home, { recursive: true }); + await writeFile(join(home, 'cloud-auth.json'), JSON.stringify(record)); + vi.stubEnv('AGENT_RELAY_HOME', home); + vi.stubEnv('FLOWS_CLOUD_TOKEN', ''); + vi.stubEnv('FLOWS_CLOUD_URL', ''); + // vi.stubEnv('', '') leaves an empty string; emulate "unset" precisely. + delete process.env['FLOWS_CLOUD_TOKEN']; + delete process.env['FLOWS_CLOUD_URL']; + } + + it('uses the login token and its API URL when no explicit credential is configured', async () => { + await loginStore({ apiUrl: 'https://login.example/cloud', accessToken: 'login-token', + accessTokenExpiresAt: new Date(Date.now() + 60_000).toISOString() }); + expect(cloudConnection({})).toEqual({ baseUrl: 'https://login.example/cloud', token: 'login-token' }); + }); + + it('lets FLOWS_CLOUD_TOKEN win over the login store', async () => { + await loginStore({ apiUrl: 'https://login.example/cloud', accessToken: 'login-token' }); + vi.stubEnv('FLOWS_CLOUD_TOKEN', 'explicit'); + expect(cloudConnection({})).toEqual({ baseUrl: 'https://agentrelay.com/cloud', token: 'explicit' }); + }); + + it('refuses an expired login with the re-login remedy, and a missing store with the configuration message', async () => { + await loginStore({ apiUrl: 'https://login.example/cloud', accessToken: 'stale', + accessTokenExpiresAt: new Date(Date.now() - 1).toISOString() }); + expect(() => cloudConnection({})).toThrow(/agent-relay cloud login/u); + vi.stubEnv('AGENT_RELAY_HOME', join(await tempDir('relay-home-empty-'), 'nope')); + expect(() => cloudConnection({})).toThrow(expect.objectContaining({ code: 'configuration' })); + }); +}); + +describe('flows undeploy', () => { + it('deletes the listener and reports it', async () => { + const calls = cloud({ '/api/v1/flows/listeners/agent-9': () => ({ listenerId: 'agent-9', status: 'deleted' }) }); + const out: string[] = []; + expect(await runCli(['undeploy', 'agent-9'], { stdout: line => out.push(line), stderr: () => {} })).toBe(0); + expect(out).toEqual(['UNDEPLOYED agent-9']); + expect(calls.map(c => [c.method, c.path])).toEqual([['DELETE', '/api/v1/flows/listeners/agent-9']]); + }); + + it.each([['undeploy'], ['undeploy', 'a', 'b'], ['undeploy', '--json', '--json', 'a'], ['undeploy', '--dir', 'a']]) + ('refuses argv %j', async (...args) => { + const fetch = vi.spyOn(globalThis, 'fetch'); + expect(await runCli(args, { stdout: () => {}, stderr: () => {} })).toBe(2); + expect(fetch).not.toHaveBeenCalled(); + }); +}); From fa9af7362407669d981cf9ef189848d0f62de87e Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 17 Sep 2026 12:23:57 -0700 Subject: [PATCH 3/4] =?UTF-8?q?fix(sdk):=20harden=20code=20sync=20per=20re?= =?UTF-8?q?view=20=E2=80=94=20modes,=20memory,=20git=20failures,=20symlink?= =?UTF-8?q?s,=20aborts,=20deletions?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review findings on #440 (Devin, Cursor, CodeRabbit), each addressed: - Executable bits were written as 0644 for every file; the archive now carries 0755 for anything with an execute bit, so a synced `./run.sh` runs on the host. Proven live (run 418dab2c, gate on the script's output). - Packing held the tree, a `Buffer.concat` copy and the gzip output at once. The ustar is now streamed through `createGzip` into a spooled temporary file that is disposed after upload; peak footprint is one file plus the compressor's window, and the only whole-archive copy is the compressed one the sized PUT needs. - Any `git ls-files` failure fell back to a plain walk, which ignores `.gitignore`: a locked index or missing `git` would have uploaded `.env`. Now only a real "not a git repository" walks; a `.git` entry git cannot read, or an `ls-files` failure inside a checkout, refuses as `sync_unsupported` with a one-line safe summary. Proven live: the gitignored `.env` was absent on the host (`test ! -e .env` passed). - Symlinks whose target resolves outside the tree (absolute, or through `..`) are dropped and reported as `sync_link_skipped` warnings. - Ctrl-C during prepare, packing or upload reported `admission_unknown`. `runInCloud` now signals `onSubmit` immediately before the one non-idempotent request; earlier interruptions report `submission_aborted` and say rerunning is safe. - `flows sync` reported files from `+++ b/` lines, so deletions vanished from the list. Paths now come from `diff --git` headers. - Overwriting control files via a patch is by design and unchanged: the patch is the user's own run's diff, applied uncommitted; the CLI now says so and the doc spells out the review-before-keeping contract. Also found by the live proof: preflight probed path-like deterministic commands against the flow file's directory, while the kernel runs them in the daemon's cwd (`exec_det.rs`). Locally those coincide; on Cloud the source sits in the state directory and the code in the mount, so a synced `./run.sh` was refused `command_missing` before running. The probe now uses `process.cwd()`, which is where the step will execute. Co-Authored-By: Claude Opus 5 (1M context) --- docs/CLOUD.md | 31 +++- packages/sdk/src/cli/check.ts | 7 +- packages/sdk/src/cli/cloud-run.ts | 16 +- packages/sdk/src/cli/cloud-sync.ts | 9 +- packages/sdk/src/cloud-run.ts | 27 +++- packages/sdk/src/cloud-sync.ts | 160 ++++++++++++++----- packages/sdk/src/index.ts | 2 +- packages/sdk/tests/check-command-cwd.test.ts | 34 ++++ packages/sdk/tests/cloud-sync.test.ts | 118 ++++++++++++-- 9 files changed, 326 insertions(+), 78 deletions(-) create mode 100644 packages/sdk/tests/check-command-cwd.test.ts diff --git a/docs/CLOUD.md b/docs/CLOUD.md index e21ebff53..80b7793a9 100644 --- a/docs/CLOUD.md +++ b/docs/CLOUD.md @@ -55,11 +55,22 @@ flows sync # apply the run's changes to this checkout run — every `f.run` and every `f.agent` — executes inside that tree, the way v1's `agent-relay cloud run --sync-code` did. Inside a Git checkout the upload is `git ls-files --cached --others --exclude-standard`: `.gitignore` governs, -untracked files ride along, `.git` and `node_modules` never do. Outside Git, -every file except those two directories. The limit is 256 MiB uncompressed. -The flow source itself is still sent in the request body, so it must stay -self-contained; sibling imports inside the tree are not resolved by the hosted -runner. +untracked files ride along, `.git` and `node_modules` never do. A checkout +whose `git` fails for any other reason (a corrupt `.git`, a locked index, no +`git` on PATH) is refused as `sync_unsupported` rather than widened to a plain +walk, so an ignored `.env` never reaches Cloud because Git was unavailable. +Outside Git, every file except those two directories — there is no ignore +rule there, so keep secrets out of such a tree. Executable bits are preserved; +symlinks whose target resolves outside the tree are dropped and listed as +`sync_link_skipped` warnings. The archive is streamed to a temporary file, so +packing costs one file plus the compressor's window, not the tree. The limit is +256 MiB uncompressed. The flow source itself is still sent in the request +body, so it must stay self-contained; sibling imports inside the tree are not +resolved by the hosted runner. + +Interrupting before the run request — during prepare, packing or upload — +reports `submission_aborted`: nothing was admitted and rerunning is safe. Only +an interrupted submission itself reports `admission_unknown`. The transport is the Cloud API only. `POST /api/v1/workflows/prepare` must answer with a `cloud-api` workflow-storage backend (Cloud's R2), the archive @@ -74,9 +85,13 @@ pointing at a tree Cloud does not hold. `flows sync ` fetches the sandbox's post-run diff from `/api/v1/workflows/runs//patch` and applies it with `git apply` after a `--check` pass, so a conflicting patch leaves the tree untouched -(`patch_conflict`, exit 2). Runs that declared several mounted paths carry one -patch per path and are refused here (`sync_unsupported`). `--dir ` -targets a checkout other than the current directory. +(`patch_conflict`, exit 2). The patch lands in the working tree uncommitted +and every touched path is listed, deletions included: what the run changed — +it is your own flow's output, but it is agent output — is reviewed with +`git diff` before any of it is kept, the same contract v1's `cloud sync` had. +Runs that declared several mounted paths carry one patch per path and are +refused here (`sync_unsupported`). `--dir ` targets a checkout other +than the current directory. A synced run and a Cloud repository grant are mutually exclusive on the server: `--sync-code` is the local-driven development loop, and diff --git a/packages/sdk/src/cli/check.ts b/packages/sdk/src/cli/check.ts index d6a2bf2c5..0aee42a86 100644 --- a/packages/sdk/src/cli/check.ts +++ b/packages/sdk/src/cli/check.ts @@ -340,7 +340,12 @@ function systemProbes(flowDirectory: string, config: ProjectConfig): PreflightPr helper: helperReady, cli: (cli, source, model) => probeCli(cli, source === 'project' ? config.directory : flowDirectory, model), executor: (trigger) => config.executors.includes(trigger.executor), - command: (binary) => executableExists(binary, flowDirectory), + // A deterministic step runs in the daemon's working directory — the + // directory `flows run` was invoked from, or Cloud's code mount — not in + // the flow file's. Probing `./x` against the flow's directory answered a + // question the kernel never asks, and refused a Cloud run whose synced + // tree held the script while its source sat in the state directory. + command: (binary) => executableExists(binary, process.cwd()), }; } diff --git a/packages/sdk/src/cli/cloud-run.ts b/packages/sdk/src/cli/cloud-run.ts index 4cbeb0ba1..6c7d38949 100644 --- a/packages/sdk/src/cli/cloud-run.ts +++ b/packages/sdk/src/cli/cloud-run.ts @@ -16,8 +16,12 @@ export async function runCloudCli( process.once('SIGINT', abort); process.once('SIGTERM', abort); let runId: string | undefined; + // Flips at the run submission. An interruption before it — during prepare, + // packing or upload — admitted nothing and is safe to retry; only an + // interrupted submission has unknown admission. + let submitting = false; try { - const options: RunInCloudOptions = { signal: controller.signal }; + const options: RunInCloudOptions = { signal: controller.signal, onSubmit: () => { submitting = true; } }; if (isAuthoredFlowPath(path)) { // Same parse as a local direct run, so a file-or-inline argument means // the same thing on both sides of `--cloud`. @@ -38,6 +42,9 @@ export async function runCloudCli( io.stdout(receipt.apiUrl); if (receipt.synced) { io.stdout(`SYNCED ${receipt.synced.files} files (${receipt.synced.bytes} bytes); pull changes with: flows sync ${receipt.runId}`); + for (const link of receipt.synced.skippedLinks) { + io.stderr(`WARNING [sync_link_skipped] ${link} points outside the synced tree and was not uploaded`); + } } } if (!wait) { @@ -50,11 +57,14 @@ export async function runCloudCli( else io.stdout(`${run.status.toUpperCase()} ${run.runId} completionReason: ${'completionReason' in run ? run.completionReason : 'unavailable'}`); return ok ? 0 : 1; } catch (error) { - const code = controller.signal.aborted ? (runId ? 'observation_aborted' : 'admission_unknown') + const code = controller.signal.aborted + ? runId ? 'observation_aborted' : submitting ? 'admission_unknown' : 'submission_aborted' : error instanceof CloudFlowError ? error.code : 'cloud_run_failed'; const message = controller.signal.aborted ? runId ? 'Stopped observing; the hosted run has not been cancelled.' - : 'Submission interrupted before a receipt was received. Admission is unknown; Cloud may have started the run. Do not resubmit blindly.' + : submitting + ? 'Submission interrupted before a receipt was received. Admission is unknown; Cloud may have started the run. Do not resubmit blindly.' + : 'Interrupted before the run was submitted; nothing was admitted. Safe to run again.' : error instanceof Error ? error.message : 'Cloud run failed.'; if (json) io.stdout(JSON.stringify({ ok: false, code, message, ...(runId ? { runId } : {}) })); else io.stderr(`${code}: ${message}${runId ? ` (run ${runId})` : ''}`); diff --git a/packages/sdk/src/cli/cloud-sync.ts b/packages/sdk/src/cli/cloud-sync.ts index 267a79fe0..f890c1ea4 100644 --- a/packages/sdk/src/cli/cloud-sync.ts +++ b/packages/sdk/src/cli/cloud-sync.ts @@ -1,5 +1,5 @@ import { CloudFlowError } from '../cloud-http.js'; -import { applyCloudPatch, downloadCloudPatch } from '../cloud-sync.js'; +import { applyCloudPatch, downloadCloudPatch, patchedPaths } from '../cloud-sync.js'; import type { CliIo } from '../cli.js'; /** @@ -20,9 +20,12 @@ export async function runCloudSyncCli( return 0; } applyCloudPatch(root, patch); - const files = [...patch.matchAll(/^\+\+\+ b\/(.+)$/gmu)].map(match => match[1]!); + const files = patchedPaths(patch); if (json) io.stdout(JSON.stringify({ ok: true, runId, hasChanges: true, applied: true, files })); - else io.stdout(`APPLIED ${runId}: ${files.length} file${files.length === 1 ? '' : 's'}${files.length ? `\n ${files.join('\n ')}` : ''}`); + else { + io.stdout(`APPLIED ${runId}: ${files.length} file${files.length === 1 ? '' : 's'}${files.length ? `\n ${files.join('\n ')}` : ''}`); + io.stdout('Applied to the working tree, uncommitted: review with git diff before keeping it.'); + } return 0; } catch (error) { const code = error instanceof CloudFlowError ? error.code : 'cloud_sync_failed'; diff --git a/packages/sdk/src/cloud-run.ts b/packages/sdk/src/cloud-run.ts index 279f9b116..ca69b876e 100644 --- a/packages/sdk/src/cloud-run.ts +++ b/packages/sdk/src/cloud-run.ts @@ -29,6 +29,13 @@ export interface RunInCloudOptions extends CloudConnectionOptions { * resolved by the hosted runner. */ syncCode?: { root: string }; + /** + * Called immediately before the one non-idempotent request, the run + * submission. Everything before it (prepare, pack, upload) is safe to + * retry, so a caller classifying an interruption can tell "nothing was + * admitted" from "admission unknown". + */ + onSubmit?: () => void; } export interface CloudRunReceipt { runId: string; @@ -38,7 +45,7 @@ export interface CloudRunReceipt { /** Authenticated run API resource; this is not a public sharing URL. */ apiUrl: string; /** Present when a working tree was synced: what was uploaded, by count and size. */ - synced?: { files: number; bytes: number }; + synced?: { files: number; bytes: number; skippedLinks: string[] }; } export interface CloudAuthoredAuthority { readonly schemaVersion: 1; @@ -133,14 +140,22 @@ export async function runInCloud( // under it, so the run request below names code Cloud already holds. A // refused backend or failed upload therefore never leaves a launched run // pointing at a tree that is not there. - let synced: { runId: string; codeKey: string; files: number; bytes: number } | undefined; + let synced: { runId: string; codeKey: string; files: number; bytes: number; skippedLinks: string[] } | undefined; if (options.syncCode !== undefined) { const prepared = await prepareCloudSync(options); options.signal?.throwIfAborted(); - const packed = packWorkingTree(options.syncCode.root); - await uploadCloudCode(prepared, packed.tarball, options); - synced = { runId: prepared.runId, codeKey: prepared.codeKey, files: packed.files.length, bytes: packed.bytes }; + const packed = await packWorkingTree(options.syncCode.root); + try { + options.signal?.throwIfAborted(); + await uploadCloudCode(prepared, packed, options); + } finally { + packed.dispose(); + } + synced = { runId: prepared.runId, codeKey: prepared.codeKey, files: packed.files.length, bytes: packed.bytes, + skippedLinks: packed.skippedLinks }; } + options.signal?.throwIfAborted(); + options.onSubmit?.(); const result = await cloudRequest('/api/v1/workflows/run', options, { // JSON is a YAML subset. Sending canonical data preserves the exact spec // while using the server's existing YAML-to-config admission path. @@ -162,7 +177,7 @@ export async function runInCloud( } return { runId, status: result.status, specHash: hash, apiUrl: `${baseUrl}/api/v1/workflows/runs/${runId}`, - ...(synced === undefined ? {} : { synced: { files: synced.files, bytes: synced.bytes } }), + ...(synced === undefined ? {} : { synced: { files: synced.files, bytes: synced.bytes, skippedLinks: synced.skippedLinks } }), }; } diff --git a/packages/sdk/src/cloud-sync.ts b/packages/sdk/src/cloud-sync.ts index ba596adcc..783ba69ee 100644 --- a/packages/sdk/src/cloud-sync.ts +++ b/packages/sdk/src/cloud-sync.ts @@ -1,7 +1,11 @@ import { spawnSync } from 'node:child_process'; -import { lstatSync, readFileSync, readdirSync, readlinkSync } from 'node:fs'; -import { join, relative, resolve, sep } from 'node:path'; -import { gzipSync } from 'node:zlib'; +import { createReadStream, createWriteStream, existsSync, lstatSync, mkdtempSync, readdirSync, readlinkSync, rmSync } from 'node:fs'; +import { readFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { dirname, isAbsolute, join, relative, resolve, sep } from 'node:path'; +import { Readable } from 'node:stream'; +import { pipeline } from 'node:stream/promises'; +import { createGzip } from 'node:zlib'; import { CloudFlowError, cloudFetch, cloudRequest, cloudRunId, isCloudRecord, type CloudConnectionOptions, } from './cloud-http.js'; @@ -56,57 +60,114 @@ export async function prepareCloudSync(options: CloudConnectionOptions): Promise } export interface PackedTree { - tarball: Buffer; + /** The gzip'd ustar, spooled to a temporary file; `dispose()` removes it. */ + archivePath: string; /** Archive members, `root`-relative POSIX paths, sorted. */ files: string[]; /** Uncompressed byte total across regular files. */ bytes: number; + /** Symlinks left out because their target resolves outside `root`. */ + skippedLinks: string[]; + dispose(): void; } /** - * gzip'd ustar of the working tree at `root`. Inside a Git checkout the member - * list is `git ls-files --cached --others --exclude-standard`, so `.gitignore` - * governs and untracked-but-not-ignored files ride along, as they do in v1. - * Outside Git, every regular file and symlink except `.git`/`node_modules`. - * Ordering and mtimes are deterministic, so the same tree packs to the same - * bytes; what Cloud extracts is exactly this list, nothing implicit. + * gzip'd ustar of the working tree at `root`, streamed to a temporary file so + * the peak footprint is one file plus the compressor's window, never the + * whole tree. Inside a Git checkout the member list is `git ls-files --cached + * --others --exclude-standard`, so `.gitignore` governs and untracked-but- + * not-ignored files ride along, as they do in v1. Outside Git, every regular + * file and in-tree symlink except `.git`/`node_modules`. A Git failure other + * than "not a repository" is refused rather than silently widened to the walk: + * an ignored `.env` must never reach Cloud because `git` was missing or the + * index was locked. Executable bits survive; symlinks whose target leaves the + * tree are dropped and reported, since on the host they would point at + * whatever happens to live at that path. Ordering and mtimes are + * deterministic, so the same tree packs to the same bytes. */ -export function packWorkingTree(root: string): PackedTree { +export async function packWorkingTree(root: string): Promise { const absoluteRoot = resolve(root); const files = listTreeFiles(absoluteRoot); - const chunks: Buffer[] = []; + const spool = mkdtempSync(join(tmpdir(), 'flows-sync-')); + const archivePath = join(spool, 'code.tar.gz'); + const dispose = (): void => rmSync(spool, { recursive: true, force: true }); + const members: string[] = []; + const skippedLinks: string[] = []; let bytes = 0; - for (const path of files) { - const absolute = join(absoluteRoot, ...path.split('/')); - const stat = lstatSync(absolute); - if (stat.isSymbolicLink()) { - chunks.push(ustarHeader(path, 0, '2', readlinkSync(absolute))); - continue; + async function* entries(): AsyncGenerator { + for (const path of files) { + const absolute = join(absoluteRoot, ...path.split('/')); + const stat = lstatSync(absolute); + if (stat.isSymbolicLink()) { + const target = readlinkSync(absolute); + const resolved = resolve(dirname(absolute), target); + if (isAbsolute(target) || relative(absoluteRoot, resolved).startsWith('..')) { + skippedLinks.push(path); + continue; + } + members.push(path); + yield ustarHeader(path, 0, '2', 0o777, target); + continue; + } + if (!stat.isFile()) continue; + bytes += stat.size; + if (bytes > MAX_SYNC_BYTES) { + throw new CloudFlowError('sync_too_large', + `The working tree exceeds the ${MAX_SYNC_BYTES}-byte code sync limit; narrow it with .gitignore.`); + } + members.push(path); + yield ustarHeader(path, stat.size, '0', stat.mode & 0o111 ? 0o755 : 0o644); + yield createReadStream(absolute); + const padding = (512 - (stat.size % 512)) % 512; + if (padding) yield Buffer.alloc(padding); } - if (!stat.isFile()) continue; - bytes += stat.size; - if (bytes > MAX_SYNC_BYTES) { - throw new CloudFlowError('sync_too_large', - `The working tree exceeds the ${MAX_SYNC_BYTES}-byte code sync limit; narrow it with .gitignore.`); - } - const content = readFileSync(absolute); - chunks.push(ustarHeader(path, content.length, '0'), content); - const padding = (512 - (content.length % 512)) % 512; - if (padding) chunks.push(Buffer.alloc(padding)); + yield Buffer.alloc(1024); + } + try { + await pipeline(Readable.from(flatten(entries())), createGzip({ level: 6 }), createWriteStream(archivePath)); + } catch (error) { + dispose(); + throw error; + } + return { archivePath, files: members, bytes, skippedLinks, dispose }; +} + +/** Yield file streams chunk by chunk so `Readable.from` never buffers a whole file. */ +async function* flatten(source: AsyncGenerator): AsyncGenerator { + for await (const part of source) { + if (Buffer.isBuffer(part)) { yield part; continue; } + for await (const chunk of part as AsyncIterable) yield chunk; } - chunks.push(Buffer.alloc(1024)); - return { tarball: gzipSync(Buffer.concat(chunks), { level: 6 }), files, bytes }; } function listTreeFiles(root: string): string[] { - const git = spawnSync('git', ['-C', root, 'ls-files', '-z', '--cached', '--others', '--exclude-standard'], - { encoding: 'utf8', maxBuffer: 64 * 1024 * 1024 }); + const inside = spawnSync('git', ['-C', root, 'rev-parse', '--is-inside-work-tree'], { encoding: 'utf8' }); + const notARepository = inside.error !== undefined + || (inside.status !== 0 && /not a git repository/iu.test(inside.stderr)); let candidates: string[]; - if (git.status === 0) { - candidates = git.stdout.split('\0').filter(Boolean); - } else { + if (notARepository) { + // A `.git` that git itself cannot read is a broken checkout, not a plain + // directory: its .gitignore was meant to apply, so walking would upload + // exactly what it excluded. + if (existsSync(join(root, '.git'))) { + throw new CloudFlowError('sync_unsupported', + `${root} has a .git entry but git cannot read it (${summarize(inside.stderr ?? inside.error?.message)}); ` + + 'refusing to upload a checkout without its .gitignore.'); + } candidates = []; walk(root, root, candidates); + } else { + if (inside.status !== 0 || inside.stdout.trim() !== 'true') { + throw new CloudFlowError('sync_unsupported', + `git could not describe ${root} (${summarize(inside.stderr)}); refusing to guess which files to upload.`); + } + const git = spawnSync('git', ['-C', root, 'ls-files', '-z', '--cached', '--others', '--exclude-standard'], + { encoding: 'utf8', maxBuffer: 256 * 1024 * 1024 }); + if (git.status !== 0) { + throw new CloudFlowError('sync_unsupported', + `git ls-files failed in ${root} (${summarize(git.stderr)}); refusing to upload without .gitignore.`); + } + candidates = git.stdout.split('\0').filter(Boolean); } const files = new Set(); for (const candidate of candidates) { @@ -120,6 +181,11 @@ function listTreeFiles(root: string): string[] { return [...files].sort(); } +function summarize(stderr: string | undefined): string { + const line = (stderr ?? '').split('\n').map(l => l.trim()).find(Boolean) ?? 'no diagnostic'; + return line.length > 160 ? `${line.slice(0, 157)}...` : line; +} + function walk(root: string, current: string, out: string[]): void { for (const entry of readdirSync(current, { withFileTypes: true })) { if (ALWAYS_SKIPPED.has(entry.name)) continue; @@ -130,7 +196,7 @@ function walk(root: string, current: string, out: string[]): void { } /** POSIX ustar header. Names beyond the 100+155 split are refused, not truncated. */ -function ustarHeader(path: string, size: number, type: '0' | '2', link = ''): Buffer { +function ustarHeader(path: string, size: number, type: '0' | '2', mode: number, link = ''): Buffer { let name = path; let prefix = ''; if (Buffer.byteLength(name) > 100) { @@ -146,7 +212,7 @@ function ustarHeader(path: string, size: number, type: '0' | '2', link = ''): Bu } const header = Buffer.alloc(512); header.write(name, 0, 100); - header.write(type === '2' ? '0000777' : '0000644', 100, 8); + header.write(mode.toString(8).padStart(7, '0'), 100, 8); header.write('0000000', 108, 8); header.write('0000000', 116, 8); header.write(size.toString(8).padStart(11, '0'), 124, 12); @@ -165,8 +231,11 @@ function ustarHeader(path: string, size: number, type: '0' | '2', link = ''): Bu /** `PUT` the archive through the Cloud API; the run-scoped credential wins when present. */ export async function uploadCloudCode( - prepared: PreparedCloudSync, tarball: Buffer, options: CloudConnectionOptions, + prepared: PreparedCloudSync, packed: PackedTree, options: CloudConnectionOptions, ): Promise { + // The request needs a sized body, so the compressed archive is read back + // once; that is the one copy that has to exist, and it is the small one. + const tarball = await readFile(packed.archivePath); await cloudFetch( `/api/v1/workflows/runs/${encodeURIComponent(prepared.runId)}/storage/${encodeURIComponent(prepared.codeKey)}`, { ...options, requestTimeoutMs: options.requestTimeoutMs ?? 300_000 }, @@ -195,7 +264,20 @@ export async function downloadCloudPatch(runId: string, options: CloudConnection return { patch: payload.patch, hasChanges: payload.hasChanges }; } -/** `git apply --check` then `git apply`; a conflict leaves the tree untouched. */ +/** Every path a unified diff touches, deletions included, in order of first appearance. */ +export function patchedPaths(patch: string): string[] { + const paths: string[] = []; + for (const match of patch.matchAll(/^diff --git a\/(.+?) b\/(.+)$/gmu)) { + for (const path of [match[1]!, match[2]!]) if (!paths.includes(path)) paths.push(path); + } + return paths; +} + +/** + * `git apply --check` then `git apply`; a conflict leaves the tree untouched. + * The patch lands in the working tree uncommitted, so what the run changed is + * reviewed with `git diff` before anything is kept — the same contract as v1. + */ export function applyCloudPatch(root: string, patch: string): void { const args = ['-C', resolve(root), 'apply', '--whitespace=nowarn']; const check = spawnSync('git', [...args, '--check'], { input: patch, encoding: 'utf8' }); diff --git a/packages/sdk/src/index.ts b/packages/sdk/src/index.ts index 497ba773b..c36942182 100644 --- a/packages/sdk/src/index.ts +++ b/packages/sdk/src/index.ts @@ -64,7 +64,7 @@ export { type CloudFlowSource, type RunInCloudOptions, type CloudRunReceipt, type CloudRunState, } from './cloud-run.js'; export { - downloadCloudPatch, applyCloudPatch, packWorkingTree, MAX_SYNC_BYTES, + downloadCloudPatch, applyCloudPatch, packWorkingTree, patchedPaths, MAX_SYNC_BYTES, type CloudPatch, type PackedTree, } from './cloud-sync.js'; export { diff --git a/packages/sdk/tests/check-command-cwd.test.ts b/packages/sdk/tests/check-command-cwd.test.ts new file mode 100644 index 000000000..378210f52 --- /dev/null +++ b/packages/sdk/tests/check-command-cwd.test.ts @@ -0,0 +1,34 @@ +import { chmod, mkdtemp, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { checkFlow } from '../src/cli/check.js'; + +const dirs: string[] = []; +afterEach(async () => { + vi.restoreAllMocks(); + for (const dir of dirs.splice(0)) await rm(dir, { recursive: true, force: true }); +}); + +describe('path-like deterministic commands are probed where the step will run', () => { + it('resolves ./script against the invoking cwd, not the flow file directory', async () => { + const sourceDir = await mkdtemp(join(tmpdir(), 'check-cwd-source-')); + const runDir = await mkdtemp(join(tmpdir(), 'check-cwd-run-')); + dirs.push(sourceDir, runDir); + const flow = join(sourceDir, 'flow.yaml'); + await writeFile(flow, JSON.stringify({ + version: '0.1.0', name: 'cwd-proof', + steps: [{ id: 'exec', type: 'deterministic', command: './run.sh' }], + })); + await writeFile(join(runDir, 'run.sh'), '#!/bin/sh\necho ok\n'); + await chmod(join(runDir, 'run.sh'), 0o755); + + vi.spyOn(process, 'cwd').mockReturnValue(runDir); + expect(checkFlow(flow).report.ok).toBe(true); + + vi.spyOn(process, 'cwd').mockReturnValue(sourceDir); + const refused = checkFlow(flow).report; + expect(refused.ok).toBe(false); + expect(refused.diagnostics).toContainEqual(expect.objectContaining({ kind: 'command_missing', stepId: 'exec' })); + }); +}); diff --git a/packages/sdk/tests/cloud-sync.test.ts b/packages/sdk/tests/cloud-sync.test.ts index 1a39c824d..b10adb409 100644 --- a/packages/sdk/tests/cloud-sync.test.ts +++ b/packages/sdk/tests/cloud-sync.test.ts @@ -1,12 +1,14 @@ import { execFileSync } from 'node:child_process'; -import { mkdir, mkdtemp, readFile, readdir, rm, symlink, writeFile } from 'node:fs/promises'; +import { chmod, mkdir, mkdtemp, readFile, readdir, rm, stat, symlink, writeFile } from 'node:fs/promises'; +import { execFile } from 'node:child_process'; +import { promisify } from 'node:util'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { gunzipSync } from 'node:zlib'; import { afterEach, describe, expect, it, vi } from 'vitest'; import { runCli } from '../src/cli.js'; import { runInCloud } from '../src/cloud-run.js'; -import { applyCloudPatch, packWorkingTree, prepareCloudSync } from '../src/cloud-sync.js'; +import { applyCloudPatch, packWorkingTree, patchedPaths, prepareCloudSync } from '../src/cloud-sync.js'; const dirs: string[] = []; afterEach(async () => { @@ -63,12 +65,22 @@ const PREPARED = { async function untar(tarball: Buffer): Promise<{ dir: string; members: string[] }> { const dir = await tempDir('cloud-sync-extract-'); await writeFile(join(dir, 'code.tar.gz'), tarball); - const members = execFileSync('tar', ['-tzf', join(dir, 'code.tar.gz')], { encoding: 'utf8' }).trim().split('\n').sort(); + const members = execFileSync('tar', ['-tzf', join(dir, 'code.tar.gz')], { encoding: 'utf8' }).trim().split('\n').filter(Boolean).sort(); await mkdir(join(dir, 'out')); execFileSync('tar', ['-xzf', join(dir, 'code.tar.gz'), '-C', join(dir, 'out')]); return { dir: join(dir, 'out'), members }; } +/** Pack, read the spooled archive once, and dispose it, as the upload path does. */ +async function pack(root: string) { + const packed = await packWorkingTree(root); + try { + return { ...packed, tarball: await readFile(packed.archivePath) }; + } finally { + packed.dispose(); + } +} + describe('packWorkingTree', () => { it('packs a Git checkout by ls-files semantics: tracked plus untracked, never ignored, .git or node_modules', async () => { const root = await tempDir('cloud-sync-git-'); @@ -82,39 +94,85 @@ describe('packWorkingTree', () => { await writeFile(join(root, 'dist/build.js'), 'ignored\n'); await writeFile(join(root, 'debug.log'), 'ignored\n'); await symlink('src/deep/a.ts', join(root, 'link.ts')); - git(root, 'add', '.gitignore', 'src'); + await symlink('/etc/hosts', join(root, 'escape-abs')); + await symlink('../../outside', join(root, 'escape-rel')); + await writeFile(join(root, 'run.sh'), '#!/bin/sh\necho ok\n'); + await chmod(join(root, 'run.sh'), 0o755); + git(root, 'add', '.gitignore', 'src', 'run.sh'); git(root, 'commit', '-q', '-m', 'baseline'); await writeFile(join(root, 'untracked.md'), 'rides along\n'); - const packed = packWorkingTree(root); - expect(packed.files).toEqual(['.gitignore', 'link.ts', 'src/deep/a.ts', 'untracked.md']); - expect(packed.bytes).toBe(Buffer.byteLength('dist/\n*.log\n') + Buffer.byteLength('export const a = 1;\n') + Buffer.byteLength('rides along\n')); + const packed = await pack(root); + expect(packed.files).toEqual(['.gitignore', 'link.ts', 'run.sh', 'src/deep/a.ts', 'untracked.md']); + expect(packed.skippedLinks).toEqual(['escape-abs', 'escape-rel']); + expect(packed.bytes).toBe(Buffer.byteLength('dist/\n*.log\n') + Buffer.byteLength('export const a = 1;\n') + + Buffer.byteLength('#!/bin/sh\necho ok\n') + Buffer.byteLength('rides along\n')); const { dir, members } = await untar(packed.tarball); - expect(members).toEqual(['.gitignore', 'link.ts', 'src/deep/a.ts', 'untracked.md']); + expect(members).toEqual(['.gitignore', 'link.ts', 'run.sh', 'src/deep/a.ts', 'untracked.md']); expect(await readFile(join(dir, 'src/deep/a.ts'), 'utf8')).toBe('export const a = 1;\n'); expect(await readFile(join(dir, 'link.ts'), 'utf8')).toBe('export const a = 1;\n'); + // Executable bits survive the round trip; plain files stay 0644. + expect((await stat(join(dir, 'run.sh'))).mode & 0o111).not.toBe(0); + expect((await stat(join(dir, 'src/deep/a.ts'))).mode & 0o111).toBe(0); + expect((await promisify(execFile)(join(dir, 'run.sh'))).stdout).toBe('ok\n'); expect(await readdir(dir)).not.toContain('node_modules'); + expect(await readdir(dir)).not.toContain('escape-abs'); + }); + + it('refuses to widen to a walk when git fails inside a checkout', async () => { + const root = await tempDir('cloud-sync-gitfail-'); + git(root, 'init', '-q'); + await writeFile(join(root, '.gitignore'), '.env\n'); + await writeFile(join(root, '.env'), 'SECRET=1\n'); + await writeFile(join(root, 'ok.txt'), 'ok\n'); + // A locked index makes ls-files fail while rev-parse still says "inside". + await writeFile(join(root, '.git/index.lock'), ''); + const broken = await tempDir('cloud-sync-gitbroken-'); + await mkdir(join(broken, '.git')); + await writeFile(join(broken, '.git/HEAD'), 'garbage\n'); + await writeFile(join(broken, '.env'), 'SECRET=1\n'); + // A corrupt .git is refused outright. A locked index is refused where git + // refuses ls-files; a git that tolerates the lock must still honour .gitignore. + await expect(pack(broken)).rejects.toMatchObject({ code: 'sync_unsupported' }); + let error: unknown; + try { await pack(root); } catch (e) { error = e; } + if (error === undefined) expect((await pack(root)).files).not.toContain('.env'); + else { + expect(error).toMatchObject({ code: 'sync_unsupported' }); + expect(String((error as Error).message)).not.toContain('SECRET'); + } }); - it('packs a plain directory by walking it, skipping only .git and node_modules', async () => { + it('packs a plain directory by walking it, skipping node_modules and nested .git entries', async () => { const root = await tempDir('cloud-sync-plain-'); await mkdir(join(root, 'lib')); await mkdir(join(root, 'node_modules')); - await mkdir(join(root, '.git')); + // A vendored checkout inside a plain directory: its .git is skipped, its files ride along. + await mkdir(join(root, 'vendor/.git'), { recursive: true }); await writeFile(join(root, 'lib/x.txt'), 'x'); await writeFile(join(root, 'node_modules/y.txt'), 'y'); - await writeFile(join(root, '.git/HEAD'), 'ref'); + await writeFile(join(root, 'vendor/.git/HEAD'), 'ref'); + await writeFile(join(root, 'vendor/v.txt'), 'v'); await writeFile(join(root, 'top.txt'), 'top'); - const packed = packWorkingTree(root); - expect(packed.files).toEqual(['lib/x.txt', 'top.txt']); - expect((await untar(packed.tarball)).members).toEqual(['lib/x.txt', 'top.txt']); + const packed = await pack(root); + expect(packed.files).toEqual(['lib/x.txt', 'top.txt', 'vendor/v.txt']); + expect((await untar(packed.tarball)).members).toEqual(['lib/x.txt', 'top.txt', 'vendor/v.txt']); }); it('is deterministic for the same tree', async () => { const root = await tempDir('cloud-sync-det-'); await writeFile(join(root, 'a'), 'a'); - expect(packWorkingTree(root).tarball.equals(packWorkingTree(root).tarball)).toBe(true); + expect((await pack(root)).tarball.equals((await pack(root)).tarball)).toBe(true); + }); + + it('disposes its spool file after packing', async () => { + const root = await tempDir('cloud-sync-spool-'); + await writeFile(join(root, 'a'), 'a'); + const packed = await packWorkingTree(root); + expect((await stat(packed.archivePath)).size).toBeGreaterThan(0); + packed.dispose(); + await expect(stat(packed.archivePath)).rejects.toMatchObject({ code: 'ENOENT' }); }); }); @@ -239,14 +297,40 @@ describe('flows run --cloud --sync-code / flows sync', () => { await writeFile(join(root, 'a.txt'), 'one\n'); git(root, 'add', 'a.txt'); git(root, 'commit', '-q', '-m', 'baseline'); + await writeFile(join(root, 'gone.txt'), 'bye\n'); + git(root, 'add', 'gone.txt'); + git(root, 'commit', '-q', '-m', 'add gone'); const patch = 'diff --git a/a.txt b/a.txt\n--- a/a.txt\n+++ b/a.txt\n@@ -1 +1,2 @@\n one\n+two\n' - + 'diff --git a/b.txt b/b.txt\nnew file mode 100644\n--- /dev/null\n+++ b/b.txt\n@@ -0,0 +1 @@\n+fresh\n'; + + 'diff --git a/b.txt b/b.txt\nnew file mode 100644\n--- /dev/null\n+++ b/b.txt\n@@ -0,0 +1 @@\n+fresh\n' + + 'diff --git a/gone.txt b/gone.txt\ndeleted file mode 100644\n--- a/gone.txt\n+++ /dev/null\n@@ -1 +0,0 @@\n-bye\n'; cloud({ '/api/v1/workflows/runs/done-run/patch': () => ({ patch, hasChanges: true }) }); const output: string[] = []; expect(await runCli(['sync', '--json', '--dir', root, 'done-run'], { stdout: line => output.push(line), stderr: line => output.push(line) })).toBe(0); - expect(JSON.parse(output[0]!)).toEqual({ ok: true, runId: 'done-run', hasChanges: true, applied: true, files: ['a.txt', 'b.txt'] }); + expect(JSON.parse(output[0]!)).toEqual({ ok: true, runId: 'done-run', hasChanges: true, applied: true, files: ['a.txt', 'b.txt', 'gone.txt'] }); expect(await readFile(join(root, 'a.txt'), 'utf8')).toBe('one\ntwo\n'); expect(await readFile(join(root, 'b.txt'), 'utf8')).toBe('fresh\n'); + await expect(stat(join(root, 'gone.txt'))).rejects.toMatchObject({ code: 'ENOENT' }); + expect(patchedPaths(patch)).toEqual(['a.txt', 'b.txt', 'gone.txt']); + }); + + it('classifies an interruption during upload as pre-admission, not admission_unknown', async () => { + const root = await tempDir('cloud-sync-abort-'); + await writeFile(join(root, 'flow.yaml'), JSON.stringify({ + version: '0.1.0', name: 'synced', steps: [{ id: 'ls', type: 'deterministic', command: 'ls' }], + })); + const { calls } = cloud({ + '/api/v1/workflows/prepare': () => { process.emit('SIGINT'); return PREPARED; }, + '/api/v1/workflows/runs/prepared-run/storage/code.tar.gz': () => ({ ok: true }), + '/api/v1/workflows/run': () => ({ runId: 'prepared-run', status: 'pending' }), + }); + const cwd = vi.spyOn(process, 'cwd').mockReturnValue(root); + const output: string[] = []; + const code = await runCli(['run', '--cloud', '--sync-code', '--json', join(root, 'flow.yaml')], + { stdout: line => output.push(line), stderr: () => {} }); + cwd.mockRestore(); + expect(code).toBe(1); + expect(JSON.parse(output[0]!)).toMatchObject({ ok: false, code: 'submission_aborted' }); + expect(calls.map(c => c.path)).not.toContain('/api/v1/workflows/run'); }); it('reports no changes without touching the tree', async () => { From 966bacd5e5534f5a07d0a78a6657ee18dea0c686 Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 17 Sep 2026 12:39:07 -0700 Subject: [PATCH 4/4] =?UTF-8?q?fix(sdk):=20second=20review=20round=20?= =?UTF-8?q?=E2=80=94=20login=20URL=20binding,=20git-less=20subtrees,=20dot?= =?UTF-8?q?-named=20links,=20local=20refusals?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - A login-store token is now bound to the deployment that issued it: an explicit `apiUrl`/`FLOWS_CLOUD_URL` naming another base refuses with `configuration` instead of sending the token there (CodeRabbit). - `listTreeFiles` looks for the nearest `.git` in `root` or any ancestor, so a subdirectory of a broken checkout, or a tree with no `git` on PATH but a `.git` above it, refuses rather than walking past `.gitignore` (Cursor). A tree with no `.git` anywhere still walks without git. - Symlink escape detection compares whole path segments, so an in-tree target like `..cache/file` is kept (CodeRabbit). - `flows deploy` reports a missing/unreadable source as `invalid_input` and an unloadable one as `unsupported_source`, exit 2, before any HTTP (CodeRabbit). - The deployment example in docs/CLOUD.md is copy-pasteable; optional flags are listed beside it (CodeRabbit). Co-Authored-By: Claude Opus 5 (1M context) --- docs/CLOUD.md | 5 +++- packages/sdk/src/cloud-deploy.ts | 21 ++++++++++++-- packages/sdk/src/cloud-http.ts | 18 +++++++++++- packages/sdk/src/cloud-sync.ts | 37 ++++++++++++++++++++----- packages/sdk/tests/cloud-deploy.test.ts | 25 +++++++++++++++++ packages/sdk/tests/cloud-sync.test.ts | 31 +++++++++++++++++++++ 6 files changed, 125 insertions(+), 12 deletions(-) diff --git a/docs/CLOUD.md b/docs/CLOUD.md index 80b7793a9..3f26353bf 100644 --- a/docs/CLOUD.md +++ b/docs/CLOUD.md @@ -116,11 +116,14 @@ produces; a deployment (CI) token gets `session_required` and the CLI says so. flows deploy issue-triage.flow.ts \ --repo AgentWorkforce/flows \ --on github:labels=agent \ - --approver khaliqgant [--agents claude,codex] [--name "Issue triage"] [--draft] + --approver khaliqgant flows deployments flows undeploy ``` +Optional flags: `--agents claude,codex`, `--name "Issue triage"`, `--draft`, +`--json`, and further `--on` sources. + `flows deploy ` is the CLI form of the agentrelay.com onboarding's deploy wizard: `POST /api/v1/flows/deploy` stores one self-contained authored source and creates a proactive listener whose watch rules match the chosen diff --git a/packages/sdk/src/cloud-deploy.ts b/packages/sdk/src/cloud-deploy.ts index 39988f833..78ee04ca4 100644 --- a/packages/sdk/src/cloud-deploy.ts +++ b/packages/sdk/src/cloud-deploy.ts @@ -114,7 +114,13 @@ export async function deployToCloud( if (!/\.flow\.ts$/iu.test(input.path)) { throw new CloudFlowError('unsupported_source', 'flows deploy takes one authored .flow.ts source.'); } - const bytes = await readFile(input.path); + let bytes: Buffer; + try { + bytes = await readFile(input.path); + } catch (error) { + throw new CloudFlowError('invalid_input', + `Cannot read ${input.path}: ${(error as NodeJS.ErrnoException).code ?? (error as Error).message}.`); + } const source = bytes.toString('utf8'); if (!bytes.length || Buffer.from(source, 'utf8').compare(bytes) !== 0) { throw new CloudFlowError('invalid_input', 'Authored source must be nonempty, lossless UTF-8.'); @@ -122,8 +128,17 @@ export async function deployToCloud( if (bytes.length > MAX_SOURCE_BYTES) { throw new CloudFlowError('invalid_input', `Authored source exceeds Cloud's ${MAX_SOURCE_BYTES}-byte deploy limit.`); } - const loaded = await loadAuthoredFlow(input.path); - const definition = loaded.getDefinition(loaded.handle); + let loaded: Awaited>; + let definition: ReturnType; + try { + loaded = await loadAuthoredFlow(input.path); + definition = loaded.getDefinition(loaded.handle); + } catch (error) { + if (error instanceof CloudFlowError) throw error; + // A source that does not load is an authoring problem, refused before HTTP. + throw new CloudFlowError('unsupported_source', + `${input.path} is not a loadable authored flow: ${error instanceof Error ? error.message : String(error)}`); + } if (loaded.graph.length !== 1) { throw new CloudFlowError('unsupported_source', 'Cloud deploys one self-contained .flow.ts source without use dependencies.'); diff --git a/packages/sdk/src/cloud-http.ts b/packages/sdk/src/cloud-http.ts index 59b61d979..9c875f98b 100644 --- a/packages/sdk/src/cloud-http.ts +++ b/packages/sdk/src/cloud-http.ts @@ -92,7 +92,23 @@ export function cloudConnection(options: CloudConnectionOptions): { baseUrl: str || url.username || url.password || url.search || url.hash || !/^\/[A-Za-z0-9/_-]*$/u.test(url.pathname)) { throw new CloudFlowError('configuration', 'Cloud URL must use HTTPS and a plain base path.'); } - return { baseUrl: `${url.origin}${url.pathname.replace(/\/+$/u, '')}`, token }; + const baseUrl = `${url.origin}${url.pathname.replace(/\/+$/u, '')}`; + // A login-store token is bound to the deployment that issued it. An explicit + // URL that names another deployment gets no token at all — set + // FLOWS_CLOUD_TOKEN for that deployment instead. + if (loginApiUrl !== undefined) { + let issued: string | undefined; + try { + const login = new URL(loginApiUrl); + issued = `${login.origin}${login.pathname.replace(/\/+$/u, '')}`; + } catch { issued = undefined; } + if (issued !== baseUrl) { + throw new CloudFlowError('configuration', + `The agent-relay cloud login was issued for ${loginApiUrl}, not ${baseUrl}. ` + + 'Set FLOWS_CLOUD_TOKEN for that deployment, or unset FLOWS_CLOUD_URL to use the login.'); + } + } + return { baseUrl, token }; } export async function cloudRequest( diff --git a/packages/sdk/src/cloud-sync.ts b/packages/sdk/src/cloud-sync.ts index 783ba69ee..f48425a91 100644 --- a/packages/sdk/src/cloud-sync.ts +++ b/packages/sdk/src/cloud-sync.ts @@ -100,8 +100,8 @@ export async function packWorkingTree(root: string): Promise { const stat = lstatSync(absolute); if (stat.isSymbolicLink()) { const target = readlinkSync(absolute); - const resolved = resolve(dirname(absolute), target); - if (isAbsolute(target) || relative(absoluteRoot, resolved).startsWith('..')) { + const relativeTarget = relative(absoluteRoot, resolve(dirname(absolute), target)); + if (isAbsolute(target) || relativeTarget === '..' || relativeTarget.startsWith(`..${sep}`)) { skippedLinks.push(path); continue; } @@ -140,18 +140,41 @@ async function* flatten(source: AsyncGenerator): } } +/** The nearest `.git` entry at `root` or any ancestor — the checkout a .gitignore would belong to. */ +function nearestGitEntry(root: string): string | undefined { + let directory = root; + while (true) { + if (existsSync(join(directory, '.git'))) return join(directory, '.git'); + const parent = dirname(directory); + if (parent === directory) return undefined; + directory = parent; + } +} + function listTreeFiles(root: string): string[] { const inside = spawnSync('git', ['-C', root, 'rev-parse', '--is-inside-work-tree'], { encoding: 'utf8' }); + if (inside.error !== undefined) { + // Without git there is no way to honour a .gitignore, so a tree that has + // one anywhere above it cannot be described honestly; a tree with none is + // a plain directory and can be walked. + const entry = nearestGitEntry(root); + if (entry !== undefined) { + throw new CloudFlowError('sync_unsupported', + `git is not available (${summarize(inside.error.message)}) but ${entry} exists; ` + + 'refusing to upload a checkout without its .gitignore.'); + } + } const notARepository = inside.error !== undefined || (inside.status !== 0 && /not a git repository/iu.test(inside.stderr)); let candidates: string[]; if (notARepository) { - // A `.git` that git itself cannot read is a broken checkout, not a plain - // directory: its .gitignore was meant to apply, so walking would upload - // exactly what it excluded. - if (existsSync(join(root, '.git'))) { + // A `.git` that git itself cannot read — here or in a parent — is a broken + // checkout, not a plain directory: its .gitignore was meant to apply, so + // walking would upload exactly what it excluded. + const entry = nearestGitEntry(root); + if (entry !== undefined) { throw new CloudFlowError('sync_unsupported', - `${root} has a .git entry but git cannot read it (${summarize(inside.stderr ?? inside.error?.message)}); ` + `${entry} exists but git cannot read it (${summarize(inside.stderr ?? inside.error?.message)}); ` + 'refusing to upload a checkout without its .gitignore.'); } candidates = []; diff --git a/packages/sdk/tests/cloud-deploy.test.ts b/packages/sdk/tests/cloud-deploy.test.ts index 7ba94402e..ff4b9f0be 100644 --- a/packages/sdk/tests/cloud-deploy.test.ts +++ b/packages/sdk/tests/cloud-deploy.test.ts @@ -105,6 +105,22 @@ describe('deployToCloud', () => { expect(deployment.sourceSha256).toMatch(/^[a-f0-9]{64}$/u); }); + it('reports a missing or unloadable source as an input refusal (exit 2), before HTTP', async () => { + const calls = cloud({}); + const base = { repository: { owner: 'o', name: 'r' }, sources: [parseTriggerSource('github')], approver: 'k' }; + await expect(deployToCloud({ ...base, path: join(await tempDir('cloud-deploy-missing-'), 'nope.flow.ts') })) + .rejects.toMatchObject({ code: 'invalid_input' }); + const dir = await tempDir('cloud-deploy-broken-'); + await symlink(join(process.cwd(), 'node_modules'), join(dir, 'node_modules'), 'dir'); + await writeFile(join(dir, 'broken.flow.ts'), 'export default 42;\n'); + await expect(deployToCloud({ ...base, path: join(dir, 'broken.flow.ts') })) + .rejects.toMatchObject({ code: 'unsupported_source' }); + const errors: string[] = []; + expect(await runCli(['deploy', join(dir, 'broken.flow.ts'), '--repo', 'o/r', '--on', 'github', '--approver', 'k'], + { stdout: () => {}, stderr: line => errors.push(line) })).toBe(2); + expect(calls).toHaveLength(0); + }); + it('refuses non-authored sources, empty sources, duplicate providers and a blank approver before HTTP', async () => { const path = await authoredFlow(); const calls = cloud({}); @@ -240,6 +256,15 @@ describe('agent-relay cloud login fallback', () => { expect(cloudConnection({})).toEqual({ baseUrl: 'https://login.example/cloud', token: 'login-token' }); }); + it('never sends the login token to a deployment other than the one that issued it', async () => { + await loginStore({ apiUrl: 'https://login.example/cloud', accessToken: 'login-token' }); + vi.stubEnv('FLOWS_CLOUD_URL', 'https://other.example/cloud'); + expect(() => cloudConnection({})).toThrow(expect.objectContaining({ code: 'configuration' })); + expect(() => cloudConnection({ apiUrl: 'https://other.example/cloud' })).toThrow(/issued for https:\/\/login\.example\/cloud/u); + // The same deployment spelled with a trailing slash is still the same deployment. + expect(cloudConnection({ apiUrl: 'https://login.example/cloud/' })).toEqual({ baseUrl: 'https://login.example/cloud', token: 'login-token' }); + }); + it('lets FLOWS_CLOUD_TOKEN win over the login store', async () => { await loginStore({ apiUrl: 'https://login.example/cloud', accessToken: 'login-token' }); vi.stubEnv('FLOWS_CLOUD_TOKEN', 'explicit'); diff --git a/packages/sdk/tests/cloud-sync.test.ts b/packages/sdk/tests/cloud-sync.test.ts index b10adb409..54c25ecdf 100644 --- a/packages/sdk/tests/cloud-sync.test.ts +++ b/packages/sdk/tests/cloud-sync.test.ts @@ -144,6 +144,37 @@ describe('packWorkingTree', () => { } }); + it('keeps in-tree symlinks whose names merely start with dots, and refuses git-less subdirectories of a checkout', async () => { + const root = await tempDir('cloud-sync-dotlink-'); + await mkdir(join(root, '..cache')); + await writeFile(join(root, '..cache/file'), 'cached'); + await symlink('..cache/file', join(root, 'dotlink')); + const packed = await pack(root); + expect(packed.files).toEqual(['..cache/file', 'dotlink']); + expect(packed.skippedLinks).toEqual([]); + + // A subdirectory of a broken checkout has no .git of its own; the parent's + // unreadable one must still stop the walk from uploading ignored files. + const repo = await tempDir('cloud-sync-broken-parent-'); + await mkdir(join(repo, '.git')); + await writeFile(join(repo, '.git/HEAD'), 'garbage\n'); + await mkdir(join(repo, 'sub')); + await writeFile(join(repo, 'sub/.env'), 'SECRET=1\n'); + await expect(pack(join(repo, 'sub'))).rejects.toMatchObject({ code: 'sync_unsupported' }); + + // Without git at all, a tree with a .git anywhere above is refused; a plain one walks. + const path = process.env['PATH']; + vi.stubEnv('PATH', '/nonexistent'); + try { + const plain = await tempDir('cloud-sync-nogit-plain-'); + await writeFile(join(plain, 'a.txt'), 'a'); + expect((await pack(plain)).files).toEqual(['a.txt']); + await expect(pack(join(repo, 'sub'))).rejects.toMatchObject({ code: 'sync_unsupported' }); + } finally { + vi.stubEnv('PATH', path ?? ''); + } + }); + it('packs a plain directory by walking it, skipping node_modules and nested .git entries', async () => { const root = await tempDir('cloud-sync-plain-'); await mkdir(join(root, 'lib'));