// Package adapter содержит кольцевой буфер последних принятых байт GPIO — // источник данных для HTTP API (/api/latest, /api/history, /api/health). package adapter import ( "sync" "sync/atomic" "time" ) // RingBuffer — потокобезопасный кольцевой буфер байт фиксированного размера. // Помимо самих данных ведёт статистику приёма: время последней записи, // суммарные счётчики и текущую скорость (байт/с и бит/с). type RingBuffer struct { mu sync.RWMutex data []byte writePtr int count int lastWrite atomic.Int64 totalBytes atomic.Uint64 // ===== SPEED ===== speedMu sync.Mutex lastCalcTime time.Time lastCalcBytes uint64 // Текущие значения скорости currentBPS atomic.Value // float64 currentBPSBits atomic.Value // float64 // Для скользящего окна recentBytes []timestampedByte recentMu sync.Mutex } type timestampedByte struct { timestamp time.Time } // NewRingBuffer создаёт буфер на size элементов. func NewRingBuffer(size int) *RingBuffer { rb := &RingBuffer{ data: make([]byte, size), lastCalcTime: time.Now(), recentBytes: make([]timestampedByte, 0, 1000), // храним последние 1000 байт } rb.currentBPS.Store(float64(0)) rb.currentBPSBits.Store(float64(0)) return rb } // Write добавляет байт в буфер (при переполнении затирается самый старый) // и обновляет статистику: счётчики, время последней записи, скорость. 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() rb.lastWrite.Store(now.UnixNano()) rb.totalBytes.Add(1) // ===== SPEED CALC (обновляем редко) ===== rb.updateSpeedPeriodically(now) // ===== Для скользящего окна ===== rb.addTimestampedByte(now) } // Периодическое обновление скорости (простое среднее) func (rb *RingBuffer) updateSpeedPeriodically(now time.Time) { rb.speedMu.Lock() defer rb.speedMu.Unlock() elapsed := now.Sub(rb.lastCalcTime) // Обновляем не чаще чем раз в 250мс if elapsed < 250*time.Millisecond { return } currentTotal := rb.totalBytes.Load() bytesInPeriod := currentTotal - rb.lastCalcBytes if elapsed.Seconds() > 0 { bps := float64(bytesInPeriod) / elapsed.Seconds() rb.currentBPS.Store(bps) rb.currentBPSBits.Store(bps * 8) } rb.lastCalcTime = now rb.lastCalcBytes = currentTotal } // addTimestampedByte ведёт скользящее окно временных меток за последнюю секунду // и пересчитывает по нему скорость (перетирая значение периодического механизма). func (rb *RingBuffer) addTimestampedByte(now time.Time) { rb.recentMu.Lock() defer rb.recentMu.Unlock() // Добавляем новую запись rb.recentBytes = append(rb.recentBytes, timestampedByte{timestamp: now}) // Удаляем старые записи (старше 1 секунды) cutoff := now.Add(-1 * time.Second) cutIdx := 0 for i, tb := range rb.recentBytes { if tb.timestamp.After(cutoff) { cutIdx = i break } } if cutIdx > 0 { rb.recentBytes = rb.recentBytes[cutIdx:] } // Вычисляем скорость на основе скользящего окна if len(rb.recentBytes) > 1 { windowDuration := rb.recentBytes[len(rb.recentBytes)-1].timestamp.Sub(rb.recentBytes[0].timestamp) if windowDuration.Seconds() > 0 { bps := float64(len(rb.recentBytes)) / windowDuration.Seconds() rb.currentBPS.Store(bps) rb.currentBPSBits.Store(bps * 8) } } } // GetWindowSpeed возвращает скорость приёма, рассчитанную по скользящему окну // в 1 секунду. Альтернатива значениям из Stats; текущим API не используется. func (rb *RingBuffer) GetWindowSpeed() (bytesPerSec, bitsPerSec float64) { rb.recentMu.Lock() defer rb.recentMu.Unlock() now := time.Now() cutoff := now.Add(-1 * time.Second) // Считаем байты за последнюю секунду count := 0 for i := len(rb.recentBytes) - 1; i >= 0; i-- { if rb.recentBytes[i].timestamp.Before(cutoff) { break } count++ } if count > 0 { // Находим время самого старого байта в окне oldest := rb.recentBytes[len(rb.recentBytes)-count].timestamp duration := now.Sub(oldest).Seconds() if duration > 0 { bytesPerSec = float64(count) / duration bitsPerSec = bytesPerSec * 8 } } return } // LastWriteTime возвращает время последней записи в буфер. func (rb *RingBuffer) LastWriteTime() time.Time { return time.Unix(0, rb.lastWrite.Load()) } // GetLatest возвращает последний записанный байт; false — если буфер пуст. 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 } // GetLast возвращает до n последних байт в порядке от старых к новым. 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 } // IsAlive сообщает, была ли запись в буфер за последние timeout. func (rb *RingBuffer) IsAlive(timeout time.Duration) bool { return time.Since(time.Unix(0, rb.lastWrite.Load())) < timeout } // Stats возвращает статистику буфера для /api/health. Ключи: size, filled, // last_write, total_bytes, total_bits, bytes_per_sec, bits_per_sec. func (rb *RingBuffer) Stats() map[string]interface{} { rb.mu.RLock() filled := rb.count rb.mu.RUnlock() bps, _ := rb.currentBPS.Load().(float64) bits, _ := rb.currentBPSBits.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, } }