An address still being associated (assigning_fip) has no fip_id in the database, so "clear queue" did not detach it, and the association then finished after the row was gone, leaving the floating IP on the validator port for good. - After clear/cancel/delete, ask the cloud which floating IPs sit on the affected validator ports (new ListFloatingIPsByPort) and detach those that this system queued (known in ip_registry); foreign ones are left. - SetFIPAssociated applies only to a row still in assigning_fip; if the address was removed meanwhile, associateFIP detaches the floating IP. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
917 lines
29 KiB
Go
917 lines
29 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())
|
|
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 _, 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)
|
|
}
|
|
}
|
|
return tx.Commit()
|
|
}
|
|
|
|
// ClaimNextQueued atomically hands the next queued IP (lowest sequence) to
|
|
// the given idle validator. It returns (nil, nil) if the validator isn't
|
|
// idle or no IP is queued. 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 id, ip_address, sequence, attempt_number, retry_count
|
|
FROM ip_queue WHERE state=? ORDER BY sequence LIMIT 1
|
|
`, IPQueued).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=?
|
|
`, 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=?, updated_at=?
|
|
WHERE id=?
|
|
`, IPChecking, timeToDB(now.Add(leaseTTL)), 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 {
|
|
state := IPDone
|
|
if result == ResultFail {
|
|
state = IPFailed
|
|
}
|
|
now := timeToDB(Now())
|
|
_, err := d.ExecContext(ctx, `
|
|
UPDATE ip_queue SET state=?, overall_result=?, aggregated_at=?, updated_at=?
|
|
WHERE id=?
|
|
`, state, result, now, now, ipID)
|
|
return err
|
|
}
|
|
|
|
// ReleaseFIP records that the floating IP has been disassociated and frees
|
|
// the owning validator back to idle, 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, `
|
|
UPDATE validators SET state=?, current_ip_id=NULL, updated_at=?
|
|
WHERE validator_id=?
|
|
`, ValidatorIdle, now, validatorID); 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, `
|
|
UPDATE validators SET state=?, current_ip_id=NULL, updated_at=?
|
|
WHERE validator_id=?
|
|
`, ValidatorIdle, now, validatorID); 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(®istryID); 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, 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 {
|
|
return err
|
|
}
|
|
|
|
if validatorID != "" {
|
|
if _, err := tx.ExecContext(ctx, `
|
|
UPDATE validators SET state=?, current_ip_id=NULL, updated_at=?
|
|
WHERE validator_id=?
|
|
`, ValidatorIdle, now, validatorID); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return tx.Commit()
|
|
}
|
|
|
|
// 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).
|
|
// - 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) {
|
|
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())
|
|
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
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `
|
|
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)
|
|
|
|
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(®istryID); 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, 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, 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())
|
|
res, err := d.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)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// 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
|
|
}
|
|
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 := 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 state=?, current_ip_id=NULL, updated_at=?
|
|
WHERE current_ip_id=?
|
|
`, 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 state=?, current_ip_id=NULL, updated_at=?
|
|
WHERE current_ip_id IS NOT NULL
|
|
`, 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 := 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 every queue row with an attached floating IP (fip_id
|
|
// set) — typically at most one per validator — so a bulk clear can
|
|
// disassociate them without loading 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<>'' ORDER BY id`)
|
|
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 without an attached 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 := make([]any, len(part))
|
|
for i, a := range part {
|
|
args[i] = a
|
|
}
|
|
rows, err := d.QueryContext(ctx,
|
|
`SELECT id, ip_address, fip_id FROM ip_queue WHERE fip_id<>'' 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)
|
|
}
|
|
|
|
// 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, 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, 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,
|
|
®istryID, &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.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
|
|
}
|