Files
wild-central/internal/natsbus/server.go
Paul Payne 428d47f876 Refactor architecture: extract reconciler, add interfaces, reduce complexity, improve test coverage
Architecture:
- Extract reconcileNetworking into internal/reconcile package with 7 consumer-side interfaces
- Add locked modifyState helper to fix state.yaml read-modify-write race condition
- Extract CrowdSec and VPN handler groups with interfaces documenting dependency surface
- Replace raw map[string]any YAML manipulation with typed AddDHCPStaticLease/RemoveDHCPStaticLease
- Extract Cloudflare API functions into cfClient struct, eliminating repeated auth boilerplate

Complexity reduction:
- haproxy.GenerateWithOpts: 50 → 6 (extracted 8 focused helpers)
- reconcile.Reconcile: 42 → 11 (extracted buildRoutes, buildDNSEntries, writeHAProxyConfig)
- AutheliaUpdateConfig: 25 → 10 (extracted enableAuthelia/disableAuthelia)

Test coverage improvements:
- reconcile: 4.5% → 56.8% (stub-based tests for route building, DNS entries, orchestration)
- dnsfilter: 10.7% → 45.6% (Manager.Compile, AddList, ToggleList, custom entries)
- config: 67.3% → 89.8% (DHCP static lease mutation tests)
2026-07-14 04:21:30 +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
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")
}