From 6b58ebb6926848794fa8864e13fda8fa31da376b Mon Sep 17 00:00:00 2001 From: Paul Payne Date: Wed, 8 Jul 2026 23:36:40 +0000 Subject: [PATCH] feat: Add embedded NATS JetStream coordination bus Central now starts an embedded NATS JetStream server as its core coordination bus. This provides: - wild-services KV bucket: service registrations from Cloud/Works - wild-presence KV bucket: node liveness with TTL-based keys - wild-events stream: events from all sources (24h retention) NATS is started in main.go before the API and shut down on SIGTERM. Cloud and Works will connect to Central's NATS (port 4222) to register services and maintain presence. Co-Authored-By: Claude Opus 4.6 (1M context) --- go.mod | 14 +++ go.sum | 23 ++++ internal/natsbus/server.go | 186 ++++++++++++++++++++++++++++++++ internal/natsbus/server_test.go | 151 ++++++++++++++++++++++++++ main.go | 20 ++++ 5 files changed, 394 insertions(+) create mode 100644 internal/natsbus/server.go create mode 100644 internal/natsbus/server_test.go 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") }