diff --git a/.gitignore b/.gitignore index 2100914..ddac4d5 100644 --- a/.gitignore +++ b/.gitignore @@ -2,38 +2,53 @@ __pycache__/ *.pyc .venv/ venv/ -.env -data/ -models/*.onnx -models/*.caffemodel -models/*.prototxt -*.log .pytest_cache/ +*.log +# macOS Finder metadata and rotated engine logs. +.DS_Store +*.log.[0-9] +node_modules/ + +# -------------------------------------------------------------------------- +# Everything below is ANCHORED with a leading slash on purpose. +# +# An unanchored pattern matches at every level, and `spool/` therefore excluded +# `agent/pkg/spool/` - the durable queue, source code - so a fresh clone did not +# compile at all. Found by cloning the repository and building it, which is the +# only way this class of mistake is ever found. `data/` and `agent.json` have +# exactly the same shape and are anchored for the same reason. +# -------------------------------------------------------------------------- + +# Camera credentials. Copied between machines by hand, never committed. +/.env + +# The engine's writable state: the biometric database, the camera list, logs. +/data/ + +# Downloaded models, ~200 MB, fetched by `setup-models`. +/models/*.onnx +/models/*.caffemodel +/models/*.prototxt # run-local.sh's working directory: the built binary, the encryption key and # the broker's password file. Nothing here belongs in a repository. -.local/ +/.local/ # The agent's local state when BEHAVISION_DATA_DIR points at a checkout. # agent.json holds this PC's broker password and its API token. -agent.json -spool/ +/agent.json +/spool/ # Build output. The Windows package is ~400 MB unpacked and is rebuilt from # source by installer/build.ps1; the WebView2 bootstrapper is Microsoft's # redistributable, fetched at build time rather than vendored into history. /dist/ /build/ -desktop/build/bin/ -installer/vendor/ -node_modules/ +/desktop/build/bin/ +/installer/vendor/ # NOT ignored, deliberately: server/internal/web/dist and -# desktop/frontend/dist. Both are `go:embed`ed at COMPILE time, so without -# them in the tree `go build ./...` fails on a fresh checkout - on a machine -# that may have no npm at all. They are ~200 KB and regenerating them is one -# command; a repository that does not compile is the more expensive problem. - -# macOS finder metadata and rotated engine logs. -.DS_Store -*.log.[0-9] +# desktop/frontend/dist. Both are `go:embed`ed at COMPILE time, so without them +# in the tree `go build ./...` fails on a fresh checkout - on a machine that may +# have no npm at all. They are ~200 KB and regenerating them is one command; a +# repository that does not compile is the more expensive problem. diff --git a/agent/pkg/spool/spool.go b/agent/pkg/spool/spool.go new file mode 100644 index 0000000..54c3c65 --- /dev/null +++ b/agent/pkg/spool/spool.go @@ -0,0 +1,222 @@ +// 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) +} diff --git a/agent/pkg/spool/spool_test.go b/agent/pkg/spool/spool_test.go new file mode 100644 index 0000000..f9272bf --- /dev/null +++ b/agent/pkg/spool/spool_test.go @@ -0,0 +1,181 @@ +package spool + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "sync" + "testing" +) + +type visit struct { + Person string `json:"person"` +} + +func open(t *testing.T, max int) (*Spool, string) { + t.Helper() + dir := t.TempDir() + s, err := Open(dir, max) + if err != nil { + t.Fatalf("open: %v", err) + } + return s, dir +} + +func TestEventsSurviveARestart(t *testing.T) { + // The whole point: the store's internet drops, the PC reboots, and the + // footfall that happened in between still arrives. + s, dir := open(t, 100) + for i := 0; i < 3; i++ { + if err := s.Append("visit", visit{Person: fmt.Sprint(i)}); err != nil { + t.Fatal(err) + } + } + reopened, err := Open(dir, 100) + if err != nil { + t.Fatal(err) + } + if got := reopened.Len(); got != 3 { + t.Fatalf("lost events across restart: %d of 3", got) + } + got, err := reopened.Peek(10) + if err != nil { + t.Fatal(err) + } + for i, e := range got { + var v visit + json.Unmarshal(e.Payload, &v) + if v.Person != fmt.Sprint(i) { + t.Fatalf("out of order at %d: %+v", i, got) + } + } +} + +func TestSequenceDoesNotRestartAfterReopen(t *testing.T) { + // A reused sequence would sort a new event ahead of an older undelivered + // one, silently reordering the customer's footfall. + s, dir := open(t, 100) + s.Append("visit", visit{"a"}) + reopened, _ := Open(dir, 100) + reopened.Append("visit", visit{"b"}) + got, _ := reopened.Peek(10) + if len(got) != 2 || got[0].Seq >= got[1].Seq { + t.Fatalf("sequence collided across restart: %+v", got) + } +} + +func TestAckIsWhatRemovesAnEvent(t *testing.T) { + // Peek must not consume: a connection that drops mid-publish would lose + // everything it had handed out. + s, _ := open(t, 100) + s.Append("visit", visit{"a"}) + if _, err := s.Peek(10); err != nil { + t.Fatal(err) + } + if s.Len() != 1 { + t.Fatal("peek consumed the event") + } + got, _ := s.Peek(1) + s.Ack(got[0].Seq) + if s.Len() != 0 { + t.Fatal("ack did not remove the event") + } +} + +func TestAckingTwiceIsNotAnError(t *testing.T) { + // A redelivered broker confirmation must not take the pump down. + s, _ := open(t, 100) + s.Append("visit", visit{"a"}) + got, _ := s.Peek(1) + if err := s.Ack(got[0].Seq); err != nil { + t.Fatal(err) + } + if err := s.Ack(got[0].Seq); err != nil { + t.Fatalf("second ack errored: %v", err) + } +} + +func TestQueueIsBoundedAndSaysWhatItDropped(t *testing.T) { + // A store offline for a week must not fill its own disk. Dropping is + // acceptable; dropping silently is the same class of bug as a headcount + // that is wrong in a way nobody can detect. + s, _ := open(t, 3) + for i := 0; i < 6; i++ { + s.Append("visit", visit{Person: fmt.Sprint(i)}) + } + if s.Len() != 3 { + t.Fatalf("queue grew past its cap: %d", s.Len()) + } + if s.Dropped() != 3 { + t.Fatalf("dropped %d, want 3", s.Dropped()) + } + // The OLDEST go, not the newest: recent footfall is the useful part. + got, _ := s.Peek(10) + var first visit + json.Unmarshal(got[0].Payload, &first) + if first.Person != "3" { + t.Fatalf("dropped the wrong end: head is %q", first.Person) + } +} + +func TestACorruptEntryDoesNotWedgeTheQueue(t *testing.T) { + // One unparseable file at the head would otherwise be retried forever and + // nothing behind it would ever send again. + s, dir := open(t, 100) + s.Append("visit", visit{"good"}) + head, _ := s.Peek(1) + if err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%020d.evt", head[0].Seq)), + []byte("{ truncated"), 0o600); err != nil { + t.Fatal(err) + } + s.Append("visit", visit{"later"}) + + got, err := s.Peek(10) + if err != nil { + t.Fatalf("peek gave up on a corrupt entry: %v", err) + } + if len(got) != 1 { + t.Fatalf("want the surviving event, got %d", len(got)) + } + var v visit + json.Unmarshal(got[0].Payload, &v) + if v.Person != "later" { + t.Fatalf("wrong survivor: %+v", v) + } + // Kept for diagnosis rather than deleted - it is evidence of a real bug. + if _, err := os.Stat(filepath.Join(dir, quarantineIn)); err != nil { + t.Fatal("corrupt entry was not quarantined") + } +} + +func TestPartialWritesAreNeverVisible(t *testing.T) { + // Append writes to a dot-prefixed temp name and renames, so a crash + // mid-write leaves a file the reader never picks up. + s, dir := open(t, 100) + if err := os.WriteFile(filepath.Join(dir, ".00000000000000000099.tmp"), + []byte("half"), 0o600); err != nil { + t.Fatal(err) + } + s.Append("visit", visit{"a"}) + if s.Len() != 1 { + t.Fatalf("a leftover temp file was counted: %d", s.Len()) + } +} + +func TestConcurrentAppendsDoNotCollide(t *testing.T) { + // The engine publishes from its own goroutine while the pump acks. + s, _ := open(t, 1000) + var wg sync.WaitGroup + for i := 0; i < 50; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + s.Append("visit", visit{Person: fmt.Sprint(i)}) + }(i) + } + wg.Wait() + if s.Len() != 50 { + t.Fatalf("lost events under concurrency: %d of 50", s.Len()) + } +}