// 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 } // 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, } } 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). 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) } } // 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) 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 { 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 } 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) } // 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.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) } } if err := o.DB.DeleteIP(ctx, item.ID); err != nil { return err } 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. func (o *Orchestrator) DeleteIPs(ctx context.Context, addresses []string) (db.DeleteIPsResult, error) { for _, addr := range addresses { item, err := o.DB.GetIPByAddress(ctx, addr) if err != nil { continue // unknown address — DB.DeleteIPs will report it in NotFound } 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) } } } result, err := o.DB.DeleteIPs(ctx, addresses) if err != nil { return result, err } 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, implemented as DeleteIPs over // the full current address list rather than a separate DB code path. func (o *Orchestrator) ClearQueue(ctx context.Context) (db.DeleteIPsResult, error) { items, err := o.DB.ListIPs(ctx) if err != nil { return db.DeleteIPsResult{}, fmt.Errorf("list ips: %w", err) } addresses := make([]string, len(items)) for i, item := range items { addresses[i] = item.IPAddress } for _, item := range items { if item.FIPID != "" { if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil { o.Log.Error("disassociate fip on clear queue", "ip_id", item.ID, "fip_id", item.FIPID, "err", err) } } } result, err := o.DB.DeleteIPs(ctx, addresses) if err != nil { return result, err } o.event(ctx, "control-api", "", nil, "queue_cleared", deletedAddressesPayload(result.Deleted)) return result, nil } // ScanFloatingIPs lists every floating IP in the OpenStack project, filters // to the ones not currently associated to any port (the free pool awaiting // validation before reissue), and submits that address list to the check // queue via db.SubmitIPs — the same entry point the admin API's "add // addresses" call uses, so add/requeue/reorder semantics are identical // whether the address list came from an operator or from this scan. Returns // the SubmitIPs outcome plus how many free floating IPs were found in total // (which can be larger than the sum of the SubmitIPsResult slices, since // addresses already mid-check are silently skipped — see db.SubmitIPs). func (o *Orchestrator) ScanFloatingIPs(ctx context.Context) (db.SubmitIPsResult, int, error) { fips, err := o.OS.ListFloatingIPs(ctx) if err != nil { return db.SubmitIPsResult{}, 0, fmt.Errorf("list floating ips: %w", err) } var free []string for _, f := range fips { if f.PortID == "" { free = append(free, f.Address) } } if len(free) == 0 { o.event(ctx, "control-api", "", nil, "fip_scan", `{"scanned_free":0}`) return db.SubmitIPsResult{}, 0, nil } result, err := o.DB.SubmitIPs(ctx, free) if err != nil { return result, len(free), fmt.Errorf("submit scanned ips: %w", err) } o.event(ctx, "control-api", "", nil, "fip_scan", fmt.Sprintf( `{"scanned_free":%d,"added":%d,"requeued":%d,"reordered":%d,"skipped_in_progress":%d}`, len(free), len(result.Added), len(result.Requeued), len(result.Reordered), len(result.SkippedInProgress))) return result, len(free), nil } // deletedAddressesPayload builds the event payload for the batch delete // operations — a proper JSON array via encoding/json rather than fmt's %q // slice formatting (which produces space-separated quoted strings, not // valid JSON). func deletedAddressesPayload(addresses []string) string { b, err := json.Marshal(struct { Addresses []string `json:"addresses"` }{addresses}) 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 } 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 } 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) } }