This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / eventconsumer / upgrade_test.go
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}