package db import ( "context" "database/sql" "fmt" "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 { _, 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) if 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 } func (d *DB) SetFIPAssociated(ctx context.Context, ipID int64, fipID string, leaseTTL time.Duration) error { now := Now() _, err := d.ExecContext(ctx, ` UPDATE ip_queue SET state=?, fip_id=?, fip_associated_at=?, lease_expires_at=?, updated_at=? WHERE id=? `, IPAwaitingSelfCheck, fipID, timeToDB(now), timeToDB(now.Add(leaseTTL)), timeToDB(now), ipID) return err } 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() } // 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 { _, 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, overall_result='', assigned_at=NULL, fip_associated_at=NULL, updated_at=? WHERE id=? `, nextState, retryCount, 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: 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: 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 { 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: 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, 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 { 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) 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, delete dependent checks/events, then the ip_queue row itself // — the same FK-clearing order DeleteValidator uses for owner_validator_id. 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, `DELETE FROM checks WHERE ip_id=?`, ipID); err != nil { return fmt.Errorf("delete 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, `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 } 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, 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, 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 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, ); err != nil { return nil, err } if owner.Valid { item.OwnerValidatorID = &owner.String } 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 }