dbd1808803
The server pushes each new step over server-sent events (GET /api/v1/live/stream) the moment it is sampled, so updates no longer arrive in uneven pairs. The chart slides left steadily between steps instead of jumping, and the big numbers count to their new value. Both stay still with reduced motion.
135 lines
3.3 KiB
Go
135 lines
3.3 KiB
Go
package main
|
|
|
|
import (
|
|
"log/slog"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// Speeds keeps the last few minutes of each peer's speed in memory for the
|
|
// Live page. It reads the kernel counters every speedStep, apart from the
|
|
// traffic history in Stats, and never writes to disk.
|
|
|
|
const (
|
|
speedStep = 2 * time.Second
|
|
speedPoints = 60 // 2 minutes
|
|
)
|
|
|
|
// SpeedPoint is one step: per peer ID, download and upload in bits per
|
|
// second, from the peer's point of view.
|
|
type SpeedPoint struct {
|
|
T int64 `json:"t"`
|
|
Peers map[string][2]int64 `json:"peers"`
|
|
}
|
|
|
|
type Speeds struct {
|
|
store *Store
|
|
kernel Kernel
|
|
|
|
mu sync.Mutex
|
|
last map[string][2]int64 // raw rx, tx by public key
|
|
lastAt time.Time
|
|
points []SpeedPoint
|
|
subs map[chan SpeedPoint]struct{}
|
|
done chan struct{} // closed when Run returns
|
|
}
|
|
|
|
func newSpeeds(store *Store, kernel Kernel) *Speeds {
|
|
return &Speeds{store: store, kernel: kernel, last: map[string][2]int64{}, subs: map[chan SpeedPoint]struct{}{}, done: make(chan struct{})}
|
|
}
|
|
|
|
// Subscribe returns the current points and a channel that receives each new
|
|
// one; cancel ends the subscription. A subscriber that falls behind misses
|
|
// points rather than holding up the sampler.
|
|
func (s *Speeds) Subscribe() (points []SpeedPoint, ch <-chan SpeedPoint, cancel func()) {
|
|
c := make(chan SpeedPoint, 4)
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
s.subs[c] = struct{}{}
|
|
return append([]SpeedPoint{}, s.points...), c, func() {
|
|
s.mu.Lock()
|
|
delete(s.subs, c)
|
|
s.mu.Unlock()
|
|
}
|
|
}
|
|
|
|
// Done is closed when the sampler stops, so streams can end.
|
|
func (s *Speeds) Done() <-chan struct{} { return s.done }
|
|
|
|
func (s *Speeds) sample(now time.Time) {
|
|
cfg := s.store.Get()
|
|
samples, err := s.kernel.Sample(cfg.Server.Interface)
|
|
if err != nil {
|
|
slog.Debug("speed sample failed", "err", err)
|
|
return
|
|
}
|
|
idByKey := map[string]string{}
|
|
for _, p := range cfg.Peers {
|
|
if p.hasKey() {
|
|
idByKey[p.PublicKey] = p.ID
|
|
}
|
|
}
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
secs := now.Sub(s.lastAt).Seconds()
|
|
first := s.lastAt.IsZero()
|
|
cur := map[string][2]int64{}
|
|
pt := SpeedPoint{T: now.Unix(), Peers: map[string][2]int64{}}
|
|
for _, smp := range samples {
|
|
cur[smp.PublicKey] = [2]int64{smp.RxBytes, smp.TxBytes}
|
|
id := idByKey[smp.PublicKey]
|
|
prev, ok := s.last[smp.PublicKey]
|
|
if id == "" || !ok || first {
|
|
continue
|
|
}
|
|
dRx, dTx := smp.RxBytes-prev[0], smp.TxBytes-prev[1]
|
|
if dRx < 0 || dTx < 0 { // counters were reset
|
|
continue
|
|
}
|
|
// Tx is what the server sent: the peer's download.
|
|
pt.Peers[id] = [2]int64{int64(float64(dTx*8) / secs), int64(float64(dRx*8) / secs)}
|
|
}
|
|
s.last, s.lastAt = cur, now
|
|
if first {
|
|
return
|
|
}
|
|
s.points = append(s.points, pt)
|
|
if len(s.points) > speedPoints {
|
|
s.points = s.points[len(s.points)-speedPoints:]
|
|
}
|
|
for c := range s.subs {
|
|
select {
|
|
case c <- pt:
|
|
default:
|
|
}
|
|
}
|
|
}
|
|
|
|
// Since returns the points newer than the unix time t, oldest first.
|
|
func (s *Speeds) Since(t int64) []SpeedPoint {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
out := []SpeedPoint{}
|
|
for _, p := range s.points {
|
|
if p.T > t {
|
|
out = append(out, p)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (s *Speeds) Run(stop <-chan struct{}) {
|
|
defer close(s.done)
|
|
s.sample(time.Now())
|
|
t := time.NewTicker(speedStep)
|
|
defer t.Stop()
|
|
for {
|
|
select {
|
|
case <-stop:
|
|
return
|
|
case now := <-t.C:
|
|
s.sample(now)
|
|
}
|
|
}
|
|
}
|