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