Subscribing to Cellario OS Event Service
A guide for OS customers: receiving events via Webhook or WebSocket
1. Introduction
The Cellario OS Event Service is a centralized, Azure-hosted service that captures and distributes events generated across the Cellario OS platform. As work happens in the lab — orders are scheduled, jobs run, devices report status — the components that make up Cellario OS publish structured events to the Event Service. As a customer, you can subscribe to the events you care about and receive them in your own systems in near real time.
What an event looks like
Every event is a JSON document with a common envelope plus an optional event-specific payload. Key fields include:
Field | Description |
|---|---|
id | Unique numeric identifier for the event |
type | The event type, e.g. OrderCreated, OrderCompleted, JobCompleted |
state | Optional state, e.g. pending, completed, failed |
severity | One of debug, info, warning, error, critical |
source_category | Originating system area: cellario_os, cellario_edge, integration, other |
source_component | Originating component, e.g. cellario_scheduler, cellario_agent, driver_host, job_orchestrator, data_access_api |
source_name / source_id | Human-readable name and identifier of the source |
source_created_at / received_at | When the event occurred and when the service received it |
payload_type / payload_json | MIME type and event-specific data, carried as a JSON string |
stream_id / sequence | Optional grouping and ordering of related events |
What is a "subscription"?
A subscription tells the service which events to deliver to your subscriber. Each subscription carries a filter, supplied one of two ways:
- a simple filter — a lightweight routing-key pattern (most often used to match an event type, e.g. all OrderCreated events), or
- a complex filter — a reusable, field-level filter you create separately (match on type, source, severity, or even payload contents) and then reference.
Both are shown in Section 3; the full filter field, operator, and payload-matching reference (with worked examples) is in Appendix A at the end of this guide.
Two ways to receive events
You choose one delivery mechanism per subscriber:
- Webhook — the Event Service pushes events to an HTTPS endpoint you host.
- WebSocket — your client opens and maintains a connection to the Event Service and receives events over it.
Both are described below.
2. Webhooks vs. WebSockets
| Webhook | WebSocket |
|---|---|---|
Direction | Event Service → your endpoint (outbound HTTP POST) | Your client → Event Service (inbound, persistent connection) |
You provide | A public-internet-reachable HTTPS endpoint | A long-running client process |
Connection | Stateless; one HTTP request per delivery | Stateful; a single connection stays open |
Best for | Serverless functions, external systems, low-frequency events, "set and forget" integrations | Live dashboards, real-time UIs, high-frequency streams |
Authentication of delivery | A shared API key you supply, sent by the service as the X-API-KEY header on every call | A single-use connection token issued by the service |
Infrastructure burden | You must expose and secure a public endpoint — this often requires IT involvement (firewall / ingress), and a tunnel such as ngrok for local testing | No inbound connectivity needed (the connection is client-initiated / outbound-only); you must maintain the connection (heartbeats, reconnects) |
Webhook — A webhook is a "reverse API call." Rather than you polling for data, the Event Service calls you: when a matching event occurs, it sends an HTTP POST containing the event to a URL you registered. Your endpoint replies 2xx to confirm receipt. If your endpoint is unreachable or returns an error, the service retries with exponential backoff. Your endpoint must be reachable from the public internet — behind a corporate network this usually means involving IT, or using a tunnel (e.g. ngrok) for testing. If you cannot expose an inbound endpoint, choose WebSocket.
WebSocket — A WebSocket is a persistent, bidirectional connection over a single long-lived TCP session. Your client connects to the Event Service once and then receives events as they happen, with no repeated HTTP handshakes. The connection is client-initiated (outbound-only), so you need no inbound connectivity or public endpoint. The service sends periodic ping heartbeats that your client answers with pong to keep the connection healthy.
3. Common setup: create a subscriber
Regardless of delivery mechanism, the configuration calls are REST APIs hosted in Azure. They are authenticated with a JWT bearer token issued to your OS account, and require the EVENTS_WRITE role.
Step 1 — Configure the subscriber (choose delivery type)
This is shown per-mechanism in Sections 4 and 5. The configure call creates the subscriber if it doesn't already exist, and activates delivery.
Step 2 — Create a subscription
A subscription carries a filter, supplied as exactly one of simple_filter or complex_filter_id (not both, not neither).
Option A — simple filter (match by type)
POST /v1/subscribers/{subscriberId}/subscriptions
Authorization: Bearer {jwt_token}
Content-Type: application/json
{
"simple_filter": "OrderCreated.#"
}simple_filter is a RabbitMQ topic pattern matched against the event's routing key: * matches exactly one segment, # matches zero or more.
Option B — complex filter (match by field)
For field-level matching, first create a complex filter, then reference it from the subscription.
POST /v1/subscribers/{subscriberId}/filters
Authorization: Bearer {jwt_token}
Content-Type: application/json
{
"filter": {
"type": { "in": ["OrderCreated"] },
"source_category": { "in": ["cellario_os"] }
}
}The response contains a filter_id. Use it when creating the subscription:
POST /v1/subscribers/{subscriberId}/subscriptions
Authorization: Bearer {jwt_token}
Content-Type: application/json
{
"complex_filter_id": "<filter_id from the previous response>"
}Inside a complex filter you can match multiple types ("type": { "in": ["OrderCreated", "OrderCompleted"] }), filter by severity or source_component, or filter on payload contents (payload_has_keys, payload_matches). All field names are snake_case and string operators are in / not_in. A subscriber may hold several subscriptions; an event is delivered if it matches any of them. See Appendix A for the complete field, operator, and payload-matching reference, plus worked examples.
4. Receiving events via Webhook
How it works
- You configure the subscriber with your endpoint URL and an API key.
- When an event matches your subscription, the service sends: POST https://your-endpoint/events with header X-API-KEY: <your key> and the event as a JSON body.
- Your endpoint validates the key, processes the event, and returns 200/204.
- On a 5xx, timeout, or network error, the service retries with exponential backoff. A 4xx (other than 408/429) is treated as permanent and is not retried.
Step 1 — Configure the subscriber for Webhook delivery
POST /v1/subscribers/{subscriberId}/configure
Authorization: Bearer {jwt_token}
Content-Type: application/json
{
"name": "My Webhook Subscriber",
"delivery_mode": "guaranteed",
"delivery_configuration": {
"delivery_type": "webhook",
"web_hook": {
"endpoint_url": "https://your-system.example.com/events",
"api_key": "your-secret-api-key",
"max_retries": 5,
"base_backoff_seconds": 30,
"max_backoff_seconds": 3600
}
}
}- delivery_mode: guaranteed (ordered; a failed delivery blocks later events for this subscriber until resolved) or best_effort (parallel, non-blocking).
- api_key: a shared secret you create here — any opaque string (e.g. openssl rand -hex 16). The Event Service stores it (encrypted) and presents it as the X-API-KEY header on every delivery; your endpoint compares the incoming header against the same value. It is not a credential you obtain or register elsewhere (and is unrelated to Platform-API keys). See Section 6.
- Your endpoint_url must be a publicly reachable HTTPS URL (private IP ranges and cloud metadata addresses are rejected).
Step 2 — Implement your endpoint
Python (Flask)
import os
from flask import Flask, request, jsonify
app = Flask(__name__)
EXPECTED_API_KEY = os.environ["WEBHOOK_API_KEY"]
# Use a durable store (Redis/DB) in production.
processed_events = set()
@app.route("/events", methods=["POST"])
def receive_event():
# 1. Authenticate the caller using the API key.
if request.headers.get("X-API-KEY") != EXPECTED_API_KEY:
return jsonify(error="unauthorized"), 401
# 2. De-duplicate using the event id (retries may redeliver).
event_id = request.headers.get("X-Event-Id")
if event_id in processed_events:
return jsonify(status="already_processed"), 200
# 3. Process the event.
event = request.get_json()
print(f"Received {event['type']} (attempt {event['delivery_attempt']})")
# ... your business logic here ...
# 4. Acknowledge success so the service stops retrying.
processed_events.add(event_id)
return jsonify(status="processed"), 200
if __name__ == "__main__":
app.run(port=8080)C# (ASP.NET Core minimal API)
using System.Text.Json.Serialization;
var builder = WebApplication.CreateBuilder(args);
var app = builder.Build();
var expectedApiKey = Environment.GetEnvironmentVariable("WEBHOOK_API_KEY");
var processed = new HashSet<string>(); // Use a distributed cache in production.
app.MapPost("/events", (HttpRequest req, EventEnvelope evt) =>
{
// 1. Authenticate the caller using the API key.
if (req.Headers["X-API-KEY"] != expectedApiKey)
return Results.Unauthorized();
// 2. De-duplicate using the event id.
var eventId = req.Headers["X-Event-Id"].ToString();
if (processed.Contains(eventId))
return Results.Ok(new { status = "already_processed" });
// 3. Process the event.
app.Logger.LogInformation("Received {Type} (attempt {Attempt})",
evt.Type, evt.DeliveryAttempt);
// ... your business logic here ...
// 4. Acknowledge success.
processed.Add(eventId);
return Results.Ok(new { status = "processed" });
});
app.Run();
// Configure JsonSerializerOptions with JsonNamingPolicy.SnakeCaseLower,
// or annotate each property as shown.
public record EventEnvelope
{
[JsonPropertyName("id")] public long Id { get; init; }
[JsonPropertyName("type")] public string Type { get; init; } = "";
[JsonPropertyName("state")] public string? State { get; init; }
[JsonPropertyName("severity")] public string? Severity { get; init; }
[JsonPropertyName("payload_json")] public string? PayloadJson { get; init; }
[JsonPropertyName("delivery_attempt")] public int DeliveryAttempt { get; init; }
[JsonPropertyName("subscriber_id")] public string SubscriberId { get; init; } = "";
[JsonPropertyName("subscription_id")] public string SubscriptionId { get; init; } = "";
}Delivery semantics to plan for
- Always return 2xx quickly. Do heavy processing asynchronously if needed and return 202 Accepted.
- Be idempotent. Use the X-Event-Id header to ignore duplicate deliveries.
- Reserve 4xx for permanent rejections — returning 4xx (other than 408/429) tells the service to stop retrying.
5. Receiving events via WebSocket
How it works
- You configure the subscriber for WebSocket delivery; the response includes a single-use connection token.
- Your client opens a WebSocket to the connect endpoint, presenting the token.
- The service sends connection.ready, then events.batch messages as events occur.
- The service sends periodic ping; your client must reply with pong.
- If the connection drops, you regenerate the token and reconnect.
Step 1 — Configure the subscriber for WebSocket delivery
POST /v1/subscribers/{subscriberId}/configure
Authorization: Bearer {jwt_token}
Content-Type: application/json
{
"delivery_mode": "guaranteed",
"delivery_configuration": {
"delivery_type": "websocket",
"web_socket": {
"connection_token_expiry_minutes": 60,
"acknowledgment_mode": "auto",
"heartbeat_interval_seconds": 30
}
}
}The response contains delivery_configuration.web_socket.connection_token and its connection_token_expires_at. Use that token to connect.
- acknowledgment_mode: auto (events are considered delivered when sent) or manual (you must ack each batch, or it is redelivered).
Step 2 — Connect and receive events
Python (websockets)
import asyncio, json, os, websockets
SUBSCRIBER_ID = os.environ["SUBSCRIBER_ID"]
CONNECTION_TOKEN = os.environ["CONNECTION_TOKEN"] # from the configure response
BASE = "wss://your-os-tenant.cellario.cloud/api/events"
async def run():
url = f"{BASE}/v1/subscribers/{SUBSCRIBER_ID}/connect"
# Present the connection token as a Bearer credential.
async with websockets.connect(
url, additional_headers={"Authorization": f"Bearer {CONNECTION_TOKEN}"}
) as ws:
async for raw in ws:
msg = json.loads(raw)
if msg["type"] == "connection.ready":
print("Connected:", msg["subscriber_id"])
elif msg["type"] == "events.batch":
for evt in msg["events"]:
print(f"Event: {evt['type']} - {evt.get('state')}")
# Manual ack mode only:
# await ws.send(json.dumps({"type": "ack", "batch_id": msg["batch_id"]}))
elif msg["type"] == "ping":
await ws.send(json.dumps(
{"type": "pong", "timestamp": msg["timestamp"]}))
elif msg["type"] == "error":
print("Server error:", msg["code"], msg["message"])
asyncio.run(run())C# (ClientWebSocket)
using System.Net.WebSockets;
using System.Text;
using System.Text.Json;
var subscriberId = Environment.GetEnvironmentVariable("SUBSCRIBER_ID");
var connectionToken = Environment.GetEnvironmentVariable("CONNECTION_TOKEN");
var baseUri = "wss://your-os-tenant.cellario.cloud/api/events";
using var ws = new ClientWebSocket();
// Present the connection token as a Bearer credential.
ws.Options.SetRequestHeader("Authorization", $"Bearer {connectionToken}");
await ws.ConnectAsync(
new Uri($"{baseUri}/v1/subscribers/{subscriberId}/connect"), CancellationToken.None);
var buffer = new byte[8192];
while (ws.State == WebSocketState.Open)
{
var result = await ws.ReceiveAsync(buffer, CancellationToken.None);
if (result.MessageType == WebSocketMessageType.Close) break;
var json = Encoding.UTF8.GetString(buffer, 0, result.Count);
using var doc = JsonDocument.Parse(json);
var type = doc.RootElement.GetProperty("type").GetString();
switch (type)
{
case "connection.ready":
Console.WriteLine("Connected");
break;
case "events.batch":
foreach (var evt in doc.RootElement.GetProperty("events").EnumerateArray())
Console.WriteLine($"Event: {evt.GetProperty("type").GetString()}");
// Manual ack mode only:
// await Send(ws, new { type = "ack",
// batch_id = doc.RootElement.GetProperty("batch_id").GetString() });
break;
case "ping":
await Send(ws, new { type = "pong",
timestamp = doc.RootElement.GetProperty("timestamp").GetString() });
break;
case "error":
Console.WriteLine($"Server error: {doc.RootElement.GetProperty("code").GetString()}");
break;
}
}
static async Task Send(ClientWebSocket ws, object message) =>
await ws.SendAsync(
Encoding.UTF8.GetBytes(JsonSerializer.Serialize(message)),
WebSocketMessageType.Text, true, CancellationToken.None);Connection lifecycle to plan for
- Answer every ping with a pong (echo the timestamp). If the service misses your pong for 2× the heartbeat interval, it closes the connection.
- Connection tokens are single-use. After a disconnect, request a fresh token (Section 6) and reconnect using exponential backoff.
- No events are lost across reconnects — pending events are redelivered once you reconnect, thanks to the service's outbox.
6. Configuring the API key / token for authentication
Authentication has two distinct layers. The configuration APIs (creating subscribers and subscriptions, regenerating tokens) are always protected by your OS-issued JWT bearer token with the EVENTS_WRITE role. The delivery channel is then secured differently for each mechanism, as described below.
Webhook: a revocable shared API key
For webhooks, the API key is a secret you choose (any opaque string — it is not a Platform-API key and is not provisioned elsewhere) and register in the subscriber's webhook configuration. The Event Service stores it (encrypted with AES-256-CBC) and presents it on every delivery as the X-API-KEY request header. Your endpoint compares the incoming header against the expected value and rejects mismatches with 401.
Register it in the api_key field of the configure call:
"web_hook": {
"endpoint_url": "https://your-system.example.com/events",
"api_key": "your-secret-api-key" // sent as X-API-KEY on each delivery
}Revoke or rotate it by re-running the configure call with a new api_key (or removing it). Because the key is sent on every request rather than baked into a session, rotation takes effect on the next delivery:
- Generate a new secret and deploy it to your endpoint (accept old + new briefly to avoid a gap).
- POST /v1/subscribers/{subscriberId}/configure with the new api_key.
- Retire the old secret from your endpoint.
The key is never echoed back in API responses (it shows as masked), so treat your stored copy as the source of truth.
WebSocket: a revocable single-use connection token
For WebSockets, you do not supply a long-lived key. Instead, the Event Service issues a connection token (32 bytes of cryptographic randomness, Base64, time-limited by connection_token_expiry_minutes). You present it when opening the connection:
GET /v1/subscribers/{subscriberId}/connect
Authorization: Bearer {connectionToken}
Upgrade: websocketBrowser clients, which cannot set custom WebSocket headers, may instead pass it as a query parameter: …/connect?token={connectionToken}. This is safe because tokens are single-use and expire.
Obtain it from the configure response, or regenerate on demand:
POST /v1/subscribers/{subscriberId}/connection-token
Authorization: Bearer {jwt_token}Revoke it simply by regenerating: generating a new token invalidates the previous unused one, and tokens are consumed (invalidated) the moment they are used to connect. Because tokens are short-lived and single-use, replay of a leaked token is prevented by design. To cut off access entirely, archive or reconfigure the subscriber.
At a glance
| Webhook | WebSocket |
|---|---|---|
Credential | Shared API key you choose | Single-use token the service issues |
Where configured | web_hook.api_key in the configure call | Returned from configure / connection-token |
Presented as | X-API-KEY header on each delivery | Authorization: Bearer (or ?token=) on connect |
Authenticates | Event Service → your endpoint | Your client → Event Service |
Revoke / rotate | Re-configure with a new key | Regenerate token; archive subscriber to fully revoke |
At rest | Encrypted (AES-256-CBC, Azure Key Vault) | Not stored long-term; short-lived |
7. Choosing between Webhook and WebSocket
Use a Webhook when you want the simplest possible integration, can host a public HTTPS endpoint (including a serverless function) that is reachable from the public internet, and your event volume is low-to-moderate. The service handles retries and ordering for you. Note that exposing a public inbound endpoint behind a corporate network often requires IT involvement (firewall / ingress), and a tunnel such as ngrok for local testing.
Use a WebSocket when you need the lowest latency, are powering a live UI or dashboard, expect high event volume, or cannot expose any inbound endpoint (the connection is client-initiated / outbound-only). In exchange you take on maintaining the connection (heartbeats, reconnects, token regeneration).
Appendix A — Event Filter Schema (EventFieldFilters)
Top-level structure
The filter is a JSON object whose keys are event fields (snake_case). All present conditions must match (logical AND across top-level keys).
Filterable fields
Field | Type of filter | Notes |
|---|---|---|
type | The event type, e.g. OrderCreated, ManualTask | |
state | The event state, e.g. READY | |
source_name | Publisher name, e.g. manual-task-service | |
source_id | Publisher instance id | |
stream_id | Groups related events | |
payload_type | Declared payload type, if any | |
source_category | cellario_os | cellario_edge | integration | other | |
source_component | e.g. data_access_api, driver, … (see below) | |
severity | debug | info | warning | error | critical | |
payload_matches | key/value pairs, all must match (AND) | |
payload_matches_any | per key: value matches any in list (OR per key, AND across keys) | |
payload_has_keys | keys that must all be present (AND) | |
payload_has_any_keys | keys where at least one must be present (OR) |
The complex-filter model also supports advanced fields not detailed here: time ranges (source_created_at, received_at), id ranges (id, triggering_event_id), and raw-JSONPath payload filters (payload_filter, payload_filters, payload_filters_any). See the API reference for their exact shapes.
String filters
Used by type, state, source_name, source_id, stream_id, payload_type. Several operators on one field are legal and are ANDed — {"type": {"in": ["Workflow"], "contains": "flow"}} applies both. All case-sensitive.
Operator | Shape | Meaning |
|---|---|---|
in | {"in": ["A","B"]} | field equals any listed value (exact) |
not_in | {"not_in": ["A"]} | field equals none of the listed values |
contains | {"contains": "Order"} | substring match |
starts_with | {"starts_with": "Order"} | prefix match |
ends_with | {"ends_with": "Created"} | suffix match |
is_null | {"is_null": true} | field is null (or false for not-null) |
{ "type": { "in": ["OrderCreated"] }, "state": { "in": ["pending", "completed"] } }Enum filters
Used by source_category, source_component, severity. Operators: in, not_in, is_null (no substring operators). Values must be the exact wire values below.
- source_category: cellario_os, cellario_edge, integration, other
- source_component: cellario_scheduler, cellario_scheduler_cloud, cellario_agent, driver_host, driver, data_access_api, events_api, lab_services_api, blob_storage_api, platform_api, extension_service, job_orchestrator
- severity: debug, info, warning, error, critical
{ "severity": { "in": ["error", "critical"] }, "source_component": { "in": ["driver"] } }Payload filters
Match against the event's JSON payload body. Keys support dot-notation for nested access (order.customer.tier). Keys and string values are case-sensitive. Values are matched type-aware against the JSON payload:
- string → quote it: "9b2f4c11-aaaa-4562-b3fc-2c963f66afa6" (compared with ==, case-sensitive)
- number → don't quote it: 56 (matches a JSON number; "56" would only match a JSON string "56")
- boolean → true / false; null → null
Field | Shape | Logic |
|---|---|---|
payload_matches | {"OrderAssetId": "9b2f4c11-aaaa-4562-b3fc-2c963f66afa6", "State": "FINISHED"} | every pair must match (AND) |
payload_matches_any | {"State": ["FINISHED","FAILED"]} | per key: any value in list (OR); across keys: AND |
payload_has_keys | ["WorkflowOrderId","State"] | all keys present (AND) |
payload_has_any_keys | ["a","b"] | at least one present (OR) |
{ "type": { "in": ["ManualTask"] }, "payload_matches": { "OrderAssetId": "9b2f4c11-aaaa-4562-b3fc-2c963f66afa6" } }Common event values (Cellario OS)
Use these exact strings. (Matching is case-sensitive — see the warning above.)
ManualTask events (emitted by the Guided Task / manual-task service)
Filter field | Exact value(s) |
|---|---|
type | ManualTask |
source_name | manual-task-service |
state | READY, STARTED, FINISHED, FAILED, CANCELED, SKIPPED |
ManualTask payload keys (PascalCase, dot-notation for nesting):
Key | Always present? | Notes |
|---|---|---|
State | ✅ | same value as the state field |
TaskId | ✅ | manual task instance id (GUID) |
OrderAssetId | ✅ | the order id — use this to match "my order" |
SampleOperationId | ✅ | runtime id string |
ParameterValues | ✅ | array of {ParameterId, Value} |
GeneratedDate | ✅ | event timestamp |
WorkflowOrderId | ⚠️ only for workflow-driven manual tasks | null for tasks created via direct API — prefer OrderAssetId |
WorkflowOrderStepId | ⚠️ only for workflow-driven manual tasks | as above |
Worked examples
1. Match manual tasks entering READY
{
"type": { "in": ["ManualTask"] },
"source_name": { "in": ["manual-task-service"] },
"state": { "in": ["READY"] }
}2. Match manual tasks reaching a terminal state (any of several)
{
"type": { "in": ["ManualTask"] },
"state": { "in": ["FINISHED", "FAILED", "CANCELED"] }
}3. Match a specific order's manual-task completion
Match the FINISHED event for one order by its OrderAssetId:
{
"type": { "in": ["ManualTask"] },
"state": { "in": ["FINISHED"] },
"payload_matches": { "OrderAssetId": "9b2f4c11-aaaa-4562-b3fc-2c963f66afa6" }
}4. Match any of several payload values
{
"type": { "in": ["ManualTask"] },
"payload_matches_any": { "State": ["FAILED", "CANCELED"] }
}5. Match error-or-worse events from a driver
{
"source_component": { "in": ["driver", "driver_host"] },
"severity": { "in": ["error", "critical"] }
}6. Match events whose payload contains a key
{
"type": { "in": ["ManualTask"] },
"payload_has_keys": ["WorkflowOrderId"]
}Full-syntax reference (NOT a real filter)
Every supported element shown once, for syntax reference only. Do not use this as a template — several of these conditions are redundant/contradictory together; a real filter uses a small subset.
{
"type": { "in": ["OrderCreated"] },
"state": { "not_in": ["pending"] },
"source_name": { "contains": "manual-task" },
"source_id": { "starts_with": "mts-" },
"stream_id": { "ends_with": "-v1" },
"payload_type": { "is_null": false },
"source_category": { "in": ["cellario_os"] },
"source_component": { "in": ["data_access_api"] },
"severity": { "not_in": ["debug"] },
"payload_matches": { "OrderAssetId": "9b2f4c11-aaaa-4562-b3fc-2c963f66afa6", "State": "FINISHED" },
"payload_matches_any": { "State": ["FINISHED", "FAILED"] },
"payload_has_keys": ["WorkflowOrderId", "State"],
"payload_has_any_keys": ["TaskId", "SampleOperationId"]
}