Networking con Python

Networking con Python

Python offre uno degli ecosistemi più ricchi per la programmazione di rete, ma già la sola libreria standard copre ogni livello: il modulo socket espone l'API BSD quasi senza filtri, selectors permette il multiplexing dell'I/O, asyncio fornisce un framework completo per server e client asincroni e ssl aggiunge la cifratura TLS. In questo articolo partiremo dai socket bloccanti per arrivare a server asincroni con framing, timeout e chiusura ordinata, passando per UDP, DNS e TLS.

Socket TCP bloccanti

Il modulo socket è un wrapper sottile sopra le chiamate di sistema. Un client TCP essenziale apre una connessione, invia dati e legge la risposta. La funzione socket.create_connection() è preferibile alla creazione manuale del socket perché risolve il nome, prova tutti gli indirizzi restituiti (IPv4 e IPv6) e applica il timeout.

import socket


def fetch_banner(host: str, port: int, timeout: float = 3.0) -> str:
    """Si connette a un servizio e restituisce il banner iniziale."""
    # create_connection prova in sequenza tutti gli indirizzi risolti
    with socket.create_connection((host, port), timeout=timeout) as sock:
        try:
            data = sock.recv(1024)
        except socket.timeout:
            return ""
        return data.decode("utf-8", errors="replace").strip()


if __name__ == "__main__":
    for host, port in [("smtp.gmail.com", 587), ("github.com", 22)]:
        try:
            print(f"{host}:{port} -> {fetch_banner(host, port)}")
        except OSError as exc:
            print(f"{host}:{port} -> errore: {exc}")

Il parametro timeout si applica sia alla connessione sia alle successive operazioni sul socket. Senza di esso, recv() potrebbe restare bloccata indefinitamente in attesa di dati che non arriveranno mai.

Leggere e scrivere tutto: sendall e il framing

Due trappole classiche dei socket TCP riguardano le operazioni parziali. send() può inviare solo una parte del buffer, ed è per questo che esiste sendall(). recv(n), invece, restituisce al massimo n byte e non c'è un equivalente recvall(): va scritto a mano. Inoltre TCP non conserva i confini dei messaggi, quindi serve un protocollo di framing. Usiamo un prefisso di lunghezza a 32 bit codificato con struct:

import socket
import struct

HEADER = struct.Struct("!I")  # intero senza segno a 32 bit, network byte order
MAX_FRAME_SIZE = 16 * 1024 * 1024


class FrameError(Exception):
    pass


def send_frame(sock: socket.socket, payload: bytes) -> None:
    if len(payload) > MAX_FRAME_SIZE:
        raise FrameError(f"payload troppo grande: {len(payload)} byte")
    # sendall continua a inviare finché tutti i byte sono stati trasmessi
    sock.sendall(HEADER.pack(len(payload)) + payload)


def recv_exactly(sock: socket.socket, size: int) -> bytes:
    buffer = bytearray()
    while len(buffer) < size:
        chunk = sock.recv(size - len(buffer))
        if not chunk:
            # Una lettura vuota indica che il peer ha chiuso la connessione
            raise ConnectionError("connessione chiusa durante la lettura")
        buffer.extend(chunk)
    return bytes(buffer)


def recv_frame(sock: socket.socket) -> bytes:
    (length,) = HEADER.unpack(recv_exactly(sock, HEADER.size))
    if length > MAX_FRAME_SIZE:
        raise FrameError(f"frame troppo grande: {length} byte")
    return recv_exactly(sock, length)

La verifica su MAX_FRAME_SIZE protegge il ricevitore da frame che dichiarano dimensioni enormi, siano essi malevoli o frutto di dati corrotti.

Un server TCP con socketserver

Per server semplici il modulo socketserver elimina gran parte del codice ripetitivo. Con ThreadingTCPServer ogni connessione viene gestita in un thread separato:

import socketserver


class EchoHandler(socketserver.StreamRequestHandler):
    # Timeout di inattività applicato al socket di ogni client
    timeout = 30

    def handle(self) -> None:
        peer = f"{self.client_address[0]}:{self.client_address[1]}"
        print(f"nuova connessione da {peer}")

        try:
            # rfile è un file bufferizzato: readline gestisce il framing a righe
            for raw_line in self.rfile:
                line = raw_line.decode("utf-8", errors="replace").rstrip("\r\n")
                self.wfile.write(f"echo: {line}\n".encode("utf-8"))
        except TimeoutError:
            print(f"{peer}: timeout di inattività")
        finally:
            print(f"{peer} disconnesso")


class ReusableServer(socketserver.ThreadingTCPServer):
    # Permette di riavviare il server senza attendere il rilascio della porta
    allow_reuse_address = True
    daemon_threads = True


if __name__ == "__main__":
    with ReusableServer(("0.0.0.0", 9000), EchoHandler) as server:
        print("server in ascolto sulla porta 9000")
        try:
            server.serve_forever()
        except KeyboardInterrupt:
            print("arresto in corso...")

L'opzione allow_reuse_address imposta SO_REUSEADDR: senza di essa, riavviando il server subito dopo l'arresto si otterrebbe l'errore Address already in use, perché le connessioni chiuse restano per qualche tempo nello stato TIME_WAIT.

Multiplexing con selectors

Un thread per connessione funziona bene fino a qualche centinaio di client. Oltre, conviene un modello a evento singolo: un unico thread che sorveglia tutti i socket e reagisce solo a quelli pronti. Il modulo selectors sceglie automaticamente il meccanismo più efficiente del sistema (epoll su Linux, kqueue su macOS e BSD).

import selectors
import socket

selector = selectors.DefaultSelector()
buffers: dict[socket.socket, bytearray] = {}


def accept(server: socket.socket) -> None:
    client, address = server.accept()
    client.setblocking(False)
    buffers[client] = bytearray()
    selector.register(client, selectors.EVENT_READ, handle_client)
    print(f"connesso {address}")


def handle_client(client: socket.socket) -> None:
    try:
        data = client.recv(4096)
    except ConnectionResetError:
        data = b""

    if not data:
        close_client(client)
        return

    buffer = buffers[client]
    buffer.extend(data)

    # Si elaborano solo le righe complete, il resto rimane nel buffer
    while (index := buffer.find(b"\n")) != -1:
        line = bytes(buffer[:index]).rstrip(b"\r")
        del buffer[: index + 1]
        client.sendall(b"echo: " + line + b"\n")


def close_client(client: socket.socket) -> None:
    selector.unregister(client)
    buffers.pop(client, None)
    client.close()


def main() -> None:
    server = socket.create_server(("0.0.0.0", 9000), reuse_port=False)
    server.setblocking(False)
    selector.register(server, selectors.EVENT_READ, accept)

    print("server con selectors in ascolto sulla porta 9000")

    while True:
        # select restituisce solo i socket pronti per l'operazione richiesta
        for key, _ in selector.select(timeout=1):
            callback = key.data
            callback(key.fileobj)


if __name__ == "__main__":
    main()

Per semplicità l'esempio usa sendall() anche sui socket non bloccanti; in un server reale le scritture andrebbero bufferizzate e completate quando il socket risulta pronto in scrittura (EVENT_WRITE). È proprio questa complessità che asyncio nasconde.

Server asincrono con asyncio

asyncio offre un'API ad alto livello basata su stream: asyncio.start_server() invoca una coroutine per ogni connessione, passandole uno StreamReader e uno StreamWriter. Il codice sembra sequenziale, ma ogni await cede il controllo all'event loop.

import asyncio
import json
import signal
import struct

HEADER = struct.Struct("!I")
MAX_FRAME_SIZE = 1024 * 1024
IDLE_TIMEOUT = 30


async def read_frame(reader: asyncio.StreamReader) -> bytes:
    # readexactly solleva IncompleteReadError se lo stream termina prima
    header = await reader.readexactly(HEADER.size)
    (length,) = HEADER.unpack(header)
    if length > MAX_FRAME_SIZE:
        raise ValueError(f"frame troppo grande: {length} byte")
    return await reader.readexactly(length)


async def write_frame(writer: asyncio.StreamWriter, payload: bytes) -> None:
    writer.write(HEADER.pack(len(payload)) + payload)
    # drain applica la backpressure: attende se il buffer di scrittura è pieno
    await writer.drain()


def dispatch(request: dict) -> dict:
    command = request.get("command")
    params = request.get("params", {})

    if command == "ping":
        return {"result": "pong"}
    if command == "sum":
        return {"result": sum(params.get("values", []))}
    return {"error": f"comando sconosciuto: {command}"}


class JsonServer:
    def __init__(self, host: str, port: int) -> None:
        self.host = host
        self.port = port
        self.connections: set[asyncio.Task] = set()

    async def handle(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
        task = asyncio.current_task()
        self.connections.add(task)
        peer = writer.get_extra_info("peername")

        try:
            while True:
                try:
                    # Timeout di inattività sulla lettura del frame successivo
                    async with asyncio.timeout(IDLE_TIMEOUT):
                        frame = await read_frame(reader)
                except TimeoutError:
                    print(f"{peer}: timeout di inattività")
                    break

                try:
                    request = json.loads(frame)
                    response = {"id": request.get("id"), **dispatch(request)}
                except json.JSONDecodeError:
                    response = {"error": "JSON non valido"}

                await write_frame(writer, json.dumps(response).encode("utf-8"))
        except asyncio.IncompleteReadError:
            pass  # il client ha chiuso la connessione
        except (ConnectionResetError, ValueError) as exc:
            print(f"{peer}: {exc}")
        finally:
            writer.close()
            await writer.wait_closed()
            self.connections.discard(task)

    async def run(self) -> None:
        server = await asyncio.start_server(self.handle, self.host, self.port)
        stop = asyncio.Event()

        loop = asyncio.get_running_loop()
        for sig in (signal.SIGINT, signal.SIGTERM):
            loop.add_signal_handler(sig, stop.set)

        print(f"server asyncio in ascolto su {self.host}:{self.port}")

        async with server:
            await stop.wait()

            # Stop alle nuove connessioni, poi attesa limitata di quelle aperte
            server.close()
            print("arresto in corso...")

            if self.connections:
                _, pending = await asyncio.wait(self.connections, timeout=5)
                for task in pending:
                    task.cancel()


if __name__ == "__main__":
    asyncio.run(JsonServer("0.0.0.0", 9000).run())

Il context manager asyncio.timeout(), disponibile da Python 3.11, è il modo moderno di imporre un limite di tempo a qualsiasi blocco di codice asincrono. Va notato anche await writer.drain(): senza di esso, un client lento farebbe crescere senza limite il buffer di scrittura in memoria.

Client asincrono con richieste concorrenti

Il client corrispondente si connette con asyncio.open_connection(). Per mostrare la potenza del modello asincrono, ogni richiesta apre una propria connessione e tutte vengono eseguite in parallelo con un TaskGroup:

import asyncio
import json
import struct

HEADER = struct.Struct("!I")


async def call(host: str, port: int, request: dict, timeout: float = 5.0) -> dict:
    async with asyncio.timeout(timeout):
        reader, writer = await asyncio.open_connection(host, port)
        try:
            payload = json.dumps(request).encode("utf-8")
            writer.write(HEADER.pack(len(payload)) + payload)
            await writer.drain()

            (length,) = HEADER.unpack(await reader.readexactly(HEADER.size))
            return json.loads(await reader.readexactly(length))
        finally:
            writer.close()
            await writer.wait_closed()


async def main() -> None:
    requests = [
        {"id": 1, "command": "ping"},
        {"id": 2, "command": "sum", "params": {"values": [10, 20, 30]}},
        {"id": 3, "command": "unknown"},
    ]

    # Il TaskGroup cancella tutti i task se uno di essi fallisce
    async with asyncio.TaskGroup() as group:
        tasks = [group.create_task(call("127.0.0.1", 9000, request)) for request in requests]

    for task in tasks:
        print(task.result())


if __name__ == "__main__":
    asyncio.run(main())

Un port scanner asincrono

Lo stesso approccio si presta a uno strumento di diagnostica: verificare centinaia di porte in parallelo, limitando la concorrenza con un asyncio.Semaphore per non esaurire i descrittori di file.

import asyncio
import time


async def probe(host: str, port: int, semaphore: asyncio.Semaphore, timeout: float) -> tuple[int, bool, float]:
    async with semaphore:
        start = time.perf_counter()
        try:
            async with asyncio.timeout(timeout):
                _, writer = await asyncio.open_connection(host, port)
        except (OSError, TimeoutError):
            return port, False, 0.0

        latency = (time.perf_counter() - start) * 1000
        writer.close()
        await writer.wait_closed()
        return port, True, latency


async def scan(host: str, ports: range, concurrency: int = 200, timeout: float = 1.0) -> None:
    semaphore = asyncio.Semaphore(concurrency)
    results = await asyncio.gather(*(probe(host, port, semaphore, timeout) for port in ports))

    for port, is_open, latency in results:
        if is_open:
            print(f"{port:>5}  aperta  ({latency:.1f} ms)")


if __name__ == "__main__":
    asyncio.run(scan("127.0.0.1", range(1, 1025)))

UDP con asyncio

Per UDP asyncio non offre stream ma un'API basata su protocolli: si definisce una classe che implementa datagram_received(). Ecco un server che risponde a un semplice protocollo di heartbeat e un client che misura il round-trip time:

import asyncio
import json
import time


class HeartbeatServer(asyncio.DatagramProtocol):
    def connection_made(self, transport: asyncio.DatagramTransport) -> None:
        self.transport = transport

    def datagram_received(self, data: bytes, addr: tuple[str, int]) -> None:
        try:
            message = json.loads(data)
        except json.JSONDecodeError:
            return

        # La risposta riporta il timestamp originale per il calcolo dell'RTT
        reply = {"type": "pong", "sent_at": message.get("sent_at")}
        self.transport.sendto(json.dumps(reply).encode("utf-8"), addr)


class HeartbeatClient(asyncio.DatagramProtocol):
    def __init__(self) -> None:
        self.responses: asyncio.Queue[dict] = asyncio.Queue()

    def datagram_received(self, data: bytes, addr: tuple[str, int]) -> None:
        self.responses.put_nowait(json.loads(data))


async def main() -> None:
    loop = asyncio.get_running_loop()

    server_transport, _ = await loop.create_datagram_endpoint(
        HeartbeatServer, local_addr=("127.0.0.1", 9999)
    )

    client_transport, client = await loop.create_datagram_endpoint(
        HeartbeatClient, remote_addr=("127.0.0.1", 9999)
    )

    try:
        for sequence in range(5):
            payload = {"type": "ping", "seq": sequence, "sent_at": time.perf_counter()}
            client_transport.sendto(json.dumps(payload).encode("utf-8"))

            try:
                # UDP non garantisce la consegna: senza risposta si passa oltre
                async with asyncio.timeout(1):
                    reply = await client.responses.get()
                rtt = (time.perf_counter() - reply["sent_at"]) * 1000
                print(f"seq={sequence} rtt={rtt:.3f} ms")
            except TimeoutError:
                print(f"seq={sequence} perso")

            await asyncio.sleep(0.5)
    finally:
        client_transport.close()
        server_transport.close()


if __name__ == "__main__":
    asyncio.run(main())

Risoluzione DNS

La funzione fondamentale per la risoluzione è socket.getaddrinfo(), che usa il resolver di sistema e restituisce tutte le combinazioni di famiglia, tipo e indirizzo. In codice asincrono va usata la variante loop.getaddrinfo(), che la esegue in un thread separato senza bloccare l'event loop:

import asyncio
import socket


async def resolve(hostname: str) -> dict[str, list[str]]:
    loop = asyncio.get_running_loop()
    infos = await loop.getaddrinfo(hostname, None, type=socket.SOCK_STREAM)

    result: dict[str, list[str]] = {"ipv4": [], "ipv6": []}
    for family, _, _, _, sockaddr in infos:
        address = sockaddr[0]
        key = "ipv4" if family == socket.AF_INET else "ipv6"
        if address not in result[key]:
            result[key].append(address)
    return result


async def reverse(ip: str) -> str | None:
    loop = asyncio.get_running_loop()
    try:
        # gethostbyaddr è bloccante: viene eseguita in un thread del pool
        hostname, _, _ = await loop.run_in_executor(None, socket.gethostbyaddr, ip)
        return hostname
    except socket.herror:
        return None


async def main() -> None:
    hosts = ["python.org", "example.com", "dominio-inesistente.invalid"]
    results = await asyncio.gather(*(resolve(host) for host in hosts), return_exceptions=True)

    for host, result in zip(hosts, results):
        print(f"{host}: {result}")

    print(await reverse("8.8.8.8"))


if __name__ == "__main__":
    asyncio.run(main())

Per interrogare record diversi da A e AAAA, come MX, TXT o SRV, la libreria standard non basta: lo strumento di riferimento è il pacchetto dnspython, che offre anche un resolver asincrono nativo.

import dns.asyncresolver


async def mail_exchangers(domain: str) -> list[tuple[int, str]]:
    resolver = dns.asyncresolver.Resolver()
    resolver.nameservers = ["1.1.1.1", "8.8.8.8"]
    resolver.lifetime = 3.0

    answer = await resolver.resolve(domain, "MX")
    # Ordinamento per priorità: valore più basso = server preferito
    return sorted((record.preference, record.exchange.to_text()) for record in answer)

TLS con il modulo ssl

Il modulo ssl avvolge un socket esistente aggiungendo la cifratura. ssl.create_default_context() restituisce un contesto con impostazioni sicure: verifica del certificato, controllo del nome host e protocolli moderni.

import socket
import ssl
from datetime import datetime, timezone


def inspect_certificate(host: str, port: int = 443, timeout: float = 5.0) -> dict:
    context = ssl.create_default_context()
    context.minimum_version = ssl.TLSVersion.TLSv1_2

    with socket.create_connection((host, port), timeout=timeout) as raw_sock:
        # server_hostname abilita SNI e la verifica del nome nel certificato
        with context.wrap_socket(raw_sock, server_hostname=host) as tls_sock:
            cert = tls_sock.getpeercert()
            expires = datetime.fromtimestamp(
                ssl.cert_time_to_seconds(cert["notAfter"]), tz=timezone.utc
            )

            subject = dict(item[0] for item in cert["subject"])
            issuer = dict(item[0] for item in cert["issuer"])

            return {
                "host": host,
                "protocol": tls_sock.version(),
                "cipher": tls_sock.cipher()[0],
                "subject": subject.get("commonName"),
                "issuer": issuer.get("organizationName"),
                "expires": expires.isoformat(),
                "days_left": (expires - datetime.now(timezone.utc)).days,
                "san": [value for kind, value in cert.get("subjectAltName", ()) if kind == "DNS"],
            }


if __name__ == "__main__":
    for host in ["python.org", "expired.badssl.com"]:
        try:
            print(inspect_certificate(host))
        except ssl.SSLCertVerificationError as exc:
            print(f"{host}: certificato non valido ({exc.verify_message})")
        except OSError as exc:
            print(f"{host}: {exc}")

Lo stesso contesto può essere passato ad asyncio.open_connection() tramite il parametro ssl, ottenendo connessioni TLS asincrone senza modificare il resto del codice.

Buone pratiche

  • Impostare sempre un timeout: sui socket bloccanti con il parametro timeout, nel codice asincrono con asyncio.timeout().
  • Usare sendall() e una funzione recv_exactly(), mai assumere che una singola chiamata trasferisca tutti i dati.
  • Chiamare await writer.drain() dopo ogni scrittura negli stream asincroni, per rispettare la backpressure.
  • Non bloccare l'event loop: le funzioni sincrone come gethostbyaddr() vanno eseguite con run_in_executor() o asyncio.to_thread().
  • Usare ssl.create_default_context() e non disattivare mai la verifica dei certificati nel codice di produzione.

Conclusioni

Python permette di affrontare il networking a qualsiasi livello di astrazione: dalla manipolazione diretta dei socket alla scrittura di server asincroni capaci di gestire migliaia di connessioni. La libreria standard, con socket, selectors, asyncio e ssl, copre la quasi totalità delle esigenze, e i pochi vuoti, come le query DNS avanzate, sono colmati da pacchetti maturi. Come sempre nella programmazione di rete, la differenza tra un prototipo e un servizio affidabile sta nei timeout, nel framing e nella gestione accurata degli errori.