0.2.55 - Remote config/Request fixes
This commit is contained in:
@@ -2,9 +2,13 @@
|
||||
package logging
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
@@ -90,6 +94,18 @@ type Buffer struct {
|
||||
events []Event
|
||||
start int
|
||||
next atomic.Int64
|
||||
history *historyFile
|
||||
}
|
||||
|
||||
// historyFile is an append-only JSONL archive of the same structured events the admin
|
||||
// console reads. It is compacted to the ring's retained tail at startup and after each
|
||||
// further ringful, so persistence cannot become an unbounded disk cost.
|
||||
type historyFile struct {
|
||||
mu sync.Mutex
|
||||
path string
|
||||
file *os.File
|
||||
capacity int
|
||||
lastCompacted int64
|
||||
}
|
||||
|
||||
// ParseCapacity returns a non-negative log buffer capacity from configuration.
|
||||
@@ -131,6 +147,69 @@ func NewBuffered(
|
||||
return slog.New(&captureHandler{next: written, buffer: buffer, level: level}), buffer
|
||||
}
|
||||
|
||||
// NewPersistentBuffered restores retained events from path before accepting new ones and
|
||||
// appends every subsequent accepted record. The caller should close the returned Buffer.
|
||||
// An empty path keeps the in-memory behaviour used by unit tests and small embeddings.
|
||||
func NewPersistentBuffered(
|
||||
w io.Writer, level slog.Leveler, capacity int, format Format, path string,
|
||||
) (*slog.Logger, *Buffer, error) {
|
||||
logger, buffer := NewBuffered(w, level, capacity, format)
|
||||
path = strings.TrimSpace(path)
|
||||
if capacity <= 0 || path == "" {
|
||||
return logger, buffer, nil
|
||||
}
|
||||
history, restored, err := openHistory(path, capacity)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
buffer.history = history
|
||||
buffer.events = restored
|
||||
if len(restored) > 0 {
|
||||
buffer.next.Store(restored[len(restored)-1].Sequence)
|
||||
}
|
||||
return logger, buffer, nil
|
||||
}
|
||||
|
||||
func openHistory(path string, capacity int) (*historyFile, []Event, error) {
|
||||
if err := os.MkdirAll(filepath.Dir(path), 0o750); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
input, err := os.OpenFile(path, os.O_CREATE|os.O_RDONLY, 0o640)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
restored := make([]Event, 0, capacity)
|
||||
scanner := bufio.NewScanner(input)
|
||||
// Error attributes are capped upstream, but allow headroom for structured records.
|
||||
scanner.Buffer(make([]byte, 64<<10), 1<<20)
|
||||
for scanner.Scan() {
|
||||
var event Event
|
||||
if json.Unmarshal(scanner.Bytes(), &event) != nil || event.Sequence <= 0 {
|
||||
continue
|
||||
}
|
||||
restored = append(restored, event)
|
||||
if len(restored) > capacity {
|
||||
copy(restored, restored[len(restored)-capacity:])
|
||||
restored = restored[:capacity]
|
||||
}
|
||||
}
|
||||
closeErr := input.Close()
|
||||
if err := scanner.Err(); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
if closeErr != nil {
|
||||
return nil, nil, closeErr
|
||||
}
|
||||
history := &historyFile{path: path, capacity: capacity}
|
||||
if len(restored) > 0 {
|
||||
history.lastCompacted = restored[len(restored)-1].Sequence
|
||||
}
|
||||
if err := history.rewrite(restored); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
return history, restored, nil
|
||||
}
|
||||
|
||||
type captureHandler struct {
|
||||
next slog.Handler
|
||||
buffer *Buffer
|
||||
@@ -197,9 +276,77 @@ func (b *Buffer) append(event Event) {
|
||||
if len(b.events) == b.capacity {
|
||||
b.events[b.start] = event
|
||||
b.start = (b.start + 1) % b.capacity
|
||||
return
|
||||
} else {
|
||||
b.events = append(b.events, event)
|
||||
}
|
||||
b.events = append(b.events, event)
|
||||
if b.history != nil {
|
||||
ordered := b.orderedEventsLocked()
|
||||
_ = b.history.append(event, ordered)
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Buffer) orderedEventsLocked() []Event {
|
||||
ordered := make([]Event, len(b.events))
|
||||
for i := range b.events {
|
||||
ordered[i] = b.events[(b.start+i)%len(b.events)]
|
||||
}
|
||||
return ordered
|
||||
}
|
||||
|
||||
func (h *historyFile) append(event Event, retained []Event) error {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
if event.Sequence-h.lastCompacted >= int64(h.capacity) {
|
||||
if err := h.rewriteLocked(retained); err != nil {
|
||||
return err
|
||||
}
|
||||
h.lastCompacted = event.Sequence
|
||||
return nil
|
||||
}
|
||||
return json.NewEncoder(h.file).Encode(event)
|
||||
}
|
||||
|
||||
func (h *historyFile) rewrite(events []Event) error {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
return h.rewriteLocked(events)
|
||||
}
|
||||
|
||||
func (h *historyFile) rewriteLocked(events []Event) error {
|
||||
if h.file != nil {
|
||||
if err := h.file.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
file, err := os.OpenFile(h.path, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o640)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
encoder := json.NewEncoder(file)
|
||||
for _, event := range events {
|
||||
if err := encoder.Encode(event); err != nil {
|
||||
_ = file.Close()
|
||||
return err
|
||||
}
|
||||
}
|
||||
if err := file.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
h.file, err = os.OpenFile(h.path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o640)
|
||||
return err
|
||||
}
|
||||
|
||||
// Close flushes the persistent archive. It is safe for an in-memory Buffer.
|
||||
func (b *Buffer) Close() error {
|
||||
if b == nil || b.history == nil {
|
||||
return nil
|
||||
}
|
||||
b.history.mu.Lock()
|
||||
defer b.history.mu.Unlock()
|
||||
if b.history.file == nil {
|
||||
return nil
|
||||
}
|
||||
return b.history.file.Close()
|
||||
}
|
||||
|
||||
// Events returns records strictly newer than after, up to limit. If the caller fell
|
||||
|
||||
@@ -3,6 +3,7 @@ package logging
|
||||
import (
|
||||
"bytes"
|
||||
"log/slog"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
@@ -114,6 +115,32 @@ func TestBufferedLoggerRetainsStructuredEventsWithCursorPagination(t *testing.T)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPersistentBufferRestoresTheRetainedTail(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "events.jsonl")
|
||||
var output bytes.Buffer
|
||||
logger, first, err := NewPersistentBuffered(&output, slog.LevelInfo, 3, FormatConsole, path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for i := 1; i <= 5; i++ {
|
||||
logger.Info("request", "number", i)
|
||||
}
|
||||
if err := first.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
_, restored, err := NewPersistentBuffered(&output, slog.LevelInfo, 3, FormatConsole, path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer restored.Close()
|
||||
page := restored.Events(0, 10)
|
||||
if len(page.Events) != 3 || page.Events[0].Attributes["number"] != "3" ||
|
||||
page.Events[2].Attributes["number"] != "5" {
|
||||
t.Fatalf("restored events = %+v", page.Events)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseLevel(t *testing.T) {
|
||||
tests := map[string]slog.Level{
|
||||
"": slog.LevelInfo,
|
||||
|
||||
Reference in New Issue
Block a user