рабочая версия

This commit is contained in:
Tot Maxim
2026-05-07 21:30:48 +03:00
parent 649bb01645
commit 9c28b9501d
2 changed files with 189 additions and 113 deletions

View File

@@ -1,99 +1,140 @@
// internal/adapter/buffer.go
package adapter
import (
"sync"
"time"
"sync"
"time"
)
type RingBuffer struct {
mu sync.RWMutex
data []byte
writePtr int
count int
mu sync.RWMutex
data []byte
writePtr int
count int
lastWrite time.Time
lastWrite time.Time
// ===== SPEED METRICS =====
totalBytes uint64
lastTick time.Time
lastBytes uint64
bytesPerSec float64
bitsPerSec float64
}
func NewRingBuffer(size int) *RingBuffer {
return &RingBuffer{
data: make([]byte, size),
lastWrite: time.Now(),
}
return &RingBuffer{
data: make([]byte, size),
lastWrite: time.Now(),
lastTick: time.Now(),
}
}
func (rb *RingBuffer) Write(b byte) {
rb.mu.Lock()
defer rb.mu.Unlock()
rb.mu.Lock()
rb.data[rb.writePtr] = b
rb.writePtr = (rb.writePtr + 1) % len(rb.data)
// ===== WRITE =====
rb.data[rb.writePtr] = b
rb.writePtr = (rb.writePtr + 1) % len(rb.data)
if rb.count < len(rb.data) {
rb.count++
}
if rb.count < len(rb.data) {
rb.count++
}
rb.lastWrite = time.Now()
rb.lastWrite = time.Now()
rb.totalBytes++
rb.mu.Unlock()
// ===== SPEED CALC (outside lock) =====
now := rb.lastWrite
if rb.lastTick.IsZero() {
rb.lastTick = now
rb.lastBytes = rb.totalBytes
return
}
elapsed := now.Sub(rb.lastTick).Seconds()
if elapsed >= 1.0 {
diff := rb.totalBytes - rb.lastBytes
rb.bytesPerSec = float64(diff) / elapsed
rb.bitsPerSec = rb.bytesPerSec * 8
rb.lastTick = now
rb.lastBytes = rb.totalBytes
}
}
func (rb *RingBuffer) LastWriteTime() time.Time {
rb.mu.RLock()
defer rb.mu.RUnlock()
return rb.lastWrite
rb.mu.RLock()
defer rb.mu.RUnlock()
return rb.lastWrite
}
func (rb *RingBuffer) GetLatest() (byte, bool) {
rb.mu.RLock()
defer rb.mu.RUnlock()
rb.mu.RLock()
defer rb.mu.RUnlock()
if rb.count == 0 {
return 0, false
}
if rb.count == 0 {
return 0, false
}
idx := rb.writePtr - 1
if idx < 0 {
idx += len(rb.data)
}
idx := rb.writePtr - 1
if idx < 0 {
idx += len(rb.data)
}
return rb.data[idx], true
return rb.data[idx], true
}
func (rb *RingBuffer) GetLast(n int) []byte {
rb.mu.RLock()
defer rb.mu.RUnlock()
rb.mu.RLock()
defer rb.mu.RUnlock()
if n > rb.count {
n = rb.count
}
if n > rb.count {
n = rb.count
}
res := make([]byte, n)
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])
}
start := rb.writePtr - n
return res
}
if start < 0 {
start += len(rb.data)
func (rb *RingBuffer) Stats() map[string]interface{} {
rb.mu.RLock()
defer rb.mu.RUnlock()
copy(res, rb.data[start:])
copy(res[len(rb.data)-start:], rb.data[:rb.writePtr])
} else {
copy(res, rb.data[start:rb.writePtr])
}
return map[string]interface{}{
"size": len(rb.data),
"filled": rb.count,
"last_write": rb.lastWrite,
}
return res
}
func (rb *RingBuffer) IsAlive(timeout time.Duration) bool {
rb.mu.RLock()
defer rb.mu.RUnlock()
rb.mu.RLock()
defer rb.mu.RUnlock()
return time.Since(rb.lastWrite) < timeout
return time.Since(rb.lastWrite) < timeout
}
func (rb *RingBuffer) Stats() map[string]interface{} {
rb.mu.RLock()
defer rb.mu.RUnlock()
return map[string]interface{}{
"size": len(rb.data),
"filled": rb.count,
"last_write": rb.lastWrite,
"total_bytes": rb.totalBytes,
"total_bits": rb.totalBytes * 8,
"bytes_per_sec": rb.bytesPerSec,
"bits_per_sec": rb.bitsPerSec,
}
}