A failed self-check returned the address to the queue and freed the
validator in the database, but left the floating IP attached to the
validator's port. Every later association on that port then failed with
409 ("fixed IP already has a floating IP"), so one failed self-check
poisoned a validator for good; on 2026-10-01 all 20 validators were
poisoned within 23 minutes after ifconfig.me timeouts.
SelfCheckResult now disassociates the floating IP before requeueing, and
ignores a late failed report for an address the validator no longer
holds (it could belong to another validator by then).
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
854 lines
32 KiB
Go
854 lines
32 KiB
Go
// Package orchestrator implements the Control API's core scheduling loop:
|
|
// claiming queued IPs onto idle validators, driving each IP through
|
|
// FIP-association -> self-check -> checking -> aggregation -> release, and
|
|
// reclaiming work from crashed/stuck validators via a lease sweep. It has
|
|
// no HTTP dependency — internal/httpapi calls into this package, and it can
|
|
// be exercised directly in tests against an in-memory OpenStack mock and a
|
|
// temp-file SQLite database.
|
|
package orchestrator
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"sync"
|
|
"time"
|
|
|
|
"cloudipvalidator/internal/config"
|
|
"cloudipvalidator/internal/db"
|
|
"cloudipvalidator/internal/openstack"
|
|
)
|
|
|
|
// CheckConfig is the check-type/target configuration handed to a
|
|
// validator-agent once its IP has passed self-check. It mirrors
|
|
// db.ResolvedCheckType, kept as a distinct type so httpapi's DTO layer
|
|
// doesn't need to import internal/db just for this shape.
|
|
type CheckConfig struct {
|
|
Type string `json:"type"`
|
|
Targets []string `json:"targets"`
|
|
}
|
|
|
|
type Orchestrator struct {
|
|
DB *db.DB
|
|
OS openstack.FloatingIPClient
|
|
Cfg config.OrchestratorConfig
|
|
Agg config.AggregationConfig
|
|
Log *slog.Logger
|
|
|
|
// autoCycleMu serializes AutoCycleStep with StartAutoCycle/StopAutoCycle
|
|
// so an API call can never interleave with a half-finished step.
|
|
autoCycleMu sync.Mutex
|
|
|
|
// ScanPageSize is the Neutron page size used by the floating-IP scan
|
|
// (config openstack.list_page_size); 0 means openstack.DefaultListPageSize.
|
|
ScanPageSize int
|
|
|
|
// scan is the background floating-IP scan job (see scanjob.go); its zero
|
|
// value is ready to use.
|
|
scan scanJob
|
|
|
|
// Async makes Tick return without waiting for the slow per-address work
|
|
// (OpenStack calls): each address is handled in its own goroutine, so one
|
|
// slow cloud call never delays another validator. With Async=false (the
|
|
// default, used by tests) Tick runs the same work in parallel but waits
|
|
// for all of it before returning.
|
|
Async bool
|
|
|
|
// inflight holds the keys of work already running in a goroutine, so the
|
|
// next Tick does not start the same work twice.
|
|
inflight sync.Map
|
|
wg sync.WaitGroup
|
|
}
|
|
|
|
// Wait blocks until all background work started by Tick has finished.
|
|
func (o *Orchestrator) Wait() { o.wg.Wait() }
|
|
|
|
// spawn runs fn in its own goroutine unless work with the same key is
|
|
// already running. It returns immediately; callers that need the result
|
|
// wait with o.wg (see runAll).
|
|
func (o *Orchestrator) spawn(key string, fn func()) {
|
|
if _, busy := o.inflight.LoadOrStore(key, struct{}{}); busy {
|
|
return
|
|
}
|
|
o.wg.Add(1)
|
|
go func() {
|
|
defer o.wg.Done()
|
|
defer o.inflight.Delete(key)
|
|
fn()
|
|
}()
|
|
}
|
|
|
|
// New constructs an Orchestrator. Egress check types/targets, prober sites,
|
|
// and prober inbound check config (ports/icmp) are no longer taken from
|
|
// cfg — they're read from the database on every use (see
|
|
// AssignmentForValidator, expectedCheckCount, isReadyToAggregate,
|
|
// httpapi.handleProberAssignments) so admin API changes to them take effect
|
|
// without a restart. cfg.Validators/.Sites/.CheckTypes/.Targets/.Inbound/
|
|
// .IPAddresses are only consulted once, at process startup, by
|
|
// db.BootstrapFromConfig.
|
|
func New(d *db.DB, osClient openstack.FloatingIPClient, cfg *config.ControlAPI, log *slog.Logger) *Orchestrator {
|
|
return &Orchestrator{
|
|
DB: d,
|
|
OS: osClient,
|
|
Cfg: cfg.Orchestrator,
|
|
Agg: cfg.Aggregation,
|
|
Log: log,
|
|
|
|
ScanPageSize: cfg.OpenStack.ListPageSize,
|
|
}
|
|
}
|
|
|
|
func (o *Orchestrator) leaseTTL() time.Duration {
|
|
return time.Duration(o.Cfg.LeaseTTLSeconds) * time.Second
|
|
}
|
|
|
|
// Tick runs one pass of the scheduling loop: claim+associate for idle
|
|
// validators, sweep the checking window for ready-to-aggregate IPs, and
|
|
// reclaim expired leases. Intended to be called on a fixed interval
|
|
// (Cfg.PollIntervalSeconds) by the caller (cmd/control-api/main.go).
|
|
//
|
|
// The slow part of every step (calls to OpenStack) runs per address in its
|
|
// own goroutine, so validators never wait for each other. In Async mode
|
|
// Tick does not wait for those goroutines; otherwise it waits for them.
|
|
func (o *Orchestrator) Tick(ctx context.Context) {
|
|
if err := o.assignIdleValidators(ctx); err != nil {
|
|
o.Log.Error("assign idle validators", "err", err)
|
|
}
|
|
if err := o.sweepCheckingWindow(ctx); err != nil {
|
|
o.Log.Error("sweep checking window", "err", err)
|
|
}
|
|
if err := o.sweepExpiredLeases(ctx); err != nil {
|
|
o.Log.Error("sweep expired leases", "err", err)
|
|
}
|
|
if !o.Async {
|
|
o.wg.Wait()
|
|
}
|
|
}
|
|
|
|
// assignIdleValidators claims the next queued IP for every currently idle
|
|
// validator and kicks off FIP association for each newly claimed IP.
|
|
func (o *Orchestrator) assignIdleValidators(ctx context.Context) error {
|
|
idle, err := o.DB.ListIdleValidators(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("list idle validators: %w", err)
|
|
}
|
|
for _, v := range idle {
|
|
item, err := o.DB.ClaimNextQueued(ctx, v.ValidatorID, o.leaseTTL())
|
|
if err != nil {
|
|
o.Log.Error("claim next queued", "validator", v.ValidatorID, "err", err)
|
|
continue
|
|
}
|
|
if item == nil {
|
|
continue // no work available for this validator right now
|
|
}
|
|
o.Log.Info("claimed ip", "validator", v.ValidatorID, "ip", item.IPAddress, "ip_id", item.ID)
|
|
v, item := v, item
|
|
o.spawn("assign:"+v.ValidatorID, func() {
|
|
if err := o.associateFIP(ctx, v.ValidatorID, v.OSPortID, item); err != nil {
|
|
o.Log.Error("associate fip", "validator", v.ValidatorID, "ip", item.IPAddress, "err", err)
|
|
}
|
|
})
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (o *Orchestrator) associateFIP(ctx context.Context, validatorID, osPortID string, item *db.IPQueueItem) error {
|
|
fip, err := o.OS.GetFloatingIPByAddress(ctx, item.IPAddress)
|
|
if err != nil {
|
|
o.requeueOrFail(ctx, item.ID, validatorID, fmt.Sprintf("lookup floating ip: %v", err))
|
|
return err
|
|
}
|
|
// The cloud is live: an address queued as "free" (from bootstrap config
|
|
// or an admin POST) may have drifted onto another port by the time we
|
|
// actually get here, or an operator may have queued an already-occupied
|
|
// address by mistake. This is the one authoritative moment to catch it —
|
|
// checked here rather than at enqueue time because enqueue-time state
|
|
// could itself be stale by the time the claim happens. fip.PortID !=
|
|
// osPortID guards against a false positive when the FIP is already
|
|
// associated to this same validator's own port (e.g. control-api
|
|
// restarted between associating and recording it) — that's a resume, not
|
|
// a conflict.
|
|
if fip.PortID != "" && fip.PortID != osPortID {
|
|
if err := o.DB.MarkFIPOccupied(ctx, item.ID, validatorID); err != nil {
|
|
o.Log.Error("mark fip occupied", "ip_id", item.ID, "err", err)
|
|
return err
|
|
}
|
|
o.event(ctx, "control-api", "", &item.ID, "fip_occupied",
|
|
fmt.Sprintf(`{"fip_id":%q,"port_id":%q}`, fip.ID, fip.PortID))
|
|
o.Log.Info("fip already occupied by another port, skipping check cycle",
|
|
"ip", item.IPAddress, "fip_port_id", fip.PortID)
|
|
return nil
|
|
}
|
|
if err := o.OS.AssociateFloatingIP(ctx, fip.ID, osPortID); err != nil {
|
|
o.requeueOrFail(ctx, item.ID, validatorID, fmt.Sprintf("associate floating ip: %v", err))
|
|
return err
|
|
}
|
|
if err := o.DB.SetFIPAssociated(ctx, item.ID, fip.ID, o.leaseTTL()); err != nil {
|
|
if errors.Is(err, db.ErrInvalidState) {
|
|
// The address was deleted or cancelled while the cloud call was
|
|
// running: nobody owns this attachment any more, so detach it,
|
|
// otherwise it stays on the validator's port forever.
|
|
o.Log.Info("address removed during association, detaching floating ip", "ip", item.IPAddress, "fip_id", fip.ID)
|
|
if derr := o.OS.DisassociateFloatingIP(ctx, fip.ID); derr != nil {
|
|
o.Log.Error("disassociate fip of removed address", "ip", item.IPAddress, "fip_id", fip.ID, "err", derr)
|
|
}
|
|
return nil
|
|
}
|
|
return fmt.Errorf("set fip associated: %w", err)
|
|
}
|
|
o.event(ctx, "control-api", "", &item.ID, "fip_associated", fmt.Sprintf(`{"fip_id":%q,"validator_id":%q}`, fip.ID, validatorID))
|
|
return nil
|
|
}
|
|
|
|
func (o *Orchestrator) requeueOrFail(ctx context.Context, ipID int64, validatorID, reason string) {
|
|
if err := o.DB.RequeueOrFail(ctx, ipID, validatorID, o.Cfg.MaxRetries); err != nil {
|
|
o.Log.Error("requeue or fail", "ip_id", ipID, "err", err)
|
|
return
|
|
}
|
|
o.event(ctx, "control-api", "", &ipID, "retry_or_fail", fmt.Sprintf(`{"reason":%q}`, reason))
|
|
}
|
|
|
|
// SelfCheckResult is called by the httpapi layer when a validator-agent
|
|
// reports its post-association self-check outcome.
|
|
func (o *Orchestrator) SelfCheckResult(ctx context.Context, validatorID string, ipID int64, success bool, detail string) error {
|
|
o.event(ctx, "validator-agent", validatorID, &ipID, "self_check_result",
|
|
fmt.Sprintf(`{"success":%t,"detail":%q}`, success, detail))
|
|
|
|
if !success {
|
|
item, err := o.DB.GetIP(ctx, ipID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// A late report for an address this validator no longer owns (lease
|
|
// already reclaimed, address cancelled or deleted) must not touch
|
|
// the floating IP: it may be attached for another validator now.
|
|
if item.State != db.IPAwaitingSelfCheck || item.OwnerValidatorID == nil || *item.OwnerValidatorID != validatorID {
|
|
o.Log.Warn("ignoring failed self-check for an address the validator does not hold",
|
|
"validator", validatorID, "ip_id", ipID, "state", item.State)
|
|
return nil
|
|
}
|
|
// Detach the floating IP before the address goes back to the queue.
|
|
// requeueOrFail frees the validator in the database but knows nothing
|
|
// about the cloud: a floating IP left on the validator's port makes
|
|
// every later association on that port fail with 409 ("fixed IP
|
|
// already has a floating IP"). Best-effort, like the other release
|
|
// paths — the database state must be freed even if Neutron hiccups.
|
|
if item.FIPID != "" {
|
|
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil {
|
|
o.Log.Error("disassociate fip after failed self-check", "ip_id", ipID, "fip_id", item.FIPID, "err", err)
|
|
}
|
|
}
|
|
if item.RetryCount+1 > o.Cfg.MaxSelfCheckRetries {
|
|
o.requeueOrFail(ctx, ipID, validatorID, "self-check failed: "+detail)
|
|
return nil
|
|
}
|
|
// Retry association without fully requeuing: re-drive the same
|
|
// claim by cycling back through requeue/claim keeps the logic in
|
|
// one place at the cost of the IP briefly returning to `queued`.
|
|
o.requeueOrFail(ctx, ipID, validatorID, "self-check failed, retrying: "+detail)
|
|
return nil
|
|
}
|
|
|
|
return o.DB.SetChecking(ctx, ipID, o.leaseTTL())
|
|
}
|
|
|
|
// AssignmentForValidator returns the check config for a validator's current
|
|
// IP if it's ready to be worked on (awaiting_self_check or checking),
|
|
// or nil if the validator has nothing to do right now. The check config is
|
|
// read fresh from the database on every call, so admin changes to
|
|
// check_types/targets apply to the very next assignment.
|
|
func (o *Orchestrator) AssignmentForValidator(ctx context.Context, validatorID string) (*db.IPQueueItem, []CheckConfig, error) {
|
|
v, err := o.DB.GetValidator(ctx, validatorID)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
if v.CurrentIPID == nil {
|
|
return nil, nil, nil
|
|
}
|
|
item, err := o.DB.GetIP(ctx, *v.CurrentIPID)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
if item.State != db.IPAwaitingSelfCheck && item.State != db.IPChecking {
|
|
return nil, nil, nil
|
|
}
|
|
if item.State == db.IPAwaitingSelfCheck {
|
|
settled, err := o.isFIPSettled(ctx, item)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
if !settled {
|
|
return nil, nil, nil
|
|
}
|
|
}
|
|
resolved, err := o.DB.ListResolvedCheckTypes(ctx)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
checks := make([]CheckConfig, len(resolved))
|
|
for i, r := range resolved {
|
|
checks[i] = CheckConfig{Type: r.Type, Targets: r.Targets}
|
|
}
|
|
return item, checks, nil
|
|
}
|
|
|
|
// isFIPSettled reports whether an address currently awaiting_self_check
|
|
// has cleared the configured fip_settle_seconds pause since its floating
|
|
// IP was associated — withholding the assignment until then is how the
|
|
// pause is enforced, with zero changes needed to the agent's poll loop or
|
|
// the assignment endpoint's wire contract (it just keeps seeing 204s).
|
|
// Settings are read fresh on every call, same as check_types/sites
|
|
// elsewhere in this file, so an admin change applies immediately even to
|
|
// an address already mid-wait. A nil FIPAssociatedAt (an in-flight row
|
|
// from before this feature's migration) is always treated as settled —
|
|
// upgrading control-api must never newly strand an address that was
|
|
// already awaiting self-check.
|
|
func (o *Orchestrator) isFIPSettled(ctx context.Context, item *db.IPQueueItem) (bool, error) {
|
|
settings, err := o.DB.GetSettings(ctx)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
if settings.FIPSettleSeconds <= 0 || item.FIPAssociatedAt == nil {
|
|
return true, nil
|
|
}
|
|
deadline := item.FIPAssociatedAt.Add(time.Duration(settings.FIPSettleSeconds) * time.Second)
|
|
return !db.Now().Before(deadline), nil
|
|
}
|
|
|
|
// SetFIPSettleSeconds validates and persists a new fip_settle_seconds
|
|
// value. Lives here rather than internal/db because the cross-field rule
|
|
// below needs o.Cfg, which the db package has no access to: the pause plus
|
|
// self-check's own timeout must leave room inside the claim lease, or the
|
|
// lease sweep would reclaim the address before self-check ever gets a
|
|
// chance to run, producing a perpetual requeue loop.
|
|
func (o *Orchestrator) SetFIPSettleSeconds(ctx context.Context, seconds int) error {
|
|
if seconds < 0 {
|
|
return fmt.Errorf("fip_settle_seconds must be >= 0: %w", db.ErrValidation)
|
|
}
|
|
if seconds+o.Cfg.SelfCheckTimeoutSeconds >= o.Cfg.LeaseTTLSeconds {
|
|
return fmt.Errorf(
|
|
"fip_settle_seconds (%d) + self_check_timeout_seconds (%d) must be < lease_ttl_seconds (%d): %w",
|
|
seconds, o.Cfg.SelfCheckTimeoutSeconds, o.Cfg.LeaseTTLSeconds, db.ErrValidation)
|
|
}
|
|
return o.DB.SetFIPSettleSeconds(ctx, seconds)
|
|
}
|
|
|
|
// SiteIndexForID resolves a configured site_id to its 1/2/3 index, or
|
|
// (0, nil) if unconfigured.
|
|
func (o *Orchestrator) SiteIndexForID(ctx context.Context, siteID string) (int, error) {
|
|
return o.DB.GetSiteIndex(ctx, siteID)
|
|
}
|
|
|
|
// RecordCheck upserts a single check result and, if it represents a
|
|
// completion signal (egress or a given site's full port+icmp sweep),
|
|
// updates the corresponding *_complete flag.
|
|
func (o *Orchestrator) RecordCheck(ctx context.Context, c db.Check) error {
|
|
return o.DB.UpsertCheck(ctx, c)
|
|
}
|
|
|
|
func (o *Orchestrator) MarkEgressComplete(ctx context.Context, ipID int64) error {
|
|
return o.DB.SetEgressComplete(ctx, ipID)
|
|
}
|
|
|
|
func (o *Orchestrator) MarkSiteComplete(ctx context.Context, ipID int64, siteIndex int) error {
|
|
return o.DB.SetSiteComplete(ctx, ipID, siteIndex)
|
|
}
|
|
|
|
// releaseValidatorPorts detaches, in the cloud, every floating IP that sits on
|
|
// the given validators' ports and belongs to an address this system has
|
|
// queued. The cloud is asked directly because our database does not know
|
|
// about an attachment that is still being created (state assigning_fip has no
|
|
// fip_id yet) or that was lost earlier. Floating IPs this system never
|
|
// queued are left alone. Ports are processed in parallel.
|
|
func (o *Orchestrator) releaseValidatorPorts(ctx context.Context, validators []db.Validator) {
|
|
var wg sync.WaitGroup
|
|
for _, v := range validators {
|
|
if v.OSPortID == "" {
|
|
continue
|
|
}
|
|
wg.Add(1)
|
|
go func(v db.Validator) {
|
|
defer wg.Done()
|
|
fips, err := o.OS.ListFloatingIPsByPort(ctx, v.OSPortID)
|
|
if err != nil {
|
|
o.Log.Error("list floating ips of validator port", "validator", v.ValidatorID, "port", v.OSPortID, "err", err)
|
|
return
|
|
}
|
|
addrs := make([]string, len(fips))
|
|
for i, f := range fips {
|
|
addrs[i] = f.Address
|
|
}
|
|
known, err := o.DB.KnownAddresses(ctx, addrs)
|
|
if err != nil {
|
|
o.Log.Error("check known addresses", "validator", v.ValidatorID, "err", err)
|
|
return
|
|
}
|
|
for _, f := range fips {
|
|
if !known[f.Address] {
|
|
o.Log.Warn("floating ip on validator port was not queued by this system, left attached",
|
|
"validator", v.ValidatorID, "ip", f.Address)
|
|
continue
|
|
}
|
|
if err := o.OS.DisassociateFloatingIP(ctx, f.ID); err != nil {
|
|
o.Log.Error("detach floating ip from validator port", "validator", v.ValidatorID, "ip", f.Address, "err", err)
|
|
continue
|
|
}
|
|
o.Log.Info("detached floating ip from validator port", "validator", v.ValidatorID, "ip", f.Address)
|
|
}
|
|
}(v)
|
|
}
|
|
wg.Wait()
|
|
}
|
|
|
|
// validatorsByID returns the validators with the given ids (unknown ids are skipped).
|
|
func (o *Orchestrator) validatorsByID(ctx context.Context, ids ...string) []db.Validator {
|
|
var out []db.Validator
|
|
for _, id := range ids {
|
|
if v, err := o.DB.GetValidator(ctx, id); err == nil {
|
|
out = append(out, *v)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// busyValidatorsFor returns the validators currently working on one of the
|
|
// given addresses.
|
|
func (o *Orchestrator) busyValidatorsFor(ctx context.Context, addresses []string) []db.Validator {
|
|
want := make(map[string]bool, len(addresses))
|
|
for _, a := range addresses {
|
|
want[a] = true
|
|
}
|
|
all, err := o.DB.ListValidators(ctx)
|
|
if err != nil {
|
|
o.Log.Error("list validators", "err", err)
|
|
return nil
|
|
}
|
|
var out []db.Validator
|
|
for _, v := range all {
|
|
if v.CurrentIPID == nil {
|
|
continue
|
|
}
|
|
if ip, err := o.DB.GetIP(ctx, *v.CurrentIPID); err == nil && want[ip.IPAddress] {
|
|
out = append(out, v)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// ForceCancel stops an in-progress (or still-queued) check for the given
|
|
// address on admin request, even though it was never going to finish on
|
|
// its own within the checking window. If a floating IP is currently
|
|
// associated, it's disassociated best-effort (same fallthrough-on-error
|
|
// behavior as aggregateAndRelease/sweepExpiredLeases: the DB/validator
|
|
// state must still be freed even if Neutron hiccups). Returns
|
|
// db.ErrNotFound if the address is unknown, or db.ErrInvalidState if it has
|
|
// already reached done/failed.
|
|
func (o *Orchestrator) ForceCancel(ctx context.Context, ipAddress string) error {
|
|
item, err := o.DB.GetIPByAddress(ctx, ipAddress)
|
|
if err != nil {
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return fmt.Errorf("ip %q: %w", ipAddress, db.ErrNotFound)
|
|
}
|
|
return err
|
|
}
|
|
|
|
if item.FIPID != "" {
|
|
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil {
|
|
o.Log.Error("disassociate fip on force cancel", "ip_id", item.ID, "fip_id", item.FIPID, "err", err)
|
|
}
|
|
}
|
|
|
|
if err := o.DB.CancelIP(ctx, item.ID); err != nil {
|
|
return err
|
|
}
|
|
|
|
if item.OwnerValidatorID != nil {
|
|
if err := o.DB.FreeValidator(ctx, *item.OwnerValidatorID); err != nil {
|
|
return fmt.Errorf("free validator: %w", err)
|
|
}
|
|
o.releaseValidatorPorts(ctx, o.validatorsByID(ctx, *item.OwnerValidatorID))
|
|
}
|
|
|
|
o.event(ctx, "control-api", "", &item.ID, "force_cancel", "")
|
|
return nil
|
|
}
|
|
|
|
// DeleteIP disassociates the floating IP if attached, then permanently
|
|
// removes the address and its full history — differs from ForceCancel,
|
|
// which keeps a cancelled record instead of deleting it. Works from any
|
|
// state, including actively checking: it does the same resource-freeing
|
|
// (FIP disassociation, validator release) ForceCancel does, but goes
|
|
// straight to physical deletion rather than parking in `cancelled`.
|
|
// Returns db.ErrNotFound if the address is unknown.
|
|
func (o *Orchestrator) DeleteIP(ctx context.Context, ipAddress string) error {
|
|
item, err := o.DB.GetIPByAddress(ctx, ipAddress)
|
|
if err != nil {
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return fmt.Errorf("ip %q: %w", ipAddress, db.ErrNotFound)
|
|
}
|
|
return err
|
|
}
|
|
|
|
if item.FIPID != "" {
|
|
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil {
|
|
o.Log.Error("disassociate fip on delete", "ip_id", item.ID, "fip_id", item.FIPID, "err", err)
|
|
}
|
|
}
|
|
|
|
var owners []db.Validator
|
|
if item.OwnerValidatorID != nil {
|
|
owners = o.validatorsByID(ctx, *item.OwnerValidatorID)
|
|
}
|
|
if err := o.DB.DeleteIP(ctx, item.ID); err != nil {
|
|
return err
|
|
}
|
|
o.releaseValidatorPorts(ctx, owners)
|
|
|
|
o.event(ctx, "control-api", "", nil, "ip_deleted", fmt.Sprintf(`{"ip_address":%q}`, ipAddress))
|
|
return nil
|
|
}
|
|
|
|
// DeleteIPs disassociates the floating IP (best-effort) for every address
|
|
// in the list that has one attached, then deletes the whole list in one
|
|
// DB.DeleteIPs call. Addresses not currently in the queue are simply
|
|
// omitted from the disassociation pass and reported back in NotFound by
|
|
// DB.DeleteIPs — not an error. The rows holding a floating IP are found with
|
|
// a few IN (...) queries rather than one lookup per address.
|
|
func (o *Orchestrator) DeleteIPs(ctx context.Context, addresses []string) (db.DeleteIPsResult, error) {
|
|
refs, err := o.DB.ListFIPRefsByAddresses(ctx, addresses)
|
|
if err != nil {
|
|
o.Log.Error("list attached fips before delete", "err", err)
|
|
}
|
|
for _, ref := range refs {
|
|
if err := o.OS.DisassociateFloatingIP(ctx, ref.FIPID); err != nil {
|
|
o.Log.Error("disassociate fip on delete", "ip_id", ref.IPID, "fip_id", ref.FIPID, "err", err)
|
|
}
|
|
}
|
|
|
|
owners := o.busyValidatorsFor(ctx, addresses)
|
|
result, err := o.DB.DeleteIPs(ctx, addresses)
|
|
if err != nil {
|
|
return result, err
|
|
}
|
|
o.releaseValidatorPorts(ctx, owners)
|
|
o.event(ctx, "control-api", "", nil, "ips_deleted", deletedAddressesPayload(result.Deleted))
|
|
return result, nil
|
|
}
|
|
|
|
// ClearQueue deletes every address currently in the queue, regardless of
|
|
// state — the "delete everything" operation. It is set-based (see
|
|
// db.ClearAllIPs): O(1) statements however many rows there are. Floating IPs
|
|
// attached to rows are disassociated first (best-effort), and only for rows
|
|
// that actually hold one.
|
|
func (o *Orchestrator) ClearQueue(ctx context.Context) (db.DeleteIPsResult, error) {
|
|
refs, err := o.DB.ListFIPRefs(ctx)
|
|
if err != nil {
|
|
return db.DeleteIPsResult{}, fmt.Errorf("list attached fips: %w", err)
|
|
}
|
|
for _, ref := range refs {
|
|
if err := o.OS.DisassociateFloatingIP(ctx, ref.FIPID); err != nil {
|
|
o.Log.Error("disassociate fip on clear queue", "ip_id", ref.IPID, "fip_id", ref.FIPID, "err", err)
|
|
}
|
|
}
|
|
|
|
deleted, err := o.DB.ClearAllIPs(ctx)
|
|
if err != nil {
|
|
return db.DeleteIPsResult{}, err
|
|
}
|
|
// After the rows are gone, free every validator port in the cloud: this
|
|
// catches attachments that were still being created (no fip_id yet) and
|
|
// ones left over from earlier runs. An association that finishes after
|
|
// this point detaches itself (see associateFIP).
|
|
if validators, err := o.DB.ListValidators(ctx); err != nil {
|
|
o.Log.Error("list validators to release ports", "err", err)
|
|
} else {
|
|
o.releaseValidatorPorts(ctx, validators)
|
|
}
|
|
o.event(ctx, "control-api", "", nil, "queue_cleared", deletedAddressesPayload(deleted))
|
|
return db.DeleteIPsResult{Deleted: deleted}, nil
|
|
}
|
|
|
|
// maxEventAddresses caps how many addresses a batch-delete event payload
|
|
// lists: clearing 6440 addresses must not write a 100 KB event row.
|
|
const maxEventAddresses = 50
|
|
|
|
// deletedAddressesPayload builds the event payload for the batch delete
|
|
// operations — a proper JSON object via encoding/json: the total count plus
|
|
// at most the first maxEventAddresses addresses (and "truncated":true when
|
|
// the list was cut).
|
|
func deletedAddressesPayload(addresses []string) string {
|
|
p := struct {
|
|
Count int `json:"count"`
|
|
Addresses []string `json:"addresses"`
|
|
Truncated bool `json:"truncated,omitempty"`
|
|
}{Count: len(addresses), Addresses: addresses}
|
|
if p.Addresses == nil {
|
|
p.Addresses = []string{}
|
|
}
|
|
if len(p.Addresses) > maxEventAddresses {
|
|
p.Addresses = p.Addresses[:maxEventAddresses]
|
|
p.Truncated = true
|
|
}
|
|
b, err := json.Marshal(p)
|
|
if err != nil {
|
|
return "{}"
|
|
}
|
|
return string(b)
|
|
}
|
|
|
|
// sweepCheckingWindow moves IPs that have either finished reporting from
|
|
// every source, or hit the checking-window deadline, into aggregation.
|
|
func (o *Orchestrator) sweepCheckingWindow(ctx context.Context) error {
|
|
deadline := db.Now().Add(-time.Duration(o.Cfg.CheckingWindowSeconds) * time.Second)
|
|
checking, err := o.DB.ListChecking(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("list checking: %w", err)
|
|
}
|
|
sites, err := o.DB.ListSites(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("list sites: %w", err)
|
|
}
|
|
for _, item := range checking {
|
|
ready, err := o.isReadyToAggregate(ctx, item, deadline, sites)
|
|
if err != nil {
|
|
o.Log.Error("check ready to aggregate", "ip_id", item.ID, "err", err)
|
|
continue
|
|
}
|
|
if !ready {
|
|
continue
|
|
}
|
|
item := item
|
|
o.spawn(fmt.Sprintf("aggregate:%d", item.ID), func() {
|
|
if err := o.aggregateAndRelease(ctx, item); err != nil {
|
|
o.Log.Error("aggregate and release", "ip_id", item.ID, "err", err)
|
|
}
|
|
})
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// isReadyToAggregate reports whether an in-progress IP has either finished
|
|
// reporting from every source it's actually expecting, or hit the
|
|
// checking-window deadline. Which inbound sources it's expecting is driven
|
|
// entirely by the currently configured sites — inbound checks are optional:
|
|
// an empty (or partial) sites configuration means this IP is ready as soon
|
|
// as egress completes (or after the corresponding subset of sites has
|
|
// reported in ip_site_checks), with no need to wait on a prober that will
|
|
// never exist, and no cap on how many sites can be configured.
|
|
func (o *Orchestrator) isReadyToAggregate(ctx context.Context, item db.IPQueueItem, deadline time.Time, sites []db.Site) (bool, error) {
|
|
if item.AssignedAt != nil && item.AssignedAt.Before(deadline) {
|
|
return true, nil
|
|
}
|
|
if !item.EgressComplete {
|
|
return false, nil
|
|
}
|
|
completed, err := o.DB.ListCompletedSiteIndices(ctx, item.ID, item.AttemptNumber)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
for _, s := range sites {
|
|
if !completed[s.Index] {
|
|
return false, nil
|
|
}
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
func (o *Orchestrator) aggregateAndRelease(ctx context.Context, item db.IPQueueItem) error {
|
|
if err := o.DB.SetAggregating(ctx, item.ID); err != nil {
|
|
return err
|
|
}
|
|
checks, err := o.DB.ListChecksForAttempt(ctx, item.ID, item.AttemptNumber)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
expected, err := o.expectedCheckCount(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("expected check count: %w", err)
|
|
}
|
|
passCount := 0
|
|
for _, c := range checks {
|
|
if c.Success {
|
|
passCount++
|
|
}
|
|
}
|
|
missing := expected - len(checks)
|
|
if missing < 0 {
|
|
missing = 0
|
|
}
|
|
failCount := (len(checks) - passCount) + missing
|
|
|
|
var result string
|
|
switch {
|
|
case passCount > 0 && failCount == 0:
|
|
result = db.ResultPass
|
|
case passCount == 0:
|
|
result = db.ResultFail
|
|
default:
|
|
result = db.ResultPartial
|
|
}
|
|
if missing > 0 && o.Agg.MissingCountsAsFail && result == db.ResultPass {
|
|
result = db.ResultPartial
|
|
}
|
|
|
|
if err := o.DB.FinishIP(ctx, item.ID, result); err != nil {
|
|
return err
|
|
}
|
|
o.event(ctx, "control-api", "", &item.ID, "aggregated",
|
|
fmt.Sprintf(`{"result":%q,"checks":%d,"passed":%d,"missing":%d}`, result, len(checks), passCount, missing))
|
|
|
|
if settings, err := o.DB.GetSettings(ctx); err != nil {
|
|
o.Log.Error("get settings for history retention", "ip_id", item.ID, "err", err)
|
|
} else if settings.HistoryRetentionCycles > 0 {
|
|
if err := o.DB.PruneRegistryHistory(ctx, item.RegistryID, settings.HistoryRetentionCycles); err != nil {
|
|
o.Log.Error("prune registry history", "ip_id", item.ID, "registry_id", item.RegistryID, "err", err)
|
|
}
|
|
}
|
|
|
|
if item.FIPID != "" {
|
|
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil {
|
|
o.Log.Error("disassociate fip", "ip_id", item.ID, "fip_id", item.FIPID, "err", err)
|
|
// Fall through and still free the validator/DB state — the
|
|
// lease sweep or an operator can reconcile a stuck Neutron
|
|
// association separately; we must not leave the validator
|
|
// wedged as "checking" forever over a cloud API hiccup.
|
|
}
|
|
}
|
|
if item.OwnerValidatorID != nil {
|
|
if err := o.DB.ReleaseFIP(ctx, item.ID, *item.OwnerValidatorID); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// expectedCheckCount is the number of check rows a fully-reported IP should
|
|
// have: one per (egress check-type x target) plus one per (site x inbound
|
|
// port/icmp probe). Reads the current check_types/targets/sites/inbound
|
|
// checks from the database, so a config change between assignment and
|
|
// aggregation is reflected in this specific aggregation (see the "accepted
|
|
// tradeoff" note in docs/PLAN_API_CONFIG_MANAGEMENT.md).
|
|
func (o *Orchestrator) expectedCheckCount(ctx context.Context) (int, error) {
|
|
resolved, err := o.DB.ListResolvedCheckTypes(ctx)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
egress := 0
|
|
for _, c := range resolved {
|
|
egress += len(c.Targets)
|
|
}
|
|
sites, err := o.DB.ListSites(ctx)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
inbound, err := o.DB.GetInboundChecks(ctx)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
inboundPerSite := len(inbound.Ports)
|
|
if inbound.ICMP {
|
|
inboundPerSite++
|
|
}
|
|
for _, p := range inbound.Ports {
|
|
if p == 443 || p == 22 {
|
|
inboundPerSite++
|
|
}
|
|
}
|
|
return egress + inboundPerSite*len(sites), nil
|
|
}
|
|
|
|
// sweepExpiredLeases reclaims non-terminal IPs whose lease has passed —
|
|
// this is both the "stuck/crashed validator" reclaim path and, since all
|
|
// state lives in SQLite, the control-api crash-recovery path: a freshly
|
|
// restarted process finds the same expired leases and reclaims them the
|
|
// same way, with no separate recovery code required.
|
|
func (o *Orchestrator) sweepExpiredLeases(ctx context.Context) error {
|
|
expired, err := o.DB.ListExpiredLeases(ctx, db.Now())
|
|
if err != nil {
|
|
return fmt.Errorf("list expired leases: %w", err)
|
|
}
|
|
for _, item := range expired {
|
|
validatorID := ""
|
|
if item.OwnerValidatorID != nil {
|
|
validatorID = *item.OwnerValidatorID
|
|
}
|
|
item, validatorID := item, validatorID
|
|
o.spawn(fmt.Sprintf("lease:%d", item.ID), func() {
|
|
o.Log.Info("lease expired, reclaiming", "ip_id", item.ID, "ip", item.IPAddress, "validator", validatorID)
|
|
if item.FIPID != "" {
|
|
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil {
|
|
o.Log.Error("disassociate fip on lease reclaim", "ip_id", item.ID, "err", err)
|
|
}
|
|
}
|
|
o.event(ctx, "control-api", "", &item.ID, "lease_expired", fmt.Sprintf(`{"validator_id":%q}`, validatorID))
|
|
o.requeueOrFail(ctx, item.ID, validatorID, "lease expired")
|
|
})
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// sweepStaleHeartbeats marks validators unreachable if they haven't
|
|
// heartbeated within HeartbeatTimeoutSeconds. It does not itself reclaim
|
|
// their in-flight IP — that happens independently via lease expiry, so a
|
|
// validator that stops heartbeating but whose lease hasn't yet expired
|
|
// still finishes its current check window if it recovers in time.
|
|
func (o *Orchestrator) SweepStaleHeartbeats(ctx context.Context) error {
|
|
cutoff := db.Now().Add(-time.Duration(o.Cfg.HeartbeatTimeoutSeconds) * time.Second)
|
|
stale, err := o.DB.ListStaleHeartbeats(ctx, cutoff)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, v := range stale {
|
|
if err := o.DB.MarkValidatorUnreachable(ctx, v.ValidatorID); err != nil {
|
|
o.Log.Error("mark validator unreachable", "validator", v.ValidatorID, "err", err)
|
|
continue
|
|
}
|
|
o.event(ctx, "control-api", "", nil, "validator_unreachable", fmt.Sprintf(`{"validator_id":%q}`, v.ValidatorID))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SweepStaleSiteHeartbeats marks prober sites unreachable if they haven't
|
|
// heartbeated within HeartbeatTimeoutSeconds — mirrors SweepStaleHeartbeats
|
|
// exactly, for sites instead of validators.
|
|
func (o *Orchestrator) SweepStaleSiteHeartbeats(ctx context.Context) error {
|
|
cutoff := db.Now().Add(-time.Duration(o.Cfg.HeartbeatTimeoutSeconds) * time.Second)
|
|
stale, err := o.DB.ListStaleSiteHeartbeats(ctx, cutoff)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, s := range stale {
|
|
if err := o.DB.MarkSiteUnreachable(ctx, s.SiteID); err != nil {
|
|
o.Log.Error("mark site unreachable", "site_id", s.SiteID, "err", err)
|
|
continue
|
|
}
|
|
o.event(ctx, "control-api", "", nil, "site_unreachable", fmt.Sprintf(`{"site_id":%q}`, s.SiteID))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// RecordEvent is the exported entry point httpapi uses to log
|
|
// agent/prober-reported audit events (config_received, fip_changed,
|
|
// error, etc.) through the same path as internally generated events
|
|
// (fip_associated, fip_occupied, retry_or_fail, etc.).
|
|
func (o *Orchestrator) RecordEvent(ctx context.Context, sourceType, sourceID string, ipID *int64, eventType, payload string) {
|
|
o.event(ctx, sourceType, sourceID, ipID, eventType, payload)
|
|
}
|
|
|
|
func (o *Orchestrator) event(ctx context.Context, sourceType, sourceID string, ipID *int64, eventType, payload string) {
|
|
if err := o.DB.InsertEvent(ctx, db.Event{
|
|
SourceType: sourceType,
|
|
SourceID: sourceID,
|
|
IPID: ipID,
|
|
EventType: eventType,
|
|
Payload: payload,
|
|
OccurredAt: db.Now(),
|
|
}); err != nil {
|
|
o.Log.Error("insert event", "type", eventType, "err", err)
|
|
}
|
|
}
|