From 273cd5441593f20c6164b997dec445990c75cdf5 Mon Sep 17 00:00:00 2001 From: Vassiliy Yegorov Date: Mon, 21 Sep 2026 21:50:07 +0700 Subject: [PATCH] =?UTF-8?q?feat:=20traefik-selectel-dns=20=E2=80=94=20A-?= =?UTF-8?q?=D0=B7=D0=B0=D0=BF=D0=B8=D1=81=D0=B8=20Selectel=20=D0=BF=D0=BE?= =?UTF-8?q?=20=D1=80=D0=BE=D1=83=D1=82=D0=B0=D0=BC=20Traefik=20+=20CI=20?= =?UTF-8?q?=D1=81=D0=B1=D0=BE=D1=80=D0=BA=D0=B0=20=D0=BE=D0=B1=D1=80=D0=B0?= =?UTF-8?q?=D0=B7=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- .dockerignore | 13 + .gitea/workflows/build.yaml | 118 ++++++ .gitignore | 6 + Dockerfile | 22 ++ README.md | 67 ++++ cmd/traefik-selectel-dns/main.go | 135 +++++++ deploy/config.example.yaml | 36 ++ deploy/docker-compose.example.yml | 97 +++++ deploy/dynamic/api-internal.yml | 14 + deploy/dynamic/example.yml | 13 + deploy/traefik.example.yml | 56 +++ deploy/whoami.example.yml | 17 + go.mod | 5 + go.sum | 4 + internal/config/config.go | 286 ++++++++++++++ internal/config/config_test.go | 140 +++++++ internal/dockerevents/watch.go | 134 +++++++ internal/dockerevents/watch_test.go | 105 +++++ internal/health/health.go | 122 ++++++ internal/health/health_test.go | 40 ++ internal/hostparse/hostparse.go | 212 ++++++++++ internal/hostparse/hostparse_test.go | 78 ++++ internal/hostparse/zone.go | 17 + internal/httpx/retry.go | 128 ++++++ internal/httpx/retry_test.go | 99 +++++ internal/reconciler/loop.go | 61 +++ internal/reconciler/loop_test.go | 40 ++ internal/reconciler/reconciler.go | 321 +++++++++++++++ internal/reconciler/reconciler_test.go | 368 ++++++++++++++++++ internal/selectel/auth.go | 186 +++++++++ internal/selectel/client.go | 308 +++++++++++++++ internal/selectel/selectel_test.go | 288 ++++++++++++++ internal/traefik/client.go | 90 +++++ internal/traefik/client_test.go | 58 +++ .../traefik-selectel-dns-2026-09-21.md | 38 ++ 35 files changed, 3722 insertions(+) create mode 100644 .dockerignore create mode 100644 .gitea/workflows/build.yaml create mode 100644 .gitignore create mode 100644 Dockerfile create mode 100644 README.md create mode 100644 cmd/traefik-selectel-dns/main.go create mode 100644 deploy/config.example.yaml create mode 100644 deploy/docker-compose.example.yml create mode 100644 deploy/dynamic/api-internal.yml create mode 100644 deploy/dynamic/example.yml create mode 100644 deploy/traefik.example.yml create mode 100644 deploy/whoami.example.yml create mode 100644 go.mod create mode 100644 go.sum create mode 100644 internal/config/config.go create mode 100644 internal/config/config_test.go create mode 100644 internal/dockerevents/watch.go create mode 100644 internal/dockerevents/watch_test.go create mode 100644 internal/health/health.go create mode 100644 internal/health/health_test.go create mode 100644 internal/hostparse/hostparse.go create mode 100644 internal/hostparse/hostparse_test.go create mode 100644 internal/hostparse/zone.go create mode 100644 internal/httpx/retry.go create mode 100644 internal/httpx/retry_test.go create mode 100644 internal/reconciler/loop.go create mode 100644 internal/reconciler/loop_test.go create mode 100644 internal/reconciler/reconciler.go create mode 100644 internal/reconciler/reconciler_test.go create mode 100644 internal/selectel/auth.go create mode 100644 internal/selectel/client.go create mode 100644 internal/selectel/selectel_test.go create mode 100644 internal/traefik/client.go create mode 100644 internal/traefik/client_test.go create mode 100644 swarm-report/traefik-selectel-dns-2026-09-21.md diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..569d1ce --- /dev/null +++ b/.dockerignore @@ -0,0 +1,13 @@ +.git +.gitignore +.dockerignore +Dockerfile +deploy +swarm-report +README.md +secrets +*.md +**/*_test.go + +.env +.env.* diff --git a/.gitea/workflows/build.yaml b/.gitea/workflows/build.yaml new file mode 100644 index 0000000..0349c0d --- /dev/null +++ b/.gitea/workflows/build.yaml @@ -0,0 +1,118 @@ +name: Build + +on: + push: + branches: [main, master] + paths: + - "cmd/**" + - "internal/**" + - "go.mod" + - "go.sum" + - "Dockerfile" + - ".gitea/workflows/**" + workflow_dispatch: + +env: + REGISTRY: git.realmanual.ru + # Образ называется как репозиторий: git.realmanual.ru/pub/traefik-selectel + IMAGE: git.realmanual.ru/${{ gitea.repository }} + +# Секреты репозитория: TOKEN (write:package — сборка пушит образ), +# TELEGRAM_BOT_URI, TELEGRAM_BOT_TOKEN, TELEGRAM_CHAT_ID. +permissions: + contents: read + packages: write + +jobs: + # Имя без дефиса намеренно: к результату job'а обращаются как + # needs.<имя>.result, а дефис в этом выражении разбирается как минус. + tests: + name: Tests + runs-on: ubuntu-22.04 + container: catthehacker/ubuntu:act-latest + steps: + - uses: actions/checkout@v3 + + - name: Set up Go + uses: actions/setup-go@v5 + with: + go-version-file: go.mod + # Кэш модулей выключен: в self-hosted раннерах сервис кэша actions + # обычно недоступен (getCacheEntry → ETIMEDOUT). + cache: false + + - name: Run tests + run: go test ./... -count=1 + + - name: Vet + run: go vet ./... + + # gofmt не падает сам по себе, поэтому проверяем список файлов. + - name: Check gofmt + run: | + unformatted="$(gofmt -l .)" + if [ -n "$unformatted" ]; then + echo "не отформатированы:" + echo "$unformatted" + exit 1 + fi + + build: + needs: tests + name: Build image + runs-on: ubuntu-22.04 + container: catthehacker/ubuntu:act-latest + outputs: + tag: ${{ steps.meta.outputs.TAG }} + steps: + - uses: actions/checkout@v3 + + - name: Compute tag + id: meta + # Короткий SHA коммита: тег неизменяем, по нему видно, какой код + # работает на ВМ, и есть куда откатиться. + run: echo "TAG=$(echo "${{ github.sha }}" | cut -c1-8)" >> $GITHUB_OUTPUT + + - name: Log in to registry + uses: docker/login-action@v3 + with: + registry: ${{ env.REGISTRY }} + username: ${{ github.actor }} + password: ${{ secrets.TOKEN }} + + - name: Build and push + uses: docker/build-push-action@v6 + with: + context: . + file: ./Dockerfile + push: true + build-args: | + VERSION=${{ steps.meta.outputs.TAG }} + tags: | + ${{ env.IMAGE }}:${{ steps.meta.outputs.TAG }} + ${{ env.IMAGE }}:latest + + - name: Notify Telegram + if: success() + run: | + curl -s -X POST "${{ secrets.TELEGRAM_BOT_URI }}/bot${{ secrets.TELEGRAM_BOT_TOKEN }}/sendMessage" \ + -d chat_id="${{ secrets.TELEGRAM_CHAT_ID }}" \ + -d parse_mode="HTML" \ + -d text="✅ traefik-selectel образ собран%0A%0AОбраз: ${{ env.IMAGE }}:${{ steps.meta.outputs.TAG }}" + + # Отдельным job: ловит падение любой из стадий. + # failure() сам по себе не годится: он истинен только когда упал прямой + # предок, поэтому результаты перечисляем явно. + notifyfailure: + name: Notify on failure + needs: [tests, build] + if: always() && (needs.tests.result == 'failure' || needs.build.result == 'failure') + runs-on: ubuntu-22.04 + container: catthehacker/ubuntu:act-latest + steps: + - name: Notify Telegram + run: | + curl -s -X POST "${{ secrets.TELEGRAM_BOT_URI }}/bot${{ secrets.TELEGRAM_BOT_TOKEN }}/sendMessage" \ + -d chat_id="${{ secrets.TELEGRAM_CHAT_ID }}" \ + -d parse_mode="HTML" \ + -d text="❌ traefik-selectel — тесты или сборка упали%0A%0AЛоги" diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..a08304b --- /dev/null +++ b/.gitignore @@ -0,0 +1,6 @@ +/traefik-selectel-dns +/secrets/ +deploy/secrets/ + +.env +.env.* diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..4af045a --- /dev/null +++ b/Dockerfile @@ -0,0 +1,22 @@ +# syntax=docker/dockerfile:1 +FROM golang:1.27-alpine AS build +WORKDIR /src +COPY go.mod go.sum ./ +RUN go mod download +COPY cmd ./cmd +COPY internal ./internal +ARG VERSION=dev +RUN CGO_ENABLED=0 GOOS=linux go build -trimpath \ + -ldflags="-s -w -X main.version=${VERSION}" \ + -o /out/traefik-selectel-dns ./cmd/traefik-selectel-dns + +# static: CA-сертификаты (для HTTPS к Selectel) + пользователь nonroot (uid 65532), без shell. +# Для воспроизводимости закрепите образ по digest. +FROM gcr.io/distroless/static-debian12:nonroot +COPY --from=build /out/traefik-selectel-dns /traefik-selectel-dns +USER 65532:65532 +EXPOSE 9090 +HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \ + CMD ["/traefik-selectel-dns", "-healthcheck"] +ENTRYPOINT ["/traefik-selectel-dns"] +CMD ["-config", "/etc/traefik-selectel-dns/config.yaml"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..8ed991c --- /dev/null +++ b/README.md @@ -0,0 +1,67 @@ +# traefik-selectel-dns + +Демон, который создаёт A-записи в Selectel DNS (API v2) для хостов из роутеров Traefik. +Сертификаты выписывает сам Traefik (встроенный `dnsChallenge.provider: selectelv2`) — кода для этого не нужно. + +## Почему не плагин Traefik + +Плагин (Yaegi) работает внутри цепочки обработки запросов: он не видит создание/удаление роутов +и не должен хранить токен с правом записи в DNS. Отдельный демон читает Traefik API +(`/api/http/routers`), поэтому одинаково покрывает Docker labels и file-provider. + +## Как работает + +1. Раз в `interval` (и сразу после событий Docker) читает `GET /api/http/routers`. +2. Из `rule` берёт все `Host(...)` у роутеров со статусом `enabled` (`HostRegexp`, `!Host(...)` игнорируются, `*.x` пропускаются с warning). +3. Хосты, попадающие в `zones` из конфига (равны зоне или её поддомен; при вложенных зонах — самая длинная), получают A-запись `ip`/`ttl`. +4. Записи, созданные демоном, помечаются `comment: managed-by=traefik-selectel-dns`. Меняются/удаляются только они. Существующая A-запись без метки (и имя с CNAME) не трогается — только warning. +5. Удаление выключено по умолчанию. При `delete.enabled: true` запись удаляется после того, как роут отсутствует непрерывно `delete.gracePeriod` (состояние в памяти; после рестарта отсчёт начинается заново). Роутеры со статусом `warning` считаются присутствующими. При ошибке Traefik API или пустом списке (после непустого) цикл пропускается и ничего не удаляется. + +## Настройка + +Пример полного стенда: `deploy/` (`docker-compose.example.yml`, `traefik.example.yml`, `config.example.yaml`, `dynamic/`). + +1. Создайте в Selectel сервисного пользователя с минимальной ролью на DNS проекта. +2. Положите креды в файлы секретов (`deploy/secrets/*`, не коммитить) — см. комментарий в compose. +3. Поправьте `zones`, `ip`, `traefik.apiURL` в `config.example.yaml`; wildcard `tls.domains` — в `traefik.example.yml`. +4. Для первого запуска включите `dryRun: true` и проверьте логи. + +Сборка: `docker build -t traefik-selectel-dns .` (distroless, non-root, HEALTHCHECK через `-healthcheck`). + +## Конфигурация + +YAML (`-config` или `TSD_CONFIG`), неизвестные поля — ошибка. Переопределение через env: + +| YAML | env | По умолчанию | +|---|---|---| +| `zones` | `TSD_ZONES` (через запятую) | обязательно | +| `ip` | `TSD_IP` | обязательно, IPv4 | +| `ttl` | `TSD_TTL` | 300 (60..604800) | +| `interval` | `TSD_INTERVAL` | 30s | +| `dryRun` | `TSD_DRY_RUN` | false | +| `traefik.apiURL` / `.timeout` | `TSD_TRAEFIK_API_URL` / `TSD_TRAEFIK_TIMEOUT` | http://traefik:8080 / 10s | +| `docker.enabled` / `.host` | `TSD_DOCKER_ENABLED` / `TSD_DOCKER_HOST` | true / unix:///var/run/docker.sock (`tcp://` тоже) | +| `delete.enabled` / `.gracePeriod` | `TSD_DELETE_ENABLED` / `TSD_DELETE_GRACE_PERIOD` | false / 10m | +| `selectel.baseURL` / `.authURL` / `.timeout` | `TSD_SELECTEL_BASE_URL` / `_AUTH_URL` / `_TIMEOUT` | api.selectel.ru/domains/v2, cloud.api.selcloud.ru/identity/v3, 30s | +| `listen` | `TSD_LISTEN` | :9090 (`/healthz`) | +| `logLevel` / `logFormat` | `TSD_LOG_LEVEL` / `TSD_LOG_FORMAT` | info / json | + +Секреты — только env (каждая поддерживает суффикс `_FILE`; одновременно обе — ошибка): +`SELECTEL_USERNAME`, `SELECTEL_PASSWORD`, `SELECTEL_ACCOUNT_ID`, `SELECTEL_PROJECT_ID`. +Для Traefik (lego) те же значения передаются как `SELECTELV2_USERNAME_FILE`, `SELECTELV2_PASSWORD_FILE`, `SELECTELV2_ACCOUNT_ID_FILE`, `SELECTELV2_PROJECT_ID_FILE`. + +## Безопасность + +- Секреты не попадают в конфиг, образ и логи; IAM-токен (Keystone) кэшируется в памяти и обновляется заранее. +- Traefik API не публикуется наружу и ограничен `ipAllowList` подсетью сети `control`; пользовательские контейнеры в неё не подключать. +- Docker socket = root на хосте. Демону нужны только события, поэтому в примере используется `docker-socket-proxy` (`EVENTS=1`), а не прямой сокет. Без Docker events демон работает по опросу. +- Чужие записи не перезаписываются; удаление выключено по умолчанию. + +## Ограничения + +- Только A-записи, один фиксированный IPv4 для всех хостов; AAAA нет. +- Только `Host(...)` с ASCII-именами (IDN — в punycode); `HostRegexp` и wildcard-хосты не поддерживаются. +- Метка владения — exact-match комментария; если её вручную изменить/стереть, запись считается чужой. +- Grace-состояние в памяти; несколько реплик демона одновременно не поддерживаются. +- Поток Docker events без общего таймаута (переподключается при обрыве; страховка — периодический опрос). +- Реальный Selectel API не проверялся (нет кредов) — см. отчёт в `swarm-report/`. diff --git a/cmd/traefik-selectel-dns/main.go b/cmd/traefik-selectel-dns/main.go new file mode 100644 index 0000000..3a13974 --- /dev/null +++ b/cmd/traefik-selectel-dns/main.go @@ -0,0 +1,135 @@ +// Команда traefik-selectel-dns: создаёт A-записи в Selectel DNS для хостов из роутеров Traefik. +package main + +import ( + "context" + "flag" + "fmt" + "log/slog" + "os" + "os/signal" + "syscall" + "time" + + "github.com/realmanual/traefik-selectel/internal/config" + "github.com/realmanual/traefik-selectel/internal/dockerevents" + "github.com/realmanual/traefik-selectel/internal/health" + "github.com/realmanual/traefik-selectel/internal/reconciler" + "github.com/realmanual/traefik-selectel/internal/selectel" + "github.com/realmanual/traefik-selectel/internal/traefik" +) + +var version = "dev" + +const dockerDebounce = 5 * time.Second + +func main() { + os.Exit(run()) +} + +func run() int { + configPath := flag.String("config", os.Getenv("TSD_CONFIG"), "путь к YAML-конфигу (или TSD_CONFIG)") + healthcheck := flag.Bool("healthcheck", false, "проверить /healthz локального экземпляра и выйти (для HEALTHCHECK)") + showVersion := flag.Bool("version", false, "показать версию") + flag.Parse() + + if *showVersion { + fmt.Println(version) + return 0 + } + + cfg, err := config.Load(*configPath, os.Getenv, config.OSReadFile) + if err != nil { + fmt.Fprintln(os.Stderr, "ошибка конфигурации:", err) + return 2 + } + if *healthcheck { + if err := health.Probe(cfg.Listen, 3*time.Second); err != nil { + fmt.Fprintln(os.Stderr, err) + return 1 + } + return 0 + } + + log := newLogger(cfg) + log.Info("запуск", "version", version, "zones", cfg.Zones, "ip", cfg.IP, "ttl", cfg.TTL, + "interval", cfg.Interval.String(), "deleteEnabled", cfg.Delete.Enabled, "dryRun", cfg.DryRun, + "traefikAPI", cfg.Traefik.APIURL, "dockerEvents", cfg.Docker.Enabled) + if cfg.DryRun { + log.Warn("включён dry-run: изменения в DNS не выполняются") + } + + tokens, err := selectel.NewIAMTokenSource(cfg.Credentials, cfg.Selectel.AuthURL, cfg.Selectel.Timeout) + if err != nil { + log.Error("инициализация Selectel", "error", err) + return 2 + } + dns := selectel.NewClient(cfg.Selectel.BaseURL, tokens, cfg.Selectel.Timeout) + tr, err := traefik.NewClient(cfg.Traefik.APIURL, cfg.Traefik.Timeout) + if err != nil { + log.Error("инициализация Traefik-клиента", "error", err) + return 2 + } + rec := reconciler.New(reconciler.Config{ + Zones: cfg.Zones, IP: cfg.IP, TTL: cfg.TTL, Comment: config.ManagedComment, + DryRun: cfg.DryRun, DeleteEnabled: cfg.Delete.Enabled, GracePeriod: cfg.Delete.GracePeriod, + }, tr, dns, log) + + ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT) + defer stop() + + state := health.NewState(3*cfg.Interval + reconciler.CycleTimeout) + healthErr := make(chan error, 1) + go func() { healthErr <- health.Serve(ctx, cfg.Listen, state.Handler()) }() + + trigger := make(chan struct{}, 1) + if cfg.Docker.Enabled { + w := dockerevents.New(cfg.Docker.Host, log, func() { + select { + case trigger <- struct{}{}: + default: + } + }) + go w.Run(ctx) + } + + loopDone := make(chan struct{}) + go func() { + defer close(loopDone) + rec.Run(ctx, cfg.Interval, dockerDebounce, trigger, func(_ reconciler.Result, err error) { state.Record(err) }) + }() + + select { + case <-ctx.Done(): + log.Info("получен сигнал завершения, останавливаемся") + case err := <-healthErr: + if err != nil { + log.Error("health-сервер остановился", "error", err) + stop() + <-loopDone + return 1 + } + } + <-loopDone + log.Info("остановлено") + return 0 +} + +func newLogger(cfg config.Config) *slog.Logger { + var lvl slog.Level + switch cfg.LogLevel { + case "debug": + lvl = slog.LevelDebug + case "warn": + lvl = slog.LevelWarn + case "error": + lvl = slog.LevelError + default: + lvl = slog.LevelInfo + } + opts := &slog.HandlerOptions{Level: lvl} + if cfg.LogFormat == "text" { + return slog.New(slog.NewTextHandler(os.Stderr, opts)) + } + return slog.New(slog.NewJSONHandler(os.Stderr, opts)) +} diff --git a/deploy/config.example.yaml b/deploy/config.example.yaml new file mode 100644 index 0000000..ab31d0b --- /dev/null +++ b/deploy/config.example.yaml @@ -0,0 +1,36 @@ +# Конфиг демона traefik-selectel-dns. Секреты Selectel здесь НЕ хранятся: +# SELECTEL_USERNAME / SELECTEL_PASSWORD / SELECTEL_ACCOUNT_ID / SELECTEL_PROJECT_ID +# задаются только через env (или *_FILE). Любой параметр можно переопределить +# env-переменной TSD_* (см. README). + +# Разрешённые зоны (opt-in по зоне): записи создаются только для хостов, +# равных зоне или являющихся её поддоменом. Зоны должны существовать в проекте Selectel. +zones: + - example.com + +# IPv4 публичного адреса ВМ с Traefik — значение всех создаваемых A-записей. +ip: 203.0.113.10 +ttl: 300 # 60..604800 +interval: 30s # период опроса Traefik API +dryRun: false # true — только логировать, ничего не менять в DNS + +traefik: + apiURL: http://traefik:8080 # внутренний адрес API, наружу не публиковать + timeout: 10s + +docker: + enabled: true # подписка на события контейнеров (триггер немедленного reconcile) + host: tcp://docker-proxy:2375 # или unix:///var/run/docker.sock; недоступность не фатальна + +delete: + enabled: false # по умолчанию записи не удаляются + gracePeriod: 10m # удалять только после непрерывного отсутствия роута + +# selectel: # обычно не нужно +# baseURL: https://api.selectel.ru/domains/v2 +# authURL: https://cloud.api.selcloud.ru/identity/v3 +# timeout: 30s + +listen: ":9090" # /healthz +logLevel: info # debug|info|warn|error +logFormat: json # json|text diff --git a/deploy/docker-compose.example.yml b/deploy/docker-compose.example.yml new file mode 100644 index 0000000..c5233f8 --- /dev/null +++ b/deploy/docker-compose.example.yml @@ -0,0 +1,97 @@ +# Пример: Traefik + traefik-selectel-dns на одной ВМ. +# Перед запуском создайте файлы секретов (права 0400/0440, не в git): +# mkdir -p secrets && umask 077 +# printf '%s' 'service-user' > secrets/selectel_username +# printf '%s' 'password' > secrets/selectel_password +# printf '%s' '123456' > secrets/selectel_account_id +# printf '%s' 'project-uuid' > secrets/selectel_project_id +# Пользователь Selectel — сервисный, с минимальной ролью на DNS в нужном проекте. + +services: + traefik: + image: traefik:v3 # в проде закрепите конкретную версию/digest + restart: unless-stopped + command: ["--configFile=/etc/traefik/traefik.yml"] + ports: # публикуются ТОЛЬКО 80/443; 8080 (API) не публикуется + - "80:80" + - "443:443" + environment: + # Имена переменных lego (selectelv2); суффикс _FILE — чтение значения из файла + SELECTELV2_USERNAME_FILE: /run/secrets/selectel_username + SELECTELV2_PASSWORD_FILE: /run/secrets/selectel_password + SELECTELV2_ACCOUNT_ID_FILE: /run/secrets/selectel_account_id + SELECTELV2_PROJECT_ID_FILE: /run/secrets/selectel_project_id + secrets: + - selectel_username + - selectel_password + - selectel_account_id + - selectel_project_id + volumes: + - ./traefik.example.yml:/etc/traefik/traefik.yml:ro + - ./dynamic:/etc/traefik/dynamic:ro + - letsencrypt:/letsencrypt + - /var/run/docker.sock:/var/run/docker.sock:ro + networks: + - proxy # сюда подключаются пользовательские контейнеры + - control # только Traefik <-> демон + + docker-proxy: + # Демону нужен только поток событий: даём ему read-only прокси вместо docker.sock. + image: tecnativa/docker-socket-proxy:latest # закрепите версию + restart: unless-stopped + environment: + EVENTS: 1 + CONTAINERS: 0 + POST: 0 + volumes: + - /var/run/docker.sock:/var/run/docker.sock:ro + networks: + - control + + traefik-selectel-dns: + image: git.realmanual.ru/pub/traefik-selectel:latest # или build: .. + # build: + # context: .. + restart: unless-stopped + depends_on: + - traefik + environment: + SELECTEL_USERNAME_FILE: /run/secrets/selectel_username + SELECTEL_PASSWORD_FILE: /run/secrets/selectel_password + SELECTEL_ACCOUNT_ID_FILE: /run/secrets/selectel_account_id + SELECTEL_PROJECT_ID_FILE: /run/secrets/selectel_project_id + secrets: + - selectel_username + - selectel_password + - selectel_account_id + - selectel_project_id + volumes: + - ./config.example.yaml:/etc/traefik-selectel-dns/config.yaml:ro + read_only: true + cap_drop: [ALL] + security_opt: + - no-new-privileges:true + networks: + - control # доступ к Traefik API и docker-proxy; исходящий интернет — к Selectel + +secrets: + selectel_username: + file: ./secrets/selectel_username + selectel_password: + file: ./secrets/selectel_password + selectel_account_id: + file: ./secrets/selectel_account_id + selectel_project_id: + file: ./secrets/selectel_project_id + +volumes: + letsencrypt: + +networks: + proxy: + name: proxy + control: + name: traefik-control # НЕ подключайте к ней пользовательские контейнеры + ipam: + config: + - subnet: 172.30.250.0/24 # должна совпадать с sourceRange в dynamic/api-internal.yml diff --git a/deploy/dynamic/api-internal.yml b/deploy/dynamic/api-internal.yml new file mode 100644 index 0000000..b8810f4 --- /dev/null +++ b/deploy/dynamic/api-internal.yml @@ -0,0 +1,14 @@ +# Traefik API только для демона: entrypoint "traefik" (:8080) + разрешён только источник +# из подсети сети control. Пользовательские контейнеры (сеть proxy) получат 403. +http: + routers: + api-internal: + rule: PathPrefix(`/api`) + entryPoints: [traefik] + service: api@internal + middlewares: [api-internal-allow] + middlewares: + api-internal-allow: + ipAllowList: + sourceRange: + - 172.30.250.0/24 diff --git a/deploy/dynamic/example.yml b/deploy/dynamic/example.yml new file mode 100644 index 0000000..95f98a4 --- /dev/null +++ b/deploy/dynamic/example.yml @@ -0,0 +1,13 @@ +# Пример роута через file-provider: демон увидит его в Traefik API так же, как docker-роуты. +http: + routers: + legacy-app: + rule: Host(`legacy.example.com`) + entryPoints: [websecure] + service: legacy-app + tls: {} # используется wildcard-сертификат с entrypoint websecure + services: + legacy-app: + loadBalancer: + servers: + - url: http://10.0.0.5:8080 diff --git a/deploy/traefik.example.yml b/deploy/traefik.example.yml new file mode 100644 index 0000000..1d1986e --- /dev/null +++ b/deploy/traefik.example.yml @@ -0,0 +1,56 @@ +# Статический конфиг Traefik v3. +# Креды Selectel для dns-challenge (lego, провайдер selectelv2) передаются через env +# SELECTELV2_*_FILE (см. docker-compose.example.yml), в этом файле секретов нет. + +entryPoints: + web: + address: ":80" + http: + redirections: + entryPoint: + to: websecure + scheme: https + websecure: + address: ":443" + http: + tls: + certResolver: letsencrypt + # Общий wildcard-сертификат: Traefik выпустит его при старте, роуты с tls: {} + # на этом entrypoint будут использовать его без отдельных запросов в ACME. + domains: + - main: example.com + sans: + - "*.example.com" + # Внутренний entrypoint для API. Порт НЕ публикуется (ports:), а роутер api@internal + # ограничен по IP подсетью сети control (dynamic/api-internal.yml), т.к. Traefik также + # подключён к сети proxy с пользовательскими контейнерами. + traefik: + address: ":8080" + +api: + insecure: false # не открываем API автоматически — только через ограниченный роутер + dashboard: false # демону нужен только /api/http/routers + +certificatesResolvers: + letsencrypt: + acme: + email: admin@example.com + storage: /letsencrypt/acme.json + # Для отладки используйте staging: + # caServer: https://acme-staging-v02.api.letsencrypt.org/directory + dnsChallenge: + provider: selectelv2 + resolvers: + - "1.1.1.1:53" + - "8.8.8.8:53" + +providers: + docker: + exposedByDefault: false + network: proxy + file: + directory: /etc/traefik/dynamic + watch: true + +log: + level: INFO diff --git a/deploy/whoami.example.yml b/deploy/whoami.example.yml new file mode 100644 index 0000000..21d19ed --- /dev/null +++ b/deploy/whoami.example.yml @@ -0,0 +1,17 @@ +# Пример пользовательского контейнера: достаточно обычных labels Traefik, +# DNS-запись для whoami.example.com демон создаст сам (зона example.com разрешена в конфиге). +services: + whoami: + image: traefik/whoami + restart: unless-stopped + labels: + traefik.enable: "true" + traefik.http.routers.whoami.rule: Host(`whoami.example.com`) + traefik.http.routers.whoami.entrypoints: websecure + traefik.http.routers.whoami.tls: "true" + networks: + - proxy +networks: + proxy: + external: true + name: proxy diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..4a80999 --- /dev/null +++ b/go.mod @@ -0,0 +1,5 @@ +module github.com/realmanual/traefik-selectel + +go 1.27.1 + +require gopkg.in/yaml.v3 v3.0.1 diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..a62c313 --- /dev/null +++ b/go.sum @@ -0,0 +1,4 @@ +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..baaf826 --- /dev/null +++ b/internal/config/config.go @@ -0,0 +1,286 @@ +// Package config загружает конфигурацию демона: YAML-файл + переопределения из env. +// Секреты Selectel читаются только из env (с поддержкой суффикса _FILE). +package config + +import ( + "bytes" + "errors" + "fmt" + "io" + "net" + "net/netip" + "os" + "strconv" + "strings" + "time" + + "gopkg.in/yaml.v3" + + "github.com/realmanual/traefik-selectel/internal/hostparse" + "github.com/realmanual/traefik-selectel/internal/selectel" +) + +// ManagedComment — метка владения записями. +const ManagedComment = "managed-by=traefik-selectel-dns" + +const ( + MinTTL = 60 + MaxTTL = 604800 +) + +// Config — конфигурация демона. +type Config struct { + Zones []string `yaml:"zones"` + IP string `yaml:"ip"` + TTL int `yaml:"ttl"` + Interval time.Duration `yaml:"interval"` + DryRun bool `yaml:"dryRun"` + + Traefik TraefikConfig `yaml:"traefik"` + Docker DockerConfig `yaml:"docker"` + Delete DeleteConfig `yaml:"delete"` + Selectel SelectelConfig `yaml:"selectel"` + + Listen string `yaml:"listen"` + LogLevel string `yaml:"logLevel"` + LogFormat string `yaml:"logFormat"` + + // Credentials заполняется только из env, не из YAML. + Credentials selectel.Credentials `yaml:"-"` +} + +type TraefikConfig struct { + APIURL string `yaml:"apiURL"` + Timeout time.Duration `yaml:"timeout"` +} + +type DockerConfig struct { + Enabled bool `yaml:"enabled"` + Host string `yaml:"host"` +} + +type DeleteConfig struct { + Enabled bool `yaml:"enabled"` + GracePeriod time.Duration `yaml:"gracePeriod"` +} + +// SelectelConfig — несекретные настройки API. +type SelectelConfig struct { + BaseURL string `yaml:"baseURL"` + AuthURL string `yaml:"authURL"` + Timeout time.Duration `yaml:"timeout"` +} + +// Default возвращает конфигурацию со значениями по умолчанию. +func Default() Config { + return Config{ + TTL: 300, + Interval: 30 * time.Second, + Traefik: TraefikConfig{APIURL: "http://traefik:8080", Timeout: 10 * time.Second}, + Docker: DockerConfig{Enabled: true, Host: "unix:///var/run/docker.sock"}, + Delete: DeleteConfig{Enabled: false, GracePeriod: 10 * time.Minute}, + Selectel: SelectelConfig{Timeout: 30 * time.Second}, + Listen: ":9090", LogLevel: "info", LogFormat: "json", + } +} + +// Load читает YAML (path может быть пустым), применяет env-переопределения и секреты, валидирует. +func Load(path string, getenv func(string) string, readFile func(string) ([]byte, error)) (Config, error) { + cfg := Default() + if path != "" { + data, err := readFile(path) + if err != nil { + return cfg, fmt.Errorf("config: чтение %s: %w", path, err) + } + dec := yaml.NewDecoder(bytes.NewReader(data)) + dec.KnownFields(true) + if err := dec.Decode(&cfg); err != nil && !errors.Is(err, io.EOF) { + return cfg, fmt.Errorf("config: разбор %s: %w", path, err) + } + } + if err := applyEnv(&cfg, getenv); err != nil { + return cfg, err + } + creds, err := LoadCredentials(getenv, readFile) + if err != nil { + return cfg, err + } + cfg.Credentials = creds + if err := cfg.Normalize(); err != nil { + return cfg, err + } + return cfg, cfg.Validate() +} + +// LoadCredentials читает SELECTEL_* из env / *_FILE. +func LoadCredentials(getenv func(string) string, readFile func(string) ([]byte, error)) (selectel.Credentials, error) { + var c selectel.Credentials + var err error + for _, f := range []struct { + name string + dst *string + }{ + {"SELECTEL_USERNAME", &c.Username}, + {"SELECTEL_PASSWORD", &c.Password}, + {"SELECTEL_ACCOUNT_ID", &c.AccountID}, + {"SELECTEL_PROJECT_ID", &c.ProjectID}, + } { + if *f.dst, err = Secret(f.name, getenv, readFile); err != nil { + return c, err + } + } + return c, nil +} + +// Secret возвращает значение NAME либо содержимое файла из NAME_FILE. +// Одновременное задание обоих — ошибка. +func Secret(name string, getenv func(string) string, readFile func(string) ([]byte, error)) (string, error) { + direct, file := getenv(name), getenv(name+"_FILE") + if direct != "" && file != "" { + return "", fmt.Errorf("config: заданы и %s, и %s_FILE — оставьте что-то одно", name, name) + } + if file != "" { + b, err := readFile(file) + if err != nil { + return "", fmt.Errorf("config: чтение %s_FILE: %w", name, err) + } + return strings.TrimSpace(string(b)), nil + } + return strings.TrimSpace(direct), nil +} + +func applyEnv(c *Config, getenv func(string) string) error { + var errs []error + str := func(key string, dst *string) { + if v := getenv(key); v != "" { + *dst = v + } + } + boolean := func(key string, dst *bool) { + if v := getenv(key); v != "" { + b, err := strconv.ParseBool(v) + if err != nil { + errs = append(errs, fmt.Errorf("config: %s: %w", key, err)) + return + } + *dst = b + } + } + dur := func(key string, dst *time.Duration) { + if v := getenv(key); v != "" { + d, err := time.ParseDuration(v) + if err != nil { + errs = append(errs, fmt.Errorf("config: %s: %w", key, err)) + return + } + *dst = d + } + } + if v := getenv("TSD_ZONES"); v != "" { + c.Zones = nil + for _, z := range strings.Split(v, ",") { + c.Zones = append(c.Zones, z) + } + } + str("TSD_IP", &c.IP) + if v := getenv("TSD_TTL"); v != "" { + n, err := strconv.Atoi(v) + if err != nil { + errs = append(errs, fmt.Errorf("config: TSD_TTL: %w", err)) + } else { + c.TTL = n + } + } + dur("TSD_INTERVAL", &c.Interval) + boolean("TSD_DRY_RUN", &c.DryRun) + str("TSD_TRAEFIK_API_URL", &c.Traefik.APIURL) + dur("TSD_TRAEFIK_TIMEOUT", &c.Traefik.Timeout) + boolean("TSD_DOCKER_ENABLED", &c.Docker.Enabled) + str("TSD_DOCKER_HOST", &c.Docker.Host) + boolean("TSD_DELETE_ENABLED", &c.Delete.Enabled) + dur("TSD_DELETE_GRACE_PERIOD", &c.Delete.GracePeriod) + str("TSD_SELECTEL_BASE_URL", &c.Selectel.BaseURL) + str("TSD_SELECTEL_AUTH_URL", &c.Selectel.AuthURL) + dur("TSD_SELECTEL_TIMEOUT", &c.Selectel.Timeout) + str("TSD_LISTEN", &c.Listen) + str("TSD_LOG_LEVEL", &c.LogLevel) + str("TSD_LOG_FORMAT", &c.LogFormat) + return errors.Join(errs...) +} + +// Normalize приводит зоны к нормализованному виду и убирает дубликаты. +func (c *Config) Normalize() error { + seen := map[string]struct{}{} + var zones []string + for _, z := range c.Zones { + n := hostparse.Normalize(z) + if n == "" { + continue + } + if _, ok := seen[n]; ok { + continue + } + seen[n] = struct{}{} + zones = append(zones, n) + } + c.Zones = zones + c.IP = strings.TrimSpace(c.IP) + c.Traefik.APIURL = strings.TrimRight(strings.TrimSpace(c.Traefik.APIURL), "/") + c.LogLevel = strings.ToLower(c.LogLevel) + c.LogFormat = strings.ToLower(c.LogFormat) + return nil +} + +// Validate проверяет конфигурацию. +func (c *Config) Validate() error { + var errs []error + if len(c.Zones) == 0 { + errs = append(errs, errors.New("zones: нужна хотя бы одна зона")) + } + for _, z := range c.Zones { + if !hostparse.ValidHostname(z) { + errs = append(errs, fmt.Errorf("zones: некорректная зона %q", z)) + } + } + if addr, err := netip.ParseAddr(c.IP); err != nil || !addr.Is4() { + errs = append(errs, fmt.Errorf("ip: %q не является IPv4-адресом", c.IP)) + } + if c.TTL < MinTTL || c.TTL > MaxTTL { + errs = append(errs, fmt.Errorf("ttl: %d вне диапазона %d..%d", c.TTL, MinTTL, MaxTTL)) + } + if c.Interval < time.Second { + errs = append(errs, fmt.Errorf("interval: %s слишком мал (минимум 1s)", c.Interval)) + } + if c.Traefik.APIURL == "" { + errs = append(errs, errors.New("traefik.apiURL: не задан")) + } + if c.Traefik.Timeout <= 0 || c.Selectel.Timeout <= 0 { + errs = append(errs, errors.New("timeout: должен быть положительным")) + } + if c.Delete.Enabled && c.Delete.GracePeriod < c.Interval { + errs = append(errs, fmt.Errorf("delete.gracePeriod: %s меньше interval %s", c.Delete.GracePeriod, c.Interval)) + } + if c.Docker.Enabled && c.Docker.Host == "" { + errs = append(errs, errors.New("docker.host: не задан")) + } + if _, _, err := net.SplitHostPort(c.Listen); err != nil { + errs = append(errs, fmt.Errorf("listen: %q: %w", c.Listen, err)) + } + switch c.LogLevel { + case "debug", "info", "warn", "error": + default: + errs = append(errs, fmt.Errorf("logLevel: %q (debug|info|warn|error)", c.LogLevel)) + } + switch c.LogFormat { + case "json", "text": + default: + errs = append(errs, fmt.Errorf("logFormat: %q (json|text)", c.LogFormat)) + } + if err := c.Credentials.Validate(); err != nil { + errs = append(errs, err) + } + return errors.Join(errs...) +} + +// OSReadFile — os.ReadFile (для main). +var OSReadFile = os.ReadFile diff --git a/internal/config/config_test.go b/internal/config/config_test.go new file mode 100644 index 0000000..fd99df5 --- /dev/null +++ b/internal/config/config_test.go @@ -0,0 +1,140 @@ +package config + +import ( + "fmt" + "strings" + "testing" + "time" +) + +func env(m map[string]string) func(string) string { return func(k string) string { return m[k] } } + +func files(m map[string]string) func(string) ([]byte, error) { + return func(p string) ([]byte, error) { + if v, ok := m[p]; ok { + return []byte(v), nil + } + return nil, fmt.Errorf("no such file %s", p) + } +} + +var goodEnv = map[string]string{ + "SELECTEL_USERNAME": "u", "SELECTEL_PASSWORD": "p", "SELECTEL_ACCOUNT_ID": "1", "SELECTEL_PROJECT_ID": "2", +} + +func withEnv(extra map[string]string) map[string]string { + m := map[string]string{} + for k, v := range goodEnv { + m[k] = v + } + for k, v := range extra { + m[k] = v + } + return m +} + +func TestLoadYAMLAndDefaults(t *testing.T) { + yml := "zones: [Example.com., example.com, dev.example.org]\nip: 1.2.3.4\ndelete:\n enabled: true\n gracePeriod: 15m\n" + cfg, err := Load("/cfg.yaml", env(goodEnv), files(map[string]string{"/cfg.yaml": yml})) + if err != nil { + t.Fatal(err) + } + if len(cfg.Zones) != 2 || cfg.Zones[0] != "example.com" { + t.Errorf("zones = %v", cfg.Zones) + } + if cfg.TTL != 300 || cfg.Interval != 30*time.Second || cfg.Listen != ":9090" || cfg.Traefik.APIURL != "http://traefik:8080" { + t.Errorf("defaults: %+v", cfg) + } + if !cfg.Delete.Enabled || cfg.Delete.GracePeriod != 15*time.Minute { + t.Errorf("delete: %+v", cfg.Delete) + } +} + +func TestDeleteDisabledByDefault(t *testing.T) { + cfg, err := Load("", env(withEnv(map[string]string{"TSD_ZONES": "example.com", "TSD_IP": "1.2.3.4"})), files(nil)) + if err != nil { + t.Fatal(err) + } + if cfg.Delete.Enabled || cfg.Delete.GracePeriod != 10*time.Minute { + t.Errorf("delete: %+v", cfg.Delete) + } +} + +func TestEnvOverrides(t *testing.T) { + yml := "zones: [example.com]\nip: 1.2.3.4\nttl: 300\n" + cfg, err := Load("/c.yaml", env(withEnv(map[string]string{ + "TSD_IP": "5.6.7.8", "TSD_TTL": "120", "TSD_ZONES": "a.com, b.com", "TSD_DRY_RUN": "true", "TSD_INTERVAL": "1m", + })), files(map[string]string{"/c.yaml": yml})) + if err != nil { + t.Fatal(err) + } + if cfg.IP != "5.6.7.8" || cfg.TTL != 120 || len(cfg.Zones) != 2 || cfg.Zones[1] != "b.com" || !cfg.DryRun || cfg.Interval != time.Minute { + t.Errorf("cfg = %+v", cfg) + } +} + +func TestUnknownYAMLFieldRejected(t *testing.T) { + _, err := Load("/c.yaml", env(goodEnv), files(map[string]string{"/c.yaml": "zones: [a.com]\nip: 1.1.1.1\npassword: leak\n"})) + if err == nil { + t.Fatal("ожидалась ошибка на неизвестное поле") + } +} + +func TestValidation(t *testing.T) { + base := func() map[string]string { + return withEnv(map[string]string{"TSD_ZONES": "example.com", "TSD_IP": "1.2.3.4"}) + } + cases := map[string]map[string]string{ + "ip не IPv4": {"TSD_IP": "::1"}, + "ip мусор": {"TSD_IP": "example"}, + "ttl мал": {"TSD_TTL": "59"}, + "ttl велик": {"TSD_TTL": "604801"}, + "нет зон": {"TSD_ZONES": " "}, + "grace < interval": {"TSD_DELETE_ENABLED": "true", "TSD_DELETE_GRACE_PERIOD": "1s"}, + "bad loglevel": {"TSD_LOG_LEVEL": "trace"}, + "bad bool": {"TSD_DRY_RUN": "maybe"}, + "нет пароля": {"SELECTEL_PASSWORD": ""}, + } + for name, extra := range cases { + t.Run(name, func(t *testing.T) { + m := base() + for k, v := range extra { + m[k] = v + } + if _, err := Load("", env(m), files(nil)); err == nil { + t.Fatal("ожидалась ошибка") + } + }) + } + if _, err := Load("", env(base()), files(nil)); err != nil { + t.Fatalf("базовый конфиг должен быть валиден: %v", err) + } +} + +func TestSecretFileSupport(t *testing.T) { + e := map[string]string{ + "SELECTEL_USERNAME_FILE": "/run/secrets/u", "SELECTEL_PASSWORD_FILE": "/run/secrets/p", + "SELECTEL_ACCOUNT_ID": "1", "SELECTEL_PROJECT_ID_FILE": "/run/secrets/pid", + } + f := files(map[string]string{"/run/secrets/u": "user\n", "/run/secrets/p": " pass \n", "/run/secrets/pid": "pid-1\n"}) + c, err := LoadCredentials(env(e), f) + if err != nil { + t.Fatal(err) + } + if c.Username != "user" || c.Password != "pass" || c.ProjectID != "pid-1" || c.AccountID != "1" { + t.Errorf("creds = %s", c) + } +} + +func TestSecretBothSetIsError(t *testing.T) { + _, err := Secret("X", env(map[string]string{"X": "a", "X_FILE": "/f"}), files(map[string]string{"/f": "b"})) + if err == nil || !strings.Contains(err.Error(), "одно") { + t.Fatalf("err = %v", err) + } +} + +func TestSecretMissingFile(t *testing.T) { + if _, err := Secret("X", env(map[string]string{"X_FILE": "/nope"}), files(nil)); err == nil { + t.Fatal("ожидалась ошибка") + } +} diff --git a/internal/dockerevents/watch.go b/internal/dockerevents/watch.go new file mode 100644 index 0000000..499384f --- /dev/null +++ b/internal/dockerevents/watch.go @@ -0,0 +1,134 @@ +// Package dockerevents подписывается на события контейнеров Docker (raw HTTP, +// без зависимости от docker SDK) и вызывает trigger как сигнал к немедленному reconcile. +// События — только триггер; источник истины — Traefik API. +package dockerevents + +import ( + "bufio" + "context" + "encoding/json" + "fmt" + "log/slog" + "net" + "net/http" + "net/url" + "strings" + "time" +) + +const eventsFilters = `{"type":["container"],"event":["start","die","destroy"]}` + +// Watcher читает поток /events. +type Watcher struct { + host string + log *slog.Logger + trigger func() + + // для тестов + minBackoff, maxBackoff time.Duration +} + +// New создаёт Watcher. host: unix:///path/to.sock или tcp://host:port. +func New(host string, log *slog.Logger, trigger func()) *Watcher { + return &Watcher{host: host, log: log, trigger: trigger, minBackoff: time.Second, maxBackoff: time.Minute} +} + +func (w *Watcher) client() (*http.Client, string, error) { + u, err := url.Parse(w.host) + if err != nil { + return nil, "", fmt.Errorf("docker.host: %w", err) + } + dialer := &net.Dialer{Timeout: 5 * time.Second} + tr := &http.Transport{ResponseHeaderTimeout: 10 * time.Second} + switch u.Scheme { + case "unix": + path := u.Path + if path == "" { + return nil, "", fmt.Errorf("docker.host: пустой путь сокета") + } + tr.DialContext = func(ctx context.Context, _, _ string) (net.Conn, error) { + return dialer.DialContext(ctx, "unix", path) + } + return &http.Client{Transport: tr}, "http://docker", nil + case "tcp", "http": + tr.DialContext = dialer.DialContext + return &http.Client{Transport: tr}, "http://" + u.Host, nil + default: + return nil, "", fmt.Errorf("docker.host: неподдерживаемая схема %q (unix|tcp)", u.Scheme) + } +} + +// Run блокируется до отмены ctx, переподключаясь при обрывах. Ошибки не фатальны. +func (w *Watcher) Run(ctx context.Context) { + hc, base, err := w.client() + if err != nil { + w.log.Error("Docker events отключены", "error", err) + return + } + backoff := w.minBackoff + failing := false + for ctx.Err() == nil { + connected, err := w.stream(ctx, hc, base) + if ctx.Err() != nil { + return + } + if connected { + backoff = w.minBackoff + failing = false + } + if !failing { + w.log.Warn("Docker events недоступны, работаем по опросу; повторные попытки в фоне", "error", err) + failing = true + } else { + w.log.Debug("Docker events: повторная попытка не удалась", "error", err) + } + select { + case <-ctx.Done(): + return + case <-time.After(backoff): + } + if backoff *= 2; backoff > w.maxBackoff { + backoff = w.maxBackoff + } + } +} + +// stream возвращает connected=true, если соединение было установлено (ответ 200). +func (w *Watcher) stream(ctx context.Context, hc *http.Client, base string) (bool, error) { + q := url.Values{"filters": {eventsFilters}} + req, err := http.NewRequestWithContext(ctx, http.MethodGet, base+"/events?"+q.Encode(), nil) + if err != nil { + return false, err + } + resp, err := hc.Do(req) + if err != nil { + return false, err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return false, fmt.Errorf("docker events: статус %d", resp.StatusCode) + } + w.log.Info("подписка на Docker events установлена") + sc := bufio.NewScanner(resp.Body) + sc.Buffer(make([]byte, 64*1024), 1<<20) + for sc.Scan() { + line := strings.TrimSpace(sc.Text()) + if line == "" { + continue + } + var ev struct { + Type string `json:"Type"` + Action string `json:"Action"` + } + if err := json.Unmarshal([]byte(line), &ev); err != nil { + w.log.Debug("Docker events: не удалось разобрать событие", "error", err) + continue + } + w.log.Debug("событие Docker", "type", ev.Type, "action", ev.Action) + w.trigger() + } + if err := sc.Err(); err != nil { + return true, err + } + return true, fmt.Errorf("docker events: поток закрыт") +} diff --git a/internal/dockerevents/watch_test.go b/internal/dockerevents/watch_test.go new file mode 100644 index 0000000..7e34cc1 --- /dev/null +++ b/internal/dockerevents/watch_test.go @@ -0,0 +1,105 @@ +package dockerevents + +import ( + "context" + "io" + "log/slog" + "net" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "sync/atomic" + "testing" + "time" +) + +func quiet() *slog.Logger { return slog.New(slog.NewTextHandler(io.Discard, nil)) } + +func waitFor(t *testing.T, cond func() bool) { + t.Helper() + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + if cond() { + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatal("условие не выполнено за отведённое время") +} + +func TestTriggerOnEventsTCP(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/events" || !strings.Contains(r.URL.Query().Get("filters"), "container") { + t.Errorf("url = %s", r.URL) + } + _, _ = w.Write([]byte(`{"Type":"container","Action":"start"}` + "\n" + `{"Type":"container","Action":"die"}` + "\n")) + w.(http.Flusher).Flush() + <-r.Context().Done() + })) + defer srv.Close() + + var n int32 + w := New("tcp://"+strings.TrimPrefix(srv.URL, "http://"), quiet(), func() { atomic.AddInt32(&n, 1) }) + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { w.Run(ctx); close(done) }() + waitFor(t, func() bool { return atomic.LoadInt32(&n) >= 2 }) + cancel() + <-done +} + +func TestUnixSocketAndReconnect(t *testing.T) { + sock := filepath.Join(t.TempDir(), "d.sock") + l, err := net.Listen("unix", sock) + if err != nil { + t.Skipf("unix сокеты недоступны: %v", err) + } + var conns int32 + srv := &http.Server{Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt32(&conns, 1) + _, _ = w.Write([]byte(`{"Type":"container","Action":"start"}` + "\n")) + // поток закрывается сразу — ждём переподключения + })} + go srv.Serve(l) + defer srv.Close() + + var n int32 + w := New("unix://"+sock, quiet(), func() { atomic.AddInt32(&n, 1) }) + w.minBackoff, w.maxBackoff = 10*time.Millisecond, 20*time.Millisecond + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { w.Run(ctx); close(done) }() + waitFor(t, func() bool { return atomic.LoadInt32(&conns) >= 2 && atomic.LoadInt32(&n) >= 2 }) + cancel() + <-done +} + +func TestUnavailableSocketDoesNotPanicOrExit(t *testing.T) { + w := New("unix:///nonexistent/docker.sock", quiet(), func() { t.Error("trigger не ожидался") }) + w.minBackoff, w.maxBackoff = 5*time.Millisecond, 10*time.Millisecond + ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) + defer cancel() + done := make(chan struct{}) + go func() { w.Run(ctx); close(done) }() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("Run не завершился после отмены контекста") + } +} + +func TestInvalidHostScheme(t *testing.T) { + w := New("ftp://x", quiet(), func() {}) + if _, _, err := w.client(); err == nil { + t.Fatal("ожидалась ошибка") + } + // Run на невалидном host возвращается сразу + done := make(chan struct{}) + go func() { w.Run(context.Background()); close(done) }() + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("Run должен вернуться при невалидном host") + } +} diff --git a/internal/health/health.go b/internal/health/health.go new file mode 100644 index 0000000..00a400d --- /dev/null +++ b/internal/health/health.go @@ -0,0 +1,122 @@ +// Package health — HTTP-endpoint /healthz. +package health + +import ( + "context" + "encoding/json" + "fmt" + "net" + "net/http" + "sync" + "time" +) + +// State хранит состояние цикла reconcile. +type State struct { + maxAge time.Duration + now func() time.Time + + mu sync.Mutex + started time.Time + lastAttempt time.Time + lastSuccess time.Time + lastError string +} + +// NewState: сервис считается нездоровым, если цикл не отрабатывал дольше maxAge. +func NewState(maxAge time.Duration) *State { + s := &State{maxAge: maxAge, now: time.Now} + s.started = s.now() + return s +} + +// Record фиксирует результат цикла. +func (s *State) Record(err error) { + s.mu.Lock() + defer s.mu.Unlock() + s.lastAttempt = s.now() + if err == nil { + s.lastSuccess = s.lastAttempt + s.lastError = "" + return + } + s.lastError = err.Error() +} + +type status struct { + Status string `json:"status"` + LastAttempt string `json:"lastAttempt,omitempty"` + LastSuccess string `json:"lastSuccess,omitempty"` + LastError string `json:"lastError,omitempty"` +} + +// Handler отдаёт 200, если цикл жив (попытки идут), иначе 503. +// Ошибки внешних систем (Traefik/Selectel) не делают сервис unhealthy — +// рестарт контейнера их не исправит; они видны в lastError. +func (s *State) Handler() http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + s.mu.Lock() + ref := s.lastAttempt + if ref.IsZero() { + ref = s.started + } + ok := s.now().Sub(ref) <= s.maxAge + st := status{Status: "ok", LastError: s.lastError} + if !s.lastAttempt.IsZero() { + st.LastAttempt = s.lastAttempt.UTC().Format(time.RFC3339) + } + if !s.lastSuccess.IsZero() { + st.LastSuccess = s.lastSuccess.UTC().Format(time.RFC3339) + } + s.mu.Unlock() + code := http.StatusOK + if !ok { + st.Status, code = "stalled", http.StatusServiceUnavailable + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(code) + _ = json.NewEncoder(w).Encode(st) + }) +} + +// Serve запускает HTTP-сервер и завершает его при отмене ctx. +func Serve(ctx context.Context, addr string, h http.Handler) error { + mux := http.NewServeMux() + mux.Handle("/healthz", h) + srv := &http.Server{ + Addr: addr, Handler: mux, + ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 10 * time.Second, + WriteTimeout: 10 * time.Second, IdleTimeout: 30 * time.Second, + } + errc := make(chan error, 1) + go func() { errc <- srv.ListenAndServe() }() + select { + case err := <-errc: + return err + case <-ctx.Done(): + sctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + return srv.Shutdown(sctx) + } +} + +// Probe выполняет GET /healthz на локальном адресе listen (для флага -healthcheck). +func Probe(listen string, timeout time.Duration) error { + host, port, err := net.SplitHostPort(listen) + if err != nil { + return err + } + if host == "" || host == "0.0.0.0" || host == "::" { + host = "127.0.0.1" + } + c := &http.Client{Timeout: timeout} + resp, err := c.Get("http://" + net.JoinHostPort(host, port) + "/healthz") + if err != nil { + return err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("healthz: статус %d", resp.StatusCode) + } + return nil +} diff --git a/internal/health/health_test.go b/internal/health/health_test.go new file mode 100644 index 0000000..13c71f0 --- /dev/null +++ b/internal/health/health_test.go @@ -0,0 +1,40 @@ +package health + +import ( + "errors" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" +) + +func TestHandler(t *testing.T) { + now := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + s := NewState(time.Minute) + s.now = func() time.Time { return now } + s.started = now + + get := func() (int, string) { + rec := httptest.NewRecorder() + s.Handler().ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/healthz", nil)) + return rec.Code, rec.Body.String() + } + + if code, _ := get(); code != 200 { + t.Fatalf("сразу после старта: %d", code) + } + now = now.Add(2 * time.Minute) + if code, body := get(); code != 503 || !strings.Contains(body, "stalled") { + t.Fatalf("застрял: %d %s", code, body) + } + s.Record(errors.New("traefik down")) + code, body := get() + if code != 200 || !strings.Contains(body, "traefik down") { + t.Fatalf("ошибка внешней системы не должна делать unhealthy: %d %s", code, body) + } + s.Record(nil) + if _, body := get(); !strings.Contains(body, "lastSuccess") || strings.Contains(body, "lastError") { + t.Fatalf("body = %s", body) + } +} diff --git a/internal/hostparse/hostparse.go b/internal/hostparse/hostparse.go new file mode 100644 index 0000000..0a04972 --- /dev/null +++ b/internal/hostparse/hostparse.go @@ -0,0 +1,212 @@ +// Package hostparse извлекает имена хостов из правил роутеров Traefik +// и сопоставляет их с разрешёнными DNS-зонами. +package hostparse + +import ( + "strings" +) + +// Result — результат разбора правила роутера. +type Result struct { + // Hosts — валидные нормализованные (lower-case, без точки в конце) имена из Host(...). + Hosts []string + // Wildcards — хосты вида *.example.com (не поддерживаются). + Wildcards []string + // Invalid — значения Host(...), не являющиеся корректным DNS-именем. + Invalid []string + // RegexpCount — число HostRegexp(...) в правиле (игнорируются). + RegexpCount int +} + +// Parse разбирает rule и возвращает все «положительные» Host(...). +// Host внутри отрицания (!Host(...), !(...)) пропускается. +// Поддерживаются строки в обратных кавычках и двойных кавычках. +func Parse(rule string) Result { + var res Result + seen := map[string]struct{}{} + + s := []rune(rule) + n := len(s) + i := 0 + // стек скобок: true — группа под отрицанием + var stack []bool + negDepth := 0 + pendingNot := false + + skipSpaces := func() { + for i < n && isSpace(s[i]) { + i++ + } + } + + for i < n { + c := s[i] + switch { + case isSpace(c): + i++ + case c == '!': + pendingNot = true + i++ + case c == '`' || c == '"': + _, next := readString(s, i) + i = next + pendingNot = false + case c == '(': + neg := pendingNot + stack = append(stack, neg) + if neg { + negDepth++ + } + pendingNot = false + i++ + case c == ')': + if len(stack) > 0 { + if stack[len(stack)-1] { + negDepth-- + } + stack = stack[:len(stack)-1] + } + pendingNot = false + i++ + case isIdentStart(c): + start := i + for i < n && isIdentPart(s[i]) { + i++ + } + ident := string(s[start:i]) + skipSpaces() + negated := pendingNot || negDepth > 0 + pendingNot = false + if i >= n || s[i] != '(' { + continue + } + switch ident { + case "Host": + i++ // '(' + args, next := readArgs(s, i) + i = next + if negated { + continue + } + for _, a := range args { + classify(a, &res, seen) + } + case "HostRegexp": + i++ + _, next := readArgs(s, i) + i = next + if !negated { + res.RegexpCount++ + } + default: + // прочие функции: обычный вход в скобку, строки внутри пропускаются основным циклом + stack = append(stack, negated) + if negated { + negDepth++ + } + i++ + } + default: + i++ + } + } + return res +} + +func classify(raw string, res *Result, seen map[string]struct{}) { + h := Normalize(raw) + if h == "" { + res.Invalid = append(res.Invalid, raw) + return + } + if strings.HasPrefix(h, "*.") || h == "*" { + res.Wildcards = append(res.Wildcards, h) + return + } + if !ValidHostname(h) { + res.Invalid = append(res.Invalid, raw) + return + } + if _, ok := seen[h]; ok { + return + } + seen[h] = struct{}{} + res.Hosts = append(res.Hosts, h) +} + +// Normalize приводит имя к нижнему регистру и убирает пробелы и точку в конце. +func Normalize(name string) string { + name = strings.TrimSpace(name) + name = strings.TrimSuffix(name, ".") + return strings.ToLower(name) +} + +// ValidHostname проверяет, что имя — корректный ASCII hostname (LDH, допускается '_'). +func ValidHostname(h string) bool { + if h == "" || len(h) > 253 { + return false + } + for _, label := range strings.Split(h, ".") { + if label == "" || len(label) > 63 { + return false + } + if label[0] == '-' || label[len(label)-1] == '-' { + return false + } + for _, r := range label { + switch { + case r >= 'a' && r <= 'z', r >= '0' && r <= '9', r == '-', r == '_': + default: + return false + } + } + } + return true +} + +// readArgs читает список строк до закрывающей ')' (начиная сразу после '('). +func readArgs(s []rune, i int) ([]string, int) { + var args []string + n := len(s) + for i < n { + c := s[i] + switch { + case c == ')': + return args, i + 1 + case c == '`' || c == '"': + v, next := readString(s, i) + args = append(args, v) + i = next + default: + i++ + } + } + return args, i +} + +// readString читает строку в кавычках начиная с s[i] (кавычка), возвращает значение и позицию после неё. +func readString(s []rune, i int) (string, int) { + q := s[i] + i++ + var b strings.Builder + for i < len(s) { + c := s[i] + if q == '"' && c == '\\' && i+1 < len(s) { + b.WriteRune(s[i+1]) + i += 2 + continue + } + if c == q { + return b.String(), i + 1 + } + b.WriteRune(c) + i++ + } + return b.String(), i +} + +func isSpace(r rune) bool { return r == ' ' || r == '\t' || r == '\n' || r == '\r' } +func isIdentStart(r rune) bool { + return r == '_' || (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') +} +func isIdentPart(r rune) bool { return isIdentStart(r) || (r >= '0' && r <= '9') } diff --git a/internal/hostparse/hostparse_test.go b/internal/hostparse/hostparse_test.go new file mode 100644 index 0000000..8666b46 --- /dev/null +++ b/internal/hostparse/hostparse_test.go @@ -0,0 +1,78 @@ +package hostparse + +import ( + "reflect" + "testing" +) + +func TestParse(t *testing.T) { + tests := []struct { + name string + rule string + hosts []string + wildcards []string + invalid int + regexp int + }{ + {"single", "Host(`app.example.com`)", []string{"app.example.com"}, nil, 0, 0}, + {"multi args v3", "Host(`a.example.com`, `b.example.com`)", []string{"a.example.com", "b.example.com"}, nil, 0, 0}, + {"or", "Host(`a.example.com`) || Host(`b.example.com`)", []string{"a.example.com", "b.example.com"}, nil, 0, 0}, + {"with path and", "Host(`a.example.com`) && PathPrefix(`/api`)", []string{"a.example.com"}, nil, 0, 0}, + {"double quotes", `Host("a.example.com")`, []string{"a.example.com"}, nil, 0, 0}, + {"spaces", " Host ( `A.Example.COM.` ) ", []string{"a.example.com"}, nil, 0, 0}, + {"dedupe", "Host(`a.example.com`) || Host(`a.example.com`)", []string{"a.example.com"}, nil, 0, 0}, + {"regexp ignored", "HostRegexp(`^.+\\.example\\.com$`)", nil, nil, 0, 1}, + {"regexp and host", "Host(`a.example.com`) || HostRegexp(`{sub:.+}.example.com`)", []string{"a.example.com"}, nil, 0, 1}, + {"wildcard", "Host(`*.example.com`)", nil, []string{"*.example.com"}, 0, 0}, + {"negated host", "!Host(`a.example.com`) && Host(`b.example.com`)", []string{"b.example.com"}, nil, 0, 0}, + {"negated group", "!(Host(`a.example.com`) || Host(`c.example.com`)) && Host(`b.example.com`)", []string{"b.example.com"}, nil, 0, 0}, + {"group positive", "(Host(`a.example.com`) || Host(`b.example.com`)) && Path(`/x`)", []string{"a.example.com", "b.example.com"}, nil, 0, 0}, + {"host inside other string", "PathPrefix(`/Host(`)", nil, nil, 0, 0}, + {"hostsni not host", "HostSNI(`a.example.com`)", nil, nil, 0, 0}, + {"invalid", "Host(`bad host`, `-x.example.com`, `ünï.example.com`)", nil, nil, 3, 0}, + {"empty", "", nil, nil, 0, 0}, + {"api internal", "PathPrefix(`/api`) || PathPrefix(`/dashboard`)", nil, nil, 0, 0}, + {"unterminated", "Host(`a.example.com", []string{"a.example.com"}, nil, 0, 0}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + got := Parse(tc.rule) + if !reflect.DeepEqual(got.Hosts, tc.hosts) { + t.Errorf("hosts = %v, want %v", got.Hosts, tc.hosts) + } + if !reflect.DeepEqual(got.Wildcards, tc.wildcards) { + t.Errorf("wildcards = %v, want %v", got.Wildcards, tc.wildcards) + } + if len(got.Invalid) != tc.invalid { + t.Errorf("invalid = %v, want %d", got.Invalid, tc.invalid) + } + if got.RegexpCount != tc.regexp { + t.Errorf("regexp = %d, want %d", got.RegexpCount, tc.regexp) + } + }) + } +} + +func TestMatchZone(t *testing.T) { + zones := []string{"example.com", "dev.example.com", "other.org"} + tests := []struct { + host, zone string + ok bool + }{ + {"example.com", "example.com", true}, + {"app.example.com", "example.com", true}, + {"a.b.example.com", "example.com", true}, + {"app.dev.example.com", "dev.example.com", true}, + {"dev.example.com", "dev.example.com", true}, + {"badexample.com", "", false}, + {"example.com.evil.net", "", false}, + {"x.other.org", "other.org", true}, + {"unknown.net", "", false}, + } + for _, tc := range tests { + z, ok := MatchZone(tc.host, zones) + if z != tc.zone || ok != tc.ok { + t.Errorf("MatchZone(%q) = %q,%v want %q,%v", tc.host, z, ok, tc.zone, tc.ok) + } + } +} diff --git a/internal/hostparse/zone.go b/internal/hostparse/zone.go new file mode 100644 index 0000000..e5d09a9 --- /dev/null +++ b/internal/hostparse/zone.go @@ -0,0 +1,17 @@ +package hostparse + +import "strings" + +// MatchZone возвращает наиболее специфичную зону из zones, которой принадлежит host +// (host равен зоне или является её поддоменом). Зоны и host ожидаются нормализованными. +func MatchZone(host string, zones []string) (string, bool) { + best := "" + for _, z := range zones { + if host == z || strings.HasSuffix(host, "."+z) { + if len(z) > len(best) { + best = z + } + } + } + return best, best != "" +} diff --git a/internal/httpx/retry.go b/internal/httpx/retry.go new file mode 100644 index 0000000..9c6ac9e --- /dev/null +++ b/internal/httpx/retry.go @@ -0,0 +1,128 @@ +// Package httpx содержит общий HTTP-хелпер: ретраи с экспоненциальным backoff на 429/5xx. +package httpx + +import ( + "context" + "io" + "net/http" + "strconv" + "time" +) + +// MaxBodyBytes — верхняя граница читаемого тела ответа. +const MaxBodyBytes = 16 << 20 + +// Retrier повторяет запросы при сетевых ошибках, 429 и 5xx (кроме 501). +type Retrier struct { + MaxAttempts int // общее число попыток (>=1) + BaseDelay time.Duration // задержка перед второй попыткой + MaxDelay time.Duration // верхняя граница задержки (в т.ч. Retry-After) + Sleep func(ctx context.Context, d time.Duration) error +} + +// DefaultRetrier — разумные значения по умолчанию. +func DefaultRetrier() Retrier { + return Retrier{MaxAttempts: 4, BaseDelay: 500 * time.Millisecond, MaxDelay: 15 * time.Second} +} + +// Do выполняет запрос, собираемый build (вызывается на каждой попытке). +// Возвращённый ответ (в т.ч. с кодом 5xx после исчерпания попыток) должен быть закрыт вызывающим. +func (r Retrier) Do(ctx context.Context, c *http.Client, build func() (*http.Request, error)) (*http.Response, error) { + attempts := r.MaxAttempts + if attempts < 1 { + attempts = 1 + } + sleep := r.Sleep + if sleep == nil { + sleep = sleepCtx + } + var lastErr error + for attempt := 1; ; attempt++ { + req, err := build() + if err != nil { + return nil, err + } + req = req.WithContext(ctx) + resp, err := c.Do(req) + var delay time.Duration + switch { + case err != nil: + if ctx.Err() != nil { + return nil, ctx.Err() + } + lastErr = err + delay = r.backoff(attempt) + case retryable(resp.StatusCode): + if attempt >= attempts { + return resp, nil + } + delay = r.backoff(attempt) + if ra := retryAfter(resp); ra > 0 { + delay = ra + if r.MaxDelay > 0 && delay > r.MaxDelay { + delay = r.MaxDelay + } + } + _, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 1<<20)) + _ = resp.Body.Close() + default: + return resp, nil + } + if attempt >= attempts { + return nil, lastErr + } + if err := sleep(ctx, delay); err != nil { + return nil, err + } + } +} + +func (r Retrier) backoff(attempt int) time.Duration { + d := r.BaseDelay + for i := 1; i < attempt; i++ { + d *= 2 + if r.MaxDelay > 0 && d >= r.MaxDelay { + return r.MaxDelay + } + } + return d +} + +func retryable(code int) bool { + return code == http.StatusTooManyRequests || (code >= 500 && code != http.StatusNotImplemented) +} + +func retryAfter(resp *http.Response) time.Duration { + v := resp.Header.Get("Retry-After") + if v == "" { + return 0 + } + if sec, err := strconv.Atoi(v); err == nil && sec >= 0 { + return time.Duration(sec) * time.Second + } + return 0 +} + +func sleepCtx(ctx context.Context, d time.Duration) error { + t := time.NewTimer(d) + defer t.Stop() + select { + case <-ctx.Done(): + return ctx.Err() + case <-t.C: + return nil + } +} + +// ReadBody читает тело ответа с ограничением по размеру. +func ReadBody(resp *http.Response) ([]byte, error) { + return io.ReadAll(io.LimitReader(resp.Body, MaxBodyBytes)) +} + +// Snippet возвращает обрезанное тело для сообщений об ошибках. +func Snippet(b []byte, n int) string { + if len(b) > n { + b = b[:n] + } + return string(b) +} diff --git a/internal/httpx/retry_test.go b/internal/httpx/retry_test.go new file mode 100644 index 0000000..eeb9206 --- /dev/null +++ b/internal/httpx/retry_test.go @@ -0,0 +1,99 @@ +package httpx + +import ( + "context" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" +) + +func noSleep(context.Context, time.Duration) error { return nil } + +func TestRetrierRetriesOn5xxAnd429(t *testing.T) { + var calls int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch atomic.AddInt32(&calls, 1) { + case 1: + w.WriteHeader(http.StatusTooManyRequests) + case 2: + w.WriteHeader(http.StatusBadGateway) + default: + w.WriteHeader(http.StatusOK) + } + })) + defer srv.Close() + r := Retrier{MaxAttempts: 4, BaseDelay: time.Millisecond, Sleep: noSleep} + resp, err := r.Do(context.Background(), srv.Client(), func() (*http.Request, error) { + return http.NewRequest(http.MethodGet, srv.URL, nil) + }) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + if resp.StatusCode != 200 || calls != 3 { + t.Fatalf("status=%d calls=%d", resp.StatusCode, calls) + } +} + +func TestRetrierNoRetryOn4xx(t *testing.T) { + var calls int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt32(&calls, 1) + w.WriteHeader(http.StatusNotFound) + })) + defer srv.Close() + r := Retrier{MaxAttempts: 4, BaseDelay: time.Millisecond, Sleep: noSleep} + resp, err := r.Do(context.Background(), srv.Client(), func() (*http.Request, error) { + return http.NewRequest(http.MethodGet, srv.URL, nil) + }) + if err != nil { + t.Fatal(err) + } + resp.Body.Close() + if calls != 1 || resp.StatusCode != 404 { + t.Fatalf("calls=%d status=%d", calls, resp.StatusCode) + } +} + +func TestRetrierExhaustedReturnsLastResponse(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusServiceUnavailable) + })) + defer srv.Close() + r := Retrier{MaxAttempts: 2, BaseDelay: time.Millisecond, Sleep: noSleep} + resp, err := r.Do(context.Background(), srv.Client(), func() (*http.Request, error) { + return http.NewRequest(http.MethodGet, srv.URL, nil) + }) + if err != nil { + t.Fatal(err) + } + resp.Body.Close() + if resp.StatusCode != 503 { + t.Fatalf("status=%d", resp.StatusCode) + } +} + +func TestRetrierNetworkError(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {})) + url := srv.URL + srv.Close() + r := Retrier{MaxAttempts: 3, BaseDelay: time.Millisecond, Sleep: noSleep} + _, err := r.Do(context.Background(), http.DefaultClient, func() (*http.Request, error) { + return http.NewRequest(http.MethodGet, url, nil) + }) + if err == nil { + t.Fatal("ожидалась ошибка") + } +} + +func TestBackoffCapped(t *testing.T) { + r := Retrier{BaseDelay: time.Second, MaxDelay: 5 * time.Second} + if got := r.backoff(1); got != time.Second { + t.Errorf("attempt1 = %v", got) + } + if got := r.backoff(10); got != 5*time.Second { + t.Errorf("attempt10 = %v", got) + } +} diff --git a/internal/reconciler/loop.go b/internal/reconciler/loop.go new file mode 100644 index 0000000..b4a8249 --- /dev/null +++ b/internal/reconciler/loop.go @@ -0,0 +1,61 @@ +package reconciler + +import ( + "context" + "time" +) + +// CycleTimeout — верхняя граница одного цикла reconcile. +const CycleTimeout = 5 * time.Minute + +// Run запускает цикл: сразу, затем каждые interval, а также по сигналу trigger +// (после debounce, чтобы схлопнуть пачку событий). Блокируется до отмены ctx. +// report вызывается после каждого цикла. +func (r *Reconciler) Run(ctx context.Context, interval, debounce time.Duration, trigger <-chan struct{}, report func(Result, error)) { + cycle := func() { + cctx, cancel := context.WithTimeout(ctx, CycleTimeout) + defer cancel() + res, err := r.Reconcile(cctx) + if err != nil { + r.log.Error("цикл reconcile завершён с ошибками", "error", err) + } else { + r.log.Info("цикл reconcile завершён", + "created", res.Created, "updated", res.Updated, "deleted", res.Deleted, + "unchanged", res.Unchanged, "skippedForeign", res.SkippedForeign, + "skippedConflict", res.SkippedConflict, "pendingDelete", res.PendingDelete) + } + if report != nil { + report(res, err) + } + } + + cycle() + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + cycle() + case <-trigger: + if debounce > 0 { + t := time.NewTimer(debounce) + select { + case <-ctx.Done(): + t.Stop() + return + case <-t.C: + } + } + // схлопываем накопившиеся сигналы + select { + case <-trigger: + default: + } + r.log.Debug("reconcile по событию Docker") + cycle() + ticker.Reset(interval) + } + } +} diff --git a/internal/reconciler/loop_test.go b/internal/reconciler/loop_test.go new file mode 100644 index 0000000..6a3e7f6 --- /dev/null +++ b/internal/reconciler/loop_test.go @@ -0,0 +1,40 @@ +package reconciler + +import ( + "context" + "sync/atomic" + "testing" + "time" + + "github.com/realmanual/traefik-selectel/internal/traefik" +) + +func TestRunTicksAndTriggers(t *testing.T) { + e := newEnv(nil) + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + trigger := make(chan struct{}, 1) + var cycles int32 + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { + e.rec.Run(ctx, time.Hour, 5*time.Millisecond, trigger, func(Result, error) { atomic.AddInt32(&cycles, 1) }) + close(done) + }() + waitCycles := func(n int32) { + deadline := time.Now().Add(3 * time.Second) + for atomic.LoadInt32(&cycles) < n { + if time.Now().After(deadline) { + t.Fatalf("ожидали %d циклов, было %d", n, atomic.LoadInt32(&cycles)) + } + time.Sleep(5 * time.Millisecond) + } + } + waitCycles(1) // стартовый + trigger <- struct{}{} + waitCycles(2) + cancel() + <-done + if len(e.dns.created) != 1 { + t.Fatalf("created = %v", e.dns.created) + } +} diff --git a/internal/reconciler/reconciler.go b/internal/reconciler/reconciler.go new file mode 100644 index 0000000..954c597 --- /dev/null +++ b/internal/reconciler/reconciler.go @@ -0,0 +1,321 @@ +// Package reconciler приводит DNS-записи Selectel в соответствие с роутерами Traefik. +package reconciler + +import ( + "context" + "errors" + "fmt" + "log/slog" + "sort" + "sync" + "time" + + "github.com/realmanual/traefik-selectel/internal/hostparse" + "github.com/realmanual/traefik-selectel/internal/selectel" + "github.com/realmanual/traefik-selectel/internal/traefik" +) + +// DNSClient — операции с DNS, необходимые reconciler-у. +// Имена — нормализованные (lower-case, без точки в конце). +type DNSClient interface { + ListRRSets(ctx context.Context, zone string) ([]selectel.RRSet, error) + CreateRRSet(ctx context.Context, zone string, rr selectel.RRSet) error + UpdateRRSet(ctx context.Context, zone, rrsetID string, ttl int, records []selectel.Record, comment string) error + DeleteRRSet(ctx context.Context, zone, rrsetID string) error +} + +// Config — параметры reconcile. +type Config struct { + Zones []string // нормализованные разрешённые зоны + IP string // IPv4 для A-записей + TTL int + Comment string // метка владения + DryRun bool + DeleteEnabled bool + GracePeriod time.Duration +} + +// Result — итоги одного цикла. +type Result struct { + Created, Updated, Deleted, Unchanged, SkippedForeign, SkippedConflict int + PendingDelete int +} + +// Reconciler — не потокобезопасен по вызовам Reconcile (вызывайте последовательно). +type Reconciler struct { + cfg Config + source traefik.RouterSource + dns DNSClient + log *slog.Logger + now func() time.Time + + hadRouters bool + pending map[string]pendingDelete // ключ — имя записи + warnedMu sync.Mutex + warned map[string]struct{} +} + +type pendingDelete struct { + zone string + since time.Time +} + +// New создаёт reconciler. +func New(cfg Config, source traefik.RouterSource, dns DNSClient, log *slog.Logger) *Reconciler { + if log == nil { + log = slog.Default() + } + return &Reconciler{ + cfg: cfg, source: source, dns: dns, log: log, now: time.Now, + pending: map[string]pendingDelete{}, warned: map[string]struct{}{}, + } +} + +// warnOnce логирует warning один раз на ключ (дальше — debug), чтобы не шуметь каждый цикл. +func (r *Reconciler) warnOnce(key, msg string, args ...any) { + r.warnedMu.Lock() + _, seen := r.warned[key] + r.warned[key] = struct{}{} + r.warnedMu.Unlock() + if seen { + r.log.Debug(msg, args...) + return + } + r.log.Warn(msg, args...) +} + +type desiredSet struct { + create map[string]string // имя -> зона (enabled-роутеры) + present map[string]string // имя -> зона (enabled + warning): защита от удаления +} + +// Reconcile выполняет один цикл. +func (r *Reconciler) Reconcile(ctx context.Context) (Result, error) { + var res Result + + routers, err := r.source.Routers(ctx) + if err != nil { + // Источник недоступен: ничего не создаём и, главное, ничего не удаляем. + return res, fmt.Errorf("получение роутеров: %w", err) + } + if len(routers) == 0 && r.hadRouters { + r.log.Warn("Traefik API вернул пустой список роутеров при ранее непустом — цикл пропущен, удаления заблокированы") + return res, errors.New("пустой список роутеров при ранее непустом") + } + if len(routers) > 0 { + r.hadRouters = true + } + + desired := r.collect(routers) + + // Какие зоны надо читать: с желаемыми хостами, а при включённом удалении — все. + zoneSet := map[string]struct{}{} + for _, z := range desired.create { + zoneSet[z] = struct{}{} + } + if r.cfg.DeleteEnabled { + for _, z := range r.cfg.Zones { + zoneSet[z] = struct{}{} + } + } + zones := make([]string, 0, len(zoneSet)) + for z := range zoneSet { + zones = append(zones, z) + } + sort.Strings(zones) + + var errs []error + listed := map[string][]selectel.RRSet{} + for _, zone := range zones { + rrsets, err := r.dns.ListRRSets(ctx, zone) + if err != nil { + errs = append(errs, fmt.Errorf("зона %s: список записей: %w", zone, err)) + r.log.Error("не удалось получить записи зоны", "zone", zone, "error", err) + continue + } + listed[zone] = rrsets + } + + // Создание/обновление. + names := make([]string, 0, len(desired.create)) + for n := range desired.create { + names = append(names, n) + } + sort.Strings(names) + for _, name := range names { + zone := desired.create[name] + rrsets, ok := listed[zone] + if !ok { + continue + } + if err := r.ensure(ctx, zone, name, rrsets, &res); err != nil { + errs = append(errs, err) + } + } + + if r.cfg.DeleteEnabled { + errs = append(errs, r.deleteStale(ctx, listed, desired, &res)...) + } + res.PendingDelete = len(r.pending) + return res, errors.Join(errs...) +} + +// collect извлекает желаемые имена из роутеров. +func (r *Reconciler) collect(routers []traefik.Router) desiredSet { + d := desiredSet{create: map[string]string{}, present: map[string]string{}} + for _, rt := range routers { + if rt.Status != "enabled" && rt.Status != "warning" { + continue + } + p := hostparse.Parse(rt.Rule) + if p.RegexpCount > 0 { + r.log.Debug("HostRegexp игнорируется", "router", rt.Name, "count", p.RegexpCount) + } + for _, w := range p.Wildcards { + r.warnOnce("wildcard|"+rt.Name+"|"+w, "wildcard-хосты в Host() не поддерживаются, пропущено", "router", rt.Name, "host", w) + } + for _, inv := range p.Invalid { + r.warnOnce("invalid|"+rt.Name+"|"+inv, "некорректное имя хоста в Host(), пропущено", "router", rt.Name, "host", inv) + } + for _, h := range p.Hosts { + zone, ok := hostparse.MatchZone(h, r.cfg.Zones) + if !ok { + r.log.Debug("хост вне разрешённых зон, пропущен", "router", rt.Name, "host", h) + continue + } + d.present[h] = zone + if rt.Status == "enabled" { + d.create[h] = zone + } + } + } + return d +} + +func isManaged(rr selectel.RRSet, comment string) bool { return rr.Comment == comment } + +func (r *Reconciler) ensure(ctx context.Context, zone, name string, rrsets []selectel.RRSet, res *Result) error { + var existingA *selectel.RRSet + var cname *selectel.RRSet + for i := range rrsets { + rr := &rrsets[i] + if rr.Name != name { + continue + } + switch rr.Type { + case "A": + existingA = rr + case "CNAME": + cname = rr + } + } + + // Хост снова востребован — отменяем отложенное удаление. + delete(r.pending, name) + + switch { + case existingA != nil && !isManaged(*existingA, r.cfg.Comment): + res.SkippedForeign++ + r.warnOnce("foreign|"+name, "A-запись существует без метки владения, не трогаем", "name", name, "zone", zone) + return nil + case existingA != nil: + if r.matches(*existingA) { + res.Unchanged++ + r.log.Debug("запись актуальна", "name", name) + return nil + } + if r.cfg.DryRun { + r.log.Info("dry-run: обновил бы A-запись", "name", name, "ip", r.cfg.IP, "ttl", r.cfg.TTL) + return nil + } + if err := r.dns.UpdateRRSet(ctx, zone, existingA.ID, r.cfg.TTL, []selectel.Record{{Content: r.cfg.IP}}, r.cfg.Comment); err != nil { + r.log.Error("не удалось обновить A-запись", "name", name, "error", err) + return fmt.Errorf("обновление %s: %w", name, err) + } + res.Updated++ + r.log.Info("A-запись обновлена", "name", name, "ip", r.cfg.IP, "ttl", r.cfg.TTL) + return nil + case cname != nil: + res.SkippedConflict++ + r.warnOnce("cname|"+name, "для имени существует CNAME, A-запись не создаём", "name", name, "zone", zone) + return nil + } + + if r.cfg.DryRun { + r.log.Info("dry-run: создал бы A-запись", "name", name, "ip", r.cfg.IP, "ttl", r.cfg.TTL) + return nil + } + err := r.dns.CreateRRSet(ctx, zone, selectel.RRSet{ + Name: name, Type: "A", TTL: r.cfg.TTL, + Records: []selectel.Record{{Content: r.cfg.IP}}, Comment: r.cfg.Comment, + }) + if err != nil { + r.log.Error("не удалось создать A-запись", "name", name, "error", err) + return fmt.Errorf("создание %s: %w", name, err) + } + res.Created++ + r.log.Info("A-запись создана", "name", name, "ip", r.cfg.IP, "ttl", r.cfg.TTL) + return nil +} + +func (r *Reconciler) matches(rr selectel.RRSet) bool { + return rr.TTL == r.cfg.TTL && len(rr.Records) == 1 && + rr.Records[0].Content == r.cfg.IP && !rr.Records[0].Disabled +} + +// deleteStale удаляет managed A-записи, для которых роутер отсутствует дольше GracePeriod. +func (r *Reconciler) deleteStale(ctx context.Context, listed map[string][]selectel.RRSet, d desiredSet, res *Result) []error { + var errs []error + now := r.now() + + zones := make([]string, 0, len(listed)) + for z := range listed { + zones = append(zones, z) + } + sort.Strings(zones) + + stale := map[string]struct{}{} + for _, zone := range zones { + for _, rr := range listed[zone] { + if rr.Type != "A" || !isManaged(rr, r.cfg.Comment) { + continue + } + if _, ok := d.present[rr.Name]; ok { + continue + } + stale[rr.Name] = struct{}{} + p, ok := r.pending[rr.Name] + if !ok { + r.pending[rr.Name] = pendingDelete{zone: zone, since: now} + r.log.Info("роут исчез, запись будет удалена после grace-периода", "name", rr.Name, "grace", r.cfg.GracePeriod.String()) + continue + } + if now.Sub(p.since) < r.cfg.GracePeriod { + continue + } + if r.cfg.DryRun { + r.log.Info("dry-run: удалил бы A-запись", "name", rr.Name) + continue + } + if err := r.dns.DeleteRRSet(ctx, zone, rr.ID); err != nil { + r.log.Error("не удалось удалить A-запись", "name", rr.Name, "error", err) + errs = append(errs, fmt.Errorf("удаление %s: %w", rr.Name, err)) + continue + } + delete(r.pending, rr.Name) + res.Deleted++ + r.log.Info("A-запись удалена", "name", rr.Name) + } + } + + // Чистим состояние для успешно прочитанных зон: запись уже не stale (вернулась/удалена извне). + for name, p := range r.pending { + if _, zoneOK := listed[p.zone]; !zoneOK { + continue + } + if _, ok := stale[name]; !ok { + delete(r.pending, name) + } + } + return errs +} diff --git a/internal/reconciler/reconciler_test.go b/internal/reconciler/reconciler_test.go new file mode 100644 index 0000000..c44964b --- /dev/null +++ b/internal/reconciler/reconciler_test.go @@ -0,0 +1,368 @@ +package reconciler + +import ( + "context" + "errors" + "io" + "log/slog" + "sort" + "testing" + "time" + + "github.com/realmanual/traefik-selectel/internal/selectel" + "github.com/realmanual/traefik-selectel/internal/traefik" +) + +const marker = "managed-by=traefik-selectel-dns" + +type fakeSource struct { + routers []traefik.Router + err error +} + +func (f *fakeSource) Routers(context.Context) ([]traefik.Router, error) { return f.routers, f.err } + +type fakeDNS struct { + zones map[string][]selectel.RRSet + listErr map[string]error + created []selectel.RRSet + updated []string + deleted []string + nextID int +} + +func (f *fakeDNS) ListRRSets(_ context.Context, zone string) ([]selectel.RRSet, error) { + if err := f.listErr[zone]; err != nil { + return nil, err + } + return append([]selectel.RRSet(nil), f.zones[zone]...), nil +} + +func (f *fakeDNS) CreateRRSet(_ context.Context, zone string, rr selectel.RRSet) error { + f.created = append(f.created, rr) + f.nextID++ + rr.ID = "new-" + string(rune('0'+f.nextID)) + f.zones[zone] = append(f.zones[zone], rr) + return nil +} + +func (f *fakeDNS) UpdateRRSet(_ context.Context, zone, id string, ttl int, recs []selectel.Record, comment string) error { + f.updated = append(f.updated, id) + for i, rr := range f.zones[zone] { + if rr.ID == id { + f.zones[zone][i].TTL, f.zones[zone][i].Records, f.zones[zone][i].Comment = ttl, recs, comment + } + } + return nil +} + +func (f *fakeDNS) DeleteRRSet(_ context.Context, zone, id string) error { + f.deleted = append(f.deleted, id) + var keep []selectel.RRSet + for _, rr := range f.zones[zone] { + if rr.ID != id { + keep = append(keep, rr) + } + } + f.zones[zone] = keep + return nil +} + +func router(name, rule, status string) traefik.Router { + return traefik.Router{Name: name, Rule: rule, Status: status} +} + +func aRec(id, name, ip, comment string) selectel.RRSet { + return selectel.RRSet{ID: id, Name: name, Type: "A", TTL: 300, Comment: comment, Records: []selectel.Record{{Content: ip}}} +} + +type env struct { + src *fakeSource + dns *fakeDNS + rec *Reconciler + now time.Time +} + +func newEnv(cfgMod func(*Config)) *env { + cfg := Config{Zones: []string{"example.com"}, IP: "1.2.3.4", TTL: 300, Comment: marker, GracePeriod: 10 * time.Minute} + if cfgMod != nil { + cfgMod(&cfg) + } + e := &env{ + src: &fakeSource{}, + dns: &fakeDNS{zones: map[string][]selectel.RRSet{"example.com": nil}, listErr: map[string]error{}}, + now: time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC), + } + e.rec = New(cfg, e.src, e.dns, slog.New(slog.NewTextHandler(io.Discard, nil))) + e.rec.now = func() time.Time { return e.now } + return e +} + +func (e *env) run(t *testing.T) Result { + t.Helper() + res, err := e.rec.Reconcile(context.Background()) + if err != nil { + t.Fatalf("reconcile: %v", err) + } + return res +} + +func TestCreate(t *testing.T) { + e := newEnv(nil) + e.src.routers = []traefik.Router{ + router("a@docker", "Host(`app.example.com`)", "enabled"), + router("b@file", "Host(`x.other.net`)", "enabled"), // вне зон + router("c@docker", "Host(`off.example.com`)", "disabled"), // не enabled + router("d@docker", "HostRegexp(`^.+\\.example\\.com$`)", "enabled"), + router("e@docker", "Host(`*.example.com`)", "enabled"), + } + res := e.run(t) + if res.Created != 1 || len(e.dns.created) != 1 { + t.Fatalf("res=%+v created=%v", res, e.dns.created) + } + c := e.dns.created[0] + if c.Name != "app.example.com" || c.Type != "A" || c.TTL != 300 || c.Comment != marker || c.Records[0].Content != "1.2.3.4" { + t.Errorf("created = %+v", c) + } +} + +func TestZoneApexAndIdempotency(t *testing.T) { + e := newEnv(nil) + e.src.routers = []traefik.Router{router("a", "Host(`example.com`) || Host(`www.example.com`)", "enabled")} + if res := e.run(t); res.Created != 2 { + t.Fatalf("res=%+v", res) + } + res := e.run(t) + if res.Created != 0 || res.Updated != 0 || res.Unchanged != 2 || len(e.dns.created) != 2 { + t.Fatalf("второй цикл не идемпотентен: %+v", res) + } +} + +func TestUpdateManagedWhenDiffers(t *testing.T) { + e := newEnv(nil) + e.dns.zones["example.com"] = []selectel.RRSet{aRec("r1", "app.example.com", "9.9.9.9", marker)} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + res := e.run(t) + if res.Updated != 1 || len(e.dns.updated) != 1 || e.dns.zones["example.com"][0].Records[0].Content != "1.2.3.4" { + t.Fatalf("res=%+v", res) + } +} + +func TestUpdateOnTTLChangeAndDisabledRecord(t *testing.T) { + e := newEnv(nil) + rr := aRec("r1", "app.example.com", "1.2.3.4", marker) + rr.TTL = 60 + rr2 := aRec("r2", "b.example.com", "1.2.3.4", marker) + rr2.Records[0].Disabled = true + e.dns.zones["example.com"] = []selectel.RRSet{rr, rr2} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`, `b.example.com`)", "enabled")} + if res := e.run(t); res.Updated != 2 { + t.Fatalf("res=%+v", res) + } +} + +func TestSkipForeignRecord(t *testing.T) { + e := newEnv(nil) + e.dns.zones["example.com"] = []selectel.RRSet{ + aRec("r1", "app.example.com", "9.9.9.9", ""), + aRec("r2", "b.example.com", "9.9.9.9", "created by hand"), + } + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`) || Host(`b.example.com`)", "enabled")} + res := e.run(t) + if res.SkippedForeign != 2 || len(e.dns.created)+len(e.dns.updated)+len(e.dns.deleted) != 0 { + t.Fatalf("res=%+v dns=%+v", res, e.dns) + } +} + +func TestSkipWhenCNAMEExists(t *testing.T) { + e := newEnv(nil) + e.dns.zones["example.com"] = []selectel.RRSet{{ID: "c1", Name: "app.example.com", Type: "CNAME", TTL: 300, + Records: []selectel.Record{{Content: "other.net."}}}} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + if res := e.run(t); res.SkippedConflict != 1 || len(e.dns.created) != 0 { + t.Fatalf("res=%+v", res) + } +} + +func TestOtherTypesDoNotBlockCreate(t *testing.T) { + e := newEnv(nil) + e.dns.zones["example.com"] = []selectel.RRSet{{ID: "t1", Name: "app.example.com", Type: "TXT", TTL: 300}} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + if res := e.run(t); res.Created != 1 { + t.Fatalf("res=%+v", res) + } +} + +func TestDryRunDoesNotMutate(t *testing.T) { + e := newEnv(func(c *Config) { c.DryRun = true }) + e.dns.zones["example.com"] = []selectel.RRSet{aRec("r1", "old.example.com", "9.9.9.9", marker)} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled"), router("b", "Host(`old.example.com`)", "enabled")} + e.run(t) + if len(e.dns.created)+len(e.dns.updated)+len(e.dns.deleted) != 0 { + t.Fatalf("dry-run изменил DNS: %+v", e.dns) + } +} + +func TestDeleteDisabledByDefault(t *testing.T) { + e := newEnv(nil) + e.dns.zones["example.com"] = []selectel.RRSet{aRec("r1", "gone.example.com", "1.2.3.4", marker)} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + e.run(t) + e.now = e.now.Add(24 * time.Hour) + e.run(t) + if len(e.dns.deleted) != 0 { + t.Fatalf("удаление должно быть выключено: %v", e.dns.deleted) + } +} + +func TestDeleteAfterGracePeriod(t *testing.T) { + e := newEnv(func(c *Config) { c.DeleteEnabled = true }) + e.dns.zones["example.com"] = []selectel.RRSet{ + aRec("r1", "gone.example.com", "1.2.3.4", marker), + aRec("r2", "foreign.example.com", "1.2.3.4", ""), // чужая — никогда не удаляем + } + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + + res := e.run(t) // старт отсчёта + if len(e.dns.deleted) != 0 || res.PendingDelete != 1 { + t.Fatalf("res=%+v deleted=%v", res, e.dns.deleted) + } + e.now = e.now.Add(9 * time.Minute) + e.run(t) + if len(e.dns.deleted) != 0 { + t.Fatal("удалено раньше grace-периода") + } + e.now = e.now.Add(2 * time.Minute) + res = e.run(t) + if res.Deleted != 1 || len(e.dns.deleted) != 1 || e.dns.deleted[0] != "r1" { + t.Fatalf("res=%+v deleted=%v", res, e.dns.deleted) + } + for _, rr := range e.dns.zones["example.com"] { + if rr.ID == "r2" { + return + } + } + t.Fatal("чужая запись удалена") +} + +func TestGraceResetWhenRouterReturns(t *testing.T) { + e := newEnv(func(c *Config) { c.DeleteEnabled = true }) + e.dns.zones["example.com"] = []selectel.RRSet{aRec("r1", "app.example.com", "1.2.3.4", marker)} + other := router("x", "Host(`other.example.com`)", "enabled") + e.src.routers = []traefik.Router{other} + e.run(t) + e.now = e.now.Add(8 * time.Minute) + e.src.routers = []traefik.Router{other, router("a", "Host(`app.example.com`)", "enabled")} // вернулся + e.run(t) + e.now = e.now.Add(8 * time.Minute) + e.src.routers = []traefik.Router{other} // снова исчез — отсчёт заново + e.run(t) + e.now = e.now.Add(8 * time.Minute) + e.run(t) + if len(e.dns.deleted) != 0 { + t.Fatalf("grace не сброшен: %v", e.dns.deleted) + } +} + +func TestWarningRouterProtectsFromDelete(t *testing.T) { + e := newEnv(func(c *Config) { c.DeleteEnabled = true }) + e.dns.zones["example.com"] = []selectel.RRSet{aRec("r1", "app.example.com", "1.2.3.4", marker)} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "warning")} + e.run(t) + e.now = e.now.Add(time.Hour) + e.run(t) + if len(e.dns.deleted) != 0 { + t.Fatal("router со статусом warning не должен приводить к удалению") + } +} + +func TestSafeGuardTraefikError(t *testing.T) { + e := newEnv(func(c *Config) { c.DeleteEnabled = true }) + e.dns.zones["example.com"] = []selectel.RRSet{aRec("r1", "gone.example.com", "1.2.3.4", marker)} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + e.run(t) // gone -> pending + e.now = e.now.Add(time.Hour) + e.src.err = errors.New("connection refused") + if _, err := e.rec.Reconcile(context.Background()); err == nil { + t.Fatal("ожидалась ошибка") + } + if len(e.dns.deleted) != 0 { + t.Fatal("удаление при ошибке Traefik API") + } +} + +func TestSafeGuardEmptyListAfterNonEmpty(t *testing.T) { + e := newEnv(func(c *Config) { c.DeleteEnabled = true }) + e.dns.zones["example.com"] = []selectel.RRSet{aRec("r1", "app.example.com", "1.2.3.4", marker)} + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + e.run(t) + e.src.routers = nil + e.now = e.now.Add(time.Hour) + if _, err := e.rec.Reconcile(context.Background()); err == nil { + t.Fatal("ожидалась ошибка на пустой список") + } + e.now = e.now.Add(time.Hour) + _, _ = e.rec.Reconcile(context.Background()) + if len(e.dns.deleted) != 0 { + t.Fatal("удаление при пустом списке роутеров") + } +} + +func TestZoneListErrorSkipsZoneOnly(t *testing.T) { + e := newEnv(func(c *Config) { c.Zones = []string{"example.com", "other.org"}; c.DeleteEnabled = true }) + e.dns.zones["other.org"] = []selectel.RRSet{aRec("r9", "old.other.org", "1.2.3.4", marker)} + e.dns.listErr["other.org"] = errors.New("boom") + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`) || Host(`x.other.org`)", "enabled")} + res, err := e.rec.Reconcile(context.Background()) + if err == nil { + t.Fatal("ожидалась ошибка зоны") + } + if res.Created != 1 || e.dns.created[0].Name != "app.example.com" { + t.Fatalf("res=%+v created=%v", res, e.dns.created) + } + e.now = e.now.Add(time.Hour) + _, _ = e.rec.Reconcile(context.Background()) + if len(e.dns.deleted) != 0 { + t.Fatal("удаление в зоне с ошибкой чтения") + } +} + +func TestLongestZoneWins(t *testing.T) { + e := newEnv(func(c *Config) { c.Zones = []string{"example.com", "dev.example.com"} }) + e.dns.zones["dev.example.com"] = nil + e.src.routers = []traefik.Router{router("a", "Host(`app.dev.example.com`)", "enabled")} + e.run(t) + if len(e.dns.zones["dev.example.com"]) != 1 || len(e.dns.zones["example.com"]) != 0 { + t.Fatalf("zones = %+v", e.dns.zones) + } +} + +func TestCreateErrorReported(t *testing.T) { + e := newEnv(nil) + e.dns.zones["example.com"] = nil + e.src.routers = []traefik.Router{router("a", "Host(`app.example.com`)", "enabled")} + fd := &failingCreate{fakeDNS: e.dns} + e.rec.dns = fd + if _, err := e.rec.Reconcile(context.Background()); err == nil { + t.Fatal("ожидалась ошибка") + } +} + +type failingCreate struct{ *fakeDNS } + +func (f *failingCreate) CreateRRSet(context.Context, string, selectel.RRSet) error { + return errors.New("422") +} + +func TestNamesDeterministic(t *testing.T) { + e := newEnv(nil) + e.src.routers = []traefik.Router{router("a", "Host(`c.example.com`,`a.example.com`,`b.example.com`)", "enabled")} + e.run(t) + var got []string + for _, c := range e.dns.created { + got = append(got, c.Name) + } + if !sort.StringsAreSorted(got) { + t.Fatalf("порядок создания не детерминирован: %v", got) + } +} diff --git a/internal/selectel/auth.go b/internal/selectel/auth.go new file mode 100644 index 0000000..cdfb8f6 --- /dev/null +++ b/internal/selectel/auth.go @@ -0,0 +1,186 @@ +package selectel + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" + "strings" + "sync" + "time" + + "github.com/realmanual/traefik-selectel/internal/httpx" +) + +// DefaultAuthURL — Keystone (identity v3) Selectel. +const DefaultAuthURL = "https://cloud.api.selcloud.ru/identity/v3" + +const ( + tokenRefreshMargin = 10 * time.Minute + // Если expires_at не удалось разобрать: документация заявляет 24 часа, берём консервативно 1 час. + fallbackTokenTTL = time.Hour +) + +// Credentials — учётные данные сервисного пользователя Selectel. +// Значения секретов не должны попадать в логи: String() их скрывает. +type Credentials struct { + Username string + Password string + AccountID string // имя домена Keystone (номер аккаунта) + ProjectID string +} + +// String скрывает секреты. +func (c Credentials) String() string { + return fmt.Sprintf("Credentials{username=%q account=%q project=%q password=}", c.Username, c.AccountID, c.ProjectID) +} + +// Validate проверяет, что все поля заданы. +func (c Credentials) Validate() error { + var missing []string + if c.Username == "" { + missing = append(missing, "SELECTEL_USERNAME") + } + if c.Password == "" { + missing = append(missing, "SELECTEL_PASSWORD") + } + if c.AccountID == "" { + missing = append(missing, "SELECTEL_ACCOUNT_ID") + } + if c.ProjectID == "" { + missing = append(missing, "SELECTEL_PROJECT_ID") + } + if len(missing) > 0 { + return fmt.Errorf("selectel: не заданы учётные данные: %s", strings.Join(missing, ", ")) + } + return nil +} + +// TokenSource выдаёт IAM-токен проекта. +type TokenSource interface { + Token(ctx context.Context) (string, error) + // Invalidate сбрасывает кэш (например, после 401). + Invalidate() +} + +// IAMTokenSource получает токен через Keystone и кэширует его в памяти. +type IAMTokenSource struct { + creds Credentials + authURL string + http *http.Client + retry httpx.Retrier + now func() time.Time + + mu sync.Mutex + token string + expires time.Time +} + +// NewIAMTokenSource создаёт источник токенов. authURL пустой — DefaultAuthURL. +func NewIAMTokenSource(creds Credentials, authURL string, timeout time.Duration) (*IAMTokenSource, error) { + if err := creds.Validate(); err != nil { + return nil, err + } + if authURL == "" { + authURL = DefaultAuthURL + } + return &IAMTokenSource{ + creds: creds, + authURL: strings.TrimRight(authURL, "/"), + http: &http.Client{Timeout: timeout}, + retry: httpx.DefaultRetrier(), + now: time.Now, + }, nil +} + +// Token возвращает актуальный токен, обновляя его за tokenRefreshMargin до истечения. +func (s *IAMTokenSource) Token(ctx context.Context) (string, error) { + s.mu.Lock() + defer s.mu.Unlock() + if s.token != "" && s.now().Add(tokenRefreshMargin).Before(s.expires) { + return s.token, nil + } + tok, exp, err := s.fetch(ctx) + if err != nil { + return "", err + } + s.token, s.expires = tok, exp + return tok, nil +} + +// Invalidate сбрасывает кэшированный токен. +func (s *IAMTokenSource) Invalidate() { + s.mu.Lock() + s.token, s.expires = "", time.Time{} + s.mu.Unlock() +} + +type authRequest struct { + Auth struct { + Identity struct { + Methods []string `json:"methods"` + Password struct { + User struct { + Name string `json:"name"` + Domain map[string]string `json:"domain"` + Password string `json:"password"` + } `json:"user"` + } `json:"password"` + } `json:"identity"` + Scope struct { + Project struct { + ID string `json:"id"` + Domain map[string]string `json:"domain"` + } `json:"project"` + } `json:"scope"` + } `json:"auth"` +} + +func (s *IAMTokenSource) fetch(ctx context.Context) (string, time.Time, error) { + var ar authRequest + ar.Auth.Identity.Methods = []string{"password"} + ar.Auth.Identity.Password.User.Name = s.creds.Username + ar.Auth.Identity.Password.User.Password = s.creds.Password + ar.Auth.Identity.Password.User.Domain = map[string]string{"name": s.creds.AccountID} + ar.Auth.Scope.Project.ID = s.creds.ProjectID + ar.Auth.Scope.Project.Domain = map[string]string{"name": s.creds.AccountID} + payload, err := json.Marshal(ar) + if err != nil { + return "", time.Time{}, err + } + resp, err := s.retry.Do(ctx, s.http, func() (*http.Request, error) { + req, err := http.NewRequest(http.MethodPost, s.authURL+"/auth/tokens", bytes.NewReader(payload)) + if err != nil { + return nil, err + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Accept", "application/json") + return req, nil + }) + if err != nil { + return "", time.Time{}, fmt.Errorf("selectel: получение IAM-токена: %w", err) + } + defer resp.Body.Close() + body, _ := httpx.ReadBody(resp) + if resp.StatusCode != http.StatusCreated && resp.StatusCode != http.StatusOK { + // Тело ответа Keystone намеренно не логируется. + return "", time.Time{}, fmt.Errorf("selectel: получение IAM-токена: статус %d", resp.StatusCode) + } + tok := resp.Header.Get("X-Subject-Token") + if tok == "" { + return "", time.Time{}, fmt.Errorf("selectel: в ответе нет заголовка X-Subject-Token") + } + exp := s.now().Add(fallbackTokenTTL) + var parsed struct { + Token struct { + ExpiresAt string `json:"expires_at"` + } `json:"token"` + } + if json.Unmarshal(body, &parsed) == nil && parsed.Token.ExpiresAt != "" { + if t, err := time.Parse(time.RFC3339Nano, parsed.Token.ExpiresAt); err == nil { + exp = t + } + } + return tok, exp, nil +} diff --git a/internal/selectel/client.go b/internal/selectel/client.go new file mode 100644 index 0000000..15d403b --- /dev/null +++ b/internal/selectel/client.go @@ -0,0 +1,308 @@ +// Package selectel — клиент Selectel DNS API v2 (https://api.selectel.ru/domains/v2). +// +// Аутентификация: IAM-токен проекта в заголовке X-Auth-Token. +// Имена в API — FQDN с точкой в конце; клиент снаружи оперирует +// нормализованными именами (lower-case, без точки) и добавляет точку сам. +package selectel + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "net/http" + "net/url" + "strconv" + "strings" + "sync" + "time" + + "github.com/realmanual/traefik-selectel/internal/httpx" +) + +// DefaultBaseURL — DNS-хостинг v2 (единый для всех регионов). +const DefaultBaseURL = "https://api.selectel.ru/domains/v2" + +const ( + pageLimit = 1000 + maxPages = 1000 + zoneCacheTTL = 10 * time.Minute +) + +// ErrZoneNotFound — зона не найдена в проекте. +var ErrZoneNotFound = errors.New("selectel: зона не найдена") + +// Record — значение записи RRSet. +type Record struct { + Content string `json:"content"` + Disabled bool `json:"disabled"` +} + +// RRSet — набор записей одного имени и типа. +type RRSet struct { + ID string `json:"id"` + Name string `json:"name"` // нормализованное: lower-case, без точки в конце + Type string `json:"type"` + TTL int `json:"ttl"` + Records []Record `json:"records"` + Comment string `json:"comment"` + ManagedBy string `json:"managed_by,omitempty"` +} + +// Client — клиент DNS API. +type Client struct { + baseURL string + tokens TokenSource + http *http.Client + retry httpx.Retrier + now func() time.Time + + mu sync.Mutex + zones map[string]zoneEntry +} + +type zoneEntry struct { + id string + expires time.Time +} + +// NewClient создаёт клиент. baseURL пустой — DefaultBaseURL. +func NewClient(baseURL string, tokens TokenSource, timeout time.Duration) *Client { + if baseURL == "" { + baseURL = DefaultBaseURL + } + return &Client{ + baseURL: strings.TrimRight(baseURL, "/"), + tokens: tokens, + http: &http.Client{Timeout: timeout}, + retry: httpx.DefaultRetrier(), + now: time.Now, + zones: map[string]zoneEntry{}, + } +} + +// APIError — неуспешный ответ API. +type APIError struct { + Method string + Path string + Status int + Body string +} + +func (e *APIError) Error() string { + return fmt.Sprintf("selectel: %s %s: статус %d: %s", e.Method, e.Path, e.Status, e.Body) +} + +// do выполняет запрос с токеном; при 401 один раз сбрасывает токен и повторяет. +// Возвращает тело ответа (уже прочитанное) при 2xx. +func (c *Client) do(ctx context.Context, method, path string, query url.Values, payload any) ([]byte, error) { + var raw []byte + if payload != nil { + var err error + if raw, err = json.Marshal(payload); err != nil { + return nil, err + } + } + full := c.baseURL + path + if len(query) > 0 { + full += "?" + query.Encode() + } + for attempt := 0; attempt < 2; attempt++ { + token, err := c.tokens.Token(ctx) + if err != nil { + return nil, err + } + resp, err := c.retry.Do(ctx, c.http, func() (*http.Request, error) { + var body *bytes.Reader + if raw != nil { + body = bytes.NewReader(raw) + } else { + body = bytes.NewReader(nil) + } + req, err := http.NewRequest(method, full, body) + if err != nil { + return nil, err + } + req.Header.Set("X-Auth-Token", token) + req.Header.Set("Accept", "application/json") + if raw != nil { + req.Header.Set("Content-Type", "application/json") + } + return req, nil + }) + if err != nil { + return nil, fmt.Errorf("selectel: %s %s: %w", method, path, err) + } + data, rerr := httpx.ReadBody(resp) + resp.Body.Close() + if rerr != nil { + return nil, fmt.Errorf("selectel: %s %s: чтение ответа: %w", method, path, rerr) + } + if resp.StatusCode == http.StatusUnauthorized && attempt == 0 { + c.tokens.Invalidate() + continue + } + if resp.StatusCode < 200 || resp.StatusCode > 299 { + return nil, &APIError{Method: method, Path: path, Status: resp.StatusCode, Body: httpx.Snippet(data, 300)} + } + return data, nil + } + return nil, &APIError{Method: method, Path: path, Status: http.StatusUnauthorized, Body: "токен отклонён"} +} + +func normalize(name string) string { + return strings.ToLower(strings.TrimSuffix(strings.TrimSpace(name), ".")) +} + +func fqdn(name string) string { return normalize(name) + "." } + +type listResponse[T any] struct { + Count int `json:"count"` + NextOffset int `json:"next_offset"` + Result []T `json:"result"` +} + +type zoneDTO struct { + ID string `json:"id"` + Name string `json:"name"` +} + +// zoneID возвращает id зоны по имени (с кэшем). +func (c *Client) zoneID(ctx context.Context, zone string, force bool) (string, error) { + zone = normalize(zone) + c.mu.Lock() + e, ok := c.zones[zone] + c.mu.Unlock() + if ok && !force && c.now().Before(e.expires) { + return e.id, nil + } + offset := 0 + for page := 0; page < maxPages; page++ { + q := url.Values{} + q.Set("limit", strconv.Itoa(pageLimit)) + q.Set("offset", strconv.Itoa(offset)) + q.Set("filter", zone) + data, err := c.do(ctx, http.MethodGet, "/zones", q, nil) + if err != nil { + return "", err + } + var lr listResponse[zoneDTO] + if err := json.Unmarshal(data, &lr); err != nil { + return "", fmt.Errorf("selectel: разбор списка зон: %w", err) + } + for _, z := range lr.Result { + if normalize(z.Name) == zone { + c.mu.Lock() + c.zones[zone] = zoneEntry{id: z.ID, expires: c.now().Add(zoneCacheTTL)} + c.mu.Unlock() + return z.ID, nil + } + } + if len(lr.Result) < pageLimit { + break + } + offset += len(lr.Result) + } + return "", fmt.Errorf("%w: %s", ErrZoneNotFound, zone) +} + +func (c *Client) forgetZone(zone string) { + c.mu.Lock() + delete(c.zones, normalize(zone)) + c.mu.Unlock() +} + +// withZone выполняет fn с id зоны; при 404 сбрасывает кэш и повторяет один раз. +func (c *Client) withZone(ctx context.Context, zone string, fn func(id string) error) error { + id, err := c.zoneID(ctx, zone, false) + if err != nil { + return err + } + err = fn(id) + var apiErr *APIError + if errors.As(err, &apiErr) && apiErr.Status == http.StatusNotFound { + c.forgetZone(zone) + if id, err = c.zoneID(ctx, zone, true); err != nil { + return err + } + return fn(id) + } + return err +} + +// ListRRSets возвращает все RRSet зоны (все типы; фильтрация — на стороне вызывающего). +func (c *Client) ListRRSets(ctx context.Context, zone string) ([]RRSet, error) { + var out []RRSet + err := c.withZone(ctx, zone, func(id string) error { + out = out[:0] + offset := 0 + for page := 0; page < maxPages; page++ { + q := url.Values{} + q.Set("limit", strconv.Itoa(pageLimit)) + q.Set("offset", strconv.Itoa(offset)) + data, err := c.do(ctx, http.MethodGet, "/zones/"+url.PathEscape(id)+"/rrset", q, nil) + if err != nil { + return err + } + var lr listResponse[RRSet] + if err := json.Unmarshal(data, &lr); err != nil { + return fmt.Errorf("selectel: разбор списка rrset: %w", err) + } + for _, rr := range lr.Result { + rr.Name = normalize(rr.Name) + rr.Type = strings.ToUpper(rr.Type) + out = append(out, rr) + } + if len(lr.Result) < pageLimit { + return nil + } + offset += len(lr.Result) + } + return fmt.Errorf("selectel: слишком много страниц rrset") + }) + if err != nil { + return nil, err + } + return out, nil +} + +type rrsetBody struct { + Name string `json:"name,omitempty"` + Type string `json:"type,omitempty"` + TTL int `json:"ttl"` + Records []Record `json:"records"` + Comment string `json:"comment"` +} + +// CreateRRSet создаёт RRSet (name — без точки на конце, точка добавляется клиентом). +func (c *Client) CreateRRSet(ctx context.Context, zone string, rr RRSet) error { + body := rrsetBody{Name: fqdn(rr.Name), Type: strings.ToUpper(rr.Type), TTL: rr.TTL, Records: rr.Records, Comment: rr.Comment} + return c.withZone(ctx, zone, func(id string) error { + _, err := c.do(ctx, http.MethodPost, "/zones/"+url.PathEscape(id)+"/rrset", nil, body) + return err + }) +} + +// UpdateRRSet обновляет ttl, records и comment существующего RRSet. +func (c *Client) UpdateRRSet(ctx context.Context, zone, rrsetID string, ttl int, records []Record, comment string) error { + body := rrsetBody{TTL: ttl, Records: records, Comment: comment} + return c.withZone(ctx, zone, func(id string) error { + _, err := c.do(ctx, http.MethodPatch, "/zones/"+url.PathEscape(id)+"/rrset/"+url.PathEscape(rrsetID), nil, body) + return err + }) +} + +// DeleteRRSet удаляет RRSet. +// Ответ 404 считается успехом (идемпотентность). +func (c *Client) DeleteRRSet(ctx context.Context, zone, rrsetID string) error { + err := c.withZone(ctx, zone, func(id string) error { + _, err := c.do(ctx, http.MethodDelete, "/zones/"+url.PathEscape(id)+"/rrset/"+url.PathEscape(rrsetID), nil, nil) + return err + }) + var apiErr *APIError + if errors.As(err, &apiErr) && apiErr.Status == http.StatusNotFound { + return nil + } + return err +} diff --git a/internal/selectel/selectel_test.go b/internal/selectel/selectel_test.go new file mode 100644 index 0000000..af3bb18 --- /dev/null +++ b/internal/selectel/selectel_test.go @@ -0,0 +1,288 @@ +package selectel + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/http/httptest" + "strings" + "sync" + "sync/atomic" + "testing" + "time" +) + +type staticTokens struct { + tok string + invalidated int32 +} + +func (s *staticTokens) Token(context.Context) (string, error) { return s.tok, nil } +func (s *staticTokens) Invalidate() { atomic.AddInt32(&s.invalidated, 1) } + +func testCreds() Credentials { + return Credentials{Username: "u", Password: "secret-pass", AccountID: "123456", ProjectID: "proj-1"} +} + +func TestCredentialsStringHidesPassword(t *testing.T) { + s := testCreds().String() + if strings.Contains(s, "secret-pass") { + t.Fatalf("пароль утёк в String(): %s", s) + } + if fmt := (Credentials{}).Validate(); fmt == nil { + t.Fatal("ожидалась ошибка валидации") + } +} + +func TestIAMTokenSourceRequestAndCache(t *testing.T) { + var calls int32 + var gotBody map[string]any + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/identity/v3/auth/tokens" || r.Method != http.MethodPost { + t.Errorf("%s %s", r.Method, r.URL.Path) + } + atomic.AddInt32(&calls, 1) + b, _ := io.ReadAll(r.Body) + _ = json.Unmarshal(b, &gotBody) + w.Header().Set("X-Subject-Token", "tok-"+fmt.Sprint(atomic.LoadInt32(&calls))) + w.WriteHeader(http.StatusCreated) + _, _ = w.Write([]byte(`{"token":{"expires_at":"2030-01-01T00:00:00.000000Z"}}`)) + })) + defer srv.Close() + + ts, err := NewIAMTokenSource(testCreds(), srv.URL+"/identity/v3/", time.Second) + if err != nil { + t.Fatal(err) + } + now := time.Date(2029, 1, 1, 0, 0, 0, 0, time.UTC) + ts.now = func() time.Time { return now } + + tok, err := ts.Token(context.Background()) + if err != nil || tok != "tok-1" { + t.Fatalf("tok=%q err=%v", tok, err) + } + if tok2, _ := ts.Token(context.Background()); tok2 != "tok-1" || calls != 1 { + t.Fatalf("токен не закэширован: %q calls=%d", tok2, calls) + } + + // структура запроса + auth := gotBody["auth"].(map[string]any) + user := auth["identity"].(map[string]any)["password"].(map[string]any)["user"].(map[string]any) + if user["name"] != "u" || user["password"] != "secret-pass" || user["domain"].(map[string]any)["name"] != "123456" { + t.Errorf("user = %v", user) + } + proj := auth["scope"].(map[string]any)["project"].(map[string]any) + if proj["id"] != "proj-1" || proj["domain"].(map[string]any)["name"] != "123456" { + t.Errorf("project = %v", proj) + } + + // обновление за margin до истечения + now = time.Date(2029, 12, 31, 23, 55, 0, 0, time.UTC) + if tok3, _ := ts.Token(context.Background()); tok3 != "tok-2" || calls != 2 { + t.Fatalf("токен не обновлён: %q calls=%d", tok3, calls) + } + + ts.Invalidate() + if tok4, _ := ts.Token(context.Background()); tok4 != "tok-3" { + t.Fatalf("после Invalidate: %q", tok4) + } +} + +func TestIAMTokenSourceErrorDoesNotLeakBody(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusUnauthorized) + _, _ = w.Write([]byte("secret-pass echoed")) + })) + defer srv.Close() + ts, _ := NewIAMTokenSource(testCreds(), srv.URL, time.Second) + _, err := ts.Token(context.Background()) + if err == nil || strings.Contains(err.Error(), "secret-pass") { + t.Fatalf("err = %v", err) + } +} + +// fakeDNS — минимальный эмулятор DNS API. +type fakeDNS struct { + mu sync.Mutex + t *testing.T + requests []string + bodies map[string]string + rrsets []map[string]any + fail401 int32 +} + +func (f *fakeDNS) handler(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + if r.Header.Get("X-Auth-Token") != "tok" { + w.WriteHeader(http.StatusUnauthorized) + return + } + if atomic.LoadInt32(&f.fail401) > 0 { + atomic.AddInt32(&f.fail401, -1) + w.WriteHeader(http.StatusUnauthorized) + return + } + f.requests = append(f.requests, r.Method+" "+r.URL.Path) + b, _ := io.ReadAll(r.Body) + if len(b) > 0 { + if f.bodies == nil { + f.bodies = map[string]string{} + } + f.bodies[r.Method+" "+r.URL.Path] = string(b) + } + switch { + case r.Method == http.MethodGet && r.URL.Path == "/zones": + _ = json.NewEncoder(w).Encode(map[string]any{"count": 2, "next_offset": 0, "result": []map[string]any{ + {"id": "z-other", "name": "notexample.com."}, + {"id": "z1", "name": "example.com."}, + }}) + case r.Method == http.MethodGet && r.URL.Path == "/zones/z1/rrset": + _ = json.NewEncoder(w).Encode(map[string]any{"count": len(f.rrsets), "next_offset": 0, "result": f.rrsets}) + case r.Method == http.MethodPost && r.URL.Path == "/zones/z1/rrset": + _, _ = w.Write([]byte(`{"id":"new"}`)) + case r.Method == http.MethodPatch && strings.HasPrefix(r.URL.Path, "/zones/z1/rrset/"): + w.WriteHeader(http.StatusNoContent) + case r.Method == http.MethodDelete && r.URL.Path == "/zones/z1/rrset/gone": + w.WriteHeader(http.StatusNotFound) + case r.Method == http.MethodDelete && strings.HasPrefix(r.URL.Path, "/zones/z1/rrset/"): + w.WriteHeader(http.StatusNoContent) + default: + w.WriteHeader(http.StatusNotFound) + } +} + +func newTestClient(t *testing.T, f *fakeDNS) (*Client, *staticTokens) { + srv := httptest.NewServer(http.HandlerFunc(f.handler)) + t.Cleanup(srv.Close) + tk := &staticTokens{tok: "tok"} + c := NewClient(srv.URL, tk, time.Second) + c.retry.Sleep = func(context.Context, time.Duration) error { return nil } + return c, tk +} + +func TestListRRSetsNormalizesAndResolvesZone(t *testing.T) { + f := &fakeDNS{t: t, rrsets: []map[string]any{ + {"id": "r1", "name": "App.Example.com.", "type": "A", "ttl": 300, "comment": "managed-by=x", + "records": []map[string]any{{"content": "1.2.3.4", "disabled": false}}}, + }} + c, _ := newTestClient(t, f) + got, err := c.ListRRSets(context.Background(), "Example.com.") + if err != nil { + t.Fatal(err) + } + if len(got) != 1 || got[0].Name != "app.example.com" || got[0].Type != "A" || got[0].Records[0].Content != "1.2.3.4" || got[0].Comment != "managed-by=x" { + t.Fatalf("got = %+v", got) + } + // зона закэширована: второй вызов не ходит в /zones + if _, err := c.ListRRSets(context.Background(), "example.com"); err != nil { + t.Fatal(err) + } + zoneCalls := 0 + for _, r := range f.requests { + if r == "GET /zones" { + zoneCalls++ + } + } + if zoneCalls != 1 { + t.Fatalf("GET /zones вызван %d раз", zoneCalls) + } +} + +func TestZoneNotFound(t *testing.T) { + f := &fakeDNS{t: t} + c, _ := newTestClient(t, f) + _, err := c.ListRRSets(context.Background(), "missing.org") + if !errors.Is(err, ErrZoneNotFound) { + t.Fatalf("err = %v", err) + } +} + +func TestCreateUpdateDelete(t *testing.T) { + f := &fakeDNS{t: t} + c, _ := newTestClient(t, f) + ctx := context.Background() + + err := c.CreateRRSet(ctx, "example.com", RRSet{Name: "app.example.com", Type: "a", TTL: 300, + Records: []Record{{Content: "1.2.3.4"}}, Comment: "managed-by=traefik-selectel-dns"}) + if err != nil { + t.Fatal(err) + } + var body map[string]any + if err := json.Unmarshal([]byte(f.bodies["POST /zones/z1/rrset"]), &body); err != nil { + t.Fatal(err) + } + if body["name"] != "app.example.com." || body["type"] != "A" || body["ttl"].(float64) != 300 || body["comment"] != "managed-by=traefik-selectel-dns" { + t.Errorf("create body = %v", body) + } + recs := body["records"].([]any) + if len(recs) != 1 || recs[0].(map[string]any)["content"] != "1.2.3.4" { + t.Errorf("records = %v", recs) + } + + if err := c.UpdateRRSet(ctx, "example.com", "r1", 600, []Record{{Content: "5.6.7.8"}}, "managed-by=traefik-selectel-dns"); err != nil { + t.Fatal(err) + } + var upd map[string]any + _ = json.Unmarshal([]byte(f.bodies["PATCH /zones/z1/rrset/r1"]), &upd) + if _, has := upd["name"]; has || upd["ttl"].(float64) != 600 { + t.Errorf("patch body = %v", upd) + } + + if err := c.DeleteRRSet(ctx, "example.com", "r1"); err != nil { + t.Fatal(err) + } + if err := c.DeleteRRSet(ctx, "example.com", "gone"); err != nil { + t.Fatalf("404 на delete должен быть успехом: %v", err) + } +} + +func TestUnauthorizedInvalidatesTokenAndRetriesOnce(t *testing.T) { + f := &fakeDNS{t: t, fail401: 1} + c, tk := newTestClient(t, f) + if _, err := c.ListRRSets(context.Background(), "example.com"); err != nil { + t.Fatal(err) + } + if tk.invalidated != 1 { + t.Fatalf("invalidated = %d", tk.invalidated) + } +} + +func TestAPIErrorOn422(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/zones" { + _, _ = w.Write([]byte(`{"result":[{"id":"z1","name":"example.com."}]}`)) + return + } + w.WriteHeader(http.StatusUnprocessableEntity) + _, _ = w.Write([]byte(`{"error":"bad"}`)) + })) + defer srv.Close() + c := NewClient(srv.URL, &staticTokens{tok: "tok"}, time.Second) + err := c.CreateRRSet(context.Background(), "example.com", RRSet{Name: "a.example.com", Type: "A", TTL: 300, Records: []Record{{Content: "1.1.1.1"}}}) + var apiErr *APIError + if !errors.As(err, &apiErr) || apiErr.Status != 422 { + t.Fatalf("err = %v", err) + } +} + +func TestRetryOn503(t *testing.T) { + var n int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if atomic.AddInt32(&n, 1) == 1 { + w.WriteHeader(http.StatusServiceUnavailable) + return + } + _, _ = w.Write([]byte(`{"result":[{"id":"z1","name":"example.com."}]}`)) + })) + defer srv.Close() + c := NewClient(srv.URL, &staticTokens{tok: "tok"}, time.Second) + c.retry.Sleep = func(context.Context, time.Duration) error { return nil } + id, err := c.zoneID(context.Background(), "example.com", true) + if err != nil || id != "z1" || n != 2 { + t.Fatalf("id=%q err=%v n=%d", id, err, n) + } +} diff --git a/internal/traefik/client.go b/internal/traefik/client.go new file mode 100644 index 0000000..681e422 --- /dev/null +++ b/internal/traefik/client.go @@ -0,0 +1,90 @@ +// Package traefik — клиент Traefik API (GET /api/http/routers). +package traefik + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/url" + "strconv" + "strings" + "time" + + "github.com/realmanual/traefik-selectel/internal/httpx" +) + +// Router — подмножество полей роутера Traefik. +type Router struct { + Name string `json:"name"` + Rule string `json:"rule"` + Status string `json:"status"` + Provider string `json:"provider"` +} + +// RouterSource — источник роутов (интерфейс для тестов). +type RouterSource interface { + Routers(ctx context.Context) ([]Router, error) +} + +// Client читает роутеры из Traefik API. +type Client struct { + baseURL string + http *http.Client + retry httpx.Retrier +} + +// NewClient создаёт клиент. timeout — таймаут одного HTTP-запроса. +func NewClient(apiURL string, timeout time.Duration) (*Client, error) { + u, err := url.Parse(apiURL) + if err != nil || (u.Scheme != "http" && u.Scheme != "https") || u.Host == "" { + return nil, fmt.Errorf("traefik: некорректный apiURL %q", apiURL) + } + return &Client{ + baseURL: strings.TrimRight(u.String(), "/"), + http: &http.Client{Timeout: timeout}, + retry: httpx.DefaultRetrier(), + }, nil +} + +const perPage = 500 +const maxPages = 100 + +// Routers возвращает все HTTP-роутеры (с учётом пагинации X-Next-Page). +func (c *Client) Routers(ctx context.Context) ([]Router, error) { + var all []Router + page := 1 + for i := 0; i < maxPages; i++ { + u := fmt.Sprintf("%s/api/http/routers?page=%d&per_page=%d", c.baseURL, page, perPage) + resp, err := c.retry.Do(ctx, c.http, func() (*http.Request, error) { + req, err := http.NewRequest(http.MethodGet, u, nil) + if err != nil { + return nil, err + } + req.Header.Set("Accept", "application/json") + return req, nil + }) + if err != nil { + return nil, fmt.Errorf("traefik: запрос роутеров: %w", err) + } + body, err := httpx.ReadBody(resp) + resp.Body.Close() + if err != nil { + return nil, fmt.Errorf("traefik: чтение ответа: %w", err) + } + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("traefik: статус %d: %s", resp.StatusCode, httpx.Snippet(body, 200)) + } + var chunk []Router + if err := json.Unmarshal(body, &chunk); err != nil { + return nil, fmt.Errorf("traefik: разбор JSON: %w", err) + } + all = append(all, chunk...) + next, _ := strconv.Atoi(resp.Header.Get("X-Next-Page")) + if next <= page { + return all, nil + } + page = next + } + return nil, fmt.Errorf("traefik: слишком много страниц роутеров") +} diff --git a/internal/traefik/client_test.go b/internal/traefik/client_test.go new file mode 100644 index 0000000..0a2adbd --- /dev/null +++ b/internal/traefik/client_test.go @@ -0,0 +1,58 @@ +package traefik + +import ( + "context" + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestRoutersPagination(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/api/http/routers" { + t.Errorf("path = %s", r.URL.Path) + } + switch r.URL.Query().Get("page") { + case "1": + w.Header().Set("X-Next-Page", "2") + _, _ = w.Write([]byte(`[{"name":"a@docker","rule":"Host(` + "`a.example.com`" + `)","status":"enabled","provider":"docker"}]`)) + case "2": + w.Header().Set("X-Next-Page", "2") + _, _ = w.Write([]byte(`[{"name":"b@file","rule":"Host(` + "`b.example.com`" + `)","status":"disabled","provider":"file"}]`)) + default: + t.Errorf("unexpected page %q", r.URL.Query().Get("page")) + } + })) + defer srv.Close() + c, err := NewClient(srv.URL, time.Second) + if err != nil { + t.Fatal(err) + } + rs, err := c.Routers(context.Background()) + if err != nil { + t.Fatal(err) + } + if len(rs) != 2 || rs[0].Name != "a@docker" || rs[1].Status != "disabled" { + t.Fatalf("routers = %+v", rs) + } +} + +func TestRoutersErrorStatus(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusForbidden) + })) + defer srv.Close() + c, _ := NewClient(srv.URL, time.Second) + if _, err := c.Routers(context.Background()); err == nil { + t.Fatal("ожидалась ошибка") + } +} + +func TestNewClientInvalidURL(t *testing.T) { + for _, u := range []string{"", "traefik:8080", "ftp://x"} { + if _, err := NewClient(u, time.Second); err == nil { + t.Errorf("%q: ожидалась ошибка", u) + } + } +} diff --git a/swarm-report/traefik-selectel-dns-2026-09-21.md b/swarm-report/traefik-selectel-dns-2026-09-21.md new file mode 100644 index 0000000..b86747d --- /dev/null +++ b/swarm-report/traefik-selectel-dns-2026-09-21.md @@ -0,0 +1,38 @@ +# traefik-selectel-dns — отчёт (2026-09-21) + +## Что сделано +Go-демон `traefik-selectel-dns` (модуль `github.com/realmanual/traefik-selectel`, Go 1.27.1, единственная зависимость — `gopkg.in/yaml.v3`): +- `internal/hostparse` — разбор `Host(...)` (несколько аргументов, `||`, `&&`, кавычки, отрицания `!Host`/`!(...)`, `HostRegexp`/`HostSNI` не считаются), валидация DNS-имён, сопоставление зон (longest match). +- `internal/traefik` — клиент `GET /api/http/routers` с пагинацией (`page`/`per_page`, `X-Next-Page`). +- `internal/selectel` — Keystone-токен (кэш в памяти, обновление за 10 минут до `expires_at`), DNS API v2 (зоны с кэшем id, rrset list/create/PATCH/delete), повтор при 401 со сбросом токена, 404 на delete = успех. +- `internal/reconciler` — create/update/skip-foreign/skip-CNAME/delete-grace/safe-guard, dry-run, цикл с debounce по событиям Docker. +- `internal/dockerevents` — raw HTTP `/events` по unix-сокету или tcp (без docker SDK), переподключение с backoff, недоступность не фатальна. +- `internal/config` (YAML strict + env `TSD_*` + секреты `*_FILE`), `internal/httpx` (ретраи 429/5xx с backoff и Retry-After), `internal/health` (`/healthz`, флаг `-healthcheck`). +- Dockerfile (multi-stage, distroless static nonroot 65532), `.dockerignore`, `deploy/` (compose, traefik.yml, config, dynamic/, whoami-пример), README. + +Проверено: `go vet ./...` — чисто; `go test -race ./...` — 79 тестов проходят; `docker build` собирается, образ запускается (`-version`), USER 65532; `docker compose config -q` валиден. + +## Сверка с документацией +- DNS API v2: base URL `https://api.selectel.ru/domains/v2` (docs.selectel.ru/api/urls), `X-Auth-Token` (по исходникам lego selectelv2), пути `/zones`, `/zones/{id}/rrset[/{rrset_id}]`, тело `{name,type,ttl,records:[{content,disabled}],comment}`, PATCH — `ttl`, `records`, `comment`; имена — FQDN с точкой. +- Токен: `POST https://cloud.api.selcloud.ru/identity/v3/auth/tokens`, токен в заголовке `X-Subject-Token`, срок 24 ч. +- Переменные lego `selectelv2`: `SELECTELV2_USERNAME/PASSWORD/ACCOUNT_ID/PROJECT_ID` (+`_FILE`) — сверено через context7 (/go-acme/lego). + +## Принятые решения / отступления +- Роутеры со статусом `warning` не создают записи, но считаются «присутствующими» (не запускают grace-удаление) — чтобы временный warning не удалил запись. +- Список rrset читается целиком (все типы) и фильтруется на клиенте: формат query-параметра `rrset_types` в доке не уточнён; так же позволяет обнаружить CNAME-конфликт. +- Зона ищется по `filter` + точное сравнение имени на клиенте (семантика filter не описана). +- Traefik API: вместо `api.insecure: true` — роутер `api@internal` на entrypoint `traefik` с `ipAllowList` подсети control, т.к. Traefik также сидит в сети proxy с пользовательскими контейнерами. Docker events демон получает через `docker-socket-proxy`, а не прямой сокет. +- Добавлены сверх ТЗ: `dryRun`, `-healthcheck`, детекция CNAME-конфликта. + +## Риски +- Остаточный: подсеть `172.30.250.0/24` дублируется в compose и dynamic/api-internal.yml — при расхождении демон получит 403. +- Docker socket, смонтированный в Traefik (штатно для Docker provider) — root-эквивалент; в примере ro-монтирование не ограничивает API. +- POST при 5xx может повториться (возможен дубль) — reconcile идемпотентен, но при дубле API вернёт ошибку до следующего цикла. +- Метка владения — exact-match комментария (до 255 символов ок). +- Образы `traefik:v3`, `docker-socket-proxy:latest`, `golang:1.27-alpine`, distroless — не закреплены по digest; в проде закрепить. +- Имя образа `ghcr.io/realmanual/traefik-selectel-dns` в compose — заглушка. + +## Не проверено +- Реальный Selectel API (нет кредов): формат ответа Keystone (`expires_at`; при ошибке разбора берётся 1 час), точное поведение `filter` у зон, поведение POST на дубликат, признак конца пагинации (используется `len(result) < limit`), значение `managed_by` у записей, созданных через API. +- Реальный Traefik (структура роутеров, пагинация, статус `warning`), реальное Docker events, выпуск сертификатов и wildcard `tls.domains` — не запускалось, деплой не выполнялся. +- Не проверялись права сервисного пользователя Selectel (минимальная роль на DNS).