Skip to main content

Run Infrastructure

How workflow runs execute, stream output, pause for human input, and resume.


Starting a run​

Workflows can be triggered in two ways:

MethodWhereWho
Run panelInside the Agent Builder editor (▶ button)Editors testing drafts
REST APIPOST /api/v1/agent-workflows/{id}/runApplications, 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:

LimitWindowBehavior when exceeded
20 runs1 minuteThe 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​

StateMeaning
runningActively executing nodes
waiting_decisionPaused at a review-mode agent node or a user_approval node; an inbox task exists and awaits a reviewer's decision
ready_to_resumeThe reviewer has decided; the run is claimable by GET /runs/{run_id}/stream to continue executing
completedAll paths reached an end node
failedAn unhandled error occurred; run terminated
rejectedA reviewer rejected the run at a pause point; terminated without completing
cancelledPOST /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:

EventWhen emittedKey fields
node_startWhen a node begins executingnode_id, node_type, label
node_outputFor each LLM token or partial outputnode_id, text
node_completeWhen a node finishesnode_id, output, latency_ms
node_errorWhen a node fails (guardrails, timeout, etc.)node_id, error
node_pausedWhen a user_approval node pauses the runnode_id, run_id, task_id, review_group, expires_at
run_completeWhen all nodes finishfinal_output, total_ms
run_errorUnhandled error or rate limit exceedederror
run_cancelledThe run honored a POST /runs/{run_id}/cancel requestmessage
node_complete ordering

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:

ValueBehavior
halt (default, or unset)A node failure aborts the run — no further nodes execute
skipThe 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 failures output.

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):

KeyMeaning
_auto_escalate_after_failuresNumber of consecutive failed runs that triggers escalation
_auto_escalate_groupSpace 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:

  1. Creates an inbox task for the reviewer.
  2. Emits a run_pause SSE event with the inbox_task_id.
  3. 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_approval node (on_approve or on_reject label).
  • 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:

ModeBehavior
auto (default)Node executes immediately; output moves to the next node
reviewNode 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=1800000000

Response:

{
  "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=1800000000

Response:

{
  "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}/cancel

Same 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 paused run isn't watched while it's paused

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:

  1. Click + New Test Case and give it a name.
  2. Paste or type the input value.
  3. 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:

  1. Claims a job from the stream with XREADGROUP.
  2. Acknowledges it immediately with XACK to prevent double-processing.
  3. Executes the workflow graph.
  4. Checkpoints progress to the database after each node completes.
  5. On success: marks the run completed.
  6. 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_worker

Or with the start script:

./backend/worker/start_worker.sh

Worker 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 variableDefaultDescription
REDIS_URLredis://localhost:6379/0Redis connection for the job stream
WORKER_CONCURRENCY5Max concurrent runs per worker process
WORKER_HEALTH_PORT8081Port for the worker health endpoint
MAX_WORKFLOW_DELIVERY_ATTEMPTS3Max delivery attempts before a job goes to the dead-letter queue
WORKFLOW_STALE_MESSAGE_MS300000Milliseconds a claimed job may sit idle before another worker reclaims it (crash recovery)
WORKFLOW_QUEUE_PARTITION_BY_CATALOGfalseUse a separate queue partition per LLM catalog
WORKFLOW_MAX_TIMEOUT_SECONDS600Upper bound on a workflow run's total timeout
Hrida.ai is proprietary software of Zlabs Innovation. See the license for terms. © 2026 Zlabs Innovation.