package audit import ( "bufio" "crypto/sha256" "encoding/hex" "encoding/json" "fmt" "os" "path/filepath" "sync" "github.com/mathiasbq/hyperguild/ingestion/internal/capture" ) // FileBuffer is a durable, restart-surviving audit buffer backed by a // JSONL file: one {id, entry} record per line. It is the internal/public // tier fallback when loki is unreachable. Confirm rewrites the file // without the confirmed record, so a record is cleared only after its // central write is confirmed. // // Access is serialised by a mutex; the buffer is low-throughput (only // written during a loki outage), so a whole-file rewrite on Confirm is // acceptable and keeps the on-disk format trivially correct. type FileBuffer struct { path string mu sync.Mutex } type bufferLine struct { ID string `json:"id"` Entry capture.AuditEntry `json:"entry"` } // NewFileBuffer returns a buffer backed by path. The parent directory is // created if needed. The file itself is created lazily on first Append. func NewFileBuffer(path string) (*FileBuffer, error) { if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { return nil, fmt.Errorf("create buffer dir: %w", err) } return &FileBuffer{path: path}, nil } // Writable reports whether the buffer file can be appended to. It probes // by opening the file for append (creating it if absent) — the same // operation Append performs — so Reserve's check matches Append's reality. func (b *FileBuffer) Writable() error { b.mu.Lock() defer b.mu.Unlock() f, err := os.OpenFile(b.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o644) if err != nil { return err } return f.Close() } // Append durably writes one audit record. The ID is derived from the // content + timestamp so it is stable and unique per record. func (b *FileBuffer) Append(e capture.AuditEntry) error { b.mu.Lock() defer b.mu.Unlock() line := bufferLine{ID: recordID(e), Entry: e} data, err := json.Marshal(line) if err != nil { return fmt.Errorf("marshal buffer line: %w", err) } f, err := os.OpenFile(b.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o644) if err != nil { return err } defer func() { _ = f.Close() }() if _, err := f.Write(append(data, '\n')); err != nil { return err } return f.Sync() } // Pending reads all buffered records. A missing file means none. func (b *FileBuffer) Pending() ([]Buffered, error) { b.mu.Lock() defer b.mu.Unlock() return b.readAllLocked() } func (b *FileBuffer) readAllLocked() ([]Buffered, error) { f, err := os.Open(b.path) if err != nil { if os.IsNotExist(err) { return nil, nil } return nil, err } defer func() { _ = f.Close() }() var out []Buffered sc := bufio.NewScanner(f) sc.Buffer(make([]byte, 0, 64*1024), 1024*1024) for sc.Scan() { raw := sc.Bytes() if len(raw) == 0 { continue } var l bufferLine if err := json.Unmarshal(raw, &l); err != nil { return nil, fmt.Errorf("parse buffer line: %w", err) } out = append(out, Buffered(l)) } return out, sc.Err() } // Confirm removes a single record after its central write is confirmed, by // rewriting the file without it. Unknown IDs are a no-op. func (b *FileBuffer) Confirm(id string) error { b.mu.Lock() defer b.mu.Unlock() all, err := b.readAllLocked() if err != nil { return err } tmp := b.path + ".tmp" f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o644) if err != nil { return err } w := bufio.NewWriter(f) kept := 0 for _, rec := range all { if rec.ID == id { continue } data, _ := json.Marshal(bufferLine(rec)) if _, err := w.Write(append(data, '\n')); err != nil { _ = f.Close() return err } kept++ } if err := w.Flush(); err != nil { _ = f.Close() return err } if err := f.Sync(); err != nil { _ = f.Close() return err } if err := f.Close(); err != nil { return err } // Empty buffer → remove the file entirely so Pending sees nothing. if kept == 0 { _ = os.Remove(tmp) return os.Remove(b.path) } return os.Rename(tmp, b.path) } // recordID is a stable per-record identifier: sha256 of the principal, // timestamp, and item list. Distinct captures never collide; the same // buffered record always hashes the same. func recordID(e capture.AuditEntry) string { h := sha256.New() _, _ = fmt.Fprintf(h, "%s|%s|%v|%s", e.Principal, e.Timestamp.UTC().Format("2006-01-02T15:04:05.000000000Z07:00"), e.Items, e.SessionRef) return hex.EncodeToString(h.Sum(nil))[:16] }