File
Blob: overlay/reaper_test.go
| 1 | package overlay |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "crypto/ed25519" |
| 6 | "crypto/rand" |
| 7 | "crypto/tls" |
| 8 | "crypto/x509" |
| 9 | "io" |
| 10 | "math/big" |
| 11 | "testing" |
| 12 | "time" |
| 13 | |
| 14 | rttinstrumentation "go.miragespace.co/specter/rtt" |
| 15 | "go.miragespace.co/specter/spec/protocol" |
| 16 | "go.miragespace.co/specter/spec/rtt" |
| 17 | |
| 18 | "github.com/quic-go/quic-go" |
| 19 | "github.com/stretchr/testify/require" |
| 20 | "github.com/zhangyunhao116/skipmap" |
| 21 | "go.uber.org/zap/zaptest" |
| 22 | ) |
| 23 | |
| 24 | func TestReapPeerPreservesReplacement(t *testing.T) { |
| 25 | ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) |
| 26 | defer cancel() |
| 27 | dial := newReaperTestListener(t, ctx) |
| 28 | old, _ := dial() |
| 29 | replacement, remote := dial() |
| 30 | peer := &protocol.Node{Address: replacement.RemoteAddr().String(), Id: 1} |
| 31 | recorder := rttinstrumentation.NewInstrumentation(16) |
| 32 | transport := NewQUIC(TransportConfig{ |
| 33 | Logger: zaptest.NewLogger(t), |
| 34 | RTTRecorder: recorder, |
| 35 | }) |
| 36 | key := transport.makeCachedKey(peer) |
| 37 | |
| 38 | // Cache a replacement before the old connection's delayed cleanup. |
| 39 | current := &nodeConnection{peer: peer, quic: replacement} |
| 40 | transport.cachedConnections.Store(key, current) |
| 41 | pending := skipmap.NewUint64[int64]() |
| 42 | pending.Store(1, time.Now().UnixNano()) |
| 43 | transport.rttMap.Store(key, pending) |
| 44 | measurementKey := rtt.MakeMeasurementKey(peer) |
| 45 | recorder.RecordSent(measurementKey) |
| 46 | recorder.RecordLatency(measurementKey, float64(time.Millisecond)) |
| 47 | before := recorder.Snapshot(measurementKey, time.Minute) |
| 48 | require.NotNil(t, before) |
| 49 | |
| 50 | transport.reapPeer(old, peer) |
| 51 | |
| 52 | cached, ok := transport.cachedConnections.Load(key) |
| 53 | require.True(t, ok, "stale cleanup removed the replacement") |
| 54 | require.Same(t, current, cached) |
| 55 | require.NoError(t, replacement.Context().Err(), "stale cleanup closed the replacement") |
| 56 | require.Error(t, old.Context().Err(), "cleanup must close its own connection") |
| 57 | mapping, ok := transport.rttMap.Load(key) |
| 58 | require.True(t, ok, "stale cleanup removed the replacement's pending RTT probes") |
| 59 | require.Same(t, pending, mapping) |
| 60 | require.Equal(t, before, recorder.Snapshot(measurementKey, time.Minute)) |
| 61 | |
| 62 | // The cached replacement must still carry traffic, not just remain |
| 63 | // present in the cache. |
| 64 | outgoing, err := replacement.OpenStreamSync(ctx) |
| 65 | require.NoError(t, err) |
| 66 | deadline, _ := ctx.Deadline() |
| 67 | require.NoError(t, outgoing.SetWriteDeadline(deadline)) |
| 68 | _, err = outgoing.Write([]byte("replacement")) |
| 69 | require.NoError(t, err) |
| 70 | require.NoError(t, outgoing.Close()) |
| 71 | incoming, err := remote.AcceptStream(ctx) |
| 72 | require.NoError(t, err) |
| 73 | require.NoError(t, incoming.SetReadDeadline(deadline)) |
| 74 | payload, err := io.ReadAll(incoming) |
| 75 | require.NoError(t, err) |
| 76 | require.Equal(t, "replacement", string(payload)) |
| 77 | |
| 78 | transport.reapPeer(replacement, peer) |
| 79 | |
| 80 | _, ok = transport.cachedConnections.Load(key) |
| 81 | require.False(t, ok, "current cleanup must remove the cached connection") |
| 82 | require.Error(t, replacement.Context().Err()) |
| 83 | _, ok = transport.rttMap.Load(key) |
| 84 | require.False(t, ok, "current cleanup must remove pending RTT probes") |
| 85 | require.Nil(t, recorder.Snapshot(measurementKey, time.Minute)) |
| 86 | } |
| 87 | |
| 88 | func newReaperTestListener(t *testing.T, ctx context.Context) func() (*quic.Conn, *quic.Conn) { |
| 89 | t.Helper() |
| 90 | publicKey, privateKey, err := ed25519.GenerateKey(rand.Reader) |
| 91 | require.NoError(t, err) |
| 92 | template := &x509.Certificate{ |
| 93 | SerialNumber: big.NewInt(1), |
| 94 | DNSNames: []string{"localhost"}, |
| 95 | NotBefore: time.Now().Add(-time.Minute), |
| 96 | NotAfter: time.Now().Add(time.Hour), |
| 97 | KeyUsage: x509.KeyUsageDigitalSignature, |
| 98 | ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth}, |
| 99 | } |
| 100 | der, err := x509.CreateCertificate(rand.Reader, template, template, publicKey, privateKey) |
| 101 | require.NoError(t, err) |
| 102 | cert, err := x509.ParseCertificate(der) |
| 103 | require.NoError(t, err) |
| 104 | roots := x509.NewCertPool() |
| 105 | roots.AddCert(cert) |
| 106 | listener, err := quic.ListenAddr("127.0.0.1:0", &tls.Config{ |
| 107 | Certificates: []tls.Certificate{{Certificate: [][]byte{der}, PrivateKey: privateKey}}, |
| 108 | NextProtos: []string{"specter-reaper-test"}, |
| 109 | }, nil) |
| 110 | require.NoError(t, err) |
| 111 | t.Cleanup(func() { listener.Close() }) |
| 112 | |
| 113 | return func() (*quic.Conn, *quic.Conn) { |
| 114 | t.Helper() |
| 115 | local, err := quic.DialAddr(ctx, listener.Addr().String(), &tls.Config{ |
| 116 | ServerName: "localhost", |
| 117 | RootCAs: roots, |
| 118 | NextProtos: []string{"specter-reaper-test"}, |
| 119 | }, nil) |
| 120 | require.NoError(t, err) |
| 121 | t.Cleanup(func() { local.CloseWithError(0, "test cleanup") }) |
| 122 | remote, err := listener.Accept(ctx) |
| 123 | require.NoError(t, err) |
| 124 | t.Cleanup(func() { remote.CloseWithError(0, "test cleanup") }) |
| 125 | return local, remote |
| 126 | } |
| 127 | } |