package pzserver import "sync" // LogLine — одна строка консоли сервера с монотонным номером, по которому // веб-клиент понимает, что он пропустил, и откуда догружать. type LogLine struct { Seq int64 `json:"seq"` Time int64 `json:"time"` // Unix-миллисекунды Stream string `json:"stream"` Text string `json:"text"` } // logBuffer — кольцевой буфер последних строк консоли плюс рассылка их // живым подписчикам (SSE-соединениям веб-панели). type logBuffer struct { mu sync.RWMutex lines []LogLine capacity int nextSeq int64 subs map[int]chan LogLine nextSub int } func newLogBuffer(capacity int) *logBuffer { return &logBuffer{ lines: make([]LogLine, 0, capacity), capacity: capacity, nextSeq: 1, subs: make(map[int]chan LogLine), } } // append кладёт строку в буфер и рассылает подписчикам. func (b *logBuffer) append(line LogLine) { b.mu.Lock() line.Seq = b.nextSeq b.nextSeq++ if len(b.lines) == b.capacity { copy(b.lines, b.lines[1:]) b.lines = b.lines[:len(b.lines)-1] } b.lines = append(b.lines, line) subs := make([]chan LogLine, 0, len(b.subs)) for _, ch := range b.subs { subs = append(subs, ch) } b.mu.Unlock() for _, ch := range subs { // Неблокирующая отправка: подписчик, который не успевает читать, // теряет строки, но не тормозит вывод сервера. select { case ch <- line: default: } } } // since возвращает строки с номером больше seq (seq<=0 — весь буфер). func (b *logBuffer) since(seq int64) []LogLine { b.mu.RLock() defer b.mu.RUnlock() out := make([]LogLine, 0, len(b.lines)) for _, l := range b.lines { if l.Seq > seq { out = append(out, l) } } return out } // subscribe открывает канал живых строк и функцию отписки. func (b *logBuffer) subscribe(bufferSize int) (<-chan LogLine, func()) { ch := make(chan LogLine, bufferSize) b.mu.Lock() id := b.nextSub b.nextSub++ b.subs[id] = ch b.mu.Unlock() return ch, func() { b.mu.Lock() if existing, ok := b.subs[id]; ok { delete(b.subs, id) close(existing) } b.mu.Unlock() } }