The Blueprints & Workflows system provides server-side workflow persistence, a reusable blueprint catalog, and an event-driven execution engine for multi-agent DAG workflows. It uses topological sorting to determine execution order, EventBridge for async node invocation, and DynamoDB for durable execution state. The system supports sequential and parallel execution, conditional branching, per-node retry policies with exponential backoff, and real-time execution progress via GraphQL subscriptions. A separate Step Runner Lambda handles workflow execution independently from the existing Supervisor, preserving backward compatibility while adding DAG-based orchestration. For a task-oriented walkthrough aimed at users and operators, see WORKFLOW_USER_GUIDE.md.
- Blueprints & Workflows
- Table of Contents
- Architecture Overview
- Core Concepts
- Component Map
- How It Works
- Data Flows
- Naming Conventions
- How Access and Permissions Work
- Retry and Resilience
- Compensation Actions: A Worked Example
- Error Handling
- Testing Strategy
- Architectural Decisions
- Best Practice Alignment
- Adding a New Workflow Node Type
┌─────────────────────────────────────────────────────────────────────┐
│ Frontend (React) │
│ ┌──────────────┐ ┌──────────────┐ ┌───────────────────────────┐ │
│ │ Blueprint │ │ Workflow │ │ Execution Controls & │ │
│ │ Catalog │ │ Canvas │ │ History Panel │ │
│ └──────┬───────┘ └──────┬───────┘ └────────────┬──────────────┘ │
│ │ │ │ │
│ └─────────────────┼───────────────────────┘ │
│ │ GraphQL + Subscriptions │
└───────────────────────────┼─────────────────────────────────────────┘
│
┌───────────────────────────┼─────────────────────────────────────────┐
│ AWS AppSync API │
│ ┌─────────────────┼───────────────────────┐ │
│ │ │ │ │
│ ┌──────▼───────┐ ┌──────▼───────┐ ┌────────────▼─────────────┐ │
│ │ Workflow │ │ App │ │ Execution │ │
│ │ Resolver λ │ │ Resolver λ │ │ Resolver λ │ │
│ └──────┬───────┘ └──────┬───────┘ └────────────┬─────────────┘ │
└─────────┼─────────────────┼───────────────────────┼─────────────────┘
│ │ │
┌─────────┼─────────────────┼───────────────────────┼─────────────────┐
│ │ DynamoDB Tables │ │
│ ┌──────▼───────┐ ┌──────▼───────┐ ┌────────────▼─────────────┐ │
│ │ citadel- │ │ citadel- │ │ citadel- │ │
│ │ workflows │ │ apps │ │ executions │ │
│ └──────────────┘ └──────────────┘ └──────────────────────────┘ │
└─────────────────────────────────────────────────────────────────────┘
│ │
┌─────────▼─────────────────────────────────────────▼─────────────────┐
│ EventBridge (citadel-agents-{env}) │
│ │
│ workflow.created │ workflow.published │ workflow.node.invoke │
│ workflow.started │ workflow.node.completed │ workflow.failed │
│ │
│ ┌──────────────────────┐ ┌──────────────────────────────────┐ │
│ │ Step Runner Rule │ │ Subscription Fan-out Rule │ │
│ └──────────┬───────────┘ └──────────────┬───────────────────┘ │
│ │ │ │
│ ┌──────────▼───────────┐ ┌──────────────▼───────────────────┐ │
│ │ Step Runner λ │ │ Subscription Fan-out λ │ │
│ │ (Python 3.14) │ │ (Node.js 24.x) │ │
│ └──────────────────────┘ └──────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────────┘
A DynamoDB item keyed by workflowId representing a directed acyclic graph (DAG) of agent nodes and edges. Each workflow has an orgId for access scoping, a status (DRAFT or PUBLISHED), a serialized definition containing nodes and edges, an optional configuration for integration endpoints and agent properties, and a monotonically increasing version for optimistic locking. Workflows can be bound to Agent Apps via the appId field.
A Workflow item with isBlueprint: true — a reusable, org-agnostic template that can be imported into an App as a customizable Workflow. Blueprints are read-only once published. The system ships with five seed blueprints loaded via a CDK Custom Resource during deployment: four placeholder-agent templates ("Sequential Agent Pipeline", "Parallel Fan-Out", "Conditional Router", "Data Processing Pipeline") whose placeholder- agent IDs must be remapped to real agents before publishing, and one runnable demo ("Echo Demo Workflow", category demo) — two nodes referencing the real seeded demo-echo-agent, which echoes its input, so it passes publish validation and executes end to end.
A DynamoDB item keyed by executionId representing a single run of a Workflow. Tracks overall status (pending, running, completed, failed, cancelled), per-node results in a nodeResults map, the workflowVersion snapshot at start time, and input/output payloads. Executions are immutable with respect to the workflow definition — in-flight executions are not affected by concurrent edits.
Nodes represent agent invocations within the DAG. Each node has a nodeId, agentId, optional retryPolicy, and optional constraints for governance. Edges define directed connections between nodes, determining data flow and execution order. The Step Runner uses Kahn's algorithm for topological sorting to determine execution order.
A WorkflowEdge with an optional condition field containing an expression evaluated against the source node's output. Conditions support operators: equals, notEquals, contains, greaterThan, lessThan, exists. When a condition evaluates to false, the target node and its downstream subgraph are marked as skipped.
A per-node configuration specifying maxRetries, backoffBase (seconds), backoffMax (seconds), and retryableErrors (array of error type strings). The Step Runner retries failed nodes using exponential backoff with full jitter: delay = uniform(0, min(backoffBase × 2^attempt, backoffMax)). When retries are exhausted, the node is marked as failed.
An optional, opt-in per-node compensation block ({tool, args, sideEffecting?}) that describes how to undo a node's side effect. Compensation is inert by default — a workflow with no compensation block on any node, or with configuration.compensation.enabled unset/false, behaves byte-identically to a workflow with no compensation feature at all. When enabled and a terminal failure occurs, the Step Runner walks the completed, side-effecting, compensation-bearing nodes in strict reverse-topological ("unwind") order and dispatches each compensation through the same governed worker seam (deny-list → circuit breaker → approval → CIT-121 idempotency) that governs every other tool call — never a separate, ungoverned path. See Compensation Actions: A Worked Example for the full mechanics, the operator-visible states, and the honest limits of what does and does not compensate today.
Per-workflow settings stored in the configuration field, containing:
| Key | Description |
|---|---|
integrations |
Map of integration endpoint configs keyed by integrationId |
credentials |
Map of credential references keyed by resource identifier |
agentProperties |
Map of agent-specific config keyed by agentId |
parameters |
Map of custom key-value parameters |
The Step Runner passes the Workflow Configuration to each agent node as execution context, allowing agents to use workflow-specific integration endpoints and credentials.
Per-node execution overrides are consumed by the Worker Wrapper via exactly two configuration keys: systemPromptAddition (appended to the agent's system prompt) and modelOverride (Bedrock model ID for the node). At dispatch time, node configuration merges over workflow configuration per key — a key set on the node wins over the same key set at workflow level. Both override keys are size-capped (decision 67caf7b0): systemPromptAddition at 4000 characters by default — configurable through the Worker Wrapper's WORKER_MAX_PROMPT_ADDITION_CHARS environment variable (falls back to 4000 when missing or invalid) — and modelOverride at a fixed 256-character hygiene cap. An over-cap value is skipped entirely with a WARN log carrying the offending length and the effective cap; it is never truncated, and the node still executes without the override. Users set these overrides through the node configuration drawer — see Node Configuration and Execution Overrides.
| Component | File | Purpose |
|---|---|---|
| Workflow Resolver | backend/src/lambda/workflow-resolver.ts |
CRUD for workflows and blueprints, publish validation, import/export, version history |
| App Resolver | backend/src/lambda/app-resolver.ts |
App CRUD, workflow bind/unbind, org-scoped access |
| Execution Resolver | backend/src/lambda/execution-resolver.ts |
Start/cancel execution, execution queries |
| Subscription Fan-out | backend/src/lambda/workflow-progress-fanout.ts |
EventBridge → AppSync mutation bridge for real-time subscriptions |
| Seed Blueprints | backend/src/lambda/seed-blueprints/ |
CDK Custom Resource that loads seed blueprint definitions on deployment |
| Workflows Table | citadel-workflows-{env} |
PK=workflowId, GSIs: OrgStatusIndex (orgId/status), BlueprintIndex (isBlueprint/updatedAt) |
| Apps Table | citadel-apps-{env} |
PK=appId, GSI: OrgIndex (orgId/createdAt) |
| Executions Table | citadel-executions-{env} |
PK=executionId, GSI: WorkflowIndex (workflowId/startedAt) |
| GraphQL Schema | backend/src/schema/schema.graphql |
Workflow, AgentApp, Execution types, WorkflowStatus/AppStatus enums, subscriptions |
| CDK — BackendStack | backend/lib/backend-stack.ts |
DynamoDB tables, Workflow/App/Execution Resolver Lambdas, AppSync data sources |
| CDK — ArbiterStack | backend/lib/arbiter-stack.ts |
Step Runner Lambda, EventBridge rules |
| Component | File | Purpose |
|---|---|---|
| Step Runner | arbiter/stepRunner/index.py |
Lambda handler — event routing for execution lifecycle |
| DAG Module | arbiter/stepRunner/dag.py |
Pure functions: topological_sort, find_root_nodes, find_ready_nodes, find_convergence_nodes, find_downstream_subgraph |
| Condition Module | arbiter/stepRunner/condition.py |
Pure functions: evaluate_condition, resolve_field_path |
| Retry Module | arbiter/stepRunner/retry.py |
Pure functions: calculate_backoff, should_retry |
| Executor | arbiter/stepRunner/executor.py |
Orchestration: start_execution, invoke_node, handle_node_completion, handle_node_failure, cancel_execution |
| Events Module | arbiter/stepRunner/events.py |
EventBridge event publishing helpers |
| Component | File | Purpose |
|---|---|---|
| Blueprint Catalog | frontend/src/components/BlueprintCatalog.tsx |
Blueprint grid with search, category filter tabs, "Use in App" action |
| Workflow Toolbar | frontend/src/components/WorkflowToolbar.tsx |
Canvas toolbar: save to catalog, load from catalog, import/export JSON, validate, clear |
| Node Configuration Panel | frontend/src/components/NodeConfigurationPanel.tsx |
Node drawer: rename, model override, system prompt addition, schema parameters |
| Blueprint Card | frontend/src/components/BlueprintCard.tsx |
Individual blueprint card with name, description, agent count, category tags |
| Blueprint Preview Dialog | frontend/src/components/BlueprintPreviewDialog.tsx |
Read-only ReactFlow canvas showing blueprint node/edge layout |
| Import Blueprint Dialog | frontend/src/components/ImportBlueprintDialog.tsx |
App selection dialog for importing a blueprint |
| Execution Overlay | frontend/src/components/ExecutionOverlay.tsx |
Per-node status overlay on canvas (pending/running/completed/failed/skipped) |
| Execution History Panel | frontend/src/components/ExecutionHistoryPanel.tsx |
Side panel with past executions, per-node results, error details |
| Workflow Config Panel | frontend/src/components/WorkflowConfigPanel.tsx |
Workflow-level configuration editor (integrations, credentials, agent properties, parameters) |
| Condition Editor Panel | frontend/src/components/ConditionEditorPanel.tsx |
Edge condition editor (field, operator, value) |
| Workflow Persistence Hook | frontend/src/hooks/useWorkflowPersistence.ts |
Auto-save with debounce (3s), conflict resolution, offline fallback to localStorage |
| Execution Subscription Hook | frontend/src/hooks/useExecutionSubscription.ts |
onWorkflowProgress subscription hook for real-time node status updates |
| Workflow API Service | frontend/src/services/workflowApiService.ts |
GraphQL client for workflow CRUD |
| App API Service | frontend/src/services/appApiService.ts |
GraphQL client for app CRUD |
| Execution API Service | frontend/src/services/executionApiService.ts |
GraphQL client for execution operations |
- User creates a workflow via
createWorkflowmutation — the Workflow Resolver generates a UUID, setsstatus=DRAFT,version=1,isBlueprint=false(unless explicitly set), and persists to DynamoDB - The WorkflowCanvas auto-saves edits via
updateWorkflowwith optimistic locking (versioncondition), debounced to one save per 3 seconds - On version conflict, the persistence hook reloads the latest version and presents a conflict resolution dialog
- User publishes via
publishWorkflow— the resolver validates the definition (no disconnected nodes, no cycles, all nodes have validagentIdreferences) and updatesstatus=PUBLISHED - Published workflows can be executed; draft workflows cannot be deleted while published
- Each update stores the previous definition in
versionHistoryfor audit and rollback - All CRUD operations emit EventBridge events (
workflow.created,workflow.updated,workflow.deleted,workflow.published) with sourcecitadel.workflows
The canvas (Agentic Studio → Create Agent Blueprints) carries a toolbar implemented in frontend/src/components/WorkflowToolbar.tsx:
| Action | Behavior |
|---|---|
| Save | Saves the canvas to the blueprint catalog — a dialog collects a name and optional category, then creates the blueprint and publishes it so it is immediately usable (loadable and importable) |
| Load | Picks a published blueprint from the catalog, with search; loading replaces the current canvas after a confirmation |
| Import | Loads a workflow from a local JSON file |
| Export | Downloads the current workflow as formatted JSON |
| Validate | Checks the workflow for errors and warnings |
| Clear | Removes all nodes and edges from the canvas |
Autosave persists the canvas to the server continuously (see Workflow CRUD Lifecycle); Save is specifically the save-to-catalog action. The run-controls bar carries Publish, which moves the workflow DRAFT → PUBLISHED and unlocks Run, and History, which shows past executions.
Double-clicking a node (or using its configure action) opens the node configuration drawer (frontend/src/components/NodeConfigurationPanel.tsx):
- Rename the node
- Set execution overrides — a Model override (catalog-driven select) and a System prompt addition (up to 4000 characters, with a live character counter)
- Fill in agent-declared schema parameters, rendered from the agent config's parameter schema when present
At runtime, node configuration merges over workflow configuration per key, and the Worker Wrapper honours only modelOverride and systemPromptAddition, subject to the size caps described in Workflow Configuration — oversized values are skipped with a warning, never truncated, and the node still runs.
- User browses the Blueprint Catalog, which queries the
BlueprintIndexGSI forisBlueprint="true"items sorted byupdatedAt - User clicks "Use in App" on a blueprint card, opening a dialog to select or create an Agent App. Agent slots whose
agentIdcarries theplaceholder-prefix must be remapped to real agents in this dialog — publish validation rejectsplaceholder-references - The
importBlueprintmutation deep-copies the blueprint'sdefinitioninto a new Workflow named<blueprint> (Copy)withstatus=DRAFT,isBlueprint=false, a newworkflowId, and the targetappId - The new workflow's
workflowIdis appended to the target app'sworkflowIdsarray - The imported workflow is fully editable — the blueprint remains unchanged
- Only published blueprints can be imported; draft blueprints are rejected
- User clicks "Run Workflow" on the canvas toolbar (enabled only for
PUBLISHEDworkflows) - The
startExecutionmutation creates an Execution item withstatus=pending, initializes allnodeResultsaspending, snapshots theworkflowVersion, and publishes anexecution.start.requestedEventBridge event - The Step Runner picks up the event, performs topological sort on the DAG, and identifies root nodes (in-degree 0)
- Root nodes are invoked by publishing
workflow.node.invokeevents — the Worker Wrapper picks these up and runs the agent - On node completion, the Worker Wrapper publishes
workflow.node.completed— the Step Runner picks this up, evaluates conditional edges on outgoing connections, and identifies the next ready nodes - For convergence nodes (in-degree > 1), the Step Runner waits for all predecessors to reach
completedorskippedstatus before invoking - Independent branches execute concurrently — multiple
workflow.node.invokeevents are published simultaneously - When all nodes complete, the Step Runner marks the execution as
completedand publishesworkflow.completed - If a node fails and retries are exhausted, the execution is marked as
failedwith the error recorded in thenodeResultsmap
- The Subscription Fan-out Lambda subscribes to all
workflow.*EventBridge events - On each event, it calls the
publishWorkflowProgressAppSync mutation (IAM auth) to trigger theonWorkflowProgresssubscription - Frontend clients subscribe via
onWorkflowProgress(executionId)to receive filtered events for the execution they are monitoring - The Execution Overlay updates node status badges in real-time: gray (pending) → blue spinner (running) → green checkmark (completed) or red X (failed)
- The overlay fades out 10 seconds after execution completes
startExecution mutation
→ Execution Resolver creates Execution item (status=pending)
→ Publishes execution.start.requested to EventBridge
→ Step Runner picks up event
→ Topological sort on DAG
→ Find root nodes (in-degree 0)
→ Invoke root nodes via workflow.node.invoke events
→ Worker Wrapper runs agent, publishes workflow.node.completed
→ Step Runner picks up completion
→ Evaluate conditional edges
→ Find ready downstream nodes
→ Check convergence barriers
→ Invoke ready nodes
→ Cycle repeats until all nodes complete or failure halts execution
→ Step Runner publishes workflow.completed or workflow.failed
Step Runner → workflow.node.invoke (EventBridge)
→ Worker Wrapper picks up event
→ Load agent config from citadel-agents-{env}
→ Resolve tool bindings and scoped credentials
→ Apply workflow configuration as execution context
→ Run agent subprocess
→ On success: publish workflow.node.completed with output
→ On failure: publish workflow.node.failed with error
→ Step Runner picks up result, advances execution
Node A completes with output: { "result": { "status": "approved" } }
→ Step Runner evaluates outgoing edges from Node A:
Edge A→B: condition { field: "result.status", operator: "equals", value: "approved" }
→ resolve_field_path(output, "result.status") → "approved"
→ "approved" equals "approved" → true → invoke Node B
Edge A→C: condition { field: "result.status", operator: "equals", value: "rejected" }
→ "approved" equals "rejected" → false → skip Node C and downstream subgraph
| Entity | Pattern | Example |
|---|---|---|
| Workflows DynamoDB table | citadel-workflows-{env} |
citadel-workflows-dev |
| Apps DynamoDB table | citadel-apps-{env} |
citadel-apps-dev |
| Executions DynamoDB table | citadel-executions-{env} |
citadel-executions-dev |
| Workflow Resolver Lambda | citadel-workflow-resolver-{env} |
citadel-workflow-resolver-dev |
| App Resolver Lambda | citadel-app-resolver-{env} |
citadel-app-resolver-dev |
| Execution Resolver Lambda | citadel-execution-resolver-{env} |
citadel-execution-resolver-dev |
| Step Runner Lambda | citadel-step-runner-{env} |
citadel-step-runner-dev |
| EventBridge bus | citadel-agents-{env} |
citadel-agents-dev |
| EventBridge source (workflows) | citadel.workflows |
citadel.workflows |
| EventBridge source (apps) | citadel.apps |
citadel.apps |
| GraphQL enums | PascalCase | WorkflowStatus, AppStatus |
| GraphQL enum values | UPPER_SNAKE_CASE | DRAFT, PUBLISHED, ACTIVE, ARCHIVED |
| Execution statuses | lowercase | pending, running, completed, failed, cancelled |
| Node result statuses | lowercase | pending, running, completed, failed, skipped |
Every resolver operation follows the same org-scoped access control pattern:
- Extract
userIdfrom AppSync identity (Cognitosub) - Call
AdminGetUserto get thecustom:organizationattribute - Compare against the resource's
orgId - Throw "Access denied" if mismatch
Exception: listBlueprints is org-agnostic — blueprints are shared templates accessible to all organizations.
Optimistic locking prevents concurrent modification conflicts. All state-mutating operations (updateWorkflow, updateApp) use a DynamoDB conditional expression requiring version = :currentVersion and increment version on success. On conflict, the resolver throws "Conflict: workflow was modified concurrently. Please retry.".
Lambda functions receive least-privilege IAM policies:
| Lambda | Permissions |
|---|---|
| Workflow Resolver | Read/write workflows table, read apps table, read agent config table, PutEvents on event bus, AdminGetUser on user pool |
| App Resolver | Read/write apps table, read/write workflows table (bind/unbind), PutEvents on event bus, AdminGetUser on user pool |
| Execution Resolver | Read/write executions table, read workflows table, PutEvents on event bus, AdminGetUser on user pool |
| Step Runner | Read/write executions table, read workflows table, read agent config table, read tools config table, PutEvents on event bus |
| Fan-out Lambda | AppSync invoke (IAM auth) for publishWorkflowProgress mutation |
The onWorkflowProgress subscription uses both @aws_iam (for the Step Runner fan-out) and @aws_cognito_user_pools (for frontend clients). The publishWorkflowProgress mutation is @aws_iam only — only backend Lambdas can trigger subscription events.
All state-mutating operations use version-based optimistic locking. DynamoDB conditional writes check version = :expectedVersion and increment on success. The frontend persistence hook retries on conflict by reloading the latest version and presenting a conflict resolution dialog.
Each workflow node can declare a retryPolicy with maxRetries, backoffBase, backoffMax, and retryableErrors. The Step Runner uses exponential backoff with full jitter:
delay = uniform(0, min(backoffBase × 2^attempt, backoffMax))
The backoff result is always bounded: 0 ≤ delay ≤ backoffMax. When retries are exhausted, the node is marked as failed with the final error and retryCount recorded in the nodeResults map. A workflow.node.retrying event is published on each retry attempt.
The Step Runner is idempotent — re-invoking for the same executionId checks DynamoDB state first and resumes from the last incomplete node rather than restarting. Duplicate EventBridge deliveries produce the same execution state without duplicate node invocations.
Independent branches execute concurrently. If one branch fails, other independent branches continue executing. A convergence node is marked as failed only if a required upstream node failed. Promise.allSettled-style semantics ensure partial failures don't cascade across independent paths.
The Workflow item supports an optional timeout field (seconds). When total execution time exceeds the timeout, the Step Runner cancels all running nodes, marks the execution as failed with a timeout error, and publishes a workflow.failed event. Per-node execution timeout defaults to 60 seconds.
The Fan-out Lambda is triggered by EventBridge rules matching workflow.* events. If the AppSync mutation call fails, the event is retried by EventBridge's built-in retry policy. The frontend auto-reconnects via Amplify on subscription disconnect, showing stale status badges until reconnected.
This section walks a single concrete scenario — a workflow node that files a support ticket — end to end: the definition, what happens when a later node fails, what the operator sees in the execution detail sheet, and the honest limits of what compensation does and does not cover today (CIT-123). For the underlying data model, see Compensation Actions (CIT-123); for the governed-execution seam every compensation runs through, see Node Invocation Flow.
A three-node workflow: create-ticket files a support ticket, notify-oncall pages on-call with the ticket reference, close-out performs a final reconciliation step that can fail (e.g. a downstream system is unreachable or a governance policy denies the call).
{
"configuration": {
"compensation": {
"enabled": true,
"trigger": { "mode": "on_terminal_failure", "minCompletedNodes": 1 },
"onFailure": "stop"
}
},
"nodes": [
{
"id": "create-ticket",
"agentId": "ticketing-agent",
"compensation": {
"tool": "close_ticket",
"args": { "ticketId": "${output.ticketId}", "reason": "workflow rolled back" },
"sideEffecting": true
}
},
{ "id": "notify-oncall", "agentId": "paging-agent" },
{ "id": "close-out", "agentId": "reconciliation-agent" }
],
"edges": [
{ "source": "create-ticket", "target": "notify-oncall" },
{ "source": "notify-oncall", "target": "close-out" }
]
}configuration.compensation.enabled: true is the workflow-level opt-in — without it (or with the field absent entirely), this workflow behaves byte-identically to one with no compensation feature: no unwind, no #comp rows, no new nodeResults keys. minCompletedNodes: 1 means the unwind ceremony is skipped if the workflow fails before even create-ticket completes (nothing to roll back yet).
create-ticket's own compensation (close_ticket) references ${output.ticketId} — a path into create-ticket's own recorded output, not any other node's. This is deliberate: a compensation undoes the side effect of the node it is attached to, using that node's own result (e.g. {"ticketId": "TCK-4471"}), never a cross-node reference. The template grammar is a restricted, non-Turing substitution (output.<dotted.path> / output.items[0]) evaluated by a hand-written resolver — never eval/.format()/Jinja — so there is no code-execution surface in a rollback argument. If ticketId is missing from the recorded output (a malformed or truncated result), the renderer fails closed: the compensation is never dispatched with a fabricated value, and the failure is written to the interim sink below as an unresolved-template failure.
close-out fails terminally (say, a policy DENY — POLICY_DENIED in the failure taxonomy). Because the failure disposition is not in the RETRY_AFTER_HUMAN/CIRCUIT_OPEN carve-out (see Honest limits below), the unwind fires:
- The Step Runner computes the reverse-topological plan over completed, side-effecting, compensation-bearing nodes: only
create-ticketqualifies (notify-oncallhas nocompensationblock, so it is skipped — paging on-call is not something we "undo";close-outitself never completed, so it has no recorded output to compensate from). create-ticket#compis writtencompensatingand dispatched to the worker over the same governed seam every forward tool call uses (deny-list → circuit breaker → approval → CIT-121 idempotency reserve/finalize) — never a separate, ungoverned rollback path.close_ticketruns withticketIdresolved fromcreate-ticket's recorded output. On success,create-ticket#compbecomescompensatedand the execution'scompensationStatusbecomescompleted(this was the only qualifying node, so the unwind is done after one step).- If
close_ticketitself fails (e.g. the ticketing system also denies it, or a breaker is open),create-ticket#compbecomescompensation_failed,compensationStatusbecomespartial, and the unwind stops — seeonFailure: 'stop'below.
The execution's top-level status stays failed throughout — compensation is additive observability on a failed run, never a reclassification of it as successful.
Opening the execution in the execution detail sheet shows:
- The
FAILEDstatus badge, plus a second badge reading "failed · rolled back" (full unwind,compensationStatus: completed) or "failed · rollback incomplete" (compensationStatus: partial— a compensation itself failed and the unwind stopped). - A Compensations section, separate from the node Steps list, listing
create-ticket's compensation entry in unwind order with its own status (Compensating/Compensated/Compensation failed) shown as both a coloured dot and a text label — status is never conveyed by colour alone. - Expanding a failed compensation entry shows the raw error plus the
failureClassandrecommendedActionmirrored from the same failure-taxonomy classification the interim sink recorded (see below) — never a second, independently-derived classification.
RETRY_AFTER_HUMANandCIRCUIT_OPENdo not compensate. Ifclose-outfails with a disposition ofRETRY_AFTER_HUMAN(e.g.APPROVAL_ABSENT— a human approval is still pending) or a classification ofCIRCUIT_OPEN(the target is known-bad right now, not permanently), the unwind never fires at all. Both are "leave it for a human/the target to recover" outcomes, not "give up and roll back" outcomes — compensating here would undocreate-ticket's side effect while a human could still complete the approval, or while the circuit could still close on its own. This is enforced in code (_maybe_trigger_compensation_unwind), not just documented intent.onFailure: 'stop'halts the unwind — there is no partial-continue mode today. If a plan has multiple qualifying compensations and the first one fails, the remaining ones are never dispatched. This is the only supported mode; anonFailure: 'continue'(best-effort, run every compensation regardless of earlier failures) was considered in the design and explicitly deferred, not implemented.- CIT-126 (the recovery queue) does not exist yet. A failed or stopped unwind writes three durable, never-swallowed records — the
#comppseudo-node row, an off-frontier escalation event, and aGovernanceFinding— but nothing automatically retries, drains, or resolves them.compensationSummary.entries[](withfailureClassand a fixedrecommendedActionvocabulary:escalate_to_human/retry_after_target_recovery/manual_review_required) is shaped so a future CIT-126 consumer can drain it directly, but until CIT-126 ships, resolving a stuck compensation is a manual, off-system operation — check the escalation event and the Governance Ledger finding for the classified failure, then act by hand. - Only the failing node's predecessors are compensated — never the failing node itself.
close-outin this example has no recorded output (it failed, it didn't complete), so there is nothing to compensate it with even if it had acompensationblock. - Compensation is sequential, one at a time, in strict reverse-topological order — never parallel. A workflow with a long completed prefix and several qualifying compensations will unwind them one after another, not concurrently.
| Scenario | Behavior |
|---|---|
| Workflow not found | Error('Workflow not found') |
| Access denied (org mismatch) | Error('Access denied') |
| Optimistic lock conflict | Error('Conflict: workflow was modified concurrently. Please retry.') |
| Delete published workflow | Error('Cannot delete a published workflow. Unpublish it first.') |
| Publish validation fails | Returns validation errors (disconnected nodes, cycles, missing agent refs) without changing status |
| Invalid definition JSON | ValidationError with structure errors |
| Blueprint not published | Error('Only published blueprints can be imported') |
| Import target app not found | Error('App not found') |
| Import target app org mismatch | Error('Access denied') |
| Scenario | Behavior |
|---|---|
| App not found | Error('App not found') |
| Access denied (org mismatch) | Error('Access denied') |
| Optimistic lock conflict | Error('Conflict: app was modified concurrently. Please retry.') |
| Workflow already bound to another app | Error('Workflow is already bound to another app') |
| Bind already-bound workflow (same app) | Returns app unchanged (idempotent) |
| Org mismatch on bind | Error('Access denied') — both app and workflow must share orgId |
| Scenario | Behavior |
|---|---|
| Workflow not published | Error('Only published workflows can be executed') |
| Execution not found | Error('Execution not found') |
| Cancel non-running execution | Returns execution unchanged |
| Scenario | Behavior |
|---|---|
| Agent not found in config table | Node marked as failed with "agent not found" error |
| Node execution timeout (>60s) | Node marked as failed with timeout error |
| Workflow-level timeout exceeded | All running nodes cancelled, execution marked as failed |
| Retryable error within policy | Node retried with exponential backoff, workflow.node.retrying event published |
| Retries exhausted | Node marked as failed, execution marked as failed |
| Convergence node — upstream failed | Convergence node marked as failed |
| All conditional edges evaluate false | Downstream subgraph marked as skipped |
| Lambda timeout mid-execution | Execution remains in running with accurate nodeResults; manual retry resumes from last incomplete node |
| Duplicate event delivery | Idempotent — checks DynamoDB state before acting |
| Component | Error Behavior |
|---|---|
| Blueprint Catalog | API error → error message with retry button; empty results → empty state message |
| Execution Overlay | Failed node → error tooltip; subscription disconnect → auto-reconnect |
| Execution History Panel | API error → error message with retry; failed node click → expandable error details |
| Workflow Config Panel | Save failure → error toast, form state preserved |
| Workflow Persistence | Network error → fallback to localStorage, retry on reconnect; version conflict → reload + dialog |
All backend components emit structured JSON log entries:
{
"level": "INFO",
"component": "StepRunner",
"executionId": "exec-abc123",
"workflowId": "wf-xyz789",
"nodeId": "node-001",
"agentId": "agent-007",
"action": "invoke_node",
"timestamp": "2025-01-25T12:15:00Z"
}Retry events include additional fields:
{
"level": "WARN",
"component": "StepRunner",
"executionId": "exec-abc123",
"nodeId": "node-001",
"action": "retry_node",
"retryCount": 2,
"error": "Bedrock throttling",
"nextRetryDelay": 4.7,
"timestamp": "2025-01-25T12:15:30Z"
}All implementation follows strict Test-Driven Development (TDD). Property-based tests use fast-check (TypeScript) and Hypothesis (Python), each with a minimum of 100 iterations per property. Tests are written and verified to fail (red phase) before implementation code is created (green phase).
| # | Property | Test File | What It Validates |
|---|---|---|---|
| P1 | Workflow Definition Round-Trip | backend/src/lambda/__tests__/workflow-definition.test.ts |
JSON.parse(JSON.stringify(JSON.parse(d))) ≡ JSON.parse(d) for all valid definitions |
| P2 | Topological Sort Ordering Invariant | arbiter/stepRunner/__tests__/test_dag_properties.py |
For every edge (u, v): indexOf(u, order) < indexOf(v, order) |
| P3 | Condition Evaluation Determinism | arbiter/stepRunner/__tests__/test_condition_properties.py |
Same inputs always produce same boolean result; operator semantics correct |
| P4 | Backoff Bounds | arbiter/stepRunner/__tests__/test_retry_properties.py |
0 ≤ calculate_backoff(attempt, base, max_delay) ≤ max_delay |
| P5 | Optimistic Lock Conflict Detection | backend/src/lambda/__tests__/workflow-resolver.test.ts |
Update with stale version always fails with ConditionalCheckFailedException |
| P6 | Convergence Node Barrier | arbiter/stepRunner/__tests__/test_dag_properties.py |
Node ready iff all predecessors are completed or skipped |
| P7 | Import/Export Round-Trip | backend/src/lambda/__tests__/workflow-definition.test.ts |
export(import(export(w))) ≡ export(w) excluding server-generated fields |
| P8 | Idempotent Execution Start | arbiter/stepRunner/__tests__/test_executor_properties.py |
Calling start_execution twice produces same state, no duplicate invocations |
| Test File | Covers | Type |
|---|---|---|
workflow-resolver.test.ts |
All CRUD operations, org access, optimistic locking, publish validation | Unit + PBT |
app-resolver.test.ts |
App CRUD, bind/unbind, org access | Unit + PBT |
execution-resolver.test.ts |
Start/cancel execution, state initialization | Unit |
workflow-validation.test.ts |
Publish validation (disconnected nodes, cycles, agent refs) | PBT |
test_dag_properties.py |
Topological sort, root nodes, ready nodes, convergence, downstream subgraph | PBT |
test_condition_properties.py |
Condition evaluation, field path resolution, operator semantics | PBT |
test_retry_properties.py |
Backoff calculation, retry decision logic | PBT |
test_executor_properties.py |
Execution flow, idempotency, parallel branches | Unit + PBT |
test_events_properties.py |
EventBridge event construction and field completeness | PBT |
BlueprintCatalog.test.tsx |
Search, filter, empty state, error state, loading | Unit |
ExecutionOverlay.test.tsx |
Status rendering, transitions, fade-out | Unit |
ExecutionHistoryPanel.test.tsx |
List rendering, expand/collapse, pagination | Unit |
useWorkflowPersistence.test.ts |
Debounce, conflict resolution, offline fallback | Unit |
useExecutionSubscription.test.ts |
Event accumulation, status tracking | Unit |
# Backend unit + property tests
cd backend && npm test
# Backend property tests only
cd backend && npx jest --testPathPatterns="workflow-definition|workflow-resolver|workflow-validation" --no-coverage
# Step Runner property tests (Python)
cd arbiter/stepRunner && python -m pytest __tests__/ -v
# Step Runner with reproducible seeds
cd arbiter/stepRunner && python -m pytest __tests__/ -v --hypothesis-seed=0
# Frontend tests
cd frontend && npm test
# All arbiter tests
pytestbackend/src/lambda/__tests__/
├── workflow-resolver.test.ts # CRUD, org access, optimistic locking
├── workflow-definition.test.ts # Properties P1, P7 — round-trip
├── workflow-validation.test.ts # Publish validation PBT
├── app-resolver.test.ts # App CRUD, bind/unbind
└── execution-resolver.test.ts # Start/cancel execution
arbiter/stepRunner/__tests__/
├── test_dag_properties.py # Properties P2, P6 — topological sort, convergence
├── test_condition_properties.py # Property P3 — condition evaluation
├── test_retry_properties.py # Property P4 — backoff bounds
├── test_executor_properties.py # Property P8 — idempotent execution
└── test_events_properties.py # Event construction
frontend/src/
├── components/__tests__/
│ ├── BlueprintCatalog.test.tsx
│ ├── ExecutionOverlay.test.tsx
│ └── ExecutionHistoryPanel.test.tsx
└── hooks/__tests__/
├── useWorkflowPersistence.test.ts
└── useExecutionSubscription.test.ts
| Decision | Rationale |
|---|---|
| Separate Step Runner Lambda vs extending the Supervisor | The Supervisor is a single-turn orchestrator using Bedrock Converse + SQS agent dispatch. The Step Runner is a multi-step DAG executor with different lifecycle, timeout (5 min vs 30s), and state management needs. Keeping them separate preserves backward compatibility and follows Single Responsibility. |
| EventBridge for node invocation | Event-driven architecture means no single Lambda invocation runs for the entire workflow duration. Each step is independently retryable, the execution state in DynamoDB is the source of truth, and the system naturally handles Lambda timeouts without losing progress. |
| DynamoDB for execution state | Execution state needs durable, low-latency reads and writes with per-item conditional updates. DynamoDB's nodeResults map allows atomic per-node status updates without read-modify-write cycles. PAY_PER_REQUEST billing matches variable execution workloads. |
isBlueprint stored as String in GSI |
DynamoDB GSI partition keys must be String, Number, or Binary — not Boolean. Storing as "true"/"false" enables the BlueprintIndex GSI for efficient blueprint listing. |
| Topological sort via Kahn's algorithm | Deterministic, O(V+E) complexity, naturally detects cycles (raises ValueError), and produces a stable ordering for reproducible execution. Pure function with no side effects, enabling thorough property-based testing. |
| Subscription fan-out via separate Lambda | Decouples the Step Runner (Python) from AppSync mutation calls (Node.js). The fan-out Lambda is lightweight (30s timeout) and only bridges EventBridge events to AppSync subscriptions. |
| Workflow version snapshot on execution | In-flight executions must not be affected by concurrent edits. Snapshotting workflowVersion at start time ensures the Step Runner executes the definition that was current when the execution began. |
| Seed blueprints via CDK Custom Resource | Follows the existing seed-organizations pattern. Blueprints are loaded on deployment, ensuring every environment has the same starting templates without manual setup. |
| Condition evaluation as pure functions | evaluate_condition and resolve_field_path have no side effects, making them trivially testable with Hypothesis property-based tests. Operators are deterministic and composable. |
| Exponential backoff with full jitter | Full jitter (uniform(0, calculated_delay)) provides better spread than equal jitter, reducing thundering herd effects when multiple nodes retry simultaneously. Consistent with the existing CircuitBreaker pattern in arbiter/supervisor/circuit_breaker.py. |
| Auto-save with server-side persistence + localStorage fallback | Server-side persistence via updateWorkflow makes workflows durable and accessible from any device. localStorage fallback handles offline editing gracefully, retrying the server save when connectivity is restored. |
| Pillar | Implementation |
|---|---|
| Security | Org-scoped access control on all resolver operations, least-privilege IAM policies per Lambda, @aws_iam auth on subscription trigger mutations, input validation at resolver layer (defense in depth beyond GraphQL schema) |
| Reliability | Idempotent execution (re-processing same event checks DynamoDB state), optimistic locking for concurrent modifications, partial failure handling in parallel branches, per-node retry with exponential backoff and jitter, workflow-level timeout |
| Operational Excellence | Structured JSON logging with executionId/workflowId/nodeId/agentId fields, X-Ray tracing on all Lambdas, correlation IDs (executionId) across EventBridge events and DynamoDB writes, CloudWatch Logs Insights queries for execution debugging |
| Performance Efficiency | PAY_PER_REQUEST DynamoDB billing, event-driven execution (no long-running Lambdas), parallel branch execution for independent paths, BatchGetItem for agent config lookups |
| Cost Optimization | On-demand DynamoDB, right-sized Lambda memory (1024MB for Step Runner, default for resolvers), no always-on infrastructure, event-driven architecture avoids idle compute |
| Principle | Implementation |
|---|---|
| Single Responsibility | Workflow Resolver handles CRUD, App Resolver handles app management, Execution Resolver handles execution lifecycle, Step Runner handles DAG execution — no module takes on responsibilities belonging to another |
| Open/Closed | New node types can be added without modifying the Step Runner's core DAG traversal logic; new condition operators can be added to the condition module without changing the evaluation framework |
| Interface Segregation | Step Runner's pure function modules (dag.py, condition.py, retry.py) have minimal interfaces — each function takes only the data it needs |
| Dependency Inversion | The executor depends on abstract event publishing and DynamoDB interfaces, not concrete AWS SDK calls — enabling thorough unit testing with mocks |
To add a new type of node to the workflow system:
Add the new node type to the WorkflowNodeDefinition interface in frontend/src/types/workflow.ts:
interface WorkflowNodeDefinition {
id: string;
type: 'agent' | 'condition' | 'your_new_type'; // extend the union
// ... existing fields ...
yourNewTypeConfig?: YourNewTypeConfig;
}Create frontend/src/components/YourNewTypeNode.tsx following the AgentNode.tsx pattern. Register it in the ReactFlow nodeTypes map in WorkflowCanvas.tsx.
In arbiter/stepRunner/executor.py, extend invoke_node to handle the new type:
def invoke_node(execution_id, node, input_data, configuration):
if node['type'] == 'your_new_type':
# Custom invocation logic
result = execute_your_new_type(node, input_data, configuration)
publish_node_completed(execution_id, node['id'], result)
else:
# Existing agent invocation via EventBridge
publish_node_invoke_event(execution_id, node, input_data, configuration)Write Hypothesis tests in arbiter/stepRunner/__tests__/test_your_new_type_properties.py covering:
- The new node type integrates correctly with topological sort
- Retry policies apply to the new node type
- Conditional edges work with the new node type's output format
Extend the publish validation logic in workflow-resolver.ts to validate the new node type's required fields.
If the new node type requires additional input/output types, add them to backend/src/schema/schema.graphql.
Create a seed blueprint demonstrating the new node type in backend/src/lambda/seed-blueprints/.