Files
Behavision/agent/pkg/spool/spool.go
Suriyakumarvijayanayagam 2b69a0be3d Anchor the gitignore patterns; a fresh clone did not compile
`spool/` is unanchored, so it matched `agent/pkg/spool/` - the durable
queue, source code - and the repository excluded it. Cloning and building
was what found it; nothing in the working tree ever would, because the
files are right there.

`data/` and `agent.json` have the same shape and are anchored too.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HViLj9gYNRtSr7YVZmW5sn
2026-09-04 12:08:01 +05:30

223 lines
5.8 KiB
Go

// Package spool is the durable event queue between the store and the server.
//
// A store's internet will drop. Footfall events produced while it is down must
// still arrive, so they are written to disk before anything tries to send
// them and removed only once the broker has confirmed delivery.
//
// One file per event, named by a monotonic sequence, acked by deletion. That
// is deliberately boring: it needs no database, no CGO, and a crash can leave
// at most one partial file, which the temp-then-rename write makes invisible.
// Footfall is a handful of events a minute, so the per-file overhead is
// irrelevant next to being able to reason about the failure modes.
package spool
import (
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"sync"
"time"
)
// Entry is one queued message.
type Entry struct {
Seq uint64 `json:"seq"`
Topic string `json:"topic"`
Created time.Time `json:"created"`
Payload json.RawMessage `json:"payload"`
}
const (
ext = ".evt"
quarantineIn = "corrupt"
)
// Spool is safe for concurrent use.
type Spool struct {
dir string
max int
mu sync.Mutex
seq uint64
dropped uint64
}
// Open prepares a spool directory. max caps the number of queued events: a
// store offline for a week must not fill its own disk, and silently growing
// without limit until it does is a worse failure than dropping the oldest
// events and saying so.
func Open(dir string, max int) (*Spool, error) {
if max <= 0 {
max = 10000
}
if err := os.MkdirAll(dir, 0o700); err != nil {
return nil, err
}
s := &Spool{dir: dir, max: max}
// Resume the sequence past anything already on disk, so a restart cannot
// reuse a number and have a later event sort before an earlier one.
names, err := s.names()
if err != nil {
return nil, err
}
if len(names) > 0 {
if n, err := seqOf(names[len(names)-1]); err == nil {
s.seq = n
}
}
return s, nil
}
// Append durably queues one event. It returns after the data is on disk, so a
// crash immediately afterwards still delivers it.
func (s *Spool) Append(topic string, payload any) error {
raw, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("spool: marshal payload: %w", err)
}
s.mu.Lock()
defer s.mu.Unlock()
s.seq++
e := Entry{Seq: s.seq, Topic: topic, Created: time.Now().UTC(), Payload: raw}
blob, err := json.Marshal(e)
if err != nil {
return err
}
// Written to a temp name and renamed: rename is atomic, so a reader never
// sees a half-written event and one bad crash cannot wedge the queue.
tmp := filepath.Join(s.dir, fmt.Sprintf(".%020d.tmp", e.Seq))
if err := os.WriteFile(tmp, blob, 0o600); err != nil {
return err
}
if err := os.Rename(tmp, s.path(e.Seq)); err != nil {
os.Remove(tmp)
return err
}
return s.trimLocked()
}
// Peek returns up to n oldest entries without removing them. They are removed
// only by Ack, after the broker confirms — anything else loses events on a
// connection that drops mid-publish.
func (s *Spool) Peek(n int) ([]Entry, error) {
s.mu.Lock()
defer s.mu.Unlock()
names, err := s.names()
if err != nil {
return nil, err
}
out := make([]Entry, 0, n)
for _, name := range names {
if len(out) >= n {
break
}
blob, err := os.ReadFile(filepath.Join(s.dir, name))
if err != nil {
if errors.Is(err, os.ErrNotExist) {
continue
}
return out, err
}
var e Entry
if err := json.Unmarshal(blob, &e); err != nil {
// A single unparseable file at the head would otherwise be retried
// forever and nothing behind it would ever send. Move it aside and
// keep going: one lost event beats a permanently stuck queue.
s.quarantineLocked(name)
continue
}
out = append(out, e)
}
return out, nil
}
// Ack removes entries the broker has confirmed.
func (s *Spool) Ack(seqs ...uint64) error {
s.mu.Lock()
defer s.mu.Unlock()
for _, seq := range seqs {
if err := os.Remove(s.path(seq)); err != nil && !errors.Is(err, os.ErrNotExist) {
return err
}
}
return nil
}
// Len is the number of events waiting to be delivered.
func (s *Spool) Len() int {
s.mu.Lock()
defer s.mu.Unlock()
names, _ := s.names()
return len(names)
}
// Dropped counts events discarded because the queue was full. Non-zero means
// the site lost footfall, which must be reported rather than inferred.
func (s *Spool) Dropped() uint64 {
s.mu.Lock()
defer s.mu.Unlock()
return s.dropped
}
// -- internals -------------------------------------------------------------
func (s *Spool) path(seq uint64) string {
return filepath.Join(s.dir, fmt.Sprintf("%020d%s", seq, ext))
}
// names returns queued filenames in delivery order. Fixed-width zero-padded
// sequence numbers make lexical order the same as numeric order, so this needs
// no parsing and cannot disagree with itself past 2^64.
func (s *Spool) names() ([]string, error) {
entries, err := os.ReadDir(s.dir)
if err != nil {
return nil, err
}
names := make([]string, 0, len(entries))
for _, e := range entries {
if e.IsDir() || !strings.HasSuffix(e.Name(), ext) {
continue
}
names = append(names, e.Name())
}
sort.Strings(names)
return names, nil
}
func (s *Spool) trimLocked() error {
names, err := s.names()
if err != nil {
return err
}
for len(names) > s.max {
if err := os.Remove(filepath.Join(s.dir, names[0])); err != nil &&
!errors.Is(err, os.ErrNotExist) {
return err
}
s.dropped++
names = names[1:]
}
return nil
}
func (s *Spool) quarantineLocked(name string) {
dir := filepath.Join(s.dir, quarantineIn)
if err := os.MkdirAll(dir, 0o700); err != nil {
os.Remove(filepath.Join(s.dir, name))
return
}
if err := os.Rename(filepath.Join(s.dir, name),
filepath.Join(dir, name)); err != nil {
os.Remove(filepath.Join(s.dir, name))
}
}
func seqOf(name string) (uint64, error) {
return strconv.ParseUint(strings.TrimSuffix(name, ext), 10, 64)
}