Preparar um ensaio limitado
Guarda o código completo abaixo como run.py e executa python3 run.py --output evidence.json numa pasta de trabalho. A execução registada usou Python 3.13.1 em macOS, com sockets IPv4 em 127.0.0.1. A documentação consultada pertence à linha 3.13 e identificava 3.13.16. O script cria portas temporárias, ajusta buffers apenas nos seus sockets e fecha os recursos no fim. Antes de executar, prevê quais as operações que podem progredir com o consumidor parado. Não alteres configurações globais nem uses endpoints de negócio para este exercício.
Distinguir disponibilidade de trabalho concluído
O primeiro grupo tenta ler num socket não bloqueante sem dados disponíveis. BlockingIOError não é EOF: o peer continua ligado e ainda não enviou a resposta. Depois, o selector indica escrita mesmo sem qualquer leitura pela aplicação recetora. Regista separadamente disponibilidade local, bytes aceites pelo envio e dados consumidos remotamente. Num feed fictício de posições, uma métrica verde de ligações ou de disponibilidade de escrita não demonstra que as posições chegaram ao processamento funcional. Procura a observação que falta antes de fechar o incidente.
Retomar pelo byte certo
O produtor prepara 2 MiB de bytes sintéticos e tenta enviá-los enquanto o consumidor está parado. Um send pode aceitar apenas um prefixo; uma tentativa seguinte pode levantar BlockingIOError. Mantém o payload e soma ao offset apenas os retornos positivos confirmados. Quando o consumidor retoma, envia o sufixo ainda pendente. O laboratório compara todos os bytes recebidos com o original. Reiniciar no offset zero duplicaria o prefixo na mesma stream. Aplicar um offset de bytes ao texto antes de codificar em UTF-8 também pode corromper a sequência.
Gerir interesse no selector
Durante a drenagem, o script acompanha leitura no consumidor e escrita no produtor. Quando todos os bytes forem aceites, retira o interesse de escrita. Manter esse interesse sem dados pendentes pode gerar despertares repetidos sem trabalho. Outro grupo demonstra que disponibilidade de leitura também pode sinalizar EOF: recv devolve bytes vazios e o evento pode repetir-se. Trata o estado terminal e retira o interesse dessa direção. Se o protocolo ainda permitir enviar uma resposta, conserva separadamente o estado da direção de escrita em vez de confundir EOF com cancelamento de toda a sessão.
Conter a fila de aplicação
A classe BoundedQueue é um exercício separado em memória, com capacidade de 64 bytes de payload. Aceita dois blocos de 32, rejeita outro byte sem alterar a fila e volta a aceitar 32 depois de consumir 32. Não é uma medição da janela TCP nem um limite para todo o RSS. Numa gateway fictícia, define o que significa aceitar trabalho, quem o conserva e como se comunica sobrecarga. Não confirmes persistência só porque uma fila volátil aceitou bytes. A política de admissão precisa de respeitar o contrato de recuperação do produtor.
Ligar a observação ao diagnóstico
Preenche uma grelha com taxa de entrada, taxa de consumo, bytes pendentes, idade do item mais antigo e tempo em dependências. O ensaio apenas mostra pressão com o leitor parado e recuperação da sequência quando volta a ler. Não identifica uma base de dados lenta, perda de pacotes ou janela anunciada zero. Num incidente fictício, estas são hipóteses a distinguir com observações autorizadas. Aumentar buffers pode adiar o sintoma sem corrigir o desequilíbrio. Define recuperação através da drenagem e dos resultados funcionais, incluindo o destino dos registos admitidos antes da mitigação.
"""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']}))
Exercício: send devolve 4096 para um payload de 2097152 bytes e a chamada seguinte não progride. Escreve o offset de retoma, os bytes que devem ficar retidos e a evidência necessária para confirmar consumo remoto.
Armadilhas comuns
Reenviar o prefixo na mesma stream, contar caracteres em vez de bytes, manter escrita subscrita sem trabalho e chamar zero-window a um BlockingIOError sem captura.
Tópicos relacionados: Transporte, confirmação e mensagens · Estados, filas e controlo de fluxo · Diagnóstico no contexto da aplicação
Disponibilidade permite tentar I/O. O offset preserva bytes, a admissão limita trabalho e o resultado funcional exige evidência adicional.
Referência: Python socket interface · DR TCP/IP 2026-09; TCP RFC 9293; IPv6 RFC 8200 with RFC 9673 update; Linux socket and iproute2 guidance