Files
cloud-ip-validator/internal/db/queries_ipqueue.go
T
ayurishchevandClaude Sonnet 5.5 e95b5eb7d5 Retry a failed self-check on another validator; add the self-check failure ceiling
A validator that failed the self-check of an address no longer gets that address
again in the current round (ClaimNextQueued skips it); the validator itself stays
in service and takes all other addresses. The verdict fail is set when the number
of failed self-checks of an address reaches settings.self_check_max_attempts
(1..50, default 5, independent of the number of validators); max_retries and
retry_count are no longer used for self-check. If every working validator has
already failed the address, a new round starts and the exclusions lapse.

Migration 0012: ip_self_check_failures (permanent history per registry address),
ip_queue.sc_failures and sc_round_start_cycle (cycle_id is used instead of
attempt_number, which restarts when a queue row is recreated), the setting.
db.FailSelfCheck does it in one transaction; re-submission starts a new series.
API: self_check_max_attempts in GET/PUT /admin/config/orchestrator,
self_check_failed_on in /admin/ips/{ip} and /admin/registry/{ip}. Dashboard: the
field on /settings and the line "Self-check не прошёл на: ..." on the address
pages. Docs, plan and summary in docs/changes/.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2026-10-04 09:50:12 +03:00

1197 lines
40 KiB
Go

package db
import (
"context"
"database/sql"
"fmt"
"strings"
"time"
)
// SeedQueue inserts the configured IP address list in order, assigning each
// a stable sequence number. Re-running with the same list is a no-op for
// addresses already present (ON CONFLICT DO NOTHING keyed by the UNIQUE
// ip_address column), so restarting control-api against the same config
// never re-queues already-processed addresses.
func (d *DB) SeedQueue(ctx context.Context, addresses []string) error {
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
now := timeToDB(Now())
runID := int64(0)
for i, addr := range addresses {
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 runID == 0 {
if runID, err = openRunTx(ctx, tx, RunManual, now); err != nil {
return err
}
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO ip_queue (ip_address, sequence, state, registry_id, cycle_id, run_id, sc_round_start_cycle, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(ip_address) DO NOTHING
`, addr, i, IPQueued, registryID, cycle, runID, cycle, now, now); err != nil {
return fmt.Errorf("seed %s: %w", addr, err)
}
}
return tx.Commit()
}
// ClaimNextQueued atomically hands the next queued IP (lowest sequence) to
// the given idle validator, skipping addresses this validator failed the
// self-check of in the current round (see FailSelfCheck): such an address
// stays queued for the other validators and does not block the ones behind
// it. It returns (nil, nil) if the validator isn't idle or no IP is
// claimable for it. The DB connection pool is capped at one physical
// connection (see Open), so this transaction already has exclusive access
// to the database for its duration — no other claim, requeue, or update can
// interleave — which combined with the conditional UPDATEs (checked via
// RowsAffected) guarantees a single IP is never claimed by two validators.
func (d *DB) ClaimNextQueued(ctx context.Context, validatorID string, leaseTTL time.Duration) (*IPQueueItem, error) {
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return nil, err
}
defer tx.Rollback()
var state string
err = tx.QueryRowContext(ctx, `SELECT state FROM validators WHERE validator_id=?`, validatorID).Scan(&state)
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
if state != ValidatorIdle {
return nil, nil
}
var item IPQueueItem
err = tx.QueryRowContext(ctx, `
SELECT q.id, q.ip_address, q.sequence, q.attempt_number, q.retry_count
FROM ip_queue q WHERE q.state=? AND NOT EXISTS (
SELECT 1 FROM ip_self_check_failures f
WHERE f.registry_id=q.registry_id AND f.validator_id=? AND f.cycle_id>=q.sc_round_start_cycle)
ORDER BY q.sequence LIMIT 1
`, IPQueued, validatorID).Scan(&item.ID, &item.IPAddress, &item.Sequence, &item.AttemptNumber, &item.RetryCount)
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
now := Now()
lease := now.Add(leaseTTL)
res, err := tx.ExecContext(ctx, `
UPDATE ip_queue SET state=?, owner_validator_id=?, assigned_at=?, lease_expires_at=?, updated_at=?
WHERE id=? AND state=?
`, IPAssigningFIP, validatorID, timeToDB(now), timeToDB(lease), timeToDB(now), item.ID, IPQueued)
if err != nil {
return nil, err
}
if n, _ := res.RowsAffected(); n != 1 {
return nil, nil
}
res, err = tx.ExecContext(ctx, `
UPDATE validators SET state=?, current_ip_id=?, updated_at=?
WHERE validator_id=? AND state=? AND current_ip_id IS NULL
`, ValidatorAssigned, item.ID, timeToDB(now), validatorID, ValidatorIdle)
if err != nil {
return nil, err
}
if n, _ := res.RowsAffected(); n != 1 {
return nil, nil
}
if err := tx.Commit(); err != nil {
return nil, err
}
item.State = IPAssigningFIP
ownerID := validatorID
item.OwnerValidatorID = &ownerID
item.AssignedAt = &now
item.LeaseExpiresAt = &lease
return &item, nil
}
// SetFIPAssociated records that the floating IP is now attached. It only
// applies to a row still in assigning_fip: if the address was deleted or
// cancelled while the cloud call was in flight, it returns ErrInvalidState
// and the caller must detach the floating IP again.
func (d *DB) SetFIPAssociated(ctx context.Context, ipID int64, fipID string, leaseTTL time.Duration) error {
now := Now()
res, err := d.ExecContext(ctx, `
UPDATE ip_queue SET state=?, fip_id=?, fip_associated_at=?, lease_expires_at=?, updated_at=?
WHERE id=? AND state=?
`, IPAwaitingSelfCheck, fipID, timeToDB(now), timeToDB(now.Add(leaseTTL)), timeToDB(now), ipID, IPAssigningFIP)
if err != nil {
return err
}
if n, _ := res.RowsAffected(); n == 0 {
return fmt.Errorf("ip_id %d is no longer assigning_fip: %w", ipID, ErrInvalidState)
}
return nil
}
// KnownAddresses returns the subset of addresses present in ip_registry,
// i.e. addresses this system has ever queued.
func (d *DB) KnownAddresses(ctx context.Context, addresses []string) (map[string]bool, error) {
const chunk = 500
known := make(map[string]bool, len(addresses))
for start := 0; start < len(addresses); start += chunk {
end := start + chunk
if end > len(addresses) {
end = len(addresses)
}
part := addresses[start:end]
args := make([]any, len(part))
for i, a := range part {
args[i] = a
}
rows, err := d.QueryContext(ctx, `SELECT ip_address FROM ip_registry WHERE ip_address IN (`+placeholders(len(part))+`)`, args...)
if err != nil {
return nil, err
}
for rows.Next() {
var a string
if err := rows.Scan(&a); err != nil {
rows.Close()
return nil, err
}
known[a] = true
}
if err := rows.Err(); err != nil {
rows.Close()
return nil, err
}
rows.Close()
}
return known, nil
}
func (d *DB) SetChecking(ctx context.Context, ipID int64, leaseTTL time.Duration) error {
now := Now()
_, err := d.ExecContext(ctx, `
UPDATE ip_queue SET state=?, lease_expires_at=?, checking_started_at=?, updated_at=?
WHERE id=?
`, IPChecking, timeToDB(now.Add(leaseTTL)), timeToDB(now), timeToDB(now), ipID)
return err
}
func (d *DB) SetAggregating(ctx context.Context, ipID int64) error {
_, err := d.ExecContext(ctx, `UPDATE ip_queue SET state=?, updated_at=? WHERE id=?`,
IPAggregating, timeToDB(Now()), ipID)
return err
}
// FinishIP records the aggregated result and marks the IP done or failed.
func (d *DB) FinishIP(ctx context.Context, ipID int64, result string) error {
return d.FinishIPExpected(ctx, ipID, result, -1)
}
// FinishIPExpected is FinishIP that also records, in the address's run, how
// many checks were expected at the verdict (expected < 0: unknown) and
// finalizes the run if this was its last open address.
func (d *DB) FinishIPExpected(ctx context.Context, ipID int64, result string, expected int) error {
state := IPDone
if result == ResultFail {
state = IPFailed
}
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
now := timeToDB(Now())
if _, err := tx.ExecContext(ctx, `
UPDATE ip_queue SET state=?, overall_result=?, aggregated_at=?, updated_at=?
WHERE id=?
`, state, result, now, now, ipID); err != nil {
return err
}
if err := upsertRunResultTx(ctx, tx, ipID, result, expected, now); err != nil {
return err
}
if err := finalizeRunsTx(ctx, tx, now); err != nil {
return err
}
return tx.Commit()
}
// ReleaseFIP records that the floating IP has been disassociated and frees
// the owning validator (if this address is still its current one, see
// freeValidatorSQL), in one transaction.
func (d *DB) ReleaseFIP(ctx context.Context, ipID int64, validatorID string) error {
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
now := timeToDB(Now())
if _, err := tx.ExecContext(ctx, `UPDATE ip_queue SET fip_released_at=?, updated_at=? WHERE id=?`, now, now, ipID); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, freeValidatorSQL, now, validatorID, ipID); err != nil {
return err
}
return tx.Commit()
}
// MarkFIPOccupied terminates the current attempt immediately (no retry)
// because the floating IP was found already associated to a different port
// at claim time — the check cycle never starts for this attempt. Unlike
// RequeueOrFail, there is no retry branch: the cloud won't free the address
// on its own, and leaving it in `queued` would let it be reclaimed again
// next tick, starving the rest of the queue behind it. The owning validator
// (if any) is freed in the same transaction, same as RequeueOrFail/CancelIP.
func (d *DB) MarkFIPOccupied(ctx context.Context, ipID int64, validatorID string) error {
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
now := timeToDB(Now())
if _, err := tx.ExecContext(ctx, `
UPDATE ip_queue SET
state=?, owner_validator_id=NULL, fip_id='',
lease_expires_at=NULL, aggregated_at=?, updated_at=?
WHERE id=?
`, IPOccupied, now, now, ipID); err != nil {
return err
}
if validatorID != "" {
if _, err := tx.ExecContext(ctx, freeValidatorSQL, now, validatorID, ipID); err != nil {
return err
}
}
if err := finalizeRunsTx(ctx, tx, now); err != nil {
return err
}
return tx.Commit()
}
// RequeueOrFail is used by both the retry path (association/self-check
// failure) and the lease-sweep reclaim path. It clears ownership and
// per-attempt progress, bumps attempt_number and retry_count, and either
// sends the IP back to the queue or marks it permanently failed once
// maxRetries is exceeded. The owning validator (if any) is freed in the
// same transaction.
func (d *DB) RequeueOrFail(ctx context.Context, ipID int64, validatorID string, maxRetries int) error {
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
var retryCount int
if err := tx.QueryRowContext(ctx, `SELECT retry_count FROM ip_queue WHERE id=?`, ipID).Scan(&retryCount); err != nil {
return err
}
retryCount++
now := timeToDB(Now())
nextState := IPQueued
if retryCount > maxRetries {
nextState = IPFailed
}
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,
cycle_id=?, lease_expires_at=NULL, egress_complete=0,
overall_result='', assigned_at=NULL, checking_started_at=NULL, fip_associated_at=NULL, updated_at=?
WHERE id=?
`, nextState, retryCount, cycle, now, ipID)
} else {
_, err = tx.ExecContext(ctx, `
UPDATE ip_queue SET
state=?, retry_count=?, overall_result=?, aggregated_at=?, updated_at=?
WHERE id=?
`, nextState, retryCount, ResultFail, now, now, ipID)
if err == nil {
err = upsertRunResultTx(ctx, tx, ipID, ResultFail, -1, now)
}
if err == nil {
err = finalizeRunsTx(ctx, tx, now)
}
}
if err != nil {
return err
}
if validatorID != "" {
if _, err := tx.ExecContext(ctx, freeValidatorSQL, now, validatorID, ipID); err != nil {
return err
}
}
return tx.Commit()
}
// SelfCheckFailure is the outcome of FailSelfCheck.
type SelfCheckFailure struct {
// Failures is the number of failed self-checks in the address's current
// series, this one included.
Failures int
// Failed: the ceiling was reached and the address got the verdict fail.
Failed bool
// Validators lists the validators that failed the self-check in this
// series, oldest first (with repeats if one failed it more than once).
Validators []string
// NewRound: the address was queued again, and every working validator had
// already failed it in the round, so the round was reset — the exclusions
// no longer apply (always false when Failed).
NewRound bool
}
// FailSelfCheck handles a failed self-check of the address ipID held by
// validatorID, in one transaction: it records the failure in
// ip_self_check_failures and in ip_queue.sc_failures, then
//
// - if sc_failures reached maxAttempts: marks the address failed with the
// verdict fail;
// - otherwise sends it back to the queue without touching retry_count (a
// failed self-check has its own ceiling, unlike a failed association or
// an expired lease, see RequeueOrFail). The validator is excluded from
// this address for the rest of the round (ClaimNextQueued). If no
// working validator (idle, assigned, checking) is left without a failure
// in the round, a new round starts: the exclusions lapse and the retry
// follows the usual rules, so with fewer validators than maxAttempts the
// address never gets stuck in the queue.
//
// The validator is freed in the same transaction. Returns ErrInvalidState if
// the address is not awaiting_self_check on this validator (a late report).
func (d *DB) FailSelfCheck(ctx context.Context, ipID int64, validatorID, detail string, maxAttempts int) (SelfCheckFailure, error) {
var out SelfCheckFailure
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return out, err
}
defer tx.Rollback()
var registryID, cycle, roundStart int64
var runID sql.NullInt64
var attempt, failures int
err = tx.QueryRowContext(ctx, `
SELECT registry_id, run_id, cycle_id, attempt_number, sc_failures, sc_round_start_cycle
FROM ip_queue WHERE id=? AND state=? AND owner_validator_id=?
`, ipID, IPAwaitingSelfCheck, validatorID).Scan(&registryID, &runID, &cycle, &attempt, &failures, &roundStart)
if err == sql.ErrNoRows {
return out, fmt.Errorf("ip_id %d is not awaiting self-check on %s: %w", ipID, validatorID, ErrInvalidState)
}
if err != nil {
return out, err
}
now := timeToDB(Now())
if _, err := tx.ExecContext(ctx, `
INSERT INTO ip_self_check_failures (registry_id, run_id, cycle_id, attempt_number, validator_id, failed_at, detail)
VALUES (?, ?, ?, ?, ?, ?, ?)
`, registryID, runID, cycle, attempt, validatorID, now, detail); err != nil {
return out, fmt.Errorf("record self-check failure: %w", err)
}
failures++
out.Failures = failures
rows, err := tx.QueryContext(ctx, `
SELECT validator_id FROM (
SELECT id, validator_id FROM ip_self_check_failures WHERE registry_id=? ORDER BY id DESC LIMIT ?
) ORDER BY id
`, registryID, failures)
if err != nil {
return out, err
}
for rows.Next() {
var v string
if err := rows.Scan(&v); err != nil {
rows.Close()
return out, err
}
out.Validators = append(out.Validators, v)
}
if err := rows.Err(); err != nil {
rows.Close()
return out, err
}
rows.Close()
if failures >= maxAttempts {
out.Failed = true
if _, err := tx.ExecContext(ctx, `
UPDATE ip_queue SET
state=?, sc_failures=?, overall_result=?, aggregated_at=?, updated_at=?
WHERE id=?
`, IPFailed, failures, ResultFail, now, now, ipID); err != nil {
return out, err
}
if err := upsertRunResultTx(ctx, tx, ipID, ResultFail, -1, now); err != nil {
return out, err
}
if err := finalizeRunsTx(ctx, tx, now); err != nil {
return out, err
}
} else {
newCycle, err := nextRegistryCycleTx(ctx, tx, registryID, now)
if err != nil {
return out, err
}
var left int
if err := tx.QueryRowContext(ctx, `
SELECT COUNT(*) FROM validators v
WHERE v.state IN (?, ?, ?) AND NOT EXISTS (
SELECT 1 FROM ip_self_check_failures f
WHERE f.registry_id=? AND f.validator_id=v.validator_id AND f.cycle_id>=?)
`, ValidatorIdle, ValidatorAssigned, ValidatorChecking, registryID, roundStart).Scan(&left); err != nil {
return out, err
}
if left == 0 {
out.NewRound = true
roundStart = int64(newCycle)
}
if _, err := tx.ExecContext(ctx, `
UPDATE ip_queue SET
state=?, owner_validator_id=NULL, fip_id='', attempt_number=attempt_number+1,
cycle_id=?, sc_failures=?, sc_round_start_cycle=?, lease_expires_at=NULL, egress_complete=0,
overall_result='', assigned_at=NULL, checking_started_at=NULL, fip_associated_at=NULL, updated_at=?
WHERE id=?
`, IPQueued, newCycle, failures, roundStart, now, ipID); err != nil {
return out, err
}
}
if _, err := tx.ExecContext(ctx, freeValidatorSQL, now, validatorID, ipID); err != nil {
return out, err
}
return out, tx.Commit()
}
// GetIPRunID returns the check run the queue row belongs to (0 if none).
func (d *DB) GetIPRunID(ctx context.Context, ipID int64) (int64, error) {
var runID sql.NullInt64
if err := d.QueryRowContext(ctx, `SELECT run_id FROM ip_queue WHERE id=?`, ipID).Scan(&runID); err != nil {
return 0, err
}
return runID.Int64, nil
}
// ListSelfCheckFailedOn returns the distinct validators whose self-check of
// the address failed, in alphabetical order: over the whole history of the
// address, or, if runID > 0, only the failures that happened in that run.
func (d *DB) ListSelfCheckFailedOn(ctx context.Context, registryID, runID int64) ([]string, error) {
q := `SELECT DISTINCT validator_id FROM ip_self_check_failures WHERE registry_id=?`
args := []any{registryID}
if runID > 0 {
q += ` AND run_id=?`
args = append(args, runID)
}
rows, err := d.QueryContext(ctx, q+` ORDER BY validator_id`, args...)
if err != nil {
return nil, err
}
defer rows.Close()
out := []string{}
for rows.Next() {
var v string
if err := rows.Scan(&v); err != nil {
return nil, err
}
out = append(out, v)
}
return out, rows.Err()
}
// SubmitIPs is the single admin entry point for both "add new addresses to
// the queue" and "force a re-check of an already-finished address" — the
// same list can freely mix both. Addresses are processed in one
// transaction, in the order given:
//
// - unknown address: inserted as a new queued row.
// - address currently done/failed/occupied: reset to queued (new attempt,
// retry_count cleared — this is a deliberate admin-triggered restart,
// not a system retry). It also starts a new series of self-check
// failures (sc_failures=0, validators that failed it before are no
// longer excluded); the history in ip_self_check_failures stays.
// - address currently queued (not yet claimed): left in state=queued,
// only its sequence is updated.
// - address currently mid-check (assigning_fip / awaiting_self_check /
// checking / aggregating): left untouched entirely — never start a
// second concurrent check for the same address.
//
// Every touched/inserted address (new, requeued, or merely reordered) gets
// a sequence assigned in list order, continuing after the current max
// sequence, so a batch's relative order is preserved and, critically,
// resubmitting the same list later reproduces the same relative order.
func (d *DB) SubmitIPs(ctx context.Context, addresses []string) (SubmitIPsResult, error) {
return d.SubmitIPsAs(ctx, addresses, RunManual)
}
// SubmitIPsAs is SubmitIPs that names the kind of run it opens when no run is
// open (RunManual or RunAuto); an already open run is joined whatever the kind.
func (d *DB) SubmitIPsAs(ctx context.Context, addresses []string, kind string) (SubmitIPsResult, error) {
var result SubmitIPsResult
if len(addresses) == 0 {
return result, fmt.Errorf("addresses must not be empty: %w", ErrValidation)
}
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return result, err
}
defer tx.Rollback()
var base int
if err := tx.QueryRowContext(ctx, `SELECT COALESCE(MAX(sequence), -1) + 1 FROM ip_queue`).Scan(&base); err != nil {
return result, err
}
now := timeToDB(Now())
runID := int64(0)
ensureRun := func() (int64, error) {
if runID == 0 {
id, err := openRunTx(ctx, tx, kind, now)
if err != nil {
return 0, err
}
runID = id
}
return runID, nil
}
for i, addr := range addresses {
seq := base + i
var state string
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
}
rid, rErr := ensureRun()
if rErr != nil {
return result, rErr
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO ip_queue (ip_address, sequence, state, registry_id, cycle_id, run_id, sc_round_start_cycle, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
`, addr, seq, IPQueued, registryID, cycle, rid, cycle, now, now); err != nil {
return result, fmt.Errorf("insert %s: %w", addr, err)
}
result.Added = append(result.Added, addr)
case err != nil:
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
}
rid, rErr := ensureRun()
if rErr != nil {
return result, rErr
}
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, cycle_id=?, lease_expires_at=NULL, egress_complete=0,
overall_result='', run_id=?, sc_failures=0, sc_round_start_cycle=?,
assigned_at=NULL, checking_started_at=NULL, fip_associated_at=NULL, aggregated_at=NULL, fip_released_at=NULL, updated_at=?
WHERE ip_address=?
`, IPQueued, seq, cycle, rid, cycle, now, addr); err != nil {
return result, fmt.Errorf("requeue %s: %w", addr, err)
}
result.Requeued = append(result.Requeued, addr)
case state == IPQueued:
if _, err := tx.ExecContext(ctx, `
UPDATE ip_queue SET sequence=?, updated_at=? WHERE ip_address=?
`, seq, now, addr); err != nil {
return result, fmt.Errorf("reorder %s: %w", addr, err)
}
result.Reordered = append(result.Reordered, addr)
default:
// Actively being processed (assigning_fip / awaiting_self_check
// / checking / aggregating) — leave it alone, don't duplicate.
result.SkippedInProgress = append(result.SkippedInProgress, addr)
}
}
if err := tx.Commit(); err != nil {
return result, err
}
return result, nil
}
// CancelIP force-stops a non-terminal IP: marks it failed with
// overall_result=cancelled. It does not disassociate the floating IP or
// free the owning validator — that requires the OpenStack client, so it's
// the caller's (orchestrator's) job to do that before/after calling this.
// Returns ErrInvalidState if the IP is already done/failed (including the
// race where aggregation finishes between the caller's read and this call —
// closed by the single-connection transactional UPDATE below).
func (d *DB) CancelIP(ctx context.Context, ipID int64) error {
now := timeToDB(Now())
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
res, err := tx.ExecContext(ctx, `
UPDATE ip_queue SET
state=?, overall_result=?, aggregated_at=?, owner_validator_id=NULL, fip_id='',
lease_expires_at=NULL, updated_at=?
WHERE id=? AND state NOT IN (?, ?, ?)
`, IPFailed, ResultCancelled, now, now, ipID, IPDone, IPFailed, IPOccupied)
if err != nil {
return fmt.Errorf("cancel ip: %w", err)
}
if n, _ := res.RowsAffected(); n == 0 {
return fmt.Errorf("ip_id %d already finished: %w", ipID, ErrInvalidState)
}
if err := upsertRunResultTx(ctx, tx, ipID, ResultCancelled, -1, now); err != nil {
return err
}
if err := finalizeRunsTx(ctx, tx, now); err != nil {
return err
}
return tx.Commit()
}
// DeleteIP permanently removes an ip_queue row, along with its full check
// and event history, in one transaction — frees the owning validator (if
// any) back to idle first, same as ForceCancel does for the DB side.
// Unlike CancelIP, this leaves nothing behind: the row and its history are
// gone, not marked cancelled. Disassociating a currently-attached floating
// IP is the caller's (orchestrator's) job, same division as CancelIP.
// Returns ErrNotFound if the address is unknown.
func (d *DB) DeleteIP(ctx context.Context, ipID int64) error {
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
if err := deleteIPTx(ctx, tx, ipID); err != nil {
return err
}
if err := finalizeRunsTx(ctx, tx, timeToDB(Now())); err != nil {
return err
}
return tx.Commit()
}
// DeleteIPs deletes a specific list of addresses in one transaction,
// tolerating unknown addresses the same way SubmitIPs does: each is
// resolved to an id and either deleted (added to Deleted) or, if unknown,
// added to NotFound rather than aborting the whole call. Also the
// implementation behind "clear queue" — call it with every address
// currently in the queue.
func (d *DB) DeleteIPs(ctx context.Context, addresses []string) (DeleteIPsResult, error) {
var result DeleteIPsResult
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return result, err
}
defer tx.Rollback()
for _, addr := range addresses {
var ipID int64
err := tx.QueryRowContext(ctx, `SELECT id FROM ip_queue WHERE ip_address=?`, addr).Scan(&ipID)
if err == sql.ErrNoRows {
result.NotFound = append(result.NotFound, addr)
continue
}
if err != nil {
return result, err
}
if err := deleteIPTx(ctx, tx, ipID); err != nil {
return result, fmt.Errorf("delete %s: %w", addr, err)
}
result.Deleted = append(result.Deleted, addr)
}
if err := finalizeRunsTx(ctx, tx, timeToDB(Now())); err != nil {
return result, err
}
if err := tx.Commit(); err != nil {
return result, err
}
return result, nil
}
// deleteIPTx is the shared body of DeleteIP/DeleteIPs: free the owning
// 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, `
UPDATE validators SET current_ip_id=NULL,
state = CASE WHEN state=? THEN state ELSE ? END,
updated_at=?
WHERE current_ip_id=?
`, ValidatorUnreachable, ValidatorIdle, now, ipID); err != nil {
return fmt.Errorf("free owning validator: %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, `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)
}
res, err := tx.ExecContext(ctx, `DELETE FROM ip_queue WHERE id=?`, ipID)
if err != nil {
return fmt.Errorf("delete ip_queue row: %w", err)
}
if n, _ := res.RowsAffected(); n == 0 {
return fmt.Errorf("ip_id %d: %w", ipID, ErrNotFound)
}
return nil
}
// ClearAllIPs deletes every ip_queue row in one short, set-based
// transaction (five statements regardless of the queue size) and returns the
// deleted addresses in queue order. It does exactly what deleteIPTx does per
// row: frees every validator that still points at a doomed row, detaches
// checks/events (their registry_id keeps the history), drops the ephemeral
// ip_site_checks progress flags and finally the rows themselves.
// Disassociating attached floating IPs is the caller's (orchestrator's) job.
func (d *DB) ClearAllIPs(ctx context.Context) ([]string, error) {
tx, err := d.BeginTx(ctx, nil)
if err != nil {
return nil, err
}
defer tx.Rollback()
rows, err := tx.QueryContext(ctx, `SELECT ip_address FROM ip_queue ORDER BY sequence`)
if err != nil {
return nil, err
}
deleted := []string{}
for rows.Next() {
var addr string
if err := rows.Scan(&addr); err != nil {
rows.Close()
return nil, err
}
deleted = append(deleted, addr)
}
if err := rows.Err(); err != nil {
rows.Close()
return nil, err
}
rows.Close()
now := timeToDB(Now())
if _, err := tx.ExecContext(ctx, `
UPDATE validators SET current_ip_id=NULL,
state = CASE WHEN state=? THEN state ELSE ? END,
updated_at=?
WHERE current_ip_id IS NOT NULL
`, ValidatorUnreachable, ValidatorIdle, now); err != nil {
return nil, fmt.Errorf("free owning validators: %w", err)
}
if _, err := tx.ExecContext(ctx, `UPDATE checks SET ip_id=NULL WHERE ip_id IS NOT NULL`); err != nil {
return nil, fmt.Errorf("detach checks: %w", err)
}
if _, err := tx.ExecContext(ctx, `UPDATE events SET ip_id=NULL WHERE ip_id IS NOT NULL`); err != nil {
return nil, fmt.Errorf("detach events: %w", err)
}
if _, err := tx.ExecContext(ctx, `DELETE FROM ip_site_checks`); err != nil {
return nil, fmt.Errorf("delete ip_site_checks: %w", err)
}
if _, err := tx.ExecContext(ctx, `DELETE FROM ip_queue`); err != nil {
return nil, fmt.Errorf("delete ip_queue rows: %w", err)
}
if err := finalizeRunsTx(ctx, tx, now); err != nil {
return nil, err
}
if err := tx.Commit(); err != nil {
return nil, err
}
return deleted, nil
}
// FIPRef identifies a queue row that currently holds a Neutron floating-IP
// association.
type FIPRef struct {
IPID int64
IPAddress string
FIPID string
}
// ListFIPRefs returns the queue rows that may still hold a floating IP: an
// fip_id is set and the row is not finished. A finished row (done, failed,
// occupied) keeps its fip_id for display, but its floating IP was already
// disassociated before the final state was written, so listing it would only
// make a bulk clear issue thousands of pointless cloud calls. At most one row
// per validator qualifies, so a clear does not need to load the whole queue.
func (d *DB) ListFIPRefs(ctx context.Context) ([]FIPRef, error) {
rows, err := d.QueryContext(ctx, `SELECT id, ip_address, fip_id FROM ip_queue WHERE fip_id<>'' AND state NOT IN (?, ?, ?) ORDER BY id`,
IPDone, IPFailed, IPOccupied)
if err != nil {
return nil, err
}
defer rows.Close()
var out []FIPRef
for rows.Next() {
var r FIPRef
if err := rows.Scan(&r.IPID, &r.IPAddress, &r.FIPID); err != nil {
return nil, err
}
out = append(out, r)
}
return out, rows.Err()
}
// ListFIPRefsByAddresses is ListFIPRefs restricted to the given addresses
// (unknown addresses and rows that hold no floating IP are simply absent), using a handful of IN (...) queries instead of one lookup per
// address.
func (d *DB) ListFIPRefsByAddresses(ctx context.Context, addresses []string) ([]FIPRef, error) {
const chunk = 500
var out []FIPRef
for start := 0; start < len(addresses); start += chunk {
end := start + chunk
if end > len(addresses) {
end = len(addresses)
}
part := addresses[start:end]
args := []any{IPDone, IPFailed, IPOccupied}
for _, a := range part {
args = append(args, a)
}
rows, err := d.QueryContext(ctx,
`SELECT id, ip_address, fip_id FROM ip_queue WHERE fip_id<>'' AND state NOT IN (?, ?, ?) AND ip_address IN (`+placeholders(len(part))+`)`, args...)
if err != nil {
return nil, err
}
for rows.Next() {
var r FIPRef
if err := rows.Scan(&r.IPID, &r.IPAddress, &r.FIPID); err != nil {
rows.Close()
return nil, err
}
out = append(out, r)
}
if err := rows.Err(); err != nil {
rows.Close()
return nil, err
}
rows.Close()
}
return out, nil
}
func placeholders(n int) string {
if n <= 0 {
return ""
}
return strings.TrimSuffix(strings.Repeat("?,", n), ",")
}
// CountIPsByState returns how many ip_queue rows are in each state, plus the
// grand total, via a single GROUP BY (no row loading).
func (d *DB) CountIPsByState(ctx context.Context) (map[string]int, int, error) {
rows, err := d.QueryContext(ctx, `SELECT state, COUNT(*) FROM ip_queue GROUP BY state`)
if err != nil {
return nil, 0, err
}
defer rows.Close()
counts := map[string]int{}
total := 0
for rows.Next() {
var state string
var n int
if err := rows.Scan(&state, &n); err != nil {
return nil, 0, err
}
counts[state] = n
total += n
}
return counts, total, rows.Err()
}
// CountIPsByResult returns how many ip_queue rows carry each non-empty
// overall_result (pass/partial/fail/cancelled), via a single GROUP BY.
func (d *DB) CountIPsByResult(ctx context.Context) (map[string]int, error) {
rows, err := d.QueryContext(ctx, `
SELECT overall_result, COUNT(*) FROM ip_queue WHERE overall_result<>'' GROUP BY overall_result
`)
if err != nil {
return nil, err
}
defer rows.Close()
counts := map[string]int{}
for rows.Next() {
var res string
var n int
if err := rows.Scan(&res, &n); err != nil {
return nil, err
}
counts[res] = n
}
return counts, rows.Err()
}
// AnyNonTerminalIP reports whether any ip_queue row is still unfinished (its
// state is none of done/failed/occupied).
func (d *DB) AnyNonTerminalIP(ctx context.Context) (bool, error) {
var any bool
err := d.QueryRowContext(ctx, `SELECT EXISTS(SELECT 1 FROM ip_queue WHERE state NOT IN (?, ?, ?))`,
IPDone, IPFailed, IPOccupied).Scan(&any)
return any, err
}
// Order values for IPFilter.Order.
const (
IPOrderSequence = "sequence"
IPOrderAggregatedAtDesc = "aggregated_at_desc"
)
// IPFilter narrows ListIPsPage. The zero value matches everything, ordered by
// queue sequence.
type IPFilter struct {
States []string // any of these states (empty = all)
Query string // substring of ip_address
Result string // overall_result equals this (pass|partial|fail|cancelled)
Order string // IPOrderSequence (default) | IPOrderAggregatedAtDesc
}
// ListIPsPage returns one page (limit/offset) of ip_queue rows matching f,
// plus the total number of matching rows. limit <= 0 means no limit.
func (d *DB) ListIPsPage(ctx context.Context, f IPFilter, limit, offset int) ([]IPQueueItem, int, error) {
var conds []string
var args []any
if len(f.States) > 0 {
conds = append(conds, "state IN ("+placeholders(len(f.States))+")")
for _, s := range f.States {
args = append(args, s)
}
}
if f.Query != "" {
conds = append(conds, "instr(ip_address, ?) > 0")
args = append(args, f.Query)
}
if f.Result != "" {
conds = append(conds, "overall_result = ?")
args = append(args, f.Result)
}
where := ""
if len(conds) > 0 {
where = "WHERE " + strings.Join(conds, " AND ") + " "
}
var total int
if err := d.QueryRowContext(ctx, `SELECT COUNT(*) FROM ip_queue `+where, args...).Scan(&total); err != nil {
return nil, 0, err
}
order := "ORDER BY sequence"
if f.Order == IPOrderAggregatedAtDesc {
// Timestamps are stored as RFC3339Nano, which trims trailing zeros
// and so does not sort correctly as text; strftime normalizes them to
// a fixed millisecond layout. NULL aggregated_at sorts last in DESC.
order = "ORDER BY strftime('%Y-%m-%d %H:%M:%f', aggregated_at) DESC, sequence"
}
q := ipQueueSelect + where + order
qargs := append([]any(nil), args...)
if limit > 0 {
q += " LIMIT ? OFFSET ?"
qargs = append(qargs, limit, offset)
}
rows, err := d.QueryContext(ctx, q, qargs...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
items, err := scanIPQueueItems(rows)
if err != nil {
return nil, 0, err
}
if items == nil {
items = []IPQueueItem{}
}
return items, total, nil
}
func (d *DB) SetEgressComplete(ctx context.Context, ipID int64) error {
_, err := d.ExecContext(ctx, `UPDATE ip_queue SET egress_complete=1, updated_at=? WHERE id=?`, timeToDB(Now()), ipID)
return err
}
func (d *DB) GetIP(ctx context.Context, ipID int64) (*IPQueueItem, error) {
row := d.QueryRowContext(ctx, ipQueueSelect+`WHERE id=?`, ipID)
return scanIPQueueItem(row)
}
func (d *DB) GetIPByAddress(ctx context.Context, address string) (*IPQueueItem, error) {
row := d.QueryRowContext(ctx, ipQueueSelect+`WHERE ip_address=?`, address)
return scanIPQueueItem(row)
}
func (d *DB) ListIPs(ctx context.Context) ([]IPQueueItem, error) {
rows, err := d.QueryContext(ctx, ipQueueSelect+`ORDER BY sequence`)
if err != nil {
return nil, err
}
defer rows.Close()
return scanIPQueueItems(rows)
}
// ListChecking returns all IPs currently in the checking state — the set a
// prober should be actively probing.
func (d *DB) ListChecking(ctx context.Context) ([]IPQueueItem, error) {
rows, err := d.QueryContext(ctx, ipQueueSelect+`WHERE state=? ORDER BY sequence`, IPChecking)
if err != nil {
return nil, err
}
defer rows.Close()
return scanIPQueueItems(rows)
}
// ListCheckingForSite returns the addresses in the checking state that the
// given prober site still has to probe: those for which it has not yet
// reported completion in the current attempt. Without this filter a prober
// would re-probe every checking address on every poll until the verdict.
func (d *DB) ListCheckingForSite(ctx context.Context, siteIndex int) ([]IPQueueItem, error) {
rows, err := d.QueryContext(ctx, ipQueueSelect+`
WHERE state=? AND NOT EXISTS (
SELECT 1 FROM ip_site_checks s
WHERE s.ip_id=ip_queue.id AND s.attempt_number=ip_queue.attempt_number
AND s.site_idx=? AND s.complete=1)
ORDER BY sequence`, IPChecking, siteIndex)
if err != nil {
return nil, err
}
defer rows.Close()
return scanIPQueueItems(rows)
}
// ListExpiredLeases returns non-terminal IPs whose lease has expired —
// candidates for the lease sweep (crash recovery + stuck-validator reclaim).
func (d *DB) ListExpiredLeases(ctx context.Context, now time.Time) ([]IPQueueItem, error) {
rows, err := d.QueryContext(ctx, ipQueueSelect+`
WHERE state NOT IN (?, ?, ?) AND lease_expires_at IS NOT NULL AND lease_expires_at < ?
`, IPDone, IPFailed, IPOccupied, timeToDB(now))
if err != nil {
return nil, err
}
defer rows.Close()
return scanIPQueueItems(rows)
}
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, checking_started_at, fip_associated_at, aggregated_at, fip_released_at,
registry_id, cycle_id, created_at, updated_at
FROM ip_queue
`
func scanIPQueueItems(rows *sql.Rows) ([]IPQueueItem, error) {
var out []IPQueueItem
for rows.Next() {
item, err := scanIPQueueItem(rows)
if err != nil {
return nil, err
}
out = append(out, *item)
}
return out, rows.Err()
}
func scanIPQueueItem(row rowScanner) (*IPQueueItem, error) {
var item IPQueueItem
var owner sql.NullString
var leaseExpires, assignedAt, checkingStartedAt, 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, &checkingStartedAt, &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
}
if item.AssignedAt, err = nullStringToTimePtr(assignedAt); err != nil {
return nil, err
}
if item.CheckingStartedAt, err = nullStringToTimePtr(checkingStartedAt); err != nil {
return nil, err
}
if item.FIPAssociatedAt, err = nullStringToTimePtr(fipAssociatedAt); err != nil {
return nil, err
}
if item.AggregatedAt, err = nullStringToTimePtr(aggregatedAt); err != nil {
return nil, err
}
if item.FIPReleasedAt, err = nullStringToTimePtr(fipReleasedAt); 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
}