Skip the check cycle for a Floating IP already occupied by another port
The cloud is live: the address list submitted as "free" (bootstrap config
or POST /api/v1/admin/ips) can drift by the time the orchestrator claims
it, or an operator can queue an already-occupied address by mistake.
Neutron's floating-IP association is a blind "last write wins" PUT with
no conflict error to catch, so associateFIP now checks the FIP's PortID
(already fetched via GetFloatingIPByAddress) before associating, guarded
against the false-positive of the FIP already belonging to this same
validator's own port.
A match routes the address straight to a new terminal ip_queue.state
("occupied", distinct from failed/fail) via db.MarkFIPOccupied — no
retries, since Neutron won't free it on its own and requeuing would let
it be reclaimed again next tick, starving the rest of the queue — plus a
dedicated fip_occupied audit event. Resubmitting the address later (once
the conflict is resolved) resets it to queued via the existing
POST /api/v1/admin/ips resubmit path (CancelIP/ListExpiredLeases updated
to treat occupied as terminal too). admin-dashboard gets its own "занят"
badge, distinct from fail/partial/cancelled.
Rebuilt bin/{control-api,admin-dashboard,prober,validator-agent} and
bin/SHA256SUMS per docs/SETUP.md's documented build recipe, since
control-api and admin-dashboard source changed.
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NeVbMVEiE7XQAkBd7HQgj6
This commit is contained in:
1 parent
2246369b64
commit
93c79b63ea
17 files changed
+326
-17
No files matched your search
@@ -25,6 +25,11 @@ const (
|
||||
IPAggregating = "aggregating"
|
||||
IPDone = "done"
|
||||
IPFailed = "failed"
|
||||
// IPOccupied is a terminal state distinct from IPFailed: the floating IP
|
||||
// was found already associated to a different port at claim time, so the
|
||||
// check cycle never started for this attempt. See
|
||||
// Orchestrator.associateFIP and db.MarkFIPOccupied.
|
||||
IPOccupied = "occupied"
|
||||
|
||||
ResultPass = "pass"
|
||||
ResultPartial = "partial"
|
||||
|
||||
@@ -167,6 +167,40 @@ func (d *DB) ReleaseFIP(ctx context.Context, ipID int64, validatorID string) 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
|
||||
@@ -228,7 +262,7 @@ func (d *DB) RequeueOrFail(ctx context.Context, ipID int64, validatorID string,
|
||||
// transaction, in the order given:
|
||||
//
|
||||
// - unknown address: inserted as a new queued row.
|
||||
// - address currently done/failed: reset to queued (new attempt,
|
||||
// - 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,
|
||||
@@ -277,7 +311,7 @@ func (d *DB) SubmitIPs(ctx context.Context, addresses []string) (SubmitIPsResult
|
||||
case err != nil:
|
||||
return result, err
|
||||
|
||||
case state == IPDone || state == IPFailed:
|
||||
case state == IPDone || state == IPFailed || state == IPOccupied:
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
UPDATE ip_queue SET
|
||||
state=?, sequence=?, owner_validator_id=NULL, fip_id='', retry_count=0,
|
||||
@@ -324,8 +358,8 @@ func (d *DB) CancelIP(ctx context.Context, ipID int64) error {
|
||||
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)
|
||||
WHERE id=? AND state NOT IN (?, ?, ?)
|
||||
`, IPFailed, ResultCancelled, now, now, ipID, IPDone, IPFailed, IPOccupied)
|
||||
if err != nil {
|
||||
return fmt.Errorf("cancel ip: %w", err)
|
||||
}
|
||||
@@ -461,8 +495,8 @@ func (d *DB) ListChecking(ctx context.Context) ([]IPQueueItem, error) {
|
||||
// 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))
|
||||
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
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
@@ -38,6 +39,110 @@ func TestSetFIPAssociatedStampsTimestamp(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestMarkFIPOccupiedIsTerminalAndFreesValidator(t *testing.T) {
|
||||
d, ctx := newTestDB(t)
|
||||
|
||||
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
|
||||
t.Fatalf("seed queue: %v", err)
|
||||
}
|
||||
if err := d.AdminCreateValidator(ctx, "validator-1", "port-1"); err != nil {
|
||||
t.Fatalf("create validator: %v", err)
|
||||
}
|
||||
claimed, err := d.ClaimNextQueued(ctx, "validator-1", time.Minute)
|
||||
if err != nil || claimed == nil {
|
||||
t.Fatalf("claim: item=%+v err=%v", claimed, err)
|
||||
}
|
||||
|
||||
if err := d.MarkFIPOccupied(ctx, claimed.ID, "validator-1"); err != nil {
|
||||
t.Fatalf("mark fip occupied: %v", err)
|
||||
}
|
||||
|
||||
item, err := d.GetIP(ctx, claimed.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("get ip: %v", err)
|
||||
}
|
||||
if item.State != IPOccupied {
|
||||
t.Fatalf("expected occupied, got state=%s", item.State)
|
||||
}
|
||||
if item.OwnerValidatorID != nil {
|
||||
t.Fatalf("expected owner cleared, got %v", *item.OwnerValidatorID)
|
||||
}
|
||||
if item.FIPID != "" {
|
||||
t.Fatalf("expected fip_id cleared, got %q", item.FIPID)
|
||||
}
|
||||
if item.LeaseExpiresAt != nil {
|
||||
t.Fatalf("expected lease cleared, got %v", item.LeaseExpiresAt)
|
||||
}
|
||||
|
||||
v, err := d.GetValidator(ctx, "validator-1")
|
||||
if err != nil {
|
||||
t.Fatalf("get validator: %v", err)
|
||||
}
|
||||
if v.State != ValidatorIdle || v.CurrentIPID != nil {
|
||||
t.Fatalf("expected validator freed to idle, got state=%s current_ip=%v", v.State, v.CurrentIPID)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubmitIPsResubmitsOccupiedAddress(t *testing.T) {
|
||||
d, ctx := newTestDB(t)
|
||||
|
||||
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
|
||||
t.Fatalf("seed queue: %v", err)
|
||||
}
|
||||
if err := d.AdminCreateValidator(ctx, "validator-1", "port-1"); err != nil {
|
||||
t.Fatalf("create validator: %v", err)
|
||||
}
|
||||
claimed, err := d.ClaimNextQueued(ctx, "validator-1", time.Minute)
|
||||
if err != nil || claimed == nil {
|
||||
t.Fatalf("claim: item=%+v err=%v", claimed, err)
|
||||
}
|
||||
if err := d.MarkFIPOccupied(ctx, claimed.ID, "validator-1"); err != nil {
|
||||
t.Fatalf("mark fip occupied: %v", err)
|
||||
}
|
||||
|
||||
result, err := d.SubmitIPs(ctx, []string{"1.2.3.4"})
|
||||
if err != nil {
|
||||
t.Fatalf("submit ips: %v", err)
|
||||
}
|
||||
if len(result.Requeued) != 1 || result.Requeued[0] != "1.2.3.4" {
|
||||
t.Fatalf("expected address requeued, got %+v", result)
|
||||
}
|
||||
if len(result.SkippedInProgress) != 0 {
|
||||
t.Fatalf("expected nothing skipped as in-progress, got %+v", result.SkippedInProgress)
|
||||
}
|
||||
|
||||
item, err := d.GetIP(ctx, claimed.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("get ip: %v", err)
|
||||
}
|
||||
if item.State != IPQueued {
|
||||
t.Fatalf("expected queued after resubmit, got state=%s", item.State)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCancelIPRejectsAlreadyOccupied(t *testing.T) {
|
||||
d, ctx := newTestDB(t)
|
||||
|
||||
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
|
||||
t.Fatalf("seed queue: %v", err)
|
||||
}
|
||||
if err := d.AdminCreateValidator(ctx, "validator-1", "port-1"); err != nil {
|
||||
t.Fatalf("create validator: %v", err)
|
||||
}
|
||||
claimed, err := d.ClaimNextQueued(ctx, "validator-1", time.Minute)
|
||||
if err != nil || claimed == nil {
|
||||
t.Fatalf("claim: item=%+v err=%v", claimed, err)
|
||||
}
|
||||
if err := d.MarkFIPOccupied(ctx, claimed.ID, "validator-1"); err != nil {
|
||||
t.Fatalf("mark fip occupied: %v", err)
|
||||
}
|
||||
|
||||
err = d.CancelIP(ctx, claimed.ID)
|
||||
if !errors.Is(err, ErrInvalidState) {
|
||||
t.Fatalf("expected ErrInvalidState cancelling an occupied ip, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRequeueClearsFIPAssociatedAt(t *testing.T) {
|
||||
d, ctx := newTestDB(t)
|
||||
|
||||
|
||||
Reference in new issue
Block a user