Prepare a bounded experiment
Save the complete code below as run.py and run python3 run.py --output evidence.json in a working directory. The recorded execution used Python 3.13.1 on macOS, with IPv4 sockets on 127.0.0.1. Consulted documentation belongs to the 3.13 line and displayed 3.13.16. The script creates temporary ports, adjusts buffers only on its sockets and closes resources afterwards. Before running, predict which operations can progress with the consumer paused. Do not change global settings or use business endpoints for this exercise.
Distinguish readiness from completed work
The first group attempts a nonblocking read without available data. BlockingIOError is not EOF: the peer remains connected and has not yet sent a response. Next, the selector reports writing even though the receiving application has made no reads. Record local readiness, bytes accepted by sending and remotely consumed data separately. In a fictional positions feed, a green connection or write-readiness metric does not establish that positions reached business processing. Identify the missing observation before closing the incident.
Resume at the correct byte
The producer prepares 2 MiB of synthetic bytes and attempts to send them while the consumer is paused. A send can accept only a prefix; a later attempt can raise BlockingIOError. Retain the payload and add only confirmed positive returns to the offset. When the consumer resumes, send the outstanding suffix. The lab compares every received byte with the original. Restarting at offset zero would duplicate the prefix on the same stream. Applying a byte offset to text before UTF-8 encoding can also corrupt the sequence.
Manage selector interest
During drainage, the script monitors consumer reading and producer writing. Once all bytes are accepted, it removes write interest. Retaining that interest without pending data can cause repeated wakeups without work. Another group shows that read readiness can also signal EOF: recv returns empty bytes and the event can recur. Handle terminal state and remove interest in that direction. If the protocol still permits sending a response, separately retain write-direction state instead of confusing EOF with cancellation of the entire session.
Bound the application queue
BoundedQueue is a separate in-memory exercise with capacity for 64 payload bytes. It accepts two blocks of 32, rejects another byte without changing the queue and accepts 32 again after consuming 32. This is neither a TCP window measurement nor a limit on all RSS. In a fictional gateway, define what accepting work means, who retains it and how overload is communicated. Do not acknowledge persistence merely because a volatile queue accepted bytes. Admission policy must respect the producer’s recovery contract.
Connect observation to diagnosis
Fill a worksheet with arrival rate, consumption rate, pending bytes, oldest-item age and dependency time. The experiment only shows pressure with the reader paused and sequence recovery when reading resumes. It does not identify a slow database, packet loss or an advertised zero window. In a fictional incident these are hypotheses to distinguish through authorized observations. Enlarging buffers can delay symptoms without correcting imbalance. Define recovery through drainage and business results, including the fate of records admitted before mitigation.
"""Original bounded TCP application experiments; only synthetic IPv4 loopback traffic.
python3 run.py --output evidence.json
No packet capture, route changes, global kernel tuning, remote endpoints or business effects.
"""
import argparse, hashlib, json, pathlib, platform, selectors, socket, threading, time
from contextlib import ExitStack, contextmanager
from datetime import datetime, timezone
class DeadlineExpired(TimeoutError):
def __init__(self, received, reads):
super.__init__('operation deadline expired')
self.received, self.reads = bytes(received), reads
def receive_exact(sock, size, deadline):
data=bytearray;reads=0
while len(data)<size:
remaining=deadline-time.monotonic
if remaining<=0:raise DeadlineExpired(data,reads)
sock.settimeout(remaining);reads+=1
try:part=sock.recv(size-len(data))
except TimeoutError:raise DeadlineExpired(data,reads) from None
if not part:raise EOFError('incomplete application response')
data.extend(part)
return bytes(data),reads
class BoundedQueue:
"""A toy admission guard; bytes are counted, no durable job storage exists."""
def __init__(self,capacity):self.capacity=capacity;self.data=bytearray
def admit(self,data):
if len(data)>self.capacity-len(self.data):return False
self.data.extend(data);return True
def consume(self,count):
if not 0<=count<=len(self.data):raise ValueError('invalid consumption')
del self.data[:count]
def run:
endpoints=[];checks={};metrics={};sockets=[];threads=[]
with selectors.DefaultSelector as selector:
selector_class=type(selector).__name__
@contextmanager
def pair(label):
with ExitStack as stack:
listen=stack.enter_context(socket.socket);listen.settimeout(3);listen.bind(('127.0.0.1',0));listen.listen(1)
client=stack.enter_context(socket.socket);client.settimeout(3);client.connect(listen.getsockname)
server,_=listen.accept;stack.enter_context(server);server.settimeout(3);sockets.extend([listen,client,server])
endpoints.append({'label':label,'client':client.getsockname,'server':server.getsockname})
yield client,server
def check(label,details,valid):
checks[label]={'passed':bool(valid),**details}
if not valid:raise AssertionError(label+': '+json.dumps(details))
with pair('nonblocking-read') as (client,server):
client.setblocking(False)
try:client.recv(1)
except BlockingIOError:blocked=True
else:blocked=False
with selectors.DefaultSelector as sel:
sel.register(client,selectors.EVENT_READ);readable=bool(sel.select(.02))
check('would-block-is-not-stream-eof',{'wouldBlock':blocked,'readReadyWithoutData':readable,'peerClosed':False},blocked and not readable)
with pair('writability-and-pressure') as (client,server):
client.setsockopt(socket.SOL_SOCKET,socket.SO_SNDBUF,4096);server.setsockopt(socket.SOL_SOCKET,socket.SO_RCVBUF,4096)
client.setblocking(False);server.setblocking(False)
with selectors.DefaultSelector as sel:
sel.register(client,selectors.EVENT_WRITE);first=bool(sel.select(1));second=bool(sel.select(0))
check('writability-does-not-prove-peer-application-progress',{'writeReady':first,'readyAgainWithoutWork':second,'peerApplicationReads':0},first and second)
payload=bytes(range(251))*(2*1024*1024//251+1);payload=payload[:2*1024*1024]
offset=0;partial=False;blocked=False;writes=[]
while offset<len(payload):
try:n=client.send(memoryview(payload)[offset:])
except BlockingIOError:blocked=True;break
if n<=0:raise RuntimeError('no send progress')
partial|=n<len(payload)-offset;writes.append(n);offset+=n
metrics['pressure']={'bytesAcceptedBeforeReader':offset,'sendReturns':writes.copy,'effectiveSendBuffer':client.getsockopt(socket.SOL_SOCKET,socket.SO_SNDBUF),'effectiveReceiveBuffer':server.getsockopt(socket.SOL_SOCKET,socket.SO_RCVBUF)}
check('paused-reader-produces-bounded-send-pressure',{'wouldBlock':blocked,'someButNotAllBytesAccepted':0<offset<len(payload),'partialSendObserved':partial,'peerApplicationReads':0},blocked and 0<offset<len(payload) and partial)
received=bytearray;deadline=time.monotonic+15
with selectors.DefaultSelector as sel:
sel.register(client,selectors.EVENT_WRITE,'write');sel.register(server,selectors.EVENT_READ,'read')
while len(received)<len(payload):
if time.monotonic>=deadline:raise TimeoutError('drain did not complete')
for key,mask in sel.select(.1):
if key.data=='write':
try:n=client.send(memoryview(payload)[offset:])
except BlockingIOError:continue
if n<=0:raise RuntimeError('no send progress')
writes.append(n);offset+=n
if offset==len(payload):sel.unregister(client)
else:
try:part=server.recv(65536)
except BlockingIOError:continue
if not part:raise EOFError('unexpected EOF while draining')
received.extend(part)
check('resume-from-send-offset-preserves-exact-bytes',{'sentBytes':offset,'receivedBytes':len(received),'bodyMatches':bytes(received)==payload,'restartedFromZero':False},offset==len(payload) and bytes(received)==payload)
metrics['pressure'].update({'totalSendCalls':len(writes),'receivedSha256':hashlib.sha256(received).hexdigest})
queue=BoundedQueue(64);a=queue.admit(b'A'*32);b=queue.admit(b'B'*32);snapshot=bytes(queue.data);rejected=not queue.admit(b'C');unchanged=bytes(queue.data)==snapshot;queue.consume(32);resumed=queue.admit(b'C'*32)
check('application-admission-has-an-explicit-byte-bound',{'firstTwoAccepted':a and b,'thirdRejected':rejected,'rejectionPreservesQueue':unchanged,'resumedAfterConsumption':resumed,'pendingBytes':len(queue.data),'capacityBytes':64,'scope':'Pure in-memory admission guard, not TCP receive-window measurement.'},a and b and rejected and unchanged and resumed and len(queue.data)==64)
@contextmanager
def drip(server):
stop=threading.Event;errors=[]
def send:
try:
for byte in b'abcdefghij':
if stop.is_set:break
server.sendall(bytes([byte]))
if stop.wait(.08):break
except OSError as error:
if not stop.is_set:errors.append(str(error))
worker=threading.Thread(target=send,daemon=True);threads.append(worker);worker.start
try:yield
finally:
stop.set;worker.join(timeout=3)
if worker.is_alive:raise RuntimeError('fixture sender did not stop')
if errors:raise RuntimeError(errors)
with pair('per-read-timeout') as (client,server):
client.settimeout(.4);started=time.monotonic;body=bytearray
with drip(server):
while len(body)<10:
part=client.recv(1)
if not part:raise EOFError('drip response ended early')
body.extend(part)
elapsed=time.monotonic-started;metrics['perRead']={'elapsedSeconds':elapsed,'readTimeoutSeconds':.4,'intervalSeconds':.08}
check('per-read-timeout-does-not-bound-total-read-loop',{'complete':bytes(body)==b'abcdefghij','totalExceededSingleReadTimeout':elapsed>.4,'timeoutRaised':False},bytes(body)==b'abcdefghij' and elapsed>.4)
with pair('operation-deadline') as (client,server):
started=time.monotonic;expired=None
with drip(server):
try:receive_exact(client,10,started+.25)
except DeadlineExpired as error:expired=error
elapsed=time.monotonic-started
if expired is None:raise AssertionError('operation unexpectedly completed')
metrics['deadline']={'elapsedSeconds':elapsed,'budgetSeconds':.25,'partialBytes':len(expired.received),'recvCalls':expired.reads}
check('shared-deadline-rejects-incomplete-response',{'deadlineExpired':True,'someButNotAllBytes':0<len(expired.received)<10,'applicationResponseCommitted':False,'fixtureStopIsProtocolCancellation':False},0<len(expired.received)<10)
with pair('expired-before-read') as (client,server):
try:receive_exact(client,1,time.monotonic-1)
except DeadlineExpired as error:reads=error.reads;length=len(error.received)
else:raise AssertionError('expired deadline was ignored')
check('spent-budget-starts-no-new-receive',{'recvCalls':reads,'receivedBytes':length},reads==0 and length==0)
with pair('eof-readiness') as (client,server):
client.setblocking(False);server.shutdown(socket.SHUT_WR)
with selectors.DefaultSelector as sel:
sel.register(client,selectors.EVENT_READ);ready=bool(sel.select(1));data=client.recv(1);again=bool(sel.select(0));sel.unregister(client);registered=len(sel.get_map)
check('read-readiness-can-report-eof',{'readReady':ready,'receivedBytes':len(data),'readyAgainAfterEOF':again,'registeredAfterCleanup':registered},ready and data==b'' and again and registered==0)
def line(sock):
out=bytearray
while len(out)<64:
byte=sock.recv(1)
if not byte:raise EOFError('line unfinished')
out.extend(byte)
if byte==b'\n':return bytes(out)
raise ValueError('line exceeds fixture limit')
with pair('late-response-correlation') as (client,server):
client.sendall(b'R1\n');assert line(server)==b'R1\n'client.settimeout(.05)
try:client.recv(1)
except TimeoutError:pass
else:raise AssertionError('unexpected early response')
client.settimeout(3);client.sendall(b'R2\n');assert line(server)==b'R2\n'
server.sendall(b'R1|OK\n');late=line(client);late_id=late.split(b'|')[0];applied_to_r2=late_id==b'R2'
server.sendall(b'R2|OK\n');current=line(client);current_id=current.split(b'|')[0]
check('late-response-is-not-assigned-to-next-request',{'lateResponseId':late_id.decode,'lateAppliedToR2':applied_to_r2,'nextResponseId':current_id.decode,'businessEffects':0},not applied_to_r2 and current_id==b'R2')
return {'executedAt':datetime.now(timezone.utc).isoformat,'pythonVersion':platform.python_version,'platform':platform.platform,'selectorClass':selector_class,'checks':checks,'passed':sum(c['passed'] for c in checks.values),'failed':sum(not c['passed'] for c in checks.values),'metrics':metrics,'endpoints':endpoints,'allSocketsClosed':all(s.fileno==-1 for s in sockets),'allFixtureThreadsStopped':all(not t.is_alive for t in threads),'scriptSha256':hashlib.sha256(pathlib.Path(__file__).read_bytes).hexdigest,'scope':'Actual synthetic IPv4 TCP sockets on one loopback host plus an explicitly separate in-memory queue guard. No packet capture, TCP zero-window observation, TCP RTO measurement, global kernel tuning, WAN, DNS, TLS, IPv6, durable jobs, production workload or business effects. Read/send sizes are application observations, not captured segments. Fixture thread stopping is cleanup, not a network cancellation protocol.'}
if __name__=='__main__':
p=argparse.ArgumentParser;p.add_argument('--output',required=True);args=p.parse_args;result=run;pathlib.Path(args.output).write_text(json.dumps(result,indent=2)+'\n');print(json.dumps({k:result[k] for k in ['passed','failed','allSocketsClosed','allFixtureThreadsStopped']}))
Exercise: send returns 4096 for a 2097152-byte payload and the next call makes no progress. Write the resume offset, the bytes that must be retained and the evidence needed to confirm remote consumption.
Common pitfalls
Resending the prefix on the same stream, counting characters instead of bytes, subscribing to writing without work and calling BlockingIOError a zero window without a capture.
Related topics: Transport, acknowledgement, and messages · States, queues, and flow control · Diagnosis in the application context
Readiness permits attempting I/O. The offset preserves bytes, admission bounds work and business completion requires additional evidence.
Reference: Python socket interface · BigSavant TCP/IP 2026-09; TCP RFC 9293; IPv6 RFC 8200 with RFC 9673 update; Linux socket and iproute2 guidance