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) <noreply@anthropic.com>
This commit is contained in:
14
go.mod
14
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
|
||||
)
|
||||
|
||||
23
go.sum
23
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=
|
||||
|
||||
186
internal/natsbus/server.go
Normal file
186
internal/natsbus/server.go
Normal file
@@ -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")
|
||||
}
|
||||
151
internal/natsbus/server_test.go
Normal file
151
internal/natsbus/server_test.go
Normal file
@@ -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())
|
||||
}
|
||||
}
|
||||
20
main.go
20
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")
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user