Files
wild-central/internal/natsbus/server_test.go
2026-07-10 20:46:22 +00:00

152 lines
3.2 KiB
Go

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())
}
}