Fleet automation: subprocess fan-out, SSH and HTTP at scale
Bounded fan-out, process-group kills, pinned SSH host keys and partial failure.
tar -xzf scr-py-systems.tar.gz, which creates scr-py-systems/. SHA-256: 40e5cb3d361b7ed74c7a6e514cbbea0c7919c499abb7aa9e5b91cf727bee431cMost fleet automation is glue: your Python reaches out and drives many machines at once. In this lesson you will fan a probe out across many targets with bounded concurrency, a per-target timeout and an output cap, and stop it cleanly on SIGTERM without orphaning a single child process. You will drive SSH with paramiko against local targets whose host keys are pinned, so an unexpected key fails closed, quote every value that reaches a remote shell, and read a command's output in a way that cannot stall. Last, you will reuse one pooled HTTP client with connection limits. The clients are pinned in a venv (paramiko 5.0.0, httpx 0.28.1); the behaviour is the same on Ubuntu's Python 3.14.4 and upstream 3.14.7.
Owners of the basics: py-sec "Running commands safely" (subprocess list form, start_new_session, os.killpg), "Concurrency on Python 3.14" earlier in this course (asyncio, TaskGroup, asyncio.timeout) and py-sec "Calling HTTP APIs safely" (timeouts, 429s, requests.Session). Unpack the lesson files (the box at the top of this page) in your home directory: they create scr-py-systems/, and every command below runs inside ~/scr-py-systems. The terminals show the lab account deploy; the scripts use whoever runs them, so you will see your own user name. First the venv with the two pinned clients:
# Pinned clients for driving a fleet. Installed into the lab's venv.paramiko==5.0.0httpx==0.28.1
Bounded fan-out that leaves nothing behind
Running a command against 10,000 hosts at once would exhaust file descriptors and memory, and a few hung targets would hold the run open forever. A semaphore caps how many probes run at once and a per-target timeout bounds each one. Two failures are easier to miss. A target that prints without end (a log stuck in a loop, a hostile box) fills your memory if you keep everything it sends. And a probe that has started children (an SSH call, a pipeline) leaves them running if you kill only its main process, which is what happens when the controller itself is stopped.
The probe stays a Bash script: it is one host's check in a few commands, the kind of job "When a script has outgrown Bash" leaves in Bash. The fan-out, the deadlines and the per-target results are the part that moved to Python. This stand-in picks its behaviour from the target name:
#!/usr/bin/env bash# A stand-in for a remote probe. The target name picks the behaviour, so the fan-out meets fast# hosts, a failing host, one that never stops printing, one that prints non-UTF-8 bytes, and one# that hangs while holding a child process.set -euo pipefailcase "$1" inslow-*)# a child that outlives a naive kill of the parent: sleep in the background, then wait on our ownsleep 300 &sleep 300;;fail-*)echo "unreachable"exit 3;;big-*)yes "$1 log line" # endless output: only a byte cap stops it before the timeout;;bin-*)printf 'bad banner \377\376\n' # not valid UTF-8exit 2;;*)echo "ok $1";;esac
"""Run one probe in its own process group, with a deadline and a cap on how much output we hold.Whatever ends the run early (the timeout, too much output, a cancellation from the caller or aSIGTERM), the finally block kills the probe's whole process group, so no child outlives it."""import asyncioimport contextlibimport osimport signalTIMEOUT = 2.0 # seconds one probe may takeMAX_OUTPUT = 64_000 # bytes of output we are willing to hold per probeasync def read_capped(stream: asyncio.StreamReader, limit: int) -> bytes | None:"""Read to end of file, or return None as soon as the output passes limit bytes."""data = bytearray()while chunk := await stream.read(65536):data += chunkif len(data) > limit:return Nonereturn bytes(data)async def run_probe(target: str) -> tuple[str, str]:proc = await asyncio.create_subprocess_exec("./probe.sh", target,stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.STDOUT,start_new_session=True, # own process group, so killpg reaches every child it starts)finished = Falsetry:async with asyncio.timeout(TIMEOUT):out = await read_capped(proc.stdout, MAX_OUTPUT)if out is None:return "too much output", f"over {MAX_OUTPUT} bytes, probe killed"rc = await proc.wait()finished = Truetext = out.decode(errors="replace").strip() # a target's bytes may not be UTF-8return ("ok" if rc == 0 else f"exit {rc}"), textexcept TimeoutError:return "timed out", f"after {TIMEOUT:.0f}s, group killed"finally:if not finished: # timeout, output cap, cancellation or SIGTERM: stop the whole groupwith contextlib.suppress(ProcessLookupError): # the group may be gone alreadyos.killpg(proc.pid, signal.SIGKILL) # signal the whole process group# read what is left in the pipe (its writers are dead, so it ends) and reap the probe;# proc.wait() alone can hang here, because it also waits for the pipe to closeawait proc.communicate()
start_new_session=True makes each probe the leader of its own process group, so one os.killpg reaches every child it started (the mechanism is py-sec's; the new part is when it runs). read_capped stops reading at 64,000 bytes instead of holding whatever arrives, and decode(errors="replace") turns bytes that are not UTF-8 into U+FFFD instead of raising UnicodeDecodeError and ending the run. The finally block is the important part. finished becomes true only when the output was read to the end and the probe exited, so every other way out (the timeout, the output cap, a cancellation from the caller) kills the whole group. contextlib.suppress(ProcessLookupError) covers a group that is already gone. The block then calls communicate(), not wait(): asyncio's wait() also waits for the output pipe to close, and a pipe the capped reader stopped draining never closes, so with wait() there a probe stopped by the output cap hung the whole run until it was killed by hand.
"""Fan a probe out across many targets: at most CONCURRENCY at once, a result per target, a non-zeroexit on partial failure, and a clean stop on SIGTERM that leaves no probe running."""import asyncioimport signalimport sysfrom runprobe import run_probeCONCURRENCY = 5 # probes in flight at once, however many targets there areTARGETS = ([f"web{i:02d}" for i in range(12)]+ ["fail-db1", "fail-db2", "big-log1", "bin-fw1", "slow-gw1", "slow-gw2"])async def bounded(target: str, sem: asyncio.Semaphore) -> tuple[str, str, str]:async with sem:return (target, *await run_probe(target))async def main() -> int:# SIGTERM (systemctl stop, timeout, kill) or Ctrl-C cancels this task. The TaskGroup then cancels# every probe task, and each one's finally kills its process group.main_task, received = asyncio.current_task(), []def stop(sig: signal.Signals) -> None:received.append(sig)main_task.cancel()for sig in (signal.SIGTERM, signal.SIGINT):asyncio.get_running_loop().add_signal_handler(sig, stop, sig)sem = asyncio.Semaphore(CONCURRENCY)tasks = []try:async with asyncio.TaskGroup() as tg: # an unexpected error in one task cancels the resttasks = [tg.create_task(bounded(t, sem)) for t in TARGETS]except asyncio.CancelledError:if not received:raisestopped = sum(t.cancelled() for t in tasks)print(f"stopped by {received[0].name}: {stopped} unfinished probe(s) cancelled", file=sys.stderr)return 128 + received[0] # 143 for SIGTERM, 130 for SIGINT, as a shell reports themresults = [t.result() for t in tasks]bad = [r for r in results if r[1] != "ok"]print(f"{len(results) - len(bad)} ok, {len(bad)} not ok")for target, status, detail in bad:print(f" {target:9} {status}: {detail}")return 1 if bad else 0 # partial failure is a non-zero exitif __name__ == "__main__":sys.exit(asyncio.run(main()))
asyncio.Semaphore(5) is the bound: at most five probes run at once, so the cost is fixed however many targets there are. The TaskGroup waits for every task, and if one raises something unexpected it cancels the rest, which runs their finally blocks. add_signal_handler turns SIGTERM (from systemctl stop, timeout or kill) and SIGINT (Ctrl-C) into a cancellation of the main task, and the script then exits 128 plus the signal number, the status a shell reports for a job that signal ended: 143 for SIGTERM, 130 for SIGINT.
Twelve targets answered and six did not, each with its reason. big-log1 was cut off at the cap long before its timeout, bin-fw1's invalid bytes came through as replacement characters instead of a traceback, and the two slow- probes hit the 2-second timeout. The exit status is 1, the partial-failure signal a scheduler needs. pgrep found no probe and no sleep 300: the group kills took the background children too.
Now stop the run halfway, the way a service manager does. The first version of this script killed groups only in except TimeoutError. Stopped with timeout -s TERM 1, Python died at once (SIGTERM's default action) and both slow- probes and their four sleep 300 processes kept running: start_new_session had put them outside the controller's session, so no signal meant for the controller reached them. The version above, under the same stop:
The handler cancelled the main task, the TaskGroup cancelled the two unfinished probes, their finally blocks killed both groups, and the script exited 143 (--preserve-status makes timeout pass that on instead of its own 124). Ctrl-C takes the same path and exits 130; the lab checked that no probe survives it either. Under systemd, the default KillMode=control-group is the backstop: it kills every process in the unit's cgroup, whichever session it is in.
SSH with pinned host keys: an unexpected key fails closed
A host key is a server's cryptographic identity. The one line that decides whether SSH automation is safe is the host-key policy: what your code does when it meets a key it has not seen. paramiko offers RejectPolicy (its default), WarningPolicy and AutoAddPolicy (trust whatever answers, the wrong choice for anything that matters). Pin the keys you trust in a known_hosts file, load only those, and reject everything else, so a rebuilt server or a machine-in-the-middle offering a different key fails before any command runs.
You need two SSH servers to try this, and neither may be the machine's real sshd on port 22. sshd-targets.sh builds two throwaway ones in the lesson directory: a host key each, a client key, an authorized_keys, the target list and one config per server. It also writes the deliberate mistake: known_hosts pins host A's key for both ports.
#!/usr/bin/env bash# Build two throwaway SSH targets in the current directory. Run it as yourself, not with sudo.# Host A will listen on 127.0.0.1:18823 and host B on 127.0.0.1:18824, each with its own host key.# known_hosts pins host A's key for BOTH ports on purpose: host B then presents a key that does# not match its pin, which is the failure ssh_fleet.py must refuse.set -euo pipefaildir=$PWDuser=$(id -un)if [[ -e hostkey_a ]]; thenecho "targets already set up in $dir (remove hostkey_* and id_lab* to start again)" >&2exit 1fissh-keygen -q -t ed25519 -f hostkey_a -N '' -C lab-host-assh-keygen -q -t ed25519 -f hostkey_b -N '' -C lab-host-bssh-keygen -q -t ed25519 -f id_lab -N '' -C lab-client # lab only: no passphrase# restrict: no pty, no forwarding; from=: the key is refused from anywhere but this machineprintf 'restrict,from="127.0.0.1" %s\n' "$(cat id_lab.pub)" > authorized_keyschmod 0600 hostkey_a hostkey_b id_lab authorized_keyspin_a=$(cut -d' ' -f1,2 hostkey_a.pub)printf '[127.0.0.1]:18823 %s\n[127.0.0.1]:18824 %s\n' "$pin_a" "$pin_a" > known_hostsprintf '127.0.0.1 18823\n127.0.0.1 18824\n' > targets.txtfor target in a:18823 b:18824; doid=${target%%:*} port=${target##*:}cat > "sshd_$id.conf" <<CONFPort $portListenAddress 127.0.0.1HostKey $dir/hostkey_$idPidFile $dir/sshd_$id.pidAuthorizedKeysFile $dir/authorized_keysAllowUsers $userPasswordAuthentication noKbdInteractiveAuthentication noUsePAM yesStrictModes noPrintMotd noCONFdoneecho "keys, known_hosts, targets.txt, sshd_a.conf and sshd_b.conf are in $dir"
ssh-keygen -lf prints each key's fingerprint. Both pins carry host A's fingerprint, and host B's own key has another one, so host B will present a key that does not match its pin. UsePAM yes with password and keyboard-interactive logins switched off lets a key login work even for an account whose password is locked, as deploy's is; StrictModes no is there because the keys live in the lesson directory instead of ~/.ssh. sshd has to start as root, because after a login it switches to your user; these are the only two commands in the lesson that need sudo, and the last section stops both servers:
"""How every script in this lesson connects: pinned host keys, the lab key, bounded connect."""import getpassimport paramikodef connect(host: str, port: int) -> paramiko.SSHClient:client = paramiko.SSHClient()client.load_host_keys("known_hosts") # only these keys are trustedclient.set_missing_host_key_policy(paramiko.RejectPolicy()) # an unknown host key fails closedtry:client.connect(host, port=port, username=getpass.getuser(),key_filename="id_lab", # lab only: an unencrypted key file on diskallow_agent=False, look_for_keys=False,timeout=5, banner_timeout=5, auth_timeout=5)except BaseException:client.close()raisereturn client
"""Run one command across an SSH fleet: a bounded thread pool, a deadline on every host, cappedoutput, a result per host and a non-zero exit on partial failure."""import shleximport sysfrom concurrent.futures import ThreadPoolExecutorfrom pathlib import Pathimport paramikofrom labssh import connectCONCURRENCY = 10 # hosts in flight at once; paramiko blocks, so each one gets a threadDEADLINE = 10 # seconds the command may run, enforced on the host by timeout(1)MAX_OUTPUT = 64_000 # bytes of output kept per hostdef check_host(host: str, port: int, command: str) -> tuple[bool, str]:name = f"{host}:{port}"try:with connect(host, port) as client: # closes the connection on every pathremote = f"timeout {DEADLINE} sh -c {shlex.quote(command)}"_, stdout, _ = client.exec_command(remote, timeout=DEADLINE + 5)stdout.channel.set_combine_stderr(True) # one stream, so an unread stderr cannot stall usout = stdout.read(MAX_OUTPUT + 1) # drain the output before the exit statusif len(out) > MAX_OUTPUT:return False, f"{name}: ERROR output over {MAX_OUTPUT} bytes"rc = stdout.channel.recv_exit_status()except paramiko.BadHostKeyException as err:return False, f"{name}: ERROR host key mismatch ({err.expected_key.get_name()})"except Exception as err: # refused, timed out, auth failed: a result for this host, not a crashreturn False, f"{name}: ERROR {type(err).__name__}: {err}"text = out.decode(errors="replace").strip() # a host's bytes may not be UTF-8if rc != 0:return False, f"{name}: exit {rc} {text}".rstrip()return True, f"{name}: {text}"def main() -> int:command = sys.argv[1]targets = [line.split() for line in Path("targets.txt").read_text().splitlines() if line.strip()]with ThreadPoolExecutor(max_workers=CONCURRENCY) as pool:results = list(pool.map(lambda t: check_host(t[0], int(t[1]), command), targets))for _, line in results:print(line)ok = sum(good for good, _ in results)print(f"{ok}/{len(results)} hosts ok")return 0 if ok == len(results) else 1 # any unreachable, mismatched or failing host is a non-zero exitif __name__ == "__main__":sys.exit(main())
load_host_keys loads the curated file and RejectPolicy refuses any key not in it. The three connect timeouts cover the three ways a connect can stall: the TCP handshake, the SSH greeting and the login. They bound silence, not duration: a command that prints a line every few seconds would never trip a 5-second read timeout. So ssh_fleet.py runs the command under timeout 10 on the host, quoted with shlex.quote (next section), and keeps at most 64,000 bytes of output. paramiko blocks, so the hosts run in a ThreadPoolExecutor of ten threads, the same bound as the semaphore above; a loop over 1,000 hosts would otherwise take as long as all their timeouts added up.
Host A answered with the user name. Host B raised BadHostKeyException, caught and reported as a mismatch, and the fleet exited 1 because not every host succeeded. The wrong key stops the command: nothing prompts and nothing is trusted on first use. AutoAddPolicy would have recorded host B's key and run the command, which is how a machine-in-the-middle gets handed your session. In the second run sleep 60 hit the host-side deadline: timeout ended it after ten seconds with status 124, and the whole run took just over ten seconds.
Pinning at fleet scale needs a way to update the pins, or someone under pressure switches to AutoAddPolicy the first time a host is rebuilt. Generate known_hosts from your provisioning inventory, the keys your build system recorded when it created each host, and ship the file with the tool; rotation means updating the inventory, never trusting first contact. OpenSSH's own client can trust a host certificate authority with a @cert-authority line; paramiko 5.0 does not parse those lines, so for paramiko the inventory file is the way. The client key here is lab only: an unencrypted file on disk. For real hosts use an agent or a short-lived SSH certificate, and restrict the key in authorized_keys as the setup script does (restrict drops pty and forwarding, from= limits where it is accepted).
Remote commands run through a shell: quote, then drain
exec_command hands its string to the remote user's login shell, so the list-form protection that makes local subprocess safe does not exist over SSH: a shell metacharacter in an interpolated value starts a second command on the far side. It is the problem "Trust boundaries" solved for Bash's ssh with ${var@Q}; Python's answer is shlex.quote, whose single quotes also suit a remote sh.
"""exec_command runs through the remote login shell, so the list-form protection does not apply.Any untrusted value put into a remote command must be quoted with shlex.quote, or a shellmetacharacter in it starts a second command on the far side. Usage: remote_quote.py VALUE unsafe|safe"""import shleximport sysfrom labssh import connectvalue, mode = sys.argv[1], sys.argv[2]if mode == "unsafe":command = f"wc -l {value}" # value pasted straight into the remote command lineelse:command = f"wc -l {shlex.quote(value)}" # value folded into a single, inert argumentwith connect("127.0.0.1", 18823) as client:_, stdout, stderr = client.exec_command(command, timeout=5)out = stdout.read().decode(errors="replace") # a few bytes each: reading one stream, then theerr = stderr.read().decode(errors="replace") # other, is safe only for small output (ssh_drain.py)rc = stdout.channel.recv_exit_status()print(f"sent : {command}")print(f"result : rc={rc} out={out.strip()!r} err={err.strip()!r}")
Unsafe, the ; ended the wc command and echo INJECTED-REMOTE ran as a second command on the server. Quoted, shlex.quote folded the whole value into one argument, so the remote wc looked for a single oddly-named file and failed cleanly.
The second habit is reading output before the exit status, and reading all of it. SSH flow control gives each channel a window: the server may send only that many bytes (2 MiB by default in paramiko) before your side says it has room for more, and one window covers stdout and stderr together (RFC 4254, section 5.2). paramiko opens the window again only as your code reads. ssh_drain.py runs a command that writes about 6.9 MB to stderr and one line to stdout, and reads it three ways:
"""Read a remote command's output before its exit status, without stalling on stderr.stdout and stderr share one SSH flow-control window, and paramiko reopens it only as you read.The command writes about 6.9 MB to stderr and one line to stdout. Usage:ssh_drain.py status-first | stdout-first | combined"""import sysfrom labssh import connectMODE = sys.argv[1]COMMAND = "seq 1 1000000 >&2; echo done"with connect("127.0.0.1", 18823) as client:_, stdout, stderr = client.exec_command(COMMAND, timeout=5) # 5 s without data raiseschan = stdout.channeltry:if MODE == "status-first":rc = chan.recv_exit_status() # waits for the command, which waits for a readerout = stdout.read()elif MODE == "stdout-first":out = stdout.read() + stderr.read() # stdout waits for EOF while stderr fills the windowrc = chan.recv_exit_status()else:chan.set_combine_stderr(True) # stderr joins stdout: one stream to drainout = stdout.read()rc = chan.recv_exit_status() # only after the output has been readexcept TimeoutError:held = len(chan.recv_stderr(8 * 2**20)) # what arrived on stderr and was never readprint(f"{MODE}: stalled, no data for 5 s; {held} bytes of stderr were waiting unread")sys.exit(1)print(f"{MODE}: read {len(out)} bytes, exit status {rc}")
Asking for the status first hung until timeout 8 ended it with 124: recv_exit_status() has no timeout, and the command could not exit while nobody read its output. paramiko's documentation warns about exactly this. Reading stdout to the end before stderr looks correct and still stalled: stdout never ended, because the remote command was blocked on stderr, and the unread stderr held the whole 2 MiB window. The fix is to drain both streams at once. set_combine_stderr(True) merges them into one, which ssh_fleet.py does because it reports them together; when you need them apart, read stderr in a second thread. Only then ask for the status.
A pooled HTTP client with connection limits
Calling an API per host, or one host many times, should reuse connections instead of opening a new socket each time. One client with a connection pool holds the TLS settings and keeps sockets warm; connection limits keep a burst from opening thousands at once. This uses httpx instead of requests for two reasons: it exposes the pool size directly as httpx.Limits, and it has one API for the sync and async clients, so the same shape fits the asyncio fan-out from the first section. (requests.Session pools too, but the pool size is set through a mounted adapter.) server.py from the lesson files is the local API: /ok, /slow (0.2 s), /fail (503) and /stats, which counts the TCP connections and requests it has seen. Start it in the background; the until loop waits at most 20 seconds for it to answer:
"""A pooled HTTP client with connection limits and a timeout, reusing sockets across many requests.httpx over requests here for two reasons: it exposes the pool size directly as httpx.Limits, and ithas one API for the sync client and the async client, so the same code shape fits the asyncio fan-outabove. requests.Session also pools, but the pool size is set indirectly through a mounted adapter."""import httpxBASE = "http://127.0.0.1:18825"def main() -> None:limits = httpx.Limits(max_connections=5, max_keepalive_connections=5)timeout = httpx.Timeout(5.0, connect=2.0) # a timeout on every request; never wait foreverwith httpx.Client(base_url=BASE, limits=limits, timeout=timeout) as client:ok = sum(client.get("/ok").status_code == 200 for _ in range(50))fail = client.get("/fail") # a 5xx is surfaced here; retrying it belongs to "Failure design"stats = client.get("/stats").json()print(f"{ok}/50 ok, /fail -> {fail.status_code}")print(f"server saw {stats['connections']} connection(s) for {stats['requests']} requests")if __name__ == "__main__":main()
Fifty requests succeeded, /fail returned 503 (surfaced, not retried here, because retry policy belongs to "Failure design"), and the server saw two connections: the curl readiness check in the start step, and one reused connection for all 52 of the client's requests (the 53 the server counted include curl's). Every request had a timeout, and the pool cap means a fan-out to many URLs never opens more than five sockets at once. The same Limits on httpx.AsyncClient bounds an async fan-out.
Try this
The pooled client above is synchronous. Rewrite it as apool.py with httpx.AsyncClient and an asyncio.Semaphore(5), firing 50 requests at /slow (each sleeps 0.2 s) with at most five in flight. Predict the wall-clock time first: 50 requests, 5 at a time, is ten rounds of 0.2 s, about two seconds, not the ten a sequential client would take. The lab measured about 2.7 s and printed 50/50 ok. Then explain why raising the semaphore to 50 while leaving max_connections=5 would not make it ten times faster.
When you are done, stop the HTTP server and both sshd targets; the wait loop confirms that nothing listens on the three lesson ports any more. rm -rf ~/scr-py-systems then removes every key, config and log the lesson created.
Takeaway
Every early exit of a fan-out, including SIGTERM, must kill the process groups it started, and every wait and every buffer needs a bound. Over SSH, pin host keys from inventory and reject the unknown, quote what reaches the remote shell, and drain both streams before asking for the exit status.
start_new_session=True and calls os.killpg in except TimeoutError. Under systemd it works, but when the same script is stopped with timeout -s TERM from cron, probe processes are left running. Why, and what fixes it?client.set_missing_host_key_policy(paramiko.AutoAddPolicy()) and a curated known_hosts. A reviewer flags it. What is the concrete risk, and the fix?rc = stdout.channel.recv_exit_status() first, then out = stdout.read(). It works for months, then hangs on one host running a command with large output. Why?