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
7 changes: 7 additions & 0 deletions docs/CLOUD.md
Original file line number Diff line number Diff line change
Expand Up @@ -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/<id>/steps`
answers per-step rows carrying state, attempts, timing, gate verdicts, spend
Expand Down
12 changes: 12 additions & 0 deletions docs/SURFACE.md
Original file line number Diff line number Diff line change
Expand Up @@ -965,6 +965,18 @@ effect refs, so their absence is structural. `LEASE OVERDUE by <t>` 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 <n>` 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/<run-id>/steps/<step-id>/attempt-<n>.{stdout,stderr}.tail`:
Expand Down
58 changes: 46 additions & 12 deletions packages/sdk/src/authored-step-index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<unknown>,
index: Map<string, AuthoredStepRecord> = new Map(),
): Map<string, AuthoredStepRecord> {
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 } : {}),
};
}
Comment thread
cursor[bot] marked this conversation as resolved.

function isAuthoredStepRecord(value: unknown): value is AuthoredStepRecord {
if (typeof value !== 'object' || value === null || Array.isArray(value)) return false;
const record = value as Partial<AuthoredStepRecord>;
Expand Down
51 changes: 50 additions & 1 deletion packages/sdk/src/cli/status.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -179,14 +180,62 @@ 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);
}
if (taken.failure !== undefined) io.stderr(`FAILED [${taken.failure.code}] ${taken.failure.message}`);
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<RunView, 'steps'> & { this_step: string | null; steps: PresentedStep[] };

Expand Down
65 changes: 64 additions & 1 deletion packages/sdk/tests/authored-step-graph-live.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown> }> {
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<string, unknown> }) };
}

async function statusJson(dataDir: string, runId: string): Promise<Record<string, unknown>> {
const { code, view } = await readStatus(dataDir, runId);
expect(code).toBe(0);
return view!;
}

const closes: Array<() => Promise<void>> = [];
// Every close runs even when an earlier one rejects, so a failed agent close never leaks the daemon.
afterEach(async () => {
Expand Down Expand Up @@ -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<void>((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<string, unknown> | 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<Record<string, unknown>>)
.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(
Expand Down Expand Up @@ -90,5 +139,19 @@ describe('the authored step DAG through the live kernel', () => {
state: 'completed', completionReason: 'success', ...edges,
}])),
);

// `flows status --json <root>` 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);
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.
52 changes: 52 additions & 0 deletions packages/sdk/tests/authored-step-index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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', {
Expand Down
Loading
Loading