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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
76 changes: 76 additions & 0 deletions packages/sdk/src/agent-artifacts.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
import { createHash } from 'node:crypto';
import { readFile, readdir, stat } from 'node:fs/promises';
import { join, relative, sep } from 'node:path';

function isEnoent(error: unknown): boolean {
return error instanceof Error && (error as NodeJS.ErrnoException).code === 'ENOENT';
}

const SKIPPED_DIR_NAMES = new Set(['node_modules']);

/**
* A recursive snapshot of every regular file under `dir`, keyed by its
* `dir`-relative POSIX path, valued by a content signature (size + sha256).
* Content, not mtime: a fast rewrite inside the filesystem's timestamp
* resolution must still count as a change — the same rule
* `examples/research/shims/agent-cli.ts`'s `snapshot()` uses for the same
* "did the agent actually write something" question. Dotfiles/dotdirs
* (`.git`, the kernel's own `.relayflowd` data dir, editor swap files, ...)
* are never agent-authored content and are skipped, as is `node_modules`.
* 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.
*/
export async function snapshotWorkspaceFiles(dir: string): Promise<Map<string, string>> {
const out = new Map<string, string>();
await walk(dir, dir, out);
return out;
}

async function walk(root: string, current: string, out: Map<string, string>): Promise<void> {
let entries;
try {
entries = await readdir(current, { withFileTypes: true });
} catch (error) {
if (isEnoent(error)) return;
throw error;
}
for (const entry of entries) {
if (entry.name.startsWith('.') || SKIPPED_DIR_NAMES.has(entry.name)) continue;
const path = join(current, entry.name);
if (entry.isDirectory()) {
await walk(root, path, out);
continue;
}
if (!entry.isFile()) continue;
let info;
try {
info = await stat(path);
} catch (error) {
// A file can vanish between readdir and stat (a CLI's own temp file);
// that is not an artifact, not a scan failure.
if (isEnoent(error)) continue;
throw error;
}
let bytes;
try {
bytes = await readFile(path);
} catch (error) {
if (isEnoent(error)) continue;
throw error;
}
out.set(relative(root, path).split(sep).join('/'), `${info.size}:${createHash('sha256').update(bytes).digest('hex')}`);
}
}

/** `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[] = [];
for (const [path, signature] of after) {
if (before.get(path) !== signature) changed.push(path);
}
return changed.sort();
}
20 changes: 19 additions & 1 deletion packages/sdk/src/authored-worker-step.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import { isSurfaceCompletionReason, readCompletedStepOutput, readSuccessfulOutpu
import type { AuthoredFlowJournalStep } from './authored-flow-executor.js';
import { snapshotJsonValue } from './json-value.js';
import { authoredChildAdmissionKey } from './authored-admission.js';
import { diffWorkspaceFiles, snapshotWorkspaceFiles } from './agent-artifacts.js';

const WORKSPACE_PERMISSION_ANNOTATION = /:\s*(readonly|readwrite)\s*$/i;

Expand Down Expand Up @@ -129,6 +130,20 @@ export function authoredWorkerRunner(
`f.agent options.transport must be 'direct' or 'relay' (got ${JSON.stringify(options.transport)}).`,
);
}
// Artifact detection only tells the truth for the local-agent DIRECT
// path: that is the only case that runs in this same process, on this
// same filesystem, so `options.cwd` (or `process.cwd()`) is provably
// where the CLI actually wrote — a workspace-scoped step never reaches
// here with a local agent attached (refused above). `transport: 'relay'`
// dispatches to agent-relay, which executes on a remote host even
// though a local agent stream is still attached, so it gets no local
// snapshot either. Any other worker attachment may execute on a
// different host entirely; snapshotting this process's filesystem for
// that case would be a guess, not a fact, so `artifacts` stays `[]`
// there, exactly as before this fix.
const artifactRoot = localAgentStream === undefined || options.transport === 'relay'
? undefined : options.cwd ?? process.cwd();
const before = artifactRoot === undefined ? undefined : await snapshotWorkspaceFiles(artifactRoot);
Comment thread
cursor[bot] marked this conversation as resolved.
const output = await run({
id, type: 'agent', instruction: options.task,
...(localAgentStream === undefined ? {} : { surfaces: { streams: [{ stream: localAgentStream }] } }),
Expand All @@ -143,7 +158,10 @@ export function authoredWorkerRunner(
throw new AuthoredFlowExecutionError('journal_protocol_violation', `step "${id}" produced a non-object output`);
}
const stdout = 'stdout_tail' in output ? output.stdout_tail : undefined;
return { summary: typeof stdout === 'string' ? stdout : JSON.stringify(output), artifacts: [] };
const artifacts = before === undefined || artifactRoot === undefined
? []
: diffWorkspaceFiles(before, await snapshotWorkspaceFiles(artifactRoot));
return { summary: typeof stdout === 'string' ? stdout : JSON.stringify(output), artifacts };
},
async llm(id: string, prompt: string, options?: LlmOptions, verification?: NamedGate): Promise<unknown> {
if (options !== undefined && (typeof options !== 'object' || options === null || options.output === undefined)) {
Expand Down
73 changes: 73 additions & 0 deletions packages/sdk/tests/agent-artifacts.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
import { mkdtempSync, mkdirSync, rmSync, writeFileSync } from 'node:fs';
import { chmod } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { afterEach, describe, expect, it } from 'vitest';
import { diffWorkspaceFiles, snapshotWorkspaceFiles } from '../src/agent-artifacts.js';

const dirs: string[] = [];
afterEach(() => { for (const dir of dirs.splice(0)) rmSync(dir, { recursive: true, force: true }); });

function tempDir(): string {
const dir = mkdtempSync(join(tmpdir(), 'agent-artifacts-'));
dirs.push(dir);
return dir;
}

describe('snapshotWorkspaceFiles / diffWorkspaceFiles', () => {
it('reports a newly written nested file as an artifact', async () => {
const dir = tempDir();
const before = await snapshotWorkspaceFiles(dir);
mkdirSync(join(dir, 'research'));
writeFileSync(join(dir, 'research', 'notes.md'), 'facts');
const after = await snapshotWorkspaceFiles(dir);
expect(diffWorkspaceFiles(before, after)).toEqual(['research/notes.md']);
});

it('reports a same-size content change to an existing file', async () => {
const dir = tempDir();
writeFileSync(join(dir, 'post.md'), 'aaaa');
const before = await snapshotWorkspaceFiles(dir);
writeFileSync(join(dir, 'post.md'), 'bbbb');
const after = await snapshotWorkspaceFiles(dir);
expect(diffWorkspaceFiles(before, after)).toEqual(['post.md']);
});

it('ignores files that are unchanged', async () => {
const dir = tempDir();
writeFileSync(join(dir, 'unrelated.txt'), 'same');
const before = await snapshotWorkspaceFiles(dir);
const after = await snapshotWorkspaceFiles(dir);
expect(diffWorkspaceFiles(before, after)).toEqual([]);
});

it('skips dotdirs and node_modules', async () => {
const dir = tempDir();
const before = await snapshotWorkspaceFiles(dir);
mkdirSync(join(dir, '.relayflowd'));
writeFileSync(join(dir, '.relayflowd', 'kernel.log'), 'noise');
mkdirSync(join(dir, 'node_modules'));
writeFileSync(join(dir, 'node_modules', 'pkg.js'), 'noise');
const after = await snapshotWorkspaceFiles(dir);
expect(diffWorkspaceFiles(before, after)).toEqual([]);
});

it('yields an empty snapshot for a directory that does not exist yet', async () => {
const dir = join(tempDir(), 'does-not-exist');
expect(await snapshotWorkspaceFiles(dir)).toEqual(new Map());
});

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();
const locked = join(dir, 'locked');
mkdirSync(locked);
writeFileSync(join(locked, 'secret.txt'), 'x');
try {
await chmod(locked, 0o000);
await expect(snapshotWorkspaceFiles(dir)).rejects.toMatchObject({ code: 'EACCES' });
} finally {
await chmod(locked, 0o755);
}
});
});
162 changes: 162 additions & 0 deletions packages/sdk/tests/authored-agent-artifacts.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
import { chmodSync, mkdirSync, mkdtempSync, rmSync, symlinkSync, writeFileSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join, resolve } from 'node:path';
import type { Server } from 'node:net';
import { afterEach, describe, expect, it } from 'vitest';
import { flow } from '@relayflows/surface';
import { executeAuthoredFlow } from '../src/authored-flow-executor.js';
import { JournalClient } from '../src/journal-client.js';
import { sendOk, sendResult, sockPath, startLoopback } from './journal-client-loopback.js';

// A minimal `relayflows-agent-cli-v1` wrapper: enough for `checkAuthoredFlow`'s
// real preflight probe (auth status + model round-trip) to resolve a CLI.
// The step's actual execution never reaches this wrapper in these tests — no
// local-agent worker is attached, so the fake journal server below completes
// every run directly, exactly like `authored-flow.test.ts` does.
function writeWrapper(root: string): string {
const wrapper = join(root, 'adapter.mjs');
writeFileSync(wrapper, `#!/usr/bin/env node
import { receiveWrapperRequest } from ${JSON.stringify(resolve('../../testdata/preflight/wrapper-session.mjs'))};
if (process.argv[2] === 'auth') process.exit(0);
await receiveWrapperRequest();
process.stdout.write('unused');
`);
chmodSync(wrapper, 0o755);
return wrapper;
}

function setupProject(): { root: string } {
const root = mkdtempSync(join(tmpdir(), 'authored-artifacts-'));
const wrapper = writeWrapper(root);
writeFileSync(join(root, 'flows.json'), JSON.stringify({ cli: wrapper, models: ['test-model'] }));
writeFileSync(join(root, 'package.json'), '{"type":"module"}');
symlinkSync(resolve('node_modules'), join(root, 'node_modules'));
return { root };
}

/**
* A fake kernel that completes every run it is sent immediately: an agent
* step "writes" its file as a side effect of `run.start` (standing in for a
* real CLI, which nothing here spawns — no local-agent worker is attached),
* and every other step type (the executor's own synthetic `f.done()` marker
* step included) completes as a trivial deterministic success.
*/
function startArtifactServer(sock: string, root: string): Server {
let nextRun = 1;
const stepByRun = new Map<string, { id: string; type: string }>();
return startLoopback(sock, {
hello: ctx => sendOk(ctx),
'run.start': (ctx, params) => {
const spec = params.spec as Record<string, unknown>;
const step = (spec['steps'] as Record<string, unknown>[])[0]!;
const runId = `authored-artifact-run-${nextRun++}`;
stepByRun.set(runId, { id: step['id'] as string, type: step['type'] as string });
if (step['type'] === 'agent') {
mkdirSync(join(root, 'research'), { recursive: true });
writeFileSync(join(root, 'research', 'notes.md'), 'facts, with sources');
}
sendResult(ctx, { run_id: runId, status: 'completed', completion_reason: 'success', completed_steps: 1 });
},
'journal.read': (ctx, params) => {
const step = stepByRun.get(params.run_id as string)!;
sendResult(ctx, {
entries: [{
entry_type: 'step.completed',
step_id: step.id,
payload: {
completionReason: 'success',
disposition: 'step_done',
output: step.type === 'agent'
? { exit_code: 0, stdout_tail: 'wrote research/notes.md', stderr_tail: '' }
: { exit_code: 0, stdout_tail: '', stderr_tail: '' },
},
}],
});
},
});
}

describe('f.agent artifacts (local-agent path)', () => {
let server: Server | undefined;
let path: string | undefined;
let root: string | undefined;

afterEach(async () => {
if (server !== undefined) await new Promise<void>(resolve => server!.close(() => resolve()));
if (path !== undefined) rmSync(path, { force: true });
if (root !== undefined) rmSync(root, { recursive: true, force: true });
server = undefined; path = undefined; root = undefined;
});

it('reports a file the agent step wrote under its cwd', async () => {
({ root } = setupProject());
path = sockPath();
server = startArtifactServer(path, root);
const client = new JournalClient(path, { requestTimeoutMs: 2000 });
await client.connect();
await client.hello('authored-artifacts-test');
try {
const handle = flow('artifact-test', async (f) => {
const result = await f.agent('writer', { task: 'write research/notes.md', cwd: root! });
expect(result.artifacts).toEqual(['research/notes.md']);
f.done('success');
});
const result = await executeAuthoredFlow(handle, client, undefined, {
flowPath: join(root!, 'artifact-test.flow.ts'),
localAgentStream: 'test-stream',
});
expect(result.completionReason).toBe('success');
} finally {
client.close();
}
});

it('reports no artifacts for transport: relay even with a local agent stream attached', async () => {
({ root } = setupProject());
path = sockPath();
server = startArtifactServer(path, root);
const client = new JournalClient(path, { requestTimeoutMs: 2000 });
await client.connect();
await client.hello('authored-artifacts-test-relay');
try {
const handle = flow('artifact-test-relay', async (f) => {
// The fake server still performs its "the agent wrote a file" side
// effect for any agent-type step, but a relay-transport step ran on
// a remote host — this process's filesystem is not where it wrote,
// so the local snapshot must not be trusted here.
const result = await f.agent('writer', { task: 'write research/notes.md', cwd: root!, transport: 'relay' });
expect(result.artifacts).toEqual([]);
f.done('success');
});
const result = await executeAuthoredFlow(handle, client, undefined, {
flowPath: join(root!, 'artifact-test-relay.flow.ts'),
localAgentStream: 'test-stream',
});
expect(result.completionReason).toBe('success');
} finally {
client.close();
}
});

it('reports no artifacts when no local agent is attached', async () => {
({ root } = setupProject());
path = sockPath();
server = startArtifactServer(path, root);
const client = new JournalClient(path, { requestTimeoutMs: 2000 });
await client.connect();
await client.hello('authored-artifacts-test-2');
try {
const handle = flow('artifact-test-no-local', async (f) => {
const result = await f.agent('writer', { task: 'write research/notes.md', cwd: root!, workspace: 'research' });
expect(result.artifacts).toEqual([]);
f.done('success');
});
const result = await executeAuthoredFlow(handle, client, undefined, {
flowPath: join(root!, 'artifact-test-no-local.flow.ts'),
});
expect(result.completionReason).toBe('success');
} finally {
client.close();
}
});
});
Loading