Files
wild-central/internal/natsbus/server.go
Paul Payne 6b58ebb692 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>
2026-07-08 23:36:40 +00:00

187 lines
4.6 KiB
Go

// 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")
}