Networking con Go
Go è stato progettato con il networking al centro: il pacchetto net della libreria standard, unito alle goroutine e al runtime che integra un poller di rete basato su epoll e kqueue, permette di scrivere server capaci di gestire decine di migliaia di connessioni con codice sequenziale e leggibile. In questo articolo vedremo come costruire server e client TCP, gestire deadline e cancellazione con context, lavorare con UDP e DNS, configurare correttamente il client HTTP e implementare uno shutdown ordinato.
Il modello di Go: codice bloccante, runtime non bloccante
In molti linguaggi l'I/O non bloccante richiede callback, promise o event loop espliciti. In Go si scrive codice apparentemente bloccante: una chiamata a conn.Read() sospende la goroutine corrente finché non arrivano dati. Dietro le quinte, però, il runtime registra il descrittore nel poller di rete e libera il thread del sistema operativo per eseguire altre goroutine. Il risultato è la semplicità del modello "una goroutine per connessione" con l'efficienza di un server event-driven.
Un server TCP con una goroutine per connessione
package main
import (
"bufio"
"fmt"
"log"
"net"
"strings"
"time"
)
const idleTimeout = 30 * time.Second
func main() {
listener, err := net.Listen("tcp", ":9000")
if err != nil {
log.Fatalf("impossibile avviare il listener: %v", err)
}
defer listener.Close()
log.Printf("server in ascolto su %s", listener.Addr())
for {
conn, err := listener.Accept()
if err != nil {
log.Printf("errore in accept: %v", err)
continue
}
// Ogni connessione viene gestita in una goroutine dedicata
go handleConnection(conn)
}
}
func handleConnection(conn net.Conn) {
defer conn.Close()
peer := conn.RemoteAddr().String()
log.Printf("nuova connessione da %s", peer)
scanner := bufio.NewScanner(conn)
// Limite massimo della singola riga, per difendersi da input malevoli
scanner.Buffer(make([]byte, 0, 4096), 64*1024)
for {
// La deadline viene rinnovata prima di ogni lettura
if err := conn.SetReadDeadline(time.Now().Add(idleTimeout)); err != nil {
return
}
if !scanner.Scan() {
if err := scanner.Err(); err != nil {
log.Printf("%s: %v", peer, err)
}
log.Printf("%s disconnesso", peer)
return
}
line := strings.TrimSpace(scanner.Text())
fmt.Fprintf(conn, "echo: %s\n", line)
}
}
In Go i timeout sui socket si esprimono come deadline, cioè istanti assoluti e non durate. Per realizzare un timeout di inattività la deadline va quindi spostata in avanti prima di ogni lettura. Quando scade, Read restituisce un errore che soddisfa os.ErrDeadlineExceeded, e lo scanner termina.
Framing binario con prefisso di lunghezza
Il protocollo a righe è comodo per il testo, ma per payload binari il prefisso di lunghezza è più adatto. Il pacchetto encoding/binary e la funzione io.ReadFull rendono il codec molto compatto:
package framing
import (
"encoding/binary"
"errors"
"fmt"
"io"
)
const (
headerSize = 4
maxFrameSize = 16 << 20 // 16 MiB
)
var ErrFrameTooLarge = errors.New("frame troppo grande")
// WriteFrame scrive la lunghezza in big-endian seguita dal payload
func WriteFrame(w io.Writer, payload []byte) error {
if len(payload) > maxFrameSize {
return ErrFrameTooLarge
}
var header [headerSize]byte
binary.BigEndian.PutUint32(header[:], uint32(len(payload)))
// Un'unica scrittura riduce il numero di segmenti TCP inviati
buffer := make([]byte, 0, headerSize+len(payload))
buffer = append(buffer, header[:]...)
buffer = append(buffer, payload...)
_, err := w.Write(buffer)
return err
}
// ReadFrame legge esattamente un frame completo dallo stream
func ReadFrame(r io.Reader) ([]byte, error) {
var header [headerSize]byte
// io.ReadFull continua a leggere finché il buffer non è pieno
if _, err := io.ReadFull(r, header[:]); err != nil {
return nil, err
}
length := binary.BigEndian.Uint32(header[:])
if length > maxFrameSize {
return nil, fmt.Errorf("%w: %d byte", ErrFrameTooLarge, length)
}
payload := make([]byte, length)
if _, err := io.ReadFull(r, payload); err != nil {
return nil, err
}
return payload, nil
}
io.ReadFull risolve in modo elegante il problema delle letture parziali: se la connessione si chiude a metà frame restituisce io.ErrUnexpectedEOF, distinguendo chiaramente un messaggio troncato da una chiusura pulita tra un messaggio e l'altro (io.EOF).
Un server con shutdown ordinato
Un server di produzione deve poter terminare senza troncare le richieste in corso, ad esempio quando un orchestratore invia SIGTERM. Combinando signal.NotifyContext, un sync.WaitGroup e la chiusura del listener si ottiene uno shutdown pulito:
package main
import (
"context"
"encoding/json"
"errors"
"io"
"log"
"net"
"os"
"os/signal"
"sync"
"syscall"
"time"
"example.com/netdemo/framing"
)
type Request struct {
ID int `json:"id"`
Command string `json:"command"`
Params json.RawMessage `json:"params,omitempty"`
}
type Response struct {
ID int `json:"id"`
Result any `json:"result,omitempty"`
Error string `json:"error,omitempty"`
}
type Server struct {
listener net.Listener
wg sync.WaitGroup
mu sync.Mutex
conns map[net.Conn]struct{}
}
func NewServer(address string) (*Server, error) {
listener, err := net.Listen("tcp", address)
if err != nil {
return nil, err
}
return &Server{listener: listener, conns: make(map[net.Conn]struct{})}, nil
}
func (s *Server) Serve(ctx context.Context) error {
// Alla cancellazione del contesto si chiude il listener per sbloccare Accept
go func() {
<-ctx.Done()
s.listener.Close()
}()
for {
conn, err := s.listener.Accept()
if err != nil {
if errors.Is(err, net.ErrClosed) {
return nil
}
log.Printf("errore in accept: %v", err)
continue
}
s.track(conn, true)
s.wg.Add(1)
go func() {
defer s.wg.Done()
defer s.track(conn, false)
s.handle(ctx, conn)
}()
}
}
func (s *Server) Shutdown(timeout time.Duration) {
done := make(chan struct{})
go func() {
s.wg.Wait()
close(done)
}()
select {
case <-done:
log.Println("tutte le connessioni sono state chiuse")
case <-time.After(timeout):
// Tempo scaduto: si chiudono forzatamente le connessioni residue
s.mu.Lock()
for conn := range s.conns {
conn.Close()
}
s.mu.Unlock()
log.Println("connessioni residue chiuse forzatamente")
}
}
func (s *Server) track(conn net.Conn, add bool) {
s.mu.Lock()
defer s.mu.Unlock()
if add {
s.conns[conn] = struct{}{}
} else {
delete(s.conns, conn)
conn.Close()
}
}
func (s *Server) handle(ctx context.Context, conn net.Conn) {
for {
// Durante lo shutdown non si accettano nuove richieste sulla connessione
if ctx.Err() != nil {
return
}
conn.SetReadDeadline(time.Now().Add(30 * time.Second))
payload, err := framing.ReadFrame(conn)
if err != nil {
if !errors.Is(err, io.EOF) && !errors.Is(err, os.ErrDeadlineExceeded) {
log.Printf("%s: %v", conn.RemoteAddr(), err)
}
return
}
var request Request
response := Response{}
if err := json.Unmarshal(payload, &request); err != nil {
response.Error = "JSON non valido"
} else {
response = s.dispatch(request)
}
encoded, _ := json.Marshal(response)
conn.SetWriteDeadline(time.Now().Add(5 * time.Second))
if err := framing.WriteFrame(conn, encoded); err != nil {
return
}
}
}
func (s *Server) dispatch(request Request) Response {
switch request.Command {
case "ping":
return Response{ID: request.ID, Result: "pong"}
case "time":
return Response{ID: request.ID, Result: time.Now().UTC().Format(time.RFC3339)}
default:
return Response{ID: request.ID, Error: "comando sconosciuto: " + request.Command}
}
}
func main() {
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
server, err := NewServer(":9000")
if err != nil {
log.Fatal(err)
}
log.Println("server in ascolto sulla porta 9000")
if err := server.Serve(ctx); err != nil {
log.Fatal(err)
}
log.Println("arresto in corso...")
server.Shutdown(10 * time.Second)
}
Il punto chiave è che Accept è una chiamata bloccante che non accetta un contesto: per interromperla si chiude il listener da un'altra goroutine, e l'errore risultante viene riconosciuto con errors.Is(err, net.ErrClosed).
Il client: dial con contesto
Lato client, net.Dialer permette di configurare timeout, keep-alive e cancellazione tramite DialContext. Per un host con più indirizzi IP, il dialer prova automaticamente gli indirizzi in sequenza e applica l'algoritmo Happy Eyeballs tra IPv4 e IPv6.
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"net"
"time"
"example.com/netdemo/framing"
)
type Client struct {
conn net.Conn
nextID int
}
func Dial(ctx context.Context, address string) (*Client, error) {
dialer := net.Dialer{
Timeout: 3 * time.Second,
KeepAlive: 30 * time.Second,
}
conn, err := dialer.DialContext(ctx, "tcp", address)
if err != nil {
return nil, fmt.Errorf("connessione a %s fallita: %w", address, err)
}
return &Client{conn: conn}, nil
}
func (c *Client) Call(ctx context.Context, command string) (json.RawMessage, error) {
c.nextID++
// La deadline della connessione segue quella del contesto, se presente
if deadline, ok := ctx.Deadline(); ok {
c.conn.SetDeadline(deadline)
} else {
c.conn.SetDeadline(time.Now().Add(5 * time.Second))
}
request, _ := json.Marshal(map[string]any{"id": c.nextID, "command": command})
if err := framing.WriteFrame(c.conn, request); err != nil {
return nil, err
}
payload, err := framing.ReadFrame(c.conn)
if err != nil {
return nil, err
}
var response struct {
Result json.RawMessage `json:"result"`
Error string `json:"error"`
}
if err := json.Unmarshal(payload, &response); err != nil {
return nil, err
}
if response.Error != "" {
return nil, fmt.Errorf("errore remoto: %s", response.Error)
}
return response.Result, nil
}
func (c *Client) Close() error {
return c.conn.Close()
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
client, err := Dial(ctx, "127.0.0.1:9000")
if err != nil {
log.Fatal(err)
}
defer client.Close()
for _, command := range []string{"ping", "time", "unknown"} {
result, err := client.Call(ctx, command)
if err != nil {
log.Printf("%s: %v", command, err)
continue
}
log.Printf("%s: %s", command, result)
}
}
Un port scanner concorrente con limite di parallelismo
Le goroutine rendono banale lanciare migliaia di connessioni in parallelo, ma farlo senza limiti esaurisce i descrittori di file e può essere scambiato per un attacco. Un canale bufferizzato usato come semaforo limita il parallelismo:
package main
import (
"context"
"fmt"
"net"
"sort"
"strconv"
"sync"
"time"
)
type PortResult struct {
Port int
Open bool
Latency time.Duration
}
func ScanPorts(ctx context.Context, host string, ports []int, concurrency int, timeout time.Duration) []PortResult {
results := make([]PortResult, 0, len(ports))
semaphore := make(chan struct{}, concurrency)
var mu sync.Mutex
var wg sync.WaitGroup
dialer := net.Dialer{Timeout: timeout}
for _, port := range ports {
wg.Add(1)
go func(port int) {
defer wg.Done()
// Acquisizione di uno slot del semaforo
select {
case semaphore <- struct{}{}:
case <-ctx.Done():
return
}
defer func() { <-semaphore }()
address := net.JoinHostPort(host, strconv.Itoa(port))
start := time.Now()
conn, err := dialer.DialContext(ctx, "tcp", address)
result := PortResult{Port: port, Latency: time.Since(start)}
if err == nil {
result.Open = true
conn.Close()
}
mu.Lock()
results = append(results, result)
mu.Unlock()
}(port)
}
wg.Wait()
sort.Slice(results, func(i, j int) bool { return results[i].Port < results[j].Port })
return results
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
ports := make([]int, 0, 1024)
for port := 1; port <= 1024; port++ {
ports = append(ports, port)
}
for _, result := range ScanPorts(ctx, "127.0.0.1", ports, 100, 500*time.Millisecond) {
if result.Open {
fmt.Printf("%5d aperta (%v)\n", result.Port, result.Latency.Round(time.Microsecond))
}
}
}
Da notare l'uso di net.JoinHostPort invece di una semplice concatenazione: è l'unico modo corretto di costruire un indirizzo che funzioni anche con IPv6, dove l'host va racchiuso tra parentesi quadre.
UDP
Con UDP si usa net.ListenUDP, che restituisce un *net.UDPConn. Ogni chiamata a ReadFromUDP legge un datagramma completo insieme all'indirizzo del mittente. Ecco un server di metriche in stile StatsD che riceve contatori e li aggrega:
package main
import (
"log"
"net"
"strconv"
"strings"
"sync"
"time"
)
type Aggregator struct {
mu sync.Mutex
counters map[string]int64
}
func (a *Aggregator) Add(name string, value int64) {
a.mu.Lock()
a.counters[name] += value
a.mu.Unlock()
}
// Flush restituisce i contatori correnti e li azzera
func (a *Aggregator) Flush() map[string]int64 {
a.mu.Lock()
defer a.mu.Unlock()
snapshot := a.counters
a.counters = make(map[string]int64)
return snapshot
}
func main() {
address, _ := net.ResolveUDPAddr("udp", ":8125")
conn, err := net.ListenUDP("udp", address)
if err != nil {
log.Fatal(err)
}
defer conn.Close()
aggregator := &Aggregator{counters: make(map[string]int64)}
// Stampa periodica dei valori aggregati
go func() {
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for range ticker.C {
for name, value := range aggregator.Flush() {
log.Printf("%s = %d", name, value)
}
}
}()
buffer := make([]byte, 65535)
for {
n, _, err := conn.ReadFromUDP(buffer)
if err != nil {
log.Printf("errore in lettura: %v", err)
continue
}
// Un datagramma può contenere più metriche separate da newline
for _, line := range strings.Split(string(buffer[:n]), "\n") {
// Formato atteso: nome:valore|c
name, rest, found := strings.Cut(line, ":")
if !found {
continue
}
valueText, metricType, _ := strings.Cut(rest, "|")
if metricType != "c" {
continue
}
value, err := strconv.ParseInt(valueText, 10, 64)
if err != nil {
continue
}
aggregator.Add(name, value)
}
}
}
Si può provare il server inviando datagrammi con nc:
echo -n "api.requests:1|c" | nc -u -w0 127.0.0.1 8125
Risoluzione DNS con net.Resolver
net.Resolver espone metodi per ogni tipo di record e, come tutto il pacchetto net, accetta un contesto per timeout e cancellazione. Impostando PreferGo si usa il resolver scritto in Go invece di quello di sistema basato su cgo, e con un Dial personalizzato si può puntare a un server DNS specifico:
package main
import (
"context"
"fmt"
"log"
"net"
"time"
)
func newResolver(server string) *net.Resolver {
return &net.Resolver{
PreferGo: true,
// Tutte le query vengono inviate al server DNS indicato
Dial: func(ctx context.Context, network, address string) (net.Conn, error) {
dialer := net.Dialer{Timeout: 2 * time.Second}
return dialer.DialContext(ctx, network, server)
},
}
}
func main() {
resolver := newResolver("1.1.1.1:53")
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
domain := "go.dev"
addresses, err := resolver.LookupIPAddr(ctx, domain)
if err != nil {
log.Fatalf("lookup fallito: %v", err)
}
for _, address := range addresses {
fmt.Println("IP:", address.IP)
}
if records, err := resolver.LookupMX(ctx, "gmail.com"); err == nil {
for _, mx := range records {
fmt.Printf("MX: %s (priorità %d)\n", mx.Host, mx.Pref)
}
}
if records, err := resolver.LookupTXT(ctx, domain); err == nil {
for _, txt := range records {
fmt.Println("TXT:", txt)
}
}
if names, err := resolver.LookupAddr(ctx, "8.8.8.8"); err == nil {
fmt.Println("PTR:", names)
}
// Gli errori DNS possono essere ispezionati in dettaglio
_, err = resolver.LookupHost(ctx, "dominio-inesistente.invalid")
if dnsErr, ok := err.(*net.DNSError); ok {
fmt.Printf("errore DNS: not found=%v, timeout=%v\n", dnsErr.IsNotFound, dnsErr.IsTimeout)
}
}
Il client HTTP: mai usare quello predefinito in produzione
http.DefaultClient non ha alcun timeout: una richiesta verso un server che smette di rispondere può bloccare una goroutine per sempre. Un client di produzione va configurato esplicitamente, sia a livello complessivo sia a livello di Transport:
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"net"
"net/http"
"time"
)
func newHTTPClient() *http.Client {
transport := &http.Transport{
Proxy: http.ProxyFromEnvironment,
DialContext: (&net.Dialer{
Timeout: 3 * time.Second,
KeepAlive: 30 * time.Second,
}).DialContext,
TLSHandshakeTimeout: 5 * time.Second,
ResponseHeaderTimeout: 10 * time.Second,
ExpectContinueTimeout: 1 * time.Second,
// Pool di connessioni riutilizzabili verso lo stesso host
MaxIdleConns: 100,
MaxIdleConnsPerHost: 20,
IdleConnTimeout: 90 * time.Second,
ForceAttemptHTTP2: true,
}
return &http.Client{
Transport: transport,
// Limite complessivo, inclusa la lettura del corpo
Timeout: 30 * time.Second,
}
}
type Release struct {
Version string `json:"version"`
Stable bool `json:"stable"`
}
func fetchReleases(ctx context.Context, client *http.Client) ([]Release, error) {
request, err := http.NewRequestWithContext(ctx, http.MethodGet, "https://go.dev/dl/?mode=json", nil)
if err != nil {
return nil, err
}
request.Header.Set("Accept", "application/json")
response, err := client.Do(request)
if err != nil {
return nil, err
}
// Il corpo va sempre chiuso per restituire la connessione al pool
defer response.Body.Close()
if response.StatusCode != http.StatusOK {
return nil, fmt.Errorf("status inatteso: %s", response.Status)
}
var releases []Release
if err := json.NewDecoder(response.Body).Decode(&releases); err != nil {
return nil, err
}
return releases, nil
}
func main() {
client := newHTTPClient()
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
releases, err := fetchReleases(ctx, client)
if err != nil {
log.Fatal(err)
}
for _, release := range releases {
fmt.Printf("%s (stabile: %v)\n", release.Version, release.Stable)
}
}
Il valore predefinito di MaxIdleConnsPerHost è 2: un servizio che effettua molte richieste concorrenti verso la stessa API apre e chiude continuamente connessioni, pagando ogni volta l'handshake TCP e TLS. Aumentarlo è una delle ottimizzazioni più efficaci e meno conosciute.
Buone pratiche
- Usare deadline su ogni connessione e rinnovarle prima di ogni operazione per realizzare timeout di inattività.
- Propagare il
contextin tutte le funzioni di rete, usandoDialContext,NewRequestWithContexte i metodi delResolver. - Non usare mai
http.DefaultClientnéhttp.Getin codice di produzione. - Chiudere sempre il corpo delle risposte HTTP, anche quando non lo si legge.
- Limitare il parallelismo con semafori o worker pool: le goroutine sono economiche, i descrittori di file no.
- Costruire gli indirizzi con
net.JoinHostPortper supportare IPv6 senza sorprese.
Conclusioni
Go offre probabilmente l'esperienza più lineare per la programmazione di rete tra i linguaggi mainstream: il codice resta sequenziale, la concorrenza è economica e la libreria standard copre TCP, UDP, DNS, TLS e HTTP con API coerenti e componibili. La disciplina richiesta riguarda soprattutto i dettagli: deadline, contesti, limiti al parallelismo e configurazione esplicita dei client. Curati questi aspetti, si ottengono servizi di rete che scalano con naturalezza.