Skip to content
File

Blob: overlay/rtt.go

go145 lines
1package overlay
2 
3import (
4 "context"
5 "encoding/binary"
6 "time"
7 
8 "go.miragespace.co/specter/spec/protocol"
9 "go.miragespace.co/specter/spec/rtt"
10 "go.miragespace.co/specter/spec/transport"
11 "go.miragespace.co/specter/util"
12 
13 "github.com/quic-go/quic-go"
14 "github.com/zhangyunhao116/skipmap"
15 "go.uber.org/zap"
16)
17 
18func (t *QUIC) sendRTTSyn(ctx context.Context, q *quic.Conn, peer *protocol.Node) {
19 l := t.Logger.With(zap.String("endpoint", q.RemoteAddr().String()), zap.Object("peer", peer))
20 
21 var (
22 qKey = t.makeCachedKey(peer)
23 mapper = skipmap.NewUint64[int64]()
24 counter uint64 = 1
25 rttBuf protocol.Datagram = protocol.Datagram{
26 Type: protocol.Datagram_RTT_SYN,
27 Data: make([]byte, 8),
28 }
29 err error
30 buf []byte
31 )
32 
33 t.rttMap.Store(qKey, mapper)
34 
35 go t.handleRTTLost(ctx, q, peer)
36 
37 for {
38 select {
39 case <-ctx.Done():
40 return
41 case <-q.Context().Done():
42 return
43 default:
44 binary.BigEndian.PutUint64(rttBuf.Data, counter)
45 buf, err = rttBuf.MarshalVT()
46 if err != nil {
47 l.Error("error encoding rtt syn datagram to proto", zap.Error(err))
48 continue
49 }
50 mapper.Store(counter, time.Now().UnixNano())
51 if err = q.SendDatagram(buf); err != nil {
52 l.Error("failed to send rtt syn", zap.Error(err))
53 mapper.Delete(counter)
54 continue
55 }
56 if t.RTTRecorder != nil {
57 t.RTTRecorder.RecordSent(rtt.MakeMeasurementKey(peer))
58 }
59 counter++
60 time.Sleep(util.RandomTimeRange(transport.RTTMeasureInterval))
61 }
62 }
63}
64 
65func (t *QUIC) handleRTTAck(ctx context.Context) {
66 var (
67 l *zap.Logger
68 qKey string
69 received int64
70 sent int64
71 counter uint64
72 ok bool
73 mapper *skipmap.Uint64Map[int64]
74 d *transport.DatagramDelegate
75 )
76 
77 for {
78 select {
79 case <-ctx.Done():
80 t.Logger.Info("Exiting RTT ack goroutine")
81 return
82 case d = <-t.rttChan:
83 l = t.Logger.With(zap.Object("peer", d.Identity))
84 qKey = t.makeCachedKey(d.Identity)
85 mapper, ok = t.rttMap.Load(qKey)
86 if !ok {
87 l.Warn("No mapping found for processing rtt ack")
88 continue
89 }
90 if len(d.Buffer) != 8 {
91 l.Warn("invalid length for counter timestamp", zap.Int("got", len(d.Buffer)))
92 continue
93 }
94 received = time.Now().UnixNano()
95 counter = binary.BigEndian.Uint64(d.Buffer)
96 sent, ok = mapper.LoadAndDelete(counter)
97 if !ok {
98 l.Warn("no timestamp found for counter", zap.Uint64("counter", counter))
99 continue
100 }
101 if t.RTTRecorder != nil {
102 t.RTTRecorder.RecordLatency(rtt.MakeMeasurementKey(d.Identity), float64(received-sent))
103 }
104 }
105 }
106}
107 
108func (t *QUIC) handleRTTLost(ctx context.Context, q *quic.Conn, peer *protocol.Node) {
109 if t.RTTRecorder == nil {
110 return
111 }
112 
113 var (
114 qKey = t.makeCachedKey(peer)
115 mapper *skipmap.Uint64Map[int64]
116 ok bool
117 )
118 
119 mapper, ok = t.rttMap.Load(qKey)
120 if !ok {
121 t.Logger.Error("No mapper found for peer", zap.Object("peer", peer))
122 return
123 }
124 
125 ticker := time.NewTicker(transport.RTTMeasureInterval)
126 defer ticker.Stop()
127 
128 for {
129 select {
130 case <-ctx.Done():
131 return
132 case <-q.Context().Done():
133 return
134 case <-ticker.C:
135 mapper.Range(func(counter uint64, sent int64) bool {
136 if time.Since(time.Unix(0, sent)) > 2*transport.RTTMeasureInterval {
137 t.RTTRecorder.RecordLost(rtt.MakeMeasurementKey(peer))
138 mapper.Delete(counter)
139 }
140 return true
141 })
142 }
143 }
144}