mirror of
https://github.com/luxfi/zap.git
synced 2026-07-26 22:53:31 +00:00
WIP: session checkpoint 2026-06-05 (codec rip in flight)
This commit is contained in:
@@ -152,21 +152,20 @@ by default).
|
||||
|
||||
Wire format on each per-Call stream is one length-prefixed ZAP frame
|
||||
in each direction — no correlation header needed since each stream
|
||||
carries exactly one request + one response. Byte-identical with the
|
||||
control-stream Call format minus the 8-byte `(reqID, flag)` preamble.
|
||||
carries exactly one request + one response.
|
||||
|
||||
`ZAP_DISABLE_CALL_STREAMS=1` reverts to the legacy serialized path
|
||||
for A/B benchmarking. Production must leave it unset.
|
||||
Per-stream is the only Call path on QUIC; there is no fallback. The
|
||||
control stream carries one-way Sends only.
|
||||
|
||||
End-to-end Call throughput (2 ms simulated handler):
|
||||
|
||||
| Concurrency | Per-stream | Serialized | Speedup |
|
||||
| Concurrency | Per-stream | Reference | Speedup |
|
||||
|-------------|-------------|-------------|---------|
|
||||
| 1 worker | 3.19 ms/op | 3.17 ms/op | 1.0x |
|
||||
| 10 workers | 1.24 ms/op | 3.01 ms/op | 2.4x |
|
||||
| 100 workers | 1.06 ms/op | 3.01 ms/op | 2.85x |
|
||||
|
||||
The serialized path plateaus at the handler latency since every Call
|
||||
waits for the prior one to clear the control stream. Per-stream
|
||||
Reference = handler latency on a serialized control-stream path
|
||||
(pre-optimization measurement, retained for context). Per-stream
|
||||
overlaps handlers in flight (peak inflight = 32 observed in
|
||||
`TestQUICCallConcurrent`).
|
||||
|
||||
@@ -36,13 +36,11 @@ type Node struct {
|
||||
discovery *mdns.Discovery
|
||||
|
||||
// Network
|
||||
listener net.Listener
|
||||
transports map[string]TransportConn // peerID -> transport conn (QUIC path)
|
||||
transClose func() error // closer for the QUIC listener
|
||||
reqIDQuic uint32 // QUIC-path request-ID counter
|
||||
reqIDQuicMu sync.Mutex
|
||||
conns map[string]*Conn
|
||||
connsMu sync.RWMutex
|
||||
listener net.Listener
|
||||
transports map[string]TransportConn // peerID -> transport conn (QUIC path)
|
||||
transClose func() error // closer for the QUIC listener
|
||||
conns map[string]*Conn
|
||||
connsMu sync.RWMutex
|
||||
|
||||
// Handlers
|
||||
handlers map[uint16]Handler
|
||||
|
||||
+37
-180
@@ -5,25 +5,15 @@ package zap
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/luxfi/mdns"
|
||||
)
|
||||
|
||||
// disableCallStreams forces Node.Call onto the legacy serialized
|
||||
// control-stream path even when the transport supports per-Call
|
||||
// streams. Set by the ZAP_DISABLE_CALL_STREAMS environment variable
|
||||
// — used ONLY for A/B benchmarking. Production deployments must
|
||||
// leave it unset.
|
||||
var disableCallStreams = os.Getenv("ZAP_DISABLE_CALL_STREAMS") == "1"
|
||||
|
||||
// startQUIC bootstraps the QUIC transport path. It is the
|
||||
// Transport == TransportQUIC fork of (*Node).Start.
|
||||
//
|
||||
@@ -99,22 +89,22 @@ func (n *Node) onAcceptedTransportConn(peerID string, tc TransportConn) {
|
||||
go n.serveTransportConn(peerID, tc)
|
||||
}
|
||||
|
||||
// serveTransportConn is the QUIC equivalent of handleConn — it
|
||||
// drives the read loops and dispatches incoming frames to handlers.
|
||||
// serveTransportConn drives the read loops on a QUIC peer connection.
|
||||
//
|
||||
// Two parallel read loops run on each connection:
|
||||
//
|
||||
// 1. Control-stream loop (tc.Recv): handles one-way Sends and the
|
||||
// legacy correlated-frame Call path for transports that do not
|
||||
// implement TransportStreamer.
|
||||
//
|
||||
// 2. Per-Call stream accept loop (AcceptCallStream): handles the
|
||||
// stream-per-Call optimized path. Each accepted stream carries
|
||||
// exactly one request + one response, so the dispatch goroutine
|
||||
// reads, handles, writes, and closes the stream in one shot.
|
||||
// 1. Control-stream loop (tc.Recv): handles one-way Sends only.
|
||||
// 2. Per-Call stream accept loop (AcceptCallStream): handles every
|
||||
// Call. Each accepted stream carries exactly one request + one
|
||||
// response, so the dispatch goroutine reads, handles, writes,
|
||||
// and closes the stream in one shot.
|
||||
//
|
||||
// The two loops are independent: each runs in its own goroutine, and
|
||||
// both terminate when the connection closes or n.ctx is cancelled.
|
||||
//
|
||||
// The QUIC transport always implements TransportStreamer; if a
|
||||
// connection somehow lacks it we close it — the per-stream Call path
|
||||
// is the only Call path.
|
||||
func (n *Node) serveTransportConn(peerID string, tc TransportConn) {
|
||||
defer n.wg.Done()
|
||||
defer func() {
|
||||
@@ -127,15 +117,13 @@ func (n *Node) serveTransportConn(peerID string, tc TransportConn) {
|
||||
n.logger.Info("Peer disconnected (QUIC)", "peerID", peerID)
|
||||
}()
|
||||
|
||||
// Spawn the per-Call stream accept loop when the transport
|
||||
// supports it. We let this goroutine outlive serveTransportConn
|
||||
// only briefly — AcceptCallStream returns once the underlying
|
||||
// connection is closed, which the deferred tc.Close above
|
||||
// triggers.
|
||||
if s, ok := tc.(TransportStreamer); ok {
|
||||
n.wg.Add(1)
|
||||
go n.acceptCallStreams(peerID, s)
|
||||
streamer, ok := tc.(TransportStreamer)
|
||||
if !ok {
|
||||
n.logger.Error("QUIC TransportConn missing TransportStreamer", "peerID", peerID)
|
||||
return
|
||||
}
|
||||
n.wg.Add(1)
|
||||
go n.acceptCallStreams(peerID, streamer)
|
||||
|
||||
for {
|
||||
select {
|
||||
@@ -155,14 +143,20 @@ func (n *Node) serveTransportConn(peerID string, tc TransportConn) {
|
||||
n.logger.Debug("QUIC Recv error", "peerID", peerID, "error", err)
|
||||
return
|
||||
}
|
||||
n.dispatchFrame(peerID, data, tc)
|
||||
// Control stream carries one-way Sends only. Calls ride
|
||||
// per-Call streams.
|
||||
msg, err := Parse(data)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
n.invokeHandlerOneWay(peerID, msg)
|
||||
}
|
||||
}
|
||||
|
||||
// acceptCallStreams runs in its own goroutine for each QUIC peer and
|
||||
// processes inbound per-Call streams. Each stream gets its own handler
|
||||
// goroutine so multiple in-flight Calls from the same peer execute
|
||||
// concurrently — the headline parallelism win of the optimization.
|
||||
// concurrently — the headline parallelism win.
|
||||
//
|
||||
// Termination: AcceptCallStream returns a non-nil err when the
|
||||
// connection is torn down. The loop exits without logging — the
|
||||
@@ -223,39 +217,6 @@ func (n *Node) handleCallStream(peerID string, stream TransportStream) {
|
||||
}
|
||||
}
|
||||
|
||||
// dispatchFrame routes a single ZAP frame to either the pending-Call
|
||||
// channel (if it has a response correlation header) or to the
|
||||
// registered handler. Mirrors the TCP path in node.go's handleConn.
|
||||
func (n *Node) dispatchFrame(peerID string, data []byte, tc TransportConn) {
|
||||
if len(data) >= 8 {
|
||||
reqFlag := binary.LittleEndian.Uint32(data[4:8])
|
||||
if reqFlag == ReqFlagResp {
|
||||
reqID := binary.LittleEndian.Uint32(data[0:4])
|
||||
msg, err := Parse(data[8:])
|
||||
if err == nil {
|
||||
n.routeQUICResponse(peerID, reqID, msg)
|
||||
}
|
||||
return
|
||||
}
|
||||
if reqFlag == ReqFlagReq {
|
||||
reqID := binary.LittleEndian.Uint32(data[0:4])
|
||||
msg, err := Parse(data[8:])
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
n.handleQUICRequest(peerID, tc, reqID, msg)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Regular one-way message.
|
||||
msg, err := Parse(data)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
n.invokeHandlerOneWay(peerID, msg)
|
||||
}
|
||||
|
||||
// invokeHandlerOneWay runs the registered handler and discards the
|
||||
// response (matches TCP one-way semantics).
|
||||
func (n *Node) invokeHandlerOneWay(peerID string, msg *Message) {
|
||||
@@ -269,118 +230,26 @@ func (n *Node) invokeHandlerOneWay(peerID string, msg *Message) {
|
||||
_, _ = handler(n.ctx, peerID, msg)
|
||||
}
|
||||
|
||||
// handleQUICRequest dispatches a Call request and writes back the
|
||||
// correlation-tagged response on the same transport conn.
|
||||
func (n *Node) handleQUICRequest(peerID string, tc TransportConn, reqID uint32, msg *Message) {
|
||||
msgType := msg.Flags() >> 8
|
||||
n.handlersMu.RLock()
|
||||
handler, ok := n.handlers[msgType]
|
||||
n.handlersMu.RUnlock()
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
resp, err := handler(n.ctx, peerID, msg)
|
||||
if err != nil {
|
||||
n.logger.Error("Handler error", "peerID", peerID, "msgType", msgType, "error", err)
|
||||
return
|
||||
}
|
||||
if resp == nil {
|
||||
return
|
||||
}
|
||||
respBytes := resp.Bytes()
|
||||
// TransportConn.Send takes one slice; we still allocate the
|
||||
// 8-byte-prefixed buffer once, but reuse a per-message slab when
|
||||
// the body fits a common size class. Anything ≥ 64 KiB falls
|
||||
// through to a one-off alloc (rare on the response path; matches
|
||||
// the read-side 10 MiB ceiling).
|
||||
wrapped := makeCorrelatedFrame(reqID, ReqFlagResp, respBytes)
|
||||
if err := tc.Send(wrapped); err != nil {
|
||||
n.logger.Debug("QUIC Send error", "peerID", peerID, "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
// quicPendingMu / quicPending hold the per-peer pending Call response
|
||||
// channels for the QUIC transport. We keep them on the Node (not on
|
||||
// each TransportConn) because TransportConn is an interface and we
|
||||
// don't want to require implementors to carry the pending map.
|
||||
var (
|
||||
quicPendingMu sync.Mutex
|
||||
quicPending = map[string]map[uint32]chan *Message{}
|
||||
)
|
||||
|
||||
func (n *Node) routeQUICResponse(peerID string, reqID uint32, msg *Message) {
|
||||
quicPendingMu.Lock()
|
||||
defer quicPendingMu.Unlock()
|
||||
if peerMap, ok := quicPending[n.nodeID+":"+peerID]; ok {
|
||||
if ch, ok := peerMap[reqID]; ok {
|
||||
select {
|
||||
case ch <- msg:
|
||||
default:
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// quicCall is the QUIC path for Node.Call.
|
||||
//
|
||||
// When the transport implements TransportStreamer (QUIC always does),
|
||||
// every Call gets its own fresh bidirectional stream. The request goes
|
||||
// out on the new stream, the response comes back on the same stream,
|
||||
// and the stream is torn down. This is the headline parallel-Call
|
||||
// optimization — concurrent Calls no longer serialize on the shared
|
||||
// control stream's write mutex (ctrlMu), so throughput scales linearly
|
||||
// with goroutine count up to the peer-advertised stream limit.
|
||||
// Every Call gets its own fresh bidirectional QUIC stream. The
|
||||
// request goes out on the new stream, the response comes back on the
|
||||
// same stream, and the stream is torn down. Concurrent Calls run in
|
||||
// parallel up to the peer-advertised stream limit (1024 for ZAP's
|
||||
// QUIC config) — no shared write mutex, no correlation header.
|
||||
//
|
||||
// The legacy control-stream-with-correlation-header path is kept as a
|
||||
// fallback when the transport does NOT implement TransportStreamer
|
||||
// (notional non-QUIC transports). Today the QUIC transport is the only
|
||||
// streaming transport, so the fallback is effectively dead code — but
|
||||
// we leave it so a future transport without per-stream multiplexing can
|
||||
// still be plugged in without breaking the Node.Call API.
|
||||
//
|
||||
// Setting ZAP_DISABLE_CALL_STREAMS=1 in the environment forces every
|
||||
// Call onto the legacy serialized path — useful ONLY for A/B
|
||||
// benchmarking the optimization. Never set this in production.
|
||||
// QUIC's TransportConn always implements TransportStreamer; if it
|
||||
// somehow does not, that is a programmer error and we surface it.
|
||||
func (n *Node) quicCall(ctx context.Context, peerID string, msg *Message) (*Message, error) {
|
||||
tc, err := n.getOrConnectQUIC(ctx, peerID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if streamer, ok := tc.(TransportStreamer); ok && !disableCallStreams {
|
||||
return n.quicCallStream(ctx, streamer, msg)
|
||||
}
|
||||
|
||||
// Fallback: serialized correlated-frame Call on the control stream.
|
||||
reqID := nextReqID(n)
|
||||
|
||||
respCh := make(chan *Message, 1)
|
||||
key := n.nodeID + ":" + peerID
|
||||
quicPendingMu.Lock()
|
||||
if quicPending[key] == nil {
|
||||
quicPending[key] = make(map[uint32]chan *Message)
|
||||
}
|
||||
quicPending[key][reqID] = respCh
|
||||
quicPendingMu.Unlock()
|
||||
|
||||
defer func() {
|
||||
quicPendingMu.Lock()
|
||||
delete(quicPending[key], reqID)
|
||||
quicPendingMu.Unlock()
|
||||
}()
|
||||
|
||||
wrapped := makeCorrelatedFrame(reqID, ReqFlagReq, msg.Bytes())
|
||||
|
||||
if err := tc.Send(wrapped); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
select {
|
||||
case resp := <-respCh:
|
||||
return resp, nil
|
||||
case <-ctx.Done():
|
||||
return nil, ctx.Err()
|
||||
streamer, ok := tc.(TransportStreamer)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("zap: QUIC TransportConn for %s missing TransportStreamer", peerID)
|
||||
}
|
||||
return n.quicCallStream(ctx, streamer, msg)
|
||||
}
|
||||
|
||||
// quicCallStream runs one Call on a fresh per-Call QUIC stream.
|
||||
@@ -392,10 +261,8 @@ func (n *Node) quicCall(ctx context.Context, peerID string, msg *Message) (*Mess
|
||||
// 4. Close — stream ID returns to the QUIC pool.
|
||||
//
|
||||
// No correlation header is needed on the wire: each stream carries
|
||||
// exactly one request + one response, so the (reqID, flag) preamble
|
||||
// from the legacy control-stream layout is implicit in the stream
|
||||
// itself. The response payload is the raw ZAP frame; we Parse it
|
||||
// directly.
|
||||
// exactly one request + one response. The response payload is the raw
|
||||
// ZAP frame; we Parse it directly.
|
||||
//
|
||||
// Each stream gets cancelled if ctx fires before the response
|
||||
// arrives; the QUIC layer will surface ErrStreamCanceled into
|
||||
@@ -513,13 +380,3 @@ func (n *Node) getOrConnectQUIC(ctx context.Context, peerID string) (TransportCo
|
||||
return newTC, nil
|
||||
}
|
||||
|
||||
// nextReqID returns the next request ID for a QUIC peer. We piggy-
|
||||
// back on the existing Conn.reqID counter pattern but use a per-node
|
||||
// global because the TransportConn interface is opaque about its own
|
||||
// state. Atomicity is provided by the Node struct's reqIDMu.
|
||||
func nextReqID(n *Node) uint32 {
|
||||
n.reqIDQuicMu.Lock()
|
||||
defer n.reqIDQuicMu.Unlock()
|
||||
n.reqIDQuic++
|
||||
return n.reqIDQuic
|
||||
}
|
||||
|
||||
@@ -56,38 +56,6 @@ func TestQUICCallRoundTrip(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestQUICCallFallback exercises the ZAP_DISABLE_CALL_STREAMS env var
|
||||
// path. This is the legacy serialized-on-control-stream path that we
|
||||
// kept as a fallback for transports without TransportStreamer support.
|
||||
// We don't toggle the env in-process (the var is read at init time);
|
||||
// we just confirm the round-trip works through the same Node API.
|
||||
//
|
||||
// Real fallback testing happens in subprocess CI runs with the env
|
||||
// var set.
|
||||
func TestQUICCallFallback(t *testing.T) {
|
||||
// This test is identical to RoundTrip; the fallback path is
|
||||
// only chosen at init() so we exercise it via the env-var
|
||||
// process boundary. Here we just confirm the API is stable.
|
||||
srvAddr, srvNode, cliNode := newQUICTestPair(t, "srv-fb", "cli-fb")
|
||||
defer srvNode.Stop()
|
||||
defer cliNode.Stop()
|
||||
srvNode.Handle(0x42, func(ctx context.Context, from string, msg *zap.Message) (*zap.Message, error) {
|
||||
return buildEchoMessage(msg.Root().Uint64(0) + 7), nil
|
||||
})
|
||||
if err := cliNode.ConnectDirect(srvAddr); err != nil {
|
||||
t.Fatalf("ConnectDirect: %v", err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
resp, err := cliNode.Call(ctx, "srv-fb", buildEchoMessage(42))
|
||||
if err != nil {
|
||||
t.Fatalf("Call: %v", err)
|
||||
}
|
||||
if got := resp.Root().Uint64(0); got != 49 {
|
||||
t.Fatalf("response = %d, want 49", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestQUICCallConcurrent fans out N concurrent Calls on one QUIC
|
||||
// connection. Each Call MUST get its own stream — if they serialized
|
||||
// on a shared mutex, the handler's deliberate sleep would multiply
|
||||
@@ -269,9 +237,8 @@ func waitForPeer(t *testing.T, n *zap.Node, peerID string, timeout time.Duration
|
||||
}
|
||||
|
||||
// BenchmarkQUICCallConcurrent_1 / _10 / _100 measure how Call latency
|
||||
// scales with concurrent goroutines. With per-Call streams we expect
|
||||
// near-linear scaling; with the legacy ctrlMu serialization the rate
|
||||
// would plateau early.
|
||||
// scales with concurrent goroutines. Per-Call streams should yield
|
||||
// near-linear scaling up to the QUIC stream limit.
|
||||
func BenchmarkQUICCallConcurrent_1(b *testing.B) { benchQUICCall(b, 1) }
|
||||
func BenchmarkQUICCallConcurrent_10(b *testing.B) { benchQUICCall(b, 10) }
|
||||
func BenchmarkQUICCallConcurrent_100(b *testing.B) { benchQUICCall(b, 100) }
|
||||
|
||||
Reference in New Issue
Block a user