Skip to content
File

Blob: rtt/rtt.go

go124 lines
1package rtt
2 
3import (
4 "sync"
5 "time"
6 
7 "go.miragespace.co/specter/spec/rtt"
8 "go.miragespace.co/specter/util"
9 
10 "github.com/montanaflynn/stats"
11 "github.com/zhangyunhao116/skipmap"
12)
13 
14type point struct {
15 time time.Time
16 value float64
17}
18 
19type container struct {
20 mu sync.RWMutex
21 data []point
22 sent uint64
23 lost uint64
24}
25 
26type Instrumentation struct {
27 measurement *skipmap.StringMap[*container]
28 length int
29}
30 
31var _ rtt.Recorder = (*Instrumentation)(nil)
32 
33func NewInstrumentation(max int) *Instrumentation {
34 return &Instrumentation{
35 measurement: skipmap.NewString[*container](),
36 length: max,
37 }
38}
39 
40func (i Instrumentation) getContainer(key string) *container {
41 c, _ := i.measurement.LoadOrStoreLazy(key, func() *container {
42 return &container{
43 data: make([]point, 0),
44 }
45 })
46 return c
47}
48 
49func (i *Instrumentation) RecordLatency(key string, value float64) {
50 if value < 0 {
51 return
52 }
53 
54 c := i.getContainer(key)
55 
56 c.mu.Lock()
57 if len(c.data) > i.length {
58 c.data = c.data[1:]
59 }
60 c.data = append(c.data, point{
61 time: time.Now(),
62 value: value,
63 })
64 c.mu.Unlock()
65}
66 
67func (i *Instrumentation) RecordSent(key string) {
68 c := i.getContainer(key)
69 c.mu.Lock()
70 c.sent++
71 c.mu.Unlock()
72}
73 
74func (i *Instrumentation) RecordLost(key string) {
75 c := i.getContainer(key)
76 c.mu.Lock()
77 c.lost++
78 c.mu.Unlock()
79}
80 
81func (i *Instrumentation) Snapshot(key string, last time.Duration) *rtt.Statistics {
82 c, ok := i.measurement.Load(key)
83 if !ok {
84 return nil
85 }
86 var (
87 values = make([]float64, 0)
88 since time.Time
89 until time.Time
90 sent uint64
91 lost uint64
92 )
93 c.mu.RLock()
94 for _, p := range c.data {
95 if time.Since(p.time) <= last {
96 if since.IsZero() {
97 since = p.time
98 }
99 until = p.time
100 values = append(values, p.value)
101 }
102 }
103 sent = c.sent
104 lost = c.lost
105 c.mu.RUnlock()
106 if len(values) < 1 {
107 return nil
108 }
109 return &rtt.Statistics{
110 Since: since,
111 Until: until,
112 Min: time.Duration(util.Must(stats.Min(values))),
113 Average: time.Duration(util.Must(stats.Mean(values))),
114 Max: time.Duration(util.Must(stats.Max(values))),
115 StandardDeviation: time.Duration(util.Must(stats.StandardDeviation(values))),
116 Sent: sent,
117 Lost: lost,
118 }
119}
120 
121func (i *Instrumentation) Drop(key string) {
122 i.measurement.Delete(key)
123}