Forward verified Sume webhooks to Kafka, keyed by job_id

Verify the sume-v1 signature, refuse an empty secret, then produce each job event to Kafka with job_id as the key so retries land in one partition.

6 min readSume
All posts

Verify the delivery first, then call producer.produce(topic, key=job_id, value=raw_body) and return 204 only after the broker has acknowledged it. Keying by job_id sends every retry of the same job to one partition, so a consumer can dedupe in order.

Why key by job_id

Sume retries a failed delivery up to 10 times at 30 second spacing, so the same event can arrive more than once. The webhook body carries job_id, which the docs name as the dedupe handle. A Kafka key does not dedupe by itself; it only keeps duplicates adjacent for your consumer.

Webhook headers and what to do with them (read 2026-10-04)
HeaderMeaningAction
x-sume-webhook-timestampUnix secondsReject outside 300 s
x-sume-webhook-signaturesume-v1=<hex>, comma separated during rotationAccept if any entry matches
x-sume-webhook-secret-fingerprintFingerprint of the signing secretLog for support

Receiver

The Confluent client documents Producer(conf) with bootstrap.servers, produce(topic, key=..., value=..., callback=...), poll() to serve delivery callbacks, and flush() before shutdown. The handler below waits for the callback before it answers Sume. It uses only the standard library HTTP server.

import asyncio, hashlib, hmac, os, threading, time
from http.server import BaseHTTPRequestHandler, HTTPServer
from confluent_kafka import Producer
SECRET = os.environ.get("SUME_COM_WEBHOOK_SIGNING_SECRET", "")
if not SECRET:
    raise SystemExit("SUME_COM_WEBHOOK_SIGNING_SECRET is empty; refusing to start")
producer = Producer({"bootstrap.servers": os.environ["KAFKA_BOOTSTRAP"]})
def verify(raw, ts, header):
    if abs(time.time() - int(ts)) > 300: return False
    want = hmac.new(SECRET.encode(), ts.encode() + b"." + raw, hashlib.sha256).hexdigest()
    return any(hmac.compare_digest(p.strip().removeprefix("sume-v1="), want) for p in header.split(","))
class H(BaseHTTPRequestHandler):
    def do_POST(self):
        raw = self.rfile.read(int(self.headers.get("content-length", 0)))
        try:
            ok = verify(raw, self.headers.get("x-sume-webhook-timestamp", ""), self.headers.get("x-sume-webhook-signature", ""))
        except ValueError:
            ok = False
        if not ok:
            return self.send_response(401) or self.end_headers()
        import json
        job_id = json.loads(raw)["job_id"]
        err = []
        producer.produce("sume.job.events", key=job_id, value=raw, callback=lambda e, m: err.append(e))
        producer.flush(10)
        self.send_response(204 if err and err[0] is None else 503); self.end_headers()
async def main():
    srv = HTTPServer(("0.0.0.0", 8080), H)
    await asyncio.to_thread(srv.serve_forever)
asyncio.run(main())

Status codes

A 503 makes Sume retry, which is what you want when Kafka is down. A 204 after flush means the event is on the broker. Never answer 204 before the callback fires, or a crash loses the event with no retry.

Keep a polling fallback

Keep polling GET /v1/jobs/{id}/status as a backup. The docs call a webhook a delivery optimization, not your only recovery path. See also how to route on the event name.

Sources

Related posts

More in Developers

All Developers posts

Written by Sume