PidokuInfra

Networking and the Netpoller

Advanced 55 min Difficulty 3/5 Topic 05 of 05

Prerequisites IV.01, IV.03, III.05

The idea in one minute#

In Go you write network code that looks blocking — conn.Read(buf) simply waits — and it is cheap anyway. Under the surface every socket is non-blocking and registered with the operating system’s event notification system (epoll, kqueue, IOCP). When a read would block, the runtime parks the goroutine and the thread runs something else; when data arrives, the netpoller makes the goroutine runnable again.

That gives you one goroutine per connection with the simplicity of threads and the scalability of an event loop. net/http is built on it, and so is the thing an AI service needs most: a long-lived response that streams tokens as they are produced.

An analogy#

A restaurant with one buzzer per table. A waiter takes an order and, instead of standing at the table until the food is ready, moves on to other tables. When a dish is up, that table’s buzzer goes and whichever waiter is free handles it. Each table still feels like it has a dedicated waiter; the restaurant employs six.

A picture#

flowchart TB
  C1["client"] <--> S1["socket fd 7<br/>non-blocking"]
  C2["client"] <--> S2["socket fd 8"]
  S1 --> EP["OS event queue<br/>epoll / kqueue / IOCP"]
  S2 --> EP
  EP --> NP["netpoller<br/>checked by the scheduler and sysmon"]
  NP -->|"fd 7 readable: make G12 runnable"| RQ["run queue"]
  RQ --> G12["G12: handler for fd 7<br/>conn.Read returns"]
  G12 -->|"Read would block: park G12, register interest"| NP
  class C1,C2 neutral
  class S1,S2,EP io
  class NP,RQ queue
  class G12 compute

How it really works#

One goroutine per connection#

Go
ln, _ := net.Listen("tcp", ":8080")
for {
    conn, err := ln.Accept()
    if err != nil { break }
    go handle(conn)             // blocking-style code inside, parked when idle
}

An idle connection costs one parked goroutine — a few kilobytes — and no thread. Deadlines (conn.SetReadDeadline) are implemented with runtime timers that wake the goroutine with a timeout error.

File I/O is different: ordinary files are not pollable on most systems, so a blocking file read occupies a thread (IV.01).

net/http server#

Go
mux := http.NewServeMux()
mux.HandleFunc("POST /v1/embed", embedHandler)        // method and path patterns (Go 1.22)
mux.HandleFunc("GET /models/{name}", modelHandler)    // r.PathValue("name")

srv := &http.Server{
    Addr:              ":8080",
    Handler:           mux,
    ReadHeaderTimeout: 5 * time.Second,
    ReadTimeout:       30 * time.Second,
    IdleTimeout:       120 * time.Second,
    // WriteTimeout: leave unset for streaming endpoints; bound them with the request context
}
  • Each connection gets a goroutine; each request runs its handler on it. HTTP/2 multiplexes many requests over one connection, each on its own goroutine.
  • Set timeouts. The zero-value server has none, so a slow or malicious client can hold connections forever.
  • The handler’s r.Context() is cancelled when the client disconnects (IV.03).
  • A panic in a handler is recovered by the server and kills only that request.
  • srv.Shutdown(ctx) stops accepting and waits for in-flight requests: graceful shutdown.
  • Bound work explicitly. The server will happily start a goroutine for every connection it is sent; a semaphore or a bounded queue in front of expensive work is your admission control (IV.06).

net/http client#

  • Reuse one http.Client. Its Transport keeps a pool of connections; a client per request means a TCP and TLS handshake per request.
  • Always close and drain the response body, or the connection cannot be reused: defer resp.Body.Close() and read to EOF (or io.Copy(io.Discard, resp.Body)).
  • Set a timeout: Client.Timeout for simple calls; for streaming, use a context deadline for the whole call plus your own per-chunk idle timer, because Client.Timeout would cut a long stream off.
  • Tune Transport.MaxIdleConnsPerHost (default 2) when you send many concurrent requests to one host, such as a model server.

Streaming responses#

LLM APIs stream with Server-Sent Events (SSE): a normal HTTP response with Content-Type: text/event-stream whose body is a sequence of data: ... lines separated by blank lines, flushed one event at a time.

Go
func stream(w http.ResponseWriter, r *http.Request) {
    w.Header().Set("Content-Type", "text/event-stream")
    w.Header().Set("Cache-Control", "no-cache")
    rc := http.NewResponseController(w)
    for tok := range tokens(r.Context()) {
        fmt.Fprintf(w, "data: %s\n\n", tok)
        if err := rc.Flush(); err != nil {       // push it to the client now
            return                               // client went away
        }
    }
    fmt.Fprint(w, "data: [DONE]\n\n")
}

What matters:

  • Flush after every event. The response writer buffers; without a flush the client receives nothing until the buffer fills or the handler returns. Proxies buffer too — disable it for streaming routes.
  • Watch the request context between tokens so that a disconnected client stops generation.
  • Backpressure is built in: Write blocks when the client’s TCP window is full, so a slow reader slows the handler. Decide whether a slow client should hold a GPU slot; usually the answer is a per-write deadline (rc.SetWriteDeadline).
  • On the client, read with a bufio.Scanner or bufio.Reader line by line, and raise the scanner’s buffer limit — a single event can exceed the default 64 KB.

WebSockets (bidirectional) and gRPC streaming (HTTP/2, typed messages) follow the same goroutine-per-stream model.

Reducing per-request cost#

CostRemedy
JSON encode/decode with reflectionencoding/json/v2 (Go 1.27) is substantially faster at decoding; decode straight from the body with a Decoder; avoid map[string]any
A new buffer per requestsync.Pool of buffers (III.05)
fmt.Fprintf per tokenstrconv.Append* into a reused []byte, one Write
Many tiny writesA bufio.Writer, flushed per event
A handshake per requestKeep-alive and HTTP/2; one shared client
Copying between files and socketsio.Copy uses sendfile/splice where it can

Observing#

Export request rate, errors and duration per route, the number of in-flight requests and the queue depth in front of your workers — the RED metrics from Observability II.02. net/http/pprof and net/http/httptrace (DNS, connect, TLS and first-byte timings on the client) cover the rest.

Code#

A streaming server and client in one process, using an in-memory listener from httptest.

Go
// stream.go — an SSE token stream: flush per event, stop when the client leaves.
package main

import (
	"bufio"
	"context"
	"fmt"
	"net/http"
	"net/http/httptest"
	"strings"
	"sync/atomic"
	"time"
)

var generated atomic.Int64 // tokens the "model" actually produced

func handler(w http.ResponseWriter, r *http.Request) {
	w.Header().Set("Content-Type", "text/event-stream")
	w.Header().Set("Cache-Control", "no-cache")
	rc := http.NewResponseController(w)
	for i := 0; i < 50; i++ {
		select {
		case <-time.After(5 * time.Millisecond): // one decode step
		case <-r.Context().Done(): // the client disconnected: stop spending compute
			return
		}
		generated.Add(1)
		fmt.Fprintf(w, "data: tok%d\n\n", i)
		if err := rc.Flush(); err != nil {
			return
		}
	}
	fmt.Fprint(w, "data: [DONE]\n\n")
}

// consume reads events until done or until it has seen `limit` tokens, then hangs up.
func consume(client *http.Client, url string, limit int) (n int, ttft, total time.Duration, err error) {
	ctx, cancel := context.WithCancel(context.Background())
	defer cancel()
	req, _ := http.NewRequestWithContext(ctx, "GET", url, nil)
	start := time.Now()
	resp, err := client.Do(req)
	if err != nil {
		return 0, 0, 0, err
	}
	defer resp.Body.Close()
	sc := bufio.NewScanner(resp.Body)
	for sc.Scan() {
		line := sc.Text()
		data, ok := strings.CutPrefix(line, "data: ")
		if !ok {
			continue // blank separator lines
		}
		if data == "[DONE]" {
			break
		}
		if n == 0 {
			ttft = time.Since(start)
		}
		n++
		if n == limit {
			cancel() // walk away mid-stream
			break
		}
	}
	return n, ttft, time.Since(start), nil
}

func main() {
	srv := httptest.NewServer(http.HandlerFunc(handler))
	defer srv.Close()
	client := srv.Client()

	n, ttft, total, err := consume(client, srv.URL, -1)
	fmt.Printf("full stream:   %2d tokens, first after %3d ms, all after %3d ms, err=%v\n",
		n, ttft.Milliseconds(), total.Milliseconds(), err)
	fmt.Println("               the first token arrived long before the last: flushing works")

	generated.Store(0)
	n, _, total, _ = consume(client, srv.URL, 10)
	time.Sleep(60 * time.Millisecond) // give the server time to notice
	fmt.Printf("client leaves: read %d tokens in %d ms; server generated %d of 50 and stopped\n",
		n, total.Milliseconds(), generated.Load())
}

Remember this#

  • Blocking-style network code is cheap: the netpoller parks goroutines instead of threads.
  • Set server timeouts; reuse one client; always close and drain response bodies.
  • Streaming = write an event, flush, check the request context, repeat.
  • A slow client applies backpressure through Write; decide how long it may hold resources.

Try it#

  1. Run stream.go. Remove the rc.Flush() call. What happens to time-to-first-token?
  2. Remove the r.Context().Done() case. How many tokens does the server generate after the client leaves?
  3. Add a per-write deadline of 50 ms and a client that sleeps 200 ms between reads. What does the server do?

Check yourself#

  1. What happens to a goroutine when conn.Read has no data?
  2. Why must a streaming handler flush after each event?
  3. Why should an http.Client be shared rather than created per request?

↑↓ navigate↵ openesc close