feat(go): advance validation sidecar toward production
This commit is contained in:
parent
6a06eba255
commit
f85c2e831a
92 changed files with 8881 additions and 91 deletions
87
go-agent/internal/eventlog/read.go
Normal file
87
go-agent/internal/eventlog/read.go
Normal file
|
|
@ -0,0 +1,87 @@
|
|||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
package eventlog
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
)
|
||||
|
||||
const maxRecordBytes = 1 << 20
|
||||
|
||||
// ReadRecent returns at most limit records in chronological order, reading the
|
||||
// older rotated segment before the active file. Every returned row is valid
|
||||
// JSON; corruption is reported instead of passed to an operator as evidence.
|
||||
func ReadRecent(path string, limit int) ([][]byte, error) {
|
||||
if path == "" {
|
||||
return nil, fmt.Errorf("event log path is required")
|
||||
}
|
||||
if limit <= 0 {
|
||||
return nil, fmt.Errorf("event limit must be positive")
|
||||
}
|
||||
buffer := recentBuffer{limit: limit}
|
||||
if _, err := os.Stat(path + ".1"); err == nil {
|
||||
if err := scanRecords(path+".1", buffer.add); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
} else if !os.IsNotExist(err) {
|
||||
return nil, fmt.Errorf("stat rotated event log: %w", err)
|
||||
}
|
||||
if err := scanRecords(path, buffer.add); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return buffer.values(), nil
|
||||
}
|
||||
|
||||
func scanRecords(path string, consume func([]byte)) error {
|
||||
file, err := os.Open(path)
|
||||
if err != nil {
|
||||
return fmt.Errorf("open event log %s: %w", path, err)
|
||||
}
|
||||
defer file.Close()
|
||||
scanner := bufio.NewScanner(file)
|
||||
scanner.Buffer(make([]byte, 64*1024), maxRecordBytes)
|
||||
line := 0
|
||||
for scanner.Scan() {
|
||||
line++
|
||||
raw := scanner.Bytes()
|
||||
if len(raw) == 0 {
|
||||
continue
|
||||
}
|
||||
if !json.Valid(raw) {
|
||||
return fmt.Errorf("invalid JSON in %s at line %d", path, line)
|
||||
}
|
||||
consume(append([]byte(nil), raw...))
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
return fmt.Errorf("read event log %s: %w", path, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
type recentBuffer struct {
|
||||
limit int
|
||||
rows [][]byte
|
||||
next int
|
||||
}
|
||||
|
||||
func (b *recentBuffer) add(raw []byte) {
|
||||
if len(b.rows) < b.limit {
|
||||
b.rows = append(b.rows, raw)
|
||||
return
|
||||
}
|
||||
b.rows[b.next] = raw
|
||||
b.next = (b.next + 1) % b.limit
|
||||
}
|
||||
|
||||
func (b *recentBuffer) values() [][]byte {
|
||||
if len(b.rows) < b.limit || b.next == 0 {
|
||||
return b.rows
|
||||
}
|
||||
result := make([][]byte, 0, len(b.rows))
|
||||
result = append(result, b.rows[b.next:]...)
|
||||
result = append(result, b.rows[:b.next]...)
|
||||
return result
|
||||
}
|
||||
81
go-agent/internal/eventlog/read_test.go
Normal file
81
go-agent/internal/eventlog/read_test.go
Normal file
|
|
@ -0,0 +1,81 @@
|
|||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
package eventlog
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestReadRecentSpansRotationInChronologicalOrder(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "events.jsonl")
|
||||
writeLines(t, path+".1", 1, 2, 3)
|
||||
writeLines(t, path, 4, 5)
|
||||
rows, err := ReadRecent(path, 3)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := rowSequences(t, rows); len(got) != 3 || got[0] != 3 || got[1] != 4 || got[2] != 5 {
|
||||
t.Fatalf("sequences=%v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadRecentReturnsAllWhenUnderLimit(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "events.jsonl")
|
||||
writeLines(t, path, 7, 8)
|
||||
rows, err := ReadRecent(path, 10)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := rowSequences(t, rows); len(got) != 2 || got[0] != 7 || got[1] != 8 {
|
||||
t.Fatalf("sequences=%v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadRecentRejectsCorruptionAndInvalidArguments(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "events.jsonl")
|
||||
if err := os.WriteFile(path, []byte("{not-json}\n"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := ReadRecent(path, 1); err == nil {
|
||||
t.Fatal("corrupt log accepted")
|
||||
}
|
||||
if _, err := ReadRecent("", 1); err == nil {
|
||||
t.Fatal("empty path accepted")
|
||||
}
|
||||
if _, err := ReadRecent(path, 0); err == nil {
|
||||
t.Fatal("zero limit accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func writeLines(t *testing.T, path string, sequences ...int) {
|
||||
t.Helper()
|
||||
file, err := os.Create(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer file.Close()
|
||||
encoder := json.NewEncoder(file)
|
||||
for _, sequence := range sequences {
|
||||
if err := encoder.Encode(map[string]any{"sequence": sequence}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func rowSequences(t *testing.T, rows [][]byte) []int {
|
||||
t.Helper()
|
||||
result := make([]int, 0, len(rows))
|
||||
for _, row := range rows {
|
||||
var record struct {
|
||||
Sequence int `json:"sequence"`
|
||||
}
|
||||
if err := json.Unmarshal(row, &record); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
result = append(result, record.Sequence)
|
||||
}
|
||||
return result
|
||||
}
|
||||
111
go-agent/internal/eventlog/writer.go
Normal file
111
go-agent/internal/eventlog/writer.go
Normal file
|
|
@ -0,0 +1,111 @@
|
|||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
// Package eventlog retains a bounded JSONL copy of emitted event envelopes.
|
||||
package eventlog
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// Writer appends records to Path and keeps one rotated segment at Path + ".1".
|
||||
type Writer struct {
|
||||
Path string
|
||||
MaxBytes int64
|
||||
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
// New preflights the parent directory and destination permissions. A service
|
||||
// can call it before signaling readiness so a broken retention path fails the
|
||||
// start instead of silently dropping records.
|
||||
func New(path string, maxBytes int64) (*Writer, error) {
|
||||
if path == "" {
|
||||
return nil, fmt.Errorf("event log path is required")
|
||||
}
|
||||
if maxBytes <= 0 {
|
||||
return nil, fmt.Errorf("event log max bytes must be positive")
|
||||
}
|
||||
if err := os.MkdirAll(filepath.Dir(path), 0o750); err != nil {
|
||||
return nil, fmt.Errorf("create event log directory: %w", err)
|
||||
}
|
||||
file, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("open event log: %w", err)
|
||||
}
|
||||
if err := file.Chmod(0o600); err != nil {
|
||||
file.Close()
|
||||
return nil, fmt.Errorf("protect event log: %w", err)
|
||||
}
|
||||
if err := file.Close(); err != nil {
|
||||
return nil, fmt.Errorf("close event log: %w", err)
|
||||
}
|
||||
return &Writer{Path: path, MaxBytes: maxBytes}, nil
|
||||
}
|
||||
|
||||
// Append writes exactly one JSON record. Alert and incident records are synced
|
||||
// before success is reported; routine status records rely on normal kernel
|
||||
// writeback to avoid forcing a disk flush every sweep.
|
||||
func (w *Writer) Append(record map[string]any) error {
|
||||
if w == nil {
|
||||
return fmt.Errorf("event log writer is required")
|
||||
}
|
||||
raw, err := json.Marshal(record)
|
||||
if err != nil {
|
||||
return fmt.Errorf("encode event log record: %w", err)
|
||||
}
|
||||
raw = append(raw, '\n')
|
||||
if len(raw) > maxRecordBytes {
|
||||
return fmt.Errorf("event log record exceeds %d bytes", maxRecordBytes)
|
||||
}
|
||||
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
if err := w.rotateBefore(int64(len(raw))); err != nil {
|
||||
return err
|
||||
}
|
||||
file, err := os.OpenFile(w.Path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600)
|
||||
if err != nil {
|
||||
return fmt.Errorf("open event log: %w", err)
|
||||
}
|
||||
if err := file.Chmod(0o600); err != nil {
|
||||
file.Close()
|
||||
return fmt.Errorf("protect event log: %w", err)
|
||||
}
|
||||
if _, err := file.Write(raw); err != nil {
|
||||
file.Close()
|
||||
return fmt.Errorf("append event log: %w", err)
|
||||
}
|
||||
if eventType, _ := record["event_type"].(string); eventType != "status" {
|
||||
if err := file.Sync(); err != nil {
|
||||
file.Close()
|
||||
return fmt.Errorf("sync event log: %w", err)
|
||||
}
|
||||
}
|
||||
if err := file.Close(); err != nil {
|
||||
return fmt.Errorf("close event log: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (w *Writer) rotateBefore(incoming int64) error {
|
||||
info, err := os.Stat(w.Path)
|
||||
if os.IsNotExist(err) {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("stat event log: %w", err)
|
||||
}
|
||||
// Always allow one record into an empty file, even when that individual
|
||||
// envelope is larger than the configured bound.
|
||||
if info.Size() == 0 || info.Size()+incoming <= w.MaxBytes {
|
||||
return nil
|
||||
}
|
||||
if err := os.Rename(w.Path, w.Path+".1"); err != nil {
|
||||
return fmt.Errorf("rotate event log: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
131
go-agent/internal/eventlog/writer_test.go
Normal file
131
go-agent/internal/eventlog/writer_test.go
Normal file
|
|
@ -0,0 +1,131 @@
|
|||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
package eventlog
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestWriterAppendsJSONLWithPrivatePermissions(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "state", "events.jsonl")
|
||||
writer, err := New(path, 1<<20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, record := range []map[string]any{
|
||||
{"event_type": "status", "sequence": 1},
|
||||
{"event_type": "alert", "sequence": 2},
|
||||
} {
|
||||
if err := writer.Append(record); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
rows := readRows(t, path)
|
||||
if len(rows) != 2 || rows[0]["sequence"] != float64(1) || rows[1]["sequence"] != float64(2) {
|
||||
t.Fatalf("rows=%#v", rows)
|
||||
}
|
||||
info, err := os.Stat(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := info.Mode().Perm(); got != 0o600 {
|
||||
t.Fatalf("mode=%#o", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriterRotatesBeforeCrossingBound(t *testing.T) {
|
||||
directory := t.TempDir()
|
||||
path := filepath.Join(directory, "events.jsonl")
|
||||
first := map[string]any{"event_type": "status", "marker": "first"}
|
||||
encoded, err := json.Marshal(first)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
writer, err := New(path, int64(len(encoded)+1))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := writer.Append(first); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := writer.Append(map[string]any{"event_type": "alert", "marker": "second"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
current := readRows(t, path)
|
||||
rotated := readRows(t, path+".1")
|
||||
if len(current) != 1 || current[0]["marker"] != "second" {
|
||||
t.Fatalf("current=%#v", current)
|
||||
}
|
||||
if len(rotated) != 1 || rotated[0]["marker"] != "first" {
|
||||
t.Fatalf("rotated=%#v", rotated)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriterSerializesConcurrentAppends(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "events.jsonl")
|
||||
writer, err := New(path, 1<<20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var wait sync.WaitGroup
|
||||
for index := range 32 {
|
||||
wait.Add(1)
|
||||
go func() {
|
||||
defer wait.Done()
|
||||
if err := writer.Append(map[string]any{"event_type": "status", "sequence": index}); err != nil {
|
||||
t.Errorf("append: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
wait.Wait()
|
||||
if rows := readRows(t, path); len(rows) != 32 {
|
||||
t.Fatalf("row count=%d", len(rows))
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewRejectsInvalidConfiguration(t *testing.T) {
|
||||
if _, err := New("", 100); err == nil {
|
||||
t.Fatal("empty path accepted")
|
||||
}
|
||||
if _, err := New(filepath.Join(t.TempDir(), "events.jsonl"), 0); err == nil {
|
||||
t.Fatal("zero bound accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriterRejectsRecordLargerThanReaderLimit(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "events.jsonl")
|
||||
writer, err := New(path, 2<<20)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := writer.Append(map[string]any{"event_type": "alert", "payload": make([]byte, maxRecordBytes)}); err == nil {
|
||||
t.Fatal("oversized record accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func readRows(t *testing.T, path string) []map[string]any {
|
||||
t.Helper()
|
||||
file, err := os.Open(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer file.Close()
|
||||
var rows []map[string]any
|
||||
scanner := bufio.NewScanner(file)
|
||||
for scanner.Scan() {
|
||||
var row map[string]any
|
||||
if err := json.Unmarshal(scanner.Bytes(), &row); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rows = append(rows, row)
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return rows
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue