`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
182 lines
4.7 KiB
Go
182 lines
4.7 KiB
Go
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())
|
|
}
|
|
}
|