package natsbus import ( "context" "testing" "time" "github.com/nats-io/nats.go/jetstream" ) func TestStartAndShutdown(t *testing.T) { s, err := Start(Config{ Port: -1, // random available port DataDir: t.TempDir(), }) if err != nil { t.Fatalf("Start failed: %v", err) } defer s.Shutdown() if s.ClientURL() == "" { t.Error("expected non-empty client URL") } } func TestBucketsCreated(t *testing.T) { s, err := Start(Config{ Port: -1, DataDir: t.TempDir(), }) if err != nil { t.Fatalf("Start failed: %v", err) } defer s.Shutdown() ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() // Verify domains bucket exists kv, err := s.js.KeyValue(ctx, BucketDomains) if err != nil { t.Fatalf("domains bucket not found: %v", err) } // Write and read a value _, err = kv.Put(ctx, "test/svc1", []byte(`{"name":"svc1"}`)) if err != nil { t.Fatalf("Put failed: %v", err) } entry, err := kv.Get(ctx, "test/svc1") if err != nil { t.Fatalf("Get failed: %v", err) } if string(entry.Value()) != `{"name":"svc1"}` { t.Errorf("unexpected value: %s", entry.Value()) } // Verify presence bucket exists with TTL pKV, err := s.js.KeyValue(ctx, BucketPresence) if err != nil { t.Fatalf("presence bucket not found: %v", err) } status, err := pKV.Status(ctx) if err != nil { t.Fatalf("presence status failed: %v", err) } if status.TTL() != PresenceTTL { t.Errorf("expected TTL %v, got %v", PresenceTTL, status.TTL()) } } func TestEventsStream(t *testing.T) { s, err := Start(Config{ Port: -1, DataDir: t.TempDir(), }) if err != nil { t.Fatalf("Start failed: %v", err) } defer s.Shutdown() ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() // Publish an event _, err = s.js.Publish(ctx, "wild.events.test.service-started", []byte(`{"service":"my-api"}`)) if err != nil { t.Fatalf("Publish failed: %v", err) } // Verify stream has the message stream, err := s.js.Stream(ctx, StreamEvents) if err != nil { t.Fatalf("Stream not found: %v", err) } info, err := stream.Info(ctx) if err != nil { t.Fatalf("Stream info failed: %v", err) } if info.State.Msgs != 1 { t.Errorf("expected 1 message in stream, got %d", info.State.Msgs) } } func TestPresenceTTL(t *testing.T) { tmpDir := t.TempDir() // Use a very short TTL for testing origTTL := PresenceTTL // We can't modify the const, but we can verify the bucket was created with the right TTL _ = origTTL s, err := Start(Config{ Port: -1, DataDir: tmpDir, }) if err != nil { t.Fatalf("Start failed: %v", err) } defer s.Shutdown() ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() kv, err := s.js.KeyValue(ctx, BucketPresence) if err != nil { t.Fatalf("presence bucket not found: %v", err) } // Write a presence key _, err = kv.Put(ctx, "nodes/test-node", []byte(`{"ip":"192.168.8.60"}`)) if err != nil { t.Fatalf("Put failed: %v", err) } // Verify it exists entry, err := kv.Get(ctx, "nodes/test-node") if err != nil { t.Fatalf("Get failed: %v", err) } if entry.Operation() != jetstream.KeyValuePut { t.Errorf("expected Put operation, got %v", entry.Operation()) } }