163 lines
2.7 KiB
Go
163 lines
2.7 KiB
Go
package adapter
|
|
|
|
import (
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
)
|
|
|
|
type RingBuffer struct {
|
|
mu sync.RWMutex
|
|
data []byte
|
|
writePtr int
|
|
count int
|
|
|
|
lastWrite atomic.Int64
|
|
|
|
totalBytes atomic.Uint64
|
|
|
|
// ===== SPEED =====
|
|
lastTick atomic.Int64
|
|
lastBytes atomic.Uint64
|
|
|
|
bytesPerSec atomic.Value // float64
|
|
bitsPerSec atomic.Value // float64
|
|
}
|
|
|
|
// =====================
|
|
func NewRingBuffer(size int) *RingBuffer {
|
|
rb := &RingBuffer{
|
|
data: make([]byte, size),
|
|
}
|
|
|
|
now := time.Now().UnixNano()
|
|
|
|
rb.lastWrite.Store(now)
|
|
rb.lastTick.Store(now)
|
|
|
|
rb.bytesPerSec.Store(float64(0))
|
|
rb.bitsPerSec.Store(float64(0))
|
|
|
|
return rb
|
|
}
|
|
|
|
// =====================
|
|
func (rb *RingBuffer) Write(b byte) {
|
|
rb.mu.Lock()
|
|
|
|
rb.data[rb.writePtr] = b
|
|
rb.writePtr = (rb.writePtr + 1) % len(rb.data)
|
|
|
|
if rb.count < len(rb.data) {
|
|
rb.count++
|
|
}
|
|
|
|
rb.mu.Unlock()
|
|
|
|
// ===== ATOMIC STATS =====
|
|
now := time.Now().UnixNano()
|
|
rb.lastWrite.Store(now)
|
|
|
|
total := rb.totalBytes.Add(1)
|
|
|
|
// ===== SPEED CALC =====
|
|
lastTick := rb.lastTick.Load()
|
|
if lastTick == 0 {
|
|
rb.lastTick.Store(now)
|
|
rb.lastBytes.Store(total)
|
|
return
|
|
}
|
|
|
|
elapsed := time.Duration(now - lastTick).Seconds()
|
|
if elapsed < 1.0 {
|
|
return
|
|
}
|
|
|
|
last := rb.lastBytes.Load()
|
|
diff := total - last
|
|
|
|
bps := float64(diff) / elapsed
|
|
|
|
rb.bytesPerSec.Store(bps)
|
|
rb.bitsPerSec.Store(bps * 8)
|
|
|
|
rb.lastTick.Store(now)
|
|
rb.lastBytes.Store(total)
|
|
}
|
|
|
|
// =====================
|
|
func (rb *RingBuffer) LastWriteTime() time.Time {
|
|
return time.Unix(0, rb.lastWrite.Load())
|
|
}
|
|
|
|
// =====================
|
|
func (rb *RingBuffer) GetLatest() (byte, bool) {
|
|
rb.mu.RLock()
|
|
defer rb.mu.RUnlock()
|
|
|
|
if rb.count == 0 {
|
|
return 0, false
|
|
}
|
|
|
|
idx := rb.writePtr - 1
|
|
if idx < 0 {
|
|
idx += len(rb.data)
|
|
}
|
|
|
|
return rb.data[idx], true
|
|
}
|
|
|
|
// =====================
|
|
func (rb *RingBuffer) GetLast(n int) []byte {
|
|
rb.mu.RLock()
|
|
defer rb.mu.RUnlock()
|
|
|
|
if n > rb.count {
|
|
n = rb.count
|
|
}
|
|
|
|
res := make([]byte, n)
|
|
|
|
start := rb.writePtr - n
|
|
|
|
if start < 0 {
|
|
start += len(rb.data)
|
|
copy(res, rb.data[start:])
|
|
copy(res[len(rb.data)-start:], rb.data[:rb.writePtr])
|
|
} else {
|
|
copy(res, rb.data[start:rb.writePtr])
|
|
}
|
|
|
|
return res
|
|
}
|
|
|
|
// =====================
|
|
func (rb *RingBuffer) IsAlive(timeout time.Duration) bool {
|
|
return time.Since(time.Unix(0, rb.lastWrite.Load())) < timeout
|
|
}
|
|
|
|
// =====================
|
|
func (rb *RingBuffer) Stats() map[string]interface{} {
|
|
rb.mu.RLock()
|
|
filled := rb.count
|
|
rb.mu.RUnlock()
|
|
|
|
bps, _ := rb.bytesPerSec.Load().(float64)
|
|
bits, _ := rb.bitsPerSec.Load().(float64)
|
|
|
|
total := rb.totalBytes.Load()
|
|
|
|
return map[string]interface{}{
|
|
"size": len(rb.data),
|
|
"filled": filled,
|
|
|
|
"last_write": time.Unix(0, rb.lastWrite.Load()),
|
|
|
|
"total_bytes": total,
|
|
"total_bits": total * 8,
|
|
|
|
"bytes_per_sec": bps,
|
|
"bits_per_sec": bits,
|
|
}
|
|
}
|