This repository has no description
2.3 kB
105 lines
1package eventconsumer
2
3import (
4 "context"
5 "io"
6 "log/slog"
7 "path/filepath"
8 "sync"
9 "testing"
10 "time"
11
12 "tangled.org/core/eventconsumer/cursor"
13 "tangled.org/core/eventstream"
14)
15
16func sqliteCursorStore(t *testing.T) cursor.Store {
17 t.Helper()
18 store, err := cursor.NewSQLiteStore(filepath.Join(t.TempDir(), "spindle.db"))
19 if err != nil {
20 t.Fatalf("new sqlite cursor store: %v", err)
21 }
22 return store
23}
24
25func drainProcessed(t *testing.T, store cursor.Store, source Source, expected int) []int64 {
26 t.Helper()
27
28 var mu sync.Mutex
29 var seen []int64
30
31 c := NewConsumer(ConsumerConfig{
32 ProcessFunc: func(_ context.Context, _ Source, msg eventstream.Event) error {
33 mu.Lock()
34 seen = append(seen, msg.Created)
35 mu.Unlock()
36 return nil
37 },
38 WorkerCount: 1,
39 QueueSize: 16,
40 ConnectionTimeout: 2 * time.Second,
41 CursorStore: store,
42 Logger: slog.New(slog.NewTextHandler(io.Discard, nil)),
43 })
44
45 ctx, cancel := context.WithCancel(context.Background())
46 defer cancel()
47 c.Start(ctx)
48 c.AddSource(ctx, source)
49
50 deadline := time.Now().Add(3 * time.Second)
51 for time.Now().Before(deadline) {
52 mu.Lock()
53 n := len(seen)
54 mu.Unlock()
55 if n >= expected {
56 break
57 }
58 time.Sleep(20 * time.Millisecond)
59 }
60
61 c.Stop()
62
63 mu.Lock()
64 defer mu.Unlock()
65 return append([]int64(nil), seen...)
66}
67
68func TestSpindleUpgrade_OrphanedCursorReplaysFromZero(t *testing.T) {
69 src := &memSrc{}
70 for i := range 8 {
71 src.add(mkEv(i))
72 }
73 source, _ := startEventServer(t, src)
74
75 store := sqliteCursorStore(t)
76 store.Set(source.Host, 5)
77
78 seen := drainProcessed(t, store, source, 8)
79
80 if len(seen) != 8 {
81 t.Fatalf("orphaned bare-host cursor processed %d events, want a full replay of 8: %v", len(seen), seen)
82 }
83}
84
85func TestSpindleUpgrade_MigratedCursorResumesNoReplay(t *testing.T) {
86 src := &memSrc{}
87 for i := range 8 {
88 src.add(mkEv(i))
89 }
90 source, _ := startEventServer(t, src)
91
92 store := sqliteCursorStore(t)
93 store.Set(source.Host, 5)
94
95 MigrateLegacyCursor(store, source)
96
97 seen := drainProcessed(t, store, source, 3)
98
99 if len(seen) != 3 {
100 t.Fatalf("migrated cursor processed %d events, want a resume of 3: %v", len(seen), seen)
101 }
102 if seen[0] != 6 || seen[2] != 8 {
103 t.Fatalf("resumed events = %v, want [6 7 8]", seen)
104 }
105}