Files
cloud-ip-validator/internal/orchestrator/orchestrator.go
T
ayurishchevandClaude Sonnet 5.5 0532baff09 Keep one address per validator; fix heartbeat handling and queue clear
A mass check on 2026-10-02 stalled 7 of 20 validators and sent 42
addresses to fail without a single check. A validator busy with slow
checks went silent, was marked unreachable, and its next heartbeat put it
back to idle while it still held the address; it was handed a second one,
whose association never ran (the in-flight guard was keyed by validator),
and both waited for their leases to expire.

- Heartbeat/re-register return an unreachable validator to assigned when
  it still holds an address, else idle.
- A validator is released only from the address it currently holds
  (ReleaseFIP, RequeueOrFail, MarkFIPOccupied, FreeValidator); an
  unreachable validator stays unreachable until its next heartbeat, so a
  dead validator is no longer handed a new address every lease period.
- ClaimNextQueued refuses a validator that still has an address; a
  ReconcileValidators pass on every tick repairs rows that disagree with
  the queue.
- Association guard is keyed by address, not validator.
- The agent sends heartbeats from their own goroutine.
- Clear queue / delete: detach only floating IPs of unfinished rows (done,
  failed and occupied rows kept their fip_id and made a clear issue >1000
  sequential cloud calls: 256 s), at most 8 in parallel; the operation no
  longer dies with the client connection (10 minute limit).

Includes the incident analysis and the plan under analysis/ and
docs/changes/, and rebuilt bin/control-api and bin/validator-agent.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2026-10-02 14:42:12 +03:00

879 lines
33 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 n, err := o.DB.ReconcileValidators(ctx); err != nil {
o.Log.Error("reconcile validators", "err", err)
} else if n > 0 {
o.Log.Warn("repaired validators that disagreed with the queue", "count", n)
}
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
// Keyed by address, not validator: the key only guards against
// starting the same association twice. A validator-wide key made a
// second address claimed while the first was still associating skip
// its association and wait for the lease to expire.
o.spawn(fmt.Sprintf("assign:%d", item.ID), 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, item.ID); 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)
}
o.disassociateAll(ctx, refs, "delete")
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
}
// maxParallelDetach bounds how many floating IPs a bulk delete or clear
// detaches at the same time (each is one Neutron call).
const maxParallelDetach = 8
// disassociateAll detaches the given floating IPs best-effort, at most
// maxParallelDetach at a time; a failure is logged and does not stop the rest.
func (o *Orchestrator) disassociateAll(ctx context.Context, refs []db.FIPRef, what string) {
var wg sync.WaitGroup
sem := make(chan struct{}, maxParallelDetach)
for _, ref := range refs {
ref := ref
sem <- struct{}{}
wg.Add(1)
go func() {
defer wg.Done()
defer func() { <-sem }()
if err := o.OS.DisassociateFloatingIP(ctx, ref.FIPID); err != nil {
o.Log.Error("disassociate fip on "+what, "ip_id", ref.IPID, "fip_id", ref.FIPID, "err", err)
}
}()
}
wg.Wait()
}
// 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)
}
o.disassociateAll(ctx, refs, "clear queue")
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)
}
}