// 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" "os" "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 BucketDomains = "wild-domains" // domain 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 it fails (stale lock files, corrupt // store), clean the data directory and retry once. if !ns.ReadyForConnections(10 * time.Second) { ns.Shutdown() slog.Warn("NATS server failed to start, cleaning store and retrying", "dataDir", cfg.DataDir) os.RemoveAll(cfg.DataDir) os.MkdirAll(cfg.DataDir, 0755) ns, err = natsserver.NewServer(opts) if err != nil { return nil, fmt.Errorf("creating NATS server (retry): %w", err) } go ns.Start() if !ns.ReadyForConnections(10 * time.Second) { ns.Shutdown() return nil, fmt.Errorf("NATS server failed to start after retry") } } 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: BucketDomains, 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", BucketDomains) // 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") }