Skip to content
File

Blob: overlay/reaper.go

go69 lines
1package overlay
2 
3import (
4 "context"
5 "time"
6 
7 "go.miragespace.co/specter/spec/protocol"
8 "go.miragespace.co/specter/spec/rtt"
9 "go.miragespace.co/specter/util"
10 
11 "github.com/quic-go/quic-go"
12 "go.uber.org/zap"
13)
14 
15func (t *QUIC) reapPeer(q *quic.Conn, peer *protocol.Node) {
16 qKey := t.makeCachedKey(peer)
17 unlock := t.cachedMutex.Lock(qKey)
18 defer unlock()
19 
20 // A delayed close callback may arrive after a replacement was cached.
21 // Only the current connection owns the cached entry and its RTT state.
22 if cached, loaded := t.cachedConnections.Load(qKey); loaded && cached.quic == q {
23 t.Logger.Debug("reaping cached QUIC connection to peer", zap.String("key", qKey))
24 t.cachedConnections.Delete(qKey)
25 t.rttMap.Delete(qKey)
26 if t.RTTRecorder != nil {
27 t.RTTRecorder.Drop(rtt.MakeMeasurementKey(peer))
28 }
29 }
30 q.CloseWithError(401, "Gone")
31}
32 
33// TODO: investigate if reaper is now deprecated
34func (t *QUIC) reaper(ctx context.Context) {
35 timer := time.NewTimer(util.RandomTimeRange(quicConfig.HandshakeIdleTimeout))
36 
37 d := &protocol.Datagram{
38 Type: protocol.Datagram_ALIVE,
39 }
40 alive, err := d.MarshalVT()
41 if err != nil {
42 panic(err)
43 }
44 
45 for {
46 select {
47 case <-ctx.Done():
48 t.Logger.Info("Exiting reaper goroutine")
49 return
50 case <-timer.C:
51 candidate := make([]*nodeConnection, 0)
52 t.cachedConnections.Range(func(key string, value *nodeConnection) bool {
53 if err := value.quic.SendDatagram(alive); err != nil {
54 candidate = append(candidate, value)
55 }
56 return true
57 })
58 
59 if len(candidate) > 0 {
60 for _, c := range candidate {
61 t.reapPeer(c.quic, c.peer)
62 }
63 }
64 
65 timer.Reset(util.RandomTimeRange(quicConfig.MaxIdleTimeout))
66 }
67 }
68}