Skip to content

AsyncClient locks up for good after a node restart under load: every call fails with PoolTimeout while the cluster sits idle #146

Description

@alangmartini

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 AsyncClient and default settings, restarting one node of a three node cluster while another node is slow can leave the client permanently broken. Every call fails with httpx.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

  • A three node Typesense 31.0.rc14 cluster in Docker, one CPU per node, 300k documents.
  • nginx in front, laid out like a managed deployment: one load balanced endpoint (nearest_node) and one endpoint per node (nodes).
  • One typesense.AsyncClient with 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.
  • At t=60s node 2 restarts. The slow trigger 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 plain restart, 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. hang freezes node 2 for 10 seconds (it accepts connections but never answers), then restarts it.

Clients compared:

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.

Client Trigger Peak server requests per app call (10s window) Calls stuck at the end Server requests in the last 90s Calls failed in the last 90s
stock restart (run 1) 1.39 2,944 0 all
stock restart (run 2) 1.45 1,751 9 all
patched restart 1.00 (did not trip) 0 2,353 0
noretry restart 1.00 (did not trip) 0 2,354 0
stock hang 1.12 0 2,352 0
stock slow 1.82 2,007 0 all
patched slow 1.36 2,341 0 all
noretry slow 1.00 0 2,353 567 (24%)
capped slow 1.44 1,151 2,048 175 of the 917 that finished; the rest still waiting in the app

A plain restart trips the client only sometimes, so the restart rows that did not trip say nothing about those clients. The slow rows 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, so infl reads 0 and lb is blank in those windows.

== out/stock-slow: arm=stock, 26.0/s open loop, 2810 app calls ==
events (s from start): {'restart_begin': 60.3, 'container_up': 62.8, 'node2_healthy': 73.0, 'node1_restored': 83.2, 'cut': 223.5}
   t calls fail  att  srv     x   n1   n2   n3 direct   p50   p95 infl lb
   0   260    0  260  260  1.00   87   87   86      0   0.0   0.2    3 
  10   260    0  260  260  1.00   86   87   87      0   0.0   0.2    3 
  20   260    0  260  260  1.00   87   86   87      0   0.0   0.2    2 
  30   260    0  260  260  1.00   87   87   86      0   0.0   0.2    2 
  40   260    0  260  260  1.00   86   87   87      0   0.0   0.2    4 
  50   260    0  260  260  1.00   87   86   87      0   0.0   0.2    4 
  60   260   24  356  356  1.37  122   51  183    246   1.4   9.4   86 part <restart_begin> <container_up>
  70   260   48  474  474  1.82   86  273  115    474   2.9  16.2  107 DOWN <node2_healthy>
  80   260  161  530  399  1.53   67  232  100    399  18.6  33.2  216 DOWN <node1_restored>
  90   234  222  635  165  0.71   46    0  119    165  33.0  46.1  329 DOWN
 100   236  236  611    0  0.00    0    0    0      0  75.6  80.8  467 DOWN
 110     0    0  465    0  0.00    0    0    0      0   0.0   0.0  548 DOWN
 120     0    0  216    0  0.00    0    0    0      0   0.0   0.0  692 DOWN
 130     0    0  158    0  0.00    0    0    0      0   0.0   0.0  718 DOWN
 140     0    0   57    0  0.00    0    0    0      0   0.0   0.0  932 DOWN
 150     0    0  106    0  0.00    0    0    0      0   0.0   0.0    0 
 160     0    0    0    0  0.00    0    0    0      0   0.0   0.0    0 
 170     0    0    0    0  0.00    0    0    0      0   0.0   0.0    0 
 180     0    0    0    0  0.00    0    0    0      0   0.0   0.0 1298 DOWN
 190     0    0    0    0  0.00    0    0    0      0   0.0   0.0    0 
 200     0    0    0    0  0.00    0    0    0      0   0.0   0.0    0 
 210     0    0    0    0  0.00    0    0    0      0   0.0   0.0    0 
 220     0    0    0    0  0.00    0    0    0      0   0.0   0.0    0  <cut>

What happens

  1. One slow search drops the load balancer for 60 seconds. connection_timeout_seconds (3s) is also the read timeout. One search through the load balancer that takes longer than that marks it unhealthy for healthcheck_interval_seconds (60s), and the client sends everything straight to the nodes.
  2. Traffic herds onto whichever nodes are still marked healthy. Each node that answers slowly once is also dropped for 60 seconds, so its share lands on the nodes left. In stock run 2, node 2 was the only node still marked healthy and received all 260 requests per 10 seconds, against a capacity of about 18 per second. It timed out, was dropped, and the pile moved on. After a successful request the client marks the next node healthy, not the one that answered, and round-robin skips nodes #144 makes the spread worse.
  3. Pool exhaustion turns into more retries. Timeouts plus immediate retries (retry_interval_seconds is stored but never used, so retries fire with no delay (regression from 0.21.0) #140) push the number of requests in flight past httpx's pool of 100 connections. Waiting longer than 3 seconds for a connection raises 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).
  4. The pool never recovers. At the end of stock run 2, all 100 connections reported busy and 1,109 requests were queued, while the cluster received 9 requests in 90 seconds. This matches Cancelling requests under pool contention permanently leaks connection slots (AsyncConnectionPool → PoolTimeout) encode/httpcore#1093: requests cancelled while the pool is under contention leak their connection slot, and httpx implements timeouts as cancellations. That issue's own script reproduces on httpcore 1.0.9 in the same image. The client process also sat at 100% CPU while all three nodes were idle; we did not profile it. No run recovered, and the longest one was watched for 10 minutes.

The hang row 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

  • The pause between retries plus correct node health (patched) does not prevent it. That client locked up exactly like stock.
  • num_retries=0 avoids 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.
  • Letting at most 50 searches into the client (capped) keeps the pool alive. Nothing failed with PoolTimeout, 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 with ReadTimeout. 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.
  • In that late window, the nodes' own processing time (nginx upstream time) had a median of about 3 seconds, the point where the client gives up, against 0.04 seconds in steady state. At the same time they were receiving fewer searches than in steady state: about 7.6 per second per node, against 8.7. That would fit searches the client abandons at its timeout continuing to run on the server, but we did not verify it.

Suggestions

Fixing #143 removes the step that turns pool pressure into retries. Beyond that, the client could:

  • offer an optional limit on requests in flight, below the pool size (on its own that keeps the pool alive but, as the capped run shows, does not bring the client back to normal while traffic continues);
  • set the pool wait timeout separately from connection_timeout_seconds, and let users pass httpx.Limits;
  • notice a wedged pool (repeated PoolTimeout with no responses coming back) and replace the httpx client;
  • reconsider whether a single timeout should take nearest_node out of rotation for the full healthcheck_interval_seconds.

Reproducer

Requires Docker and bash. ./reproduce.sh all brings the cluster up, seeds 300k documents, runs all four clients with the slow trigger, 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 calibrate measures yours.

reproduce.sh
#!/bin/bash
# Issue: typesense-python 2.0.0 AsyncClient locks up after a node restart under load
# Typesense Version: 31.0.rc14, typesense-python 2.0.0
# Description:
#   A real 3 node HA cluster behind nginx, wired like Typesense Cloud: the
#   client's nearest_node is the load balancer (port 8000) and its nodes are
#   per-node endpoints (8001..8003). A typesense-python AsyncClient sends
#   searches on a fixed schedule (open loop, like web traffic) while node 2 is
#   restarted. nginx logs every request each node receives, so the run shows
#   how many server requests one app call turned into, second by second, and
#   whether that settles once node 2 is back.
#
#   ./reproduce.sh all                 # up, seed, all four arms with the slow trigger, down
#   ./reproduce.sh up | seed | calibrate | down
#   ./reproduce.sh run <arm> <restart|hang|slow|wipe> <rate>   # arm: stock|patched|noretry|capped

set -u
export MSYS_NO_PATHCONV=1

VERSION=31.0.rc14
NET=ts-storm
SUBNET=10.231.0.0/24
NODES_CONF="10.231.0.11:8107:8108,10.231.0.12:8107:8108,10.231.0.13:8107:8108"
CPUS=${CPUS:-1}            # per node: small, known capacity so overload is reachable
DOCS=${DOCS:-300000}
STEADY=${STEADY:-60}       # seconds of normal traffic before node 2 goes away
AFTER=${AFTER:-150}        # seconds of traffic after node 2 is healthy again

HERE="$(cd "$(dirname "$0")" && pwd)"
case "$(uname -s)" in MINGW*|MSYS*|CYGWIN*) HERE_D="$(cd "$HERE" && pwd -W)";; *) HERE_D="$HERE";; esac
OUT="$HERE/out"; OUT_D="$HERE_D/out"

lbexec() { docker exec storm-lb "$@"; }
now() { docker exec storm-gen python3 -c 'import time; print(time.time())'; }  # same clock as loadgen and nginx
node_health() { lbexec wget -qO- -T 2 "http://10.231.0.1$1:8108/health" 2>/dev/null; }

start_node() {
  local i=$1
  docker volume create "storm-n$i" >/dev/null
  docker run --rm -v "storm-n$i:/data" nginx:1.27-alpine sh -c "printf '%s' '$NODES_CONF' > /data/nodes"
  docker run -d --name "ts$i" --network "$NET" --ip "10.231.0.1$i" --cpus "$CPUS" \
    -v "storm-n$i:/data" "typesense/typesense:$VERSION" \
    --data-dir /data --api-key=xyz --api-port 8108 --peering-port 8107 \
    --peering-address "10.231.0.1$i" --nodes /data/nodes >/dev/null
}

wait_healthy() {  # node index, timeout seconds
  local i=$1 t=${2:-300} n=0
  until node_health "$i" | grep -q '"ok":true'; do
    n=$((n + 1)); [ $n -ge $t ] && { echo "node $i not healthy after ${t}s"; return 1; }
    sleep 1
  done
}

up() {
  down >/dev/null 2>&1
  mkdir -p "$OUT/logs"
  docker network create --subnet="$SUBNET" "$NET" >/dev/null
  for i in 1 2 3; do start_node "$i"; done
  docker run -d --name storm-lb --network "$NET" --ip 10.231.0.10 --network-alias nginx \
    -v "$HERE_D/nginx.conf:/etc/nginx/nginx.conf:ro" -v "$OUT_D/logs:/logs" nginx:1.27-alpine >/dev/null
  docker run -d --name storm-gen --network "$NET" -v "$HERE_D:/w" -w /w python:3.12-slim sleep infinity >/dev/null
  docker exec storm-gen pip install -q --disable-pip-version-check --root-user-action=ignore typesense==2.0.0
  for i in 1 2 3; do wait_healthy "$i" 120 || return 1; done
  echo "cluster up"
}

down() {
  docker rm -f ts1 ts2 ts3 storm-lb storm-gen >/dev/null 2>&1
  docker volume rm storm-n1 storm-n2 storm-n3 >/dev/null 2>&1
  docker network rm "$NET" >/dev/null 2>&1
  return 0
}

seed() { docker exec storm-gen python3 -u loadgen.py seed "$DOCS"; }

calibrate() {
  for c in 2 4 8; do docker exec storm-gen python3 -u loadgen.py calibrate 8001 "$c" 20; done
  for c in 6 12 24; do docker exec storm-gen python3 -u loadgen.py calibrate 8000 "$c" 20; done
}

run() {  # arm variant rate
  local arm=$1 variant=$2 rate=$3 tag="$1-$2"
  local ev="$OUT/$tag.events"; : > "$ev"
  for i in 1 2 3; do wait_healthy "$i" 600 || return 1; done
  sleep 20  # let LB fail_timeouts and caches settle
  lbexec sh -c ': > /logs/access.log'
  local dur=$((STEADY + 360))   # generous; the run is cut once AFTER has elapsed
  docker exec -d storm-gen sh -c "python3 -u loadgen.py run $arm $rate $dur out/$tag.jsonl > out/$tag.gen.log 2>&1"
  sleep "$STEADY"
  echo "restart_begin $(now)" >> "$ev"
  if [ "$variant" = wipe ]; then
    docker rm -f ts2 >/dev/null; docker volume rm storm-n2 >/dev/null; start_node 2
  elif [ "$variant" = hang ]; then
    # Node 2 stops answering but keeps accepting connections (frozen, like a
    # wedged node or one stuck in shutdown), then restarts. Deterministic,
    # unlike a plain restart, which only sometimes leaves requests hanging.
    docker pause ts2 >/dev/null; sleep "${HANG:-10}"; docker unpause ts2 >/dev/null
    echo "unfrozen $(now)" >> "$ev"
    docker restart ts2 >/dev/null
  elif [ "$variant" = slow ]; then
    # Node 2 restarts and, at the same moment, node 1 slows down for 20s (CPU
    # cut to 0.3), the way a leader slows while it streams a snapshot to the
    # node rejoining. Both are back to normal ~20s later; the question is
    # whether the client is.
    docker update --cpus 0.3 ts1 >/dev/null
    docker restart ts2 >/dev/null
    ( sleep 20; docker update --cpus "$CPUS" ts1 >/dev/null; echo "node1_restored $(now)" >> "$ev" ) &
  else
    docker restart ts2 >/dev/null
  fi
  echo "container_up $(now)" >> "$ev"
  wait_healthy 2 900
  echo "node2_healthy $(now)" >> "$ev"
  sleep "$AFTER"
  echo "cut $(now)" >> "$ev"
  touch "$OUT/$tag.jsonl.stop"
  until grep -q '"k": "end"' "$OUT/$tag.jsonl" 2>/dev/null; do sleep 2; done
  rm -f "$OUT/$tag.jsonl.stop"; sleep 2
  cp "$OUT/logs/access.log" "$OUT/$tag.access.log"
  docker exec storm-gen python3 analyze.py "out/$tag"
}

case "${1:-all}" in
  up) up ;; down) down ;; seed) seed ;; calibrate) calibrate ;;
  run) run "$2" "$3" "$4" ;;
  all)
    # 26/s was ~65% of the 3 node capacity on the machine this was built on;
    # run calibrate and adjust RATE if yours differs.
    RATE=${RATE:-26}
    up && seed && for arm in stock patched noretry capped; do run "$arm" slow "$RATE"; done
    down ;;
esac
nginx.conf
worker_processes 2;
events { worker_connections 8192; }

http {
    # One line per request. $upstream_addr lists every backend nginx tried, so
    # a request the LB re-sent to a second node counts as two server hits.
    log_format hits '$msec $server_port $request_method $uri $status '
                    '"$upstream_addr" "$upstream_status" $request_time "$upstream_response_time"';
    access_log /logs/access.log hits buffer=64k flush=1s;

    proxy_http_version 1.1;
    proxy_set_header Connection "";
    proxy_read_timeout 120s;
    proxy_connect_timeout 300ms;  # a stopped container has no IP: fail fast like a refused port
    client_max_body_size 200m;

    # Load-balanced endpoint, the client's nearest_node. Behaves like a decent
    # managed LB: a backend that refuses connections, times out or answers 503
    # is skipped for 10s and the request is re-sent to another node, so the
    # client rarely sees a node-level 503 through it. Anything that still marks
    # the LB unhealthy in the client is then the client's own timeout.
    upstream ts_all {
        server 10.231.0.11:8108 max_fails=1 fail_timeout=10s;
        server 10.231.0.12:8108 max_fails=1 fail_timeout=10s;
        server 10.231.0.13:8108 max_fails=1 fail_timeout=10s;
        keepalive 128;
    }
    server {
        listen 8000;
        location / {
            proxy_pass http://ts_all;
            proxy_next_upstream error timeout http_502 http_503 http_504;
        }
    }

    # Per-node hostnames, the client's nodes list. A down node either refuses
    # the connection or answers 503 "Not Ready", and the client treats both as
    # a node failure. nginx would turn a refused connection into a 502,
    # which the client does NOT retry, so map it to 503 here.
    server {
        listen 8001;
        location / { proxy_pass http://10.231.0.11:8108; }
        error_page 502 504 =503 /__down;
        location = /__down { internal; default_type application/json; return 503 '{"message":"Not Ready or Lagging"}'; }
    }
    server {
        listen 8002;
        location / { proxy_pass http://10.231.0.12:8108; }
        error_page 502 504 =503 /__down;
        location = /__down { internal; default_type application/json; return 503 '{"message":"Not Ready or Lagging"}'; }
    }
    server {
        listen 8003;
        location / { proxy_pass http://10.231.0.13:8108; }
        error_page 502 504 =503 /__down;
        location = /__down { internal; default_type application/json; return 503 '{"message":"Not Ready or Lagging"}'; }
    }
}
loadgen.py
"""
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]))
analyze.py
"""
Line up one run's app calls (out/<tag>.jsonl), the nginx access log
(out/<tag>.access.log) and the restart events (out/<tag>.events).

  python3 analyze.py out/<tag>

Server executions are counted from nginx: every backend address nginx sent a
request to, except connection failures (upstream status 502/504, which nginx
makes up). A request the client abandoned at its 3s timeout still counts,
because Typesense keeps running a search after the caller hangs up.
"""

import collections
import json
import shlex
import sys

tag = sys.argv[1]
BUCKET = 10
NODE = {"10.231.0.11:8108": "n1", "10.231.0.12:8108": "n2", "10.231.0.13:8108": "n3"}

calls, health, t0 = [], [], None
for line in open(f"{tag}.jsonl"):
    r = json.loads(line)
    if r["k"] == "start":
        t0, arm, rate = r["t"], r["arm"], r["rate"]
    elif r["k"] == "call":
        calls.append(r)
    elif r["k"] == "health":
        health.append(r)

events = {}
for line in open(f"{tag}.events"):
    name, t = line.split()
    events[name] = float(t) - t0

hits = []  # (start_rel, node, via)
for line in open(f"{tag}.access.log"):
    parts = shlex.split(line)
    msec, port, rt = float(parts[0]), parts[1], float(parts[7])
    if "/search" not in parts[3]:
        continue
    addrs = [a.strip() for a in parts[5].split(",")]
    stats = [s.strip() for s in parts[6].split(",")]
    for a, s in zip(addrs, stats + ["-"] * len(addrs)):
        if a in NODE and s not in ("502", "504"):
            hits.append((msec - rt - t0, NODE[a], "lb" if port == "8000" else "direct"))

end = events.get("cut", max(c["t"] for c in calls) - t0)
nb = int(end // BUCKET) + 1
B = lambda t: int(t // BUCKET)

started = collections.Counter(B(c["t"] - t0) for c in calls)
failed = collections.Counter(B(c["t"] - t0) for c in calls if c["out"] != "ok")
attempts = collections.Counter(B(a[2] - t0) for c in calls for a in c["att"])
lat = collections.defaultdict(list)
for c in calls:
    lat[B(c["t"] - t0)].append(c["end"] - c["t"])
srv = collections.defaultdict(collections.Counter)
for t, n, via in hits:
    if t >= 0:
        srv[B(t)][n] += 1
        srv[B(t)][via] += 1
lb_down = collections.defaultdict(list)
inflight = collections.defaultdict(int)
for h in health:
    b = B(h["t"] - t0)
    lb_down[b].append(not h["lb"])
    inflight[b] = max(inflight[b], h["inflight"])

print(f"== {tag}: arm={arm}, {rate}/s open loop, {len(calls)} app calls ==")
print("events (s from start):", {k: round(v, 1) for k, v in events.items()})
print(f"{'t':>4} {'calls':>5} {'fail':>4} {'att':>4} {'srv':>4} {'x':>5} {'n1':>4} {'n2':>4} {'n3':>4} {'direct':>6} {'p50':>5} {'p95':>5} {'infl':>4} lb")
for b in range(nb):
    s = srv[b]
    tot = s["n1"] + s["n2"] + s["n3"]
    l = sorted(lat[b]) or [0]
    mark = ""
    for k, v in events.items():
        if B(v) == b:
            mark += f" <{k}>"
    lbd = lb_down[b]
    lbs = "DOWN" if lbd and all(lbd) else ("part" if any(lbd) else "")
    x = tot / started[b] if started[b] else 0
    print(f"{b * BUCKET:>4} {started[b]:>5} {failed[b]:>4} {attempts[b]:>4} {tot:>4} {x:>5.2f} {s['n1']:>4} {s['n2']:>4} {s['n3']:>4} "
          f"{s['direct']:>6} {l[len(l) // 2]:>5.1f} {l[int(len(l) * .95)]:>5.1f} {inflight[b]:>4} {lbs}{mark}")

def phase(a, b):
    cs = [c for c in calls if a <= c["t"] - t0 < b]
    hs = [h for h in hits if a <= h[0] < b]
    if not cs:
        return None
    outs = collections.Counter(c["out"] for c in cs)
    l = sorted(c["end"] - c["t"] for c in cs)
    return {"secs": round(b - a), "app_calls": len(cs), "server_execs": len(hs),
            "amplification": round(len(hs) / len(cs), 2),
            "client_attempts": sum(len(c["att"]) for c in cs),
            "failed": {k: v for k, v in outs.items() if k != "ok"},
            "p50_s": round(l[len(l) // 2], 2), "p95_s": round(l[int(len(l) * .95)], 2)}

rb, nh = events.get("restart_begin", end), events.get("node2_healthy", end)
print("\nphases:")
for name, a, b in (("steady", 10, rb), ("node2 away", rb, nh), ("first 60s back", nh, nh + 60), ("after that", nh + 60, end)):
    print(f"  {name:<15}", json.dumps(phase(a, b)))

Environment

  • typesense-python 2.0.0, httpx 0.28.1, httpcore 1.0.9, Python 3.12 (python:3.12-slim).
  • Typesense 31.0.rc14, nginx 1.27-alpine, Docker with Linux containers.

No activity

Activity on this issue will appear here.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions