← TCP/IP: fundamentos e diagnóstico
09 / 12 · 60 MIN

Pressão de envio e consumidores lentos

Usa sockets reais para observar envios parciais, retomar pelo offset certo e limitar a admissão de trabalho.

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']}))
NA PRÁTICA

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

Leva esta ideia contigo

Disponibilidade permite tentar I/O. O offset preserva bytes, a admissão limita trabalho e o resultado funcional exige evidência adicional.

Criar conta

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