Files
wild-central/internal/natsbus/server.go
Paul Payne 53a640c34f fix: Auto-recover NATS JetStream from stale store on startup
If the NATS server fails to start (stale lock files, corrupt store
from a crash), automatically clean the data directory and retry once.
This prevents the frustrating "NATS server failed to start" error
after unclean shutdowns.

The NATS KV data (registered services) is also stored as YAML files
in the services directory, so cleaning the NATS store loses nothing
— services are re-registered on the next startup or API call.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-07-09 14:12:29 +00:00

201 lines
5.0 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"
"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
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 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: 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")
}