// Command control-api is the orchestrator/control-plane binary: it reads // the deployment config, opens the SQLite database, connects to OpenStack // (or a mock, per config), seeds the IP work queue, and serves the HTTP API // that validator-agents and probers poll against. package main import ( "context" "flag" "fmt" "log/slog" "net/http" "os" "os/signal" "syscall" "time" "cloudipvalidator/internal/config" "cloudipvalidator/internal/db" "cloudipvalidator/internal/httpapi" "cloudipvalidator/internal/openstack" "cloudipvalidator/internal/orchestrator" ) func main() { configPath := flag.String("config", "configs/control-api.yaml", "path to control-api config file") flag.Parse() log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) if err := run(*configPath, log); err != nil { log.Error("fatal", "err", err) os.Exit(1) } } func run(configPath string, log *slog.Logger) error { cfg, err := config.LoadControlAPI(configPath) if err != nil { return fmt.Errorf("load config: %w", err) } ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() database, err := db.Open(ctx, cfg.Database.Path) if err != nil { return fmt.Errorf("open database: %w", err) } defer database.Close() if err := database.BootstrapFromConfig(ctx, cfg); err != nil { return fmt.Errorf("bootstrap database: %w", err) } osClient, err := newOpenStackClient(ctx, cfg) if err != nil { return fmt.Errorf("init openstack client: %w", err) } orch := orchestrator.New(database, osClient, cfg, log) srv := httpapi.New(database, orch, log) httpServer := &http.Server{Addr: cfg.Server.ListenAddr, Handler: srv.Handler()} go runOrchestratorLoop(ctx, orch, cfg, log) errCh := make(chan error, 1) go func() { log.Info("listening", "addr", cfg.Server.ListenAddr) if err := httpServer.ListenAndServe(); err != nil && err != http.ErrServerClosed { errCh <- err } }() select { case <-ctx.Done(): log.Info("shutting down") shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() return httpServer.Shutdown(shutdownCtx) case err := <-errCh: return err } } func runOrchestratorLoop(ctx context.Context, orch *orchestrator.Orchestrator, cfg *config.ControlAPI, log *slog.Logger) { interval := time.Duration(cfg.Orchestrator.PollIntervalSeconds) * time.Second ticker := time.NewTicker(interval) defer ticker.Stop() heartbeatTicker := time.NewTicker(interval * 2) defer heartbeatTicker.Stop() // Periodic floating-IP scanning is optional: a zero interval leaves // scanTickerC nil, and a select on a nil channel simply never fires, so // the loop falls back to manual-only scanning (POST // /api/v1/admin/ips/scan) without a special-cased branch below. var scanTickerC <-chan time.Time if cfg.Orchestrator.FIPScanIntervalSeconds > 0 { scanTicker := time.NewTicker(time.Duration(cfg.Orchestrator.FIPScanIntervalSeconds) * time.Second) defer scanTicker.Stop() scanTickerC = scanTicker.C } for { select { case <-ctx.Done(): return case <-ticker.C: orch.Tick(ctx) orch.AutoCycleStep(ctx) case <-heartbeatTicker.C: if err := orch.SweepStaleHeartbeats(ctx); err != nil { log.Error("sweep stale heartbeats", "err", err) } if err := orch.SweepStaleSiteHeartbeats(ctx); err != nil { log.Error("sweep stale site heartbeats", "err", err) } case <-scanTickerC: // The auto-cycle owns the queue while enabled: a periodic scan // would add addresses in the middle of a cycle. ac, err := orch.GetAutoCycle(ctx) if err != nil { log.Error("read auto-cycle before periodic scan, scanning anyway", "err", err) } else if ac.Enabled { continue } if _, _, err := orch.ScanFloatingIPs(ctx); err != nil { log.Error("scan floating ips", "err", err) } } } } func newOpenStackClient(ctx context.Context, cfg *config.ControlAPI) (openstack.FloatingIPClient, error) { if cfg.OpenStack.Mode == "real" { return newRealOpenStackClient(ctx, cfg) } // Mock mode: pre-seed one synthetic floating-ip resource per // configured address, standing in for the pre-allocated Neutron // floating IPs a real deployment's service project would already have. mock := openstack.NewMockClient() for i, addr := range cfg.IPAddresses { mock.Seed(fmt.Sprintf("mock-fip-%d", i), addr, "mock-project") } return mock, nil } func newRealOpenStackClient(ctx context.Context, cfg *config.ControlAPI) (openstack.FloatingIPClient, error) { clientCfg := openstack.ClientConfig{ AuthURL: os.Getenv(cfg.OpenStack.AuthURLEnv), ProjectID: os.Getenv(cfg.OpenStack.ProjectIDEnv), Region: os.Getenv(cfg.OpenStack.RegionEnv), Interface: os.Getenv(cfg.OpenStack.InterfaceEnv), } switch cfg.OpenStack.AuthMethod { case "password": clientCfg.Method = openstack.AuthMethodPassword clientCfg.Username = os.Getenv(cfg.OpenStack.UsernameEnv) clientCfg.Password = os.Getenv(cfg.OpenStack.PasswordEnv) clientCfg.UserDomainName = os.Getenv(cfg.OpenStack.UserDomainNameEnv) if clientCfg.Username == "" || clientCfg.Password == "" { return nil, fmt.Errorf("openstack.auth_method=password requires %s and %s to be set in the environment", cfg.OpenStack.UsernameEnv, cfg.OpenStack.PasswordEnv) } case "token", "": clientCfg.Method = openstack.AuthMethodToken clientCfg.Token = os.Getenv(cfg.OpenStack.TokenEnv) if clientCfg.Token == "" { return nil, fmt.Errorf("openstack.auth_method=token requires %s to be set in the environment", cfg.OpenStack.TokenEnv) } default: return nil, fmt.Errorf("openstack.auth_method: unknown value %q (expected \"token\" or \"password\")", cfg.OpenStack.AuthMethod) } if clientCfg.AuthURL == "" || clientCfg.ProjectID == "" { return nil, fmt.Errorf("openstack.mode=real requires %s and %s to be set in the environment", cfg.OpenStack.AuthURLEnv, cfg.OpenStack.ProjectIDEnv) } return openstack.NewClient(ctx, clientCfg) }