Skip to content
File

Blob: overlay/reaper_test.go

go128 lines
1package overlay
2 
3import (
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 
24func 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 
88func 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}