-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathretransmission_queue.go
102 lines (90 loc) · 2.02 KB
/
retransmission_queue.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
package spectral
import (
"slices"
"sync"
"time"
)
const retransmissionAttempts = 3
type retransmissionEntry struct {
sequenceID uint32
payload []byte
sent time.Time
attempts int
}
type retransmissionQueue struct {
queue []*retransmissionEntry
mu sync.RWMutex
}
func newRetransmissionQueue() *retransmissionQueue {
return &retransmissionQueue{}
}
func (r *retransmissionQueue) add(now time.Time, sequenceID uint32, p []byte) {
r.mu.Lock()
r.queue = append(r.queue, &retransmissionEntry{sequenceID: sequenceID, payload: p, sent: now})
r.sort()
r.mu.Unlock()
}
func (r *retransmissionQueue) remove(sequenceID uint32) *retransmissionEntry {
r.mu.Lock()
defer r.mu.Unlock()
index := slices.IndexFunc(r.queue, func(e *retransmissionEntry) bool { return e.sequenceID == sequenceID })
if index >= 0 {
entry := r.queue[index]
r.queue = slices.Delete(r.queue, index, index+1)
return entry
}
return nil
}
func (r *retransmissionQueue) next(rto time.Duration) (t time.Time) {
r.mu.RLock()
defer r.mu.RUnlock()
if len(r.queue) > 0 {
return r.queue[0].sent.Add(rto)
}
return
}
func (r *retransmissionQueue) shift(now time.Time, rto time.Duration) (p []byte, t time.Time) {
r.mu.Lock()
defer r.mu.Unlock()
if len(r.queue) == 0 {
return
}
entry := r.queue[0]
if now.Sub(entry.sent) >= rto {
sent := entry.sent
entry.sent = now
entry.attempts++
if entry.attempts >= retransmissionAttempts {
r.queue[0] = nil
r.queue = r.queue[1:]
} else {
r.queue = append(r.queue[1:], entry)
r.sort()
}
return entry.payload, sent
}
return
}
func (r *retransmissionQueue) clear() {
r.mu.Lock()
for i, entry := range r.queue {
entry.payload = entry.payload[:0]
entry.payload = nil
r.queue[i] = nil
}
r.queue = r.queue[:0]
r.queue = nil
r.mu.Unlock()
}
func (r *retransmissionQueue) sort() {
if len(r.queue) > 1 {
slices.SortFunc(r.queue, func(a, b *retransmissionEntry) int {
if a.sent.Before(b.sent) {
return -1
} else if b.sent.Before(a.sent) {
return 1
}
return 0
})
}
}