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