diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index 2293b87..7582eb7 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -118e09f98944ac59593af603c04a1d5f0e7cafb6549b5a562300bc31d9c769da control-api -48c9b99fa88be751d9badfba8b7d80f894e326a743a90e85478f1d2251ce4ab6 validator-agent -c7b4c2464a5974237b1f98f1565884f09b1ee91a1efadb4b8e1e7d9aaf0b9fe2 prober -b0339b9b341d9ab2d4f6835fda5d9f927a18ed8a831ec9f24dba05624664a417 admin-dashboard +8f2ad37b7819131a9b25eee9d35a57880c7aa233bbb741f4fd8dc0af33f6241c control-api +3288930ce4e09e6793048103e25991f025f71dbedced39dee9bccf821eb72a3d validator-agent +586cb533631dfcef2246cdb69e20777f882506755d15afe475399356c1ed2d02 prober +4729c33a2bce584cba29a8833b0701c0a1cd98bb17bf331aabffc27c4cde940c admin-dashboard diff --git a/bin/admin-dashboard b/bin/admin-dashboard index 25938cb..4510b8d 100755 Binary files a/bin/admin-dashboard and b/bin/admin-dashboard differ diff --git a/bin/control-api b/bin/control-api index eb03d87..81270de 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/bin/prober b/bin/prober index 022012a..03b4eba 100755 Binary files a/bin/prober and b/bin/prober differ diff --git a/bin/validator-agent b/bin/validator-agent index 5a1bb75..ce183ed 100755 Binary files a/bin/validator-agent and b/bin/validator-agent differ diff --git a/docs/API.md b/docs/API.md index a51fe47..9adc8f3 100644 --- a/docs/API.md +++ b/docs/API.md @@ -598,6 +598,21 @@ queued ──(control-api сам, без вызова API)──▶ assigning_fi `orchestrator.poll_interval_seconds`), явного HTTP-метода для их запуска нет — это фоновый цикл (`Tick`), а не запрос/ответ. +Из `assigning_fip` есть и второй, терминальный исход: если на момент +попытки ассоциации Floating IP уже привязан к чужому порту (облако живое — +список адресов мог разойтись с реальностью с момента постановки в +очередь, либо адрес был ошибочно передан занятым), control-api переводит +адрес в состояние `occupied` вместо продолжения в `awaiting_self_check` — +цикл проверки для этой попытки не запускается вовсе. Это отдельное +терминальное состояние, а не `failed`: `failed` означает «проверка +стартовала и не прошла», `occupied` — «проверка не стартовала, потому что +адрес занят кем-то другим». В аудит-логе адреса (`events`) фиксируется +строка `fip_occupied`. `overall_result` для этого состояния остаётся +пустым. Как и `done`/`failed`, `occupied` сбрасывается обратно в `queued` +повторной постановкой через `POST /api/v1/admin/ips` — этим способом +оператор возвращает адрес в работу, убедившись, что конфликт в облаке +разрешился. + Вход в `awaiting_self_check` не означает мгновенную видимость агенту: если настроена пауза (`fip_settle_seconds`, см. [«Настройки оркестратора»](#настройки-оркестратора-apiv1adminconfigorchestrator) @@ -616,8 +631,8 @@ queued ──(control-api сам, без вызова API)──▶ assigning_fi `/api/v1/admin/ips*`, а не самим оркестратором: - **любое нетерминальное состояние → `failed` (`overall_result: "cancelled"`)** — `POST /api/v1/admin/ips/{ip}/cancel`; -- **`done`/`failed` → `queued` (новая попытка)** — `POST - /api/v1/admin/ips` с уже завершённым адресом в списке; +- **`done`/`failed`/`occupied` → `queued` (новая попытка)** — `POST + /api/v1/admin/ips` с уже завершённым (или занятым) адресом в списке; - **любое состояние → адрес физически исчезает из очереди**, вместе со всей историей — `DELETE /api/v1/admin/ips/{ip}`, `POST /api/v1/admin/ips/delete`, `POST /api/v1/admin/ips/clear` (см. diff --git a/docs/DASHBOARD.md b/docs/DASHBOARD.md index 01d499a..8dc5e71 100644 --- a/docs/DASHBOARD.md +++ b/docs/DASHBOARD.md @@ -49,7 +49,7 @@ admin-dashboard -config /etc/cloud-ip-validator/admin-dashboard.yaml | Страница | Назначение | |---|---| | `/overview` | Сводная статистика: счётчики по состояниям, «текущая проверка» (live-снимок всех IP не в терминальном состоянии) и «последние N завершённых» (по умолчанию 20, `overview.last_completed_count`) с разбивкой pass/partial/fail/cancelled. Обновляется каждые `overview.poll_interval_seconds` секунд без перезагрузки страницы. | -| `/ips` | Полная очередь. Форма сверху принимает список адресов (по одному на строке или через запятую) и отправляет их в `POST /api/v1/admin/ips` — **один и тот же вызов** добавляет новые адреса и принудительно перезапускает уже завершённые (см. ниже). У каждого адреса — кнопка «Перепроверить» (для `done`/`failed`) или «Отменить» (для активных состояний), и всегда — «Удалить» (безвозвратно, в отличие от «Отменить», см. ниже). Чекбоксы у строк + кнопка «Удалить выбранные» удаляют список одним вызовом; «Очистить всё» удаляет вообще всё, включая активные проверки — обе операции требуют явного подтверждения. Пока не истекла настроенная на `/settings` пауза (`fip_settle_seconds`), только что привязавший Floating IP адрес показывает отдельный бейдж «прогрев FIP» вместо обычного статуса. | +| `/ips` | Полная очередь. Форма сверху принимает список адресов (по одному на строке или через запятую) и отправляет их в `POST /api/v1/admin/ips` — **один и тот же вызов** добавляет новые адреса и принудительно перезапускает уже завершённые (см. ниже). У каждого адреса — кнопка «Перепроверить» (для `done`/`failed`) или «Отменить» (для активных состояний), и всегда — «Удалить» (безвозвратно, в отличие от «Отменить», см. ниже). Чекбоксы у строк + кнопка «Удалить выбранные» удаляют список одним вызовом; «Очистить всё» удаляет вообще всё, включая активные проверки — обе операции требуют явного подтверждения. Пока не истекла настроенная на `/settings` пауза (`fip_settle_seconds`), только что привязавший Floating IP адрес показывает отдельный бейдж «прогрев FIP» вместо обычного статуса. Если на момент попытки привязки Floating IP оказался уже занят другим портом (дрейф состояния облака или ошибочно переданный адрес), цикл проверки для него не запускается — адрес показывает отдельный бейдж «занят» (отличный от «fail») и строку `fip_occupied` в списке событий на его странице; кнопка «Перепроверить» ставит его в очередь заново. | | `/ips/{ip}` | Детали одного адреса: все проверки текущей попытки и вся история событий. | | `/validators` | Список валидаторов + создание/изменение `os_port_id`/удаление. | | `/sites` | Площадки — число слотов не ограничено, форма сверху добавляет новый слот, назначить/сменить/освободить `site_id` в каждой строке; колонка «Статус» показывает бейдж подключения пробера (`unregistered`/`idle`/`unreachable`, по аналогии с `/validators`), см. [USAGE.md](USAGE.md#состояния-площадки). | diff --git a/docs/USAGE.md b/docs/USAGE.md index f0ecf6a..82e0ebc 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -117,7 +117,7 @@ curl -s http://:8080/api/v1/admin/ips \ |---|---| | `IPAddress` | Проверяемый адрес | | `Sequence` | Позиция в очереди (порядок постановки — из конфига при первом старте либо из последнего вызова `POST /api/v1/admin/ips`) | -| `State` | Текущий этап: `queued`, `assigning_fip`, `awaiting_self_check`, `checking`, `aggregating`, `done`, `failed` | +| `State` | Текущий этап: `queued`, `assigning_fip`, `awaiting_self_check`, `checking`, `aggregating`, `done`, `failed`, `occupied` (см. ниже) | | `OwnerValidatorID` | Какой валидатор сейчас (или последним) занимался этим адресом | | `FIPID` | Идентификатор Floating IP в OpenStack, к которому привязан адрес (пусто, если ещё/уже не привязан) | | `AttemptNumber` | Номер попытки — растёт при каждом requeue (сбой привязки, сбой self-check, реклейм по таймауту) | @@ -154,6 +154,21 @@ curl -s http://:8080/api/v1/admin/ips \ система по итогам проверок. Отличать от обычного `fail` полезно, чтобы не путать «адрес не прошёл проверку» с «проверку прервали вручную». +Отдельно от `OverallResult` стоит состояние **`State: "occupied"`** — +облако живое, и список адресов, переданный как «свободные» (из конфига +или через `POST /api/v1/admin/ips`), мог с тех пор разойтись с +реальностью, либо адрес мог быть передан на проверку по ошибке уже +занятым. Если при попытке привязки Floating IP control-api видит, что тот +уже привязан к чужому порту, адрес переводится в `occupied` **до начала** +цикла проверки — `OverallResult` при этом остаётся пустым, это не `fail`: +`fail` означает «проверка стартовала и не прошла», `occupied` — «проверка +не стартовала, адрес занят кем-то другим». В `events` по адресу +появляется строка `fip_occupied`. Автоматических повторных попыток нет +(Neutron сам не освобождает адрес) — верните адрес в работу вручную через +`POST /api/v1/admin/ips`, когда убедитесь, что конфликт в облаке +разрешился; в дашборде для таких адресов также показывается кнопка +«Перепроверить» вместо «Отменить». + Отсутствие ответа от источника (площадка не прислала результат до истечения `checking_window_seconds`) засчитывается как провал — это управляется настройкой `aggregation.missing_counts_as_fail` в конфиге diff --git a/internal/dashboard/render.go b/internal/dashboard/render.go index 1de9d5b..39f6b52 100644 --- a/internal/dashboard/render.go +++ b/internal/dashboard/render.go @@ -20,7 +20,10 @@ type Badge struct{ Class, Label string } // elapsed since the floating IP was associated, the row gets a distinct // "прогрев FIP" badge instead of looking identical to a row just waiting // on the agent's next poll. The time math happens here, not in the -// template, same as fmtTime's existing precedent. +// template, same as fmtTime's existing precedent. state=="occupied" gets +// its own pill, deliberately separate from the done/failed result switch +// below — it means the check cycle never ran (the floating IP was already +// bound to another port at claim time), not that a check failed. func ipBadge(state, result string, fipAssociatedAt *time.Time, settleSeconds int) Badge { switch state { case "done", "failed": @@ -36,6 +39,8 @@ func ipBadge(state, result string, fipAssociatedAt *time.Time, settleSeconds int } case "queued": return Badge{"pill-neutral", "queued"} + case "occupied": + return Badge{"pill-occupied", "занят"} case "awaiting_self_check": if settleSeconds > 0 && fipAssociatedAt != nil && time.Now().Before(fipAssociatedAt.Add(time.Duration(settleSeconds)*time.Second)) { return Badge{"pill-warning", "прогрев FIP"} diff --git a/internal/dashboard/static/dashboard.css b/internal/dashboard/static/dashboard.css index bd71402..1b74096 100644 --- a/internal/dashboard/static/dashboard.css +++ b/internal/dashboard/static/dashboard.css @@ -38,6 +38,7 @@ --info: #0E6FA8; --info-soft: #DDEEF8; --neutral: #5B6576; --neutral-soft: #E4E8EE; --cancel: #6D4FC4; --cancel-soft: #ECE7FA; + --occupied: #0F7B72; --occupied-soft: #DCF3F0; --overlay: rgba(10, 14, 20, .32); @@ -82,6 +83,7 @@ --info: #5AC0FF; --info-soft: #10222E; --neutral: #A3ACBE; --neutral-soft: #1B212B; --cancel: #B79CFF; --cancel-soft: #221B3A; + --occupied: #4FD8C9; --occupied-soft: #102B28; --overlay: rgba(0, 0, 0, .55); @@ -114,6 +116,7 @@ --info: #5AC0FF; --info-soft: #10222E; --neutral: #A3ACBE; --neutral-soft: #1B212B; --cancel: #B79CFF; --cancel-soft: #221B3A; + --occupied: #4FD8C9; --occupied-soft: #102B28; --overlay: rgba(0, 0, 0, .55); @@ -411,6 +414,7 @@ td.num { font-family: var(--font-mono); font-variant-numeric: tabular-nums; colo .pill-warning { background: var(--warning-soft); color: var(--warning); } .pill-danger { background: var(--danger-soft); color: var(--danger); } .pill-cancel { background: var(--cancel-soft); color: var(--cancel); } +.pill-occupied { background: var(--occupied-soft); color: var(--occupied); } /* ---------- alerts ---------- */ .alert { diff --git a/internal/dashboard/templates/ips.html b/internal/dashboard/templates/ips.html index 9963769..42bfbd2 100644 --- a/internal/dashboard/templates/ips.html +++ b/internal/dashboard/templates/ips.html @@ -56,7 +56,7 @@ {{range .Items}} {{$b := ipBadge .State .OverallResult .FIPAssociatedAt $.FIPSettleSeconds}} -{{$terminal := or (eq .State "done") (eq .State "failed")}} +{{$terminal := or (eq .State "done") (eq .State "failed") (eq .State "occupied")}} {{.IPAddress}} diff --git a/internal/db/models.go b/internal/db/models.go index dd32acd..da1d922 100644 --- a/internal/db/models.go +++ b/internal/db/models.go @@ -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" diff --git a/internal/db/queries_ipqueue.go b/internal/db/queries_ipqueue.go index 09638a6..97eb2c6 100644 --- a/internal/db/queries_ipqueue.go +++ b/internal/db/queries_ipqueue.go @@ -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 } diff --git a/internal/db/queries_ipqueue_test.go b/internal/db/queries_ipqueue_test.go index 13cae25..6d31017 100644 --- a/internal/db/queries_ipqueue_test.go +++ b/internal/db/queries_ipqueue_test.go @@ -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) diff --git a/internal/openstack/mock.go b/internal/openstack/mock.go index ea458d1..6502514 100644 --- a/internal/openstack/mock.go +++ b/internal/openstack/mock.go @@ -35,6 +35,16 @@ func (m *MockClient) Seed(id, address, projectID string) { m.byIP[address] = id } +// SeedWithPort registers a floating IP as already associated to portID — +// for simulating a FIP that's occupied (e.g. by another VM's port) before a +// test's Tick runs. +func (m *MockClient) SeedWithPort(id, address, projectID, portID string) { + m.mu.Lock() + defer m.mu.Unlock() + m.fips[id] = &FloatingIP{ID: id, Address: address, ProjectID: projectID, PortID: portID} + m.byIP[address] = id +} + func (m *MockClient) GetFloatingIPByAddress(ctx context.Context, address string) (*FloatingIP, error) { m.mu.Lock() defer m.mu.Unlock() diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 9083817..836e2bd 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -106,6 +106,27 @@ func (o *Orchestrator) associateFIP(ctx context.Context, validatorID, osPortID s o.requeueOrFail(ctx, item.ID, validatorID, fmt.Sprintf("lookup floating ip: %v", err)) return err } + // The cloud is live: an address queued as "free" (from bootstrap config + // or an admin POST) may have drifted onto another port by the time we + // actually get here, or an operator may have queued an already-occupied + // address by mistake. This is the one authoritative moment to catch it — + // checked here rather than at enqueue time because enqueue-time state + // could itself be stale by the time the claim happens. fip.PortID != + // osPortID guards against a false positive when the FIP is already + // associated to this same validator's own port (e.g. control-api + // restarted between associating and recording it) — that's a resume, not + // a conflict. + if fip.PortID != "" && fip.PortID != osPortID { + if err := o.DB.MarkFIPOccupied(ctx, item.ID, validatorID); err != nil { + o.Log.Error("mark fip occupied", "ip_id", item.ID, "err", err) + return err + } + o.event(ctx, "control-api", "", &item.ID, "fip_occupied", + fmt.Sprintf(`{"fip_id":%q,"port_id":%q}`, fip.ID, fip.PortID)) + o.Log.Info("fip already occupied by another port, skipping check cycle", + "ip", item.IPAddress, "fip_port_id", fip.PortID) + return nil + } if err := o.OS.AssociateFloatingIP(ctx, fip.ID, osPortID); err != nil { o.requeueOrFail(ctx, item.ID, validatorID, fmt.Sprintf("associate floating ip: %v", err)) return err @@ -608,7 +629,8 @@ func (o *Orchestrator) SweepStaleSiteHeartbeats(ctx context.Context) error { // RecordEvent is the exported entry point httpapi uses to log // agent/prober-reported audit events (config_received, fip_changed, -// error, etc.) through the same path as internally generated events. +// error, etc.) through the same path as internally generated events +// (fip_associated, fip_occupied, retry_or_fail, etc.). func (o *Orchestrator) RecordEvent(ctx context.Context, sourceType, sourceID string, ipID *int64, eventType, payload string) { o.event(ctx, sourceType, sourceID, ipID, eventType, payload) } diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go index 45892f2..585a037 100644 --- a/internal/orchestrator/orchestrator_test.go +++ b/internal/orchestrator/orchestrator_test.go @@ -165,6 +165,100 @@ func TestHappyPath(t *testing.T) { } } +// TestFIPAlreadyOccupiedByAnotherPortSkipsCheckCycle covers the defense +// against a Floating IP that turns out to already be attached to some other +// VM's port at claim time — the cloud is live, so the "free" list supplied +// at bootstrap/via the admin API can drift, or an operator can mistakenly +// queue an already-occupied address. The address must be terminated as +// `occupied` immediately, without ever entering the check cycle, without +// stealing the port, and with the validator freed back to idle so the rest +// of the queue isn't starved behind it. +func TestFIPAlreadyOccupiedByAnotherPortSkipsCheckCycle(t *testing.T) { + ctx := context.Background() + o, d, mock := newTestOrchestrator(t, 180) + + mock.SeedWithPort("fip-1", "1.2.3.4", "svc-project", "someone-elses-port") + if err := d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1"); err != nil { + t.Fatalf("register validator: %v", err) + } + if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil { + t.Fatalf("seed queue: %v", err) + } + + o.Tick(ctx) + + ip, err := d.GetIPByAddress(ctx, "1.2.3.4") + if err != nil { + t.Fatalf("get ip: %v", err) + } + if ip.State != db.IPOccupied { + t.Fatalf("expected occupied, got %s", ip.State) + } + if ip.OverallResult != "" { + t.Fatalf("expected empty overall_result, got %q", ip.OverallResult) + } + if ip.OwnerValidatorID != nil { + t.Fatalf("expected no owning validator, got %v", *ip.OwnerValidatorID) + } + if ip.FIPID != "" { + t.Fatalf("expected no fip_id recorded, got %q", ip.FIPID) + } + + v, err := d.GetValidator(ctx, "validator-1") + if err != nil { + t.Fatalf("get validator: %v", err) + } + if v.State != db.ValidatorIdle || v.CurrentIPID != nil { + t.Fatalf("expected validator freed back to idle, got state=%s current_ip=%v", v.State, v.CurrentIPID) + } + + if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "someone-elses-port" { + t.Fatalf("expected fip's port untouched (not stolen), got %q", fip.PortID) + } + + events, err := d.ListEventsForIP(ctx, ip.ID) + if err != nil { + t.Fatalf("list events: %v", err) + } + found := false + for _, e := range events { + if e.EventType == "fip_occupied" { + found = true + } + } + if !found { + t.Fatalf("expected a fip_occupied event, got %+v", events) + } +} + +// TestFIPAssociatedToOwnValidatorPortIsNotOccupied guards against a false +// positive: a Floating IP already attached to the very validator we're +// about to associate it with (e.g. control-api restarted between the +// OpenStack call and recording it in the DB) is a resume, not a conflict — +// it must proceed through the normal happy path, not be flagged occupied. +func TestFIPAssociatedToOwnValidatorPortIsNotOccupied(t *testing.T) { + ctx := context.Background() + o, d, mock := newTestOrchestrator(t, 180) + + mock.SeedWithPort("fip-1", "1.2.3.4", "svc-project", "port-1") + if err := d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1"); err != nil { + t.Fatalf("register validator: %v", err) + } + if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil { + t.Fatalf("seed queue: %v", err) + } + + o.Tick(ctx) + + ip, err := d.GetIPByAddress(ctx, "1.2.3.4") + if err != nil { + t.Fatalf("get ip: %v", err) + } + if ip.State != db.IPAwaitingSelfCheck { + t.Fatalf("expected awaiting_self_check (not occupied), got %s", ip.State) + } +} + func TestPartialResult(t *testing.T) { ctx := context.Background() o, d, mock := newTestOrchestrator(t, 180)