Carrier CPS limits: the approaches, animated

Plivo and Vobiz reject (they don't queue) SIP INVITEs that go over the account's calls-per-second limit. Plivo returns hangup 5180 cps_limit_reached and Vobiz returns 429 / SIP 503. Each tab below runs the same campaign burst through our pipeline under a different strategy: scheduler → orch-dispatch → worker → CreateSIPParticipant → carrier. Set up the campaign in Providers: add or remove carriers (up to 5) and set each one's calls and CPS. Each carrier enforces its own account limit, independently of the others. The default is 10 calls on carrier A (2 CPS) and 10 on carrier B (4 CPS). With a single provider, tab 4 shows the per-tenant pitfall instead of the global-key one.

Providers

sim time 0.00 s
untracked INVITE on same account ✓ accepted by carrier ✗ rejected → call FAILED ↻ rejected → retried
Reserved dial times (what the rate gate has promised, relative to now)
What the carrier sees: every INVITE on the account (ours + untracked) in the trailing 1 s, and each INVITE's outcome

Side by side, with the current settings

The same burst, the same random warm-up times, run through every approach. Change the sliders above and this table updates.

How the gate decides: "defer to the next free slot"

1 · Ask Redis for a slot

Right before CreateSIPParticipant, the worker runs one atomic script on key {voicebot:cps}:<carrier account>. Redis keeps one number, tat: the next moment an INVITE may leave.

2 · Get your wait, move tat forward

wait = max(tat, now) − now, then tat = max(tat, now) + 1/CPS × 1.15. Each caller takes a different slot, so they don't all wake together at the next second boundary. The 15 % margin absorbs LiveKit→carrier jitter.

3 · Sleep, then dial

await asyncio.sleep(wait), then dial. If wait would exceed a cap (e.g. 30 s), don't take a slot. Hand the call back as DEFERRED(CPS_LIMIT) without using up an attempt.

Redis script A · worker gate

-- KEYS[1] = {voicebot:cps}:<carrier_account>
-- ARGV[1] = interval_ms (1000/CPS × 1.15)
-- ARGV[2] = max_wait_ms
local t   = redis.call('TIME')  -- Redis clock, not pod clocks
local now = t[1]*1000 + math.floor(t[2]/1000)
local tat = tonumber(redis.call('GET', KEYS[1]) or now)
if tat < now then tat = now end
local wait = tat - now
if wait > tonumber(ARGV[2]) then return -1 end
local iv = tonumber(ARGV[1])
redis.call('SET', KEYS[1], tat + iv, 'PX', wait + iv + 1000)
return wait

Worker call site proposed

async def wait_for_dial_slot(redis, key, cps, max_wait_s=30):
    interval_ms = int(1000 / cps * 1.15)
    wait_ms = await redis.evalsha(GCRA_SHA, 1, key,
                                  interval_ms, int(max_wait_s * 1000))
    if wait_ms < 0:
        raise CpsBackpressure(key)   # → DEFERRED(CPS_LIMIT), no attempt used
    await asyncio.sleep(wait_ms / 1000)

# sip_dispatch.create_sip_participant, just before the RPC
await wait_for_dial_slot(
    redis, f"{{voicebot:cps}}:{telephony.rate_limit_key}",
    telephony.cps_limit)             # new fields on the telephony blob
info = await lkapi.sip.create_sip_participant(request)

Dispatch pacing B · orchestration

// DispatchCall: plan a dial slot per provider before CreateAgentDispatch
slot := pacer.Reserve(ctx, provider.RateLimitKey, provider.CPS) // same GCRA
if due := slot.Add(-warmLead); due.After(now) {
    // per-message delay (x-delay plugin or TTL+DLX), not the 30 s DEFERRED path
    return publisher.PublishCallQueuedAfter(ctx, call, due.Sub(now))
}
// otherwise dispatch now; the worker gate (A) still has the final say

The provider is resolved from from_number in GetDispatchConfig today. B needs that lookup done at dispatch time.

Treat CPS rejects as retryable C · worker + status updater

except SipDispatchError as e:
    if is_cps_reject(e) and attempt < 3:   # 429 / 503 / 5180, confirm from logs
        await asyncio.sleep(random.uniform(0.6, 1.6))
        continue                            # re-reserve a slot, re-dial
    raise

Before hard-coding codes, check which sip_status_code / reason each carrier actually sends for a CPS reject. The agent.sip_dispatch.failed log line already records it. A plain 503 can also be a real carrier failure.

Model assumptions (illustrative): dispatch 0.25–0.45 s per message, 10 dispatch consumers, worker warm-up (job start + GetCallExecutionConfig) 0.8–1.8 s, LiveKit→carrier latency 40–120 ms, 30 worker slots. The carrier meter is modelled as a sliding 1 s window over accepted INVITEs; neither carrier publishes its exact algorithm. Untracked traffic ≈ 0.6 INVITE/s. Random draws are seeded, so every approach sees the same warm-up times. With several providers, their calls are interleaved in the backlog, and untracked traffic hits the first provider's account. Total calls are capped at 30 (the modelled governor capacity).