Files
cloud-ip-validator/cmd/control-api/main.go
T
2026-08-21 09:58:08 +03:00

164 lines
5.0 KiB
Go

// 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()
for _, v := range cfg.Validators {
if err := database.RegisterValidator(ctx, v.ValidatorID, "", v.OSPortID, ""); err != nil {
return fmt.Errorf("seed validator %s: %w", v.ValidatorID, err)
}
}
if err := database.SeedQueue(ctx, cfg.IPAddresses); err != nil {
return fmt.Errorf("seed ip queue: %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()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
orch.Tick(ctx)
case <-heartbeatTicker.C:
if err := orch.SweepStaleHeartbeats(ctx); err != nil {
log.Error("sweep stale heartbeats", "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)
}