diff --git a/go.mod b/go.mod index 022b3b1..abf24ee 100644 --- a/go.mod +++ b/go.mod @@ -5,6 +5,20 @@ go 1.25.0 require ( github.com/google/uuid v1.6.0 github.com/gorilla/mux v1.8.1 + github.com/nats-io/nats-server/v2 v2.14.3 + github.com/nats-io/nats.go v1.52.0 golang.org/x/time v0.15.0 gopkg.in/yaml.v3 v3.0.1 ) + +require ( + github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op // indirect + github.com/google/go-tpm v0.9.8 // indirect + github.com/klauspost/compress v1.18.6 // indirect + github.com/minio/highwayhash v1.0.4 // indirect + github.com/nats-io/jwt/v2 v2.8.2 // indirect + github.com/nats-io/nkeys v0.4.16 // indirect + github.com/nats-io/nuid v1.0.1 // indirect + golang.org/x/crypto v0.53.0 // indirect + golang.org/x/sys v0.46.0 // indirect +) diff --git a/go.sum b/go.sum index 01b158b..9cfdd24 100644 --- a/go.sum +++ b/go.sum @@ -1,7 +1,30 @@ +github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op h1:Z/MZK75wC/NSrkgqeNIa7jexam9uWzhLmFTSCPI/kn0= +github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op/go.mod h1:FQyySiasQQM8735Ddel3MRojmy4dA1IqCeyJ5jmPMbI= +github.com/google/go-tpm v0.9.8 h1:slArAR9Ft+1ybZu0lBwpSmpwhRXaa85hWtMinMyRAWo= +github.com/google/go-tpm v0.9.8/go.mod h1:h9jEsEECg7gtLis0upRBQU+GhYVH6jMjrFxI8u6bVUY= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/gorilla/mux v1.8.1 h1:TuBL49tXwgrFYWhqrNgrUNEY92u81SPhu7sTdzQEiWY= github.com/gorilla/mux v1.8.1/go.mod h1:AKf9I4AEqPTmMytcMc0KkNouC66V3BtZ4qD5fmWSiMQ= +github.com/klauspost/compress v1.18.6 h1:2jupLlAwFm95+YDR+NwD2MEfFO9d4z4Prjl1XXDjuao= +github.com/klauspost/compress v1.18.6/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/minio/highwayhash v1.0.4 h1:asJizugGgchQod2ja9NJlGOWq4s7KsAWr5XUc9Clgl4= +github.com/minio/highwayhash v1.0.4/go.mod h1:GGYsuwP/fPD6Y9hMiXuapVvlIUEhFhMTh0rxU3ik1LQ= +github.com/nats-io/jwt/v2 v2.8.2 h1:XXRgB60MSTnqsRwejQurVDs/hcv2dkt+86GjI+I/bMc= +github.com/nats-io/jwt/v2 v2.8.2/go.mod h1:Ag/56sq9OblL4JgdYufDd16Egb17Kr/8WwwuO/forVc= +github.com/nats-io/nats-server/v2 v2.14.3 h1:+xjydPt7rkit67G+04TN0mcO2n+8nveZE7tK/PPV53A= +github.com/nats-io/nats-server/v2 v2.14.3/go.mod h1:5IlCtBzfwyzQzPMjmoJ9W2/LKmnJRtNyuOs/OT+NHDY= +github.com/nats-io/nats.go v1.52.0 h1:n3avV4VBsCgsdwh71TppsTwtv+QdPs7ntSKM8qJLGsc= +github.com/nats-io/nats.go v1.52.0/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno= +github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg= +github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs= +github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= +github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= +golang.org/x/crypto v0.53.0 h1:QZ4Muo8THX6CizN2vPPd5fBGHyogrdK9fG4wLPFUsto= +golang.org/x/crypto v0.53.0/go.mod h1:DNLU434OwVakk9PzuwV8w62mAJpRJL3vsgcfp4Qnsio= +golang.org/x/sys v0.21.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= +golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= diff --git a/internal/natsbus/server.go b/internal/natsbus/server.go new file mode 100644 index 0000000..3345004 --- /dev/null +++ b/internal/natsbus/server.go @@ -0,0 +1,186 @@ +// Package natsbus manages the embedded NATS JetStream server that provides +// the coordination bus for Wild Central. Cloud and Works connect to this NATS +// to register services, publish events, and maintain presence. +package natsbus + +import ( + "context" + "fmt" + "log/slog" + "time" + + natsserver "github.com/nats-io/nats-server/v2/server" + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" +) + +const ( + // KV bucket names + BucketServices = "wild-services" // service registrations + BucketPresence = "wild-presence" // node liveness (TTL keys) + + // Stream names + StreamEvents = "wild-events" // events from all sources + + // Presence TTL — keys expire if not renewed within this window + PresenceTTL = 60 * time.Second +) + +// Server wraps an embedded NATS server with JetStream enabled. +type Server struct { + ns *natsserver.Server + nc *nats.Conn + js jetstream.JetStream + port int +} + +// Config for the embedded NATS server. +type Config struct { + Port int // listen port (default 4222) + DataDir string // JetStream storage directory +} + +// Start launches the embedded NATS server with JetStream and creates +// the required KV buckets and streams. +func Start(cfg Config) (*Server, error) { + if cfg.Port == 0 { + cfg.Port = 4222 + } + + opts := &natsserver.Options{ + Port: cfg.Port, + JetStream: true, + StoreDir: cfg.DataDir, + // Quiet logging — NATS logs at its own level + NoLog: true, + NoSigs: true, + } + + ns, err := natsserver.NewServer(opts) + if err != nil { + return nil, fmt.Errorf("creating NATS server: %w", err) + } + + // Start server in background + go ns.Start() + + // Wait for server to be ready + if !ns.ReadyForConnections(10 * time.Second) { + ns.Shutdown() + return nil, fmt.Errorf("NATS server failed to start within 10s") + } + + slog.Info("NATS JetStream server started", "port", cfg.Port, "dataDir", cfg.DataDir) + + // Connect as internal client + nc, err := nats.Connect(ns.ClientURL()) + if err != nil { + ns.Shutdown() + return nil, fmt.Errorf("connecting to embedded NATS: %w", err) + } + + js, err := jetstream.New(nc) + if err != nil { + nc.Close() + ns.Shutdown() + return nil, fmt.Errorf("creating JetStream context: %w", err) + } + + s := &Server{ + ns: ns, + nc: nc, + js: js, + port: cfg.Port, + } + + if err := s.ensureBuckets(); err != nil { + s.Shutdown() + return nil, fmt.Errorf("creating KV buckets: %w", err) + } + + if err := s.ensureStreams(); err != nil { + s.Shutdown() + return nil, fmt.Errorf("creating streams: %w", err) + } + + return s, nil +} + +// ensureBuckets creates the required KV buckets if they don't exist. +func (s *Server) ensureBuckets() error { + ctx := context.Background() + + // Services bucket — no TTL, entries persist until explicitly deleted + _, err := s.js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ + Bucket: BucketServices, + Description: "Service registrations from Wild Cloud, Wild Works, and manual entries", + }) + if err != nil { + return fmt.Errorf("services bucket: %w", err) + } + slog.Info("NATS KV bucket ready", "bucket", BucketServices) + + // Presence bucket — TTL-based, keys expire when nodes stop renewing + _, err = s.js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ + Bucket: BucketPresence, + Description: "Node liveness via TTL-based keys", + TTL: PresenceTTL, + }) + if err != nil { + return fmt.Errorf("presence bucket: %w", err) + } + slog.Info("NATS KV bucket ready", "bucket", BucketPresence, "ttl", PresenceTTL) + + return nil +} + +// ensureStreams creates the required streams if they don't exist. +func (s *Server) ensureStreams() error { + ctx := context.Background() + + _, err := s.js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{ + Name: StreamEvents, + Description: "Events from all Wild Central consumers", + Subjects: []string{"wild.events.>"}, + MaxAge: 24 * time.Hour, // retain events for 24h + Storage: jetstream.FileStorage, + }) + if err != nil { + return fmt.Errorf("events stream: %w", err) + } + slog.Info("NATS stream ready", "stream", StreamEvents) + + return nil +} + +// JetStream returns the JetStream context for internal use. +func (s *Server) JetStream() jetstream.JetStream { + return s.js +} + +// Conn returns the internal NATS connection. +func (s *Server) Conn() *nats.Conn { + return s.nc +} + +// ClientURL returns the URL clients should use to connect. +func (s *Server) ClientURL() string { + return s.ns.ClientURL() +} + +// Port returns the port the server is listening on. +func (s *Server) Port() int { + return s.port +} + +// Shutdown gracefully stops the NATS server. +func (s *Server) Shutdown() { + if s.nc != nil { + s.nc.Close() + } + if s.ns != nil { + s.ns.Shutdown() + s.ns.WaitForShutdown() + } + slog.Info("NATS server stopped") +} diff --git a/internal/natsbus/server_test.go b/internal/natsbus/server_test.go new file mode 100644 index 0000000..7b25d12 --- /dev/null +++ b/internal/natsbus/server_test.go @@ -0,0 +1,151 @@ +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 services bucket exists + kv, err := s.js.KeyValue(ctx, BucketServices) + if err != nil { + t.Fatalf("services 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()) + } +} diff --git a/main.go b/main.go index c60552a..e282f9a 100644 --- a/main.go +++ b/main.go @@ -7,6 +7,7 @@ import ( "net/http" "os" "os/signal" + "path/filepath" "strings" "syscall" "time" @@ -16,6 +17,7 @@ import ( v1 "github.com/wild-cloud/wild-central/internal/api/v1" "github.com/wild-cloud/wild-central/internal/frontend" "github.com/wild-cloud/wild-central/internal/logging" + "github.com/wild-cloud/wild-central/internal/natsbus" ) var startTime time.Time @@ -84,11 +86,28 @@ func main() { slog.Info("configured directories", "dataDir", dataDir) + // Start embedded NATS JetStream server + natsDataDir := filepath.Join(dataDir, "nats") + natsPort := 4222 + if v := os.Getenv("WILD_CENTRAL_NATS_PORT"); v != "" { + fmt.Sscanf(v, "%d", &natsPort) + } + + natsSrv, err := natsbus.Start(natsbus.Config{ + Port: natsPort, + DataDir: natsDataDir, + }) + if err != nil { + slog.Error("failed to start NATS server", "error", err) + os.Exit(1) + } + allowedOrigins := buildAllowedOrigins() api, err := v1.NewAPI(dataDir, Version, allowedOrigins) if err != nil { slog.Error("failed to initialize API", "error", err) + natsSrv.Shutdown() os.Exit(1) } @@ -144,6 +163,7 @@ func main() { sig := <-sigChan slog.Info("shutdown signal received", "signal", sig) cancel() + natsSrv.Shutdown() slog.Info("wild-central stopped") }