diff --git a/packages/sdk/src/agent-artifacts.ts b/packages/sdk/src/agent-artifacts.ts new file mode 100644 index 00000000..13bcfda9 --- /dev/null +++ b/packages/sdk/src/agent-artifacts.ts @@ -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> { + const out = new Map(); + await walk(dir, dir, out); + return out; +} + +async function walk(root: string, current: string, out: Map): Promise { + 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, after: Map): string[] { + const changed: string[] = []; + for (const [path, signature] of after) { + if (before.get(path) !== signature) changed.push(path); + } + return changed.sort(); +} diff --git a/packages/sdk/src/authored-worker-step.ts b/packages/sdk/src/authored-worker-step.ts index cb7e670f..defa8ff5 100644 --- a/packages/sdk/src/authored-worker-step.ts +++ b/packages/sdk/src/authored-worker-step.ts @@ -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; @@ -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); const output = await run({ id, type: 'agent', instruction: options.task, ...(localAgentStream === undefined ? {} : { surfaces: { streams: [{ stream: localAgentStream }] } }), @@ -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 { if (options !== undefined && (typeof options !== 'object' || options === null || options.output === undefined)) { diff --git a/packages/sdk/tests/agent-artifacts.test.ts b/packages/sdk/tests/agent-artifacts.test.ts new file mode 100644 index 00000000..59c0d6c6 --- /dev/null +++ b/packages/sdk/tests/agent-artifacts.test.ts @@ -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); + } + }); +}); diff --git a/packages/sdk/tests/authored-agent-artifacts.test.ts b/packages/sdk/tests/authored-agent-artifacts.test.ts new file mode 100644 index 00000000..04e59852 --- /dev/null +++ b/packages/sdk/tests/authored-agent-artifacts.test.ts @@ -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(); + return startLoopback(sock, { + hello: ctx => sendOk(ctx), + 'run.start': (ctx, params) => { + const spec = params.spec as Record; + const step = (spec['steps'] as Record[])[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(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(); + } + }); +});