`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
223 lines
5.8 KiB
Go
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)
|
|
}
|