← Load balancing: traffic and resilience
09 / 12 · 60 MIN

Concurrency, selection and limits

Compare algorithms with in-flight requests and observe connection limits, backup use and resumption after capacity release.

Run a bounded experiment

Save the code below as run.py. Set DR_NGINX_BIN to a local NGINX executable and run python3 run.py --output evidence.json. Recorded execution uses NGINX 1.30.5 and Python 3.13.1 on macOS. It creates an NGINX process with one worker and two fictional HTTP backends on IPv4 loopback. It explicitly disables upstream keepalive and configures no shared zone. Held requests wait for script events under a bounded deadline. It performs no financial operations, external traffic or throughput benchmark. Processes, sockets and the temporary directory are closed afterwards.

Compare selection with a held request

In the least_conn group, the first request remains on A. While it stays pending, six short requests go to B. In the equal-weight round-robin group with a similar hold, the six requests split three to A and three to B. Predict these outcomes before execution. The comparison demonstrates the selection signal without measuring CPU or business cost. Artificial waiting on A is not evidence of processor saturation. To choose an algorithm for an actual portal, measure latency, errors, occupancy and resources using a representative workload.

Occupy both limits

The cap group defines max_conns=1 per backend. The script holds one request on each target and sends a third. That request receives 502 without a new application event, and the proxy log indicates no eligible upstream. No queue is configured. After releasing the held requests, both complete with 200 and a new probe also passes. The new response does not change the 502 already delivered to the earlier client. In a fictional incident, correlate occupancy, attempts and outcomes before attributing a pre-forwarding error to the backend.

Observe backup with a healthy primary

In the reserve group, A has max_conns=1 and B is marked backup. A held request occupies A; the next receives 200 from B while the first remains pending. When released by the script, A completes normally. Backup use in this case does not establish process or application failure. The primary lacked eligible connection capacity. During a fictional maintenance window, use this distinction in reporting: health, eligibility and occupancy need their own signals. Also establish backup capacity, because changing destination creates no additional resources.

Interpret limits with workers and keepalive

This experiment uses only one worker, so it does not establish a global limit in a multi-worker deployment. Without a shared zone, max_conns applies per worker. Idle connections retained by keepalive should also not be confused with active connections or executing operations. Executed configuration removes that variable by disabling upstream keepalive. Before applying an actual value, identify shared state, worker count, reuse pattern and the resource to protect. A limit of one was chosen to make the case observable rather than to recommend production sizing.

Prepare the capacity decision

In a fictional project, A will be removed during batch processing and B will become the only available target. The team needs to demonstrate headroom and rejection behavior under expected demand, with stop and recovery criteria. Doubling max_conns can admit more work without increasing processing, worsening queues or dependencies. The technical manager should connect the change to metrics, owners and remaining capacity. Use this lesson’s cases to prepare go/no-go questions, then obtain representative evidence. If that evidence is absent, record the gap; six local requests do not replace service sizing.

"""Original DR NGINX concurrency and passive-state fixture. IPv4 loopback only."""
import argparse
import collections
import concurrent.futures
import datetime
import hashlib
import http.client
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
import json
import os
from pathlib import Path
import platform
import signal
import socket
import subprocess
import tempfile
import threading
import time

parser = argparse.ArgumentParser
parser.add_argument('--output', default='evidence.json')
args = parser.parse_args
nginx = os.environ.get('DR_NGINX_BIN')
if not nginx or not Path(nginx).is_file:
 parser.error('DR_NGINX_BIN must identify a local NGINX executable')
nginx = str(Path(nginx).resolve)
events, checks, details, commands = [], {}, {}, []
lock = threading.Lock
held = {key: threading.Event for key in ['lc', 'rr', 'cap1', 'cap2', 'backup']}
release = {key: threading.Event for key in held}
recover = threading.Event
servers, workers = [], []
process = None

def check(name, condition, **facts):
 checks[name] = {'passed': bool(condition), **facts}

def snapshot:
 with lock:
 return list(events)

def handler(name):
 class Backend(BaseHTTPRequestHandler):
 protocol_version = 'HTTP/1.1'
 def log_message(self, *unused):
 pass
 def do_GET(self):
 route = self.path.split('?')[0]
 token = self.path.split('hold=')[-1] if 'hold=' in self.path else None
 with lock:
 events.append({'backend': name, 'path': self.path, 'host': self.headers.get('Host')})
 if token in held:
 held[token].set
 if not release[token].wait(15):
 raise RuntimeError('Bounded held-request deadline exceeded')
 status = 200
 if name == 'A':
 if route == '/passive' and not recover.is_set:
 status = 503
 if route in ['/noaccount', '/single']:
 status = 503
 if route == '/notfound':
 status = 404
 if route.startswith('/health') and self.headers.get('Host') == 'app.fund.test':
 status = 503
 body = name.encode
 self.send_response(status)
 self.send_header('Content-Length', str(len(body)))
 self.send_header('Connection', 'close')
 self.end_headers
 self.wfile.write(body)
 self.close_connection = True
 return Backend

for name in ['A', 'B']:
 server = ThreadingHTTPServer(('127.0.0.1', 0), handler(name))
 thread = threading.Thread(target=server.serve_forever)
 thread.start
 servers.append(server)
 workers.append(thread)
pa, pb = [s.server_port for s in servers]
with socket.socket as reservation:
 reservation.bind(('127.0.0.1', 0))
 port = reservation.getsockname[1]

def request(route):
 conn = http.client.HTTPConnection('127.0.0.1', port, timeout=20)
 try:
 conn.request('GET', route, headers={'Connection': 'close'})
 response = conn.getresponse
 return {'status': response.status, 'body': response.read.decode}
 finally:
 conn.close

def run(*options):
 result = subprocess.run([nginx, *options], capture_output=True, text=True, timeout=15)
 commands.append({'args': list(options), 'exitCode': result.returncode,
 'stdout': result.stdout, 'stderr': result.stderr})
 return result

try:
 with tempfile.TemporaryDirectory(prefix='dr-lb-capacity-') as directory:
 root = Path(directory)
 (root/'logs').mkdir
 config = f'''daemon off;
master_process on;
worker_processes 1;
pid {root}/nginx.pid;
error_log {root}/error.log notice;
events {{ worker_connections 128; }}
http {{
 log_format evidence escape=json '{{"path":"$request_uri","status":"$status","upstream":"$upstream_addr","upstreamStatus":"$upstream_status","worker":"$pid"}}'
 access_log {root}/access.jsonl evidence;
 proxy_http_version 1.1;
 proxy_set_header Connection close;
 proxy_read_timeout 20s;
 proxy_next_upstream_tries 2;
 upstream lc {{ least_conn; keepalive 0; server 127.0.0.1:{pa}; server 127.0.0.1:{pb}; }}
 upstream rr {{ keepalive 0; server 127.0.0.1:{pa}; server 127.0.0.1:{pb}; }}
 upstream cap {{ keepalive 0; server 127.0.0.1:{pa} max_conns=1; server 127.0.0.1:{pb} max_conns=1; }}
 upstream reserve {{ keepalive 0; server 127.0.0.1:{pa} max_conns=1; server 127.0.0.1:{pb} backup; }}
 upstream passive {{ keepalive 0; server 127.0.0.1:{pa} max_fails=1 fail_timeout=1s; server 127.0.0.1:{pb} backup; }}
 upstream noaccount {{ keepalive 0; server 127.0.0.1:{pa} max_fails=0; server 127.0.0.1:{pb} backup; }}
 upstream single {{ keepalive 0; server 127.0.0.1:{pa} max_fails=1 fail_timeout=60s; }}
 upstream notfound {{ keepalive 0; server 127.0.0.1:{pa} max_fails=1 fail_timeout=60s; server 127.0.0.1:{pb} backup; }}
 upstream health {{ keepalive 0; server 127.0.0.1:{pa}; }}
 server {{ listen 127.0.0.1:{port};
 location /ready {{ proxy_pass http://health; }}
 location /lc {{ proxy_pass http://lc; }}
 location /rr {{ proxy_pass http://rr; }}
 location /cap {{ proxy_pass http://cap; }}
 location /backup {{ proxy_pass http://reserve; }}
 location /passive {{ proxy_pass http://passive; proxy_next_upstream http_503; }}
 location /noaccount {{ proxy_pass http://noaccount; proxy_next_upstream http_503; }}
 location /single {{ proxy_pass http://single; proxy_next_upstream http_503; }}
 location /notfound {{ proxy_pass http://notfound; proxy_next_upstream http_404; }}
 location /health-generic {{ proxy_pass http://health; proxy_set_header Host default.fund.test; }}
 location /health-app {{ proxy_pass http://health; proxy_set_header Host app.fund.test; }}
 }}
}}
'''
 (root/'nginx.conf').write_text(config)
 version = run('-V')
 assert version.returncode == 0
 validation = run('-p', directory+'/', '-c', str(root/'nginx.conf'), '-t')
 assert validation.returncode == 0, validation.stderr
 with (root/'console.log').open('w') as console:
 process = subprocess.Popen([nginx, '-p', directory+'/', '-c', str(root/'nginx.conf')], stdout=console, stderr=console)
 for attempt in range(100):
 assert process.poll is None
 try:
 if request('/ready')['status'] == 200:
 break
 except OSError:
 pass
 time.sleep(.03)
 else:
 raise RuntimeError('Local NGINX did not become ready')
 with concurrent.futures.ThreadPoolExecutor(max_workers=2) as pool:
 for label in ['lc', 'rr']:
 pending = pool.submit(request, '/'+label+'?hold='+label)
 assert held[label].wait(3)
 holder = next(e['backend'] for e in snapshot if e['path']=='/'+label+'?hold='+label)
 samples = [request('/'+label)['body'] for _ in range(6)]
 before_release = not pending.done
 release[label].set
 completion = pending.result(timeout=3)
 details[label] = {'heldBackend': holder, 'samples': samples, 'completion': completion}
 if label == 'lc':
 check('least-conn-avoids-the-held-backend', holder == 'A' and samples == ['B']*6 and before_release,
 heldBackend=holder, otherBackendResponses=samples.count('B'), heldStillPending=before_release)
 else:
 counts = dict(collections.Counter(samples))
 check('round-robin-retains-its-selection-pattern-during-held-work', counts == {'A':3,'B':3} and before_release,
 counts=counts, heldStillPending=before_release)
 pending1 = pool.submit(request, '/cap?hold=cap1'); assert held['cap1'].wait(3)
 pending2 = pool.submit(request, '/cap?hold=cap2'); assert held['cap2'].wait(3)
 before = len(snapshot); overflow = request('/cap'); after = len(snapshot)
 check('max-conns-exhaustion-is-not-an-automatic-queue', overflow['status']==502 and before==after,
 status=overflow['status'], newBackendRequests=after-before, activeHeldRequests=2, queueConfigured=False)
 release['cap1'].set; release['cap2'].set
 completions = [pending1.result(timeout=3),pending2.result(timeout=3)]
 after_release = request('/cap')
 check('released-capacity-accepts-a-new-request', all(x['status']==200 for x in completions) and after_release['status']==200,
 completedHeldRequests=sum(x['status']==200 for x in completions), nextStatus=after_release['status'])
 pending = pool.submit(request, '/backup?hold=backup'); assert held['backup'].wait(3)
 reserve = request('/backup'); still_held = not pending.done; release['backup'].set; primary = pending.result(timeout=3)
 check('backup-can-serve-when-primary-reaches-connection-limit', reserve=={'status':200,'body':'B'} and primary['body']=='A' and still_held,
 primaryHeld=still_held, backupBody=reserve['body'], backupStatus=reserve['status'])
 before = len(snapshot); first=request('/passive'); first_events=snapshot[before:]
 before = len(snapshot); second=request('/passive'); second_events=snapshot[before:]
 check('passive-failure-excludes-primary-for-subsequent-request', first['body']=='B' and second['body']=='B' and [x['backend'] for x in first_events]==['A','B'] and [x['backend'] for x in second_events]==['B'],
 firstBackends=[x['backend'] for x in first_events], nextBackends=[x['backend'] for x in second_events])
 recover.set; before = len(snapshot); start=time.monotonic; time.sleep(2.2); elapsed=time.monotonic-start
 check('passive-idle-period-does-not-run-an-active-health-probe', len(snapshot)==before,
 newBackendRequests=len(snapshot)-before, configuredFailTimeoutSeconds=1, waitedBeyondTimeout=elapsed>2)
 before = len(snapshot); reentry=request('/passive'); reentry_events=snapshot[before:]
 check('recovered-primary-is-retried-after-passive-timeout', reentry=={'status':200,'body':'A'} and [x['backend'] for x in reentry_events]==['A'],
 status=reentry['status'], backend=reentry['body'], attempts=len(reentry_events))
 for route, expected_status, backends in [('noaccount',200,['A','B','A','B']),('single',503,['A','A']),('notfound',200,['A','B','A','B'])]:
 before=len(snapshot); values=[request('/'+route),request('/'+route)]; observed=[x['backend'] for x in snapshot[before:]]
 name={'noaccount':'max-fails-zero-disables-exclusion-accounting','single':'single-server-group-is-not-passively-excluded','notfound':'retrying-404-does-not-count-it-as-passive-failure'}[route]
 check(name, observed==backends and all(x['status']==expected_status for x in values),
 backends=observed, statuses=[x['status'] for x in values])
 generic=request('/health-generic'); representative=request('/health-app')
 check('default-host-health-can-hide-application-failure', generic['status']==200 and representative['status']==503,
 genericStatus=generic['status'], applicationHostStatus=representative['status'], activeHealthModuleUsed=False)
 process.send_signal(signal.SIGQUIT); process.wait(timeout=10); assert process.returncode==0
 access=[json.loads(line) for line in (root/'access.jsonl').read_text.splitlines]
 error=(root/'error.log').read_text
 assert 'no live upstreams' in error
 result={'executedAt':datetime.datetime.now(datetime.timezone.utc).isoformat,
 'pythonVersion':platform.python_version,'platform':platform.platform,
 'nginxVersion':version.stderr.strip,'binarySha256':hashlib.sha256(Path(nginx).read_bytes).hexdigest,
 'scriptSha256':hashlib.sha256(Path(__file__).read_bytes).hexdigest,
 'checks':checks,'details':details,'commands':commands,'configuration':config,
 'backendEvents':snapshot,'accessLog':access,'errorLog':error,
 'scope':'Actual NGINX with one worker, no shared upstream zone, upstream keepalive explicitly disabled, and two synthetic HTTP backends on IPv4 loopback. Bounded concurrency observations, not a throughput benchmark. No TLS, cloud load balancer, Kubernetes, HAProxy, active health module, multi-worker limit validation, real business writes, production integration or human workshop.'}
finally:
 for flag in release.values:flag.set
 if process and process.poll is None:
 process.send_signal(signal.SIGQUIT)
 try:process.wait(timeout=5)
 except subprocess.TimeoutExpired:process.kill;process.wait
 for server in servers:server.shutdown;server.server_close
 for thread in workers:thread.join(timeout=3)
result['childExited']=process.returncode==0
result['backendThreadsStopped']=all(not t.is_alive for t in workers)
result['temporaryDirectoryRemoved']=not root.exists
result['passed']=sum(v['passed'] for v in checks.values)
result['failed']=len(checks)-result['passed']
Path(args.output).write_text(json.dumps(result,indent=2)+'\n')
print(json.dumps({'passed':result['passed'],'failed':result['failed'],'failedChecks':[k for k,v in checks.items if not v['passed']],'childExited':result['childExited'],'backendThreadsStopped':result['backendThreadsStopped']}))
raise SystemExit(bool(result['failed']))
IN PRACTICE

Exercise: predict backend and status for six short requests with A held, then occupy both limits and explain the 502 without inventing an application failure.

Common pitfalls

Confusing counts with cost, backup use with a dead process, max_conns with a queue, or one local worker with a global production limit.

Related topics: Algorithms, affinity, and state · Health checks and readiness · Capacity and cascading failures

Take this idea with you

Selection and admission use concrete signals. Actual capacity needs representative measurement and observed recovery in addition to valid configuration.

Create account

Reference: NGINX upstream selection and limits · BigSavant load balancing 2026-09; selected NGINX, HAProxy 3.2, Kubernetes and AWS ALB behavior