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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 7 additions & 2 deletions docs/CLOUD.md
Original file line number Diff line number Diff line change
Expand Up @@ -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` |
Expand Down Expand Up @@ -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
Expand Down
35 changes: 29 additions & 6 deletions packages/sdk/src/agent-artifacts.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand All @@ -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<Map<string, string>> {
const out = new Map<string, string>();
Expand Down Expand Up @@ -68,12 +76,27 @@ async function walk(root: string, current: string, out: Map<string, string>): 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<string, string>, after: Map<string, string>): string[] {
const changed: string[] = [];
Expand Down
5 changes: 3 additions & 2 deletions packages/sdk/src/cli/cloud-mirror-session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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
Expand Down
185 changes: 185 additions & 0 deletions packages/sdk/src/cloud-auth-store.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown> {
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<string, unknown> {
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<void> {
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<CloudLoginRefresh> {
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 }; }
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Finally overrides successful token refresh

Medium Severity

A return in the finally of refreshCloudLogin replaces the try/catch result when lock cleanup fails. A completed persist can surface as auth_store_unwritable, and a cancellation can be swallowed, so the caller refuses instead of using the already-rotated tokens.

Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 445ee43. Configure here.

}
}
Loading
Loading