Handler execution model & blocking contract¶
This documents how SIPhon runs Python script handlers, what may block, and the
elasticity / backpressure / liveness guarantees. It is the user-facing companion
to the source doc-comments in
src/script/py_executor.rs and
src/script/async_pool.rs.
The two handler pools¶
Every inbound SIP message that reaches a script handler runs on one of two pools of OS threads, each with a persistently-attached free-threaded Python interpreter (the persistent attach avoids per-handler mimalloc heap churn / a heap leak):
| Pool | Runs | Size config | Default |
|---|---|---|---|
Sync executor (PyExecutor) |
sync @proxy.on_request, @proxy.on_reply, @registrar.on_change, @rtpengine.on_dtmf, timers, … |
script.sync_pool_size / script.sync_pool_max |
core max(8, 2×CPUs), max max(32, 4×core) |
Async driver pool (AsyncPool) |
async def handlers + their asyncio.create_task work |
script.async_pool_size |
CPUs |
The sync pool is elastic¶
The sync pool starts at sync_pool_size (the core, always-on workers) and a
background grower adds workers on demand — up to sync_pool_max — whenever the
job queue has more work than the idle workers can take. It never shrinks:
workers are never reaped, which is exactly what keeps the persistent
free-threaded-CPython attach from leaking (reaping a persistently-attached
thread orphans ~2 MB of heap). Growth-on-demand restores the burst headroom that
blocking handlers need; never-reaping keeps the no-leak property.
Why elastic: an earlier change moved inbound dispatch off tokio's elastic
spawn_blockingpool (which grew threads on demand) onto a fixed pool to stop the heap leak — but that removed the burst valve. A blocking handler pins a worker for the whole call, so on a small box a couple of concurrent blocking REGISTERs exhausted the fixed pool and wedged the engine with no recovery. The elastic pool is the proper fix: it grows like the oldspawn_blockingpool but never reaps, so it neither wedges nor leaks. The regression is locked down bypool_grows_under_blocking_loadinpy_executor.rs.
The queue feeding the pool is bounded (script.executor_queue_capacity,
default 1024): once the pool is at its thread cap and the queue is full, new
jobs are shed.
Which handlers run, and which decision survives¶
@proxy.on_request takes an optional method filter. A filtered handler does
not replace the unfiltered one — both run, in registration order, and an
unfiltered handler matches every method.
That much is ordinary middleware. The part that surprises people is what happens to the routing decision.
One action slot, last writer wins¶
Every handler that runs for a request shares one request object with a single
action slot. reply(), relay(), fork() and reject() do not send
anything — they assign that slot. The dispatcher executes its final value
once, after every handler has returned.
So handlers do not each act. They take turns overwriting one decision:
@proxy.on_request("OPTIONS")
def probe(request):
request.reply(200, "OK") # ← discarded
@proxy.on_request
def route(request):
request.relay(NEXT_HOP) # ← this is what happens
The health probe is not answered and then relayed. It is only relayed —
the reply(200) is gone, silently, because a later handler assigned over it.
Registration order decides, which is rarely what the author meant.
Side effects are different: set_header(), record_route(), log, metrics and
cache writes all happen, from every handler that runs. It is only the routing
decision that is last-writer-wins.
Keeping your decision: stop_propagation()¶
@proxy.on_request("OPTIONS")
def probe(request):
request.reply(200, "OK")
request.stop_propagation() # nothing may overwrite this
The dispatcher stops the chain and executes what this handler chose.
It is opt-in on purpose. Answering is not on its own a request to stop — a metrics or logging handler running afterwards is legitimate, and stopping by default would silently drop it. Idempotent, and it leaves the chosen action alone.
The alternative, if you would rather not register two handlers at all, is to branch inside one:
@proxy.on_request
def route(request):
if request.method == "OPTIONS" and not request.in_dialog:
request.reply(200, "OK")
return
request.relay(NEXT_HOP)
Coming from Kamailio or OpenSIPS?
There is no equivalent of this in either. Both have exactly one
automatic entry point — Kamailio's request_route, OpenSIPS's route{} —
and everything else (route(NAME), branch_route, failure_route) is
invoked or armed explicitly. You branch inside the one route with
if (is_method("INVITE")), so two blocks can never both claim a request.
siphon's model is closer to HTTP middleware: several handlers may match, and
stop_propagation() is the next()-style control that model needs.
@proxy.on_request("INVITE") is therefore not a drop-in for
is_method("INVITE") — the Kamailio form is an exclusive branch, the
decorator is additive.
@diameter.on_request is the other way round
The Diameter decorator takes a filter of the same shape
(@diameter.on_request("ULR"), "ULR|AIR", "S6a:ULR") but dispatches
one handler per request: the most specific filter that matches wins, and
an unfiltered @diameter.on_request is the lowest-specificity fallback that
runs only when nothing more specific matched.
The difference is deliberate rather than accidental. A Diameter request needs exactly one answer, so one handler must own it. A SIP request can legitimately interest several handlers at once — metrics, lawful intercept, authentication, routing — so they compose. Worth knowing all the same if you are porting a routing pattern between the two namespaces.
What @b2bua.on_failure decides¶
A B2BUA call that could not be connected runs @b2bua.on_failure(call, code,
reason) once, before the caller hears anything. code is what the call failed
on: the best of its branches' failures (RFC 3261 §16.7), 408 for the ring
timeout, 503 when the B-leg INVITE never left or no LCR carrier was routable,
and 500 when @b2bua.on_answer raised or ended an answered call. The Call
has the same single action slot, and its final value is carried out:
| The handler leaves | What happens |
|---|---|
nothing, or call.terminate() |
the caller gets code |
call.reject(code, reason) |
the caller gets that response instead (3xx-6xx) |
call.dial(), call.fork(), call.route() |
the call is routed again, still unanswered, and the handler runs again if that fails too |
call.handover(app) |
the unanswered call goes to a control app |
call.answer() |
siphon answers the caller itself, and the call lives on |
Re-routes are capped at 10 per call, so a handler that keeps dialling a target
that keeps failing still ends the call. A decision that cannot apply (a
reject() with a 1xx or 2xx, a REFER decision, anything on a call siphon placed
itself) is logged at warn, and the call ends with its failure. So does a
handler that raises: a decision taken by a failing handler is not one to route a
call on.
The blocking contract — what script authors must know¶
A handler may call Rust APIs that block the worker thread on I/O:
auth.require_digest with the HTTP/Diameter backend,
proxy.send_request(wait_for_response=True), cache.fetch, diameter.*,
RTPEngine control, DNS/TLS connect during relay(), etc. While a handler
blocks, it occupies one pool worker.
The pool grows to absorb concurrent blocking handlers up to sync_pool_max, so
short blocking bursts are fine. But sustained blocking beyond the cap still
queues, and the maximum sustainable rate of a blocking handler is roughly:
Design accordingly:
- Cache hot lookups. For HTTP digest auth, set
auth.http.cache_ttl_secsso a registration storm for the same subscribers reuses a cached HA1 instead of making a blocking fetch per REGISTER — the pool then rarely needs to grow. - Fire-and-forget slow side-effects. Do contact-change notifications, CDR
posts, webhooks, etc. with
asyncio.create_task(...)from anasynchandler — don't block the SIP path on them. Ahttpx.Clientis not safe to share across threads. - Size for your backends and your memory. Raise
sync_pool_maxfor many slow blocking backends; lower it on memory-constrained NFs (peak memory ≈sync_pool_max × ~2 MB). - Never
time.sleep()in a sync handler. Ringing before answering is a real thing to want — alert the caller, then decide — but atime.sleep()in adefhandler pins a pool worker for the whole ring, on every inbound call, which is a worse problem than the one it solves. Use anasync defhandler andawait asyncio.sleep(...): awaiting mid-handler is the supported shape, and it is what the same handler already does betweenrtpengine.offerand the call.progress()andcall.answer()are imperative — the response goes out where you call them — so the wait sits between them:
@b2bua.on_invite
async def route(call):
call.progress(180, "Ringing") # on the wire here
await asyncio.sleep(2) # ring, without holding a worker
call.handover("ai-app", answer=True, profile="voice_ai")
handover, dial and reject are the deferred ones: they are applied when
the handler returns. If the wait depends on something only an out-of-process
application knows — a queue position, an agent becoming free, a model
finishing its load — hand the call over un-answered and let the controller
ring with the control plane's ring verb and then connect it with an anchored
answer (see the control-plane reference).
Blocking calls must release the interpreter (free-threaded GC safety)¶
On free-threaded CPython (3.14t) the cyclic GC performs a stop-the-world:
it pauses every thread that is attached to the interpreter at a safe point. A
thread that performs a blocking call (an HTTP/Diameter auth fetch, a DNS lookup,
…) while still attached can never reach that safe point, so for the duration
of that block every other handler that allocates cyclic garbage — which Python
does constantly — stalls behind the GC. This is verified: a thread blocking
while attached hangs a concurrent gc.collect() for the whole block; detached,
it returns immediately. Under blocking-heavy load (auth/Diameter storms) the
result is periodic engine-wide latency spikes, each lasting as long as the
blocking call — and it bites even at low concurrency (one blocked-while-attached
handler plus one GC trigger).
The worst case has been seen on cpus ≈ 1 nodes (a single-worker runtime,
where available_parallelism() reports 1), where the engine wedged outright
rather than just stalling. The exact escalation from a transient stall to a
permanent wedge on such nodes is not reproduced here, so treat it as an observed
correlation; the deadlock-aware watchdog (below) is the recovery backstop for
it. Either way the fix is the same.
The fix, and the rule for any blocking Rust-side work, is to release the
interpreter for the blocking window — in pyo3 terms,
Python::attach(|py| py.detach(|| block_in_place(block_on(future)))). siphon's
built-in blocking APIs (e.g. the HTTP digest-auth backend) do this internally;
the deadlock-aware watchdog above is the backstop for any path that doesn't.
This GC hazard is specific to Rust-side blocking calls, which hold the attach
for their whole duration. A blocking call made from Python in a handler — a
synchronous httpx/urllib/requests request in @registrar.on_change, say —
does not stall the GC: CPython releases the interpreter around the blocking
socket syscall, so the handler is detached for the wait. It still pins a pool
worker for the duration, though, so prefer asyncio.create_task(...) for slow
side-effects (throughput, not safety). Both paths are exercised by
run-tests.sh --http-auth: one scenario drives a REGISTER storm through the
blocking HA1-fetch (Rust) path, another through a blocking on_change notify
(Python) path, each on a single-worker (cpus: 0.5) siphon; the engine must
complete the registrations.
Backpressure & liveness guarantees¶
Beyond elasticity, the pool is defended on two more fronts so a misbehaving handler degrades gracefully instead of taking the node down silently:
- Bounded queue + load-shed. When the pool is at its cap and the queue is
full, new jobs are dropped (the SIP client retransmits) rather than growing
memory without bound. Counted by
siphon_pyexec_jobs_shed_total. - Liveness watchdog / fail-fast. A dedicated thread (immune to any lock a
wedged handler holds) aborts the process when there is work pending (a
handler in-flight or jobs queued) yet zero completions for
script.handler_stall_abort_secs(default 30 s;0disables). The trigger is independent of pool fill, so it catches a low-concurrency deadlock (a handful of handlers stuck on a lock/await while the pool sits far below its cap) as well as full saturation — an earlier "at the thread cap + fully busy" condition could never see the former, since the pool never grew to the cap. A healthy pool advances completions every tick; a genuinely idle pool has no pending work — neither trips. Aborting is deliberate: a hung-but-alive SIP engine never recovers on its own, so arestart: always/ systemd policy never fires — the abort turns an indefinite outage into a seconds-long restart and leaves a core for post-mortem.
Metrics (/metrics)¶
| Metric | Meaning |
|---|---|
siphon_pyexec_pool_size |
live worker threads (grows core→max under load) |
siphon_pyexec_pool_max |
configured thread ceiling |
siphon_pyexec_inflight |
handlers currently executing |
siphon_pyexec_queue_depth |
handler jobs waiting in the queue |
siphon_pyexec_jobs_completed_total |
completed handler jobs |
siphon_pyexec_jobs_shed_total |
jobs dropped because the queue was full |
siphon_auth_ha1_cache_hits_total |
HTTP-auth lookups served from cache |
Alert on: a sustained rate(siphon_pyexec_jobs_shed_total) > 0, or
siphon_pyexec_pool_size == siphon_pyexec_pool_max with
siphon_pyexec_inflight == siphon_pyexec_pool_size held for minutes — both mean
the pool is fully grown and saturated, approaching the watchdog's abort condition.