Run Infrastructure
How workflow runs execute, stream output, pause for human input, and resume.
Starting a run
Workflows can be triggered in two ways:
| Method | Where | Who |
|---|---|---|
| Run panel | Inside the Agent Builder editor (▶ button) | Editors testing drafts |
| REST API | POST /api/v1/agent-workflows/{id}/run | Applications, automations |
Run request body:
{
"input": "Summarize the following ticket: ...",
"state": {}
}input is the starting value for the workflow's input state key. state allows you to pre-seed additional keys before execution begins.
Input length limit: input is capped at 50,000 characters. Requests exceeding this limit are rejected with HTTP 422 before the workflow executes.
Rate limiting
To protect shared infrastructure, a sliding-window rate limit is enforced per user:
| Limit | Window | Behavior when exceeded |
|---|---|---|
| 20 runs | 1 minute | The run endpoint returns a run_error SSE event immediately; no workflow executes |
The limit is tracked in Redis. If Redis is unavailable the check is skipped (fail-open). Operators can adjust the limit with the WORKFLOW_RUN_RATE_LIMIT_PER_MINUTE environment variable (set to 0 to disable).
Run states
| State | Meaning |
|---|---|
running | Actively executing nodes |
waiting_decision | Paused at a review-mode agent node or a user_approval node; an inbox task exists and awaits a reviewer's decision |
ready_to_resume | The reviewer has decided; the run is claimable by GET /runs/{run_id}/stream to continue executing |
completed | All paths reached an end node |
failed | An unhandled error occurred; run terminated |
rejected | A reviewer rejected the run at a pause point; terminated without completing |
cancelled | POST /runs/{run_id}/cancel was honored |
SSE event streaming
While a run is in progress, output is streamed in real time using Server-Sent Events (SSE). The run endpoint (POST /api/v1/agent-workflows/{id}/run) returns a streaming response directly — no separate stream connection is needed:
POST /api/v1/agent-workflows/{workflow_id}/run
Authorization: Bearer <token>
Accept: text/event-stream
Content-Type: application/json
{"input": "Summarize this ticket: ..."}The response includes an X-Run-Id header with the run's ID, available as soon as the response starts — before any node has necessarily completed (unlike learning the run ID from a node_paused event, which only happens if the run actually pauses). Capture it if you need to reattach to this run later via GET /runs/{run_id}/stream (e.g. after a dropped connection) or cancel it via POST /runs/{run_id}/cancel.
Each line in the response is a JSON-encoded event:
data: {"event": "node_start", "node_id": "classify_1", "node_type": "classify", "label": "Classify Ticket"}
data: {"event": "node_output", "node_id": "classify_1", "text": "billing"}
data: {"event": "node_complete", "node_id": "classify_1", "output": "billing", "latency_ms": 412}
data: {"event": "node_start", "node_id": "agent_billing", "node_type": "agent", "label": "Billing Agent"}
data: {"event": "node_output", "node_id": "agent_billing", "text": "Your invoice shows"}
data: {"event": "node_output", "node_id": "agent_billing", "text": " a discrepancy in..."}
data: {"event": "node_complete", "node_id": "agent_billing", "output": "Your invoice shows a discrepancy in...", "latency_ms": 1823}
data: {"event": "run_complete", "final_output": "Your invoice shows a discrepancy in...", "total_ms": 2401}
Event types:
| Event | When emitted | Key fields |
|---|---|---|
node_start | When a node begins executing | node_id, node_type, label |
node_output | For each LLM token or partial output | node_id, text |
node_complete | When a node finishes | node_id, output, latency_ms |
node_error | When a node fails (guardrails, timeout, etc.) | node_id, error |
node_paused | When a user_approval node pauses the run | node_id, run_id, task_id, review_group, expires_at |
run_complete | When all nodes finish | final_output, total_ms |
run_error | Unhandled error or rate limit exceeded | error |
run_cancelled | The run honored a POST /runs/{run_id}/cancel request | message |
node_complete is always emitted before node_paused or node_error on the same node. This guarantees the run timeline has a complete record for every node that executed, even when the run halts.
In the editor, the Run panel connects to this stream automatically and renders token output in real time, with a node execution timeline on the side.
Heartbeat
To keep long-running streams alive through reverse proxies and load balancers (which typically close idle connections after 60 seconds), the server emits an SSE comment line every 20 seconds when no real events are pending:
: heartbeat
SSE comment lines (starting with :) are ignored by SSE parsers but reset the idle timer at the network layer. No action is needed on the client side — EventSource / fetch automatically handle this.
Connection drop detection
If the SSE connection drops mid-run (network interruption, proxy timeout), the editor's Run panel shows a "Connection lost — Reconnect" banner. The underlying workflow run continues unaffected on the server regardless of whether you reconnect.
Clicking Reconnect reattaches to the same run via GET /runs/{run_id}/stream — it does not start a new run. (An earlier version of this button called the run endpoint again, silently duplicating every LLM/tool call the original, still-executing run had already made; the editor now captures the run's ID from the initial X-Run-Id response header specifically so Reconnect can target it.) Buffered events since the connection dropped are not replayed — Reconnect resumes streaming from wherever the run currently is, so you won't see node events that already completed before the drop, only what happens from that point on.
Resilience & fault handling
Production workflows call external, occasionally-unreliable dependencies — LLM providers, MCP servers, knowledge bases. These mechanisms keep a single flaky dependency from taking down an entire run.
Per-node error handling: on_error
Every node accepts an on_error config field:
| Value | Behavior |
|---|---|
halt (default, or unset) | A node failure aborts the run — no further nodes execute |
skip | The failure is still fully reported (node_error SSE event, recorded in the run timeline), but the node completes with empty output and the rest of the graph continues to run |
Use skip on nodes whose output other branches don't depend on — e.g. a best-effort logging/notification step that shouldn't take down the main path if it fails.
Circuit breakers
LLM providers, mcp servers, file_search, and knowledge-base retrieval (used by Skill knowledge injection) are each tracked by an independent circuit breaker, keyed per provider/server/knowledge-base:
- Opens after 5 consecutive failures for that key.
- Auto-resets to half-open after 60 seconds, so a recovered dependency is retried automatically without manual intervention.
- LLM provider breaker state is shared across worker pods via Redis in a Tier-2 deployment — one pod tripping the breaker is immediately visible to the others, instead of each pod having to independently accumulate its own 5 failures.
- While open, the node doesn't hang or crash the run: LLM calls fall back to the run's default provider (see below); other node types complete with a clear
circuit open — too many recent failuresoutput.
Retry with backoff
Transient LLM errors — rate limiting, timeouts, 503/529, "overloaded" — on the main agent / classify / guardrails dispatch path are retried automatically: up to 3 attempts, with exponential backoff (2s, then 4s). Only failures caught before any real output reached the stream are retried — once a response has started streaming to the client, a failure mid-stream is surfaced as-is rather than retried (retrying would mean duplicating or discarding output the client already received).
LLM fallback model
If a node's own configured LLM provider has its circuit breaker open, it automatically falls back once to the run's default LLM credentials (the workflow/catalog-level default) instead of failing outright. There are no fallback chains — the fallback call itself cannot fall back again — so a fully-down default provider still surfaces as a real error rather than looping.
Workflow-level auto-escalation
A workflow that fails repeatedly can automatically raise an inbox task instead of failing silently forever. Configure it via two keys stored on the workflow graph (alongside the graph's existing metadata):
| Key | Meaning |
|---|---|
_auto_escalate_after_failures | Number of consecutive failed runs that triggers escalation |
_auto_escalate_group | Space role or admin group the escalation task is assigned to (default: lifecycle_manager) |
The check runs at the end of every failed run. It fires exactly once per failure streak — the moment the consecutive-failure count first reaches the threshold — so a chronically-broken workflow doesn't spam a new inbox task on every subsequent run after that.
Workflow-level timeout ceiling
Every run is bounded by WORKFLOW_MAX_TIMEOUT_SECONDS (default 600s) as a last-resort guard against a stuck while_loop or hung call. Hitting this ceiling is tracked as its own OTel metric, distinct from an ordinary node failure — see Logging & Tracing → Metrics for hridaai.workflow.timeout_ceiling_hits.total.
Pause & resume
How pause works
When execution reaches a user_approval node, the run transitions to waiting state and the runtime:
- Creates an inbox task for the reviewer.
- Emits a
run_pauseSSE event with theinbox_task_id. - Suspends execution. No further nodes run until a decision is made.
The reviewer sees the task in Dashboard > Inbox (or in any configured notification channel). The task shows the message text and the content from the node's input_key.
Deciding via inbox
The reviewer clicks Approve or Reject in the inbox UI. They can optionally add a comment.
Internally this calls:
POST /api/v1/inbox/{task_id}/decide
{
"decision": "approve",
"comment": "Looks good"
}What happens after the decision
- The decision and comment are written into workflow state (
approval_decision,approval_comment). - The run transitions back to
running. - The runtime follows the matching outgoing edge from the
user_approvalnode (on_approveoron_rejectlabel). - SSE streaming resumes and continues until the next pause or the run ends.
Execution mode (review nodes)
agent nodes support an execution_mode config field:
| Mode | Behavior |
|---|---|
auto (default) | Node executes immediately; output moves to the next node |
review | Node output is held and an inbox task is created for a reviewer to approve before the workflow continues |
review mode is different from user_approval node: it pauses on the output of an existing node, whereas user_approval is a dedicated step you insert explicitly in the graph.
Run history and analytics
All runs are stored with their full state, timeline, and cost metadata.
Viewing runs in the editor
Click the Runs tab in the right panel of the editor to see the run history for the current workflow, including status, latency, and token cost.
Per-workflow analytics API
These endpoints are accessible to any agent_developer role — no admin access required:
List runs:
GET /api/v1/agent-analytics/workflows/{workflow_id}/runs
?status=completed # filter by status
&limit=25 # page size (max 200)
&offset=0 # pagination offset
&start_ts=1700000000 # Unix timestamp filter
&end_ts=1800000000Response:
{
"items": [
{
"id": "run_abc123",
"workflow_id": "wf_xyz",
"status": "completed",
"started_at": 1700001000,
"ended_at": 1700001003,
"duration_ms": 3186,
"triggered_by": "user_id",
"error_message": null,
"prompt_tokens": 512,
"completion_tokens": 128,
"total_tokens": 640,
"cost_usd": 0.0012
}
],
"total": 42
}Aggregate summary:
GET /api/v1/agent-analytics/workflows/{workflow_id}/summary
?start_ts=1700000000
&end_ts=1800000000Response:
{
"total_runs": 42,
"completed_runs": 38,
"failed_runs": 3,
"paused_runs": 1,
"success_rate": 0.9048,
"avg_duration_ms": 3186.0,
"p95_duration_ms": 7842.0,
"total_prompt_tokens": 21504,
"total_completion_tokens": 5376,
"total_cost_usd": 0.0504
}For admin-level analytics across all workflows, catalogs, and node types, see the Admin > Analytics section of the documentation.
For structured logs, distributed traces, and OTel metrics for runs and nodes, see Logging & Tracing.
Cancelling a run
POST /api/v1/agent-workflows/runs/{run_id}/cancelSame access rule as the run-status endpoints: only the run's owner or an admin. Returns 404 if the run doesn't exist, 409 if it's already in a terminal state (completed/failed/rejected/cancelled — nothing to cancel), 503 if Redis is unavailable (cancellation isn't silently accepted as a no-op).
This requests cancellation — it doesn't cancel synchronously. The endpoint sets a Redis flag (hrida:cancel:{run_id}); the process actually executing the run (the API process for a direct/SSE run, or a Tier-2 worker pod) checks that flag on a timer — the same cadence as the WORKFLOW_MAX_TIMEOUT_SECONDS guard, throttled to once every WORKFLOW_CANCEL_CHECK_INTERVAL_SECONDS (default 2s) — and marks the run cancelled itself once it sees the flag. In practice this means a cancel typically takes effect within a couple of seconds, not instantly mid-node.
A run sitting in waiting_decision (paused at a review/user_approval node) has no active execution loop checking the cancel flag — pausing exits that loop entirely. Cancelling a paused run sets the flag, but the run's status won't flip to cancelled until it's actually resumed (POST /runs/{run_id}/stream after an inbox decision), at which point resume_workflow sees the flag immediately and stops before running any further nodes. If you need a paused run gone for good, reject it via its inbox task instead of cancelling it.
Inbox tasks are not automatically marked cancelled as a side effect of cancelling their run — an inbox reviewer will still see the task until it's separately resolved or expires.
Test cases
From the editor's Test Cases panel, you can save named input fixtures and replay them on demand:
- Click + New Test Case and give it a name.
- Paste or type the input value.
- Click Run on the test case to execute the workflow with that input.
Test cases are stored against the workflow (not a specific version) and can be run against any draft or published version. They are useful for regression testing after edits.
Distributed execution (Tier-2)
For production deployments handling many concurrent workflow runs, hrida-ai-studio supports a distributed worker model built on Redis Streams.
How it works
API server → Redis Stream (XADD) → Worker node 1
→ Worker node 2
→ Worker node N
When a workflow run is started, the API server enqueues a job into a Redis Stream (hrida:workflow_jobs). Worker nodes run as separate processes (or containers) and compete to claim jobs. Each worker:
- Claims a job from the stream with
XREADGROUP. - Acknowledges it immediately with
XACKto prevent double-processing. - Executes the workflow graph.
- Checkpoints progress to the database after each node completes.
- On success: marks the run
completed. - On crash: the checkpoint allows another worker to resume from the last completed node.
Crash recovery
Each node's output is saved to the workflow_run table as it completes. If a worker crashes mid-run, a background recovery task (XAUTOCLAIM) detects the orphaned job and re-queues it. The new worker resumes from the last checkpoint — only nodes after the checkpoint are re-executed.
Dead-letter queue
Jobs that fail repeatedly (more than the configured retry limit) are moved to hrida:workflow_jobs_dlq. Admins can inspect dead-letter runs in the database and re-queue them manually if needed.
Per-catalog queue partitioning
In multi-tenant deployments, each LLM catalog can have its own queue partition (e.g. hrida:workflow_jobs:{catalog_id}). This prevents a high-volume catalog from starving others and allows per-catalog worker pools.
Per-node LLM timeout
Each agent node has a hard timeout enforced by asyncio.timeout(). If an LLM call takes longer than the configured limit (default: 120 seconds), the node fails with a timeout error and the run is marked failed. The timeout is configurable per node in the node properties panel under Advanced → LLM Timeout (s).
Scaling out
Add more workers by starting additional worker processes (or replicas in Kubernetes). Workers are stateless — they read from Redis, write to Postgres, and share nothing else. No coordination layer is required.
# Start a worker (one per process/container)
python -m hrida_ai_studio.worker.hrida_workflow_workerOr with the start script:
./backend/worker/start_worker.shWorker count can be scaled independently of the API server. A typical production setup runs 2–4 API server replicas and 4–8 workers.
Configuration
| Environment variable | Default | Description |
|---|---|---|
REDIS_URL | redis://localhost:6379/0 | Redis connection for the job stream |
WORKER_CONCURRENCY | 5 | Max concurrent runs per worker process |
WORKER_HEALTH_PORT | 8081 | Port for the worker health endpoint |
MAX_WORKFLOW_DELIVERY_ATTEMPTS | 3 | Max delivery attempts before a job goes to the dead-letter queue |
WORKFLOW_STALE_MESSAGE_MS | 300000 | Milliseconds a claimed job may sit idle before another worker reclaims it (crash recovery) |
WORKFLOW_QUEUE_PARTITION_BY_CATALOG | false | Use a separate queue partition per LLM catalog |
WORKFLOW_MAX_TIMEOUT_SECONDS | 600 | Upper bound on a workflow run's total timeout |