diff --git a/docs/CLOUD.md b/docs/CLOUD.md index ff834127..9571cef3 100644 --- a/docs/CLOUD.md +++ b/docs/CLOUD.md @@ -494,7 +494,8 @@ Every refusal is one `REFUSED [code] message` line naming what to do next: | code | when | | --- | --- | | `cloud_auth_missing` | no credential anywhere; names `agent-relay cloud login` | -| `cloud_auth_expired` | the stored login expired; refused before any request | +| `cloud_configuration` (`auth_store_unwritable`) | cannot persist refreshed credentials; check the named path, filesystem error, permissions and free space | +| `cloud_auth_expired` | the stored login expired and cannot be refreshed | | `cloud_auth_rejected` | Cloud answered 401: the token is unknown or revoked | | `cloud_forbidden` | 403: authenticated, but not allowed to read that run or log | | `cloud_run_not_found` | 404: no such run for this credential; points at `flows runs` | @@ -524,7 +525,11 @@ 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. +access token is renewed using the stored refresh token, and both rotated tokens +are written atomically before use. Missing or expired refresh tokens and rejected +refresh credentials require re-login; network and server failures retain their +transport and HTTP classifications. Explicit credentials never read or write the +store. Renewal is triggered by the stored expiry, not by a request returning 401. Running and syncing work with either kind of token. Deploying, listing and removing listeners need the interactive `cli:auth` credential the login diff --git a/packages/sdk/src/agent-artifacts.ts b/packages/sdk/src/agent-artifacts.ts index ba98f915..cf248fcb 100644 --- a/packages/sdk/src/agent-artifacts.ts +++ b/packages/sdk/src/agent-artifacts.ts @@ -7,6 +7,12 @@ function isEnoent(error: unknown): boolean { return error instanceof Error && (error as NodeJS.ErrnoException).code === 'ENOENT'; } +/** A file present on disk whose bytes this process is not permitted to read. */ +function isUnreadable(error: unknown): boolean { + const code = error instanceof Error ? (error as NodeJS.ErrnoException).code : undefined; + return code === 'EACCES' || code === 'EPERM'; +} + /** * A recursive snapshot of every regular file under `dir`, keyed by its * `dir`-relative POSIX path, valued by a content signature (size + sha256). @@ -27,10 +33,12 @@ function isEnoent(error: unknown): boolean { * Only regular files are signed; symlinks are not followed and directories are * descended, not recorded. A missing `dir` (an agent step whose cwd does not * exist yet) yields an empty snapshot rather than throwing. Only a vanished - * path (`ENOENT`) is ever swallowed this way; any other filesystem error - * (permissions, `ENOTDIR`, `EISDIR`, ...) propagates, because a step whose - * artifact scan silently dropped files it could not read must not report a - * successful, incomplete `artifacts` list as if it were the truth. + * path (`ENOENT`) is ever swallowed this way; a file whose bytes cannot be read + * is recorded by size and reason instead of being skipped, and any other + * filesystem error (`ENOTDIR`, `EIO`, an unreadable *directory*, ...) + * propagates, because a step whose artifact scan silently dropped files it + * could not read must not report a successful, incomplete `artifacts` list as + * if it were the truth. */ export async function snapshotWorkspaceFiles(dir: string): Promise> { const out = new Map(); @@ -68,12 +76,27 @@ async function walk(root: string, current: string, out: Map): Pr bytes = await readFile(path); } catch (error) { if (isEnoent(error)) continue; - throw error; + if (!isUnreadable(error)) throw error; + // A file whose bytes are unreadable is still accounted for — by size + // and reason, which can never collide with a content hash — rather + // than dropped or made to fail the scan. `bun build --compile` creates + // its output in its cwd with `O_CREAT|O_EXCL` and mode 000, writes the + // (tens of megabytes) executable into it, then renames it onto the + // --outfile; `bundle-typescript.ts` runs exactly that with cwd set to + // the flow's own directory. So any tree a build is running in holds an + // unreadable file for seconds at a time, and failing closed here lets + // a neighbouring process fail an agent step that has done its work. + out.set(signedPath(root, path), `${info.size}:unreadable:${(error as NodeJS.ErrnoException).code}`); + continue; } - out.set(relative(root, path).split(sep).join('/'), `${info.size}:${createHash('sha256').update(bytes).digest('hex')}`); + out.set(signedPath(root, path), `${info.size}:${createHash('sha256').update(bytes).digest('hex')}`); } } +function signedPath(root: string, path: string): string { + return relative(root, path).split(sep).join('/'); +} + /** `dir`-relative paths present in `after` that are new or changed since `before`, sorted. */ export function diffWorkspaceFiles(before: Map, after: Map): string[] { const changed: string[] = []; diff --git a/packages/sdk/src/cli/cloud-mirror-session.ts b/packages/sdk/src/cli/cloud-mirror-session.ts index 2c73c11f..aaa98b95 100644 --- a/packages/sdk/src/cli/cloud-mirror-session.ts +++ b/packages/sdk/src/cli/cloud-mirror-session.ts @@ -25,7 +25,7 @@ import { readFile } from 'node:fs/promises'; import type { CliIo } from '../cli.js'; import { canonicalize } from '../canonical.js'; -import { cloudConnection, CloudFlowError } from '../cloud-http.js'; +import { resolveCloudConnection, CloudFlowError } from '../cloud-http.js'; import { recordMirroredRun, readMirroredRun } from '../cloud-mirror-ledger.js'; import { createRunMirror, readJournalEvents, type RunMirror } from '../cloud-mirror.js'; import { MirrorClient, registerLocalRun, type MirrorRunSource } from '../cloud-mirror-transport.js'; @@ -231,7 +231,8 @@ export function mirrorSourceFromJournal(dataDir: string): (runId: string) => Pro // the deployment this invocation will actually register with, so a run // mirrored to staging never claims to continue an id that means something // else in production. - const resumedFromRunId = await readMirroredRun(dataDir, runId, cloudConnection({}).baseUrl) + const { baseUrl } = await resolveCloudConnection({}); + const resumedFromRunId = await readMirroredRun(dataDir, runId, baseUrl) .catch(() => undefined); // Through the retrying reader: a resume reads this journal while the // daemon is writing to it, and a single `journal_busy` used to abandon diff --git a/packages/sdk/src/cloud-auth-store.ts b/packages/sdk/src/cloud-auth-store.ts new file mode 100644 index 00000000..21dd90fc --- /dev/null +++ b/packages/sdk/src/cloud-auth-store.ts @@ -0,0 +1,185 @@ +import { randomUUID } from 'node:crypto'; +import { readFileSync } from 'node:fs'; +import fs from 'node:fs/promises'; +import { homedir } from 'node:os'; +import { basename, dirname, join } from 'node:path'; +import { setTimeout as delay } from 'node:timers/promises'; + +// @agent-relay/cloud@12.4.1: types.js AUTH_FILE_PATH / DEFAULT_REFRESH_TIMEOUT_MS; +// auth.js writeStoredAuth, AUTH_LOCK_*, requestStoredAuthRefresh. Keep together +// for contract updates. flows additionally honours AGENT_RELAY_HOME on writes. +const AUTH_DIRECTORY = '.agentworkforce/relay'; +const AUTH_FILE = 'cloud-auth.json'; +const DIRECTORY_MODE = 0o700; +const FILE_MODE = 0o600; +const LOCK_RETRY_MS = 50; +// Interop constant: changing this can steal a live relay lock or delay recovery. +const LOCK_STALE_MS = 30_000; +const REFRESH_TIMEOUT_MS = 10_000; +const REFRESH_ROUTE = '/api/v1/auth/token/refresh'; +// Deliberately shorter than relay's 30s acquire timeout for interactive reads. +const LOCK_ACQUIRE_MS = 5_000; + +export interface AgentRelayCloudLogin extends Record { + apiUrl: string; + accessToken: string; + accessTokenExpiresAt?: string; + refreshToken?: string; + refreshTokenExpiresAt?: string; +} + +export function agentRelayCloudAuthPath(env: NodeJS.ProcessEnv = process.env): string { + return join(env['AGENT_RELAY_HOME'] ?? join(homedir(), AUTH_DIRECTORY), AUTH_FILE); +} + +function record(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value); +} + +/** Explicit credentials bypass this store; expired fallback logins may be refreshed. */ +export function readAgentRelayCloudLogin( + env: NodeJS.ProcessEnv = process.env, + read: (path: string) => string = path => readFileSync(path, 'utf8'), +): AgentRelayCloudLogin | undefined { + try { + const parsed: unknown = JSON.parse(read(agentRelayCloudAuthPath(env))); + if (!record(parsed) || typeof parsed.apiUrl !== 'string' || typeof parsed.accessToken !== 'string' + || !parsed.accessToken.trim()) return undefined; + const login = { ...parsed } as AgentRelayCloudLogin; + for (const key of ['accessTokenExpiresAt', 'refreshToken', 'refreshTokenExpiresAt'] as const) { + if (typeof parsed[key] !== 'string') delete login[key]; + } + return login; + } catch { return undefined; } +} + +export function expiredCloudLogin(login: AgentRelayCloudLogin, now = Date.now()): boolean { + const expiry = Date.parse(login.accessTokenExpiresAt ?? ''); + return Number.isFinite(expiry) && expiry <= now; +} + +export function canRefreshCloudLogin(login: AgentRelayCloudLogin, now = Date.now()): boolean { + return typeof login.refreshToken === 'string' && !!login.refreshToken.trim() + && !(Date.parse(login.refreshTokenExpiresAt ?? '') <= now); +} + +/** Shared transport policy; undefined lets the synchronous resolver keep its refusals. */ +export function normalizeCloudBaseUrl(raw: string): string | undefined { + try { + const url = new URL(raw); + if (url.protocol !== 'https:' || url.username || url.password || url.search || url.hash + || !/^\/[A-Za-z0-9/_-]*$/u.test(url.pathname)) return undefined; + return `${url.origin}${url.pathname.replace(/\/+$/u, '')}`; + } catch { return undefined; } +} + +export type CloudLoginRefresh = + | { kind: 'refreshed' | 'current' } + | { kind: 'refused'; status?: number } + | { kind: 'blocked' } + | { kind: 'transport'; error: unknown } + | { kind: 'unwritable'; error: unknown; path: string }; + +interface RefreshOptions { + signal?: AbortSignal; + timeoutMs?: number; + env?: NodeJS.ProcessEnv; + now?: () => number; + fetch?: typeof globalThis.fetch; + refreshTimeoutMs?: number; + lock?: { retryMs?: number; staleMs?: number; acquireTimeoutMs?: number }; +} + +async function writeLogin(path: string, login: AgentRelayCloudLogin): Promise { + await fs.mkdir(dirname(path), { recursive: true, mode: DIRECTORY_MODE }); + const temporaryPath = join(dirname(path), `.${basename(path)}.${process.pid}.${Date.now()}.${randomUUID()}.tmp`); + try { + await fs.writeFile(temporaryPath, `${JSON.stringify(login, null, 2)}\n`, { encoding: 'utf8', mode: FILE_MODE }); + await fs.chmod(temporaryPath, FILE_MODE); + await fs.rename(temporaryPath, path); + } finally { + await fs.rm(temporaryPath, { force: true }); + } +} + +async function requestRefresh(login: AgentRelayCloudLogin, baseUrl: string, options: RefreshOptions): Promise<{ kind: 'ready'; login: AgentRelayCloudLogin } | CloudLoginRefresh> { + const deadline = AbortSignal.timeout(Math.min(options.timeoutMs ?? REFRESH_TIMEOUT_MS, + options.refreshTimeoutMs ?? REFRESH_TIMEOUT_MS)); + try { + const response = await (options.fetch ?? globalThis.fetch)(`${baseUrl}${REFRESH_ROUTE}`, { + method: 'POST', headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ refreshToken: login.refreshToken }), redirect: 'error', + signal: options.signal ? AbortSignal.any([options.signal, deadline]) : deadline, + }); + if (!response.ok) return { kind: 'refused', status: response.status }; + const payload: unknown = await response.json(); + if (!record(payload) || typeof payload.accessToken !== 'string' || !payload.accessToken.trim() + || /[\r\n]/u.test(payload.accessToken) || /^(?:rk|ot)_live_/u.test(payload.accessToken.trim()) + || typeof payload.refreshToken !== 'string' || !payload.refreshToken.trim() + || typeof payload.accessTokenExpiresAt !== 'string' + || !(Date.parse(payload.accessTokenExpiresAt) > (options.now ?? Date.now)())) return { kind: 'refused' }; + const apiUrl = typeof payload.apiUrl === 'string' && payload.apiUrl.trim() ? payload.apiUrl.trim() : login.apiUrl; + const refreshTokenExpiresAt = typeof payload.refreshTokenExpiresAt === 'string' && payload.refreshTokenExpiresAt.trim() + ? payload.refreshTokenExpiresAt.trim() : login.refreshTokenExpiresAt; + if (normalizeCloudBaseUrl(apiUrl) !== baseUrl + || (refreshTokenExpiresAt !== undefined && !Number.isFinite(Date.parse(refreshTokenExpiresAt)))) return { kind: 'refused' }; + return { kind: 'ready', login: { ...login, apiUrl, accessToken: payload.accessToken, refreshToken: payload.refreshToken, + accessTokenExpiresAt: payload.accessTokenExpiresAt, refreshTokenExpiresAt } }; + } catch (error) { + options.signal?.throwIfAborted(); + if (error instanceof SyntaxError && !deadline.aborted) return { kind: 'refused' }; + return { kind: 'transport', error: deadline.aborted ? deadline.reason : error }; + } +} + +/** Serialize rotations with relay and re-source credentials after acquiring its lock. */ +export async function refreshCloudLogin( + login: AgentRelayCloudLogin, baseUrl: string, options: RefreshOptions = {}, +): Promise { + const path = agentRelayCloudAuthPath(options.env); + const lockPath = `${path}.lock`; + const now = options.now ?? Date.now; + const started = now(); + const waitMs = Math.min(options.timeoutMs ?? LOCK_ACQUIRE_MS, options.lock?.acquireTimeoutMs ?? LOCK_ACQUIRE_MS); + let acquired = false; + try { + options.signal?.throwIfAborted(); + await fs.mkdir(dirname(path), { recursive: true, mode: DIRECTORY_MODE }); + while (!acquired) { + options.signal?.throwIfAborted(); + try { + await fs.mkdir(lockPath, { mode: DIRECTORY_MODE }); + acquired = true; + } catch (error) { + if (!record(error) || error.code !== 'EEXIST') throw error; + try { + if (now() - (await fs.stat(lockPath)).mtimeMs > (options.lock?.staleMs ?? LOCK_STALE_MS)) { + await fs.rm(lockPath, { recursive: true, force: true }); + continue; + } + } catch (error) { + if (!record(error) || error.code !== 'ENOENT') throw error; + } + const remaining = waitMs - (now() - started); + if (remaining <= 0) return { kind: 'blocked' }; + await delay(Math.min(options.lock?.retryMs ?? LOCK_RETRY_MS, remaining), undefined, { signal: options.signal }); + } + } + const latest = readAgentRelayCloudLogin(options.env); + // A replaced/deleted login belongs to the other process; never overwrite it. + if (!latest || latest.apiUrl !== login.apiUrl || !expiredCloudLogin(latest, now())) return { kind: 'current' }; + if (!canRefreshCloudLogin(latest, now())) return { kind: 'current' }; + const result = await requestRefresh(latest, baseUrl, options); + if (result.kind !== 'ready') return result; + await writeLogin(path, result.login); + return { kind: 'refreshed' }; + } catch (error) { + options.signal?.throwIfAborted(); + return { kind: 'unwritable', error, path }; + } finally { + if (acquired) { + try { await fs.rm(lockPath, { recursive: true, force: true }); } + catch (error) { return { kind: 'unwritable', error, path }; } + } + } +} diff --git a/packages/sdk/src/cloud-http.ts b/packages/sdk/src/cloud-http.ts index 390faa88..5bc01a1e 100644 --- a/packages/sdk/src/cloud-http.ts +++ b/packages/sdk/src/cloud-http.ts @@ -1,6 +1,5 @@ -import { readFileSync } from 'node:fs'; -import { homedir } from 'node:os'; -import { join } from 'node:path'; +import { readAgentRelayCloudLogin, expiredCloudLogin, canRefreshCloudLogin, normalizeCloudBaseUrl, refreshCloudLogin } from './cloud-auth-store.js'; +export { agentRelayCloudAuthPath, readAgentRelayCloudLogin, type AgentRelayCloudLogin } from './cloud-auth-store.js'; export interface CloudConnectionOptions { /** Cloud application base URL; defaults to https://agentrelay.com/cloud. */ @@ -33,6 +32,7 @@ export interface CloudRefusal { export type CloudConfigurationReason = | 'auth_missing' | 'auth_expired' + | 'auth_store_unwritable' | 'url_invalid' | 'url_mismatch' | 'timeout_invalid'; @@ -61,50 +61,13 @@ function configurationError(reason: CloudConfigurationReason, message: string): return 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 } { - let rawToken = options.token ?? process.env['FLOWS_CLOUD_TOKEN']; + let rawToken = explicitCloudToken(options); 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()) { + if (expiredCloudLogin(login)) { throw configurationError('auth_expired', 'The agent-relay cloud login has expired. Run `agent-relay cloud login` again, or set FLOWS_CLOUD_TOKEN.'); } @@ -124,11 +87,10 @@ export function cloudConnection(options: CloudConnectionOptions): { baseUrl: str } catch { throw configurationError('url_invalid', 'FLOWS_CLOUD_URL must be an absolute Cloud application base URL.'); } - if (url.protocol !== 'https:' - || url.username || url.password || url.search || url.hash || !/^\/[A-Za-z0-9/_-]*$/u.test(url.pathname)) { + const baseUrl = normalizeCloudBaseUrl(url.href); + if (baseUrl === undefined) { throw configurationError('url_invalid', 'Cloud URL must use HTTPS and a plain base path.'); } - 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. @@ -147,6 +109,49 @@ export function cloudConnection(options: CloudConnectionOptions): { baseUrl: str return { baseUrl, token }; } +function explicitCloudToken(options: CloudConnectionOptions): string | undefined { + return options.token ?? process.env['FLOWS_CLOUD_TOKEN']; +} + +function requestTimeout(options: CloudConnectionOptions): number { + const timeout = options.requestTimeoutMs ?? 30_000; + if (!Number.isSafeInteger(timeout) || timeout < 1 || timeout > 2_147_483_647) { + throw configurationError('timeout_invalid', 'requestTimeoutMs must be a positive 32-bit integer.'); + } + return timeout; +} + +async function ensureFreshCloudLogin(options: CloudConnectionOptions): Promise { + if (explicitCloudToken(options) !== undefined) return; + const login = readAgentRelayCloudLogin(); + if (!login || !expiredCloudLogin(login) || !canRefreshCloudLogin(login)) return; + const baseUrl = normalizeCloudBaseUrl(options.apiUrl ?? process.env['FLOWS_CLOUD_URL'] ?? login.apiUrl); + if (baseUrl === undefined || baseUrl !== normalizeCloudBaseUrl(login.apiUrl)) return; + const result = await refreshCloudLogin(login, baseUrl, { signal: options.signal, timeoutMs: requestTimeout(options) }); + switch (result.kind) { + case 'current': case 'refreshed': return; + case 'blocked': throw new CloudFlowError('transient_error', 'Another process is refreshing the agent-relay cloud login; try again.'); + case 'transport': throw transportError(result.error); + case 'unwritable': { + const errno = isCloudRecord(result.error) && typeof result.error.code === 'string' ? result.error.code : 'unknown error'; + throw configurationError('auth_store_unwritable', `Cannot update the agent-relay cloud login at ${result.path} (${errno}). Check filesystem permissions and available space.`); + } + case 'refused': + if (result.status !== undefined && [400, 401, 403].includes(result.status)) { + throw configurationError('auth_expired', 'The agent-relay cloud login has expired. Run `agent-relay cloud login` again, or set FLOWS_CLOUD_TOKEN.'); + } + if (result.status !== undefined) throw new CloudFlowError('http_error', `Cloud login refresh failed with HTTP ${result.status}.`, result.status); + throw new CloudFlowError('invalid_response', 'Cloud returned an unusable login refresh response.'); + } +} + +/** Re-read after refresh; no cache that could hide a relay-side rotation. */ +export async function resolveCloudConnection(options: CloudConnectionOptions): Promise<{ baseUrl: string; token: string }> { + requestTimeout(options); + await ensureFreshCloudLogin(options); + return cloudConnection(options); +} + export async function cloudRequest( path: string, options: CloudConnectionOptions, @@ -182,11 +187,8 @@ export async function cloudFetch( 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 configurationError('timeout_invalid', 'requestTimeoutMs must be a positive 32-bit integer.'); - } + const timeout = requestTimeout(options); + const { baseUrl, token } = await resolveCloudConnection(options); const bearer = init.bearerToken ?? token; if (/[\r\n]/u.test(bearer) || !bearer.trim()) { throw new CloudFlowError('invalid_response', 'Cloud issued an unusable storage credential.'); diff --git a/packages/sdk/src/cloud-mirror-transport.ts b/packages/sdk/src/cloud-mirror-transport.ts index 0b50fb60..074e444a 100644 --- a/packages/sdk/src/cloud-mirror-transport.ts +++ b/packages/sdk/src/cloud-mirror-transport.ts @@ -113,6 +113,7 @@ export async function registerLocalRun( || typeof callbackToken !== 'string' || callbackToken.length === 0) { throw new CloudFlowError('invalid_response', 'Cloud registered the run without a usable credential.'); } + // Registration above already refreshed the store through cloudFetch. const { baseUrl } = cloudConnection(options); const runUrl = typeof result['runUrl'] === 'string' && /^https:\/\//u.test(result['runUrl']) ? result['runUrl'] diff --git a/packages/sdk/src/cloud-run.ts b/packages/sdk/src/cloud-run.ts index 1afdc1e5..0d16952b 100644 --- a/packages/sdk/src/cloud-run.ts +++ b/packages/sdk/src/cloud-run.ts @@ -10,7 +10,7 @@ import { snapshotJsonValue, type JsonValue } from './json-value.js'; import { loadAuthoredFlow, type SurfaceModuleAuthority } from './authored-flow-loader.js'; import { assertNoUseDependencies, collectExtensionSubmissions, type FlowExtensionSubmission } from './flow-extension-submit.js'; import { - CloudFlowError, cloudConnection, cloudFetch, cloudRequest, cloudRunId, isCloudRecord, + CloudFlowError, resolveCloudConnection, cloudFetch, cloudRequest, cloudRunId, isCloudRecord, type CloudConnectionOptions, } from './cloud-http.js'; import { cloudRunState, isCloudRunActive, type CloudRunState } from './cloud-run-record.js'; @@ -210,7 +210,7 @@ export async function runInCloud( flow: CloudFlowSource, options: RunInCloudOptions = {}, ): Promise { - const { baseUrl } = cloudConnection(options); + const { baseUrl } = await resolveCloudConnection(options); const submission = await prepareCloudSubmission(flow, options); const hash = submission.specHash; // Sync before submission: `prepare` reserves the run ID and the upload lands diff --git a/packages/sdk/tests/agent-artifacts.test.ts b/packages/sdk/tests/agent-artifacts.test.ts index 52be9df6..1e3d618f 100644 --- a/packages/sdk/tests/agent-artifacts.test.ts +++ b/packages/sdk/tests/agent-artifacts.test.ts @@ -1,3 +1,4 @@ +import { createHash } from 'node:crypto'; import { mkdtempSync, mkdirSync, rmSync, writeFileSync } from 'node:fs'; import { chmod } from 'node:fs/promises'; import { tmpdir } from 'node:os'; @@ -103,6 +104,37 @@ describe('snapshotWorkspaceFiles / diffWorkspaceFiles', () => { expect(await snapshotWorkspaceFiles(dir)).toEqual(new Map()); }); + it('accounts for a file it cannot read instead of failing the whole scan', async () => { + // `bun build --compile` creates its output in its cwd with mode 000, fills + // it with the executable and only then renames it onto the --outfile, so a + // scan of any tree a build is running in meets a present, unreadable file + // for seconds at a time. That must not fail an agent step that has done + // its work, and it must not be silently dropped either. + if (process.getuid?.() === 0) return; // root reads mode 000; nothing to assert + const dir = tempDir(); + writeFileSync(join(dir, 'readable.txt'), 'kept'); + const building = join(dir, '.aabbccdd-00000000.bun-build'); + writeFileSync(building, 'partial'); + await chmod(building, 0o000); + const before = await snapshotWorkspaceFiles(dir); + expect(before.get('readable.txt')).toBe(`4:${createHash('sha256').update('kept').digest('hex')}`); + expect(before.get('.aabbccdd-00000000.bun-build')).toBe('7:unreadable:EACCES'); + + // The build writes more into it; the signature tracks what can be seen. + await chmod(building, 0o600); + writeFileSync(building, 'partial and then some'); + await chmod(building, 0o000); + const after = await snapshotWorkspaceFiles(dir); + expect(after.get('.aabbccdd-00000000.bun-build')).toBe('21:unreadable:EACCES'); + expect(diffWorkspaceFiles(before, after)).toEqual(['.aabbccdd-00000000.bun-build']); + + // An unreadable signature can never read as a content hash: the same file + // once readable is a different signature at the same size. + await chmod(building, 0o600); + expect((await snapshotWorkspaceFiles(dir)).get('.aabbccdd-00000000.bun-build')) + .toBe(`21:${createHash('sha256').update('partial and then some').digest('hex')}`); + }); + it('propagates a non-ENOENT scan failure instead of silently omitting files', async () => { if (process.getuid?.() === 0) return; // root bypasses permission bits; nothing to assert const dir = tempDir(); diff --git a/packages/sdk/tests/cloud-auth-refresh.test.ts b/packages/sdk/tests/cloud-auth-refresh.test.ts new file mode 100644 index 00000000..678bd61e --- /dev/null +++ b/packages/sdk/tests/cloud-auth-refresh.test.ts @@ -0,0 +1,258 @@ +import fs from 'node:fs/promises'; +import { join } from 'node:path'; +import { tmpdir } from 'node:os'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { cloudFetch } from '../src/cloud-http.js'; +import { refreshCloudLogin } from '../src/cloud-auth-store.js'; +import { isTransientRead, refusalFor } from '../src/cli/cloud-refusal.js'; + +const apiUrl = 'https://login.example/cloud'; +const old = { apiUrl, accessToken: 'old-access', accessTokenExpiresAt: '2020-01-01T00:00:00Z', + refreshToken: 'old-refresh', refreshTokenExpiresAt: '2099-01-01T00:00:00Z' }; +const rotated = { accessToken: 'new-access', accessTokenExpiresAt: '2098-01-01T00:00:00Z', + refreshToken: 'new-refresh' }; +let home: string; +let path: string; +const request = (options = {}) => cloudFetch('/resource', options, { method: 'GET' }); +const write = (login: unknown) => fs.writeFile(path, JSON.stringify(login)); +const bytes = () => fs.readFile(path, 'utf8'); +function mockCloud(payload: unknown = rotated, status = 200) { + return vi.spyOn(globalThis, 'fetch').mockImplementation(async (url) => + new Response(JSON.stringify(String(url).endsWith('/token/refresh') ? payload : { ok: true }), { status })); +} +// @agent-relay/cloud@12.4.1 dist/auth.js:62-73 isValidStoredAuth. +function relayValid(value: Record): boolean { + return typeof value.accessToken === 'string' && typeof value.refreshToken === 'string' + && typeof value.accessTokenExpiresAt === 'string' && typeof value.apiUrl === 'string' + && (value.refreshTokenExpiresAt === undefined || typeof value.refreshTokenExpiresAt === 'string') + && !Number.isNaN(Date.parse(value.accessTokenExpiresAt)) + && (value.refreshTokenExpiresAt === undefined || !Number.isNaN(Date.parse(value.refreshTokenExpiresAt as string))); +} + +beforeEach(async () => { + home = await fs.mkdtemp(join(tmpdir(), 'cloud-refresh-')); + path = join(home, 'cloud-auth.json'); + vi.stubEnv('AGENT_RELAY_HOME', home); + vi.stubEnv('FLOWS_CLOUD_TOKEN', undefined); + vi.stubEnv('FLOWS_CLOUD_URL', undefined); + await write(old); +}); +afterEach(async () => { + vi.restoreAllMocks(); + vi.unstubAllEnvs(); + await fs.rm(home, { recursive: true, force: true }); +}); + +describe('cloud login refresh', () => { + it('persists rotation atomically in relay format before proceeding, silently', async () => { + await write({ ...old, futureKey: { keep: true } }); + const rename = vi.spyOn(fs, 'rename'); + const mkdir = vi.spyOn(fs, 'mkdir'); + const writeFile = vi.spyOn(fs, 'writeFile'); + const stdout = vi.spyOn(process.stdout, 'write'); + const fetch = mockCloud(); + fetch.mockImplementation(async url => { + if (String(url).endsWith('/token/refresh')) return new Response(JSON.stringify(rotated)); + expect(JSON.parse(await bytes()).accessToken).toBe(rotated.accessToken); + return new Response(JSON.stringify({ ok: true })); + }); + await expect(request()).resolves.toEqual({ ok: true }); + expect(stdout).not.toHaveBeenCalled(); + const saved = JSON.parse(await bytes()); + expect(saved).toEqual({ ...old, ...rotated, futureKey: { keep: true } }); + expect(relayValid(saved)).toBe(true); + expect(await bytes()).toBe(`${JSON.stringify(saved, null, 2)}\n`); + expect((await fs.stat(path)).mode & 0o777).toBe(0o600); + expect(rename).toHaveBeenCalledWith(expect.stringMatching(/\.cloud-auth\.json\..*\.tmp$/u), path); + for (const [created] of [...mkdir.mock.calls, ...writeFile.mock.calls]) expect(String(created).startsWith(home)).toBe(true); + expect(await fs.readdir(home)).toEqual(['cloud-auth.json']); + expect(fetch.mock.calls[0]).toEqual([`${apiUrl}/api/v1/auth/token/refresh`, expect.objectContaining({ + method: 'POST', body: JSON.stringify({ refreshToken: old.refreshToken }), redirect: 'error', + headers: { 'content-type': 'application/json' }, + })]); + expect(fetch.mock.calls[1]?.[1]?.headers).toMatchObject({ authorization: 'Bearer new-access' }); + await request(); + expect(fetch).toHaveBeenCalledTimes(3); + }); + + it.each([ + { refreshToken: undefined }, { refreshToken: '' }, { refreshTokenExpiresAt: '2020-01-01' }, + ])('refuses unrefreshable logins without requests: %j', async patch => { + await write({ ...old, ...patch }); + const before = await bytes(); + const fetch = mockCloud(); + await expect(request()).rejects.toMatchObject({ code: 'configuration', reason: 'auth_expired' }); + expect(fetch).not.toHaveBeenCalled(); + expect(await bytes()).toBe(before); + }); + + it.each([400, 401, 403, 429, 503])('classifies HTTP %i without altering the store', async status => { + const before = await bytes(); + mockCloud({ error: 'do not echo secret' }, status); + const error = await request().catch(e => e); + if ([400, 401, 403].includes(status)) { + expect(error).toMatchObject({ code: 'configuration', reason: 'auth_expired' }); + } else { + expect(error).toMatchObject({ code: 'http_error', status }); + expect(isTransientRead(error)).toBe(true); + expect(error.message).not.toContain('has expired'); + } + expect(error.message).not.toContain('secret'); + expect(await bytes()).toBe(before); + }); + + it.each([ + { ...rotated, refreshToken: undefined }, { ...rotated, accessTokenExpiresAt: 'invalid' }, + { ...rotated, accessToken: 'bad\ntoken' }, { ...rotated, refreshTokenExpiresAt: 'invalid' }, + ...['http://login.example/cloud', 'https://login.example/cloud?q=1', 'https://login.example/a.b', + 'https://other.example/cloud'].map(apiUrl => ({ ...rotated, apiUrl })), + ])('rejects unusable payloads without persistence: %j', async payload => { + const before = await bytes(); + mockCloud(payload); + await expect(request()).rejects.toMatchObject({ code: 'invalid_response' }); + expect(await bytes()).toBe(before); + }); + + it.each(['token', 'env', 'empty-token', 'empty-env'])('explicit precedence bypasses all store operations: %s', async source => { + const before = await bytes(); + const mkdir = vi.spyOn(fs, 'mkdir'); + const fetch = mockCloud(); + const empty = source.startsWith('empty'); + const token = empty ? '' : 'explicit'; + if (source.endsWith('env')) vi.stubEnv('FLOWS_CLOUD_TOKEN', token); + else vi.stubEnv('FLOWS_CLOUD_TOKEN', 'lower-priority'); + const options = source.endsWith('env') ? {} : { token }; + if (empty) await expect(request(options)).rejects.toMatchObject({ reason: 'auth_missing' }); + else { + await request(options); + expect(fetch).toHaveBeenCalledTimes(1); + expect(fetch.mock.calls[0]?.[1]?.headers).toMatchObject({ authorization: 'Bearer explicit' }); + } + expect(mkdir).not.toHaveBeenCalled(); + expect(await bytes()).toBe(before); + }); + + it.each(['https://other.example/cloud', '', 'http://login.example/cloud'])('never refreshes mismatched or invalid target %s', async url => { + vi.stubEnv('FLOWS_CLOUD_URL', url); + const fetch = mockCloud(); + await expect(request()).rejects.toMatchObject({ reason: 'auth_expired' }); + expect(fetch).not.toHaveBeenCalled(); + }); + + it.each([0, -1, 1.5, NaN, 2 ** 32])('validates timeout %s before refreshing', async requestTimeoutMs => { + const before = await bytes(); + const fetch = mockCloud(); + await expect(request({ requestTimeoutMs })).rejects.toMatchObject({ reason: 'timeout_invalid' }); + expect(fetch).not.toHaveBeenCalled(); + expect(await bytes()).toBe(before); + }); + + it.each([undefined, 'invalid', '2099-01-01'])('does not renew unexpired/unknown access expiry %s', async accessTokenExpiresAt => { + await write({ ...old, accessTokenExpiresAt }); + const fetch = mockCloud(); + await request(); + expect(fetch).toHaveBeenCalledTimes(1); + expect(fetch.mock.calls[0]?.[1]?.headers).toMatchObject({ authorization: 'Bearer old-access' }); + }); + + it('lets the server decide an unparseable refresh expiry and persists a valid replacement', async () => { + await write({ ...old, refreshTokenExpiresAt: 'invalid' }); + mockCloud({ ...rotated, refreshTokenExpiresAt: old.refreshTokenExpiresAt }); + await request(); + expect(relayValid(JSON.parse(await bytes()))).toBe(true); + }); + + it('collapses concurrent refreshes with a lock and re-read', async () => { + const fetch = mockCloud(); + await Promise.all([request(), request()]); + expect(fetch.mock.calls.filter(([url]) => String(url).endsWith('/token/refresh'))).toHaveLength(1); + }); + + it('classifies lock contention as transient within the request timeout', async () => { + await fs.mkdir(`${path}.lock`); + const fetch = mockCloud(); + await expect(request({ requestTimeoutMs: 10 })).rejects.toMatchObject({ code: 'transient_error', message: expect.stringContaining('Another process') }); + expect(fetch).not.toHaveBeenCalled(); + expect(await fs.stat(`${path}.lock`)).toBeDefined(); + }); + + it('removes stale locks using the relay stale window', async () => { + await fs.mkdir(`${path}.lock`); + await fs.utimes(`${path}.lock`, new Date(0), new Date(0)); + mockCloud(); + expect(await refreshCloudLogin(old, apiUrl, { lock: { retryMs: 1, acquireTimeoutMs: 10 } })).toEqual({ kind: 'refreshed' }); + }); + + it('re-sources a newly rotated but still expired token after acquiring the lock', async () => { + const mkdir = fs.mkdir.bind(fs); + vi.spyOn(fs, 'mkdir').mockImplementation(async (target, options) => { + if (String(target) === `${path}.lock`) await write({ ...old, refreshToken: 'rotated-while-waiting' }); + return mkdir(target, options as Parameters[1]); + }); + const fetch = mockCloud(); + await refreshCloudLogin(old, apiUrl); + expect(fetch.mock.calls[0]?.[1]?.body).toBe(JSON.stringify({ refreshToken: 'rotated-while-waiting' })); + }); + + it.each(['rename', 'mkdir'] as const)('classifies %s EACCES, never uses an unpersisted token', async operation => { + const before = await bytes(); + vi.spyOn(fs, operation).mockRejectedValue(Object.assign(new Error('secret'), { code: 'EACCES' })); + const fetch = mockCloud(); + const error = await request().catch(e => e); + expect(error).toMatchObject({ code: 'configuration', reason: 'auth_store_unwritable' }); + expect(error.message).toContain(path); + expect(error.message).toContain('EACCES'); + expect(error.message).not.toContain('secret'); + expect(refusalFor(error, 'login', {})).toMatchObject({ exit: 2 }); + expect(fetch).toHaveBeenCalledTimes(operation === 'rename' ? 1 : 0); + expect(await bytes()).toBe(before); + expect(await fs.readdir(home)).toEqual(['cloud-auth.json']); + }); + + it.each([ + [new TypeError('fetch failed', { cause: { code: 'ECONNRESET' } }), 'transient_error'], + [new TypeError('TLS failed'), 'transport_error'], + ])('classifies transport errors without changing credentials', async (error, code) => { + const before = await bytes(); + vi.spyOn(globalThis, 'fetch').mockRejectedValue(error); + await expect(request()).rejects.toMatchObject({ code }); + expect(await bytes()).toBe(before); + }); + + it('bounds the refresh request with the refresh timeout', async () => { + const before = await bytes(); + const fetch = vi.fn(async (_url, init) => new Promise((_resolve, reject) => { + init!.signal!.addEventListener('abort', () => reject(init!.signal!.reason), { once: true }); + })); + const result = await refreshCloudLogin(old, apiUrl, { fetch, refreshTimeoutMs: 5 }); + expect(result).toMatchObject({ kind: 'transport', error: { name: 'TimeoutError' } }); + expect(await bytes()).toBe(before); + }); + + it('classifies a body timeout and malformed JSON without writing', async () => { + const before = await bytes(); + const fetch = vi.fn(async (_url, init) => ({ + ok: true, + json: () => new Promise((_resolve, reject) => { + init!.signal!.addEventListener('abort', () => reject(init!.signal!.reason), { once: true }); + }), + } as Response)); + expect(await refreshCloudLogin(old, apiUrl, { fetch, refreshTimeoutMs: 5 })) + .toMatchObject({ kind: 'transport', error: { name: 'TimeoutError' } }); + vi.spyOn(globalThis, 'fetch').mockResolvedValue(new Response('not JSON')); + await expect(request()).rejects.toMatchObject({ code: 'invalid_response' }); + expect(await bytes()).toBe(before); + }); + + it('propagates cancellation while waiting without removing the other process lock', async () => { + await fs.mkdir(`${path}.lock`); + const controller = new AbortController(); + const fetch = mockCloud(); + const promise = refreshCloudLogin(old, apiUrl, { signal: controller.signal, lock: { retryMs: 1 } }); + controller.abort(new Error('cancelled')); + await expect(promise).rejects.toThrow('cancelled'); + expect(fetch).not.toHaveBeenCalled(); + expect(await fs.stat(`${path}.lock`)).toBeDefined(); + }); + +}); diff --git a/packages/sdk/tests/cloud-read.test.ts b/packages/sdk/tests/cloud-read.test.ts index 5154a698..0a4ab440 100644 --- a/packages/sdk/tests/cloud-read.test.ts +++ b/packages/sdk/tests/cloud-read.test.ts @@ -649,6 +649,37 @@ describe('refusals', () => { expect(fetch).not.toHaveBeenCalled(); }); + it('refreshes an expired login before logs and keeps JSON refusals to one object', async () => { + const home = temporaryDirectory(); + const login = { apiUrl: 'https://cloud-contract.example', accessToken: 'old', + accessTokenExpiresAt: '2020-01-01', refreshToken: 'refresh', refreshTokenExpiresAt: '2099-01-01' }; + writeFileSync(join(home, 'cloud-auth.json'), JSON.stringify(login)); + vi.stubEnv('FLOWS_CLOUD_TOKEN', undefined); + vi.stubEnv('FLOWS_CLOUD_URL', undefined); + vi.stubEnv('AGENT_RELAY_HOME', home); + cloud((path, _query, auth) => { + if (path.endsWith('/token/refresh')) return { body: { + accessToken: 'renewed', refreshToken: 'rotated', accessTokenExpiresAt: '2099-01-01', + } }; + expect(auth).toBe('Bearer renewed'); + if (path.endsWith('/steps')) return { body: STEPS }; + if (path.endsWith('/logs')) return { body: logEnvelope('') }; + return { body: RUN_DETAIL }; + }); + const out = io(); + expect(await runCloudLogsCli({ command: 'logs', runId: RUN, step: undefined, raw: false, json: true }, out.io, { env: {} })).toBe(0); + expect(out.stdout).toHaveLength(1); + expect(() => JSON.parse(out.stdout[0]!)).not.toThrow(); + expect(out.stderr).toEqual([]); + writeFileSync(join(home, 'cloud-auth.json'), JSON.stringify(login)); + cloud(() => ({ status: 401, body: {} })); + const refused = io(); + expect(await runCloudLogsCli({ command: 'logs', runId: RUN, step: undefined, raw: false, json: true }, refused.io, { env: {} })).toBe(2); + expect(refused.stdout).toHaveLength(1); + expect(JSON.parse(refused.stdout[0]!)).toMatchObject({ code: 'cloud_auth_expired' }); + expect(refused.stderr).toEqual([]); + }); + it('separates a run the credential cannot see (404) from one it may not read (403)', async () => { cloud(() => ({ status: 404, body: { error: 'Run not found' } })); const absent = io(); diff --git a/packages/sdk/tests/cloud-run.test.ts b/packages/sdk/tests/cloud-run.test.ts index bc731fe3..92fdfd86 100644 --- a/packages/sdk/tests/cloud-run.test.ts +++ b/packages/sdk/tests/cloud-run.test.ts @@ -53,6 +53,26 @@ describe('hosted v2 submission', () => { expect(JSON.parse(body.workflow).steps[0]).toMatchObject({ id: 'gate', type: 'deterministic', command: 'printf verified' }); }); + it('refreshes the login before resolving the hosted submission base URL', async () => { + const home = await mkdtemp(join(tmpdir(), 'cloud-run-refresh-')); + dirs.push(home); + vi.stubEnv('AGENT_RELAY_HOME', home); + vi.stubEnv('FLOWS_CLOUD_TOKEN', undefined); + vi.stubEnv('FLOWS_CLOUD_URL', undefined); + await writeFile(join(home, 'cloud-auth.json'), JSON.stringify({ + apiUrl: 'https://cloud-contract.example', accessToken: 'expired', + accessTokenExpiresAt: '2020-01-01', refreshToken: 'refresh', + })); + await cloud((path, _body, auth) => { + if (path.endsWith('/token/refresh')) return { + accessToken: 'renewed', refreshToken: 'rotated', accessTokenExpiresAt: '2099-01-01', + }; + expect(auth).toBe('Bearer renewed'); + return { runId: 'refreshed-run', status: 'pending' }; + }); + expect(await runInCloud(flow)).toMatchObject({ runId: 'refreshed-run' }); + }); + it('refuses invalid specs and unsupported source extensions before any HTTP request', async () => { const fetch = vi.spyOn(globalThis, 'fetch'); const options = { token: 'test-token' };