// 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) }