AgentFieldbuild

Webhooks & streaming

Execution completion webhooks with HMAC-SHA256, SSE streaming, and observability forwarding

Webhook delivery for execution lifecycle events

Get notified when executions complete, stream live events via SSE, and forward everything to your observability stack.

AgentField pushes execution results to your endpoints via HMAC-signed webhooks, streams real-time events over SSE/WebSocket, and can forward all platform events to external observability systems. No polling required.

from agentfield.types import WebhookConfig

# HMAC-signed webhook — no polling, result pushed to your endpoint
execution_id = await app.client.execute_async(
    target="analyzer.process",
    input_data={"document": doc},
    webhook=WebhookConfig(
        url="https://your-api.com/hooks/done",
        secret="whsec_your_hmac_secret",              # HMAC-SHA256 signature
        headers={"Authorization": "Bearer tok_..."},   # Custom headers forwarded
    ),
)
# Control plane retries up to 3x with exponential backoff. Failures land in dead letter queue.

What just happened

  • The execution was submitted once, then completion moved to push delivery instead of polling
  • Signatures and retries were handled by the control plane, not by ad hoc app code
  • The same platform could stream live events and forward them to an external collector

Example completion webhook:

{
  "event": "execution.completed",
  "execution_id": "exec_abc",
  "workflow_id": "run_x7y8z9",
  "status": "succeeded",
  "target": "analyzer.process",
  "type": "reasoner",
  "duration_ms": 14200,
  "result": {
    "summary": "Document processed successfully"
  },
  "timestamp": "2026-03-23T10:00:00Z"
}
What you get
  • Execution completion webhooks -- register a URL at submission time; the control plane delivers the result with HMAC-SHA256 verification
  • Automatic retries -- exponential backoff with dead letter queue for persistently failing deliveries
  • SSE event streams -- real-time execution, node, and memory events over Server-Sent Events
  • Observability forwarding -- stream all platform events to an external webhook (Datadog, Splunk, your own collector)
  • WebSocket support -- memory events available over both SSE and WebSocket
Execution webhooks

Registering a Webhook

Include a webhook object when submitting an execution:

POST /api/v1/execute/async/analyzer.process

{
  "input": { "document": "..." },
  "webhook": {
    "url": "https://your-api.com/hooks/execution-done",
    "secret": "whsec_your_hmac_secret_here",
    "headers": {
      "Authorization": "Bearer tok_...",
      "X-Custom-Header": "value"
    }
  }
}
FieldTypeRequiredDescription
urlstringYesHTTPS endpoint to receive the callback
secretstringNoHMAC-SHA256 secret for signature verification
headersmap[string]stringNoCustom headers to include (max 20 headers, 512 chars each)

Webhook Delivery

When the execution reaches a terminal state (succeeded, failed, cancelled, timeout), the control plane delivers a POST request to the registered URL:

POST https://your-api.com/hooks/execution-done
Content-Type: application/json
X-AgentField-Signature: sha256=a1b2c3d4...

{
  "execution_id": "exec_a1b2c3",
  "run_id": "run_x7y8z9",
  "status": "succeeded",
  "result": { "summary": "..." },
  "duration_ms": 14200,
  "completed_at": "2026-03-23T10:00:15Z"
}

Retry Policy

Failed webhook deliveries are retried with exponential backoff:

SettingDefaultDescription
Max attempts3Total delivery attempts before moving to dead letter queue
Initial backoff1 secondWait time before first retry
Max backoff5 secondsMaximum wait between retries
Timeout10 secondsHTTP request timeout per attempt
Worker count4Parallel webhook delivery workers
Queue size256In-memory delivery queue
Response body limit16 KBMax response body captured for debugging

After all attempts are exhausted, the webhook is moved to the dead letter queue for manual inspection and retry.

Retry from UI

POST /api/ui/v1/executions/{execution_id}/webhook/retry

Manually trigger a webhook re-delivery from the AgentField dashboard.

Signature verification

When a secret is provided, the control plane signs the payload with HMAC-SHA256 and includes the signature in the X-AgentField-Signature header.

Verify it server-side:

import hmac
import hashlib

def verify_webhook(body: bytes, signature: str, secret: str) -> bool:
    # Strip "sha256=" prefix if present
    if signature.startswith("sha256="):
        signature = signature[7:]

    expected = hmac.new(
        secret.encode(),
        body,
        hashlib.sha256,
    ).hexdigest()

    return hmac.compare_digest(signature, expected)
SSE event streams

The control plane exposes real-time event streams over Server-Sent Events for dashboards and monitoring tools.

Execution Events

GET /api/ui/v1/executions/events

Streams execution lifecycle events as they happen:

event: execution.started
data: {"execution_id":"exec_abc","status":"running","agent_node_id":"analyzer","timestamp":"2026-03-23T10:00:01Z"}

event: execution.completed
data: {"execution_id":"exec_abc","status":"succeeded","duration_ms":14200,"timestamp":"2026-03-23T10:00:15Z"}

event: execution.waiting
data: {"execution_id":"exec_abc","status":"waiting","approval_request_id":"req_123","timestamp":"2026-03-23T10:01:00Z"}

event: execution.approval_resolved
data: {"execution_id":"exec_abc","decision":"approved","timestamp":"2026-03-23T10:05:00Z"}

Node Events

GET /api/ui/v1/nodes/events

Streams agent node lifecycle events:

event: node.registered
data: {"node_id":"analyzer","version":"2.1.0","timestamp":"2026-03-23T09:00:00Z"}

event: node.status_changed
data: {"node_id":"analyzer","old_status":"starting","new_status":"ready","timestamp":"2026-03-23T09:00:05Z"}

Memory Events

Available over both SSE and WebSocket:

GET /api/v1/memory/events/sse
GET /api/v1/memory/events/ws

Streams memory set/delete operations for real-time state observation.

Event History

GET /api/v1/memory/events/history

Query past memory events for replay and debugging.

Observability webhook

Forward all platform events (executions, node changes, memory operations) to an external observability system.

Configure

POST /api/v1/settings/observability-webhook
{
  "url": "https://collector.example.com/agentfield/events",
  "secret": "obs_hmac_secret",
  "enabled": true
}

How It Works

The observability forwarder subscribes to all internal event buses and batches events for delivery:

SettingDefaultDescription
Batch size10Max events per HTTP request
Batch timeout1 secondMax wait before sending a partial batch
HTTP timeout10 secondsRequest timeout
Max attempts3Retry attempts per batch
Initial backoff1 secondFirst retry delay
Max backoff30 secondsMaximum retry delay
Worker count2Parallel delivery workers
Queue size1000Internal event queue
Snapshot interval60 secondsPeriodic system state snapshots (0 to disable)

Events that fail all retries are moved to a dead letter queue.

Dead Letter Queue

GET  /api/v1/settings/observability-webhook/dlq      # List failed events
DELETE /api/v1/settings/observability-webhook/dlq     # Clear the queue
POST /api/v1/settings/observability-webhook/redrive   # Replay failed events

Management Endpoints

GET    /api/v1/settings/observability-webhook          # Get current config
POST   /api/v1/settings/observability-webhook          # Set config
DELETE /api/v1/settings/observability-webhook          # Remove config
GET    /api/v1/settings/observability-webhook/status   # Forwarder stats

The status endpoint returns delivery metrics: total forwarded, dropped count, last forward time, and last error.

Patterns

Webhook + Polling Hybrid

Use webhooks as the primary notification channel with polling as a fallback for missed deliveries:

from agentfield.types import WebhookConfig

@app.reasoner()
async def reliable_dispatch(task: str) -> dict:
    # Submit with webhook for push notification
    execution_id = await app.client.execute_async(
        target="worker.process",
        input_data={"task": task},
        webhook=WebhookConfig(
            url="https://api.example.com/hooks/done",
            secret="whsec_...",
        ),
    )

    # Fallback: poll if webhook hasn't arrived in 5 minutes
    result = await app.client.wait_for_execution_result(
        execution_id=execution_id,
        timeout=300,
    )

    if isinstance(result, dict) and "result" in result:
        return result["result"]
    return result

Observability Integration

Forward events to your existing monitoring stack:

# Configure forwarding to Datadog
curl -X POST http://localhost:8080/api/v1/settings/observability-webhook \
  -H "Content-Type: application/json" \
  -d '{
    "url": "https://http-intake.logs.datadoghq.com/api/v2/logs",
    "secret": "",
    "enabled": true,
    "headers": {
      "DD-API-KEY": "your-datadog-api-key"
    }
  }'

# Check delivery status
curl http://localhost:8080/api/v1/settings/observability-webhook/status