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 conasyncio.timeout(). - Usare
sendall()e una funzionerecv_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 conrun_in_executor()oasyncio.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.