package main import ( "context" "testing" "time" "github.com/alicebob/miniredis/v2" "github.com/redis/go-redis/v9" ) func newTestRedisBackend(t *testing.T, bufferSize int) (*RedisBackend, *miniredis.Miniredis) { t.Helper() mr := miniredis.RunT(t) client := redis.NewClient(&redis.Options{Addr: mr.Addr()}) return newRedisBackendFromClient(client, bufferSize), mr } func TestRedisBackend_publishAndSubscribe(t *testing.T) { backend, _ := newTestRedisBackend(t, 100) defer backend.Close() ctx, cancel := context.WithCancel(context.Background()) defer cancel() ch := backend.Subscribe(ctx) time.Sleep(50 * time.Millisecond) backend.Publish(&Event{ID: "e1", Path: "test", Payload: map[string]any{"ok": true}}) select { case got := <-ch: if got.ID != "e1" { t.Errorf("expected e1, got %s", got.ID) } if got.Path != "test" { t.Errorf("expected path test, got %s", got.Path) } case <-time.After(2 * time.Second): t.Fatal("timed out waiting for event") } } func TestRedisBackend_since(t *testing.T) { backend, _ := newTestRedisBackend(t, 100) defer backend.Close() backend.Publish(&Event{ID: "e1", Path: "test", Payload: map[string]any{}}) backend.Publish(&Event{ID: "e2", Path: "test", Payload: map[string]any{}}) backend.Publish(&Event{ID: "e3", Path: "other", Payload: map[string]any{}}) events := backend.Since("e1", "test") if len(events) != 1 { t.Fatalf("expected 1 event, got %d", len(events)) } if events[0].ID != "e2" { t.Errorf("expected e2, got %s", events[0].ID) } } func TestRedisBackend_sinceExpired(t *testing.T) { backend, _ := newTestRedisBackend(t, 3) defer backend.Close() backend.Publish(&Event{ID: "e1", Path: "test", Payload: map[string]any{}}) backend.Publish(&Event{ID: "e2", Path: "test", Payload: map[string]any{}}) backend.Publish(&Event{ID: "e3", Path: "test", Payload: map[string]any{}}) backend.Publish(&Event{ID: "e4", Path: "test", Payload: map[string]any{}}) events := backend.Since("e1", "test") if len(events) != 3 { t.Fatalf("expected 3 events (full buffer), got %d", len(events)) } if events[0].ID != "e2" { t.Errorf("expected e2, got %s", events[0].ID) } } func TestRedisBackend_sincePathFiltering(t *testing.T) { backend, _ := newTestRedisBackend(t, 100) defer backend.Close() backend.Publish(&Event{ID: "e1", Path: "a/b", Payload: map[string]any{}}) backend.Publish(&Event{ID: "e2", Path: "a/b/c", Payload: map[string]any{}}) backend.Publish(&Event{ID: "e3", Path: "x/y", Payload: map[string]any{}}) events := backend.Since("e1", "a/b") if len(events) != 1 { t.Fatalf("expected 1 event, got %d", len(events)) } if events[0].ID != "e2" { t.Errorf("expected e2, got %s", events[0].ID) } } func TestRedisBackend_multiReplica(t *testing.T) { mr := miniredis.RunT(t) client1 := redis.NewClient(&redis.Options{Addr: mr.Addr()}) backend1 := newRedisBackendFromClient(client1, 100) defer backend1.Close() client2 := redis.NewClient(&redis.Options{Addr: mr.Addr()}) backend2 := newRedisBackendFromClient(client2, 100) defer backend2.Close() ctx, cancel := context.WithCancel(context.Background()) defer cancel() ch2 := backend2.Subscribe(ctx) time.Sleep(50 * time.Millisecond) backend1.Publish(&Event{ID: "cross-replica", Path: "test", Payload: map[string]any{"from": "replica1"}}) select { case got := <-ch2: if got.ID != "cross-replica" { t.Errorf("expected cross-replica, got %s", got.ID) } case <-time.After(2 * time.Second): t.Fatal("timed out waiting for cross-replica event") } } func TestRedisBackend_multiReplicaReplay(t *testing.T) { mr := miniredis.RunT(t) client1 := redis.NewClient(&redis.Options{Addr: mr.Addr()}) backend1 := newRedisBackendFromClient(client1, 100) defer backend1.Close() backend1.Publish(&Event{ID: "e1", Path: "test", Payload: map[string]any{}}) backend1.Publish(&Event{ID: "e2", Path: "test", Payload: map[string]any{}}) client2 := redis.NewClient(&redis.Options{Addr: mr.Addr()}) backend2 := newRedisBackendFromClient(client2, 100) defer backend2.Close() events := backend2.Since("e1", "test") if len(events) != 1 { t.Fatalf("expected 1 event, got %d", len(events)) } if events[0].ID != "e2" { t.Errorf("expected e2, got %s", events[0].ID) } } func TestRedisBackend_closeIsIdempotent(t *testing.T) { backend, _ := newTestRedisBackend(t, 100) if err := backend.Close(); err != nil { t.Errorf("first close: %v", err) } if err := backend.Close(); err != nil { t.Errorf("second close: %v", err) } } func TestRedisBackend_bufferTrims(t *testing.T) { backend, _ := newTestRedisBackend(t, 3) defer backend.Close() for i := range 10 { backend.Publish(&Event{ID: string(rune('a' + i)), Path: "test", Payload: map[string]any{}}) } events := backend.Since("", "test") if len(events) != 3 { t.Fatalf("expected 3 events in trimmed buffer, got %d", len(events)) } } func TestRedisBackend_bufferHasTTL(t *testing.T) { backend, mr := newTestRedisBackend(t, 100) defer backend.Close() backend.Publish(&Event{ID: "e1", Path: "test", Payload: map[string]any{}}) ttl := mr.TTL(redisBufferKey) if ttl <= 0 { t.Fatalf("expected positive TTL on buffer key, got %v", ttl) } if ttl > 24*time.Hour { t.Fatalf("expected TTL <= 24h, got %v", ttl) } } func TestNewRedisBackend_success(t *testing.T) { mr := miniredis.RunT(t) backend, err := NewRedisBackend("redis://"+mr.Addr(), 100) if err != nil { t.Fatalf("expected no error, got %v", err) } defer backend.Close() } func TestNewRedisBackend_badURL(t *testing.T) { _, err := NewRedisBackend("not-a-url", 100) if err == nil { t.Fatal("expected error for bad URL") } } func TestNewRedisBackend_unreachable(t *testing.T) { _, err := NewRedisBackend("redis://127.0.0.1:1", 100) if err == nil { t.Fatal("expected error for unreachable redis") } } func TestRedisBackend_sinceWithCorruptData(t *testing.T) { backend, mr := newTestRedisBackend(t, 100) defer backend.Close() backend.Publish(&Event{ID: "e1", Path: "test", Payload: map[string]any{}}) mr.Lpush(redisBufferKey, "not-valid-json") backend.Publish(&Event{ID: "e2", Path: "test", Payload: map[string]any{}}) events := backend.Since("e1", "test") if len(events) != 1 { t.Fatalf("expected 1 event (corrupt data skipped), got %d", len(events)) } if events[0].ID != "e2" { t.Errorf("expected e2, got %s", events[0].ID) } } func TestRedisBackend_subscribeContextCancel(t *testing.T) { backend, _ := newTestRedisBackend(t, 100) defer backend.Close() ctx, cancel := context.WithCancel(context.Background()) ch := backend.Subscribe(ctx) cancel() // Channel should close after context cancel select { case _, ok := <-ch: if ok { t.Error("expected channel to be closed") } case <-time.After(2 * time.Second): t.Fatal("timed out waiting for channel close") } } func TestRedisBackend_subscribeSkipsBadJSON(t *testing.T) { backend, _ := newTestRedisBackend(t, 100) defer backend.Close() ctx, cancel := context.WithCancel(context.Background()) defer cancel() ch := backend.Subscribe(ctx) time.Sleep(50 * time.Millisecond) // Publish garbage directly to the Redis channel backend.client.Publish(context.Background(), redisChannel, "not-json") // Then a valid event backend.Publish(&Event{ID: "valid", Path: "test", Payload: map[string]any{}}) select { case got := <-ch: if got.ID != "valid" { t.Errorf("expected valid, got %s", got.ID) } case <-time.After(2 * time.Second): t.Fatal("timed out") } } func TestRedisBackend_sinceClosedClient(t *testing.T) { backend, _ := newTestRedisBackend(t, 100) backend.Publish(&Event{ID: "e1", Path: "test", Payload: map[string]any{}}) backend.client.Close() backend.closed = true events := backend.Since("", "test") if events != nil { t.Errorf("expected nil, got %v", events) } } func TestRedisBackend_subscribeNilMessage(t *testing.T) { mr := miniredis.RunT(t) client := redis.NewClient(&redis.Options{Addr: mr.Addr()}) backend := newRedisBackendFromClient(client, 100) ctx, cancel := context.WithCancel(context.Background()) defer cancel() ch := backend.Subscribe(ctx) time.Sleep(50 * time.Millisecond) // Close the client, which causes the pubsub channel to close client.Close() mr.Close() select { case _, ok := <-ch: if ok { t.Error("expected channel to be closed after client close") } case <-time.After(5 * time.Second): t.Fatal("timed out waiting for channel close") } } func TestRedisBackend_publishMarshalError(t *testing.T) { backend, _ := newTestRedisBackend(t, 100) defer backend.Close() err := backend.Publish(&Event{ID: "bad", Path: "test", Payload: make(chan int)}) if err == nil { t.Fatal("expected marshal error") } }