Compare commits

..
2 Commits
Author SHA1 Message Date
ayurishchevandClaude Sonnet 5.5 abbee9a08a Add self-check via control-api (self_check.methods)
control-api is hosted outside the cloud and validators reach it directly,
so it sees the floating IP as the connection's source address. New open
route GET /api/v1/agents/{id}/observed-ip returns that address (taken only
from the TCP peer; forwarding headers are ignored so a validator cannot
forge it).

The agent gets self_check.methods, a priority-ordered list of ip_echo
(unchanged) and control_api; the default stays [ip_echo]. The self-check
passes when any method confirms the address; the next method is tried on
no answer and on a mismatch. Each method has its own timeout so a hung
first method cannot starve the fallback, and control_api uses a new TCP
connection per call (a connection opened before the floating IP was
attached would keep reporting the old address).

Also: docker agent template/env, example config, docs, plan in
docs/changes, e2e script switch E2E_SELF_CHECK_METHODS, rebuilt
bin/control-api and bin/validator-agent.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2026-10-02 03:24:20 +03:00
ayurishchevandClaude Sonnet 5.5 146259cabb Detach the floating IP when a self-check fails
A failed self-check returned the address to the queue and freed the
validator in the database, but left the floating IP attached to the
validator's port. Every later association on that port then failed with
409 ("fixed IP already has a floating IP"), so one failed self-check
poisoned a validator for good; on 2026-10-01 all 20 validators were
poisoned within 23 minutes after ifconfig.me timeouts.

SelfCheckResult now disassociates the floating IP before requeueing, and
ignores a late failed report for an address the validator no longer
holds (it could belong to another validator by then).

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2026-10-02 03:24:19 +03:00
25 changed files with 860 additions and 35 deletions

No files matched your search

+4 -3
View File
@@ -48,7 +48,7 @@ docker compose up -d --build # весь стенд на одной
**Остальные компоненты**
| Компонент | Ключи |
|---|---|
| `validator-agent` | `validator_id`, `control_api_url`, `control_api_token_env` (`CONTROL_API_AGENT_TOKEN`), `poll_interval_seconds`, `self_check.*` (таймаут, `ip_echo_urls`), `checks.*` (таймауты HTTPS/ICMP, число ICMP-пакетов, `ssh.*`) |
| `validator-agent` | `validator_id`, `control_api_url`, `control_api_token_env` (`CONTROL_API_AGENT_TOKEN`), `poll_interval_seconds`, `self_check.*` (таймаут, `methods` — `ip_echo`/`control_api`, `ip_echo_urls`), `checks.*` (таймауты HTTPS/ICMP, число ICMP-пакетов, `ssh.*`) |
| `prober` | `site_id`, `control_api_url`, `control_api_token_env` (`CONTROL_API_AGENT_TOKEN`), `poll_interval_seconds`, `checks.*` (таймауты TCP/ICMP, число ICMP-пакетов) |
| `admin-dashboard` | `server.listen_addr` (`:8090`), `control_api.base_url`, `control_api.timeout_seconds`, `control_api.token_env` (`ADMIN_DASHBOARD_CONTROL_API_TOKEN`), `auth.username_env` / `password_env` / `session_secret_env` (`ADMIN_DASHBOARD_USERNAME` / `_PASSWORD` / `_SESSION_SECRET`), `auth.session_ttl_minutes` (480), `overview.last_completed_count` (20), `overview.poll_interval_seconds` (5) |
@@ -118,7 +118,7 @@ docs/ документация и планы доработок
## API (`/api/v1`)
| Область | Эндпоинты |
|---|---|
| Валидатор | `POST /agents/register`, `POST /agents/{id}/heartbeat`, `GET /agents/{id}/assignment`, `POST /agents/{id}/self-check\|events\|results\|complete` |
| Валидатор | `POST /agents/register`, `POST /agents/{id}/heartbeat`, `GET /agents/{id}/assignment`, `GET /agents/{id}/observed-ip`, `POST /agents/{id}/self-check\|events\|results\|complete` |
| Пробер | `POST /probers/register`, `POST /probers/{site_id}/heartbeat`, `GET /probers/{site_id}/assignments`, `POST /probers/{site_id}/results` |
| Очередь | `GET /admin/status`, `GET\|POST /admin/ips` (`limit/offset/state/q/result/order` — постранично), `GET /admin/ips/{ip}`, `POST /admin/ips/{ip}/cancel`, `DELETE /admin/ips/{ip}`, `POST /admin/ips/delete\|clear`, `POST\|GET /admin/ips/scan` (фоновый скан: `202`, `dry_run`, `wait`; статус и прогресс) |
| Реестр | `GET /admin/registry` (`limit/offset/q/last_result` — постранично), `GET /admin/registry/{ip}` |
@@ -138,7 +138,7 @@ docs/ документация и планы доработок
|---|---|---|
| admin | `CONTROL_API_ADMIN_TOKEN` | все `/api/v1/admin/*` |
| agent | `CONTROL_API_AGENT_TOKEN` | запись результатов валидатора и пробера: `self-check`, `events`, `results`, `complete` |
| открыто | — | `GET /healthz`, `register`, `heartbeat` и получение задания (`GET assignment` / `assignments`) |
| открыто | — | `GET /healthz`, `register`, `heartbeat`, получение задания (`GET assignment` / `assignments`) и `GET /agents/{id}/observed-ip` |
Токены разные: административный не открывает методы агентов, и наоборот. Валидатор и пробер получают настройку и задание без токена, но не могут отправить результат без токена агентов.
- **Пустой токен — уровень открыт** (обратная совместимость): `control-api` стартует с предупреждением в логе. На реальном стенде задайте оба токена и ограничьте доступ на уровне сети ([docs/SETUP.md](docs/SETUP.md#сетевые-доступы)); токены идут открытым текстом без TLS — публикуйте через reverse-proxy с TLS. Включать токен агентов нужно **после** его раздачи валидаторам и проберам ([порядок](docs/SETUP.md#5-аутентификация-токены-и-пароль-дашборда)).
- **Дашборд закрыт логином и паролем** (один администратор; пароль и ключ сессии — из env). Сессия — подписанная cookie (`HttpOnly`, `SameSite=Strict`, без состояния на сервере), CSRF-защита по `Origin`, 5 неудачных входов за 10 минут с одного IP → `429`. Без заданных логина/пароля дашборд открыт (с предупреждением в логе). Дашборд ходит в API с токеном администратора. Подробности — [docs/DASHBOARD.md](docs/DASHBOARD.md#вход-и-сессия).
@@ -200,6 +200,7 @@ scripts/run-local-e2e.sh # сквозной прог
| Дата | Веха | Документ |
|---|---|---|
| 2026-10-02 | Самопроверка через ручку control-api (`self_check.methods`, способ `control_api` рядом с IP-echo); отвязка Floating IP при провале self-check | [план](docs/changes/2026-10-02_03-06_self-check-control-api-plan.md) · [API](docs/API.md#get-apiv1agentsidobserved-ip) |
| 2026-10-01 | Скан Floating IP при тысячах адресов: фоновый постраничный скан с прогрессом, фаза `scanning` в автоцикле, постраничные `/ips` и `/registry`, «Обзор» на счётчиках | [план](docs/changes/2026-10-01_18-19_fip-scan-at-scale-plan.md) · [ревью и тесты](docs/changes/2026-10-01_18-59_fip-scan-at-scale-review.md) · [USAGE](docs/USAGE.md#сканирование-floating-ip-из-openstack) · [API](docs/API.md#post-apiv1adminipsscan) |
| 2026-10-01 | Аутентификация: токены администратора и агентов для API, логин и пароль для дашборда | [план](docs/changes/2026-10-01_11-12_authentication-plan.md) · [ревью и тесты](docs/changes/2026-10-01_11-31_authentication-review.md) · [API](docs/API.md#аутентификация) |
| 2026-10-01 | Автоматический цикл проверок по сценарию: очистка → скан FIP → проверка → пауза | [USAGE](docs/USAGE.md#автоматический-цикл-проверок) · [API](docs/API.md#автоматический-цикл-проверок) |
+2 -2
View File
@@ -1,4 +1,4 @@
4302c5d936b0f6567e4b3b7af3de0123afd862b3cdda53859fad0c7eec443c33 control-api
091ad94b5706b1778b181542de7a4b251cd16c421144269d0955769603d51e06 validator-agent
6fc8aef1671a6e17c60461431b90e86ababc3f154042dff205a3ad9dd11d83c3 control-api
b0e75dce3a840024af40d9544ab17524fb85431e7944582893fe3f9aae1b3401 validator-agent
3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober
d72234688eb1954dbff420ca1c8e83b83ceab561b18a0336af0f73fe4d9ac8be admin-dashboard
BIN
View File
Binary file not shown.
Binary file not shown.
+13 -1
View File
@@ -13,7 +13,19 @@ poll_interval_seconds: 5
self_check:
timeout_seconds: 10
# Must be a resource genuinely outside the cloud project — OpenStack only
# Способы самопроверки в порядке приоритета (допустимо: ip_echo,
# control_api); по умолчанию [ip_echo]. Самопроверка успешна, если адрес
# подтвердил любой способ: пробуются по порядку, остановка на первом
# успешном, к следующему переходим и при отсутствии ответа, и при
# несовпадении адреса. Таймаут timeout_seconds действует на каждый способ
# отдельно (зависший первый способ не лишает второй времени). control_api спрашивает у control-api, с какого адреса он видит
# это соединение; рекомендуется [control_api, ip_echo], когда control-api
# стоит вне облака и валидатор ходит к нему напрямую (через внешнюю сеть).
# Ограничение: если control-api достижим по внутренней сети облака, он
# увидит частный адрес валидатора и control_api всегда даст несовпадение —
# тогда оставьте только ip_echo (или он сработает вторым в списке).
methods: [control_api, ip_echo]
# Used by the ip_echo method. Must be a resource genuinely outside the cloud project — OpenStack only
# applies floating-IP SNAT to traffic leaving via the external network,
# so anything reachable over the project's internal network (including
# control-api itself, if it's on the same internal network) would report
+1
View File
@@ -100,6 +100,7 @@ services:
VALIDATOR_AGENT_CONTROL_API_URL: "${VALIDATOR_AGENT_CONTROL_API_URL:-http://control-api:8080}"
VALIDATOR_AGENT_POLL_INTERVAL_SECONDS: "${VALIDATOR_AGENT_POLL_INTERVAL_SECONDS:-5}"
VALIDATOR_AGENT_SELF_CHECK_TIMEOUT_SECONDS: "${VALIDATOR_AGENT_SELF_CHECK_TIMEOUT_SECONDS:-10}"
VALIDATOR_AGENT_SELF_CHECK_METHODS: "${VALIDATOR_AGENT_SELF_CHECK_METHODS:-[ip_echo]}"
VALIDATOR_AGENT_HTTPS_TIMEOUT_SECONDS: "${VALIDATOR_AGENT_HTTPS_TIMEOUT_SECONDS:-10}"
VALIDATOR_AGENT_ICMP_TIMEOUT_SECONDS: "${VALIDATOR_AGENT_ICMP_TIMEOUT_SECONDS:-5}"
VALIDATOR_AGENT_ICMP_COUNT: "${VALIDATOR_AGENT_ICMP_COUNT:-3}"
@@ -6,6 +6,7 @@ set -eu
export VALIDATOR_AGENT_POLL_INTERVAL_SECONDS="${VALIDATOR_AGENT_POLL_INTERVAL_SECONDS:-5}"
export VALIDATOR_AGENT_SELF_CHECK_TIMEOUT_SECONDS="${VALIDATOR_AGENT_SELF_CHECK_TIMEOUT_SECONDS:-10}"
export VALIDATOR_AGENT_SELF_CHECK_METHODS="${VALIDATOR_AGENT_SELF_CHECK_METHODS:-[ip_echo]}"
export VALIDATOR_AGENT_HTTPS_TIMEOUT_SECONDS="${VALIDATOR_AGENT_HTTPS_TIMEOUT_SECONDS:-10}"
export VALIDATOR_AGENT_ICMP_TIMEOUT_SECONDS="${VALIDATOR_AGENT_ICMP_TIMEOUT_SECONDS:-5}"
export VALIDATOR_AGENT_ICMP_COUNT="${VALIDATOR_AGENT_ICMP_COUNT:-3}"
@@ -16,7 +17,7 @@ export VALIDATOR_AGENT_SSH_TIMEOUT_SECONDS="${VALIDATOR_AGENT_SSH_TIMEOUT_SECOND
# NOT templated into the YAML: the binary reads them straight from this
# container's environment (names default in the config loader), so secrets
# never land in a file inside the container.
envsubst '${VALIDATOR_AGENT_VALIDATOR_ID} ${VALIDATOR_AGENT_CONTROL_API_URL} ${VALIDATOR_AGENT_POLL_INTERVAL_SECONDS} ${VALIDATOR_AGENT_SELF_CHECK_TIMEOUT_SECONDS} ${VALIDATOR_AGENT_HTTPS_TIMEOUT_SECONDS} ${VALIDATOR_AGENT_ICMP_TIMEOUT_SECONDS} ${VALIDATOR_AGENT_ICMP_COUNT} ${VALIDATOR_AGENT_SSH_ENABLED} ${VALIDATOR_AGENT_SSH_TIMEOUT_SECONDS}' \
envsubst '${VALIDATOR_AGENT_VALIDATOR_ID} ${VALIDATOR_AGENT_CONTROL_API_URL} ${VALIDATOR_AGENT_POLL_INTERVAL_SECONDS} ${VALIDATOR_AGENT_SELF_CHECK_TIMEOUT_SECONDS} ${VALIDATOR_AGENT_SELF_CHECK_METHODS} ${VALIDATOR_AGENT_HTTPS_TIMEOUT_SECONDS} ${VALIDATOR_AGENT_ICMP_TIMEOUT_SECONDS} ${VALIDATOR_AGENT_ICMP_COUNT} ${VALIDATOR_AGENT_SSH_ENABLED} ${VALIDATOR_AGENT_SSH_TIMEOUT_SECONDS}' \
< /etc/validator-agent/validator-agent.yaml.tmpl > /etc/validator-agent/validator-agent.yaml
exec /usr/local/bin/validator-agent -config /etc/validator-agent/validator-agent.yaml
@@ -4,6 +4,7 @@ poll_interval_seconds: ${VALIDATOR_AGENT_POLL_INTERVAL_SECONDS}
self_check:
timeout_seconds: ${VALIDATOR_AGENT_SELF_CHECK_TIMEOUT_SECONDS}
methods: ${VALIDATOR_AGENT_SELF_CHECK_METHODS}
checks:
https_timeout_seconds: ${VALIDATOR_AGENT_HTTPS_TIMEOUT_SECONDS}
+32 -4
View File
@@ -41,10 +41,10 @@ JSON, базовый префикс прикладных методов — `/ap
|---|---|---|
| **admin** | `CONTROL_API_ADMIN_TOKEN` | все `/api/v1/admin/*` (очередь, реестр, автоцикл, `config/*`) |
| **agent** | `CONTROL_API_AGENT_TOKEN` | запись результатов: `POST /agents/{id}/self-check`, `/events`, `/results`, `/complete` и `POST /probers/{site_id}/results` |
| **открыто** | — | `GET /healthz`; `POST /agents/register`, `POST /agents/{id}/heartbeat`, `GET /agents/{id}/assignment`; `POST /probers/register`, `POST /probers/{site_id}/heartbeat`, `GET /probers/{site_id}/assignments` |
| **открыто** | — | `GET /healthz`; `POST /agents/register`, `POST /agents/{id}/heartbeat`, `GET /agents/{id}/assignment`, `GET /agents/{id}/observed-ip`; `POST /probers/register`, `POST /probers/{site_id}/heartbeat`, `GET /probers/{site_id}/assignments` |
- Токены разные: токен администратора **не** подходит для методов агентов, и наоборот.
- Валидатор и пробер могут без токена зарегистрироваться, слать heartbeat и забирать задание (настройку); отправка результатов без токена агентов — `401`.
- Валидатор и пробер могут без токена зарегистрироваться, слать heartbeat и забирать задание (настройку); валидатор также может спросить, с какого адреса его видит control-api (`observed-ip`); отправка результатов без токена агентов — `401`.
- Имена переменных меняются в секции `auth` конфига control-api (`admin_token_env`, `agent_token_env`); сами значения в YAML не хранятся.
- Ответ при отказе: `401 {"error": "unauthorized"}` с заголовком `WWW-Authenticate: Bearer`.
- Токен не задан (пустая переменная) — уровень открыт; в логе control-api при старте предупреждение. Токены нужно генерировать случайными: `openssl rand -hex 32`.
@@ -145,6 +145,30 @@ curl -s -H "Authorization: Bearer $ADMIN_TOKEN" http://<control-api>:8080/api/v1
`check_config` — уже развёрнутая конфигурация проверок (тип + список
целей), агенту не нужно самому сопоставлять группы целей.
### `GET /api/v1/agents/{id}/observed-ip`
С какого адреса control-api видит соединение валидатора. Используется
self-check способом `control_api` (`self_check.methods` в
`validator-agent.yaml`) как альтернатива внешнему IP-echo сервису. Уровень
доступа — открыто: отдаётся только адрес самого вызывающего.
Ответ `200`:
```json
{"ip": "203.0.113.10", "source": "remote_addr"}
```
- Адрес берётся только из адреса TCP-соединения (`RemoteAddr`), приведённого
к каноничному виду (`::ffff:1.2.3.4` → `1.2.3.4`). Заголовки
`X-Forwarded-For` / `X-Real-IP` **не учитываются**: иначе валидатор мог бы
подделать адрес и пройти проверку. Метод рассчитан на прямое подключение
без обратного прокси.
- `404`, если `validator_id` не зарегистрирован.
- Способ корректен, только если соединение выходит через внешнюю сеть
(SNAT Floating IP). Если control-api достижим из облака по внутренней
сети, он увидит приватный адрес валидатора. Если порт control-api
опубликован через Docker, проверьте, что ручка показывает внешний адрес
клиента, а не адрес шлюза Docker.
### `POST /api/v1/agents/{id}/self-check`
Отчёт о результате self-check — подтверждение, что исходящий трафик
@@ -152,7 +176,11 @@ curl -s -H "Authorization: Bearer $ADMIN_TOKEN" http://<control-api>:8080/api/v1
определяет это **сам**, обращаясь к внешнему (снаружи облака) IP-echo
сервису (`self_check.ip_echo_urls` в `validator-agent.yaml`, например
`api.ipify.org`) и сравнивая ответ с `ip_address` из задания — control-api
в этом определении не участвует. Важно, что ресурс должен быть именно
в этом определении не участвует. Дополнительно можно включить способ
`control_api` (`self_check.methods`): агент спрашивает у control-api через
`GET /agents/{id}/observed-ip`, с какого адреса тот его видит. Способы
пробуются по приоритету, достаточно подтверждения любым из них; без
настройки работает только IP-echo. Важно, что ресурс должен быть именно
внешним: OpenStack применяет SNAT через Floating IP только к трафику,
уходящему через внешнюю сеть, поэтому обращение к чему-либо внутри
проекта (в том числе к самому control-api, если он в той же внутренней
@@ -165,7 +193,7 @@ curl -s -H "Authorization: Bearer $ADMIN_TOKEN" http://<control-api>:8080/api/v1
"ip_id": 42,
"detected_egress_ip": "203.0.113.10",
"success": true,
"detail": "matched"
"detail": "matched (control_api)"
}
```
+14 -5
View File
@@ -701,10 +701,17 @@ docker run -d --platform linux/amd64 --cap-add NET_RAW --name validator-agent \
| `VALIDATOR_AGENT_SSH_TIMEOUT_SECONDS` | нет | `5` |
`--cap-add NET_RAW` обязателен для ICMP-проверок, как и у `prober`.
`self_check.ip_echo_urls` в переменные не вынесен — при отсутствии в
конфиге агент сам подставляет дефолт (`api.ipify.org`, `ifconfig.me`);
свой список задавайте через смонтированный конфиг вместо шаблона, если
нужно переопределить.
`self_check.ip_echo_urls` и `self_check.methods` в переменные не вынесены —
при отсутствии в конфиге агент сам подставляет дефолты (`api.ipify.org`,
`ifconfig.me` и `methods: [ip_echo]`); свои значения задавайте через
смонтированный конфиг вместо шаблона, если нужно переопределить.
`methods` — способы самопроверки в порядке приоритета (`ip_echo`,
`control_api`), достаточно подтверждения любым. Способ `control_api`
спрашивает у control-api, с какого адреса он видит валидатора
(`GET /agents/{id}/observed-ip`); при внешнем размещении control-api
рекомендуется `[control_api, ip_echo]`. Ограничение: если control-api
достижим из облака по внутренней сети, он увидит приватный адрес
валидатора и этот способ всегда даст несовпадение — используйте `ip_echo`.
### Обновление образов после изменения кода
@@ -804,7 +811,9 @@ curl -s http://<control-api>:8080/api/v1/admin/validators | python3 -m json.tool
через floating IP, который в данный момент привязан к валидатору.
- `validator-agent` → внешние IP-echo сервисы из `self_check.ip_echo_urls`
(по умолчанию `api.ipify.org`, `ifconfig.me`) — **обязательно вне
облака**: это и есть механизм self-check (см.
облака**: это и есть механизм self-check способом `ip_echo` (при
`self_check.methods` с `control_api` достаточно ещё и доступа к control-api
по внешней сети; см.
[DIAGRAMS.md](DIAGRAMS.md#2-поток-данных-от-валидатора-к-целевому-серверу-egress-проверка)).
Если валидатор не может достучаться ни до одного из этих адресов,
self-check никогда не пройдёт и IP будет бесконечно возвращаться в
+9 -2
View File
@@ -742,7 +742,9 @@ curl -s -X PUT http://<control-api>:8080/api/v1/admin/config/orchestrator \
дело обычно в self-check: он запрашивает внешние (вне облака) сервисы из
`self_check.ip_echo_urls` в конфиге валидатора (по умолчанию
`api.ipify.org`, `ifconfig.me`) — если у ВМ-валидатора нет исходящего
доступа в интернет к этим адресам, запрос не проходит вообще, и агент
доступа в интернет к этим адресам, запрос не проходит вообще (при
`self_check.methods: [control_api, ip_echo]` агент сперва спросит адрес у
control-api, и проверка может пройти и без IP-echo), и агент
даже не может *сообщить* результат control-api (ни успешный, ни
неуспешный) — тогда статус реально зависает до истечения
`orchestrator.lease_ttl_seconds`, после чего адрес возвращается в
@@ -763,7 +765,12 @@ https://api.ipify.org`) и логи `journalctl -u validator-agent` на пре
облака (см. `self_check.ip_echo_urls`) — запрос к чему-либо внутри
проекта (в том числе к самому control-api, если он в той же внутренней
сети) покажет приватный адрес валидатора независимо от того, правильно
ли привязан FIP, и всегда будет давать ложный провал.
ли привязан FIP, и всегда будет давать ложный провал. Это относится и к
способу `control_api` (`self_check.methods`): он корректен только когда
валидатор ходит к control-api через внешнюю сеть; при внутреннем доступе в
`detail` будет подсказка про приватный адрес — оставьте `ip_echo`. В
`detail` события `self_check_result` указан сработавший способ
(`matched (control_api)`) либо причина по каждому способу.
**Площадка (`site-N`) никогда не отчитывается по конкретному IP.**
Сперва проверьте статус самой площадки — `GET
@@ -0,0 +1,110 @@
# План: самопроверка через control-api (дополнительный способ сверки публичного IP)
> Дата: 2026-10-02 03:06 · Статус: **реализовано и проверено** (юнит-тесты, локальный e2e с `ip_echo` и с `control_api`) (решения пользователя — в конце)
## Зачем
Самопроверка агента подтверждает, что исходящий трафик валидатора идёт через выданный Floating IP: агент спрашивает свой
публичный адрес у внешнего IP-echo сервиса и сверяет с назначенным адресом. Сейчас это единственный способ, и он хрупкий:
1 октября таймауты `https://ifconfig.me/ip` дали волну провалов self-check (33 случая за вечер).
`control-api` в текущем развёртывании стоит **во внешнем окружении** (вне облака), валидаторы подключаются к нему **напрямую**.
Значит, соединение валидатора с ним выходит наружу через Floating IP, и `control-api` сам видит публичный адрес источника.
Это даёт второй способ сверки без сторонних сервисов: агент спрашивает у `control-api`, с какого адреса тот его видит.
Требование: **существующий способ (IP-echo) сохраняется**, новый добавляется как опция агента.
## Решение в двух строках
1. `control-api` получает ручку «с какого адреса ты меня видишь».
2. Агент получает настройку `self_check.methods` — список способов в порядке приоритета; по умолчанию `[ip_echo]` (всё как сейчас).
## Конфигурация агента
```yaml
self_check:
timeout_seconds: 10
methods: [control_api, ip_echo] # по умолчанию [ip_echo]
ip_echo_urls: [...] # без изменений
```
- Допустимые значения: `ip_echo` (текущий способ), `control_api` (новый). Неизвестное значение — ошибка при старте агента.
- **Самопроверка успешна, если её подтвердил любой из способов.** Способы пробуются по порядку приоритета (первый —
главный); остановка на первом успешном. К следующему способу переходим и при отсутствии ответа (ошибка, таймаут,
404/5xx), и при несовпадении адреса. Провал — только если не подтвердил ни один способ; в `detail` попадает причина по
каждому способу.
- Внутри способа `ip_echo` поведение прежнее: URL перебираются по порядку, переход к следующему URL только при ошибке.
- Пустой список или отсутствие ключа → `[ip_echo]`. Агент без новой настройки ведёт себя ровно как раньше.
- Для внешнего размещения `control-api` в примерах и рекомендациях стоит `[control_api, ip_echo]`: способ через API в приоритете.
- Итог в `detail`: `detected_egress_ip=<ip> matched (control_api)`; видно в событиях адреса и в дашборде.
## Control API
**Новая ручка:** `GET /api/v1/agents/{id}/observed-ip` → `200 {"ip":"90.156.213.5","source":"remote_addr"}`.
- Уровень доступа — **открыто**, как `heartbeat` и `assignment`: ручка отдаёт только адрес самого вызывающего, секретов нет.
- `{id}` должен быть известным валидатором, иначе `404` (чтобы ручка не превращалась в публичный «узнай свой IP»).
- Адрес берётся только из `r.RemoteAddr`, приводится к каноничному виду (`::ffff:1.2.3.4` → `1.2.3.4`).
- Заголовки `X-Forwarded-For`/`X-Real-IP` **не учитываются**: подключение прямое, а доверие к заголовку позволило бы
валидатору подделать адрес и пройти проверку. Если появится обратный прокси, понадобится отдельная настройка
доверенных прокси — сейчас она не нужна (решение 1).
## Агент (`internal/agentcore`)
- Новая функция `detectViaControlAPI`: `GET /api/v1/agents/{id}/observed-ip` на `control_api_url`.
- **Новое TCP-соединение на каждый вызов** (отдельный `http.Transport` с `DisableKeepAlives`). Это ключевой момент:
соединение, открытое до привязки Floating IP (heartbeat, assignment), остаётся в старом NAT-состоянии и покажет
прежний адрес; общий клиент `apiclient` использовать нельзя.
- Токен агентов в этот запрос не нужен (ручка открытая) и не отправляется.
- Таймаут `self_check.timeout_seconds` (сейчас 10 с) действует **на каждый способ отдельно**: при общем дедлайне зависший
первый способ (приоритетный `control_api`) съел бы всё время, и запасной не успел бы ответить. Общий предел — таймаут × число способов.
- `detectPublicIP` заменяется перебором `methods` по приоритету; `fetchIPEcho` не меняется.
- Диагностика: если `control-api` вернул **частный** адрес (RFC 1918 и т. п.), в `detail` пишется подсказка: «control-api
доступен по внутренней сети, самопроверка через него невозможна; используйте ip_echo».
## Ограничение способа (важно для документации)
Способ `control_api` корректен только если соединение валидатора с `control-api` **выходит через внешнюю сеть**
(SNAT Floating IP). Если `control-api` достижим из облака по внутренней сети, он увидит частный адрес валидатора, и
этот способ всегда будет давать несовпадение (при `[control_api, ip_echo]` проверка пройдёт по `ip_echo`).
Если порт `control-api` опубликован через Docker, при выкладке проверить, что ручка показывает внешний адрес клиента, а
не адрес шлюза Docker.
## Откат и совместимость
- Новый агент + старый `control-api`: ручки нет (`404`), при `methods: [control_api, ip_echo]` агент переходит на `ip_echo`.
- Старый агент + новый `control-api`: ничего не меняется, ручка просто не вызывается.
- Откат: убрать `methods` из конфига агента (или вернуть `[ip_echo]`) и перезапустить агент.
- Схема БД и протокол `self-check` (`POST /agents/{id}/self-check`) не меняются.
## Затрагиваемые файлы
| Файл | Изменение |
|---|---|
| `internal/config/config.go` | `SelfCheckCfg.Methods`, дефолт `[ip_echo]` и проверка значений |
| `internal/httpapi/routes.go`, `handlers_agent.go` | маршрут и обработчик `observed-ip` |
| `internal/agentcore/agentcore.go` | `detectViaControlAPI`, перебор `methods` по приоритету |
| `configs/validator-agent.example.yaml` | пример и комментарии |
| `docs/API.md`, `docs/SETUP.md`, `docs/USAGE.md`, `README.md` | описание ручки, опции, ограничения |
| `bin/validator-agent`, `bin/control-api`, `SHA256SUMS` | пересборка |
## Тесты
- **config:** дефолт `[ip_echo]`; допустимые значения; ошибка на неизвестном способе.
- **httpapi:** прямой адрес; IPv4-mapped IPv6; заголовок `X-Forwarded-For` игнорируется; неизвестный валидатор → `404`.
- **agentcore:** порядок способов; успех второго способа после ошибки или несовпадения первого; провал, когда не
подтвердил ни один; каждый вызов открывает новое соединение (тестовый сервер считает соединения); подсказка про частный адрес.
- **e2e:** `scripts/run-local-e2e.sh` остаётся на `ip_echo` (проверка обратной совместимости).
## Выкладка
1. `control-api` с новой ручкой (поведение не меняется): пересборка образа, перезапуск.
2. Агенты на ВМ-валидаторах: новый `bin/validator-agent` и `methods: [control_api, ip_echo]` в их конфиге. Один валидатор
для начала, проверить `detail` в событиях (`matched (control_api)`), затем остальные.
## Решения пользователя (2026-10-02)
1. Подключение валидаторов к `control-api` — **напрямую** (без обратного прокси): `trusted_proxies` не нужен.
2. Ручка `observed-ip` — **открытая**.
3. Достаточно **одной успешной самопроверки любым из способов**; при внешнем размещении API способ через ручку API —
**в приоритете** (первый в списке).
+118 -11
View File
@@ -10,11 +10,13 @@ package agentcore
import (
"context"
"encoding/json"
"fmt"
"io"
"log/slog"
"net"
"net/http"
neturl "net/url"
"os"
"strings"
"time"
@@ -181,18 +183,13 @@ func (a *Agent) pollOnce(ctx context.Context) {
func (a *Agent) handleSelfCheckAndRun(ctx context.Context, assignment assignmentResp) {
a.postEvent(ctx, assignment.IPID, "config_received", "")
timeout := time.Duration(a.cfg.SelfCheck.TimeoutSeconds) * time.Second
// Each method gets the full timeout (see runSelfCheckMethods), so the
// overall budget scales with the number of methods.
timeout := time.Duration(a.cfg.SelfCheck.TimeoutSeconds) * time.Second * time.Duration(len(a.selfCheckMethods()))
selfCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
detectedIP, err := a.detectPublicIP(selfCtx)
success := err == nil && detectedIP == assignment.IPAddress
detail := "matched"
if err != nil {
detail = "ip echo request failed: " + err.Error()
} else if !success {
detail = fmt.Sprintf("egress ip %q does not match assigned fip %q", detectedIP, assignment.IPAddress)
}
detectedIP, _, detail, success := a.runSelfCheckMethods(selfCtx, assignment.IPAddress)
a.postSelfCheck(ctx, assignment.IPID, detectedIP, success, detail)
a.postEvent(ctx, assignment.IPID, "self_check_result", fmt.Sprintf(`{"success":%t}`, success))
@@ -204,7 +201,117 @@ func (a *Agent) handleSelfCheckAndRun(ctx context.Context, assignment assignment
a.runChecks(ctx, assignment)
}
// detectPublicIP asks each configured IP-echo URL, in order, for the
// selfCheckMethods returns the configured methods in priority order. Configs
// built without the loader (tests) may leave the list empty; that means the
// historical behaviour, ip_echo only.
func (a *Agent) selfCheckMethods() []string {
if len(a.cfg.SelfCheck.Methods) == 0 {
return []string{config.SelfCheckIPEcho}
}
return a.cfg.SelfCheck.Methods
}
// runSelfCheckMethods tries the configured methods in priority order and
// stops at the first one that confirms assignedIP. A method that gives no
// answer and one that reports a different address are treated alike: the
// next method is tried, since either may be a limitation of that method
// (e.g. control-api reached over the internal network sees a private
// address) rather than proof the floating IP is not attached. The check
// fails only when no method confirms, and detail then carries the reason
// from every method. detectedIP is the matching address on success, else the
// last address any method reported (may be empty).
//
// Every method runs under its own self_check.timeout_seconds: with one shared
// deadline a hung first method (priority control_api) would use it all up and
// the fallback would never get a chance to answer.
func (a *Agent) runSelfCheckMethods(ctx context.Context, assignedIP string) (detectedIP, method, detail string, ok bool) {
var reasons []string
perMethod := time.Duration(a.cfg.SelfCheck.TimeoutSeconds) * time.Second
for _, m := range a.selfCheckMethods() {
mctx, cancel := ctx, context.CancelFunc(func() {})
if perMethod > 0 {
mctx, cancel = context.WithTimeout(ctx, perMethod)
}
var ip string
var err error
switch m {
case config.SelfCheckControlAPI:
ip, err = a.detectViaControlAPI(mctx)
case config.SelfCheckIPEcho:
ip, err = a.detectViaIPEcho(mctx)
if err != nil {
err = fmt.Errorf("ip echo request failed: %w", err)
}
default:
err = fmt.Errorf("unknown self-check method")
}
cancel()
if err != nil {
reasons = append(reasons, m+": "+err.Error())
continue
}
if ip == assignedIP {
return ip, m, fmt.Sprintf("matched (%s)", m), true
}
detectedIP = ip
reason := fmt.Sprintf("%s: egress ip %q does not match assigned fip %q", m, ip, assignedIP)
if m == config.SelfCheckControlAPI && isLocalAddr(ip) {
reason += " (control-api sees a private address; it is reachable over the internal network, self-check via control_api is not possible, use ip_echo)"
}
reasons = append(reasons, reason)
}
if len(reasons) == 0 {
reasons = append(reasons, "no self-check methods configured")
}
return detectedIP, "", strings.Join(reasons, "; "), false
}
// isLocalAddr reports whether ip is a private, loopback or link-local
// address, i.e. one that can never be a floating IP.
func isLocalAddr(ip string) bool {
parsed := net.ParseIP(ip)
return parsed != nil && (parsed.IsPrivate() || parsed.IsLoopback() || parsed.IsLinkLocalUnicast())
}
// detectViaControlAPI asks control-api which source address it sees for this
// validator. It is only meaningful when control-api is reached over the
// external network, where the floating IP is the visible source (see
// config.SelfCheckCfg).
//
// Every call dials a brand-new TCP connection through a dedicated transport
// with keep-alives off: a connection opened before the floating IP was
// attached (heartbeat, assignment polling) keeps its old NAT state and would
// keep reporting the previous address, so the shared apiclient must not be
// used. The endpoint is open, so no agent token is sent.
func (a *Agent) detectViaControlAPI(ctx context.Context) (string, error) {
url := strings.TrimRight(a.cfg.ControlAPIURL, "/") + "/api/v1/agents/" + neturl.PathEscape(a.cfg.ValidatorID) + "/observed-ip"
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return "", fmt.Errorf("build request: %w", err)
}
transport := &http.Transport{DisableKeepAlives: true}
defer transport.CloseIdleConnections()
resp, err := (&http.Client{Transport: transport}).Do(req)
if err != nil {
return "", err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode > 299 {
return "", fmt.Errorf("unexpected status %d", resp.StatusCode)
}
var body struct {
IP string `json:"ip"`
}
if err := json.NewDecoder(io.LimitReader(resp.Body, 4096)).Decode(&body); err != nil {
return "", fmt.Errorf("decode response: %w", err)
}
if net.ParseIP(body.IP) == nil {
return "", fmt.Errorf("response is not a valid IP: %q", body.IP)
}
return body.IP, nil
}
// detectViaIPEcho asks each configured IP-echo URL, in order, for the
// address this validator is currently seen egressing from, returning the
// first one that answers with a parseable IP. These must be resources
// genuinely outside the cloud project (see config.SelfCheckCfg) — OpenStack
@@ -212,7 +319,7 @@ func (a *Agent) handleSelfCheckAndRun(ctx context.Context, assignment assignment
// network, so anything reachable over the project's internal network would
// report the validator's private address instead, regardless of whether
// the floating IP is correctly attached.
func (a *Agent) detectPublicIP(ctx context.Context) (string, error) {
func (a *Agent) detectViaIPEcho(ctx context.Context) (string, error) {
var lastErr error
for _, url := range a.cfg.SelfCheck.IPEchoURLs {
ip, err := fetchIPEcho(ctx, url)
+232
View File
@@ -0,0 +1,232 @@
package agentcore
import (
"context"
"fmt"
"net"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"time"
"cloudipvalidator/internal/config"
)
// echoServer answers every request with body (an IP-echo stand-in).
func echoServer(t *testing.T, body string) *httptest.Server {
t.Helper()
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
fmt.Fprint(w, body)
}))
t.Cleanup(ts.Close)
return ts
}
// controlAPIServer plays control-api's observed-ip route: it answers with ip
// (status 200) or with the given error status when ip is empty.
func controlAPIServer(t *testing.T, ip string, status int) *httptest.Server {
t.Helper()
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/api/v1/agents/val-1/observed-ip" {
http.NotFound(w, r)
return
}
if ip == "" {
w.WriteHeader(status)
return
}
fmt.Fprintf(w, `{"ip":%q,"source":"remote_addr"}`, ip)
}))
t.Cleanup(ts.Close)
return ts
}
func selfCheckAgent(controlAPIURL string, methods []string, echoURLs ...string) *Agent {
return &Agent{
cfg: &config.ValidatorAgent{
ValidatorID: "val-1",
ControlAPIURL: controlAPIURL,
SelfCheck: config.SelfCheckCfg{TimeoutSeconds: 2, Methods: methods, IPEchoURLs: echoURLs},
},
log: testLogger(),
}
}
func TestSelfCheckControlAPIFirstWins(t *testing.T) {
capi := controlAPIServer(t, "1.2.3.4", 0)
var echoCalls int32
echo := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&echoCalls, 1)
fmt.Fprint(w, "1.2.3.4")
}))
defer echo.Close()
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
ip, method, detail, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
if !ok || ip != "1.2.3.4" || method != config.SelfCheckControlAPI || detail != "matched (control_api)" {
t.Fatalf("got ip=%q method=%q detail=%q ok=%v", ip, method, detail, ok)
}
if n := atomic.LoadInt32(&echoCalls); n != 0 {
t.Fatalf("ip_echo was called %d times although control_api already confirmed", n)
}
}
// Any one confirming method is enough: control-api answers with a different
// address, ip_echo confirms.
func TestSelfCheckFallsThroughOnMismatch(t *testing.T) {
capi := controlAPIServer(t, "5.5.5.5", 0)
echo := echoServer(t, "1.2.3.4")
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
ip, method, _, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
if !ok || ip != "1.2.3.4" || method != config.SelfCheckIPEcho {
t.Fatalf("got ip=%q method=%q ok=%v, want a pass via ip_echo", ip, method, ok)
}
}
// An old control-api without the route (404) or a failing one (5xx) must not
// stop the self-check: the next method decides.
func TestSelfCheckFallsThroughOnControlAPIError(t *testing.T) {
for _, status := range []int{http.StatusNotFound, http.StatusInternalServerError} {
capi := controlAPIServer(t, "", status)
echo := echoServer(t, "1.2.3.4")
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
_, method, _, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
if !ok || method != config.SelfCheckIPEcho {
t.Fatalf("status %d: method=%q ok=%v, want a pass via ip_echo", status, method, ok)
}
}
}
func TestSelfCheckFailsWhenNoMethodConfirms(t *testing.T) {
capi := controlAPIServer(t, "10.0.0.5", 0) // private address: internal-network case
echo := echoServer(t, "6.6.6.6")
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
ip, method, detail, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
if ok || method != "" {
t.Fatalf("expected a failure, got ok=%v method=%q", ok, method)
}
if ip != "6.6.6.6" {
t.Fatalf("detected ip = %q, want the last reported address 6.6.6.6", ip)
}
for _, want := range []string{
`control_api: egress ip "10.0.0.5" does not match assigned fip "1.2.3.4"`,
"private address", // the hint for the internal-network case
`ip_echo: egress ip "6.6.6.6" does not match`,
} {
if !strings.Contains(detail, want) {
t.Fatalf("detail %q does not contain %q", detail, want)
}
}
}
func TestSelfCheckPrivateHintOnlyForControlAPI(t *testing.T) {
echo := echoServer(t, "10.1.1.1")
a := selfCheckAgent("http://unused", []string{config.SelfCheckIPEcho}, echo.URL)
_, _, detail, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
if ok || strings.Contains(detail, "private address") {
t.Fatalf("ok=%v detail=%q: the control_api hint must not appear for ip_echo", ok, detail)
}
}
// With nothing configured (a config built without the loader) the agent
// behaves as before: ip_echo only, control-api is never asked.
func TestSelfCheckDefaultsToIPEcho(t *testing.T) {
var capiCalls int32
capi := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&capiCalls, 1)
}))
defer capi.Close()
echo := echoServer(t, "1.2.3.4")
a := selfCheckAgent(capi.URL, nil, echo.URL)
_, method, _, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
if !ok || method != config.SelfCheckIPEcho || atomic.LoadInt32(&capiCalls) != 0 {
t.Fatalf("method=%q ok=%v control-api calls=%d", method, ok, atomic.LoadInt32(&capiCalls))
}
}
// A hung control-api must not use up the time of the fallback: each method
// has its own timeout.
func TestSelfCheckHungControlAPIDoesNotStarveFallback(t *testing.T) {
release := make(chan struct{})
capi := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
<-release
}))
defer capi.Close()
defer close(release)
echo := echoServer(t, "1.2.3.4")
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
a.cfg.SelfCheck.TimeoutSeconds = 1
start := time.Now()
// The outer context mirrors handleSelfCheckAndRun: timeout x methods.
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
_, method, detail, ok := a.runSelfCheckMethods(ctx, "1.2.3.4")
if !ok || method != config.SelfCheckIPEcho {
t.Fatalf("method=%q ok=%v detail=%q, want a pass via ip_echo after control_api timed out", method, ok, detail)
}
if elapsed := time.Since(start); elapsed > 1900*time.Millisecond {
t.Fatalf("took %s: the hung method consumed the fallback's time", elapsed)
}
}
// Every control-api request must use a new TCP connection: a connection
// opened before the floating IP was attached would report the old address.
func TestDetectViaControlAPIDialsNewConnectionEachTime(t *testing.T) {
var conns int32
ts := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
fmt.Fprint(w, `{"ip":"1.2.3.4","source":"remote_addr"}`)
}))
ts.Config.ConnState = func(_ net.Conn, s http.ConnState) {
if s == http.StateNew {
atomic.AddInt32(&conns, 1)
}
}
ts.Start()
defer ts.Close()
a := selfCheckAgent(ts.URL, nil)
for i := 0; i < 3; i++ {
if _, err := a.detectViaControlAPI(context.Background()); err != nil {
t.Fatalf("call %d: %v", i, err)
}
}
if n := atomic.LoadInt32(&conns); n != 3 {
t.Fatalf("3 calls opened %d connections, want 3 (no keep-alive reuse)", n)
}
}
func TestDetectViaControlAPIRejectsBadAnswers(t *testing.T) {
for name, body := range map[string]string{"not json": "oops", "not an ip": `{"ip":"abc"}`, "empty": `{}`} {
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, body) }))
a := selfCheckAgent(ts.URL, nil)
if ip, err := a.detectViaControlAPI(context.Background()); err == nil {
t.Fatalf("%s: expected an error, got %q", name, ip)
}
ts.Close()
}
}
// The agent token must never be sent on this request (the route is open and
// the token is meant for control-api writes only).
func TestDetectViaControlAPISendsNoToken(t *testing.T) {
var auth string
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
auth = r.Header.Get("Authorization")
fmt.Fprint(w, `{"ip":"1.2.3.4"}`)
}))
defer ts.Close()
a := selfCheckAgent(ts.URL, nil)
if _, err := a.detectViaControlAPI(context.Background()); err != nil {
t.Fatal(err)
}
if auth != "" {
t.Fatalf("Authorization header sent: %q", auth)
}
}
+24 -3
View File
@@ -233,17 +233,30 @@ type ValidatorAgent struct {
Checks AgentChecks `yaml:"checks"`
}
// Self-check methods accepted in SelfCheckCfg.Methods.
const (
SelfCheckIPEcho = "ip_echo"
SelfCheckControlAPI = "control_api"
)
// SelfCheckCfg configures how the agent confirms its egress actually flows
// through the newly assigned floating IP. This must query a resource
// genuinely outside the cloud project: OpenStack only applies floating-IP
// SNAT to traffic leaving via the external/provider network, so any
// in-project resource (including control-api, if it's reachable over the
// project's internal network) would see the validator's private address
// instead — a false negative that never changes. IPEchoURLs are tried in
// order (falling through to the next on error/timeout, not on a genuine
// mismatch) until one returns a parseable IP.
// instead — a false negative that never changes. The control_api method is
// therefore only valid when control-api is reached over the external network
// (hosted outside the cloud); then it sees the floating IP as the source.
//
// Methods are tried in priority order and the self-check passes as soon as
// any one confirms the address; the next method is tried both when one gives
// no answer and when it reports a different address. IPEchoURLs are tried in
// order within ip_echo (falling through to the next on error/timeout, not on
// a genuine mismatch) until one returns a parseable IP.
type SelfCheckCfg struct {
TimeoutSeconds int `yaml:"timeout_seconds"`
Methods []string `yaml:"methods"`
IPEchoURLs []string `yaml:"ip_echo_urls"`
}
@@ -273,6 +286,14 @@ func LoadValidatorAgent(path string) (*ValidatorAgent, error) {
if c.SelfCheck.TimeoutSeconds == 0 {
c.SelfCheck.TimeoutSeconds = 10
}
if len(c.SelfCheck.Methods) == 0 {
c.SelfCheck.Methods = []string{SelfCheckIPEcho}
}
for _, m := range c.SelfCheck.Methods {
if m != SelfCheckIPEcho && m != SelfCheckControlAPI {
return nil, fmt.Errorf("self_check.methods: unknown method %q (allowed: %s, %s)", m, SelfCheckIPEcho, SelfCheckControlAPI)
}
}
if len(c.SelfCheck.IPEchoURLs) == 0 {
c.SelfCheck.IPEchoURLs = []string{"https://api.ipify.org", "https://ifconfig.me/ip"}
}
+55
View File
@@ -64,3 +64,58 @@ func TestControlAPIExampleConfigs(t *testing.T) {
t.Fatalf("load docker example: %v", err)
}
}
func writeAgentConfig(t *testing.T, selfCheck string) string {
t.Helper()
path := filepath.Join(t.TempDir(), "agent.yaml")
body := "validator_id: v1\ncontrol_api_url: http://x\nself_check:\n" + selfCheck
if err := os.WriteFile(path, []byte(body), 0o600); err != nil {
t.Fatal(err)
}
return path
}
func TestLoadValidatorAgentSelfCheckMethods(t *testing.T) {
// No methods key: the historical behaviour, ip_echo only.
c, err := LoadValidatorAgent(writeAgentConfig(t, " timeout_seconds: 5\n"))
if err != nil {
t.Fatalf("load: %v", err)
}
if len(c.SelfCheck.Methods) != 1 || c.SelfCheck.Methods[0] != SelfCheckIPEcho {
t.Fatalf("default methods = %v, want [ip_echo]", c.SelfCheck.Methods)
}
// Order is the priority and must be kept.
c, err = LoadValidatorAgent(writeAgentConfig(t, " methods: [control_api, ip_echo]\n"))
if err != nil {
t.Fatalf("load: %v", err)
}
if got := c.SelfCheck.Methods; len(got) != 2 || got[0] != SelfCheckControlAPI || got[1] != SelfCheckIPEcho {
t.Fatalf("methods = %v, want [control_api ip_echo]", got)
}
// An unknown method is a configuration error, not a silent no-op.
if _, err := LoadValidatorAgent(writeAgentConfig(t, " methods: [control_api, ipecho]\n")); err == nil {
t.Fatal("expected an error for an unknown self-check method")
}
}
// The shipped agent example must load, and the rxprod copy must stay
// byte-identical to it.
func TestValidatorAgentExampleConfig(t *testing.T) {
c, err := LoadValidatorAgent("../../configs/validator-agent.example.yaml")
if err != nil {
t.Fatalf("load example: %v", err)
}
if got := c.SelfCheck.Methods; len(got) != 2 || got[0] != SelfCheckControlAPI {
t.Fatalf("example methods = %v", got)
}
a, _ := os.ReadFile("../../configs/validator-agent.example.yaml")
b, err := os.ReadFile("../../rxprod-compose/sources/validator-agent.example.yaml")
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(a, b) {
t.Fatal("rxprod-compose/sources/validator-agent.example.yaml differs from configs/validator-agent.example.yaml")
}
}
+2 -2
View File
@@ -97,8 +97,8 @@ func TestRouteTableIsClassified(t *testing.T) {
t.Fatalf("admin route %q is %s, want admin", rt.Pattern, rt.Access)
}
}
if counts["admin"] != 34 || counts["agent"] != 5 || counts["open"] != 7 {
t.Fatalf("access counts = %v, want admin=34 agent=5 open=7", counts)
if counts["admin"] != 34 || counts["agent"] != 5 || counts["open"] != 8 {
t.Fatalf("access counts = %v, want admin=34 agent=5 open=8", counts)
}
}
+5
View File
@@ -23,6 +23,11 @@ type okResponse struct {
OK bool `json:"ok"`
}
type observedIPResponse struct {
IP string `json:"ip"`
Source string `json:"source"`
}
type checkConfigDTO struct {
Type string `json:"type"`
Targets []string `json:"targets"`
+45
View File
@@ -1,7 +1,12 @@
package httpapi
import (
"database/sql"
"errors"
"fmt"
"net"
"net/http"
"net/netip"
"time"
"cloudipvalidator/internal/db"
@@ -64,6 +69,46 @@ func (s *Server) handleAgentAssignment(w http.ResponseWriter, r *http.Request) {
})
}
// handleAgentObservedIP tells a validator which source address control-api
// sees for its connection, so the agent's self-check can confirm the floating
// IP without a third-party IP-echo service. The route is open, so it is
// limited to known validators to keep it from becoming a public "what is my
// IP" service.
func (s *Server) handleAgentObservedIP(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
if _, err := s.DB.GetValidator(r.Context(), id); err != nil {
if errors.Is(err, sql.ErrNoRows) || errors.Is(err, db.ErrNotFound) {
writeError(w, http.StatusNotFound, "unknown validator: "+id)
return
}
writeError(w, http.StatusInternalServerError, err.Error())
return
}
ip, err := clientIP(r)
if err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
writeJSON(w, http.StatusOK, observedIPResponse{IP: ip, Source: "remote_addr"})
}
// clientIP returns the peer address of the TCP connection in canonical form
// (IPv4-mapped IPv6 unmapped to plain IPv4). It deliberately ignores
// X-Forwarded-For / X-Real-IP: validators connect directly, and trusting a
// client-supplied header would let a validator forge the address and pass
// the self-check.
func clientIP(r *http.Request) (string, error) {
host, _, err := net.SplitHostPort(r.RemoteAddr)
if err != nil {
return "", fmt.Errorf("parse remote address %q: %w", r.RemoteAddr, err)
}
addr, err := netip.ParseAddr(host)
if err != nil {
return "", fmt.Errorf("parse remote address %q: %w", r.RemoteAddr, err)
}
return addr.Unmap().String(), nil
}
func (s *Server) handleAgentSelfCheck(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
var req selfCheckRequest
@@ -0,0 +1,79 @@
package httpapi
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
)
func TestClientIP(t *testing.T) {
cases := []struct {
name string
remoteAddr string
headers map[string]string
want string
wantErr bool
}{
{name: "ipv4", remoteAddr: "90.156.213.5:51234", want: "90.156.213.5"},
{name: "ipv4 mapped in ipv6", remoteAddr: "[::ffff:90.156.213.5]:51234", want: "90.156.213.5"},
{name: "ipv6", remoteAddr: "[2001:db8::7]:443", want: "2001:db8::7"},
// A client-supplied header must never change the answer.
{name: "forwarded for is ignored", remoteAddr: "90.156.213.5:1",
headers: map[string]string{"X-Forwarded-For": "1.2.3.4", "X-Real-IP": "5.6.7.8"}, want: "90.156.213.5"},
{name: "no port", remoteAddr: "90.156.213.5", wantErr: true},
{name: "not an address", remoteAddr: "host:80", wantErr: true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
r := httptest.NewRequest(http.MethodGet, "/", nil)
r.RemoteAddr = tc.remoteAddr
for k, v := range tc.headers {
r.Header.Set(k, v)
}
got, err := clientIP(r)
if tc.wantErr {
if err == nil {
t.Fatalf("expected an error, got %q", got)
}
return
}
if err != nil || got != tc.want {
t.Fatalf("clientIP = %q, %v; want %q", got, err, tc.want)
}
})
}
}
// The route is open (no token), answers a known validator with the address
// control-api sees, and refuses unknown validators and a forged header.
func TestObservedIPEndpoint(t *testing.T) {
fc, d, _, _ := newConfigTestHarness(t)
if err := d.RegisterValidator(context.Background(), "validator-1", "host-1", "port-1", "v0.1"); err != nil {
t.Fatal(err)
}
req, _ := http.NewRequest(http.MethodGet, fc.base+"/api/v1/agents/validator-1/observed-ip", nil)
req.Header.Set("X-Forwarded-For", "8.8.8.8")
resp, err := fc.client.Do(req)
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
t.Fatalf("status = %d, want 200", resp.StatusCode)
}
var got observedIPResponse
if err := json.NewDecoder(resp.Body).Decode(&got); err != nil {
t.Fatal(err)
}
// httptest connects from loopback; the forged header must not leak through.
if got.IP != "127.0.0.1" || got.Source != "remote_addr" {
t.Fatalf("response = %+v, want 127.0.0.1 / remote_addr", got)
}
if resp, _ := fc.do(http.MethodGet, "/api/v1/agents/no-such-validator/observed-ip", nil); resp.StatusCode != http.StatusNotFound {
t.Fatalf("unknown validator: status = %d, want 404", resp.StatusCode)
}
}
+1
View File
@@ -39,6 +39,7 @@ func (s *Server) routeTable() []route {
{"POST /api/v1/agents/register", s.handleAgentRegister, accessOpen},
{"POST /api/v1/agents/{id}/heartbeat", s.handleAgentHeartbeat, accessOpen},
{"GET /api/v1/agents/{id}/assignment", s.handleAgentAssignment, accessOpen},
{"GET /api/v1/agents/{id}/observed-ip", s.handleAgentObservedIP, accessOpen},
{"POST /api/v1/agents/{id}/self-check", s.handleAgentSelfCheck, accessAgent},
{"POST /api/v1/agents/{id}/events", s.handleAgentEvent, accessAgent},
{"POST /api/v1/agents/{id}/results", s.handleAgentResults, accessAgent},
+19
View File
@@ -222,6 +222,25 @@ func (o *Orchestrator) SelfCheckResult(ctx context.Context, validatorID string,
if err != nil {
return err
}
// A late report for an address this validator no longer owns (lease
// already reclaimed, address cancelled or deleted) must not touch
// the floating IP: it may be attached for another validator now.
if item.State != db.IPAwaitingSelfCheck || item.OwnerValidatorID == nil || *item.OwnerValidatorID != validatorID {
o.Log.Warn("ignoring failed self-check for an address the validator does not hold",
"validator", validatorID, "ip_id", ipID, "state", item.State)
return nil
}
// Detach the floating IP before the address goes back to the queue.
// requeueOrFail frees the validator in the database but knows nothing
// about the cloud: a floating IP left on the validator's port makes
// every later association on that port fail with 409 ("fixed IP
// already has a floating IP"). Best-effort, like the other release
// paths — the database state must be freed even if Neutron hiccups.
if item.FIPID != "" {
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil {
o.Log.Error("disassociate fip after failed self-check", "ip_id", ipID, "fip_id", item.FIPID, "err", err)
}
}
if item.RetryCount+1 > o.Cfg.MaxSelfCheckRetries {
o.requeueOrFail(ctx, ipID, validatorID, "self-check failed: "+detail)
return nil
@@ -1086,3 +1086,79 @@ func TestClearQueueFreesAllValidatorPorts(t *testing.T) {
t.Fatalf("port-2: expected only the foreign 9.9.9.9 to remain, got %+v", got)
}
}
// A failed self-check must detach the floating IP in the cloud before the
// address returns to the queue; otherwise the validator's port keeps it and
// every later association on that port fails with 409.
func TestFailedSelfCheckDetachesFloatingIP(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.Seed("fip-1", "1.2.3.4", "svc-project")
if err := d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1"); err != nil {
t.Fatal(err)
}
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatal(err)
}
o.Tick(ctx)
if got := portFIPs(t, mock, "port-1"); len(got) != 1 {
t.Fatalf("setup: expected the floating ip on port-1, got %d", len(got))
}
ip, err := d.GetIPByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatal(err)
}
if err := o.SelfCheckResult(ctx, "validator-1", ip.ID, false, "ip echo timeout"); err != nil {
t.Fatalf("self-check result: %v", err)
}
if got := portFIPs(t, mock, "port-1"); len(got) != 0 {
t.Fatalf("port-1 still holds the floating ip after a failed self-check: %+v", got)
}
after, err := d.GetIPByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatal(err)
}
if after.State != db.IPQueued {
t.Fatalf("expected the address back in queued, got %s", after.State)
}
// The validator must be able to take the next claim: a second Tick
// associates the same address again without a conflict.
o.Tick(ctx)
if got := portFIPs(t, mock, "port-1"); len(got) != 1 {
t.Fatalf("expected the retry to attach the floating ip again, got %d", len(got))
}
}
// A failed self-check for an address the validator no longer holds is a late
// report: it must not detach a floating IP that now belongs to someone else.
func TestLateFailedSelfCheckIsIgnored(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.Seed("fip-1", "1.2.3.4", "svc-project")
for i, id := range []string{"validator-1", "validator-2"} {
if err := d.RegisterValidator(ctx, id, "h", fmt.Sprintf("port-%d", i+1), "v"); err != nil {
t.Fatal(err)
}
}
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatal(err)
}
o.Tick(ctx) // validator-1 (first by id) holds 1.2.3.4
ip, _ := d.GetIPByAddress(ctx, "1.2.3.4")
// validator-2 reports a failure for an address it does not own.
if err := o.SelfCheckResult(ctx, "validator-2", ip.ID, false, "late"); err != nil {
t.Fatalf("self-check result: %v", err)
}
if got := portFIPs(t, mock, "port-1"); len(got) != 1 {
t.Fatalf("a late report from another validator detached the floating ip: %+v", got)
}
cur, _ := d.GetIPByAddress(ctx, "1.2.3.4")
if cur.State != db.IPAwaitingSelfCheck {
t.Fatalf("a late report changed the address state to %s", cur.State)
}
}
@@ -13,7 +13,19 @@ poll_interval_seconds: 5
self_check:
timeout_seconds: 10
# Must be a resource genuinely outside the cloud project — OpenStack only
# Способы самопроверки в порядке приоритета (допустимо: ip_echo,
# control_api); по умолчанию [ip_echo]. Самопроверка успешна, если адрес
# подтвердил любой способ: пробуются по порядку, остановка на первом
# успешном, к следующему переходим и при отсутствии ответа, и при
# несовпадении адреса. Таймаут timeout_seconds действует на каждый способ
# отдельно (зависший первый способ не лишает второй времени). control_api спрашивает у control-api, с какого адреса он видит
# это соединение; рекомендуется [control_api, ip_echo], когда control-api
# стоит вне облака и валидатор ходит к нему напрямую (через внешнюю сеть).
# Ограничение: если control-api достижим по внутренней сети облака, он
# увидит частный адрес валидатора и control_api всегда даст несовпадение —
# тогда оставьте только ip_echo (или он сработает вторым в списке).
methods: [control_api, ip_echo]
# Used by the ip_echo method. Must be a resource genuinely outside the cloud project — OpenStack only
# applies floating-IP SNAT to traffic leaving via the external network,
# so anything reachable over the project's internal network (including
# control-api itself, if it's on the same internal network) would report
+3
View File
@@ -110,6 +110,9 @@ control_api_url: "http://127.0.0.1:28080"
poll_interval_seconds: 1
self_check:
timeout_seconds: 5
# ip_echo by default; E2E_SELF_CHECK_METHODS="[control_api]" runs the same
# scenario through control-api's observed-ip route (loopback sees 127.0.0.1).
methods: ${E2E_SELF_CHECK_METHODS:-[ip_echo]}
# Points at the local httpstub's /ip echo route rather than a real
# internet IP-echo service — self-check only produces a meaningful
# signal here because the "floating IP" under test (127.0.0.1) is