/** * FlowAdapter — live flow execution state exposed as a connector resource. * * Holds the canonical `FlowExecution` snapshot for a single run and bridges * three ownership roles: * * - The **agent** drives nodes to completion by calling tools (`start_node`, * `request_approval`, `request_input`, `complete_node`, `skip_node`, * `fail_node`, plus read-only `get_state` / `get_available` / `get_globals` / * `build_handoff`). * - The **runtime** drains functions through `NodeExecutor` and parks first-class * gates without executing a node or starting a model. * - The **host** runner applies out-of-turn * mutations (`applyApproval`, `applyInput`, `retryNode`, * `setAutonomousMode`, `cancel`, `hydrate`) in response to platform actions, * subscribes to `onChange` for persistence, and uses the session stimulus bus * to serialize permitted model turns. * * Every state mutation runs through a single `mutate()` helper that clones * the current snapshot, applies the change, recomputes derived fields via * `computeFlowStateFromSnapshots`, and fires the onChange callback with * the full FlowExecution shape. Consumers downstream (platform gateway, * frontend hook) see latest-wins snapshots. * * See `docs/flow-execution.md` for the universal model. */ import { AbstractConnector, type ConnectContext, type ConnectorDeclaration, type ConnectorHandle, type ToolFace } from "@skaile/workspaces/connectors"; import type { FlowExecution } from "@skaile/workspaces/types"; import type { BindingOrigin } from "./contract/expression.js"; import { type FlowDefinition, isBlocking } from "./engine/index.js"; import { type NodeExecutor } from "./node-executor.js"; import type { InlineSubFlowOperationTarget, InlineSubFlowRunner } from "./sub-flow-runner.js"; import type { SubpromptRunner } from "./subprompt-runner.js"; import { type TurnStimulus } from "./prompt-fragments.js"; /** * Minimal structural shape of the runner's `SessionStimulusBus`. Declared * locally so the adapter doesn't pull in a runtime dep on * `@skaile/workspaces/runner` (which would create a circular workspace import). * * The runner's `SessionStimulusBus` from `runner/src/session-stimulus.ts` * satisfies this structurally; mocks in tests can implement it directly. * * @docLink packages/factory-assets/flow#session-stimulus-bus */ export interface FlowAdapterStimulusBus { signal(connectorId: string, sig: { promptFragment: string; meta?: Record; }): Promise; } /** * Internal state stored in the `ConnectorHandle` for the flow adapter. * @docLink packages/connectors/concepts#flow-adapter-state */ export interface FlowAdapterState { /** The flow graph definition this adapter is executing. */ flowDef: FlowDefinition; /** Current execution snapshot (authoritative; mutated in place by `mutate()`). */ execution: FlowExecution; /** Callback fired after every state mutation for host persistence and event forwarding. */ onChange?: (state: FlowExecution) => void; /** * Optional `SessionStimulusBus` reference. When provided AND * `useBusForStimulus` is true, the adapter computes a `TurnStimulus` * after every mutation and fires `bus.signal()` to wake the agent. * The runner enables this for the serialized session-stimulus turn driver. */ stimulusBus?: FlowAdapterStimulusBus; /** * Connector ID used as the `stimulusBus.signal()` queue key. Defaults to * the adapter handle's declaration ID (e.g. `flow:` in serve mode, * `flow` in CLI mode). */ connectorId?: string; /** See {@link FlowAdapterConnectOptions.useBusForStimulus}. */ useBusForStimulus?: boolean; /** * Last execution snapshot the adapter signaled to the bus. Used to compute * a structural diff in `computeStimulus`, so back-to-back mutations that * leave the snapshot unchanged don't fire spurious turn kicks. */ lastSignaledExecution?: FlowExecution; /** Deterministic runtime backed by the already-materialized session workspace. */ nodeExecutor?: NodeExecutor; /** Isolated model-call runtime for subprompt nodes — a second driver instance, never the session driver. */ subpromptRunner?: SubpromptRunner; /** Host-owned inline child lifecycle, backed by the existing materialized session. */ inlineSubFlowRunner?: InlineSubFlowRunner; /** Transient child operation routes keyed by their owning parent sub-flow node. */ inlineOperationTargets: Map; /** Active disposable subprompt calls, cancelled and joined with their owning flow generation. */ subpromptOperations: Map; }>; /** Cancelled child operations that the next asynchronous lifecycle boundary must join. */ subpromptCleanup: Set>; /** Non-durable provenance overrides for values received from a parent sub-flow binding. */ flowInputOrigins?: Record; /** Non-durable ancestry used to reject recursive inline delegation. */ inlineSubFlowAncestry: string[]; /** Non-durable ownership epoch for results returned by in-flight inline children. */ executionGeneration: number; /** * Non-durable opportunistic wake timer for agent-node deadline re-evaluation. * The derived deadline (see `engine/node-deadline.ts`) is the source of truth; * this timer's only job is to trigger a re-evaluation call while nothing else * would (a genuinely hung agent turn has no other event to react to). Rearmed * after every mutation, `connect()`, and `hydrate()`; cleared on `disconnect()`. */ wakeTimer: ReturnType | null; } /** * Options the runner passes via `declaration.options` when connecting the `FlowAdapter`. * @docLink packages/connectors/concepts#flow-adapter-connect-options */ export interface FlowAdapterConnectOptions { flow: FlowDefinition; /** * Either a fresh start seed (runId + startedBy) or a full rehydration * state. If `execution` is provided, it replaces the default empty * state — used on cold-container rehydration when the runner bootstraps * the flow connector from a host `connector_mutate { id: "flow", * op: "hydrate", payload: { state, flow } }` (Protocol >= 3.5). */ execution?: FlowExecution; seed?: { runId: string; startedBy: string; autonomousMode?: boolean; }; /** * Optional `SessionStimulusBus` to which this adapter should fire * `signal()` calls on every mutation when `useBusForStimulus` is true. * Threaded through by the runner's session-construction code (`runFlow` * / `serve.ts`). */ stimulusBus?: FlowAdapterStimulusBus; /** * When `true` AND `stimulusBus` is provided, the * adapter computes a turn-kicking stimulus after every mutation and * calls `bus.signal()` (fire-and-forget). When `false` (the default), * only the `onChange` callback fires. Default: `false`. */ useBusForStimulus?: boolean; /** Runtime for function nodes. Required before deterministic nodes are drained. */ nodeExecutor?: NodeExecutor; /** Runtime for subprompt nodes. Required before subprompt nodes are drained. */ subpromptRunner?: SubpromptRunner; /** Host lifecycle used to resolve and drive inline child flows in the current session. */ inlineSubFlowRunner?: InlineSubFlowRunner; /** @internal Provenance for child flow inputs, retained only in the live adapter. */ flowInputOrigins?: Record; /** @internal Flow identities from the root through this transient child. */ inlineSubFlowAncestry?: string[]; } /** * Connector adapter that holds and drives a live `FlowExecution`. * * The agent owns agent-node tools, the runtime drains function nodes and parks first-class * gates, the host applies durable human decisions, and the runner subscribes to state changes * to decide whether a model turn is needed. Every mutation runs through one recomputation * pipeline and publishes the resulting snapshot. * @docLink packages/connectors/concepts#flow-adapter */ export declare class FlowAdapter extends AbstractConnector { readonly name = "flow"; readonly mountable = false; tools: ToolFace; constructor(); connect(declaration: ConnectorDeclaration, _ctx: ConnectContext): Promise; disconnect(handle: ConnectorHandle): Promise; private cancelInlineTargets; private cancelSubprompts; private projectInlineInteraction; onStateChange(handle: ConnectorHandle, cb: (state: FlowExecution) => void): void; /** * Replace the current execution snapshot with a host-provided one. * Used on container rehydration via the extended `configure` command. * Publishes only when compatibility preparation changes the supplied snapshot, * such as seeding a node kind activated after the snapshot was persisted. */ hydrate(handle: ConnectorHandle, execution: FlowExecution): void; private read; private list; private search; private describeOperations; executeOp(handle: ConnectorHandle, operation: string, args: Record): Promise; private opGetAvailable; private opValidateNode; private opStartNode; private opRequestApproval; private opRequestInput; private opCompleteNode; private opSkipNode; private opFailNode; private opApplyApproval; private opApplyInput; private opSetAutonomousMode; private opRetryNode; private opHydrate; private opStart; private doStartNode; private doRequestApproval; private doRequestInput; private doCompleteNode; private doSkipNode; private doFailNode; /** * Record an approval decision on a node currently `awaiting_approval`. * A first-class gate completes directly when approved and remains parked when rejected. * Agent-node approvals keep their existing resume behavior. Approving an escalated check * is a human override: it demands a reason and completes the node so the run continues, * but records only the decision — the attempt's evidence stays `passed: false`. * Rejecting a check is terminal. */ applyApproval(handle: ConnectorHandle, nodeId: string, decision: "approved" | "rejected", feedback: string | undefined, decidedBy: string): void; /** * Record a user-provided input on a node currently `awaiting_input`. * Transitions back to the status the node had before requesting input. */ applyInput(handle: ConnectorHandle, nodeId: string, response: unknown, providedBy: string): void; /** * Transition a recoverable failed node back to `available` for retry. * Preserves the failure in `errorHistory` so the next turn's agent can * see the prior failure context. * * `failed` is the only status accepted; a `skipped` node is refused * deliberately, because a skip has already released its dependents and this * operation owns no cascade machinery to unwind the ones that ran without * the producer's output. Reasoning in * `_devlog/entries/2026-08-15-optional-node-skip-semantics.md`. */ retryNode(handle: ConnectorHandle, nodeId: string): void; /** Toggle autonomous mode at the flow-execution level. */ setAutonomousMode(handle: ConnectorHandle, enabled: boolean): void; /** * Cancel the flow run. Status transitions to `cancelled`; dependents stay * put. Bumps `executionGeneration` — the sole late-settlement rejection * mechanism every in-flight await (function executor, sub-flow child) * checks after resolving, so a result that arrives after cancellation can * never write state. */ cancel(handle: ConnectorHandle): void; /** Read-only access to the current FlowExecution snapshot. */ getExecution(handle: ConnectorHandle): FlowExecution; /** @internal Whether a nested inline descendant owns the current agent turn. */ hasActiveInlineAgentTurn(handle: ConnectorHandle): boolean; /** Snapshot currently visible through agent-facing tools, including a live inline child. */ getAgentExecution(handle: ConnectorHandle): FlowExecution; /** Agent node IDs whose approval policy remains mandatory in autonomous mode. */ getMandatoryAgentApprovalNodeIds(handle: ConnectorHandle): string[]; /** @deprecated Use {@link getMandatoryAgentApprovalNodeIds}. */ getMandatoryGates(handle: ConnectorHandle): string[]; /** * Decide whether the current state and stimulus require the session agent's conversation. * Runtime-owned gates stay dark while pending; a rejected gate is the exception because the * rejection starts the revision conversation promised by the gate lifecycle. */ shouldDriveAgentTurn(handle: ConnectorHandle, stimulus: TurnStimulus): boolean; /** * Drain runtime-owned nodes in definition order until only agent work or a human decision * remains. Routers select branches, function nodes use NodeExecutor, gates park, and sub-flows * delegate to the host. */ drainRuntimeNodes(handle: ConnectorHandle): Promise; /** * Run every currently available function and check node through the * configured NodeExecutor. Newly-unblocked function/check nodes are * drained in definition order before returning control to the runner's * agent-turn decision. Routers and gates are excluded — see * {@link drainRuntimeNodes} for the full drain. */ executeAvailableFunctionNodes(handle: ConnectorHandle): Promise; private drainAvailableRuntimeNodes; /** Settle one decidable router per source ownership component before runtime side effects. */ private settleAvailableRouters; private s; /** * Central mutation pipeline: runs the mutator, recomputes derived * fields, updates the flow-level `status` if the mutation completed * or failed the flow, fires `onChange` once, and optionally signals the * {@link FlowAdapterStimulusBus}. The signal path is gated by the * `useBusForStimulus` flag set at connect time. */ private mutate; /** * Recompute derived fields (`focus`, `done`, and the rolled-up flow * `status`) from the current node map, without touching per-node * state directly. */ private recompute; private clearWakeTimer; /** * Whether `execution` counts as "actively running" for `control.timeoutSec` * evaluation. Ordinarily this is just `status === "running"`. The one * exception: a non-gate agent node whose approval was just decided * `approved` stays at `awaiting_approval` — `applyApproval` deliberately * does not restore `running` there, because `doCompleteNode`'s * `approvedPrior` check requires exactly that status (a mandatory-approval * node has no `runningAutonomous` fallback). Without this second clause the * timeout would pause for the approval wait (correctly, via * `remainingActiveExecutionMs` reading the now-resolved `approvalHistory` * entry) and then never resume, since nothing re-evaluates a node that * never looks "running" again. */ private isDeadlineEligible; /** * Recomputes the pending wake timer against the current state: cancels any * previously-armed timer, then arms a fresh one for the earliest deadline * among currently `running` agent nodes with a `control.timeoutSec`. A * no-op when nothing is running against a timeout, or the flow is already * terminal. Called after every {@link mutate}, plus `connect()` and * `hydrate()` (which bypass `mutate()`). * * `disconnect()` clears the timer and bumps `executionGeneration`, but a * `mutate()` already in flight (or one that lands after teardown for any * other reason) still calls this on its way out — bail here on * `handle.status` so a post-disconnect mutation can never re-arm a timer * that would fire against a handle nothing is watching anymore. */ private rearmWakeTimer; /** * Re-checks every `running` agent node's derived deadline and, for any that * have expired, either auto-retries (bounded by `control.retries`, composing * with the existing `doFailNode(recoverable: true)` and `retryNode` * pathways) or terminally fails the node. Safe to call at any time — the * derived deadline means this is a no-op unless something has actually * expired. Goes through {@link mutate}, so a state change here republishes * `onChange`/stimulus exactly like any other mutation, waking a fresh turn * when the node retries or the flow rolls up to `failed`. */ private evaluateAgentNodeDeadlines; /** * Records one failed attempt on `nodeId`'s execution and decides whether * the automatic bounded retry (`control.retries`) still has budget. Shared * by every automatic-failure seam (agent-node timeout, function/check * execution/timeout failure, sub-flow execution/timeout failure, and a * check's fail verdict) so the budget — derived from prior `'failed'` * entries in `outputHistory`, not a separate durable counter — is computed * once, the same way, everywhere. * * Returns the `NodeExecution` patch (error/errorHistory/outputHistory/ * output) plus whether budget remains. The terminal *failure* status is * uniform across kinds and computed by {@link resolveTerminalFailureStatus}; * the caller still writes the status because its non-failure outcomes differ * (a check's escalation parks at `awaiting_approval`, a surviving budget * returns to `available`). A terminal failure whose node has any * `control.retries` configured carries `recoverable: true` so the manual * `retryNode` operator override — which sits outside this budget entirely — * remains usable once the automatic budget runs out; a node that will be * written `skipped` never claims it, because a skip is not retryable. * * @param sourceErrorAtIso - The failure's own timestamp, when the caller has * one more precise than `completedAtIso` (e.g. a sub-flow child reports * its own `error.at`, which can predate this node's parent-side * `completedAt` by however long the result took to propagate). Defaults * to `completedAtIso` for callers (agent-timeout, function-node) whose * failure message is synthesized at the moment of detection. * @param checkEvidence - Derived evidence for a check's fail verdict, so it * rides on the same appended attempt as the failure it explains. */ private planAutomaticRetry; /** * Build the selected agent node's complete turn briefing plus upstream context. */ private buildHandoff; } export { isBlocking }; /** * Creates a new FlowAdapter instance. * @returns Configured FlowAdapter ready to be registered in the connector registry. * @docLink packages/factory-assets/api-reference#flow-connector-factory */ export declare function createConnector(): FlowAdapter; //# sourceMappingURL=adapter.d.ts.map