"""
Load generator for the typesense-python retry-storm reproducer.
Runs inside the docker network next to nginx (lb:8000, per-node 8001..8003).
python3 loadgen.py seed <docs>
python3 loadgen.py calibrate <port> <concurrency> <seconds>
python3 loadgen.py run <arm> <rate_per_s> <seconds> <out.jsonl>
Arms:
stock typesense 2.0.0 as shipped
patched 2.0.0 with the pause between retries honoured and the node that
actually answered marked healthy (what 0.21.0 did: PR #141 plus the
fix for the "wrong node marked healthy" bug)
noretry 2.0.0 with num_retries=0 (the floor: no client retries at all)
capped 2.0.0 as shipped, but the app lets at most CAP (50) searches into
the client at once; the rest wait in the app, so the client never
has more requests than its 100 connection pool can hold
"run" is open loop: calls start on a fixed schedule whatever the latency, the
way web traffic arrives, so in-flight calls pile up when the cluster slows.
Every app call and every client attempt is written to out.jsonl with epoch
timestamps so analyze.py can line them up with the nginx access log.
"""
import asyncio
import contextlib
import contextvars
import json
import os
import random
import sys
import threading
import time
import typesense
from typesense.async_.api_call import AsyncApiCall, _SERVER_ERRORS
from typesense.exceptions import TypesenseClientError
LB = {"host": "nginx", "port": "8000", "protocol": "http"}
NODES = [{"host": "nginx", "port": p, "protocol": "http"} for p in ("8001", "8002", "8003")]
COLL = "products"
CAP = 50 # "capped" arm: searches allowed inside the client at once
rnd = random.Random(42)
VOCAB = [f"w{i}" for i in range(3000)]
WEIGHTS = [1.0 / (i + 1) for i in range(len(VOCAB))] # Zipf: a few words are everywhere
def words(n):
return " ".join(rnd.choices(VOCAB, WEIGHTS, k=n))
def make_client(arm, timeout=3.0):
cfg = {"api_key": "xyz", "nearest_node": LB, "nodes": NODES, "connection_timeout_seconds": timeout}
if arm == "noretry":
cfg["num_retries"] = 0
client = typesense.AsyncClient(cfg)
if arm == "patched":
client.api_call._execute_request = _patched_execute.__get__(client.api_call, AsyncApiCall)
return client
async def _patched_execute(self, method, endpoint, entity_type, as_json=True, last_exception=None, num_retries=0, **kwargs):
if num_retries > self.config.num_retries:
if last_exception:
raise last_exception
raise TypesenseClientError("All nodes are unhealthy")
if num_retries > 0:
await asyncio.sleep(self.config.retry_interval_seconds)
node, url, request_kwargs = self._prepare_request_params(endpoint, **kwargs)
try:
resp = await self.request_handler.make_request(
method=method, url=url, as_json=as_json, entity_type=entity_type, client=self._client, **request_kwargs)
self.node_manager.set_node_health(node, is_healthy=True)
return resp
except _SERVER_ERRORS as err:
self.node_manager.set_node_health(node, is_healthy=False)
return await _patched_execute(self, method, endpoint, entity_type, as_json,
last_exception=err, num_retries=num_retries + 1, **kwargs)
async def seed(n):
client = make_client("stock", timeout=120)
try:
await client.collections[COLL].delete()
except Exception:
pass
await client.collections.create({
"name": COLL,
"fields": [
{"name": "title", "type": "string"},
{"name": "description", "type": "string"},
{"name": "category", "type": "string", "facet": True},
{"name": "brand", "type": "string", "facet": True},
{"name": "price", "type": "float"},
{"name": "popularity", "type": "int32"},
],
"default_sorting_field": "popularity",
})
batch, done = [], 0
for i in range(n):
batch.append({
"id": str(i), "title": words(rnd.randint(6, 12)), "description": words(rnd.randint(30, 60)),
"category": f"c{rnd.randint(0, 59)}", "brand": f"b{rnd.randint(0, 799)}",
"price": round(rnd.uniform(1, 500), 2), "popularity": rnd.randint(0, 100000),
})
if len(batch) == 5000 or i == n - 1:
res = await client.collections[COLL].documents.import_(batch, {"action": "create"})
bad = [r for r in res if not r.get("success")]
if bad:
print("import errors:", bad[:2], flush=True)
done += len(batch)
batch = []
if done % 50000 == 0 or done == n:
print(f"seeded {done}/{n}", flush=True)
await client.api_call.aclose()
def query():
return {
"q": words(rnd.choice((1, 1, 2))),
"query_by": "title,description",
"facet_by": "category,brand",
"max_facet_values": 20,
"per_page": 20,
}
async def calibrate(port, conc, seconds):
client = typesense.AsyncClient({"api_key": "xyz", "nodes": [{"host": "nginx", "port": port, "protocol": "http"}],
"connection_timeout_seconds": 60, "num_retries": 0})
lat, stop = [], time.time() + seconds
async def worker():
while time.time() < stop:
t = time.time()
await client.collections[COLL].documents.search(query())
lat.append(time.time() - t)
await asyncio.gather(*(worker() for _ in range(conc)))
await client.api_call.aclose()
lat.sort()
print(json.dumps({"port": port, "conc": conc, "qps": round(len(lat) / seconds, 1),
"p50_ms": round(lat[len(lat) // 2] * 1000), "p95_ms": round(lat[int(len(lat) * .95)] * 1000)}), flush=True)
CUR = contextvars.ContextVar("cur")
async def run(arm, rate, seconds, out_path):
client = make_client(arm)
rh = client.api_call.request_handler
orig = rh.make_request
def counted(**kw):
coro = orig(**kw)
port = kw["url"].split(":")[2].split("/")[0]
async def attempt():
rec, t = CUR.get(None), time.time()
try:
res = await coro
if rec is not None:
rec["att"].append([port, "ok", round(t, 3), round(time.time(), 3)])
return res
except BaseException as e:
if rec is not None:
rec["att"].append([port, type(e).__name__, round(t, 3), round(time.time(), 3)])
raise
return attempt()
rh.make_request = counted
out = open(out_path, "w")
nm = client.api_call.node_manager
t0 = time.time()
out.write(json.dumps({"k": "start", "arm": arm, "rate": rate, "t": t0}) + "\n")
inflight = set()
gate = asyncio.Semaphore(CAP) if arm == "capped" else contextlib.nullcontext()
async def one(i):
rec = {"k": "call", "i": i, "t": round(time.time(), 3), "att": []}
CUR.set(rec)
try:
async with gate: # call time includes any wait here
await client.collections[COLL].documents.search(query())
rec["out"] = "ok"
except Exception as e:
rec["out"] = type(e).__name__
rec["end"] = round(time.time(), 3)
out.write(json.dumps(rec) + "\n")
pool = client.api_call._client._transport._pool
async def health():
while True:
# Pool census: busy connections with no server traffic means leaked
# slots (httpcore #1093); idle or churning ones with a long queue
# means the event loop is just too busy to hand them out.
conns = list(pool.connections)
out.write(json.dumps({"k": "health", "t": round(time.time(), 3), "inflight": len(inflight),
"lb": client.config.nearest_node.healthy,
"nodes": [n.healthy for n in nm.nodes],
"conns": len(conns),
"busy": sum(1 for c in conns if not c.is_idle() and not c.is_closed()),
"idle": sum(1 for c in conns if c.is_idle()),
"queued": sum(1 for r in list(pool._requests) if r.is_queued())}) + "\n")
await asyncio.sleep(1)
htask = asyncio.create_task(health())
total = int(rate * seconds)
stop_path = out_path + ".stop"
def watchdog():
# A client stuck in its own pool can starve the event loop for many
# minutes. Once the run is cut, give in-flight calls 60s, then record
# how many never finished and exit from this thread.
while not os.path.exists(stop_path):
time.sleep(1)
time.sleep(60)
with open(out_path, "a") as f:
f.write(json.dumps({"k": "end", "forced": True, "t": time.time(), "stuck_inflight": len(inflight),
"lb": client.config.nearest_node.healthy,
"nodes": [n.healthy for n in nm.nodes]}) + "\n")
os._exit(0)
threading.Thread(target=watchdog, daemon=True).start()
for i in range(total):
if i % 20 == 0 and os.path.exists(stop_path):
break
delay = t0 + i / rate - time.time()
if delay > 0:
await asyncio.sleep(delay)
task = asyncio.create_task(one(i))
inflight.add(task)
task.add_done_callback(inflight.discard)
if inflight:
await asyncio.wait(inflight, timeout=60)
htask.cancel()
out.write(json.dumps({"k": "end", "t": time.time()}) + "\n")
out.close()
await client.api_call.aclose()
print(f"run {arm} done: {total} calls", flush=True)
if __name__ == "__main__":
mode = sys.argv[1]
if mode == "seed":
asyncio.run(seed(int(sys.argv[2])))
elif mode == "calibrate":
asyncio.run(calibrate(sys.argv[2], int(sys.argv[3]), int(sys.argv[4])))
elif mode == "run":
asyncio.run(run(sys.argv[2], float(sys.argv[3]), float(sys.argv[4]), sys.argv[5]))
Under load, if one node restarts while another is slow, the AsyncClient locks up for good. Every call fails with PoolTimeout, and the cluster sits idle until the process restarts.
It happens as a chain: timeouts drop nodes from rotation, traffic piles onto the remaining nodes, retries fill the 100-connection pool (#143), and a known httpcore bug (encode/httpcore#1093) leaks pool slots so the pool never frees up.
Fixing #143 removes the step that turns pool pressure into retries. Beyond that, the issue suggests capping requests in flight, giving the pool its own timeout, and replacing the httpx client when the pool is stuck.
Description
With the 2.0.0
AsyncClientand default settings, restarting one node of a three node cluster while another node is slow can leave the client permanently broken. Every call fails withhttpx.PoolTimeout, the httpx pool reports all 100 connections busy, the cluster receives no traffic at all, and only a new client (or a process restart) recovers.We first saw this in production after a node rotation. Below is a local reproduction with real Typesense nodes.
Setup
nearest_node) and one endpoint per node (nodes).typesense.AsyncClientwith default settings, sending 26 searches per second on a fixed schedule. That is about 65% of the cluster's capacity. As with web traffic, new calls keep arriving when the cluster slows down.slowtrigger also cuts node 1 to 0.3 CPU for 20 seconds at the same moment, the way a leader slows down while it sends a snapshot to a rejoining node. With a plainrestart, whether the client trips depends on luck: it did when a few searches on a surviving node crossed 3 seconds while node 2 was away.hangfreezes node 2 for 10 seconds (it accepts connections but never answers), then restarts it.Clients compared:
stock: 2.0.0 as released.patched: 2.0.0 with a pause between retries (the change in fix: wait retry_interval_seconds between retries (#140) #141) and the node that answered marked healthy (After a successful request the client marks the next node healthy, not the one that answered, and round-robin skips nodes #144).noretry: 2.0.0 withnum_retries=0.capped: 2.0.0 as released, but the app lets at most 50 searches into the client at once and the rest wait in the app.Results
The run is cut 2.5 minutes after node 2 is healthy again, and no new calls start after that. "Calls stuck at the end" counts calls still unfinished 60 seconds after the cut. "Last 90s" is the window from 60 seconds after node 2 was healthy until the cut, when about 2,350 calls were due.
A plain
restarttrips the client only sometimes, so therestartrows that did not trip say nothing about those clients. Theslowrows are the like-for-like comparison.Timeline of the stock client with the slow trigger (10 second windows)
calls= app calls started in the window that had finished by the end of the run (so it reads 0 once calls stop finishing, even though the app keeps starting 26 per second),att= attempts the client made,srv= requests the nodes received (from the nginx log),x= srv / calls,n1..n3= requests per node,direct= requests that bypassed the load balancer,p50/p95= call latency in seconds,infl= calls in flight,lb= load balancer marked unhealthy by the client. Once the client is stuck, its own once-a-second sample often does not arrive at all because the event loop is starved, soinflreads 0 andlbis blank in those windows.What happens
connection_timeout_seconds(3s) is also the read timeout. One search through the load balancer that takes longer than that marks it unhealthy forhealthcheck_interval_seconds(60s), and the client sends everything straight to the nodes.PoolTimeout, which the client counts as a node failure and retries (PoolTimeout and other client-side httpx errors mark a healthy node unhealthy and fail over #143).The
hangrow is the counterexample. With node 2 frozen on its own, the two other nodes stayed fast, traffic split between them, and the client went back to the load balancer after 60 seconds. The lockup needs a surviving node to be slow while another node is away, and a node rotation does exactly that to the leader.What helps and what doesn't
patched) does not prevent it. That client locked up exactly like stock.num_retries=0avoids the lockup, but 24% of calls kept failing for 2.5 minutes after the cluster was healthy again. In that window the nodes received normal traffic but stayed slow: median processing time was 0.8 seconds instead of 0.03. We did not pin down why.capped) keeps the pool alive. Nothing failed withPoolTimeout, and the nodes kept receiving about 22 searches per second until the end. But it does not get back to normal while traffic keeps arriving. The client kept the load balancer marked down and about one finished call in five still failed withReadTimeout. The cluster served fewer searches than arrived, so the queue in the app grew past 2,000 calls, with a median wait of about 100 seconds. The backlog only drained once new calls stopped: about 1,000 calls finished in the 60 seconds after the cut.Suggestions
Fixing #143 removes the step that turns pool pressure into retries. Beyond that, the client could:
cappedrun shows, does not bring the client back to normal while traffic continues);connection_timeout_seconds, and let users passhttpx.Limits;PoolTimeoutwith no responses coming back) and replace the httpx client;nearest_nodeout of rotation for the fullhealthcheck_interval_seconds.Reproducer
Requires Docker and bash.
./reproduce.sh allbrings the cluster up, seeds 300k documents, runs all four clients with theslowtrigger, and tears everything down. It takes about half an hour. To run a single client:./reproduce.sh up && ./reproduce.sh seed && ./reproduce.sh run stock slow 26, then./reproduce.sh down. Each run prints a 10 second timeline and a per-phase summary. 26 searches per second was about 65% of capacity on the machine we used;./reproduce.sh calibratemeasures yours.reproduce.shnginx.confloadgen.pyanalyze.pyEnvironment
python:3.12-slim).