Первая итерация заработала
This commit is contained in:
@@ -6,90 +6,94 @@ import (
|
||||
"time"
|
||||
)
|
||||
|
||||
// RingBuffer - потокобезопасный кольцевой буфер для байт GPIO
|
||||
type RingBuffer struct {
|
||||
mu sync.RWMutex
|
||||
data []byte
|
||||
writePtr int
|
||||
count int // сколько байт записано (до заполнения буфера)
|
||||
mu sync.RWMutex
|
||||
data []byte
|
||||
writePtr int
|
||||
count int
|
||||
|
||||
lastWrite time.Time
|
||||
}
|
||||
|
||||
// NewRingBuffer создает буфер заданного размера
|
||||
func NewRingBuffer(size int) *RingBuffer {
|
||||
return &RingBuffer{
|
||||
data: make([]byte, size),
|
||||
data: make([]byte, size),
|
||||
lastWrite: time.Now(),
|
||||
}
|
||||
}
|
||||
|
||||
// Write добавляет байт в буфер
|
||||
func (rb *RingBuffer) Write(b byte) {
|
||||
rb.mu.Lock()
|
||||
defer rb.mu.Unlock()
|
||||
|
||||
|
||||
rb.data[rb.writePtr] = b
|
||||
rb.writePtr = (rb.writePtr + 1) % len(rb.data)
|
||||
|
||||
if rb.count < len(rb.data) {
|
||||
rb.count++
|
||||
}
|
||||
|
||||
rb.lastWrite = time.Now()
|
||||
}
|
||||
|
||||
// GetLatest возвращает последний записанный байт
|
||||
func (rb *RingBuffer) LastWriteTime() time.Time {
|
||||
rb.mu.RLock()
|
||||
defer rb.mu.RUnlock()
|
||||
return rb.lastWrite
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
if n == 0 {
|
||||
return []byte{}
|
||||
}
|
||||
|
||||
result := make([]byte, n)
|
||||
|
||||
res := make([]byte, n)
|
||||
|
||||
start := rb.writePtr - n
|
||||
if start < 0 {
|
||||
start += len(rb.data)
|
||||
copy(result, rb.data[start:])
|
||||
copy(result[len(rb.data)-start:], rb.data[:rb.writePtr])
|
||||
copy(res, rb.data[start:])
|
||||
copy(res[len(rb.data)-start:], rb.data[:rb.writePtr])
|
||||
} else {
|
||||
copy(result, rb.data[start:rb.writePtr])
|
||||
copy(res, rb.data[start:rb.writePtr])
|
||||
}
|
||||
return result
|
||||
|
||||
return res
|
||||
}
|
||||
|
||||
// GetAll возвращает все данные из буфера
|
||||
func (rb *RingBuffer) GetAll() []byte {
|
||||
return rb.GetLast(rb.count)
|
||||
}
|
||||
|
||||
// Stats возвращает статистику буфера
|
||||
func (rb *RingBuffer) Stats() map[string]interface{} {
|
||||
rb.mu.RLock()
|
||||
defer rb.mu.RUnlock()
|
||||
|
||||
|
||||
return map[string]interface{}{
|
||||
"total_bytes": rb.count,
|
||||
"buffer_size": len(rb.data),
|
||||
"last_write": rb.lastWrite.Format(time.RFC3339),
|
||||
"is_full": rb.count == len(rb.data),
|
||||
"size": len(rb.data),
|
||||
"filled": rb.count,
|
||||
"last_write": rb.lastWrite,
|
||||
}
|
||||
}
|
||||
|
||||
func (rb *RingBuffer) IsAlive(timeout time.Duration) bool {
|
||||
rb.mu.RLock()
|
||||
defer rb.mu.RUnlock()
|
||||
|
||||
return time.Since(rb.lastWrite) < timeout
|
||||
}
|
||||
|
||||
@@ -1,60 +0,0 @@
|
||||
// internal/adapter/reader.go
|
||||
package adapter
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"io"
|
||||
"log"
|
||||
"os"
|
||||
)
|
||||
|
||||
// StdinReader читает бинарные данные из stdin и пишет в буфер
|
||||
type StdinReader struct {
|
||||
updates chan byte // канал для real-time уведомлений (0 = без уведомлений)
|
||||
}
|
||||
|
||||
// NewStdinReader создает читалку stdin
|
||||
func NewStdinReader(updates chan byte) *StdinReader {
|
||||
return &StdinReader{
|
||||
updates: updates,
|
||||
}
|
||||
}
|
||||
|
||||
// Start начинает чтение stdin в отдельной горутине
|
||||
func (r *StdinReader) Start(ctx context.Context, buf *RingBuffer) error {
|
||||
reader := bufio.NewReader(os.Stdin)
|
||||
|
||||
go func() {
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
log.Println("StdinReader: остановлен")
|
||||
return
|
||||
default:
|
||||
b, err := reader.ReadByte()
|
||||
if err != nil {
|
||||
if err == io.EOF {
|
||||
log.Println("StdinReader: stdin закрыт (EOF)")
|
||||
return
|
||||
}
|
||||
// Другие ошибки игнорируем, продолжаем читать
|
||||
continue
|
||||
}
|
||||
|
||||
buf.Write(b)
|
||||
|
||||
// Неблокирующая отправка в канал уведомлений
|
||||
if r.updates != nil {
|
||||
select {
|
||||
case r.updates <- b:
|
||||
default:
|
||||
// Канал заполнен, пропускаем уведомление
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user