Event Service — Python Example
A runnable Python example that subscribes to the Cellario OS Event Service and receives events over both delivery methods — Webhook and WebSocket. In a single run it configures a webhook subscriber and a websocket subscriber, posts a test event, waits for it to arrive on both channels, prints a summary, and deletes the subscribers it created. Read it alongside the Subscribing to the Event Service guide.
Files
File | What it is |
|---|---|
pyproject.toml | Project manifest and dependencies (installed by uv sync) |
src/config.py | Loads configuration from environment variables and derives the service URLs |
src/api_client.py | A thin synchronous wrapper around the Event Service REST API (subscribers and filters) |
src/webhook_receiver.py | A small Flask app that receives events pushed to your webhook URL |
src/websocket_client.py | An async WebSocket consumer that receives events over a persistent connection |
src/main.py | Orchestrates the run: configure subscribers, post a test event, confirm delivery on both channels, and clean up |
Prerequisites
- macOS or Linux
- uv for installing dependencies and running the example
- An HTTPS tunnel such as ngrok — needed so the Event Service can reach a webhook receiver running on your own machine (see Why a tunnel? below)
- A JWT for your Event Service with the EVENTS_WRITE role
Setup
- Install dependencies:
uv sync- Start a tunnel to your local port in a separate terminal, and copy the public https:// URL it prints:
ngrok http 8080- Export the required environment variables in the shell you will run from. No .env file is used — secrets stay in your shell:
export EVENTS_BASE_URL="https://your-instance.cellario.cloud/api/events"
export EVENTS_JWT="<your JWT with EVENTS_WRITE>"
export WEBHOOK_PUBLIC_URL="https://<your-subdomain>.ngrok-free.app"
export WEBHOOK_API_KEY="any-opaque-string"
# Optional, defaults to 8080:
# export WEBHOOK_LOCAL_PORT="8080"Run
From the project directory, in the same shell where you exported the variables:
uv run python -m src.mainA successful run configures both subscribers, posts a test event, reports it arriving on both the webhook and websocket channels, and ends with [result] PASS before cleaning up.
Why a tunnel?
Webhook delivery is the Event Service POSTing to a URL you register. A receiver running on your laptop sits behind NAT and is not reachable from a cluster-hosted Event Service (for example your-instance.cellario.cloud), so a tunnel such as ngrok gives the service a public URL that forwards to your local port. The WebSocket path needs no tunnel — the client opens that connection outbound to the Event Service.
Source
pyproject.toml
[project]
name = "events-python"
version = "0.1.0"
description = "Cellario OS Event Service Python example (websocket and webhook)"
readme = "README.md"
authors = [
{ name = "HighRes Biosolutions" }
]
requires-python = ">=3.11"
dependencies = [
"flask>=3.1.3",
"requests>=2.34.2",
"websockets>=16.0",
]
[dependency-groups]
dev = [
"pytest>=9.1.0",
]src/main.py
from __future__ import annotations
import asyncio
import queue
import sys
import time
import uuid
from src import webhook_receiver
from src.api_client import ApiClientError, EventsApiClient
from src.config import load_config
from src.websocket_client import consume as ws_consume
WAIT_TIMEOUT_SECONDS = 30.0
WS_SETTLE_SECONDS = 2.0
WS_SHUTDOWN_TIMEOUT_SECONDS = 5.0
def _try_get(q: "queue.Queue[dict]", timeout: float) -> dict | None:
try:
return q.get(timeout=timeout)
except queue.Empty:
return None
async def _wait_for_first(q: "queue.Queue[dict]", deadline: float) -> dict | None:
remaining = deadline - time.monotonic()
if remaining <= 0:
return None
return await asyncio.to_thread(_try_get, q, remaining)
async def main() -> int:
cfg = load_config()
run_id = uuid.uuid4().hex[:8]
event_type = f"SmokeTest_{run_id}"
state = "smoke"
routing_filter = f"{event_type}.#"
print(f"[main] run_id={run_id}")
print(f"[main] events_base_url={cfg.events_base_url}")
print(f"[main] events_ws_base_url={cfg.events_ws_base_url}")
print(f"[main] webhook_public_url={cfg.webhook_public_url}")
print(f"[main] webhook_local_port={cfg.webhook_local_port}")
api = EventsApiClient(cfg)
webhook_subscriber_id = str(uuid.uuid4())
ws_subscriber_id = str(uuid.uuid4())
cf_subscriber_id = str(uuid.uuid4())
ws_received: "queue.Queue[dict]" = queue.Queue()
cf_received: "queue.Queue[dict]" = queue.Queue()
stop_event = asyncio.Event()
ws_task: asyncio.Task | None = None
cf_task: asyncio.Task | None = None
webhook_configured = False
ws_configured = False
cf_configured = False
try:
webhook_receiver.start_in_background(cfg.webhook_api_key, cfg.webhook_local_port)
endpoint_url = f"{cfg.webhook_public_url}/events"
print(f"[main] configuring webhook subscriber {webhook_subscriber_id} -> {endpoint_url}")
api.configure_webhook(
subscriber_id=webhook_subscriber_id,
name=f"PythonSample Webhook {run_id}",
endpoint_url=endpoint_url,
api_key=cfg.webhook_api_key,
)
webhook_configured = True
api.create_subscription(
subscriber_id=webhook_subscriber_id,
name=f"PythonSample Webhook Filter {run_id}",
simple_filter=routing_filter,
)
print(f"[main] configuring websocket subscriber {ws_subscriber_id}")
ws_response = api.configure_websocket(
subscriber_id=ws_subscriber_id,
name=f"PythonSample WebSocket {run_id}",
)
ws_configured = True
connection_token = api.extract_connection_token(ws_response)
api.create_subscription(
subscriber_id=ws_subscriber_id,
name=f"PythonSample WebSocket Filter {run_id}",
simple_filter=routing_filter,
)
ws_url = f"{cfg.events_ws_base_url}/v1/subscribers/{ws_subscriber_id}/connect"
ws_task = asyncio.create_task(
ws_consume(ws_url, connection_token, ws_received, stop_event)
)
print(f"[main] configuring websocket+complex-filter subscriber {cf_subscriber_id}")
cf_configure_response = api.configure_websocket(
subscriber_id=cf_subscriber_id,
name=f"PythonSample WS+CF {run_id}",
)
cf_configured = True
cf_connection_token = api.extract_connection_token(cf_configure_response)
cf_create_response = api.create_complex_filter(
subscriber_id=cf_subscriber_id,
name=f"PythonSample Complex Filter {run_id}",
filter_={
"type": {"in": [event_type]},
"payload_has_keys": ["run_id"],
},
)
cf_filter_id = api.extract_filter_id(cf_create_response)
api.create_subscription(
subscriber_id=cf_subscriber_id,
name=f"PythonSample WS+CF Subscription {run_id}",
complex_filter_id=cf_filter_id,
)
cf_ws_url = f"{cfg.events_ws_base_url}/v1/subscribers/{cf_subscriber_id}/connect"
cf_task = asyncio.create_task(
ws_consume(cf_ws_url, cf_connection_token, cf_received, stop_event)
)
await asyncio.sleep(WS_SETTLE_SECONDS)
print(f"[main] posting test event type={event_type}")
post_result = api.post_event(
event_type=event_type,
state=state,
payload={"run_id": run_id, "hello": "world"},
run_id=run_id,
)
posted_id = post_result.get("id")
print(f"[main] posted event id={posted_id}, waiting up to {WAIT_TIMEOUT_SECONDS}s")
deadline = time.monotonic() + WAIT_TIMEOUT_SECONDS
async def wait_webhook() -> dict | None:
return await _wait_for_first(webhook_receiver.received_events, deadline)
async def wait_ws() -> dict | None:
return await _wait_for_first(ws_received, deadline)
async def wait_cf() -> dict | None:
return await _wait_for_first(cf_received, deadline)
webhook_event, ws_event, cf_event = await asyncio.gather(
wait_webhook(), wait_ws(), wait_cf()
)
if webhook_event is not None:
print(f"[result] Webhook received id={webhook_event.get('id')} type={webhook_event.get('type')}")
else:
print("[result] Webhook TIMEOUT (no event in 30s)")
if ws_event is not None:
print(f"[result] WebSocket received id={ws_event.get('id')} type={ws_event.get('type')}")
else:
print("[result] WebSocket TIMEOUT (no event in 30s)")
if cf_event is not None:
print(f"[result] WebSocket+CF received id={cf_event.get('id')} type={cf_event.get('type')}")
else:
print("[result] WebSocket+CF TIMEOUT (no event in 30s)")
if webhook_event is not None and ws_event is not None and cf_event is not None:
print("[result] PASS")
return 0
print("[result] FAIL")
return 1
except ApiClientError as e:
print(f"[error] API call failed: {e}")
return 1
finally:
stop_event.set()
for task, label in ((ws_task, "websocket"), (cf_task, "websocket+cf")):
if task is None:
continue
try:
await asyncio.wait_for(task, timeout=WS_SHUTDOWN_TIMEOUT_SECONDS)
except (asyncio.TimeoutError, asyncio.CancelledError):
task.cancel()
except Exception as e:
print(f"[warn] {label} task raised on shutdown: {e!r}")
for sid, configured in (
(webhook_subscriber_id, webhook_configured),
(ws_subscriber_id, ws_configured),
(cf_subscriber_id, cf_configured),
):
if not configured:
continue
try:
api.delete_subscriber(sid)
print(f"[cleanup] deleted subscriber {sid}")
except ApiClientError as e:
print(f"[cleanup] failed to delete subscriber {sid}: {e}")
if __name__ == "__main__":
try:
sys.exit(asyncio.run(main()))
except KeyboardInterrupt:
print("[main] interrupted")
sys.exit(1)src/config.py
from __future__ import annotations
import os
from dataclasses import dataclass
@dataclass(frozen=True)
class Config:
events_base_url: str
events_ws_base_url: str
jwt: str
webhook_public_url: str
webhook_api_key: str
webhook_local_port: int
def auth_headers(self) -> dict[str, str]:
return {"Authorization": f"Bearer {self.jwt}"}
def _derive_ws_url(http_url: str) -> str:
stripped = http_url.rstrip("/")
if stripped.startswith("https://"):
return "wss://" + stripped[len("https://"):]
if stripped.startswith("http://"):
return "ws://" + stripped[len("http://"):]
raise ValueError(f"URL must start with http:// or https://, got: {http_url}")
def load_config() -> Config:
required = {
"EVENTS_BASE_URL": os.environ.get("EVENTS_BASE_URL"),
"EVENTS_JWT": os.environ.get("EVENTS_JWT"),
"WEBHOOK_PUBLIC_URL": os.environ.get("WEBHOOK_PUBLIC_URL"),
"WEBHOOK_API_KEY": os.environ.get("WEBHOOK_API_KEY"),
}
missing = [name for name, value in required.items() if not value]
if missing:
raise RuntimeError(
"Missing required environment variables: " + ", ".join(missing)
+ ". Export them in your shell before running (see README)."
)
events_base_url = required["EVENTS_BASE_URL"].rstrip("/")
port = int(os.environ.get("WEBHOOK_LOCAL_PORT", "8080"))
return Config(
events_base_url=events_base_url,
events_ws_base_url=_derive_ws_url(events_base_url),
jwt=required["EVENTS_JWT"],
webhook_public_url=required["WEBHOOK_PUBLIC_URL"].rstrip("/"),
webhook_api_key=required["WEBHOOK_API_KEY"],
webhook_local_port=port,
)src/api_client.py
from __future__ import annotations
import json
from datetime import datetime, timezone
import requests
from src.config import Config
REQUEST_TIMEOUT_SECONDS = 30
class ApiClientError(RuntimeError):
pass
class EventsApiClient:
def __init__(self, config: Config) -> None:
self._config = config
self._headers = {
**config.auth_headers(),
"Content-Type": "application/json",
}
def configure_webhook(
self,
subscriber_id: str,
name: str,
endpoint_url: str,
api_key: str,
) -> dict:
url = f"{self._config.events_base_url}/v1/subscribers/{subscriber_id}/configure"
payload = {
"name": name,
"delivery_configuration": {
"delivery_type": "webhook",
"web_hook": {
"endpoint_url": endpoint_url,
"api_key": api_key,
},
},
}
return self._post_json(url, payload)
def configure_websocket(self, subscriber_id: str, name: str) -> dict:
url = f"{self._config.events_base_url}/v1/subscribers/{subscriber_id}/configure"
payload = {
"name": name,
"delivery_configuration": {
"delivery_type": "websocket",
"web_socket": {
"acknowledgment_mode": "auto",
},
},
}
return self._post_json(url, payload)
def create_subscription(
self,
subscriber_id: str,
name: str,
*,
simple_filter: str | None = None,
complex_filter_id: str | None = None,
) -> dict:
if (simple_filter is None) == (complex_filter_id is None):
raise ValueError(
"create_subscription requires exactly one of simple_filter or complex_filter_id"
)
url = f"{self._config.events_base_url}/v1/subscribers/{subscriber_id}/subscriptions"
payload: dict = {"name": name}
if simple_filter is not None:
payload["simple_filter"] = simple_filter
else:
payload["complex_filter_id"] = complex_filter_id
return self._post_json(url, payload)
def create_complex_filter(
self,
subscriber_id: str,
name: str,
filter_: dict,
) -> dict:
url = f"{self._config.events_base_url}/v1/subscribers/{subscriber_id}/filters"
payload = {
"name": name,
"filter": filter_,
}
return self._post_json(url, payload)
def extract_filter_id(self, create_filter_response: dict) -> str:
try:
filter_id = create_filter_response["filter_id"]
except (KeyError, TypeError):
raise ApiClientError(
f"Create-filter response missing filter_id: {create_filter_response}"
)
if not filter_id:
raise ApiClientError("filter_id was empty in create-filter response")
return filter_id
def post_event(
self,
event_type: str,
state: str,
payload: dict,
run_id: str,
) -> dict:
url = f"{self._config.events_base_url}/v1/events"
body = {
"type": event_type,
"state": state,
"source_category": "integration",
"source_component": "other",
"source_name": f"PythonSample_{run_id}",
"source_id": f"sample-{run_id}",
"source_created_at": datetime.now(timezone.utc).isoformat(),
"payload_type": "application/json",
"payload": json.dumps(payload),
"source_component_extended_description": "Cellario OS events sample",
}
return self._post_json(url, body)
def delete_subscriber(self, subscriber_id: str) -> None:
url = f"{self._config.events_base_url}/v1/subscribers/{subscriber_id}"
resp = requests.delete(url, headers=self._headers, timeout=REQUEST_TIMEOUT_SECONDS)
if resp.status_code == 404:
return
if not resp.ok:
raise ApiClientError(f"DELETE {url} returned {resp.status_code}: {resp.text}")
def extract_connection_token(self, configure_response: dict) -> str:
try:
token = configure_response["delivery_configuration"]["web_socket"]["connection_token"]
except (KeyError, TypeError):
raise ApiClientError(
f"Configure response missing connection_token: {configure_response}"
)
if not token:
raise ApiClientError("connection_token was empty in configure response")
return token
def _post_json(self, url: str, payload: dict) -> dict:
resp = requests.post(
url,
headers=self._headers,
json=payload,
timeout=REQUEST_TIMEOUT_SECONDS,
)
if not resp.ok:
raise ApiClientError(f"POST {url} returned {resp.status_code}: {resp.text}")
if resp.status_code == 204 or not resp.text:
return {}
return resp.json()src/webhook_receiver.py
from __future__ import annotations
import queue
import threading
from flask import Flask, jsonify, request
received_events: "queue.Queue[dict]" = queue.Queue()
_processed_ids: set[str] = set()
_processed_lock = threading.Lock()
def make_app(expected_api_key: str) -> Flask:
app = Flask(__name__)
@app.post("/events")
def receive():
if request.headers.get("X-API-KEY") != expected_api_key:
return jsonify(error="unauthorized"), 401
event_id = request.headers.get("X-Event-Id")
if event_id:
with _processed_lock:
if event_id in _processed_ids:
return jsonify(status="already_processed"), 200
_processed_ids.add(event_id)
body = request.get_json(silent=True) or {}
received_events.put(body)
return jsonify(status="processed"), 200
@app.get("/health")
def health():
return jsonify(status="ok"), 200
return app
def start_in_background(expected_api_key: str, port: int) -> threading.Thread:
app = make_app(expected_api_key)
def _run():
# Flask's built-in dev server is fine for a smoke test. Bind to all
# interfaces so ngrok can reach the process.
app.run(host="0.0.0.0", port=port, debug=False, use_reloader=False)
thread = threading.Thread(target=_run, name="webhook-receiver", daemon=True)
thread.start()
return threadsrc/websocket_client.py
from __future__ import annotations
import asyncio
import json
import queue
import websockets
async def consume(
ws_url: str,
connection_token: str,
received_events: "queue.Queue[dict]",
stop_event: asyncio.Event,
) -> None:
headers = {"Authorization": f"Bearer {connection_token}"}
try:
async with websockets.connect(ws_url, additional_headers=headers) as ws:
print(f"[ws] connected to {ws_url}")
stop_task = asyncio.create_task(stop_event.wait())
try:
while not stop_event.is_set():
recv_task = asyncio.create_task(ws.recv())
done, _ = await asyncio.wait(
{recv_task, stop_task},
return_when=asyncio.FIRST_COMPLETED,
)
if stop_task in done:
recv_task.cancel()
break
raw = recv_task.result()
msg = json.loads(raw)
msg_type = msg.get("type")
if msg_type == "connection.ready":
print(f"[ws] ready: subscriber_id={msg.get('subscriber_id')}")
elif msg_type == "events.batch":
for evt in msg.get("events", []):
received_events.put(evt)
elif msg_type == "ping":
await ws.send(json.dumps({
"type": "pong",
"timestamp": msg.get("timestamp"),
}))
elif msg_type == "error":
print(f"[ws] server error: {msg.get('code')} {msg.get('message')}")
break
else:
print(f"[ws] ignoring message type: {msg_type}")
finally:
stop_task.cancel()
except websockets.ConnectionClosed as e:
print(f"[ws] connection closed: {e.code} {e.reason}")
except Exception as e:
print(f"[ws] error: {e!r}")
raise
finally:
print("[ws] stopped")