diff --git a/docs/CLOUD.md b/docs/CLOUD.md index 06bf6623..eec52093 100644 --- a/docs/CLOUD.md +++ b/docs/CLOUD.md @@ -506,6 +506,13 @@ declaration and its `flows.tick` lowering either way. cannot learn from disk is Cloud's alone and is not guessed: its Cloud run id (a UUID Cloud may export separately), its sandbox, and which listener launched it. +- Cloud's live step view comes from polling `flows status --json` for every + journal under the run's data dir. For an authored run the root journal's + view also carries `authored_steps` (SURFACE.md §5): each step's `label` and + `after`, from the moment the step is admitted. A Cloud reporter that reads it + can name and connect the run graph's nodes while the run is in flight; one + that does not ignores the extra field, so this runtime is safe to pin under + either. - From outside the sandbox, step state *is* readable — see [Reading a hosted run](#reading-a-hosted-run). `GET /runs//steps` answers per-step rows carrying state, attempts, timing, gate verdicts, spend diff --git a/docs/SURFACE.md b/docs/SURFACE.md index 9537471c..4fc9f356 100644 --- a/docs/SURFACE.md +++ b/docs/SURFACE.md @@ -965,6 +965,18 @@ effect refs, so their absence is structural. `LEASE OVERDUE by ` in the text view is computed from the journaled lease deadline and the wall clock alone — a local dead-man that needs no daemon. +For an authored root, `--json` also carries `authored_steps`: the root's +`authored-steps` index (see *Authored bodies* above), folded exactly as +`readAuthoredStepIndex` folds it, one entry per authored step in admission +order — `step`, `run_id` (the child journal whose own `flows status` view holds +that step id), `state` (`admitted` | `completed`), and when present +`completion_reason`, `kernel_step`, `label`, `after` and `after_truncated`. +Because the admission record already carries `label` and `after`, a reader +polling the root sees each step's name and predecessors as soon as it is +admitted, not only once the run finishes. `label` is redacted like any other +free text; ids are printed as-is. The field is additive and absent for every +run that journals no index, so the existing shape is unchanged. + `--tail ` renders the last *n* lines of this attempt's stdout and stderr after redaction. Every direct agent attempt with a data dir tees its transcript into `runs//steps//attempt-.{stdout,stderr}.tail`: diff --git a/packages/sdk/src/authored-step-index.ts b/packages/sdk/src/authored-step-index.ts index 3ed6de0f..ac514b28 100644 --- a/packages/sdk/src/authored-step-index.ts +++ b/packages/sdk/src/authored-step-index.ts @@ -161,24 +161,58 @@ export async function readAuthoredStepIndex( let offset = 0; for (;;) { const page = await journal.streamRead(rootRunId, AUTHORED_STEP_STREAM, offset, 1000); - for (const message of page.messages) { - // `stream.read` returns either the envelope or the bare message, - // matching how the executor reads predicate verdicts. - const raw = (message as { message?: unknown }).message ?? message; - if (!isAuthoredStepRecord(raw)) continue; - const previous = index.get(raw.step); - // An `admitted` record arriving after a `completed` one (a resumed body - // re-admitting the same child under its stable admission key) must not - // un-complete it. - if (previous?.state === 'completed' && raw.state === 'admitted') continue; - index.set(raw.step, raw); - } + // `stream.read` returns either the envelope or the bare message, + // matching how the executor reads predicate verdicts. + foldAuthoredStepRecords(page.messages.map((message) => (message as { message?: unknown }).message ?? message), index); if (page.messages.length === 0 || page.next_offset <= offset) break; offset = page.next_offset; } return [...index.values()]; } +/** + * The fold itself, over messages already read from the stream in append + * order — by `readAuthoredStepIndex` through the daemon, or by `flows status` + * straight from the journal file, which must agree with it record for record. + */ +export function foldAuthoredStepRecords( + messages: Iterable, + index: Map = new Map(), +): Map { + for (const raw of messages) { + if (!isAuthoredStepRecord(raw)) continue; + const previous = index.get(raw.step); + if (previous === undefined) { + index.set(raw.step, raw); + continue; + } + // An `admitted` record arriving after a `completed` one (a resumed body + // re-admitting the same child under its stable admission key) must not + // un-complete it. Either way, graph fields one record lacks are taken from + // the other: a completion written without them (or a re-admission that + // adds them, over a root an older runtime began) must not erase them. + const keepPrevious = previous.state === 'completed' && raw.state === 'admitted'; + index.set(raw.step, withGraphFrom(keepPrevious ? previous : raw, keepPrevious ? raw : previous)); + } + return index; +} + +/** `primary`, with any `label` / `after` it lacks taken from `fallback`. */ +function withGraphFrom(primary: AuthoredStepRecord, fallback: AuthoredStepRecord): AuthoredStepRecord { + const label = primary.label ?? fallback.label; + // `after` and `afterTruncated` travel together: a truncated walk that found + // no predecessor journals the flag alone, and it must survive the fold. + const hasEdges = (record: AuthoredStepRecord) => record.after !== undefined || record.afterTruncated === true; + const edges = hasEdges(primary) ? primary : fallback; + const { label: _label, after: _after, afterTruncated: _truncated, ...lifecycle } = primary; + return { + ...lifecycle, + ...(label === undefined ? {} : { label }), + ...(edges.after === undefined ? {} : { after: edges.after }), + ...(edges.afterTruncated === true ? { afterTruncated: true as const } : {}), + }; +} + function isAuthoredStepRecord(value: unknown): value is AuthoredStepRecord { if (typeof value !== 'object' || value === null || Array.isArray(value)) return false; const record = value as Partial; diff --git a/packages/sdk/src/cli/status.ts b/packages/sdk/src/cli/status.ts index 5567b3de..bd61adb7 100644 --- a/packages/sdk/src/cli/status.ts +++ b/packages/sdk/src/cli/status.ts @@ -9,6 +9,7 @@ // sandbox, its listener) is Cloud's to answer, and this verb does not guess. import type { CliIo } from '../cli.js'; +import { AUTHORED_STEP_STREAM, foldAuthoredStepRecords, type AuthoredStepRecord } from '../authored-step-index.js'; import { canonicalize } from '../canonical.js'; import { DEFAULT_DATA_DIR } from '../daemon-connection.js'; import { JournalReadError, walkJournal, type JournalEvent } from '../journal-client.js'; @@ -179,7 +180,14 @@ export async function runStatus(args: StatusArgs, io: CliIo, options: StatusOpti const presented = present(view, thisStep, tails, env); if (args.json) { - io.stdout(canonicalize({ v: 1, ...presented, partial: taken.partial })); + const authored = authoredSteps(taken.events, env); + io.stdout(canonicalize({ + v: 1, ...presented, + // Additive: only an authored root journals this stream, so every other + // run's JSON is byte-for-byte what it was. + ...(authored.length === 0 ? {} : { authored_steps: authored }), + partial: taken.partial, + })); } else { for (const line of renderText(presented, taken.partial, args.tail !== undefined)) io.stdout(line); } @@ -187,6 +195,47 @@ export async function runStatus(args: StatusArgs, io: CliIo, options: StatusOpti return taken.partial.length === 0 ? 0 : 1; } +/** One authored step as the root's `authored-steps` index names it. */ +export interface PresentedAuthoredStep { + step: string; + run_id: string; + state: AuthoredStepRecord['state']; + completion_reason?: string; + kernel_step?: string; + label?: string; + after?: string[]; + after_truncated?: true; +} + +/** + * The authored root's step index, folded from the journal already in hand. + * + * An authored body's steps each run in their own child journal, so a child's + * view cannot say what the step is called or what it waited for; the root's + * `authored-steps` stream can, from the moment each child is admitted. It is + * the same fold `readAuthoredStepIndex` applies through the daemon, so the two + * readers cannot disagree. Identifiers are printed as-is, like every other id + * here; the label is author-chosen free text and is redacted like one. + */ +function authoredSteps(events: readonly JournalEvent[], env: NodeJS.ProcessEnv): PresentedAuthoredStep[] { + const messages: unknown[] = []; + for (const event of events) { + if (event.entry_type !== 'stream.appended') continue; + const payload = event.payload as { stream?: unknown; message?: unknown } | null; + if (payload?.stream === AUTHORED_STEP_STREAM) messages.push(payload.message); + } + return [...foldAuthoredStepRecords(messages).values()].map((record) => ({ + step: record.step, + run_id: record.runId, + state: record.state, + ...(record.completionReason === undefined ? {} : { completion_reason: record.completionReason }), + ...(record.kernelStep === undefined ? {} : { kernel_step: record.kernelStep }), + ...(record.label === undefined ? {} : { label: redact(record.label, env) }), + ...(record.after === undefined ? {} : { after: [...record.after] }), + ...(record.afterTruncated === true ? { after_truncated: true as const } : {}), + })); +} + type PresentedStep = StepView & { tails: StepTails }; type Presented = Omit & { this_step: string | null; steps: PresentedStep[] }; diff --git a/packages/sdk/tests/authored-step-graph-live.test.ts b/packages/sdk/tests/authored-step-graph-live.test.ts index 33fb3221..30a1ad46 100644 --- a/packages/sdk/tests/authored-step-graph-live.test.ts +++ b/packages/sdk/tests/authored-step-graph-live.test.ts @@ -4,8 +4,24 @@ import { executeAuthoredFlow } from '../src/authored-flow-executor.js'; import { attachLocalAgent } from '../src/local-agent.js'; import { AUTHORED_STEP_STREAM, readAuthoredStepIndex } from '../src/authored-step-index.js'; import type { JournalClient } from '../src/journal-client.js'; +import { parseStatusArgs, runStatus } from '../src/cli/status.js'; import { chainFixture } from './flow-chain-fixture.js'; +/** `flows status --json` for one run, read from the journal file on disk; the exit code beside it. */ +async function readStatus(dataDir: string, runId: string): Promise<{ code: number; view?: Record }> { + const stdout: string[] = []; + const code = await runStatus(parseStatusArgs(['--json', '--data-dir', dataDir, runId])!, { + stdout: (line) => stdout.push(line), stderr: () => undefined, + }, { env: {} }); + return { code, ...(stdout[0] === undefined ? {} : { view: JSON.parse(stdout[0]) as Record }) }; +} + +async function statusJson(dataDir: string, runId: string): Promise> { + const { code, view } = await readStatus(dataDir, runId); + expect(code).toBe(0); + return view!; +} + const closes: Array<() => Promise> = []; // Every close runs even when an earlier one rejects, so a failed agent close never leaks the daemon. afterEach(async () => { @@ -39,18 +55,51 @@ describe('the authored step DAG through the live kernel', () => { closes.push(() => agent.close()); const rootRunId = await openRoot(journal); + // A poller outside the body, as Cloud's reporter is: it reads the root's + // journal file while the daemon is still writing it. + let release!: () => void; + const observed = new Promise((resolve) => { release = resolve; }); const handle = flow('graph-live', async (f) => { await f.run('printf plan'); // Fan out on deterministic steps: the local agent worker serves one lease // at a time, so parallel agents would park rather than run. await Promise.all([f.run('printf left'), f.run('printf right')]); + // Hold here until the poller below has read the root mid-run. + await observed; await f.agent('writer', { task: 'write it' }); await f.run('printf done'); f.done('success'); }); - const result = await executeAuthoredFlow(handle, journal, undefined, { + const running = executeAuthoredFlow(handle, journal, undefined, { rootRunId, flowPath: fixture.flowPath, localAgentStream: agent.stream, dataDir: fixture.data, }); + let midRun: Record | undefined; + try { + const deadline = Date.now() + 30_000; + // A read that races the writer is refused (`journal_busy`) and simply + // retried, as Cloud's reporter retries on its next poll. + for (;;) { + const { code, view } = await readStatus(fixture.data, rootRunId); + if (code === 0) midRun = view; + const steps = midRun?.['authored_steps'] as Array<{ state: string }> | undefined; + if ((steps?.length === 3 && steps.every((step) => step.state === 'completed')) || Date.now() > deadline) break; + await new Promise((resolve) => setTimeout(resolve, 25)); + } + } finally { + release(); + } + const result = await running; + + expect(midRun?.['status']).toBe('running'); + // run-2 and run-3 run in parallel, so their admission order is not fixed. + expect((midRun?.['authored_steps'] as Array>) + .map(({ step, state, after }) => ({ step, state, after })) + .sort((a, b) => String(a.step).localeCompare(String(b.step)))) + .toEqual([ + { step: 'run-1', state: 'completed', after: undefined }, + { step: 'run-2', state: 'completed', after: ['run-1'] }, + { step: 'run-3', state: 'completed', after: ['run-1'] }, + ]); // Step identity is the kernel's and the admission key's: unchanged. expect(result.journalSteps.map((step) => step.id).sort()).toEqual( @@ -90,5 +139,19 @@ describe('the authored step DAG through the live kernel', () => { state: 'completed', completionReason: 'success', ...edges, }])), ); + + // `flows status --json ` reads the same index straight off the + // journal file — the offline path Cloud's live reporter polls — and each + // entry names the child journal whose own view carries that step id, so a + // reader can join the two without guessing. + const root = await statusJson(fixture.data, rootRunId); + expect(root['authored_steps']).toEqual(index.map((record) => ({ + step: record.step, run_id: record.runId, state: record.state, + completion_reason: record.completionReason, + ...(record.label === undefined ? {} : { label: record.label }), + ...(record.after === undefined ? {} : { after: record.after }), + }))); + const writer = await statusJson(fixture.data, journaled['agent-4']!.runId); + expect((writer['steps'] as Array<{ id: string }>).map((step) => step.id)).toEqual(['agent-4']); }, 60_000); }); diff --git a/packages/sdk/tests/authored-step-index.test.ts b/packages/sdk/tests/authored-step-index.test.ts index 73eac560..5bddc499 100644 --- a/packages/sdk/tests/authored-step-index.test.ts +++ b/packages/sdk/tests/authored-step-index.test.ts @@ -77,6 +77,58 @@ describe('the durable child index on the root run', () => { ]); }); + it('keeps graph fields across records, whichever record carried them', async () => { + const { journal } = streams(); + // A completion without graph fields after an admission that had them. + await recordAuthoredChild(journal, 'root-1', { + step: 'agent-1', runId: 'child-a', state: 'admitted', label: 'plan', after: [], + }); + await recordAuthoredChild(journal, 'root-1', { + step: 'agent-1', runId: 'child-a', state: 'completed', completionReason: 'success', + }); + // A root an older runtime began: a bare completion, then a resumed body's + // re-admission that now carries the fields. It stays completed and gains them. + await recordAuthoredChild(journal, 'root-1', { + step: 'agent-2', runId: 'child-b', state: 'completed', completionReason: 'success', + }); + await recordAuthoredChild(journal, 'root-1', { + step: 'agent-2', runId: 'child-b', state: 'admitted', label: 'writer', after: ['agent-1'], afterTruncated: true, + }); + // A truncated walk that found no predecessor journals the flag alone. + await recordAuthoredChild(journal, 'root-1', { + step: 'agent-4', runId: 'child-d', state: 'admitted', afterTruncated: true, + }); + await recordAuthoredChild(journal, 'root-1', { + step: 'agent-4', runId: 'child-d', state: 'completed', completionReason: 'success', afterTruncated: true, + }); + // A record's own fields win over the other's. + await recordAuthoredChild(journal, 'root-1', { + step: 'agent-3', runId: 'child-c', state: 'admitted', label: 'old', after: ['agent-1'], + }); + await recordAuthoredChild(journal, 'root-1', { + step: 'agent-3', runId: 'child-c', state: 'completed', completionReason: 'success', label: 'new', after: ['agent-2'], + }); + + expect(await readAuthoredStepIndex(journal, 'root-1')).toEqual([ + { + index: 'relayflows.authored-step.v1', step: 'agent-1', runId: 'child-a', state: 'completed', + completionReason: 'success', label: 'plan', after: [], + }, + { + index: 'relayflows.authored-step.v1', step: 'agent-2', runId: 'child-b', state: 'completed', + completionReason: 'success', label: 'writer', after: ['agent-1'], afterTruncated: true, + }, + { + index: 'relayflows.authored-step.v1', step: 'agent-4', runId: 'child-d', state: 'completed', + completionReason: 'success', afterTruncated: true, + }, + { + index: 'relayflows.authored-step.v1', step: 'agent-3', runId: 'child-c', state: 'completed', + completionReason: 'success', label: 'new', after: ['agent-2'], + }, + ]); + }); + it('keeps the kernel step id only when it differs from the authored one', async () => { const { journal } = streams(); await recordAuthoredChild(journal, 'root-1', { diff --git a/packages/sdk/tests/cli-status.test.ts b/packages/sdk/tests/cli-status.test.ts index f1c86614..42443cf6 100644 --- a/packages/sdk/tests/cli-status.test.ts +++ b/packages/sdk/tests/cli-status.test.ts @@ -120,6 +120,8 @@ describe('flows status', () => { 'last_attempt', 'lease', 'max_iterations', 'started_at_ms', 'state', 'tails', 'type', 'wait', ]); expect(view.steps[0].tails).toBeNull(); + // A run that journals no authored-step index keeps its original shape. + expect(view).not.toHaveProperty('authored_steps'); // Canonical: sorted keys, no whitespace, so two invocations diff clean. expect(output.stdout[0]).toBe(canonicalize(view)); for (const forbidden of ['instruction', 'Print DONE', 'stdout_tail', 'VERIFIED', 'REPORTED', 'wake_context', 'input', 'pins', 'idempotency_key', 'surface_path']) { @@ -191,6 +193,50 @@ describe('flows status', () => { expect(detail.length).toBe(1024); }); + it('--json carries the authored root\'s step index — label and after — from the first admission', async () => { + const secret = 'ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123'; + const last = EVENTS.at(-1)!; + let seq = last.seq; + const appended = (stream: string, message: unknown): JournalEvent => ({ + ...last, seq: ++seq, entry_type: 'stream.appended', step_id: null, attempt: null, + payload: { stream, offset: seq, producer: 'sdk', message }, + }); + const record = (fields: Record) => ({ index: 'relayflows.authored-step.v1', ...fields }); + const events = [ + ...EVENTS, + appended('authored-steps', record({ step: 'agent-1', runId: 'child-1', state: 'admitted', label: 'api-create', after: [] })), + appended('authored-steps', record({ + step: 'agent-2', runId: 'child-2', state: 'admitted', label: `deploy ${secret}`, after: ['agent-1'], + })), + appended('authored-steps', record({ + step: 'agent-1', runId: 'child-1', state: 'completed', completionReason: 'success', kernelStep: 'agent-1.verify', + label: 'api-create', after: [], afterTruncated: true, + })), + // A resumed body re-admitting a completed child must not un-complete it. + appended('authored-steps', record({ step: 'agent-1', runId: 'child-1', state: 'admitted', label: 'api-create' })), + // Not this index: another stream, another version, a malformed record. + appended('proof', record({ step: 'other', runId: 'child-9', state: 'admitted' })), + appended('authored-steps', { index: 'relayflows.authored-step.v2', step: 'future', runId: 'child-8', state: 'admitted' }), + appended('authored-steps', record({ step: 'bad', runId: 'child-7', state: 'admitted', after: 'agent-1' })), + ]; + const { dataDir, writer } = fixture(events); + writer.close(); + const output = await status(['--json', '--data-dir', dataDir, RUN_ID], { env: { GITHUB_TOKEN: secret } }); + expect(output.code).toBe(0); + expect(output.stdout[0]).not.toContain(secret); + const view = JSON.parse(output.stdout[0]!); + expect(view.authored_steps).toEqual([ + { + step: 'agent-1', run_id: 'child-1', state: 'completed', completion_reason: 'success', + kernel_step: 'agent-1.verify', label: 'api-create', after: [], after_truncated: true, + }, + { step: 'agent-2', run_id: 'child-2', state: 'admitted', label: 'deploy [redacted:GITHUB_TOKEN]', after: ['agent-1'] }, + ]); + // The per-step view is untouched: the index rides beside it. + expect(view.steps).toHaveLength(3); + expect(output.stdout[0]).toBe(canonicalize(view)); + }); + it('--tail says when no transcript is on disk for the attempt', async () => { const { dataDir, writer } = fixture(inFlight('x')); writer.close();