Monorepo for Tangled
Something went wrong. Try again.
6.2 kB · 278 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279package eventconsumer
import ( "context" "encoding/json" "fmt" "io" "log/slog" "net/http" "net/http/httptest" "strings" "sync" "testing" "time"
"tangled.org/core/eventconsumer/cursor" "tangled.org/core/eventstream" "tangled.org/core/notifier")
type memSrc struct { mu sync.Mutex events []eventstream.Event}
func (s *memSrc) add(ev eventstream.Event) { s.mu.Lock() defer s.mu.Unlock() s.events = append(s.events, ev)}
func (s *memSrc) GetEvents(cursor int64, limit int) ([]eventstream.Event, error) { s.mu.Lock() defer s.mu.Unlock() out := []eventstream.Event{} for _, ev := range s.events { if ev.Created > cursor { out = append(out, ev) if len(out) == limit { break } } } return out, nil}
func mkEv(i int) eventstream.Event { return eventstream.Event{ Rkey: fmt.Sprintf("rk-%04d", i), Nsid: "sh.tangled.test", EventJson: json.RawMessage(fmt.Sprintf(`{"i":%d}`, i)), Created: int64(i + 1), }}
func startEventServer(t *testing.T, src *memSrc) (Source, *notifier.Notifier) { t.Helper() n := notifier.New() mux := http.NewServeMux() mux.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) { _ = eventstream.Stream(w, r, eventstream.StreamConfig{ Backend: src, Notifier: &n, Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), BatchSize: 5, MaxBatchesPerDrain: 100, }) }) srv := httptest.NewServer(mux) t.Cleanup(srv.Close) addr := strings.TrimPrefix(srv.URL, "http://") return Source{Kind: "test", Host: addr, NoTLS: true}, &n}
func TestConsumer_DrainAdvancesCursor(t *testing.T) { src := &memSrc{} for i := range 8 { src.add(mkEv(i)) }
source, _ := startEventServer(t, src)
store := &cursor.MemoryStore{} seenMu := sync.Mutex{} seen := []int64{}
cfg := ConsumerConfig{ ProcessFunc: func(ctx context.Context, _ Source, msg eventstream.Event) error { seenMu.Lock() seen = append(seen, msg.Created) seenMu.Unlock() return nil }, WorkerCount: 1, QueueSize: 16, ConnectionTimeout: 2 * time.Second, CursorStore: store, Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), } c := NewConsumer(cfg)
ctx, cancel := context.WithCancel(context.Background()) defer cancel()
c.Start(ctx) c.AddSource(ctx, source)
deadline := time.Now().Add(3 * time.Second) for time.Now().Before(deadline) { seenMu.Lock() n := len(seen) seenMu.Unlock() if n >= 8 { break } time.Sleep(20 * time.Millisecond) }
seenMu.Lock() defer seenMu.Unlock() if len(seen) != 8 { t.Fatalf("processed %d events, want 8: %v", len(seen), seen) } for i, got := range seen { if got != int64(i+1) { t.Fatalf("event %d: got created=%d want %d", i, got, i+1) } }
if final := store.Get(source.Key()); final != 8 { t.Fatalf("cursor = %d, want 8", final) }}
func TestConsumer_CursorMonotonic_OutOfOrderWorkers(t *testing.T) { src := &memSrc{} for i := range 4 { src.add(mkEv(i)) }
source, _ := startEventServer(t, src)
store := &cursor.MemoryStore{}
releaseFirst := make(chan struct{}) processed := make(chan int64, 4)
cfg := ConsumerConfig{ ProcessFunc: func(ctx context.Context, _ Source, msg eventstream.Event) error { if msg.Created == 1 { <-releaseFirst } processed <- msg.Created return nil }, WorkerCount: 4, QueueSize: 16, ConnectionTimeout: 2 * time.Second, CursorStore: store, Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), } c := NewConsumer(cfg)
ctx, cancel := context.WithCancel(context.Background()) defer cancel()
c.Start(ctx) c.AddSource(ctx, source)
for range 3 { select { case <-processed: case <-time.After(3 * time.Second): t.Fatal("timed out waiting for events 2-4 to be processed") } }
if cur := store.Get(source.Key()); cur != 4 { t.Fatalf("cursor before slow worker finished = %d, want 4", cur) }
close(releaseFirst) select { case <-processed: case <-time.After(3 * time.Second): t.Fatal("timed out waiting for slow worker") }
if cur := store.Get(source.Key()); cur != 4 { t.Fatalf("cursor regressed after slow worker: %d, want 4", cur) }}
func TestConsumer_StopTerminatesWithoutCtxCancel(t *testing.T) { src := &memSrc{} source, _ := startEventServer(t, src)
cfg := ConsumerConfig{ ProcessFunc: func(ctx context.Context, _ Source, _ eventstream.Event) error { return nil }, WorkerCount: 2, QueueSize: 8, ConnectionTimeout: 2 * time.Second, CursorStore: &cursor.MemoryStore{}, Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), } c := NewConsumer(cfg)
c.Start(context.Background()) c.AddSource(context.Background(), source)
done := make(chan struct{}) go func() { c.Stop() close(done) }()
select { case <-done: case <-time.After(5 * time.Second): t.Fatal("Stop did not return within 5s") }}
func TestConsumer_ResumesFromStoredCursor(t *testing.T) { src := &memSrc{} for i := range 5 { src.add(mkEv(i)) }
source, _ := startEventServer(t, src)
store := &cursor.MemoryStore{} store.Set(source.Key(), 3)
seenMu := sync.Mutex{} seen := []int64{}
cfg := ConsumerConfig{ ProcessFunc: func(ctx context.Context, _ Source, msg eventstream.Event) error { seenMu.Lock() seen = append(seen, msg.Created) seenMu.Unlock() return nil }, WorkerCount: 1, QueueSize: 16, ConnectionTimeout: 2 * time.Second, CursorStore: store, Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), } c := NewConsumer(cfg)
ctx, cancel := context.WithCancel(context.Background()) defer cancel()
c.Start(ctx) c.AddSource(ctx, source)
deadline := time.Now().Add(3 * time.Second) for time.Now().Before(deadline) { seenMu.Lock() n := len(seen) seenMu.Unlock() if n >= 2 { break } time.Sleep(20 * time.Millisecond) }
seenMu.Lock() defer seenMu.Unlock() if len(seen) < 2 { t.Fatalf("processed %d events, want 2: %v", len(seen), seen) } if seen[0] != 4 || seen[1] != 5 { t.Fatalf("resumed events = %v, want [4 5]", seen) }}