fillRegistrySummary derived LastResult from whichever single checks row happened to have the latest checked_at, not the cycle's actual aggregated result — a cycle with a mix of passing and failing checks (e.g. one egress target timed out while the rest, including the chronologically-last check, succeeded) rendered as a green "pass" badge on /registry, disagreeing with the correct "partial" badge already shown on /ips for the same address. Now prefers ip_queue.overall_result (the orchestrator's own aggregation) when a live queue row has a finished cycle, leaves the badge blank while a cycle is still in progress, and only falls back to classifying the most recent cycle's own checks (pass/fail/partial) once the address has been deleted from the queue and overall_result is no longer available. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
273 lines
9.8 KiB
Go
273 lines
9.8 KiB
Go
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
|
|
}
|
|
|
|
// fillRegistrySummary computes the badge-level summary shown on the
|
|
// registry list/detail pages. LastResult must reflect the *aggregated*
|
|
// outcome of the most recent cycle (pass/partial/fail/cancelled), not the
|
|
// success flag of whichever individual check happens to have the latest
|
|
// checked_at — a cycle with a mix of passing and failing checks (e.g. one
|
|
// egress target timed out while the rest succeeded) is "partial", even
|
|
// though the chronologically-last check to report in might have passed.
|
|
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 lastCheckedAt sql.NullString
|
|
if err := d.QueryRowContext(ctx, `
|
|
SELECT checked_at FROM checks WHERE registry_id=? ORDER BY cycle_id DESC, checked_at DESC LIMIT 1
|
|
`, s.ID).Scan(&lastCheckedAt); err != nil && err != sql.ErrNoRows {
|
|
return err
|
|
}
|
|
t, err := nullStringToTimePtr(lastCheckedAt)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
s.LastCheckedAt = t
|
|
|
|
var state, overallResult sql.NullString
|
|
err = d.QueryRowContext(ctx, `SELECT state, overall_result FROM ip_queue WHERE registry_id=?`, s.ID).
|
|
Scan(&state, &overallResult)
|
|
if err != nil && err != sql.ErrNoRows {
|
|
return err
|
|
}
|
|
if err == nil {
|
|
s.InQueue = true
|
|
s.CurrentState = state.String
|
|
}
|
|
|
|
switch {
|
|
case overallResult.Valid && overallResult.String != "":
|
|
// The address has a live ip_queue row with a finished cycle
|
|
// (done/failed) — overall_result is the orchestrator's own
|
|
// aggregation (internal/orchestrator.aggregateAndRelease /
|
|
// db.CancelIP), the authoritative source of truth. Use it as-is
|
|
// rather than re-deriving it from raw check rows.
|
|
s.LastResult = overallResult.String
|
|
case s.InQueue:
|
|
// A live ip_queue row exists but its current cycle hasn't finished
|
|
// yet (still queued/checking/etc, overall_result not set) — no
|
|
// verdict to show yet; leave LastResult empty rather than guessing
|
|
// from a still-incomplete set of checks.
|
|
default:
|
|
// No live ip_queue row (deleted from the queue) — fall back to
|
|
// classifying the most recent recorded cycle from its own checks,
|
|
// the same pass/fail/partial rule aggregateAndRelease uses (minus
|
|
// "missing" checks, which aren't knowable after the fact).
|
|
result, err := lastCycleResultFromChecks(ctx, d, s.ID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
s.LastResult = result
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// lastCycleResultFromChecks classifies the most recent cycle recorded for
|
|
// registryID directly from its checks rows: pass if every recorded check
|
|
// succeeded, fail if every one failed, partial on a mix. Returns "" if no
|
|
// checks are recorded at all.
|
|
func lastCycleResultFromChecks(ctx context.Context, d *DB, registryID int64) (string, error) {
|
|
var cycle sql.NullInt64
|
|
if err := d.QueryRowContext(ctx, `SELECT MAX(cycle_id) FROM checks WHERE registry_id=?`, registryID).Scan(&cycle); err != nil {
|
|
return "", err
|
|
}
|
|
if !cycle.Valid {
|
|
return "", nil
|
|
}
|
|
var total, passed int
|
|
if err := d.QueryRowContext(ctx, `
|
|
SELECT COUNT(*), COALESCE(SUM(success), 0) FROM checks WHERE registry_id=? AND cycle_id=?
|
|
`, registryID, cycle.Int64).Scan(&total, &passed); err != nil {
|
|
return "", err
|
|
}
|
|
switch {
|
|
case total == 0:
|
|
return "", nil
|
|
case passed == 0:
|
|
return ResultFail, nil
|
|
case passed == total:
|
|
return ResultPass, nil
|
|
default:
|
|
return ResultPartial, 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
|
|
}
|