How to Add Backpressure and Load Shedding to a Python Service Before Overload Takes It Down
A step-by-step Python tutorial that measures an unprotected service collapsing at twice its capacity, then adds a bounded line, a wait limit, a 503 with Retry-After, and polite retries, with every...
Every service has a speed limit. Below it, requests are answered quickly. Above it, something has to give, and without protection it is often everything at once: the line of waiting requests grows, every answer gets slower, clients time out and hang up, and the service spends its effort finishing work that nobody is waiting for any more. The cruel part is that the service looks busy and healthy while most requests fail.
Table Of Content
- The ideas you need first
- Capacity, offered load, and the line
- Goodput, throughput, and wasted work
- Backpressure, load shedding, and rate limiting
- HTTP 503 and Retry-After
- Prerequisites
- Set up the lab folder
- Create the three building blocks
- The admission controller
- The load generator
- The shared settings
- Step 1: Measure your service’s real capacity
- Step 2: See backpressure in a plain asyncio queue
- Step 3: Watch an unprotected service collapse
- Step 4: Cap the line to bound the delay
- Step 5: Cap the wait too, because a count is not a time
- Step 6: Serve it over HTTP with a 503 and Retry-After
- The service
- The real HTTP load test
- What uvicorn’s built-in limit does
- Step 7: Make clients part of the solution
- Step 8: Test it, including the bug that leaks the line counter
- Reproduce the bug
- Write the tests
- Step 9: Run the load test that proves it
- Tune the wait limit with data
- Confirm everything works end to end
- Common mistakes and how to spot them
- Letting a line grow without a limit
- Sizing limits from a healthy day
- Counters that can leak
- Putting health checks behind the limiter
- Judging by the latency of the successes
- Retrying immediately
- Making rejection expensive
- Where to go next
- Sources
In this tutorial you will watch that happen on your own machine, then fix it one step at a time with a small amount of Python. You will measure an unprotected service falling off a cliff at twice its capacity, then add a bounded line, a limit on how long anyone may wait, a proper HTTP 503 response with a Retry-After header, and clients that retry politely. After each change you will measure again, so every claim in this article comes with a number you can reproduce.
The two ideas are old and simple. Backpressure means a component that is full tells whoever feeds it to slow down. Load shedding means turning some requests away on purpose so the rest can succeed. By the end you will have a short admission controller, a FastAPI service that uses it, a load generator, eight tests, and a load test you can point at your own code.
The ideas you need first
Capacity, offered load, and the line
Capacity is how many requests per second your service can finish when it is kept completely busy. Offered load is how many requests per second actually arrive. When offered load is below capacity, requests are served almost as soon as they arrive. When it is above capacity, the extra requests wait in a line, and the line gets longer every second the overload lasts.
How long is the wait? Little’s law says that in a stable system the average number of items inside equals the average arrival rate times the average time each item spends there, usually written L = λW. We will use it in one direction: a line of L requests in front of a service that finishes C requests per second makes the last person in line wait about L ÷ C seconds. A line of 33 in front of a service that finishes 65 per second means a wait of about half a second. That is the whole trick behind choosing a sensible limit in Step 4.
Goodput, throughput, and wasted work
A request that is answered after the client gave up is worth nothing. The AWS article Using load shedding to avoid overload draws the line like this: “Throughput is the total number of requests per second that is being sent to the server. Goodput is the subset of the throughput that is handled without errors and with low enough latency for the client to make use of the response.” The Google SRE book makes the same point about work finished too late in its chapter Addressing Cascading Failures: “The work the server did to respond is then wasted, and clients may retry the RPCs, leading to even more overload.” We will count both, and a column called wasted will show how much of the service’s effort went to requests whose client was already gone.
Backpressure, load shedding, and rate limiting
Backpressure is what a full stage does to the stage feeding it: make it wait, or make it stop. Inside a program, that is a bounded queue whose put call waits for room. Load shedding is the other choice: refuse the extra work. The SRE book defines it this way: “Load shedding drops some proportion of load by dropping traffic as the server approaches overload conditions.” It also states the goal: “to keep the server from running out of RAM, failing health checks, serving with extremely high latency, or any of the other symptoms associated with overload, while still doing as much useful work as it can.”
Rate limiting is a different tool. It caps how much each client may send. That is useful for fairness, but the SRE book warns about its limits: “because rate limiting often doesn’t take overall service health into account, it may not be able to stop a failure that has already begun.” Shedding based on the service’s own state is what protects it when many well-behaved clients arrive at once.
HTTP 503 and Retry-After
HTTP has a status code for exactly this situation. RFC 9110 says: “The 503 (Service Unavailable) status code indicates that the server is currently unable to handle the request due to a temporary overload or scheduled maintenance, which will likely be alleviated after some delay.” The server may add a Retry-After header “to suggest an appropriate amount of time for the client to wait before retrying the request.” In its simplest form the value is a number of seconds: “A delay-seconds value is a non-negative decimal integer, representing time in seconds.”
Prerequisites
- Python 3.11 or newer. The code catches the built-in
TimeoutError. The Python documentation notes thatasyncio.TimeoutError“was made an alias of TimeoutError” in version 3.11, so older versions would need a differentexceptclause. - A terminal on Windows, macOS, or Linux. I ran everything on Windows 11 with Python 3.13.14, FastAPI 0.143.0, uvicorn 0.54.0, httpx 0.28.1, and pytest 9.1.1.
- Basic Python and a first look at
asyncandawait. You do not need any prior knowledge of queues, servers, or load testing; each idea is explained when it first appears. - About four minutes of machine time to run every script in the lab, because the load tests run for real seconds.
- Optional:
curland a bash shell (Git Bash on Windows works) for the hand check near the end.
Set up the lab folder
Create an empty folder, make a virtual environment inside it, and install the four packages. Every file in this tutorial goes in that one folder.
mkdir shedding-lab
cd shedding-lab
python -m venv .venv
# Windows: .venv\Scripts\activate
# macOS/Linux: source .venv/bin/activate
python -m pip install fastapi uvicorn httpx pytest
python -c "import fastapi, uvicorn, httpx, pytest; print(fastapi.__version__, uvicorn.__version__, httpx.__version__, pytest.__version__)"
The last command prints the four versions. These are the ones I used:
0.143.0 0.54.0 0.28.1 9.1.1
Newer versions should behave the same way. If an import fails, the virtual environment is probably not activated; run the activation line again.
Create the three building blocks
Three small modules do the real work. Create them first; every step after this imports them.
The admission controller
This is the heart of the tutorial. It is short, so read it slowly.
# admission.py
"""Admission control for an asyncio service.
Every request gets one of three answers: run now, wait in a short line,
or be turned away immediately.
"""
import asyncio
from collections import Counter
from contextlib import asynccontextmanager
class Overloaded(Exception):
"""Raised instead of doing work when the service is full."""
def __init__(self, reason, retry_after=1):
super().__init__(reason)
self.reason = reason
self.retry_after = retry_after
class AdmissionController:
def __init__(self, max_concurrent, max_waiting=None, max_wait_s=None, retry_after=1):
self.max_waiting = max_waiting # longest allowed line (None = no limit)
self.max_wait_s = max_wait_s # longest allowed wait in line (None = no limit)
self.retry_after = retry_after # seconds we ask turned-away clients to stay away
self._slots = asyncio.Semaphore(max_concurrent)
self.waiting = 0
self.peak_waiting = 0
self.stats = Counter()
@asynccontextmanager
async def slot(self):
if self._slots.locked(): # every slot is busy, so this request would have to queue
if self.max_waiting is not None and self.waiting >= self.max_waiting:
self.stats["shed_line_full"] += 1
raise Overloaded("line_full", self.retry_after)
self.waiting += 1
self.peak_waiting = max(self.peak_waiting, self.waiting)
try:
await asyncio.wait_for(self._slots.acquire(), self.max_wait_s)
except TimeoutError:
self.stats["shed_waited_too_long"] += 1
raise Overloaded("waited_too_long", self.retry_after) from None
finally:
self.waiting -= 1
else:
await self._slots.acquire()
self.stats["admitted"] += 1
try:
yield
finally:
self._slots.release()
Here is what each part does and why it is there.
Overloadedis the exception the controller raises instead of doing work. It carries a shortreason, which is useful in logs and metrics, and aretry_afterhint in seconds that Step 6 turns into an HTTP header.asyncio.Semaphore(max_concurrent)is the set of slots. A semaphore is a counter of free permits:acquire()takes one, or waits if none is left, andrelease()gives it back. With four slots, at most four requests are being worked on at once.self._slots.locked()asks whether a request would have to wait right now, because every slot is busy. If not, the request takes a slot and runs at once. That is the fast path, and it never counts against any limit.- If every slot is busy, the request would have to wait in line. Before it joins, the controller applies the first limit,
max_waiting, the longest line it will allow. A request that finds the line full is turned away immediately with the reasonline_full. - A request that does join waits with
asyncio.wait_for(..., self.max_wait_s). That is the second limit: the longest any request may wait. When the time is up,wait_forcancels the waiting and raisesTimeoutError, and the controller turns that intoOverloaded("waited_too_long"). - The
finally: self.waiting -= 1runs when the request gets a slot, when it times out, and when it is cancelled. That single line is easy to get wrong. Step 8 shows what happens if you do. - Both limits default to
None, which means no limit. A controller with no limits is the unprotected service we are about to break. @asynccontextmanagerlets callers writeasync with ctrl.slot():. The slot is released when the block ends, even if the work inside raises an exception.
The load generator
A load generator is a program that sends requests at a chosen rate so you can measure how a service behaves. There are two kinds. A closed-loop generator has a fixed number of virtual users who each wait for an answer before sending the next request, so it slows down automatically when the service slows down, which hides overload. An open-loop generator sends requests on a schedule no matter what, the way real users on the internet behave. We need the open-loop kind. The one below sends request number i at time i / rate.
# loadgen.py
"""Open-loop load generator: requests arrive on a fixed schedule whether or not
earlier requests have finished, which is how real users behave."""
import asyncio
import time
from collections import Counter
from admission import Overloaded
def percentile(values, p):
if not values:
return float("nan")
values = sorted(values)
return values[min(len(values) - 1, round(p / 100 * (len(values) - 1)))]
def summarize(outcomes, duration, warmup):
"""outcomes is a list of (second of arrival, status, seconds the client waited)."""
counts = Counter(status for _, status, _ in outcomes)
waits = [waited for _, status, waited in outcomes if status == "ok"]
steady = sum(1 for second, status, _ in outcomes if status == "ok" and second >= warmup)
by_second = {}
for second, status, _ in outcomes:
cell = by_second.setdefault(second, [0, 0])
cell[0] += 1
cell[1] += status == "ok"
return {
"offered": len(outcomes),
"ok": counts["ok"],
"shed": counts["shed"],
"timeout": counts["timeout"],
"error": counts["error"],
"goodput": counts["ok"] / duration,
"steady_goodput": steady / (duration - warmup),
"p50_ms": percentile(waits, 50) * 1000,
"p99_ms": percentile(waits, 99) * 1000,
"by_second": dict(sorted(by_second.items())),
}
def _quiet(task):
if not task.cancelled():
task.exception() # mark any exception as retrieved so asyncio stays quiet
async def open_loop(call, rate, duration, client_timeout, cancel_on_timeout=False, warmup=2):
"""Send `rate` requests per second for `duration` seconds.
call() is an async function. It returns when the request is served and
raises Overloaded when the request is turned away. When a client gives up
after `client_timeout` seconds the work is left running unless
cancel_on_timeout is set, because a server that has started a request
usually does not know the client has left.
"""
outcomes = []
background = []
async def one(arrive_at):
started = time.monotonic()
task = asyncio.create_task(call())
task.add_done_callback(_quiet)
background.append(task)
done, _ = await asyncio.wait({task}, timeout=client_timeout)
waited = time.monotonic() - started
if not done:
if cancel_on_timeout:
task.cancel()
status = "timeout"
elif isinstance(task.exception(), Overloaded):
status = "shed"
elif task.exception() is not None:
status = "error"
else:
status = "ok"
outcomes.append((int(arrive_at), status, waited))
clients = []
start = time.monotonic()
for i in range(int(rate * duration)):
arrive_at = i / rate
pause = start + arrive_at - time.monotonic()
if pause > 0:
await asyncio.sleep(pause)
clients.append(asyncio.create_task(one(arrive_at)))
await asyncio.gather(*clients)
for task in background:
task.cancel()
await asyncio.gather(*background, return_exceptions=True)
return summarize(outcomes, duration, warmup)
- For every request,
one()startscall()as its own task, then waits at mostclient_timeoutseconds for it. That wait is the client’s patience. - If
call()finishes normally the outcome isok. If it raisesOverloadedthe outcome isshed. If the patience runs out first the outcome istimeout. - When a client gives up, the work is left running by default (
cancel_on_timeout=False). A server that has already started a request usually does not know the client has left, and that is exactly the waste we want to measure. Step 6 uses real HTTP clients that do hang up. summarize()turns the outcomes into goodput, counts, the median and 99th-percentile time of the answered requests, and a per-second table.
The shared settings
The third module holds the settings every step shares and one helper, run_policy, that builds a small in-process service with whatever limits you pass in, runs the load generator against it, and returns the results.
# lab_common.py
import asyncio
import json
import pathlib
from admission import AdmissionController
from loadgen import open_loop
OUT = pathlib.Path(__file__).parent / "out"
SERVICE_TIME = 0.05 # seconds of pretend work per request
SLOTS = 4 # requests the service works on at the same time
CLIENT_TIMEOUT = 1.0 # clients stop waiting after this many seconds
def load_capacity():
return json.loads((OUT / "capacity.json").read_text())["capacity_rps"]
def healthy():
return SERVICE_TIME
async def run_policy(rate, duration, service_time=healthy, **limits):
"""Run one open-loop test against an in-process service that uses `limits`."""
ctrl = AdmissionController(SLOTS, **limits)
done = {"work": 0}
async def call():
async with ctrl.slot():
await asyncio.sleep(service_time())
done["work"] += 1
result = await open_loop(call, rate, duration, CLIENT_TIMEOUT)
result["work"] = done["work"]
result["wasted"] = done["work"] - result["ok"]
result["peak_line"] = ctrl.peak_waiting
return result
HEADER = (f"{'policy':<34}{'goodput/s':>10}{'ok':>6}{'shed':>6}{'timeout':>8}"
f"{'p50 ms':>8}{'p99 ms':>8}{'wasted':>8}{'peak line':>10}")
def row(label, r):
return (f"{label:<34}{r['goodput']:>10.1f}{r['ok']:>6}{r['shed']:>6}{r['timeout']:>8}"
f"{r['p50_ms']:>8.0f}{r['p99_ms']:>8.0f}{r['wasted']:>8}{r['peak_line']:>10}")
- The pretend service has four slots, and each request needs about 50 ms of work (
asyncio.sleepstands in for a database call or a model call). Clients stop waiting after one second. wastedis the work the service finished minus the answers that arrived in time: effort spent on requests whose client had already gone.HEADERandrowprint the results table you will see in every step.
Check that everything loads. This command should print nothing and return to the prompt; an error here means one of the three files has a typo or is in the wrong folder.
python -c "import lab_common"
Step 1: Measure your service’s real capacity
Before you can overload a service you need to know where its limit is. Guessing is risky, as this step shows.
# step01_calibrate.py
import asyncio
import json
import time
from admission import AdmissionController
from lab_common import OUT, SERVICE_TIME, SLOTS
async def main():
samples = []
for _ in range(20):
started = time.perf_counter()
await asyncio.sleep(SERVICE_TIME)
samples.append(time.perf_counter() - started)
actual = sum(samples) / len(samples)
print(f"asked for {SERVICE_TIME * 1000:.0f} ms of work; one request really took {actual * 1000:.1f} ms")
ctrl = AdmissionController(SLOTS)
finished = 0
seconds = 3.0
stop = time.monotonic() + seconds
async def keep_busy():
nonlocal finished
while time.monotonic() < stop:
async with ctrl.slot():
await asyncio.sleep(SERVICE_TIME)
finished += 1
await asyncio.gather(*(keep_busy() for _ in range(SLOTS * 3)))
capacity = finished / seconds
print(f"{SLOTS} slots kept busy for {seconds:.0f} s finished {finished} requests")
print(f"capacity = {capacity:.1f} requests per second")
OUT.mkdir(exist_ok=True)
(OUT / "capacity.json").write_text(json.dumps({"capacity_rps": round(capacity, 1)}))
asyncio.run(main())
The script does two things. First it times the pretend work: asyncio.sleep(0.05) should take 50 ms. Then it keeps all four slots busy for three seconds, using twelve looping workers (three for each slot), and counts how many requests finish. That finish rate is the capacity. It saves the number to out/capacity.json, and every later step reads it so that the load you offer is always relative to the capacity measured on your machine.
python step01_calibrate.py
asked for 50 ms of work; one request really took 64.3 ms
4 slots kept busy for 3 s finished 196 requests
capacity = 65.3 requests per second
On my Windows machine 50 ms of work took 64.3 ms, because asyncio sleeps on Windows round up to the next system clock tick, about 16 ms apart (a 10 ms sleep took 16 ms when I timed it). Four slots working on 50 ms requests would suggest 80 requests per second, but the measured capacity was 65.3. On Linux or macOS you will probably see numbers closer to 50 ms and 80 requests per second. Both are fine. The point is that you measured instead of assuming. One caution: a short probe like this reads a little high, because requests that are still in flight when the three seconds end are counted too. Step 9 shows that the steady rate on my machine is about 62 requests per second, which is what four slots at 64.3 ms each give, so treat the Step 1 number as accurate to within about five percent.
Check before moving on: the script printed a capacity and created out/capacity.json. If a later step raises FileNotFoundError for that file, run this step first.
Step 2: See backpressure in a plain asyncio queue
The simplest place to meet backpressure is a queue between a producer and a consumer. This script runs a producer that can make 200 items per second against a consumer that needs about 20 ms per item, three times, with three kinds of queue.
# step02_queue.py
import asyncio
import time
RUN_SECONDS = 2.0
PRODUCE_PER_SECOND = 200 # the producer can make 200 items a second
CONSUME_TIME = 0.02 # the consumer needs about 20 ms per item
async def producer(queue, stats, drop_when_full):
start = time.monotonic()
n = 0
while time.monotonic() - start < RUN_SECONDS:
pause = start + n / PRODUCE_PER_SECOND - time.monotonic()
if pause > 0:
await asyncio.sleep(pause)
if drop_when_full:
try:
queue.put_nowait(time.monotonic()) # raises QueueFull instead of waiting
except asyncio.QueueFull:
stats["dropped"] += 1
else:
await queue.put(time.monotonic()) # waits here while a bounded queue is full
n += 1
stats["produced"] = n
async def consumer(queue, stats):
while True:
queued_at = await queue.get()
stats["oldest_wait"] = max(stats["oldest_wait"], time.monotonic() - queued_at)
await asyncio.sleep(CONSUME_TIME)
stats["consumed"] += 1
async def trial(label, maxsize, drop_when_full=False):
queue = asyncio.Queue(maxsize=maxsize) # maxsize=0 means no limit
stats = {"consumed": 0, "dropped": 0, "oldest_wait": 0.0}
worker = asyncio.create_task(consumer(queue, stats))
await producer(queue, stats, drop_when_full)
left = queue.qsize()
worker.cancel()
print(f"{label:<22} produced={stats['produced']:>4} consumed={stats['consumed']:>3} "
f"dropped={stats['dropped']:>4} left in queue={left:>4} "
f"longest wait={stats['oldest_wait']:.2f} s")
async def main():
await trial("unbounded queue", 0)
await trial("maxsize=10, wait", 10)
await trial("maxsize=10, drop", 10, drop_when_full=True)
asyncio.run(main())
The Python documentation describes the queue’s size limit like this: “If maxsize is less than or equal to zero, the queue size is infinite.” With a positive limit, “await put() blocks when the queue reaches maxsize until an item is removed by get()”. That blocking is backpressure. The alternative is put_nowait, which, in the documentation’s words, will “raise QueueFull” if “no free slot is immediately available”. That is shedding.
python step02_queue.py
unbounded queue produced= 398 consumed= 62 dropped= 0 left in queue= 335 longest wait=1.67 s
maxsize=10, wait produced= 74 consumed= 63 dropped= 0 left in queue= 10 longest wait=0.36 s
maxsize=10, drop produced= 399 consumed= 62 dropped= 326 left in queue= 10 longest wait=0.32 s
Read the three lines one at a time.
- Unbounded queue: the producer made 398 items, the consumer handled 62, and 335 were still waiting when the run ended. The last item the consumer picked up had waited 1.67 seconds. Memory and delay both grow for as long as the producer outruns the consumer.
- maxsize=10, wait: the producer was forced to slow down. It made only 74 items, the consumer handled 63, and nothing was lost. The longest wait dropped to 0.36 seconds.
- maxsize=10, drop: the producer kept its full pace but 326 items were thrown away when the queue was full. The longest wait was 0.32 seconds.
Notice what did not change: the consumer finished about 62 items in every case. Protecting the queue did not cost any useful work. It stopped the pile from growing, and the items that were served waited about five times less. A web service sits at the edge, where you cannot make a browser wait forever, so there the two options become waiting for a bounded time and saying no with a 503. The rest of this tutorial builds exactly that.
Check before moving on: the unbounded line leaves hundreds of items in the queue, both bounded lines leave 10 or fewer, and the consumer count is nearly the same in all three runs.
Step 3: Watch an unprotected service collapse
Now the experiment that motivates everything else. The service has no limits: every request either gets a slot or joins an unlimited line. We offer twice the measured capacity for five seconds. Clients are patient for one second and then give up.
# step03_no_limits.py
import asyncio
from lab_common import CLIENT_TIMEOUT, HEADER, load_capacity, row, run_policy
DURATION = 5
async def main():
capacity = load_capacity()
rate = 2 * capacity
print(f"capacity {capacity:.0f}/s, offering {rate:.0f}/s for {DURATION} s, "
f"clients give up after {CLIENT_TIMEOUT:.0f} s")
result = await run_policy(rate, DURATION)
print(HEADER)
print(row("no limits", result))
print()
print("second of arrival arrived answered in time")
for second, (arrived, ok) in result["by_second"].items():
print(f"{second:>17}{arrived:>10}{ok:>19}")
asyncio.run(main())
python step03_no_limits.py
capacity 65/s, offering 131/s for 5 s, clients give up after 1 s
policy goodput/s ok shed timeout p50 ms p99 ms wasted peak line
no limits 23.4 117 0 536 526 993 256 339
second of arrival arrived answered in time
0 131 117
1 131 0
2 130 0
3 131 0
4 130 0
The table columns are the same in every step:
- goodput/s is the number of answers delivered within the client’s one second of patience, divided by the five seconds of the test.
- ok, shed, and timeout count requests answered in time, turned away immediately, and abandoned by the client.
- p50 ms and p99 ms are the median and 99th-percentile time of the answered requests only.
- wasted is the number of requests the service finished after the client had gone.
- peak line is the longest line of waiting requests.
Look at what happened. Of the 653 requests that arrived, only 117 were answered in time. The service was busy every moment, yet 256 of the 373 requests it completed (about 69 percent) were for clients who had already given up. The per-second table shows the cliff: in the first second 117 of 131 requests were answered in time, and from the following second onward not one was. The line grew by about 65 requests every second, so a request arriving t seconds after the start waited about t seconds, and the one-second patience ran out at t = 1. After that, every answer was too late.
One more lesson hides in this table. The median answer time of the survivors looks acceptable (526 ms), but 82 percent of requests failed. If you only watch the latency of successful requests, a collapsing service can look fine. Always read goodput and failure counts next to the percentiles.
Check before moving on: you should see the same shape on your machine: a healthy first second, then zero requests answered in time, a large timeout count, a large wasted count, and a peak line of hundreds.
Step 4: Cap the line to bound the delay
The first fix is to limit how many requests may wait. How long should the line be? Use the same arithmetic as before, in reverse: the delay you can afford to add equals the line length divided by the capacity, so the line length is the capacity times that delay. With 65 requests per second and a budget of half a second, the line may hold 33 requests.
# step04_bounded_line.py
import asyncio
from lab_common import HEADER, load_capacity, row, run_policy
DURATION = 5
LINE_BUDGET_S = 0.5 # the extra delay we are willing to add by queueing
async def main():
capacity = load_capacity()
max_waiting = round(capacity * LINE_BUDGET_S) # Little's law: line length = rate x time
print(f"{capacity:.0f}/s x {LINE_BUDGET_S} s budget -> a line of at most {max_waiting}")
rate = 2 * capacity
print(HEADER)
print(row("no limits", await run_policy(rate, DURATION)))
print(row(f"line capped at {max_waiting}", await run_policy(rate, DURATION, max_waiting=max_waiting)))
asyncio.run(main())
python step04_bounded_line.py
65/s x 0.5 s budget -> a line of at most 33
policy goodput/s ok shed timeout p50 ms p99 ms wasted peak line
no limits 23.4 117 0 536 545 1011 256 340
line capped at 33 69.6 348 305 0 588 620 0 33
Same offered load, same five seconds, very different result. The service turned away 305 requests immediately and answered 348, with a 99th percentile of 620 ms. No client timed out and no work was wasted. The median of 588 ms is close to the arithmetic: half a second in line plus 64 ms of work. (The goodput column reads above the service’s steady rate of about 62 requests per second because the 37 requests that are in the line or in service when the last request arrives, 33 waiting plus 4 working, are still answered afterward and count too. Dividing 37 by five seconds adds about 7 per second.)
This is the tactic the SRE book describes: “one effective approach is to return an HTTP 503 (service unavailable) to any incoming request when there are more than a given number of client requests in flight.” The book also explains the trade-off in the line length: “Queued requests consume memory and increase latency.” For steady traffic it suggests small queues, “e.g., 50% or less” of the thread pool size, while bursty traffic may justify a longer one. A budget-based number like ours sits at the bursty end of that range. A shorter line sheds more during harmless bursts; a longer line adds delay. Pick the delay you can afford and compute the length from it.
Check before moving on: the no limits row repeats the collapse from Step 3, while the capped row shows zero timeouts, zero wasted work, and a peak line equal to the cap.
Step 5: Cap the wait too, because a count is not a time
A line length is a count, but your clients feel time. The same 33 requests wait half a second when the service is healthy and much longer when it slows down. Real services slow down all the time: a database gets busy, a downstream API degrades, a garbage collection pause hits. This experiment simulates exactly that. We offer a modest load of 80 percent of capacity, which a healthy service handles easily, and after three seconds every request starts taking four times as long.
# step05_slowdown.py
import asyncio
import time
from lab_common import HEADER, SERVICE_TIME, load_capacity, row, run_policy
DURATION = 8
SLOW_AFTER = 3 # seconds into the test
SLOWDOWN = 4 # a dependency becomes 4x slower
async def main():
capacity = load_capacity()
max_waiting = round(capacity * 0.5)
rate = 0.8 * capacity # below capacity, so a healthy service copes easily
print(f"offering {rate:.0f}/s (80% of capacity); after {SLOW_AFTER} s every request takes {SLOWDOWN}x longer")
print(HEADER)
policies = [
(f"line capped at {max_waiting}", {"max_waiting": max_waiting}),
(f"cap {max_waiting} + wait limit 0.3 s", {"max_waiting": max_waiting, "max_wait_s": 0.3}),
]
for label, limits in policies:
t0 = time.monotonic()
def service_time():
slow = time.monotonic() - t0 > SLOW_AFTER
return SERVICE_TIME * (SLOWDOWN if slow else 1)
print(row(label, await run_policy(rate, DURATION, service_time=service_time, **limits)))
asyncio.run(main())
python step05_slowdown.py
offering 52/s (80% of capacity); after 3 s every request takes 4x longer
policy goodput/s ok shed timeout p50 ms p99 ms wasted peak line
line capped at 33 22.9 183 131 103 64 881 84 33
cap 33 + wait limit 0.3 s 32.5 260 157 0 65 498 0 16
The first row is the line cap from Step 4. After the slowdown the four slots could finish only about 20 requests per second, so the same 33-request line meant a wait of roughly a second and a half, longer than the client’s patience. 103 requests timed out, the service finished 84 of them anyway, and the 99th percentile climbed to 881 ms. The second row adds the wait limit of 0.3 seconds, the max_wait_s argument the controller already supported. A request that has waited that long is turned away instead of being served late. There were no timeouts, no wasted work, the 99th percentile stayed at 498 ms, and goodput was higher (32.5 against 22.9) because the service spent its time on requests that could still succeed. The line only needed to reach 16 because the wait limit kept it short on its own.
The SRE book gives the reason a late answer is pointless: “If a user’s web search is slow because an RPC has been queued for 10 seconds, there’s a good chance the user has given up and refreshed their browser, issuing another request: there’s no point in responding to the first one, since it will be ignored!” The same passage points to more refined versions of this idea, last-in first-out queues and the controlled delay (CoDel) algorithm, which remove requests “that are unlikely to be worth processing”. A wait limit is the simplest member of that family.
How do you choose the number? Keep the wait limit comfortably below the client’s patience minus the time the work itself takes. Here the clients wait one second and the work takes about 65 ms when healthy and about 200 ms when slow, so 0.3 seconds leaves room for both. Keep the line cap as well: the wait limit bounds delay, and the cap bounds memory.
Check before moving on: with only the cap you should see timeouts after the slowdown, and with the wait limit added the timeouts and the wasted work should both be zero.
Step 6: Serve it over HTTP with a 503 and Retry-After
So far the service lived inside the test program. Now put the same controller in front of a real web service and run the same experiment over real HTTP.
The service
The service lives in service.py. You can start it by hand with python -m uvicorn service:app, and the load test below starts it for you.
# service.py
import asyncio
import os
from fastapi import FastAPI
from fastapi.responses import JSONResponse
from admission import AdmissionController, Overloaded
def make_app(ctrl, service_time):
app = FastAPI()
counters = {"work_done": 0}
@app.exception_handler(Overloaded)
async def turn_away(request, exc):
return JSONResponse(
{"error": "overloaded", "reason": exc.reason},
status_code=503,
headers={"Retry-After": str(exc.retry_after)},
)
@app.get("/work")
async def work():
async with ctrl.slot():
await asyncio.sleep(service_time)
counters["work_done"] += 1
return {"ok": True}
@app.get("/health")
async def health(): # liveness checks must never wait behind the limiter
return {"status": "up"}
@app.get("/stats")
async def stats():
return {"work_done": counters["work_done"], "peak_line": ctrl.peak_waiting, **ctrl.stats}
return app
def from_env(name, cast):
return cast(os.environ[name]) if name in os.environ else None
app = make_app(
AdmissionController(
max_concurrent=4,
max_waiting=from_env("MAX_WAITING", int),
max_wait_s=from_env("MAX_WAIT_S", float),
),
service_time=float(os.environ.get("SERVICE_TIME", "0.05")),
)
make_app(ctrl, service_time)builds the FastAPI application around any controller. Having a factory means the tests in Step 8 can build an app with exactly the limits they want.- The function decorated with
@app.exception_handler(Overloaded)is a custom exception handler. WheneverOverloadedescapes a route, FastAPI calls it, and it answers with status 503, a JSON body that names the reason, and aRetry-Afterheader. /workwraps the pretend work inasync with ctrl.slot():. Everything else about the route is ordinary./healthdeliberately does not use the controller. A health check that waits behind the line looks like a failure to whatever is watching it, and the SRE book lists failing health checks among the symptoms of overload. Keep liveness checks outside the limiter./statsreports the work finished, the peak line, and the controller’s counters, so a test can see what the server did.- The limits come from environment variables (
MAX_WAITING,MAX_WAIT_S) so the same file can run unprotected or protected without editing it.
The real HTTP load test
This script starts uvicorn as a child process, measures the capacity over HTTP, and then repeats the two-times overload with real clients that really hang up when they lose patience (cancel_on_timeout=True). It runs three servers in turn: unprotected, protected by our controller, and unprotected but started with uvicorn’s own --limit-concurrency flag, which we will meet in a moment.
# step06_http.py
import asyncio
import os
import subprocess
import sys
import time
import httpx
from admission import Overloaded
from lab_common import CLIENT_TIMEOUT, HEADER, OUT, row
from loadgen import open_loop
DURATION = 5
class Server:
"""Run `uvicorn service:app` in a child process and stop it afterwards."""
def __init__(self, port, env=None, flags=()):
self.port = port
self.env = {**os.environ, **(env or {})}
self.flags = list(flags)
self.base = f"http://127.0.0.1:{port}"
self.log_path = OUT / f"uvicorn_{port}.log"
def __enter__(self):
command = [sys.executable, "-m", "uvicorn", "service:app", "--port", str(self.port),
"--log-level", "warning", *self.flags]
self.log = open(self.log_path, "wb") # keep uvicorn's own log out of our output
self.process = subprocess.Popen(command, env=self.env, stdout=self.log, stderr=subprocess.STDOUT)
for _ in range(150):
try:
if httpx.get(self.base + "/health", timeout=1).status_code == 200:
return self
except httpx.TransportError:
pass
time.sleep(0.1)
raise RuntimeError("server did not start")
def __exit__(self, *exc):
if os.name == "nt": # the venv's python.exe is a launcher, so stop the whole process tree
subprocess.run(["taskkill", "/PID", str(self.process.pid), "/T", "/F"], capture_output=True)
else:
self.process.terminate()
self.process.wait()
self.log.close()
def log_lines_containing(self, text):
return self.log_path.read_text(errors="replace").count(text)
async def measure_capacity(base, seconds=3, workers=16):
"""Keep `workers` requests in flight at all times; the finish rate is the capacity."""
finished = 0
stop = time.monotonic() + seconds
async with httpx.AsyncClient(base_url=base, limits=httpx.Limits(max_connections=None)) as client:
async def worker():
nonlocal finished
while time.monotonic() < stop:
if (await client.get("/work")).status_code == 200:
finished += 1
await asyncio.gather(*(worker() for _ in range(workers)))
return finished / seconds
async def overload(base, rate):
example = []
limits = httpx.Limits(max_connections=None, max_keepalive_connections=100)
async with httpx.AsyncClient(base_url=base, limits=limits, timeout=None) as client:
async def call():
reply = await client.get("/work")
if reply.status_code == 503:
if not example:
example.append((reply.headers.get("retry-after"), reply.text[:60]))
raise Overloaded("http 503")
reply.raise_for_status()
before = (await client.get("/stats")).json()["work_done"] # ignore the capacity probe's work
result = await open_loop(call, rate, DURATION, CLIENT_TIMEOUT, cancel_on_timeout=True)
stats = (await client.get("/stats")).json()
result["work"] = stats["work_done"] - before
result["wasted"] = result["work"] - result["ok"]
result["peak_line"] = stats["peak_line"]
return result, (example[0] if example else None)
def main():
with Server(8761) as server:
capacity = asyncio.run(measure_capacity(server.base))
rate = 2 * capacity
print(f"capacity over HTTP: {capacity:.1f} requests/s; offering {rate:.0f}/s for {DURATION} s")
print(HEADER)
bare, _ = asyncio.run(overload(server.base, rate))
print(row("no limits", bare))
max_waiting = round(capacity * 0.5)
runs = [
(f"cap {max_waiting} + wait limit 0.3 s",
dict(port=8762, env={"MAX_WAITING": str(max_waiting), "MAX_WAIT_S": "0.3"})),
("uvicorn --limit-concurrency 36",
dict(port=8763, flags=["--limit-concurrency", "36"])),
]
for label, kwargs in runs:
with Server(**kwargs) as server:
result, example = asyncio.run(overload(server.base, rate))
print(row(label, result))
print(f" first 503 seen: Retry-After={example[0]!r}, body={example[1]!r}")
refused = server.log_lines_containing("Exceeded concurrency limit")
if refused:
print(f" uvicorn logged {refused} lines saying 'Exceeded concurrency limit.'")
main()
Two details are worth knowing. The Server class stops uvicorn with taskkill on Windows because a virtual environment’s python.exe is a launcher that starts a second process, and stopping only the launcher would leave the real server running; on other systems a plain terminate() is enough. And overload() asks the server for its work counter before and after the test, so the capacity probe’s work is not counted as wasted.
python step06_http.py
capacity over HTTP: 61.3 requests/s; offering 123/s for 5 s
policy goodput/s ok shed timeout p50 ms p99 ms wasted peak line
no limits 33.2 166 0 447 504 997 263 246
cap 31 + wait limit 0.3 s 73.6 368 245 0 356 372 0 31
first 503 seen: Retry-After='1', body='{"error":"overloaded","reason":"line_full"}'
uvicorn --limit-concurrency 36 78.0 390 223 0 457 533 0 31
first 503 seen: Retry-After=None, body='Service Unavailable'
uvicorn logged 223 lines saying 'Exceeded concurrency limit.'
- No limits: the real server collapsed exactly like the in-process model. Only 166 of the 613 requests (166 answered, 447 abandoned) were answered in time, and the server finished 263 requests that nobody was waiting for, even though each client cancelled its request when it gave up. A handler that has started keeps running; nothing tells it the client left.
- Our controller (line of 31 plus a 0.3 second wait limit): 368 answers in time, 245 requests turned away at once, no timeouts, no wasted work, and a 99th percentile of 372 ms. The first 503 it sent carried
Retry-After: 1and a JSON body saying why.
What uvicorn’s built-in limit does
The third row came from starting an otherwise unprotected service with uvicorn service:app --limit-concurrency 36. The uvicorn documentation describes the flag as “Maximum number of concurrent connections or tasks to allow, before issuing HTTP 503 responses” and says it is “Useful for ensuring known memory usage patterns even under over-resourced loads.” The source shows how it counts. This is the excerpt from uvicorn 0.54.0:
# uvicorn/protocols/http/h11_impl.py (uvicorn 0.54.0, excerpt)
# Handle 503 responses when 'limit_concurrency' is exceeded.
if self.limit_concurrency is not None and (
len(self.connections) >= self.limit_concurrency or len(self.tasks) >= self.limit_concurrency
):
app = service_unavailable
message = "Exceeded concurrency limit."
self.logger.warning(message)
It is a good first line of defense. It answered 390 requests in time and refused 223, close to our controller’s result. But compare the details. It counts open connections as well as running requests, it is a fixed number (the weakness Step 5 exposed), the 503 it sends has no Retry-After header, a plain-text body (Service Unavailable), and a connection: close header, so the client must reconnect (see service_unavailable in flow_control.py), and it wrote 223 log lines reading “Exceeded concurrency limit.” during the run, one for every refusal. Logging while overloaded is work, and the AWS article is blunt about it: “the act of shedding load isn’t free, so eventually the server falls prey to Amdahl’s law and goodput drops.” Use the built-in flag as a safety net if you like, and put your own controller in front of it for the signals and the control it cannot give.
Check before moving on: the protected row has no timeouts and no wasted work, its printed 503 shows a Retry-After value, and the uvicorn row shows Retry-After=None.
Step 7: Make clients part of the solution
A 503 only helps if clients treat it as a signal. If every turned-away client retries at once, the service receives the same traffic again, plus more. The AWS article spells out how this compounds across layers: if each layer performs a number of retries, “an overload in the bottom layer causes cascading retries that amplify the offered load exponentially.” The SRE book gives the rules for sensible retrying: “Always use randomized exponential backoff when scheduling retries.” “Limit retries per request. Don’t retry a given request indefinitely.” And: “Consider having a server-wide retry budget. For example, only allow 60 retries per minute in a process, and if the retry budget is exceeded, don’t retry; just fail the request.”
This step compares three client behaviors against the protected service from Step 5: no retries, immediate retries, and polite retries. Polite means three things. The client waits at least as long as Retry-After asks, because RFC 9110 describes the header as telling the client how long it “ought to wait”. It adds random jitter on top, so clients do not all return at the same instant. And a small retry budget caps retries at about a tenth of first attempts, so retries can never become the main traffic.
# step07_retries.py
import asyncio
import random
import time
from admission import AdmissionController, Overloaded
from lab_common import CLIENT_TIMEOUT, SERVICE_TIME, SLOTS, load_capacity
DURATION = 6
TRIES = 3
class RetryBudget:
"""Allow retries up to `ratio` of first attempts, plus a small burst."""
def __init__(self, ratio=0.1, burst=10):
self.ratio = ratio
self.burst = burst
self.tokens = burst
def first_attempt(self):
self.tokens = min(self.burst, self.tokens + self.ratio)
def allow_retry(self):
if self.tokens >= 1:
self.tokens -= 1
return True
return False
async def one_attempt(ctrl):
async with ctrl.slot():
await asyncio.sleep(SERVICE_TIME)
async def user(ctrl, policy, budget, tally):
tries = 1 if policy == "no retries" else TRIES
budget.first_attempt()
for n in range(tries):
tally["attempts"] += 1
try:
await asyncio.wait_for(one_attempt(ctrl), CLIENT_TIMEOUT)
tally["served"] += 1
return
except Overloaded as turned_away:
if policy == "Retry-After + budget":
if n == tries - 1 or not budget.allow_retry():
return
await asyncio.sleep(turned_away.retry_after * random.uniform(1.0, 2.0))
else:
await asyncio.sleep(0) # "immediate retries": knock again right away
async def trial(policy, rate, max_waiting):
ctrl = AdmissionController(SLOTS, max_waiting=max_waiting, max_wait_s=0.3)
budget = RetryBudget()
tally = {"attempts": 0, "served": 0}
users = []
start = time.monotonic()
count = int(rate * DURATION)
for i in range(count):
pause = start + i / rate - time.monotonic()
if pause > 0:
await asyncio.sleep(pause)
users.append(asyncio.create_task(user(ctrl, policy, budget, tally)))
await asyncio.gather(*users)
print(f"{policy:<22}{count:>7}{tally['attempts']:>10}{tally['attempts'] / count:>11.2f}"
f"{tally['served']:>9}{100 * tally['served'] / count:>8.0f}%{tally['attempts'] / DURATION:>14.0f}")
async def main():
random.seed(7)
capacity = load_capacity()
rate = 2 * capacity
print(f"service capacity {capacity:.0f}/s, {rate:.0f} new users per second for {DURATION} s")
print(f"{'client policy':<22}{'users':>7}{'attempts':>10}{'per user':>11}{'served':>9}{'':>8}{'attempts/s':>14}")
for policy in ("no retries", "immediate retries", "Retry-After + budget"):
await trial(policy, rate, round(capacity * 0.5))
asyncio.run(main())
RetryBudgetworks like a small bank account. Every first attempt deposits a tenth of a token, every retry withdraws one whole token, and a retry is refused when the balance is below one. Theburstvalue is the starting balance and the maximum balance.- The immediate policy sleeps for zero seconds, which just yields to the event loop, and tries again. The polite policy sleeps for
retry_aftertimes a random number between 1 and 2.
python step07_retries.py
service capacity 65/s, 131 new users per second for 6 s
client policy users attempts per user served attempts/s
no retries 783 783 1.00 394 50% 130
immediate retries 783 1597 2.04 408 52% 266
Retry-After + budget 783 864 1.10 414 53% 144
Read the attempts columns first. Immediate retries produced 2.04 attempts per user and 266 attempts per second hitting the front door, double the 130 per second of the no-retry run, for a gain of two percentage points of users served (52 percent against 50). The polite policy added only 10 percent more attempts (1.10 per user) and served 53 percent. I ran this step three times and got the same percentages every time, so treat the one or two point differences as small and the attempt counts as the real signal.
The deeper lesson is that retries do not create capacity. Demand here is twice the capacity, so at most about half the users can ever be served, and every policy lands near 50 percent. What a retry policy controls is how much extra work the service must refuse. Remember also that a retry is only safe when repeating the request is harmless. If the request charges a card, give it an idempotency key first; the tutorial on preventing duplicate charges with idempotency keys shows how, and the one on retrying with exponential backoff and jitter covers the client side in depth.
Check before moving on: the polite policy makes close to 1.0 attempts per user, while immediate retries make about 2.
Step 8: Test it, including the bug that leaks the line counter
Concurrency code deserves tests, and this controller has one bug that is easy to write and hard to see. The test will prove it can catch it.
Reproduce the bug
Here is the same controller with one line moved. The decrement of the line counter sits after the wait instead of in a finally block.
# leaky.py
"""The controller with one easy-to-make mistake, kept so a test can prove it catches it."""
import asyncio
from contextlib import asynccontextmanager
from admission import AdmissionController, Overloaded
class LeakyController(AdmissionController):
@asynccontextmanager
async def slot(self):
if self._slots.locked():
if self.max_waiting is not None and self.waiting >= self.max_waiting:
raise Overloaded("line_full", self.retry_after)
self.waiting += 1
try:
await asyncio.wait_for(self._slots.acquire(), self.max_wait_s)
except TimeoutError:
raise Overloaded("waited_too_long", self.retry_after) from None
self.waiting -= 1 # BUG: never reached when the wait times out or is cancelled
else:
await self._slots.acquire()
try:
yield
finally:
self._slots.release()
That version works while every request eventually gets a slot. But when a wait times out, or the waiting request is cancelled because the client disconnected, the exception skips the decrement, and the counter stays one too high forever. This script runs both versions through the same situation: one slot is held, and four requests each wait 30 ms and give up.
# step08_leak.py
import asyncio
from admission import AdmissionController, Overloaded
from leaky import LeakyController
async def ask(ctrl):
try:
async with ctrl.slot():
return "served"
except Overloaded as turned_away:
return turned_away.reason
async def scenario(controller_class):
ctrl = controller_class(1, max_waiting=2, max_wait_s=0.03)
gate = asyncio.Event()
async def hold_the_only_slot():
async with ctrl.slot():
await gate.wait()
blocker = asyncio.create_task(hold_the_only_slot())
await asyncio.sleep(0)
answers = [await ask(ctrl) for _ in range(4)] # each waits 30 ms, then gives up
print(f"{controller_class.__name__:<20} four requests got {answers}")
print(f"{'':<20} line counter afterwards: {ctrl.waiting} (nobody is waiting)")
asyncio.get_running_loop().call_later(0.01, gate.set) # the blocker finishes in 10 ms
print(f"{'':<20} a request that only has to wait 10 ms gets: {await ask(ctrl)}")
await blocker
for controller_class in (AdmissionController, LeakyController):
asyncio.run(scenario(controller_class))
python step08_leak.py
AdmissionController four requests got ['waited_too_long', 'waited_too_long', 'waited_too_long', 'waited_too_long']
line counter afterwards: 0 (nobody is waiting)
a request that only has to wait 10 ms gets: served
LeakyController four requests got ['waited_too_long', 'waited_too_long', 'line_full', 'line_full']
line counter afterwards: 2 (nobody is waiting)
a request that only has to wait 10 ms gets: line_full
The correct controller answers all four requests with waited_too_long and its counter returns to zero, so the final request, which only needs to wait 10 ms, is served. The leaky controller’s counter is stuck at 2 with nobody waiting, so the third and fourth requests are refused as line_full without waiting at all, and the final request is refused too. Each ghost in the counter permanently shrinks the line. After enough timeouts the service sheds every request that has to wait, even when the line is empty.
Write the tests
The test file below checks eight behaviors. hold_open starts a task that keeps one slot busy until a gate is opened, so each test can create exactly the situation it needs without sleeping for long. Notice test_the_check_above_catches_a_leaky_controller: it runs the leak check against the broken controller and expects it to fail. A test that has never been seen to fail may not be testing anything, so this one proves the check has teeth.
# test_admission.py
import asyncio
import time
import httpx
import pytest
from admission import AdmissionController, Overloaded
from leaky import LeakyController
from service import make_app
def hold_open(ctrl):
"""Start a task that keeps one slot busy until `gate` is set."""
gate = asyncio.Event()
async def hold():
async with ctrl.slot():
await gate.wait()
return asyncio.create_task(hold()), gate
async def ask(ctrl):
try:
async with ctrl.slot():
return "served"
except Overloaded as turned_away:
return turned_away.reason
def test_idle_service_serves_everything():
async def scenario():
ctrl = AdmissionController(2, max_waiting=0, max_wait_s=0)
return [await ask(ctrl) for _ in range(5)]
assert asyncio.run(scenario()) == ["served"] * 5
def test_full_line_turns_requests_away_immediately():
async def scenario():
ctrl = AdmissionController(1, max_waiting=0)
holder, gate = hold_open(ctrl)
await asyncio.sleep(0)
started = time.monotonic()
answer = await ask(ctrl)
elapsed = time.monotonic() - started
gate.set()
await holder
return answer, elapsed
answer, elapsed = asyncio.run(scenario())
assert answer == "line_full"
assert elapsed < 0.02
def test_wait_limit_turns_requests_away_after_the_limit():
async def scenario():
ctrl = AdmissionController(1, max_wait_s=0.05)
holder, gate = hold_open(ctrl)
await asyncio.sleep(0)
started = time.monotonic()
answer = await ask(ctrl)
elapsed = time.monotonic() - started
gate.set()
await holder
return answer, elapsed
answer, elapsed = asyncio.run(scenario())
assert answer == "waited_too_long"
assert 0.04 <= elapsed < 0.2
async def timeouts_give_back_their_place(controller_class):
ctrl = controller_class(1, max_waiting=2, max_wait_s=0.02)
holder, gate = hold_open(ctrl)
await asyncio.sleep(0)
answers = [await ask(ctrl) for _ in range(4)]
gate.set()
await holder
assert ctrl.waiting == 0, f"line counter stuck at {ctrl.waiting}; answers were {answers}"
def test_timeouts_give_back_their_place():
asyncio.run(timeouts_give_back_their_place(AdmissionController))
def test_the_check_above_catches_a_leaky_controller():
with pytest.raises(AssertionError, match="line counter stuck"):
asyncio.run(timeouts_give_back_their_place(LeakyController))
def test_cancelled_waiter_gives_back_its_place():
async def scenario():
ctrl = AdmissionController(1, max_waiting=2)
holder, gate = hold_open(ctrl)
await asyncio.sleep(0)
waiter = asyncio.create_task(ask(ctrl))
await asyncio.sleep(0.01)
assert ctrl.waiting == 1
waiter.cancel()
with pytest.raises(asyncio.CancelledError):
await waiter
assert ctrl.waiting == 0
gate.set()
await holder
return await ask(ctrl) # the slot itself was not lost either
assert asyncio.run(scenario()) == "served"
def test_slot_is_released_when_the_work_fails():
async def scenario():
ctrl = AdmissionController(1, max_waiting=0)
with pytest.raises(ValueError):
async with ctrl.slot():
raise ValueError("the work failed")
return await ask(ctrl)
assert asyncio.run(scenario()) == "served"
def test_http_503_carries_retry_after_and_health_stays_up():
async def scenario():
ctrl = AdmissionController(1, max_waiting=0, retry_after=3)
app = make_app(ctrl, service_time=0.1)
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
first = asyncio.create_task(client.get("/work"))
await asyncio.sleep(0.02) # the first request now holds the only slot
second = await client.get("/work")
health = await client.get("/health")
return second, health, await first
second, health, first = asyncio.run(scenario())
assert first.status_code == 200
assert second.status_code == 503
assert second.headers["retry-after"] == "3"
assert second.json() == {"error": "overloaded", "reason": "line_full"}
assert health.status_code == 200
python -m pytest -v
test_admission.py::test_idle_service_serves_everything PASSED [ 12%]
test_admission.py::test_full_line_turns_requests_away_immediately PASSED [ 25%]
test_admission.py::test_wait_limit_turns_requests_away_after_the_limit PASSED [ 37%]
test_admission.py::test_timeouts_give_back_their_place PASSED [ 50%]
test_admission.py::test_the_check_above_catches_a_leaky_controller PASSED [ 62%]
test_admission.py::test_cancelled_waiter_gives_back_its_place PASSED [ 75%]
test_admission.py::test_slot_is_released_when_the_work_fails PASSED [ 87%]
test_admission.py::test_http_503_carries_retry_after_and_health_stays_up PASSED [100%]
=== 8 passed in 0.60s ===
Check before moving on: eight tests pass in about a second. The HTTP test uses httpx’s ASGITransport, which calls the FastAPI app directly in the same process, so no server needs to be running.
Step 9: Run the load test that proves it
The last experiment is the one the AWS article recommends: “The ideal load test result is for goodput to plateau when the service is close to being fully utilized, and to remain flat even when more throughput is applied.” This script offers half the capacity, then 1, 1.5, 2, and 4 times the capacity, to the unprotected service and to the protected one, and measures goodput from the third second onward so the first seconds do not flatter the result. It takes about a minute.
# step09_sweep.py
import asyncio
from lab_common import load_capacity, run_policy
DURATION = 6
WARMUP = 2
async def main():
capacity = load_capacity()
max_waiting = round(capacity * 0.5)
print(f"capacity {capacity:.0f}/s; goodput counts requests answered within 1 s, from second {WARMUP} on")
print(f"{'offered':>9}{'no limits':>12}{'protected':>12}")
for multiple in (0.5, 1.0, 1.5, 2.0, 4.0):
rate = multiple * capacity
bare = await run_policy(rate, DURATION)
safe = await run_policy(rate, DURATION, max_waiting=max_waiting, max_wait_s=0.3)
print(f"{multiple:>7.1f}x{bare['steady_goodput']:>12.1f}{safe['steady_goodput']:>12.1f}")
asyncio.run(main())
python step09_sweep.py
capacity 65/s; goodput counts requests answered within 1 s, from second 2 on
offered no limits protected
0.5x 32.2 32.2
1.0x 65.0 64.8
1.5x 0.0 62.2
2.0x 0.0 62.0
4.0x 0.0 63.0
Below capacity the two services are identical, which is what you want: protection costs nothing when there is nothing to protect against. At the capacity measured in Step 1 they are still equal. At one and a half times that capacity the unprotected service has already fallen off the cliff and answers nothing in time, while the protected one holds at about 62 requests per second and stays flat at two and four times the capacity. That flat line is the goal. A sustained overload has a cliff, not a slope; a short burst would have let the unprotected service recover, but not every overload ends quickly. The 62 is the service’s true steady rate: four slots, each needing 64.3 ms, finish 4 ÷ 0.0643 = 62 requests per second. The 65.3 from Step 1 read a little high for the reason given there.
Tune the wait limit with data
How long should the wait limit be? Measure it instead of guessing. This script repeats the two-times-capacity test with wait limits of 0.3, 0.6, and 0.9 seconds and prints the goodput, the 99th-percentile time, and the counts.
# tune_wait.py
import asyncio
from lab_common import load_capacity, run_policy
async def main():
capacity = load_capacity()
max_waiting = round(capacity * 0.5)
print(f"offered {2 * capacity:.0f}/s for 6 s, line capped at {max_waiting}")
print(f"{'max_wait_s':>11}{'goodput/s':>11}{'p99 ms':>8}{'timeout':>9}{'shed':>6}")
for wait in (0.3, 0.6, 0.9):
r = await run_policy(2 * capacity, 6, max_waiting=max_waiting, max_wait_s=wait)
print(f"{wait:>11}{r['steady_goodput']:>11.1f}{r['p99_ms']:>8.0f}{r['timeout']:>9}{r['shed']:>6}")
asyncio.run(main())
python tune_wait.py
offered 131/s for 6 s, line capped at 33
max_wait_s goodput/s p99 ms timeout shed
0.3 62.0 355 0 390
0.6 62.8 612 0 371
0.9 62.2 612 0 375
Goodput did not move: 62.0, 62.8, and 62.2 requests per second. The 99th percentile did, from 355 ms to 612 ms. The line cap of 33 already limits a wait to about half a second at this service rate, so a wait limit of 0.6 or 0.9 seconds behaves like no wait limit at all and only lets requests wait longer. In Step 5 the wait limit earned its place when the service slowed down, because a count of requests does not bound a time. The two limits do different jobs: the cap bounds memory and the normal-case delay, and the wait limit bounds the delay when conditions change. Choose the wait limit by what your clients can tolerate, and check it with a sweep like this one.
The same AWS article gives the reason to run a test like this before an incident does it for you: “if they haven’t load tested their service to the point where it breaks, and far beyond the point where it breaks, they should assume that the service will fail in the least desirable way possible.”
Check before moving on: the protected column stays flat and well above zero as the offered load grows, while the unprotected column drops to zero from one and a half times the capacity onward.
Confirm everything works end to end
Run one last check by hand against the real server. This script starts the service with four slots and no waiting room, makes each request take a full second so the slots stay busy, fires 20 requests at the same moment with curl, and counts the distinct response lines.
# try_it.sh (bash: Linux, macOS, or Git Bash on Windows; run it in the folder that has service.py)
SERVICE_TIME=1 MAX_WAITING=0 python -m uvicorn service:app --port 8000 --log-level warning &
SERVER=$!
sleep 3
echo "20 requests at the same moment; 4 slots, no waiting room, each request takes a full second"
for i in $(seq 20); do
curl -si http://127.0.0.1:8000/work | tr -d '\r' | grep -iE "^(HTTP/|retry-after|\{)" &
done | sort | uniq -c
kill $SERVER
bash try_it.sh
20 requests at the same moment; 4 slots, no waiting room, each request takes a full second
16 {"error":"overloaded","reason":"line_full"}
4 {"ok":true}
4 HTTP/1.1 200 OK
16 HTTP/1.1 503 Service Unavailable
16 retry-after: 1
Exactly four requests were admitted and sixteen were turned away, each with a 503, a retry-after: 1 header, and a body that says line_full. On a slower machine you may see one or two more admitted if a request finishes before the last curl starts. Together with the earlier steps, you have now confirmed the following:
- The capacity you measured is below what you assumed (Step 1).
- Without limits, a sustained overload makes almost every request fail while the service stays busy (Step 3).
- A bounded line and a wait limit keep latency under the client’s patience and eliminate wasted work (Steps 4 and 5).
- The HTTP layer returns a 503 with
Retry-Afterand keeps/healthreachable (Steps 6 and 8). - Polite clients add little extra load, while immediate retries double it (Step 7).
- The counters do not leak, and a test proves it (Step 8).
- Goodput stays flat as offered load grows (Step 9).
Common mistakes and how to spot them
Letting a line grow without a limit
Queues hide in places you did not write: a semaphore, a thread pool’s task list, a connection pool, the operating system’s accept backlog. In Step 3 nothing in the code said “line”, yet the peak was 339 requests. If you cannot say how long each line in your service may grow, assume it is unbounded and find out.
Sizing limits from a healthy day
A line cap that is right when the service is fast is wrong when it is slow (Step 5). Pair a count with a time limit, and pick the time limit from the client’s timeout.
Counters that can leak
Any counter you increment before an await must be decremented in a finally block, because cancellation and timeouts both skip the code after the await (Step 8).
Putting health checks behind the limiter
A busy service that fails its health check can be restarted by the orchestrator at the worst possible moment. Keep /health outside the controller (Step 6).
Judging by the latency of the successes
The unprotected service in Step 3 had a median of 526 ms and failed 82 percent of requests. Always report goodput and failure counts next to percentiles.
Retrying immediately
Immediate retries doubled the load for no real gain (Step 7). Honor Retry-After, add jitter, and cap retries with a budget.
Making rejection expensive
The reject path should do almost nothing: no database call, no heavy logging, no building a large error body. uvicorn’s flag wrote a log line for each of 223 refusals in a five second test. If you log rejections, count them and log a sample.
Where to go next
- Throttle at the client. The SRE chapter Handling Overload describes adaptive throttling: once a client sees many of its recent requests rejected, it starts to fail excess requests locally, and “Requests above the cap fail locally without even reaching the network.”
- Try smarter queues. Switch the line from first-in first-out to last-in first-out, or study CoDel, as the SRE book suggests, so the oldest and least valuable requests are the ones dropped.
- Rank requests by importance. Shed background work before user-facing work. The Handling Overload chapter calls this criticality.
- Watch and tune it. Export the controller’s counters, and, in the SRE book’s words, “Monitor and alert when too many servers enter these modes.” It also advises that you “Design a way to quickly turn off complex graceful degradation or tune parameters if needed”; a feature flag is a natural switch for that.
- Combine it with its neighbors. A bulkhead limits how much of the service one dependency can use, and a circuit breaker stops calling a dependency that keeps failing. Admission control limits what the service as a whole takes on.
Sources
- Google, Site Reliability Engineering, Addressing Cascading Failures (queue management, load shedding, retries).
- Google, Site Reliability Engineering, Handling Overload (client-side throttling).
- AWS, Using load shedding to avoid overload (goodput, cascading retries, load testing).
- IETF, RFC 9110, HTTP Semantics (503 Service Unavailable and Retry-After).
- Python documentation, asyncio queues, asyncio tasks and wait_for, and asyncio exceptions.
- uvicorn, settings and the 0.54.0 source files h11_impl.py and flow_control.py.
- FastAPI, handling errors, and httpx, transports.
- Wikipedia, Little’s law.








No Comment! Be the first one.