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
9 changes: 9 additions & 0 deletions docs/SURFACE.md
Original file line number Diff line number Diff line change
Expand Up @@ -858,6 +858,15 @@ child. Helper-provider, MCP and plugin-effect children are not yet included
in this index. Once appended, the index survives process exit and is readable from
the journal on disk, including after a cooperative nonzero exit.

Each record, and the matching root `journalSteps` entry, may also carry the
step's place in the run's DAG. `after` lists the step ids it causally waited
for, transitively reduced and capped at 32. `afterTruncated: true` means that
list is incomplete. `label` is the author-chosen `f.agent` or `f.hook` name,
kept whole or omitted when it is over 256 characters, because a cut label
could end partway through a secret. `f.run` has no label, because a command
can carry literal tokens and URLs. The fields are optional; a step with none
writes the original record shape.

An authored step-failure JSON report keeps the child in `runId` and adds
`rootRunId` for the durable authored root. Consumers must use `rootRunId` for
the resume pointer and root index, and `runId` for the failing child's evidence.
Expand Down
17 changes: 14 additions & 3 deletions packages/sdk/src/authored-flow-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ import {
verifyAuthoredOperations,
} from './authored-flow-operation.js';
import { AuthoredFlowLifecycle } from './authored-flow-lifecycle.js';
import { displayLabel, type AuthoredStepEdges } from './authored-step-index.js';
import { JournalClient } from './journal-client.js';
import { PluginError } from './plugin-manifest.js';
import { createHookEvaluator } from './authored-hooks.js';
Expand Down Expand Up @@ -77,7 +78,7 @@ type EveryRunCompletionReasonIsAcceptedByDone = Assert<

export { AuthoredFlowExecutionError, type AuthoredFlowExecutionErrorCode };

export interface AuthoredFlowJournalStep {
export interface AuthoredFlowJournalStep extends AuthoredStepEdges {
readonly id: string;
readonly runId: string;
readonly completionReason: ProtocolCompletionReason;
Expand Down Expand Up @@ -239,19 +240,21 @@ export async function executeAuthoredFlow<Input = undefined>(
const journalSteps: AuthoredFlowJournalStep[] = [];
const authoredSteps: AuthoredFlowOperation<unknown>[] = [];
const lifecycle = new AuthoredFlowLifecycle();
const stepEdges = (step: string): AuthoredStepEdges | undefined => lifecycle.stepEdges(step);
let nextStep = 1;
let requestedCompletion: LoweredCompletionReason | undefined;

const lowerDeterministic = authoredDeterministicRunner(
definition.name, journal, journalSteps, budget, {
...(options.rootRunId === undefined ? {} : { rootRunId: options.rootRunId }),
...(options.dataDir === undefined ? {} : { dataDir: options.dataDir }),
stepEdges,
},
);

const worker = authoredWorkerRunner(
definition, journal, flowPath, journalSteps, waitOptions,
localAgentStream, budget, definition.header.budget, options.rootRunId, options.workerCapacity,
localAgentStream, budget, definition.header.budget, options.rootRunId, options.workerCapacity, stepEdges,
);

/**
Expand Down Expand Up @@ -365,6 +368,7 @@ export async function executeAuthoredFlow<Input = undefined>(
return worker.llm(id, text, undefined, llmOp.namedGate);
}, onProgress),
lifecycle,
{},
);
return trackStep(authoredSteps, llmOp);
}
Expand Down Expand Up @@ -442,13 +446,17 @@ export async function executeAuthoredFlow<Input = undefined>(
() => assertOperationAllowed('run', definition.name, requestedCompletion),
() => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress),
lifecycle,
// No label: a command is not a display name. It carries literal
// tokens and URLs, and no prefix of it is safe to show (displayLabel).
{},
);
return trackStep(authoredSteps, runOp);
},
llm: llmOperation,
agent(name, options) {
assertOperationAllowed('agent', definition.name, requestedCompletion);
void name; // Authored headers do not yet declare reusable named agents.
// Authored headers do not yet declare reusable named agents, so the name
// identifies nothing to the kernel; it is the step's label in the DAG.
const id = `agent-${nextStep++}`;
let agentOp!: AuthoredFlowOperation<AgentResult>;
agentOp = new AuthoredFlowOperation<AgentResult>(
Expand All @@ -457,6 +465,7 @@ export async function executeAuthoredFlow<Input = undefined>(
() => assertOperationAllowed('agent', definition.name, requestedCompletion),
() => observeStep(id, 'agent', () => worker.agent(id, options, agentOp.namedGate), onProgress),
lifecycle,
displayLabel(typeof name === 'string' ? name : undefined),
);
return trackStep(authoredSteps, agentOp);
},
Expand Down Expand Up @@ -518,6 +527,7 @@ export async function executeAuthoredFlow<Input = undefined>(
return recorded.answer;
}, options.onProgress),
lifecycle,
{},
);
return trackStep(authoredSteps, humanOp);
},
Expand All @@ -540,6 +550,7 @@ export async function executeAuthoredFlow<Input = undefined>(
return verdict;
},
lifecycle,
displayLabel(typeof name === 'string' ? name : undefined),
));
},
done(reason) {
Expand Down
47 changes: 45 additions & 2 deletions packages/sdk/src/authored-flow-lifecycle.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
import { AsyncLocalStorage, executionAsyncId } from 'node:async_hooks';
import { AuthoredFlowExecutionError } from './authored-flow-error.js';
import { AuthoredPromiseGraph } from './authored-promise-graph.js';
import { AuthoredStepGraph } from './authored-step-graph.js';
import type { AuthoredStepEdges } from './authored-step-index.js';

type OperationToken = object;
type ResolverProbe = Set<number>;
Expand Down Expand Up @@ -66,7 +68,7 @@ const observedCombinators: Record<CombinatorName, Combinator> = Object.fromEntri
const members = collectMembers(values);
if (members === undefined) return native.call(this, values);
const aggregate = native.call(this, members);
activeLifecycle.getStore()?.registerCombinator(members, aggregate);
activeLifecycle.getStore()?.registerCombinator(members, aggregate, name);
return aggregate;
};
Object.defineProperty(observed, 'name', { value: name, configurable: true });
Expand Down Expand Up @@ -117,6 +119,8 @@ export class AuthoredFlowLifecycle {
*/
applyPredicateGate: (<T>(operation: { id: string; predicateGate: unknown }, value: T) => Promise<T>) | undefined = undefined;
private readonly graph: AuthoredPromiseGraph;
private readonly stepGraph: AuthoredStepGraph;
private readonly operationPromises = new Map<OperationToken, number>();
private readonly activeResolverProbes: ResolverProbe[] = [];
private readonly invocations = new Map<OperationToken, AuthoredOperationInvocation[]>();
private readonly promiseAllAggregates = new Map<OperationToken, Set<number>>();
Expand All @@ -135,6 +139,7 @@ export class AuthoredFlowLifecycle {
for (const probe of this.activeResolverProbes) probe.add(asyncId);
},
);
this.stepGraph = new AuthoredStepGraph((asyncId) => this.graph.causesOf(asyncId));
installPromiseAllObserver();
try {
this.graph.enable();
Expand All @@ -152,16 +157,25 @@ export class AuthoredFlowLifecycle {
stepOwners.set(step, { lifecycle: this, operation });
}

registerCombinator(values: readonly unknown[], aggregate: Promise<unknown>): void {
registerCombinator(
values: readonly unknown[],
aggregate: Promise<unknown>,
combinator: CombinatorName,
): void {
const aggregateId = this.graph.idOf(aggregate);
if (aggregateId === undefined) return;
this.graph.registerRoot(aggregateId);
const memberIds = new Set<number>();
const awaitedMembers: number[] = [];
for (const value of values) {
if ((typeof value !== 'object' && typeof value !== 'function') || value === null) continue;
const memberId = this.graph.idOf(value);
if (memberId !== undefined) memberIds.add(memberId);
const owner = stepOwners.get(value);
// A Step is a thenable, not a promise: its member identity for the step
// graph is the operation's own promise.
const awaited = owner?.lifecycle === this ? this.operationPromises.get(owner.operation) : memberId;
if (awaited !== undefined) awaitedMembers.push(awaited);
if (owner?.lifecycle !== this) continue;
let operationAggregates = this.promiseAllAggregates.get(owner.operation);
if (operationAggregates === undefined) {
Expand All @@ -172,6 +186,33 @@ export class AuthoredFlowLifecycle {
}
this.promiseAllGroups.push({ aggregate: aggregateId, members: memberIds });
this.registrations++;
this.stepGraph.registerAggregate(aggregateId, aggregate, combinator, awaitedMembers);
}

/**
* Record an operation's own promise and, for a step the run's DAG draws, its
* place in that DAG. Called synchronously from the operation's constructor:
* the causal predecessors are read from the context that invoked
* `f.agent`/`f.run`/`f.llm`, which is observable only at that moment.
*
* Every operation's promise is recorded so `Promise.all([step, …])` can
* expand to it; only a `node` becomes a frontier the walk stops at. An
* operation that is not drawn (a helper effect) is walked through instead,
* so an edge never names a step the index has no record for.
*/
registerOperation(
operation: OperationToken,
promise: Promise<unknown>,
node?: { readonly step: string; readonly label?: string },
): AuthoredStepEdges | undefined {
const promiseId = this.graph.idOf(promise);
if (promiseId !== undefined) this.operationPromises.set(operation, promiseId);
if (node === undefined) return undefined;
return this.stepGraph.registerStep(node.step, promiseId, executionAsyncId(), node.label);
}

stepEdges(step: string): AuthoredStepEdges | undefined {
return this.stepGraph.edgesOf(step);
}

registerInvocation(
Expand Down Expand Up @@ -304,6 +345,8 @@ export class AuthoredFlowLifecycle {
this.closed = true;
this.graph.disable();
this.graph.clear();
this.stepGraph.release();
this.operationPromises.clear();
this.invocations.clear();
this.promiseAllAggregates.clear();
this.promiseAllGroups.length = 0;
Expand Down
10 changes: 10 additions & 0 deletions packages/sdk/src/authored-flow-operation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import {
AuthoredFlowLifecycle,
type AuthoredOperationInvocation,
} from './authored-flow-lifecycle.js';
import type { AuthoredStepEdges } from './authored-step-index.js';

/** Slice-P kinds the surface `.gate(config)` accepts and the SDK lowers. */
const NAMED_GATE_KINDS = new Set([
Expand Down Expand Up @@ -39,6 +40,11 @@ export class AuthoredFlowOperation<T> {
* function itself is never serialized. `flows check` cannot prove it.
*/
predicateGate: { predicate: (value: T) => boolean; because?: string } | undefined = undefined;
/**
* This step's label and causal predecessors in the run's DAG, captured when
* it was invoked. Undefined for an operation the step index does not draw.
*/
readonly edges: AuthoredStepEdges | undefined;
private state: OperationState = 'created';
private thenInvoked = false;
private rootFailureRecorded = false;
Expand All @@ -55,6 +61,7 @@ export class AuthoredFlowOperation<T> {
private readonly assertCanStart: () => void,
private readonly start: () => Promise<T>,
private readonly scope: AuthoredFlowLifecycle,
node?: { readonly label?: string },
) {
let resolve!: (value: T | PromiseLike<T>) => void;
let reject!: (reason?: unknown) => void;
Expand All @@ -64,6 +71,9 @@ export class AuthoredFlowOperation<T> {
});
this.resolve = resolve;
this.reject = reject;
this.edges = scope.registerOperation(
this, this.promise, node === undefined ? undefined : { step: id, ...node },
);

observeRejection(this.promise, (error) => this.recordRootFailure(error));
const operation = this;
Expand Down
Loading
Loading