Add Floating IP scanning and a durable address registry with configurable history depth

Adds POST /api/v1/admin/ips/scan (plus an optional periodic ticker) to
discover free Floating IPs in the OpenStack project and feed them straight
into the check queue. More importantly, decouples check/event history from
ip_queue's lifecycle: a new ip_registry table (migration 0007) gives every
address ever submitted a durable identity, so deleting it from the queue no
longer destroys its history — it's still reachable via the new
GET /api/v1/admin/registry[/{ip}] endpoints and the dashboard's /registry
pages, with retention depth configurable in check cycles per address
(history_retention_cycles, 0 = unlimited).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
ayurishchevandClaude Sonnet 5 committed 2026-09-23 09:52:01 +03:00
1 parent 78b20fa5be
commit 582b44f314
47 files changed
+1781 -107

No files matched your search

+4
View File
@@ -31,6 +31,9 @@ var unboundedSitesSchema string
//go:embed migrations/0006_prober_heartbeat.sql
var proberHeartbeatSchema string
//go:embed migrations/0007_ip_registry.sql
var ipRegistrySchema string
// migrations is the ordered list of schema versions. Each entry's SQL is
// applied, in order, for any version greater than the database's current
// PRAGMA user_version — so a fresh database walks the whole list and an
@@ -45,6 +48,7 @@ var migrations = []struct {
{4, inboundChecksAdminSchema},
{5, unboundedSitesSchema},
{6, proberHeartbeatSchema},
{7, ipRegistrySchema},
}
type DB struct {
@@ -0,0 +1,80 @@
-- Durable per-address registry, decoupled from ip_queue's lifecycle: today
-- deleting an address from ip_queue (DeleteIP/DeleteIPs/ClearQueue) cascades
-- to a hard DELETE of its checks/events, so history is lost forever if an
-- address is removed and later re-added. ip_registry gives every address
-- ever submitted a durable identity that check/event history attaches to
-- instead, surviving ip_queue row deletion and recreation.
CREATE TABLE ip_registry (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ip_address TEXT NOT NULL UNIQUE,
first_seen_at TIMESTAMP NOT NULL,
last_seen_at TIMESTAMP NOT NULL,
next_cycle INTEGER NOT NULL DEFAULT 1,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
-- Backfill: one registry row per address currently (or ever) in ip_queue.
-- ip_queue.ip_address is already UNIQUE, so this is a straight 1:1 copy.
-- next_cycle starts past the current attempt_number so the first future
-- resubmission of an address gets a cycle_id that has never been used
-- before, even though it's seeded from attempt_number here.
INSERT INTO ip_registry (ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at)
SELECT ip_address, created_at, updated_at, attempt_number + 1, created_at, updated_at FROM ip_queue;
ALTER TABLE ip_queue ADD COLUMN registry_id INTEGER REFERENCES ip_registry(id);
ALTER TABLE ip_queue ADD COLUMN cycle_id INTEGER NOT NULL DEFAULT 1;
UPDATE ip_queue SET
registry_id = (SELECT id FROM ip_registry r WHERE r.ip_address = ip_queue.ip_address),
cycle_id = attempt_number;
-- checks.ip_id is today NOT NULL + REFERENCES ip_queue(id), which is exactly
-- what forces the cascading DELETE on ip_queue row removal (foreign_keys=ON
-- would otherwise block the delete once ip_queue's row disappears out from
-- under a referencing row). To let history outlive its ip_queue row, ip_id
-- must become nullable and checks must carry their own durable registry_id.
-- SQLite has no ALTER to relax a column's NOT NULL/REFERENCES, so the table
-- is rebuilt.
CREATE TABLE checks_new (
id INTEGER PRIMARY KEY AUTOINCREMENT,
registry_id INTEGER NOT NULL REFERENCES ip_registry(id),
cycle_id INTEGER NOT NULL,
ip_id INTEGER REFERENCES ip_queue(id),
ip_address TEXT NOT NULL,
attempt_number INTEGER NOT NULL,
validator_id TEXT NOT NULL DEFAULT '',
source TEXT NOT NULL,
check_type TEXT NOT NULL,
target TEXT NOT NULL DEFAULT '',
success BOOLEAN NOT NULL,
latency_ms INTEGER NOT NULL DEFAULT 0,
detail TEXT NOT NULL DEFAULT '',
checked_at TIMESTAMP NOT NULL,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE(registry_id, cycle_id, source, check_type, target)
);
INSERT INTO checks_new (id, registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id,
source, check_type, target, success, latency_ms, detail, checked_at, created_at)
SELECT c.id, iq.registry_id, iq.cycle_id, c.ip_id, c.ip_address, c.attempt_number, c.validator_id,
c.source, c.check_type, c.target, c.success, c.latency_ms, c.detail, c.checked_at, c.created_at
FROM checks c JOIN ip_queue iq ON iq.id = c.ip_id;
DROP TABLE checks;
ALTER TABLE checks_new RENAME TO checks;
CREATE INDEX idx_checks_registry_cycle ON checks(registry_id, cycle_id);
CREATE INDEX idx_checks_ip_attempt ON checks(ip_id, attempt_number);
-- events.ip_id is already nullable, so no rebuild is needed there — just
-- add the durable registry_id (plus cycle_id, so retention pruning can cut
-- events at the same cycle boundary as checks) alongside it.
ALTER TABLE events ADD COLUMN registry_id INTEGER REFERENCES ip_registry(id);
ALTER TABLE events ADD COLUMN cycle_id INTEGER NOT NULL DEFAULT 0;
UPDATE events SET
registry_id = (SELECT registry_id FROM ip_queue WHERE ip_queue.id = events.ip_id),
cycle_id = (SELECT cycle_id FROM ip_queue WHERE ip_queue.id = events.ip_id)
WHERE ip_id IS NOT NULL;
CREATE INDEX idx_events_registry ON events(registry_id, cycle_id);
-- Configurable history retention depth, in check cycles per address. 0 (the
-- default, matching today's unbounded behavior) means keep everything.
ALTER TABLE settings ADD COLUMN history_retention_cycles INTEGER NOT NULL DEFAULT 0;
+37 -3
View File
@@ -95,12 +95,21 @@ type IPQueueItem struct {
FIPAssociatedAt *time.Time
AggregatedAt *time.Time
FIPReleasedAt *time.Time
RegistryID int64
CycleID int
CreatedAt time.Time
UpdatedAt time.Time
}
type Check struct {
ID int64
ID int64
RegistryID int64
CycleID int
// IPID is the ip_queue row this check was originally recorded against.
// It's cleared to 0 (SQL NULL) if that row was later deleted — history
// stays reachable via RegistryID/CycleID regardless (see
// migrations/0007_ip_registry.sql). 0 is never a valid ip_queue id
// (AUTOINCREMENT starts at 1), so it unambiguously means "orphaned."
IPID int64
IPAddress string
AttemptNumber int
@@ -115,11 +124,33 @@ type Check struct {
CreatedAt time.Time
}
// RegistryItem is a durable per-address record that survives an address
// being removed from ip_queue and later re-added — see
// migrations/0007_ip_registry.sql. NextCycle is the cycle_id that will be
// assigned the next time this address is (re)submitted; it only ever
// increases, so cycle_id stays unique for this address even across
// ip_queue row deletion/recreation.
type RegistryItem struct {
ID int64
IPAddress string
FirstSeenAt time.Time
LastSeenAt time.Time
NextCycle int
CreatedAt time.Time
UpdatedAt time.Time
}
type Event struct {
ID int64
SourceType string
SourceID string
IPID *int64
// RegistryID/CycleID are resolved from IPID at insert time (see
// InsertEvent) and stay set even after the ip_queue row IPID pointed to
// is later deleted, so retention pruning can cut events at the same
// cycle boundary as checks — see migrations/0007_ip_registry.sql.
RegistryID int64
CycleID int
EventType string
Payload string
OccurredAt time.Time
@@ -179,8 +210,11 @@ type DeleteIPsResult struct {
// admin-configurable at runtime (see queries_settings.go).
type Settings struct {
FIPSettleSeconds int
CreatedAt time.Time
UpdatedAt time.Time
// HistoryRetentionCycles caps how many recent check cycles are kept per
// registry address (see PruneRegistryHistory); 0 means unlimited.
HistoryRetentionCycles int
CreatedAt time.Time
UpdatedAt time.Time
}
// InboundChecksSettings is the singleton row describing what the prober
+52 -14
View File
@@ -2,50 +2,88 @@ package db
import (
"context"
"database/sql"
)
// UpsertCheck records (or, on retry, overwrites) a single check result. The
// UNIQUE(ip_id, attempt_number, source, check_type, target) constraint plus
// UpsertCheck records (or, on retry, overwrites) a single check result.
// registry_id/cycle_id are resolved from c.IPID's current ip_queue row at
// write time, so callers (agentcore/probercore) never need to know about
// the registry — see migrations/0007_ip_registry.sql. The
// UNIQUE(registry_id, cycle_id, source, check_type, target) constraint plus
// this upsert is what makes agent/prober result submission safely
// retryable without producing duplicate rows.
// retryable without producing duplicate rows, now scoped to the durable
// per-address cycle rather than the ip_queue row's attempt_number, so it
// survives that row being deleted and the address later resubmitted.
func (d *DB) UpsertCheck(ctx context.Context, c Check) error {
_, err := d.ExecContext(ctx, `
INSERT INTO checks (ip_id, ip_address, attempt_number, validator_id, source, check_type, target,
success, latency_ms, detail, checked_at, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(ip_id, attempt_number, source, check_type, target) DO UPDATE SET
INSERT INTO checks (registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id,
source, check_type, target, success, latency_ms, detail, checked_at, created_at)
VALUES ((SELECT registry_id FROM ip_queue WHERE id=?), (SELECT cycle_id FROM ip_queue WHERE id=?),
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(registry_id, cycle_id, source, check_type, target) DO UPDATE SET
validator_id=excluded.validator_id,
success=excluded.success,
latency_ms=excluded.latency_ms,
detail=excluded.detail,
checked_at=excluded.checked_at
`, c.IPID, c.IPAddress, c.AttemptNumber, c.ValidatorID, c.Source, c.CheckType, c.Target,
`, c.IPID, c.IPID, c.IPID, c.IPAddress, c.AttemptNumber, c.ValidatorID, c.Source, c.CheckType, c.Target,
c.Success, c.LatencyMS, c.Detail, timeToDB(c.CheckedAt), timeToDB(Now()))
return err
}
const checksSelect = `
SELECT id, registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id, source, check_type, target,
success, latency_ms, detail, checked_at, created_at
FROM checks
`
// ListChecksForAttempt returns every check recorded for an IP's current
// attempt — the input to overall-result aggregation.
func (d *DB) ListChecksForAttempt(ctx context.Context, ipID int64, attemptNumber int) ([]Check, error) {
rows, err := d.QueryContext(ctx, `
SELECT id, ip_id, ip_address, attempt_number, validator_id, source, check_type, target,
success, latency_ms, detail, checked_at, created_at
FROM checks WHERE ip_id=? AND attempt_number=?
rows, err := d.QueryContext(ctx, checksSelect+`
WHERE ip_id=? AND attempt_number=?
ORDER BY source, check_type, target
`, ipID, attemptNumber)
if err != nil {
return nil, err
}
defer rows.Close()
return scanChecks(rows)
}
// ListChecksForRegistry returns an address's check history across every
// cycle still retained (see PruneRegistryHistory), newest cycle first. A
// nil limit returns everything currently retained.
func (d *DB) ListChecksForRegistry(ctx context.Context, registryID int64, limit *int) ([]Check, error) {
query := checksSelect + `WHERE registry_id=? ORDER BY cycle_id DESC, source, check_type, target`
args := []interface{}{registryID}
if limit != nil {
query += ` LIMIT ?`
args = append(args, *limit)
}
rows, err := d.QueryContext(ctx, query, args...)
if err != nil {
return nil, err
}
defer rows.Close()
return scanChecks(rows)
}
func scanChecks(rows *sql.Rows) ([]Check, error) {
var out []Check
for rows.Next() {
var c Check
var ipID sql.NullInt64
var checkedAt, createdAt string
if err := rows.Scan(&c.ID, &c.IPID, &c.IPAddress, &c.AttemptNumber, &c.ValidatorID, &c.Source,
&c.CheckType, &c.Target, &c.Success, &c.LatencyMS, &c.Detail, &checkedAt, &createdAt); err != nil {
if err := rows.Scan(&c.ID, &c.RegistryID, &c.CycleID, &ipID, &c.IPAddress, &c.AttemptNumber,
&c.ValidatorID, &c.Source, &c.CheckType, &c.Target, &c.Success, &c.LatencyMS, &c.Detail,
&checkedAt, &createdAt); err != nil {
return nil, err
}
if ipID.Valid {
c.IPID = ipID.Int64
}
var err error
if c.CheckedAt, err = dbToTime(checkedAt); err != nil {
return nil, err
}
+19 -7
View File
@@ -1,12 +1,19 @@
package db
import "context"
import (
"context"
"database/sql"
)
// InsertEvent records an audit-trail row. When e.IPID is set, registry_id
// and cycle_id are resolved from that ip_queue row's current values, so the
// event stays traceable to its address (via registry_id) and prunable at
// the right cycle boundary even after the ip_queue row itself is deleted.
func (d *DB) InsertEvent(ctx context.Context, e Event) error {
_, err := d.ExecContext(ctx, `
INSERT INTO events (source_type, source_id, ip_id, event_type, payload, occurred_at, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?)
`, e.SourceType, e.SourceID, e.IPID, e.EventType, e.Payload, timeToDB(e.OccurredAt), timeToDB(Now()))
INSERT INTO events (source_type, source_id, ip_id, registry_id, cycle_id, event_type, payload, occurred_at, created_at)
VALUES (?, ?, ?, (SELECT registry_id FROM ip_queue WHERE id=?), COALESCE((SELECT cycle_id FROM ip_queue WHERE id=?), 0), ?, ?, ?, ?)
`, e.SourceType, e.SourceID, e.IPID, e.IPID, e.IPID, e.EventType, e.Payload, timeToDB(e.OccurredAt), timeToDB(Now()))
return err
}
@@ -14,7 +21,7 @@ func (d *DB) InsertEvent(ctx context.Context, e Event) error {
// first — used by the admin detail endpoint.
func (d *DB) ListEventsForIP(ctx context.Context, ipID int64) ([]Event, error) {
rows, err := d.QueryContext(ctx, `
SELECT id, source_type, source_id, ip_id, event_type, payload, occurred_at, created_at
SELECT id, source_type, source_id, ip_id, registry_id, cycle_id, event_type, payload, occurred_at, created_at
FROM events WHERE ip_id=? ORDER BY occurred_at DESC
`, ipID)
if err != nil {
@@ -26,7 +33,7 @@ func (d *DB) ListEventsForIP(ctx context.Context, ipID int64) ([]Event, error) {
func (d *DB) ListRecentEvents(ctx context.Context, limit int) ([]Event, error) {
rows, err := d.QueryContext(ctx, `
SELECT id, source_type, source_id, ip_id, event_type, payload, occurred_at, created_at
SELECT id, source_type, source_id, ip_id, registry_id, cycle_id, event_type, payload, occurred_at, created_at
FROM events ORDER BY id DESC LIMIT ?
`, limit)
if err != nil {
@@ -45,11 +52,16 @@ func scanEvents(rows interface {
for rows.Next() {
var e Event
var ipID *int64
var registryID sql.NullInt64
var occurredAt, createdAt string
if err := rows.Scan(&e.ID, &e.SourceType, &e.SourceID, &ipID, &e.EventType, &e.Payload, &occurredAt, &createdAt); err != nil {
if err := rows.Scan(&e.ID, &e.SourceType, &e.SourceID, &ipID, &registryID, &e.CycleID,
&e.EventType, &e.Payload, &occurredAt, &createdAt); err != nil {
return nil, err
}
e.IPID = ipID
if registryID.Valid {
e.RegistryID = registryID.Int64
}
var err error
if e.OccurredAt, err = dbToTime(occurredAt); err != nil {
return nil, err
+67 -20
View File
@@ -21,12 +21,26 @@ func (d *DB) SeedQueue(ctx context.Context, addresses []string) error {
now := timeToDB(Now())
for i, addr := range addresses {
_, err := tx.ExecContext(ctx, `
INSERT INTO ip_queue (ip_address, sequence, state, created_at, updated_at)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(ip_address) DO NOTHING
`, addr, i, IPQueued, now, now)
var exists bool
if err := tx.QueryRowContext(ctx, `SELECT EXISTS(SELECT 1 FROM ip_queue WHERE ip_address=?)`, addr).Scan(&exists); err != nil {
return fmt.Errorf("check existing %s: %w", addr, err)
}
if exists {
continue
}
registryID, err := findOrCreateRegistryTx(ctx, tx, addr, now)
if err != nil {
return err
}
cycle, err := nextRegistryCycleTx(ctx, tx, registryID, now)
if err != nil {
return err
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO ip_queue (ip_address, sequence, state, registry_id, cycle_id, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(ip_address) DO NOTHING
`, addr, i, IPQueued, registryID, cycle, now, now); err != nil {
return fmt.Errorf("seed %s: %w", addr, err)
}
}
@@ -227,13 +241,21 @@ func (d *DB) RequeueOrFail(ctx context.Context, ipID int64, validatorID string,
}
if nextState == IPQueued {
var registryID int64
if err := tx.QueryRowContext(ctx, `SELECT registry_id FROM ip_queue WHERE id=?`, ipID).Scan(&registryID); err != nil {
return err
}
cycle, cErr := nextRegistryCycleTx(ctx, tx, registryID, now)
if cErr != nil {
return cErr
}
_, err = tx.ExecContext(ctx, `
UPDATE ip_queue SET
state=?, owner_validator_id=NULL, fip_id='', retry_count=?, attempt_number=attempt_number+1,
lease_expires_at=NULL, egress_complete=0,
cycle_id=?, lease_expires_at=NULL, egress_complete=0,
overall_result='', assigned_at=NULL, fip_associated_at=NULL, updated_at=?
WHERE id=?
`, nextState, retryCount, now, ipID)
`, nextState, retryCount, cycle, now, ipID)
} else {
_, err = tx.ExecContext(ctx, `
UPDATE ip_queue SET
@@ -300,10 +322,18 @@ func (d *DB) SubmitIPs(ctx context.Context, addresses []string) (SubmitIPsResult
err := tx.QueryRowContext(ctx, `SELECT state FROM ip_queue WHERE ip_address=?`, addr).Scan(&state)
switch {
case err == sql.ErrNoRows:
registryID, rErr := findOrCreateRegistryTx(ctx, tx, addr, now)
if rErr != nil {
return result, rErr
}
cycle, cErr := nextRegistryCycleTx(ctx, tx, registryID, now)
if cErr != nil {
return result, cErr
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO ip_queue (ip_address, sequence, state, created_at, updated_at)
VALUES (?, ?, ?, ?, ?)
`, addr, seq, IPQueued, now, now); err != nil {
INSERT INTO ip_queue (ip_address, sequence, state, registry_id, cycle_id, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?)
`, addr, seq, IPQueued, registryID, cycle, now, now); err != nil {
return result, fmt.Errorf("insert %s: %w", addr, err)
}
result.Added = append(result.Added, addr)
@@ -312,14 +342,22 @@ func (d *DB) SubmitIPs(ctx context.Context, addresses []string) (SubmitIPsResult
return result, err
case state == IPDone || state == IPFailed || state == IPOccupied:
var registryID int64
if err := tx.QueryRowContext(ctx, `SELECT registry_id FROM ip_queue WHERE ip_address=?`, addr).Scan(&registryID); err != nil {
return result, err
}
cycle, cErr := nextRegistryCycleTx(ctx, tx, registryID, now)
if cErr != nil {
return result, cErr
}
if _, err := tx.ExecContext(ctx, `
UPDATE ip_queue SET
state=?, sequence=?, owner_validator_id=NULL, fip_id='', retry_count=0,
attempt_number=attempt_number+1, lease_expires_at=NULL, egress_complete=0,
attempt_number=attempt_number+1, cycle_id=?, lease_expires_at=NULL, egress_complete=0,
overall_result='',
assigned_at=NULL, fip_associated_at=NULL, aggregated_at=NULL, fip_released_at=NULL, updated_at=?
WHERE ip_address=?
`, IPQueued, seq, now, addr); err != nil {
`, IPQueued, seq, cycle, now, addr); err != nil {
return result, fmt.Errorf("requeue %s: %w", addr, err)
}
result.Requeued = append(result.Requeued, addr)
@@ -427,8 +465,11 @@ func (d *DB) DeleteIPs(ctx context.Context, addresses []string) (DeleteIPsResult
}
// deleteIPTx is the shared body of DeleteIP/DeleteIPs: free the owning
// validator, delete dependent checks/events, then the ip_queue row itself
// — the same FK-clearing order DeleteValidator uses for owner_validator_id.
// validator, detach dependent checks/events from the doomed ip_queue row
// (their registry_id already anchors them permanently — see
// migrations/0007_ip_registry.sql — so this is a hand-off, not a loss),
// delete the ephemeral per-attempt ip_site_checks progress flags, then the
// ip_queue row itself.
func deleteIPTx(ctx context.Context, tx *sql.Tx, ipID int64) error {
now := timeToDB(Now())
if _, err := tx.ExecContext(ctx, `
@@ -437,11 +478,11 @@ func deleteIPTx(ctx context.Context, tx *sql.Tx, ipID int64) error {
`, ValidatorIdle, now, ipID); err != nil {
return fmt.Errorf("free owning validator: %w", err)
}
if _, err := tx.ExecContext(ctx, `DELETE FROM checks WHERE ip_id=?`, ipID); err != nil {
return fmt.Errorf("delete checks: %w", err)
if _, err := tx.ExecContext(ctx, `UPDATE checks SET ip_id=NULL WHERE ip_id=?`, ipID); err != nil {
return fmt.Errorf("detach checks: %w", err)
}
if _, err := tx.ExecContext(ctx, `DELETE FROM events WHERE ip_id=?`, ipID); err != nil {
return fmt.Errorf("delete events: %w", err)
if _, err := tx.ExecContext(ctx, `UPDATE events SET ip_id=NULL WHERE ip_id=?`, ipID); err != nil {
return fmt.Errorf("detach events: %w", err)
}
if _, err := tx.ExecContext(ctx, `DELETE FROM ip_site_checks WHERE ip_id=?`, ipID); err != nil {
return fmt.Errorf("delete ip_site_checks: %w", err)
@@ -507,7 +548,8 @@ func (d *DB) ListExpiredLeases(ctx context.Context, now time.Time) ([]IPQueueIte
const ipQueueSelect = `
SELECT id, ip_address, sequence, state, owner_validator_id, fip_id, attempt_number, retry_count,
lease_expires_at, egress_complete, overall_result,
assigned_at, fip_associated_at, aggregated_at, fip_released_at, created_at, updated_at
assigned_at, fip_associated_at, aggregated_at, fip_released_at,
registry_id, cycle_id, created_at, updated_at
FROM ip_queue
`
@@ -527,18 +569,23 @@ func scanIPQueueItem(row rowScanner) (*IPQueueItem, error) {
var item IPQueueItem
var owner sql.NullString
var leaseExpires, assignedAt, fipAssociatedAt, aggregatedAt, fipReleasedAt sql.NullString
var registryID sql.NullInt64
var createdAt, updatedAt string
if err := row.Scan(
&item.ID, &item.IPAddress, &item.Sequence, &item.State, &owner, &item.FIPID,
&item.AttemptNumber, &item.RetryCount, &leaseExpires,
&item.EgressComplete,
&item.OverallResult, &assignedAt, &fipAssociatedAt, &aggregatedAt, &fipReleasedAt, &createdAt, &updatedAt,
&item.OverallResult, &assignedAt, &fipAssociatedAt, &aggregatedAt, &fipReleasedAt,
&registryID, &item.CycleID, &createdAt, &updatedAt,
); err != nil {
return nil, err
}
if owner.Valid {
item.OwnerValidatorID = &owner.String
}
if registryID.Valid {
item.RegistryID = registryID.Int64
}
var err error
if item.LeaseExpiresAt, err = nullStringToTimePtr(leaseExpires); err != nil {
return nil, err
+216
View File
@@ -0,0 +1,216 @@
package db
import (
"context"
"database/sql"
"fmt"
"time"
)
// findOrCreateRegistryTx returns the ip_registry row id for addr, creating
// it (next_cycle starting at 1) the first time this address is ever seen.
// Safe to call on every (re)submission of an address — ip_address is
// UNIQUE, so a second call for the same address is a no-op lookup.
func findOrCreateRegistryTx(ctx context.Context, tx *sql.Tx, addr string, now string) (int64, error) {
if _, err := tx.ExecContext(ctx, `
INSERT INTO ip_registry (ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at)
VALUES (?, ?, ?, 1, ?, ?)
ON CONFLICT(ip_address) DO NOTHING
`, addr, now, now, now, now); err != nil {
return 0, fmt.Errorf("find or create registry for %s: %w", addr, err)
}
var id int64
if err := tx.QueryRowContext(ctx, `SELECT id FROM ip_registry WHERE ip_address=?`, addr).Scan(&id); err != nil {
return 0, fmt.Errorf("lookup registry id for %s: %w", addr, err)
}
return id, nil
}
// nextRegistryCycleTx allocates the next cycle_id for registryID and
// advances last_seen_at. Call once per fresh generation of check history for
// an address — a brand new ip_queue row, an admin-triggered requeue, or a
// system retry that bumps attempt_number — so every generation's checks get
// a cycle_id no other generation, past or future (even across ip_queue row
// deletion and recreation), will ever reuse.
func nextRegistryCycleTx(ctx context.Context, tx *sql.Tx, registryID int64, now string) (int, error) {
var cycle int
if err := tx.QueryRowContext(ctx, `SELECT next_cycle FROM ip_registry WHERE id=?`, registryID).Scan(&cycle); err != nil {
return 0, fmt.Errorf("read next_cycle for registry %d: %w", registryID, err)
}
if _, err := tx.ExecContext(ctx, `
UPDATE ip_registry SET next_cycle=next_cycle+1, last_seen_at=?, updated_at=? WHERE id=?
`, now, now, registryID); err != nil {
return 0, fmt.Errorf("advance next_cycle for registry %d: %w", registryID, err)
}
return cycle, nil
}
// RegistrySummary is one row of the full IP registry, combining the durable
// ip_registry record with a rollup of its check history and, if the address
// currently has a live ip_queue row, that row's state.
type RegistrySummary struct {
RegistryItem
TotalCycles int
LastResult string
LastCheckedAt *time.Time
InQueue bool
CurrentState string
}
// ListRegistry returns every address ever submitted, newest first-seen
// last, each with a summary of its accumulated check history. Reads two
// simple queries plus one aggregate rather than a single large join, since
// the registry is expected to stay small enough (one row per distinct
// address ever seen) that this is simpler to reason about than a
// multi-way correlated subquery.
func (d *DB) ListRegistry(ctx context.Context) ([]RegistrySummary, error) {
rows, err := d.QueryContext(ctx, `
SELECT id, ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at
FROM ip_registry ORDER BY first_seen_at
`)
if err != nil {
return nil, err
}
defer rows.Close()
items, err := scanRegistryItems(rows)
if err != nil {
return nil, err
}
out := make([]RegistrySummary, len(items))
for i, item := range items {
s := RegistrySummary{RegistryItem: item}
if err := d.fillRegistrySummary(ctx, &s); err != nil {
return nil, err
}
out[i] = s
}
return out, nil
}
// GetRegistryByAddress returns the registry row (with summary) for a single
// address, or ErrNotFound if it has never been submitted.
func (d *DB) GetRegistryByAddress(ctx context.Context, address string) (*RegistrySummary, error) {
row := d.QueryRowContext(ctx, `
SELECT id, ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at
FROM ip_registry WHERE ip_address=?
`, address)
item, err := scanRegistryItem(row)
if err != nil {
if err == sql.ErrNoRows {
return nil, fmt.Errorf("ip %q: %w", address, ErrNotFound)
}
return nil, err
}
s := &RegistrySummary{RegistryItem: *item}
if err := d.fillRegistrySummary(ctx, s); err != nil {
return nil, err
}
return s, nil
}
func (d *DB) fillRegistrySummary(ctx context.Context, s *RegistrySummary) error {
if err := d.QueryRowContext(ctx, `
SELECT COUNT(DISTINCT cycle_id) FROM checks WHERE registry_id=?
`, s.ID).Scan(&s.TotalCycles); err != nil {
return err
}
var lastResult sql.NullBool
var lastCheckedAt sql.NullString
if err := d.QueryRowContext(ctx, `
SELECT success, checked_at FROM checks WHERE registry_id=? ORDER BY cycle_id DESC, checked_at DESC LIMIT 1
`, s.ID).Scan(&lastResult, &lastCheckedAt); err != nil && err != sql.ErrNoRows {
return err
}
if lastResult.Valid {
if lastResult.Bool {
s.LastResult = "pass"
} else {
s.LastResult = "fail"
}
}
t, err := nullStringToTimePtr(lastCheckedAt)
if err != nil {
return err
}
s.LastCheckedAt = t
var state sql.NullString
if err := d.QueryRowContext(ctx, `SELECT state FROM ip_queue WHERE registry_id=?`, s.ID).Scan(&state); err != nil && err != sql.ErrNoRows {
return err
}
if state.Valid {
s.InQueue = true
s.CurrentState = state.String
}
return nil
}
func scanRegistryItems(rows *sql.Rows) ([]RegistryItem, error) {
var out []RegistryItem
for rows.Next() {
item, err := scanRegistryItem(rows)
if err != nil {
return nil, err
}
out = append(out, *item)
}
return out, rows.Err()
}
func scanRegistryItem(row rowScanner) (*RegistryItem, error) {
var item RegistryItem
var firstSeenAt, lastSeenAt, createdAt, updatedAt string
if err := row.Scan(&item.ID, &item.IPAddress, &firstSeenAt, &lastSeenAt, &item.NextCycle, &createdAt, &updatedAt); err != nil {
return nil, err
}
var err error
if item.FirstSeenAt, err = dbToTime(firstSeenAt); err != nil {
return nil, err
}
if item.LastSeenAt, err = dbToTime(lastSeenAt); err != nil {
return nil, err
}
if item.CreatedAt, err = dbToTime(createdAt); err != nil {
return nil, err
}
if item.UpdatedAt, err = dbToTime(updatedAt); err != nil {
return nil, err
}
return &item, nil
}
// PruneRegistryHistory deletes all but the newest keepCycles cycles' worth
// of check history for registryID — a no-op if keepCycles <= 0 (unlimited
// retention) or the address has that many cycles or fewer. events rows tied
// to this registry are pruned to the same cutoff; ip_registry itself and its
// next_cycle counter are never touched, so pruning never risks a future
// cycle_id collision.
func (d *DB) PruneRegistryHistory(ctx context.Context, registryID int64, keepCycles int) error {
if keepCycles <= 0 {
return nil
}
var cutoff sql.NullInt64
err := d.QueryRowContext(ctx, `
SELECT MIN(cycle_id) FROM (
SELECT DISTINCT cycle_id FROM checks WHERE registry_id=? ORDER BY cycle_id DESC LIMIT ?
)
`, registryID, keepCycles).Scan(&cutoff)
if err != nil {
return fmt.Errorf("find prune cutoff for registry %d: %w", registryID, err)
}
if !cutoff.Valid {
return nil // fewer than keepCycles cycles recorded — nothing to prune
}
if _, err := d.ExecContext(ctx, `DELETE FROM checks WHERE registry_id=? AND cycle_id<?`, registryID, cutoff.Int64); err != nil {
return fmt.Errorf("prune checks for registry %d: %w", registryID, err)
}
// Cut events at the same cycle boundary as checks, so an address's audit
// trail and its check history are retained (and pruned) in lockstep.
if _, err := d.ExecContext(ctx, `DELETE FROM events WHERE registry_id=? AND cycle_id<?`, registryID, cutoff.Int64); err != nil {
return fmt.Errorf("prune events for registry %d: %w", registryID, err)
}
return nil
}
+242
View File
@@ -0,0 +1,242 @@
package db
import (
"errors"
"testing"
"cloudipvalidator/internal/config"
)
// TestRegistryHistorySurvivesIPDeletion proves that deleting an address
// from ip_queue no longer destroys its check history — the whole point of
// ip_registry (see migrations/0007_ip_registry.sql).
func TestRegistryHistorySurvivesIPDeletion(t *testing.T) {
d, ctx := newTestDB(t)
if _, err := d.SubmitIPs(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatalf("submit ips: %v", err)
}
ip, err := d.GetIPByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get ip: %v", err)
}
if err := d.UpsertCheck(ctx, Check{
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
Source: SourceEgress, CheckType: "https", Target: "https://example.test",
Success: true, CheckedAt: Now(),
}); err != nil {
t.Fatalf("upsert check: %v", err)
}
if err := d.DeleteIP(ctx, ip.ID); err != nil {
t.Fatalf("delete ip: %v", err)
}
// The live queue no longer knows about the address...
if _, err := d.GetIPByAddress(ctx, "1.2.3.4"); err == nil {
t.Fatalf("expected ip_queue row gone after delete")
}
// ...but the registry and its check history are untouched.
summary, err := d.GetRegistryByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get registry: %v", err)
}
if summary.TotalCycles != 1 {
t.Fatalf("expected 1 cycle retained, got %d", summary.TotalCycles)
}
checks, err := d.ListChecksForRegistry(ctx, summary.ID, nil)
if err != nil {
t.Fatalf("list checks for registry: %v", err)
}
if len(checks) != 1 {
t.Fatalf("expected 1 retained check, got %+v", checks)
}
if checks[0].IPID != 0 {
t.Fatalf("expected orphaned check's ip_id cleared to 0, got %d", checks[0].IPID)
}
}
// TestRegistryCycleSurvivesReaddWithoutCollision proves that deleting an
// address and resubmitting it later gives the new generation of checks a
// cycle_id distinct from the old one, so both generations' history is kept
// as separate rows rather than colliding on the UNIQUE constraint.
func TestRegistryCycleSurvivesReaddWithoutCollision(t *testing.T) {
d, ctx := newTestDB(t)
if _, err := d.SubmitIPs(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatalf("submit ips (1st time): %v", err)
}
ip1, err := d.GetIPByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get ip: %v", err)
}
if err := d.UpsertCheck(ctx, Check{
IPID: ip1.ID, IPAddress: ip1.IPAddress, AttemptNumber: ip1.AttemptNumber,
Source: SourceEgress, CheckType: "https", Target: "https://example.test",
Success: true, CheckedAt: Now(),
}); err != nil {
t.Fatalf("upsert check (1st time): %v", err)
}
if err := d.DeleteIP(ctx, ip1.ID); err != nil {
t.Fatalf("delete ip: %v", err)
}
if _, err := d.SubmitIPs(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatalf("submit ips (2nd time): %v", err)
}
ip2, err := d.GetIPByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get ip (2nd time): %v", err)
}
if ip2.ID == ip1.ID {
t.Fatalf("expected a fresh ip_queue row, got the same id %d", ip1.ID)
}
if ip2.CycleID == ip1.CycleID {
t.Fatalf("expected a fresh cycle_id, got the same value %d twice", ip1.CycleID)
}
// Same check_type/target/source as the first cycle — would collide on
// the old UNIQUE(ip_id, attempt_number, ...) key structure.
if err := d.UpsertCheck(ctx, Check{
IPID: ip2.ID, IPAddress: ip2.IPAddress, AttemptNumber: ip2.AttemptNumber,
Source: SourceEgress, CheckType: "https", Target: "https://example.test",
Success: false, CheckedAt: Now(),
}); err != nil {
t.Fatalf("upsert check (2nd time): %v", err)
}
summary, err := d.GetRegistryByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get registry: %v", err)
}
if summary.TotalCycles != 2 {
t.Fatalf("expected 2 distinct cycles retained, got %d", summary.TotalCycles)
}
checks, err := d.ListChecksForRegistry(ctx, summary.ID, nil)
if err != nil {
t.Fatalf("list checks for registry: %v", err)
}
if len(checks) != 2 {
t.Fatalf("expected 2 separate check rows (one per cycle), got %+v", checks)
}
}
// TestPruneRegistryHistoryKeepsOnlyNewestCycles proves the retention-depth
// setting actually deletes older cycles' checks once an address has
// accumulated more cycles than the configured depth.
func TestPruneRegistryHistoryKeepsOnlyNewestCycles(t *testing.T) {
d, ctx := newTestDB(t)
var registryID int64
for i := 0; i < 3; i++ {
if _, err := d.SubmitIPs(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatalf("submit ips (cycle %d): %v", i, err)
}
ip, err := d.GetIPByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get ip (cycle %d): %v", i, err)
}
registryID = ip.RegistryID
if err := d.UpsertCheck(ctx, Check{
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
Source: SourceEgress, CheckType: "https", Target: "https://example.test",
Success: true, CheckedAt: Now(),
}); err != nil {
t.Fatalf("upsert check (cycle %d): %v", i, err)
}
if i < 2 {
if err := d.DeleteIP(ctx, ip.ID); err != nil {
t.Fatalf("delete ip (cycle %d): %v", i, err)
}
}
}
summary, err := d.GetRegistryByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get registry: %v", err)
}
if summary.TotalCycles != 3 {
t.Fatalf("expected 3 cycles before pruning, got %d", summary.TotalCycles)
}
if err := d.PruneRegistryHistory(ctx, registryID, 1); err != nil {
t.Fatalf("prune registry history: %v", err)
}
summary, err = d.GetRegistryByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get registry after prune: %v", err)
}
if summary.TotalCycles != 1 {
t.Fatalf("expected 1 cycle after pruning to depth 1, got %d", summary.TotalCycles)
}
// Pruning must never touch next_cycle — otherwise a future cycle_id
// could collide with one that was just pruned away.
var nextCycle int
if err := d.QueryRowContext(ctx, `SELECT next_cycle FROM ip_registry WHERE id=?`, registryID).Scan(&nextCycle); err != nil {
t.Fatalf("read next_cycle: %v", err)
}
if nextCycle != 4 {
t.Fatalf("expected next_cycle unaffected by pruning (still 4), got %d", nextCycle)
}
}
// TestPruneRegistryHistoryNoopWhenUnderDepth proves pruning is a no-op when
// an address has fewer cycles than the configured retention depth.
func TestPruneRegistryHistoryNoopWhenUnderDepth(t *testing.T) {
d, ctx := newTestDB(t)
if _, err := d.SubmitIPs(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatalf("submit ips: %v", err)
}
ip, err := d.GetIPByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get ip: %v", err)
}
if err := d.UpsertCheck(ctx, Check{
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
Source: SourceEgress, CheckType: "https", Target: "https://example.test",
Success: true, CheckedAt: Now(),
}); err != nil {
t.Fatalf("upsert check: %v", err)
}
if err := d.PruneRegistryHistory(ctx, ip.RegistryID, 10); err != nil {
t.Fatalf("prune registry history: %v", err)
}
summary, err := d.GetRegistryByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatalf("get registry: %v", err)
}
if summary.TotalCycles != 1 {
t.Fatalf("expected the single cycle to survive a no-op prune, got %d", summary.TotalCycles)
}
}
func TestSetHistoryRetentionCyclesRejectsNegative(t *testing.T) {
d, ctx := newTestDB(t)
if err := d.BootstrapFromConfig(ctx, &config.ControlAPI{}); err != nil {
t.Fatalf("bootstrap: %v", err)
}
if err := d.SetHistoryRetentionCycles(ctx, -1); !errors.Is(err, ErrValidation) {
t.Fatalf("expected ErrValidation, got %v", err)
}
}
func TestGetSetHistoryRetentionCyclesRoundTrip(t *testing.T) {
d, ctx := newTestDB(t)
if err := d.BootstrapFromConfig(ctx, &config.ControlAPI{}); err != nil {
t.Fatalf("bootstrap: %v", err)
}
if err := d.SetHistoryRetentionCycles(ctx, 5); err != nil {
t.Fatalf("set: %v", err)
}
s, err := d.GetSettings(ctx)
if err != nil {
t.Fatalf("get settings: %v", err)
}
if s.HistoryRetentionCycles != 5 {
t.Fatalf("expected 5, got %d", s.HistoryRetentionCycles)
}
}
+18 -2
View File
@@ -13,8 +13,8 @@ func (d *DB) GetSettings(ctx context.Context) (Settings, error) {
var s Settings
var createdAt, updatedAt string
err := d.QueryRowContext(ctx, `
SELECT fip_settle_seconds, created_at, updated_at FROM settings WHERE id=1
`).Scan(&s.FIPSettleSeconds, &createdAt, &updatedAt)
SELECT fip_settle_seconds, history_retention_cycles, created_at, updated_at FROM settings WHERE id=1
`).Scan(&s.FIPSettleSeconds, &s.HistoryRetentionCycles, &createdAt, &updatedAt)
if err != nil {
return Settings{}, err
}
@@ -42,3 +42,19 @@ func (d *DB) SetFIPSettleSeconds(ctx context.Context, seconds int) error {
`, seconds, now)
return err
}
// SetHistoryRetentionCycles persists how many recent check cycles to keep
// per registry address; 0 means unlimited (the default, matching behavior
// before this setting existed). Does not retroactively prune anything
// itself — pruning happens per-address as each check cycle finishes (see
// orchestrator.aggregateAndRelease).
func (d *DB) SetHistoryRetentionCycles(ctx context.Context, cycles int) error {
if cycles < 0 {
return fmt.Errorf("history_retention_cycles must be >= 0: %w", ErrValidation)
}
now := timeToDB(Now())
_, err := d.ExecContext(ctx, `
UPDATE settings SET history_retention_cycles=?, updated_at=? WHERE id=1
`, cycles, now)
return err
}