Webhook consumers are the part of a communications platform that everyone ships and few people enjoy. The endpoint starts as a five-line Flask route. Then retries arrive. Then a burst of call and SMS events lands in the same second. Then an engineer at 3am needs to explain what happened, and the logs are not enough.
The webhook-aggregator-fanout sample is a deliberate answer to that. It is one Flask file that takes raw Telnyx webhooks and runs them through a six-stage pipeline: receive, verify, dedup, log, fan out, process. No Redis, no Celery, no message broker. The storage primitives are SQLite for the audit log and a Python dict for the dedup KV, so the whole thing runs locally without infrastructure.
The complete sample is at https://github.com/team-telnyx/telnyx-code-examples/tree/main/webhook-aggregator-fanout.
The pipeline
Telnyx webhook POST
|
v
[1] Receive — Flask route reads raw body + headers
|
v
[2] Verify — Ed25519 signature check via Telnyx SDK v4
|
v
[3] Dedup — TTL-based in-memory KV store (default 300s)
|
v
[4] Log — SQLite insert (INSERT OR IGNORE for race safety)
|
v
[5] Fan out — Route to call queue or SMS queue by event type
|
v
[6] Process — Drain queue: answer call + play greeting, or reply to SMS
Each stage has a single responsibility. The route reads the raw body and headers, then hands off. Verification is the SDK's job, not the route's. Deduplication happens before the log, so duplicates never touch SQLite. The log happens before the fanout, so the audit record exists before any side effect. The fanout routes by event type; the drain executes.
The separation matters because it makes each stage independently testable and independently replaceable. The in-memory KV store can become Redis without touching the route. The SQLite log can become Postgres without touching the dedup. The in-process queue can become Celery without touching the router.
Verify before you trust
The handler reads the raw body and the request headers and passes them to the SDK:
import telnyx
from telnyx.lib.webhooks_ed25519 import unwrap_with_ed25519
telnyx_client = telnyx.Telnyx(
api_key=os.getenv("TELNYX_API_KEY"),
public_key=os.getenv("TELNYX_PUBLIC_KEY"),
)
@app.route("/webhooks", methods=["POST"])
def webhook_handler():
raw_body = request.get_data(as_text=True)
try:
verified_event = unwrap_with_ed25519(
telnyx_client, raw_body, request.headers
)
except Exception:
return jsonify({"error": "Invalid signature"}), 401
unwrap_with_ed25519 reads the Telnyx-Ed25519-Signature and Telnyx-Ed25519-Timestamp headers, verifies the signature against the configured public key, and returns a typed UnwrapWebhookEvent. If the signature is invalid or the timestamp is outside the replay window, it raises — and the handler returns 401 before any business logic runs.
This is the first stage for a reason. A webhook that fails verification should never reach the dedup store, the database, or the action queues. The public key is not a secret; what it proves is that the payload was signed by someone holding the Telnyx private key.
Deduplicate with a TTL KV store
Telnyx retries webhook delivery if the endpoint does not respond within the timeout. The same event can arrive two, three, or five times within a few minutes. Without deduplication, every retry triggers business logic — a call gets answered twice, a customer receives three identical SMS replies.
The sample uses an in-memory dict with a TTL:
dedup_store = {} # {event_id: (timestamp, timestamp)}
DEDUP_TTL_SECONDS = int(os.getenv("DEDUP_TTL_SECONDS", "300"))
def is_duplicate(event_id):
current_time = time.time()
# Sweep expired entries on every check
expired = [k for k, (ts, _) in dedup_store.items()
if current_time - ts > DEDUP_TTL_SECONDS]
for key in expired:
del dedup_store[key]
if event_id in dedup_store:
return True
dedup_store[event_id] = (current_time, current_time)
return False
The sweep runs on every check, so the store never grows unbounded. The default TTL is 300 seconds, which covers Telnyx's retry window. The event ID comes from the Telnyx payload when present and falls back to a SHA-256 of the sorted payload when it does not.
A duplicate returns {"status": "duplicate"} with a 200. Returning a non-2xx would make Telnyx retry again — the opposite of what you want.
Log before you act
Every non-duplicate event is written to SQLite before any action fires:
def log_event(event_id, event_type, payload):
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
cursor.execute(
"INSERT OR IGNORE INTO webhook_events "
"(event_id, event_type, payload, received_at, processed_at) "
"VALUES (?, ?, ?, ?, ?)",
(event_id, event_type, json.dumps(payload),
datetime.now(timezone.utc).isoformat(),
datetime.now(timezone.utc).isoformat())
)
conn.commit()
conn.close()
The UNIQUE constraint on event_id with INSERT OR IGNORE makes the log idempotent at the storage layer. Even if two identical events slip past the in-memory dedup in a race, SQLite enforces uniqueness.
The log is operational data, not a log file. It is a table with event_id, event_type, payload, received_at, processed_at. You can query it, export it, or build a dashboard on top of it. When something goes wrong at 3am, you have a real table, not a grep target.
Fan out by event type
After verification, dedup, and logging, the event is routed to one of two in-memory queues based on its type:
ACTION_TYPES = ["call", "sms"]
action_queues = defaultdict(list)
# In webhook_handler:
if "call" in event_type.lower():
enqueue_action("call", event)
elif "message" in event_type.lower() or "sms" in event_type.lower():
enqueue_action("sms", event)
The fanout is a router, not a dispatcher. It decides which queue gets the event; it does not execute the action. The queue is drained in the same request for this sample, but in production the drain moves to a separate worker process.
The separation matters because call actions and SMS actions have different latency profiles and different failure modes. A call that takes 800ms to answer should not block an SMS reply that could fire in 50ms. Separate queues let you prioritize, retry, and monitor each action type independently.
Process with the Telnyx SDK v4
The queue drain uses the Telnyx Python SDK v4 client to take the actual action:
def process_call_action(event_data):
payload = event_data.get("data", {}).get("payload", {})
call_control_id = payload.get("call_control_id")
telnyx_client.calls.answer(call_control_id=call_control_id)
telnyx_client.calls.playback_start(
call_control_id=call_control_id,
audio_url="https://example.com/greeting.mp3"
)
def process_sms_action(event_data):
payload = event_data.get("data", {}).get("payload", {})
from_number = payload.get("from")
to_number = payload.get("to")
text = payload.get("text", "")
telnyx_client.messages.create(
from_=to_number,
to=from_number,
text=f"Thanks for your message! We received: {text[:50]}..."
)
Note the field access: event_data.get("data", {}).get("payload", {}). The Telnyx webhook envelope wraps the payload inside data.payload, and the SDK v4 typed event follows that structure. For teams migrating from SDK v2, this is the most common breakage point — v2 gave you the flat payload, v4 gives you the wrapped envelope.
Run the demo without credentials
The sample ships with a single-file demo launcher that runs the entire pipeline without Telnyx credentials or ngrok. It generates an Ed25519 keypair at startup, injects the public key into the Telnyx client, stubs the calls and messages API surfaces, and serves a dashboard at http://localhost:5555/.
git clone https://github.com/team-telnyx/telnyx-code-examples.git
cd telnyx-code-examples/webhook-aggregator-fanout
python -m venv .venv && source .venv/bin/activate
pip install -r requirements.txt
python demo/demo_server.py
Open the dashboard, click "Send Call Webhook" or "Send SMS Webhook," and watch the six-stage pipeline execute live. The dashboard shows the dedup KV, the SQLite event log, and the fanout queues in real time.
To test deduplication, send the same event twice. The first returns {"status": "success"}. The second returns {"status": "duplicate"}. Only one row appears in the events table.
Run it for real
cp .env.example .env
# Set TELNYX_API_KEY and TELNYX_PUBLIC_KEY in .env
flask --app app run --port 5000
Point your Telnyx webhook URL to https://your-public-url/webhooks. For local testing, use ngrok to expose the port.
Production hardening
Three changes before this goes to production:
- Move the queue drain to a worker process. The sample drains queues in the same request so the pipeline is visible end to end. In production, enqueue to Redis or a real queue and drain from a separate process. The webhook handler should return 200 as soon as the event is logged and routed.
- Add a dead-letter queue. If an action fails, it should not disappear. Log it, retry with backoff, and surface it for manual intervention. The SQLite log already gives you the event record; the dead-letter queue extends it with retry state.
- Monitor queue depths. If the call queue is growing faster than it drains, you have a capacity problem before you have a customer problem. Export queue depth to your monitoring system.
The core pattern is simple: verify before you trust, deduplicate before you act, log before you process, and fan out before you execute. Everything else is infrastructure.