From cca899e65904533998e06b6bb98e44eb0034031f Mon Sep 17 00:00:00 2001 From: Magnus Root Date: Fri, 17 Apr 2026 11:07:49 +0300 Subject: [PATCH] Added database_schema --- .gitignore | 46 +++++ ARCHITECTURE.md | 149 +++++++++++++++ DATABASE_SCHEMA.md | 69 +++++++ README.md | 230 ++++++++++++++++++++++++ api.go | 86 +++++++++ compare.go | 148 +++++++++++++++ config.go | 125 +++++++++++++ database.go | 440 +++++++++++++++++++++++++++++++++++++++++++++ dns.go | 24 +++ git.go | 78 ++++++++ go.mod | 10 ++ go.sum | 6 + main.go | 80 +++++++++ notify.go | 315 ++++++++++++++++++++++++++++++++ whitelist.go | 314 ++++++++++++++++++++++++++++++++ 15 files changed, 2120 insertions(+) create mode 100644 .gitignore create mode 100644 ARCHITECTURE.md create mode 100644 DATABASE_SCHEMA.md create mode 100644 README.md create mode 100644 api.go create mode 100644 compare.go create mode 100644 config.go create mode 100644 database.go create mode 100644 dns.go create mode 100644 git.go create mode 100644 go.mod create mode 100644 go.sum create mode 100644 main.go create mode 100644 notify.go create mode 100644 whitelist.go diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..c9d81b5 --- /dev/null +++ b/.gitignore @@ -0,0 +1,46 @@ +# Binaries +sp_sync +*.exe +*.exe~ +*.dll +*.so +*.dylib + +# Credentials +*.env + +# Test binary, built with `go test -c` +*.test + +# Output of the go coverage tool +*.out + +# Go workspace file +go.work + +# Dependency directories +vendor/ + +# IDE +.vscode/ +.idea/ +*.swp +*.swo +*~ + +# OS +.DS_Store +Thumbs.db + +# Logs +*.log + +# Config with sensitive data +config.yaml +config.yml +*.local.yaml +*.local.yml + +# Temporary files +*.tmp +*.temp diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md new file mode 100644 index 0000000..bd97662 --- /dev/null +++ b/ARCHITECTURE.md @@ -0,0 +1,149 @@ +# Архитектура API Sync + +## Обзор + +API Sync — однопроходная утилита (не демон), запускаемая по cron. При каждом запуске выполняет три последовательных этапа: синхронизацию данных из API в PostgreSQL, очистку устаревших записей и обновление whitelist-файла в Git-репозитории. + +--- + +## Поток выполнения + +``` +main() + │ + ├─ 1. Подключение к PostgreSQL + │ + ├─ 2. Загрузка списка клиентов (client_info) + │ + ├─ 3. Сбор SID (apps_settings WHERE mode = 'auto') + │ + ├─ 4. Параллельная обработка SID (worker pool) + │ └─ processSID() для каждого SID + │ ├─ GET /l7/resource/{sid}/global → whois + │ ├─ GET /l7/origin/global?l7ResourceId={sid} → origins + │ ├─ GET /l7/alias/global?l7ResourceId={sid} → aliases + │ ├─ DNS резолв домена → проверка AntiDDOS + │ └─ upsertAPIInfo() → сравнение + запись в API_info + алерты + │ + ├─ 5. Очистка устаревших SID из API_info (cleanupRemovedSIDs) + │ + └─ 6. Обновление WAF whitelist (updateWAFWhitelist) + ├─ git pull + ├─ Загрузка IP из API_info (auto) + ├─ Загрузка IP из manual_info (ручные) + ├─ Фильтрация WAF-сетей (таблица ips) + ├─ Сравнение с текущим whitelist-файлом + ├─ Запись файла + git commit + git push + └─ Отправка алерта об изменениях +``` + +--- + +## Модули + +### main.go — точка входа и worker pool + +Инициализирует подключение к БД, собирает все SID для обработки и запускает пул горутин. Количество воркеров задаётся через `MAX_CONCURRENT_WORKERS` (по умолчанию 5). SID передаются воркерам через буферизированный канал. + +После завершения всех воркеров последовательно запускаются очистка устаревших записей и обновление whitelist. + +### config.go — конфигурация + +Загружает переменные окружения из `/etc/API_sync/API_sync.env` через `godotenv`. Если файл недоступен — использует системные переменные окружения. Все переменные инициализируются в `init()` до старта `main()`. + +### api.go — HTTP-клиент и структуры API + +Содержит единственный HTTP-клиент с таймаутом 20 секунд и функцию `apiGet()`, которая выполняет GET-запрос к API API с Bearer-авторизацией и десериализует JSON-ответ в переданную структуру. + +Определены структуры ответов API: `originItem`, `aliasItem`, `whoisReAPIonse`, `apiList`. + +### database.go — работа с PostgreSQL + +**`loadClients()`** — загружает список клиентов из `client_info`. + +**`loadL7IDs()`** — загружает SID из `apps_settings` для конкретного клиента с `mode = 'auto'`. + +**`processSID()`** — оркестрирует обработку одного SID: делает три запроса к API, резолвит домен для проверки AntiDDOS, вызывает `upsertAPIInfo()`. + +**`upsertAPIInfo()`** — сравнивает новые данные со старыми из БД и при наличии изменений формирует алерты: +- изменение домена → алерт + флаг необходимости обновления PT AF +- изменение origins (IP, mode, weight) — через `compareOrigins()` +- изменение aliases — через `compareAliases()` +- изменение WAF-настроек (vendor, instance, enabled) — отдельный алерт + +Сохранение выполняется через `INSERT ... ON CONFLICT DO UPDATE` (upsert по `sid`). + +**`cleanupRemovedSIDs()`** — сравнивает все SID в `API_info` с актуальным списком активных SID. Записи, которых нет в активном списке (ресурс удалён из `apps_settings` или сменил `mode`), удаляются из `API_info`. По каждому удалённому ресурсу отправляется алерт с указанием WAF-вендора (берётся из `apps_settings` через LEFT JOIN). + +**`wafVendorWarning()`** — формирует строку предупреждения в зависимости от вендора: `ptaf` → PT AF, `sw` → SW, `wmx` → WMX. + +### compare.go — функции сравнения + +**`compareOrigins()`** — сравнивает два списка origins по IP-адресу. Определяет добавленные, удалённые и изменённые (mode/weight) записи. Возвращает отформатированную строку для алерта с текущим состоянием origins. + +**`compareAliases()`** — сравнивает два отсортированных списка доменов. Возвращает строку с изменениями и булев флаг наличия изменений. + +### dns.go — DNS-резолвер + +Функция `resolveDomain()` возвращает первый IPv4-адрес для заданного домена. Используется в `processSID()` для проверки AntiDDOS: если resolvedIP совпадает с protectedIP из API — AntiDDOS активен. + +### notify.go — система уведомлений + +**`sendAlert()`** — единая точка отправки алертов. Последовательно отправляет сообщение во все настроенные каналы. Ошибка одного канала не блокирует остальные. + +**Telegram** — `buildTelegramClient()` при каждой отправке выбирает транспорт в порядке приоритета: +1. HTTP-прокси (`TELEGRAM_HTTP_PROXY`) — если задан и доступен +2. SOCKS5-прокси (`TELEGRAM_SOCKS5_PROXY`) — если задан и доступен +3. Прямое подключение — fallback + +Доступность проверяется реальным GET-запросом на `api.telegram.org`. Сообщения отправляются в формате HTML (`parse_mode: HTML`). + +**Mattermost** — POST на `/api/v4/posts` с Bearer-токеном. HTML-теги Telegram конвертируются в Markdown через `htmlToMarkdown()`. Отправляется только если заполнены все три переменные: `MATTERMOST_URL`, `MATTERMOST_BOT_TOKEN`, `MATTERMOST_CHANNEL_ID`. + +**Email** — через стандартный `net/smtp`. Поддерживает несколько адресов получателей через запятую в `EMAIL_TO`. Авторизация опциональна — если `EMAIL_SMTP_USER` пуст, отправка идёт без аутентификации (relay). HTML-теги Telegram удаляются через `stripHTMLTags()`. Отправляется только если заполнены `EMAIL_SMTP_HOST`, `EMAIL_FROM`, `EMAIL_TO`. + +### git.go — Git-операции + +**`syncGitRepo()`** — если репозиторий не существует локально — клонирует (`git clone`), иначе обновляет (`git pull`). Автоматически настраивает `user.email` и `user.name` для коммитов. + +**`gitCommitAndPush()`** — выполняет `git add `, `git commit -m `, `git push`. Если нечего коммитить (`nothing to commit`) — завершается без ошибки. + +### whitelist.go — управление whitelist + +**`updateWAFWhitelist()`** — основная функция обновления whitelist: +1. Синхронизирует Git-репозиторий +2. Загружает WAF-сети из таблицы `ips` +3. Загружает origin IP из `API_info` (`loadAllOrigins`) и `manual_info` (`loadManualOrigins`), объединяет и дедуплицирует +4. Фильтрует IP попадающие в WAF-сети (`filterNonWAFOrigins`) +5. Сравнивает с текущим содержимым файла (`compareWhitelists`) +6. При наличии изменений — перезаписывает файл, коммитит и пушит, отправляет алерт + +**`loadManualOrigins()`** — загружает все записи из `manual_info` как есть, без дополнительных условий. Таблица заполняется администратором вручную и программой не изменяется. + +**`filterNonWAFOrigins()`** — исключает из списка IP-адреса, которые принадлежат WAF-сетям из таблицы `ips`. Такие адреса являются узлами самого WAF и не должны попадать в whitelist. + +--- + +## Алерты + +Все алерты отправляются через единую функцию `sendAlert()`. Типы алертов: + +| Событие | Содержание | +|---|---| +| Изменение origins | Добавленные/удалённые/изменённые backend IP + текущее состояние | +| Изменение aliases | Добавленные/удалённые домены | +| Изменение домена | Старое и новое значение + предупреждение об обновлении WAF | +| Изменение WAF-настроек | Изменения waf_enabled, waf_vendor, waf_instance | +| Удаление ресурса | SID, домен, предупреждение с указанием конкретного WAF-вендора | +| Обновление whitelist | Добавленные/удалённые IP с привязкой к SID и домену, итоговое количество | + +--- + +## Зависимости + +| Пакет | Назначение | +|---|---| +| `github.com/lib/pq` | PostgreSQL драйвер | +| `github.com/joho/godotenv` | Загрузка `.env` файла | +| `golang.org/x/net` | SOCKS5 прокси для Telegram | +| `net/smtp` | Отправка email (стандартная библиотека) | diff --git a/DATABASE_SCHEMA.md b/DATABASE_SCHEMA.md new file mode 100644 index 0000000..5921415 --- /dev/null +++ b/DATABASE_SCHEMA.md @@ -0,0 +1,69 @@ +# Схема базы данных + +База данных PostgreSQL. Имя БД задаётся переменной окружения `DB_NAME` (по умолчанию `waf_info`). + +--- + +## client_info + +Список клиентов, для которых выполняется синхронизация. + +| Колонка | Тип | Описание | +|---|---|---| +| `client_title` | text | Уникальное название клиента | + +--- + +## apps_settings + +Настройки L7-ресурсов по каждому клиенту. Определяет, какие ресурсы и в каком режиме обрабатываются. + +| Колонка | Тип | Описание | +|---|---|---| +| `client_title` | text | Название клиента (связь с `client_info`) | +| `l7resourceid` | bigint | ID ресурса в ServicePipe (SID) | +| `mode` | text | Режим обработки: `auto` — синхронизируется автоматически, остальные значения игнорируются при синхронизации | +| `waf_vendor` | text | Вендор WAF: `ptaf`, `sw`, `wmx` — используется для формирования текста алертов при удалении ресурса | + +--- + +## sp_info + +Основная таблица с данными о ресурсах, полученными из API ServicePipe. Заполняется и обновляется автоматически при каждом запуске синхронизации. + +| Колонка | Тип | Описание | +|---|---|---| +| `sid` | bigint | ID ресурса в ServicePipe (первичный ключ) | +| `domain_name` | text | Основной домен ресурса | +| `origins` | jsonb | Список backend-серверов в формате `[{"ip": "...", "mode": "...", "weight": N}]` | +| `aliases` | jsonb | Список доменных псевдонимов в формате `["domain1", "domain2"]` | +| `aliaces_in_sp_modified` | timestamptz | Дата последнего изменения алиасов на стороне SP | +| `protected_ip` | text | Защищённый IP из API (используется для проверки AntiDDOS) | +| `resolved_ip` | text | IP, полученный при DNS-резолвинге основного домена | +| `antiddos_enable` | boolean | Признак активности AntiDDOS: `true` если `resolved_ip == protected_ip` | +| `waf_enabled_sp` | integer | Признак включённости WAF на стороне SP (из API) | +| `waf_vendor` | text | Вендор WAF (из API ServicePipe) | +| `instance_sp` | text | Название WAF-инстанса на стороне SP | +| `updated_at` | timestamptz | Дата последнего обновления записи | + +--- + +## manual_info + +Ресурсы, добавляемые вручную администратором. Используется только как дополнительный источник IP-адресов для whitelist. Программой не изменяется. + +| Колонка | Тип | Описание | +|---|---|---| +| `sid` | bigint | ID ресурса в ServicePipe | +| `domain_name` | text | Основной домен ресурса | +| `origins` | jsonb | Список backend-серверов в формате `[{"ip": "...", "mode": "...", "weight": N}]` | + +--- + +## ips + +Содержит сети WAF-узлов в виде CIDR-диапазонов. Используется для фильтрации — IP из этих сетей не попадают в whitelist. + +| Колонка | Тип | Описание | +|---|---|---| +| `waf_networks` | text[] | Массив CIDR-диапазонов WAF-сетей, например `{"10.0.0.0/8", "192.168.1.0/24"}` | diff --git a/README.md b/README.md new file mode 100644 index 0000000..f0a039a --- /dev/null +++ b/README.md @@ -0,0 +1,230 @@ +# API Sync - API to PostgreSQL Synchronization Tool + +Утилита для автоматической синхронизации данных WAF из API в PostgreSQL с управлением whitelist для PT AF AntiDDOS. + +## Структура проекта + +``` +├── main.go # Точка входа, worker pool +├── config.go # Конфигурация и переменные окружения +├── api.go # HTTP клиент и API структуры +├── database.go # Работа с PostgreSQL +├── dns.go # DNS резолв для AntiDDOS проверки +├── notify.go # Уведомления (Telegram, Mattermost, Email) +├── git.go # Git операции (clone, commit, push) +├── whitelist.go # Управление WAF whitelist +├── compare.go # Функции сравнения (origins, aliases) +└── API_sync.env # Файл конфигурации (не в Git!) +``` + +## Сборка + +### Установка зависимостей + +```bash +go get github.com/joho/godotenv +go get github.com/lib/pq +go get golang.org/x/net +``` + +### Компиляция + +```bash +# Собрать бинарник +go build -o API_sync + +# Или с оптимизацией размера +go build -ldflags="-s -w" -o API_sync +``` + +Go автоматически найдёт все `.go` файлы в директории и скомпилирует их в один бинарник. + +## Конфигурация + +### Шаг 1: Создать файл конфигурации + +```bash +cp API_sync.env.example API_sync.env +nano API_sync.env +``` + +### Шаг 2: Заполнить переменные + +```env +# API API +BEARER_TOKEN=your_api_token + +# PostgreSQL +DB_HOST=localhost +DB_PORT=5432 +DB_USER=install +DB_PASSWORD=your_password +DB_NAME=waf_info + +# Telegram (обязательно) +TELEGRAM_BOT_TOKEN=123456789:ABC... +TELEGRAM_CHAT_ID=-1001234567890 + +# Telegram прокси (опционально, пробуются в порядке: HTTP → SOCKS5 → прямое подключение) +TELEGRAM_HTTP_PROXY=http://user:pass@proxy.example.com:3128 +TELEGRAM_SOCKS5_PROXY=proxy.example.com:1080 +TELEGRAM_SOCKS5_USER=user +TELEGRAM_SOCKS5_PASSWORD=password + +# Mattermost (опционально) +MATTERMOST_URL=https://mattermost.example.com +MATTERMOST_BOT_TOKEN=your_bot_token +MATTERMOST_CHANNEL_ID=channel_id + +# Email (опционально) +EMAIL_SMTP_HOST=smtp.example.com +EMAIL_SMTP_PORT=587 +EMAIL_SMTP_USER=user@example.com +EMAIL_SMTP_PASSWORD=your_password +EMAIL_FROM=API_sync@example.com +EMAIL_TO=admin@example.com,team@example.com + +# Git репозиторий whitelist +GIT_REPO_URL=https://token@svc-git.test.ru/wmx/waf_whitelist.git +GIT_REPO_PATH=/home/install/waf_whitelist +WHITELIST_FILE=whitelist_ptaf.txt + +# Параллельность +MAX_CONCURRENT_WORKERS=5 +``` + +### Шаг 3: Защитить файл конфигурации + +```bash +chmod 600 API_sync.env +``` + +## Использование + +### Запуск вручную + +```bash +./API_sync +``` + +### Запуск через cron + +```bash +# /etc/cron.d/API_sync +*/15 * * * * install cd /home/install && /home/install/API_sync >> /var/log/API_sync.log 2>&1 +``` + +## Возможности + +✅ **Параллельная обработка** — настраиваемый worker pool для запросов к API + +✅ **Синхронизация данных WAF** из API API в PostgreSQL + +✅ **Smart Diff** — детальное сравнение origins и aliases с уведомлением об изменениях + +✅ **AntiDDOS проверка** — DNS-резолв домена и сравнение с protected_ip из API + +✅ **WAF настройки** — отслеживание изменений waf_enabled, waf_vendor, waf_instance + +✅ **Мультиканальные уведомления** — Telegram (с поддержкой HTTP/SOCKS5 прокси), Mattermost, Email + +✅ **Автоматический whitelist** — объединение IP из API_info и manual_info, фильтрация WAF-сетей, Git push + +✅ **Очистка устаревших данных** — автоматическое удаление из API_info ресурсов, пропавших из apps_settings + +✅ **Режим auto/manual** — только ресурсы с `mode = 'auto'` синхронизируются через API; ресурсы из manual_info используются как дополнительный источник IP для whitelist + +## База данных + +### Необходимые таблицы + +```sql +-- Таблица с информацией из API (заполняется автоматически) +CREATE TABLE API_info ( + sid BIGINT PRIMARY KEY, + domain_name VARCHAR(255), + origins JSONB, + aliases JSONB, + aliaces_in_API_modified TIMESTAMP WITHOUT TIME ZONE, + protected_ip VARCHAR(45), + resolved_ip VARCHAR(45), + antiddos_enable BOOLEAN, + waf_enabled_API INTEGER, + waf_vendor VARCHAR(255), + instance_API VARCHAR(255), + created_at TIMESTAMP DEFAULT NOW(), + updated_at TIMESTAMP DEFAULT NOW() +); + +-- Таблица с WAF-сетями для фильтрации whitelist +CREATE TABLE ips ( + id SERIAL PRIMARY KEY, + waf_networks TEXT[] +); + +-- Пример заполнения WAF-сетей +INSERT INTO ips (waf_networks) VALUES + ('{109.238.89.0/24,89.20.63.0/24}'); + +-- Таблица ручных ресурсов (заполняется администратором вручную) +CREATE TABLE manual_info ( + sid BIGINT PRIMARY KEY, + domain_name VARCHAR(255), + origins JSONB +); + +-- Добавить столбец mode в apps_settings (если ещё не добавлен) +ALTER TABLE apps_settings ADD COLUMN mode VARCHAR(10) DEFAULT 'auto'; + +-- Добавить столбец waf_vendor в apps_settings (если ещё не добавлен) +ALTER TABLE apps_settings ADD COLUMN waf_vendor VARCHAR(50); +``` + +## Логи + +Программа выводит детальные логи: + +``` +2025/12/26 10:00:00 Start sync API_info +2025/12/26 10:00:00 Total SIDs to process: 25 +2025/12/26 10:00:00 Using 5 concurrent workers +2025/12/26 10:00:01 [Worker 0] Processing SID 10307 +2025/12/26 10:00:01 [WAF Info] SID: 10307, WAF Enabled: 1, Vendor: ptaf, Instance: PTAFd_02_03 +2025/12/26 10:00:01 [AntiDDOS Check] ✅ AntiDDOS is ENABLED for nationallottery.ru +2025/12/26 10:00:01 [Telegram] Using SOCKS5 proxy: proxy.example.com:1080 +2025/12/26 10:00:15 Finish sync API_info +2025/12/26 10:00:15 Start cleanup of removed SIDs +2025/12/26 10:00:15 [Cleanup] No stale SIDs found in API_info +2025/12/26 10:00:15 Finish cleanup of removed SIDs +2025/12/26 10:00:15 Start WAF whitelist update +2025/12/26 10:00:16 [WAF Whitelist] No changes needed +2025/12/26 10:00:16 Finish WAF whitelist update +``` + +## Troubleshooting + +### Ошибка: "API_sync.env file not found" + +Создайте файл конфигурации по пути `/etc/API_sync/API_sync.env` или используйте системные переменные окружения. + +### Ошибка: "git commit failed: Author identity unknown" + +Git автоматически настраивается при первом запуске. Если ошибка повторяется: + +```bash +cd /home/install/waf_whitelist +git config user.email "API_sync@example.com" +git config user.name "API Sync Bot" +``` + +### Медленная работа + +Увеличьте `MAX_CONCURRENT_WORKERS` в конфигурации: + +```env +MAX_CONCURRENT_WORKERS=10 +``` + +### Telegram не доступен напрямую + +Заполните переменные прокси в конфигурации. Утилита автоматически попробует HTTP-прокси, затем SOCKS5, и только если оба недоступны — прямое подключение. diff --git a/api.go b/api.go new file mode 100644 index 0000000..d41cfdc --- /dev/null +++ b/api.go @@ -0,0 +1,86 @@ +package main + +import ( + "encoding/json" + "fmt" + "net/http" + "time" +) + +/* HTTP CLIENT */ + +var httpClient = &http.Client{Timeout: 20 * time.Second} + +func apiGet(url string, target any) error { + req, err := http.NewRequest(http.MethodGet, url, nil) + if err != nil { + return err + } + + req.Header.Set("Authorization", "Bearer "+bearerToken) + req.Header.Set("Content-Type", "application/json") + + resp, err := httpClient.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("HTTP %s", resp.Status) + } + + return json.NewDecoder(resp.Body).Decode(target) +} + +/* API STRUCTS */ + +type info struct { + TotalCount int `json:"totalCount"` +} + +type apiList[T any] struct { + Data struct { + Result struct { + Items []T `json:"items"` + Info info `json:"info"` + } `json:"result"` + } `json:"data"` +} + +type originItem struct { + IP string `json:"ip"` + Weight int `json:"weight"` + Mode string `json:"mode"` +} + +type aliasItem struct { + ID int `json:"id"` + Domain string `json:"domain"` + UseCustomSsl int `json:"useCustomSsl"` + UseLetsencryptSsl int `json:"useLetsencryptSsl"` + SslExpireDate int64 `json:"sslExpireDate"` + CreatedAt int64 `json:"createdAt"` + ModifiedAt int64 `json:"modifiedAt"` + L7ResourceId int64 `json:"l7ResourceId"` +} + +type whoisResponse struct { + Data struct { + Result struct { + L7ResourceId int64 `json:"l7ResourceId"` + L7ResourceName string `json:"l7ResourceName"` + CreatedAt int64 `json:"createdAt"` + L7ResourceIsActive int `json:"l7ResourceIsActive"` + ModifiedAt int64 `json:"modifiedAt"` + ProtectedIp string `json:"protectedIp"` + Http2Support int `json:"http2Support"` + GrpcSupport int `json:"grpcSupport"` + WafMode int `json:"wafMode"` + WafEnabled int `json:"wafEnabled"` + WafProvider string `json:"wafProvider"` + WafInstance string `json:"wafInstance"` + ServiceIpbanEnabled int `json:"serviceIpbanEnabled"` + } `json:"result"` + } `json:"data"` +} diff --git a/compare.go b/compare.go new file mode 100644 index 0000000..86b933b --- /dev/null +++ b/compare.go @@ -0,0 +1,148 @@ +package main + +import ( + "fmt" + "reflect" + "sort" + "strings" +) + +/* COMPARISON FUNCTIONS */ + +func compareOrigins(old, new []originItem) string { + oldMap := make(map[string]originItem) + newMap := make(map[string]originItem) + + for _, o := range old { + oldMap[o.IP] = o + } + for _, o := range new { + newMap[o.IP] = o + } + + var added []originItem + var removed []originItem + var modified []string + + for ip, newOrigin := range newMap { + if oldOrigin, exists := oldMap[ip]; !exists { + added = append(added, newOrigin) + } else { + if oldOrigin.Mode != newOrigin.Mode || oldOrigin.Weight != newOrigin.Weight { + modified = append(modified, + fmt.Sprintf(" IP: %s\n Было: mode=%s, weight=%d\n Стало: mode=%s, weight=%d", + ip, oldOrigin.Mode, oldOrigin.Weight, newOrigin.Mode, newOrigin.Weight)) + } + } + } + + for ip, oldOrigin := range oldMap { + if _, exists := newMap[ip]; !exists { + removed = append(removed, oldOrigin) + } + } + + if len(added) == 0 && len(removed) == 0 && len(modified) == 0 { + return "" + } + + var result strings.Builder + result.WriteString("📡 Origins:") + + if len(added) > 0 { + result.WriteString("\n\n✅ Добавлено:\n") + for _, origin := range added { + result.WriteString(fmt.Sprintf(" IP: %s\n mode=%s, weight=%d\n", + origin.IP, origin.Mode, origin.Weight)) + } + } + + if len(removed) > 0 { + result.WriteString("\n❌ Удалено:\n") + for _, origin := range removed { + result.WriteString(fmt.Sprintf(" IP: %s\n mode=%s, weight=%d\n", + origin.IP, origin.Mode, origin.Weight)) + } + } + + if len(modified) > 0 { + result.WriteString("\n🔄 Изменено:\n") + result.WriteString(strings.Join(modified, "\n")) + } + + // Добавляем текущее состояние после всех изменений + result.WriteString("\n\n📋 Текущее состояние:\n") + + // Сортируем для красивого вывода + var sortedOrigins []originItem + for _, origin := range newMap { + sortedOrigins = append(sortedOrigins, origin) + } + sort.Slice(sortedOrigins, func(i, j int) bool { + return sortedOrigins[i].IP < sortedOrigins[j].IP + }) + + for _, origin := range sortedOrigins { + result.WriteString(fmt.Sprintf(" IP: %s\n mode=%s, weight=%d\n", + origin.IP, origin.Mode, origin.Weight)) + } + + return result.String() +} + +func compareAliases(old, new []string) (string, bool) { + oldSorted := make([]string, len(old)) + newSorted := make([]string, len(new)) + copy(oldSorted, old) + copy(newSorted, new) + sort.Strings(oldSorted) + sort.Strings(newSorted) + + if reflect.DeepEqual(oldSorted, newSorted) { + return "", false + } + + oldMap := make(map[string]bool) + newMap := make(map[string]bool) + + for _, a := range old { + oldMap[a] = true + } + for _, a := range new { + newMap[a] = true + } + + var added []string + var removed []string + + for _, a := range new { + if !oldMap[a] { + added = append(added, a) + } + } + + for _, a := range old { + if !newMap[a] { + removed = append(removed, a) + } + } + + var result strings.Builder + result.WriteString("🌐 Aliases:") + + if len(added) > 0 { + result.WriteString("\n\n✅ Добавлено:\n\n") + for _, alias := range added { + result.WriteString(alias + "\n") + } + } + + if len(removed) > 0 { + result.WriteString("\n❌ Удалено:\n\n") + for _, alias := range removed { + result.WriteString(alias + "\n") + } + } + + return result.String(), true +} diff --git a/config.go b/config.go new file mode 100644 index 0000000..e4f2ece --- /dev/null +++ b/config.go @@ -0,0 +1,125 @@ +package main + +import ( + "log" + "os" + "strconv" + + "github.com/joho/godotenv" +) + +/* CONFIG */ + +var ( + // API + bearerToken string + originURL = "https://api.servicepipe.ru/api/v1/l7/origin/global?limit=1000&l7ResourceId=" + aliasURL = "https://api.servicepipe.ru/api/v1/l7/alias/global?limit=1000&l7ResourceId=" + whoisURL = "https://api.servicepipe.ru/api/v1/l7/resource/" + + // DB + dbHost string + dbPort int + dbUser string + dbPassword string + dbName string + + // Telegram + telegramBotToken string + telegramChatID string + + // Telegram proxy + telegramHTTPProxy string + telegramSocks5Proxy string + telegramSocks5User string + telegramSocks5Password string + + // Mattermost + mattermostURL string + mattermostBotToken string + mattermostChannelID string + + // Email + smtpHost string + smtpPort int + smtpUser string + smtpPassword string + smtpFrom string + smtpTo string + + // Git + gitRepoURL string + gitRepoPath string + whitelistFile string + + // Concurrency + maxConcurrentWorkers int +) + +// Загружаем переменные из sp_sync.env при старте +func init() { + log.Println("[Config] Attempting to load sp_sync.env...") + if err := godotenv.Load("/etc/sp_sync/sp_sync.env"); err != nil { + log.Printf("[Config] ERROR loading sp_sync.env: %v", err) + log.Println("[Config] Using environment variables instead") + } else { + log.Println("[Config] sp_sync.env loaded successfully") + } + + // Инициализируем переменные ПОСЛЕ загрузки .env файла + bearerToken = getEnv("BEARER_TOKEN", "") + dbHost = getEnv("DB_HOST", "localhost") + dbPort = getEnvAsInt("DB_PORT", 5432) + dbUser = getEnv("DB_USER", "") + dbPassword = getEnv("DB_PASSWORD", "") + dbName = getEnv("DB_NAME", "waf_info") + telegramBotToken = getEnv("TELEGRAM_BOT_TOKEN", "") + telegramChatID = getEnv("TELEGRAM_CHAT_ID", "") + telegramHTTPProxy = getEnv("TELEGRAM_HTTP_PROXY", "") + telegramSocks5Proxy = getEnv("TELEGRAM_SOCKS5_PROXY", "") + telegramSocks5User = getEnv("TELEGRAM_SOCKS5_USER", "") + telegramSocks5Password = getEnv("TELEGRAM_SOCKS5_PASSWORD", "") + mattermostURL = getEnv("MATTERMOST_URL", "") + mattermostBotToken = getEnv("MATTERMOST_BOT_TOKEN", "") + mattermostChannelID = getEnv("MATTERMOST_CHANNEL_ID", "") + smtpHost = getEnv("EMAIL_SMTP_HOST", "") + smtpPort = getEnvAsInt("EMAIL_SMTP_PORT", 587) + smtpUser = getEnv("EMAIL_SMTP_USER", "") + smtpPassword = getEnv("EMAIL_SMTP_PASSWORD", "") + smtpFrom = getEnv("EMAIL_FROM", "") + smtpTo = getEnv("EMAIL_TO", "") + gitRepoURL = getEnv("GIT_REPO_URL", "https://token@svc-git.cirex.ru/wmx/waf_whitelist.git") + gitRepoPath = getEnv("GIT_REPO_PATH", "/home/install/waf_whitelist") + whitelistFile = getEnv("WHITELIST_FILE", "whitelist_ptaf.txt") + maxConcurrentWorkers = getEnvAsInt("MAX_CONCURRENT_WORKERS", 5) + + // Отладка загрузки переменных + log.Printf("[Config] DB_USER loaded: '%s'", dbUser) + log.Printf("[Config] DB_PASSWORD loaded: [%d chars]", len(dbPassword)) + log.Printf("[Config] DB_HOST: %s", dbHost) + log.Printf("[Config] DB_PORT: %d", dbPort) + log.Printf("[Config] DB_NAME: %s", dbName) +} + +// Вспомогательная функция для получения переменной окружения со значением по умолчанию +func getEnv(key, defaultValue string) string { + if value := os.Getenv(key); value != "" { + return value + } + return defaultValue +} + +// Вспомогательная функция для получения целочисленной переменной окружения +func getEnvAsInt(key string, defaultValue int) int { + valueStr := os.Getenv(key) + if valueStr == "" { + return defaultValue + } + + value, err := strconv.Atoi(valueStr) + if err != nil { + log.Printf("Warning: invalid integer value for %s: %s, using default %d", key, valueStr, defaultValue) + return defaultValue + } + return value +} diff --git a/database.go b/database.go new file mode 100644 index 0000000..c015286 --- /dev/null +++ b/database.go @@ -0,0 +1,440 @@ +package main + +import ( + "database/sql" + "encoding/json" + "fmt" + "log" + "strings" + "time" +) + +/* DATABASE */ + +func dbConnect() (*sql.DB, error) { + dsn := fmt.Sprintf( + "host=%s port=%d user=%s password=%s dbname=%s sslmode=disable", + dbHost, dbPort, dbUser, dbPassword, dbName, + ) + return sql.Open("postgres", dsn) +} + +func loadClients(db *sql.DB) ([]string, error) { + rows, err := db.Query(`SELECT client_title FROM client_info`) + if err != nil { + return nil, err + } + defer rows.Close() + + var result []string + for rows.Next() { + var c string + if err := rows.Scan(&c); err != nil { + return nil, err + } + result = append(result, c) + } + return result, nil +} + +func loadL7IDs(db *sql.DB, client string) ([]int64, error) { + rows, err := db.Query( + `SELECT l7resourceid FROM apps_settings WHERE client_title = $1 AND mode = 'auto'`, + client, + ) + if err != nil { + return nil, err + } + defer rows.Close() + + var ids []int64 + for rows.Next() { + var id int64 + if err := rows.Scan(&id); err != nil { + return nil, err + } + ids = append(ids, id) + } + return ids, nil +} + +func processSID(db *sql.DB, sid int64) error { + // WHOIS + var whois whoisResponse + if err := apiGet(whoisURL+fmt.Sprint(sid)+"/global", &whois); err != nil { + return err + } + + // ORIGINS + var originsResp apiList[originItem] + if err := apiGet(originURL+fmt.Sprint(sid), &originsResp); err != nil { + return err + } + + // ALIASES + var aliasResp apiList[aliasItem] + if err := apiGet(aliasURL+fmt.Sprint(sid), &aliasResp); err != nil { + return err + } + + originsJSON, _ := json.Marshal(originsResp.Data.Result.Items) + + var aliases []string + var latestModifiedAt int64 = 0 + + for _, a := range aliasResp.Data.Result.Items { + aliases = append(aliases, a.Domain) + if a.ModifiedAt > latestModifiedAt { + latestModifiedAt = a.ModifiedAt + } + } + aliasesJSON, _ := json.Marshal(aliases) + + // Проверяем AntiDDOS + antiddosEnabled := false + protectedIP := whois.Data.Result.ProtectedIp + var resolvedIP string + + if protectedIP != "" && whois.Data.Result.L7ResourceName != "" { + // Резолвим домен + resolvedIPAddr, err := resolveDomain(whois.Data.Result.L7ResourceName) + if err != nil { + log.Printf("[AntiDDOS Check] Failed to resolve %s: %v", whois.Data.Result.L7ResourceName, err) + } else { + resolvedIP = resolvedIPAddr + log.Printf("[AntiDDOS Check] Domain: %s, Resolved: %s, Protected: %s", + whois.Data.Result.L7ResourceName, resolvedIP, protectedIP) + + if resolvedIP == protectedIP { + antiddosEnabled = true + log.Printf("[AntiDDOS Check] ✅ AntiDDOS is ENABLED for %s", whois.Data.Result.L7ResourceName) + } else { + log.Printf("[AntiDDOS Check] ❌ AntiDDOS is DISABLED for %s", whois.Data.Result.L7ResourceName) + } + } + } + + // Извлекаем WAF данные из API + wafEnabledSP := whois.Data.Result.WafEnabled + wafVendor := whois.Data.Result.WafProvider + instanceSP := whois.Data.Result.WafInstance + + log.Printf("[WAF Info] SID: %d, WAF Enabled: %d, Vendor: %s, Instance: %s", + sid, wafEnabledSP, wafVendor, instanceSP) + + return upsertSPInfo( + db, + sid, + whois.Data.Result.L7ResourceName, + originsResp.Data.Result.Items, + aliases, + originsJSON, + aliasesJSON, + latestModifiedAt, + protectedIP, + resolvedIP, + antiddosEnabled, + wafEnabledSP, + wafVendor, + instanceSP, + ) +} + +func upsertSPInfo( + db *sql.DB, + sid int64, + domain string, + newOrigins []originItem, + newAliases []string, + originsJSON []byte, + aliasesJSON []byte, + aliasesModifiedAt int64, + protectedIP string, + resolvedIP string, + antiddosEnabled bool, + wafEnabledSP int, + wafVendor string, + instanceSP string, +) error { + // Проверяем существующие данные + var existingDomain string + var existingOriginsJSON, existingAliasesJSON []byte + var existingWafEnabled sql.NullInt64 + var existingWafVendor, existingInstanceSP sql.NullString + + err := db.QueryRow(` + SELECT domain_name, origins, aliases, waf_enabled_sp, waf_vendor, instance_sp + FROM sp_info + WHERE sid = $1 + `, sid).Scan(&existingDomain, &existingOriginsJSON, &existingAliasesJSON, + &existingWafEnabled, &existingWafVendor, &existingInstanceSP) + + isNewRecord := err == sql.ErrNoRows + + if err != nil && !isNewRecord { + return err + } + + hasChanges := false + var changeDetails []string + needsPTAFUpdate := false + + hasWAFChanges := false + var wafChangeDetails []string + + if !isNewRecord { + var existingOrigins []originItem + var existingAliases []string + json.Unmarshal(existingOriginsJSON, &existingOrigins) + json.Unmarshal(existingAliasesJSON, &existingAliases) + + if existingDomain != domain { + hasChanges = true + needsPTAFUpdate = true + changeDetails = append(changeDetails, + fmt.Sprintf("🔄 Домен изменён:\n Было: %s\n Стало: %s", existingDomain, domain)) + } + + originChanges := compareOrigins(existingOrigins, newOrigins) + if originChanges != "" { + hasChanges = true + changeDetails = append(changeDetails, originChanges) + } + + aliasChanges, aliasesChanged := compareAliases(existingAliases, newAliases) + if aliasesChanged { + hasChanges = true + needsPTAFUpdate = true + changeDetails = append(changeDetails, aliasChanges) + } + + // Проверяем изменения WAF настроек + oldWafEnabled := 0 + if existingWafEnabled.Valid { + oldWafEnabled = int(existingWafEnabled.Int64) + } + oldWafVendor := "" + if existingWafVendor.Valid { + oldWafVendor = existingWafVendor.String + } + oldInstanceSP := "" + if existingInstanceSP.Valid { + oldInstanceSP = existingInstanceSP.String + } + + if oldWafEnabled != wafEnabledSP { + hasWAFChanges = true + wafChangeDetails = append(wafChangeDetails, + fmt.Sprintf("WAF Enabled:\n Было: %d\n Стало: %d", oldWafEnabled, wafEnabledSP)) + } + + if oldWafVendor != wafVendor { + hasWAFChanges = true + oldDisplay := oldWafVendor + if oldDisplay == "" { + oldDisplay = "(не указан)" + } + newDisplay := wafVendor + if newDisplay == "" { + newDisplay = "(не указан)" + } + wafChangeDetails = append(wafChangeDetails, + fmt.Sprintf("WAF Provider:\n Было: %s\n Стало: %s", oldDisplay, newDisplay)) + } + + if oldInstanceSP != instanceSP { + hasWAFChanges = true + oldDisplay := oldInstanceSP + if oldDisplay == "" { + oldDisplay = "(не указан)" + } + newDisplay := instanceSP + if newDisplay == "" { + newDisplay = "(не указан)" + } + wafChangeDetails = append(wafChangeDetails, + fmt.Sprintf("WAF Instance:\n Было: %s\n Стало: %s", oldDisplay, newDisplay)) + } + } + + var aliasesModifiedTime *time.Time + if aliasesModifiedAt > 0 { + t := time.Unix(aliasesModifiedAt, 0) + aliasesModifiedTime = &t + } + + _, err = db.Exec(` + INSERT INTO sp_info (sid, domain_name, origins, aliases, aliaces_in_sp_modified, protected_ip, resolved_ip, antiddos_enable, waf_enabled_sp, waf_vendor, instance_sp) + VALUES ($1, $2, $3::jsonb, $4::jsonb, $5, $6, $7, $8, $9, $10, $11) + ON CONFLICT (sid) DO UPDATE SET + domain_name = EXCLUDED.domain_name, + origins = EXCLUDED.origins, + aliases = EXCLUDED.aliases, + aliaces_in_sp_modified = EXCLUDED.aliaces_in_sp_modified, + protected_ip = EXCLUDED.protected_ip, + resolved_ip = EXCLUDED.resolved_ip, + antiddos_enable = EXCLUDED.antiddos_enable, + waf_enabled_sp = EXCLUDED.waf_enabled_sp, + waf_vendor = EXCLUDED.waf_vendor, + instance_sp = EXCLUDED.instance_sp, + updated_at = now() + `, + sid, + domain, + originsJSON, + aliasesJSON, + aliasesModifiedTime, + protectedIP, + resolvedIP, + antiddosEnabled, + wafEnabledSP, + wafVendor, + instanceSP, + ) + + if err != nil { + return err + } + + // Отправляем алерт для обычных изменений + if hasChanges { + ptafWarning := "" + if needsPTAFUpdate { + ptafWarning = "\n\n⚠️ ВНИМАНИЕ: Необходимо внести изменения в кабинете PT AF!" + } + + message := fmt.Sprintf( + "🔔 Обновление WAF Info\n\n"+ + "SID: %d\n"+ + "Домен: %s\n\n"+ + "Изменения:\n\n%s%s", + sid, + domain, + strings.Join(changeDetails, "\n\n"), + ptafWarning, + ) + + if err := sendAlert(message); err != nil { + log.Printf("Failed to send Telegram alert for SID %d: %v", sid, err) + } else { + log.Printf("Alert sent for SID %d", sid) + } + } + + // Отправляем отдельный алерт для изменений WAF настроек + if hasWAFChanges { + wafMessage := fmt.Sprintf( + "⚙️ Изменения настроек инстанса на стороне SP\n\n"+ + "SID: %d\n"+ + "Домен: %s\n\n"+ + "Изменения:\n\n%s", + sid, + domain, + strings.Join(wafChangeDetails, "\n\n"), + ) + + if err := sendAlert(wafMessage); err != nil { + log.Printf("Failed to send WAF settings alert for SID %d: %v", sid, err) + } else { + log.Printf("WAF settings alert sent for SID %d", sid) + } + } + + return nil +} + +func wafVendorWarning(vendor string) string { + switch strings.ToLower(vendor) { + case "ptaf": + return "⚠️ Проверьте, не нужно ли внести изменения в PT AF!" + case "sw": + return "⚠️ Проверьте, не нужно ли внести изменения в SW!" + case "wmx": + return "⚠️ Проверьте, не нужно ли внести изменения в WMX!" + default: + return "⚠️ Проверьте, не нужно ли внести изменения в WAF!" + } +} + +func cleanupRemovedSIDs(db *sql.DB, activeSIDs []int64) error { + if len(activeSIDs) == 0 { + log.Println("[Cleanup] No active SIDs provided, skipping cleanup to avoid accidental data loss") + return nil + } + + rows, err := db.Query(` + SELECT s.sid, s.domain_name, a.waf_vendor + FROM sp_info s + LEFT JOIN apps_settings a ON a.l7resourceid = s.sid + `) + if err != nil { + return fmt.Errorf("failed to query sp_info: %w", err) + } + defer rows.Close() + + activeSet := make(map[int64]bool, len(activeSIDs)) + for _, sid := range activeSIDs { + activeSet[sid] = true + } + + type staleRecord struct { + sid int64 + domain string + wafVendor string + } + var stale []staleRecord + + for rows.Next() { + var sid int64 + var domain string + var wafVendor sql.NullString + if err := rows.Scan(&sid, &domain, &wafVendor); err != nil { + return err + } + if !activeSet[sid] { + vendor := "" + if wafVendor.Valid { + vendor = wafVendor.String + } + stale = append(stale, staleRecord{sid, domain, vendor}) + } + } + + if len(stale) == 0 { + log.Println("[Cleanup] No stale SIDs found in sp_info") + return nil + } + + log.Printf("[Cleanup] Found %d stale SID(s) to remove", len(stale)) + + for _, rec := range stale { + log.Printf("[Cleanup] Removing SID %d (%s) from sp_info", rec.sid, rec.domain) + + _, err := db.Exec(`DELETE FROM sp_info WHERE sid = $1`, rec.sid) + if err != nil { + log.Printf("[Cleanup] Failed to delete SID %d: %v", rec.sid, err) + continue + } + + wafWarning := wafVendorWarning(rec.wafVendor) + + message := fmt.Sprintf( + "🗑 Ресурс удалён из мониторинга\n\n"+ + "SID: %d\n"+ + "Домен: %s\n\n"+ + "Ресурс отсутствует в apps_settings и был удалён из sp_info.\n"+ + "%s", + rec.sid, + rec.domain, + wafWarning, + ) + + if err := sendAlert(message); err != nil { + log.Printf("[Cleanup] Failed to send Telegram alert for SID %d: %v", rec.sid, err) + } else { + log.Printf("[Cleanup] Alert sent for removed SID %d", rec.sid) + } + } + + return nil +} diff --git a/dns.go b/dns.go new file mode 100644 index 0000000..ae85524 --- /dev/null +++ b/dns.go @@ -0,0 +1,24 @@ +package main + +import ( + "fmt" + "net" +) + +/* DNS RESOLVER */ + +func resolveDomain(domain string) (string, error) { + ips, err := net.LookupIP(domain) + if err != nil { + return "", err + } + + // Возвращаем первый IPv4 адрес + for _, ip := range ips { + if ipv4 := ip.To4(); ipv4 != nil { + return ipv4.String(), nil + } + } + + return "", fmt.Errorf("no IPv4 address found for %s", domain) +} diff --git a/git.go b/git.go new file mode 100644 index 0000000..c85896b --- /dev/null +++ b/git.go @@ -0,0 +1,78 @@ +package main + +import ( + "fmt" + "log" + "os" + "os/exec" + "path/filepath" + "strings" +) + +/* GIT OPERATIONS */ + +func syncGitRepo() error { + if _, err := os.Stat(gitRepoPath); os.IsNotExist(err) { + log.Printf("[Git] Cloning repository to %s", gitRepoPath) + + parentDir := filepath.Dir(gitRepoPath) + if err := os.MkdirAll(parentDir, 0755); err != nil { + return fmt.Errorf("failed to create parent directory: %w", err) + } + + cmd := exec.Command("git", "clone", gitRepoURL, gitRepoPath) + output, err := cmd.CombinedOutput() + if err != nil { + return fmt.Errorf("git clone failed: %w, output: %s", err, output) + } + log.Println("[Git] Repository cloned successfully") + } else { + log.Println("[Git] Repository exists, checking for updates") + + cmd := exec.Command("git", "-C", gitRepoPath, "pull") + output, err := cmd.CombinedOutput() + if err != nil { + return fmt.Errorf("git pull failed: %w, output: %s", err, output) + } + + outputStr := strings.TrimSpace(string(output)) + if strings.Contains(outputStr, "Already up to date") { + log.Println("[Git] Repository already up to date") + } else { + log.Printf("[Git] Repository updated: %s", outputStr) + } + } + + // Настраиваем git config если не настроен + cmd := exec.Command("git", "-C", gitRepoPath, "config", "user.email", "waf-sync@cirex.ru") + cmd.CombinedOutput() + + cmd = exec.Command("git", "-C", gitRepoPath, "config", "user.name", "WAF Sync Bot") + cmd.CombinedOutput() + + return nil +} + +func gitCommitAndPush(message string) error { + cmd := exec.Command("git", "-C", gitRepoPath, "add", whitelistFile) + if output, err := cmd.CombinedOutput(); err != nil { + return fmt.Errorf("git add failed: %w, output: %s", err, output) + } + + cmd = exec.Command("git", "-C", gitRepoPath, "commit", "-m", message) + output, err := cmd.CombinedOutput() + if err != nil { + if strings.Contains(string(output), "nothing to commit") { + log.Println("[Git] Nothing to commit") + return nil + } + return fmt.Errorf("git commit failed: %w, output: %s", err, output) + } + + cmd = exec.Command("git", "-C", gitRepoPath, "push") + if output, err := cmd.CombinedOutput(); err != nil { + return fmt.Errorf("git push failed: %w, output: %s", err, output) + } + + return nil +} diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..ff96fe6 --- /dev/null +++ b/go.mod @@ -0,0 +1,10 @@ +module sp-sync + +go 1.25.4 + +require ( + github.com/lib/pq v1.10.9 + golang.org/x/net v0.38.0 +) + +require github.com/joho/godotenv v1.5.1 diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..5c31695 --- /dev/null +++ b/go.sum @@ -0,0 +1,6 @@ +github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0= +github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4= +github.com/lib/pq v1.10.9 h1:YXG7RB+JIjhP29X+OtkiDnYaXQwpS4JEWq7dtCCRUEw= +github.com/lib/pq v1.10.9/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o= +golang.org/x/net v0.38.0 h1:vRMAPTMaeGqVhG5QyLJHqNDwecKTomGeqbnfZyKlBI8= +golang.org/x/net v0.38.0/go.mod h1:ivrbrMbzFq5J41QOQh0siUuly180yBYtLp+CKbEaFx8= diff --git a/main.go b/main.go new file mode 100644 index 0000000..8632c3d --- /dev/null +++ b/main.go @@ -0,0 +1,80 @@ +package main + +import ( + "log" + "sync" + + _ "github.com/lib/pq" // PostgreSQL driver +) + +func main() { + log.Println("Start sync sp_info") + + db, err := dbConnect() + if err != nil { + log.Fatal(err) + } + defer db.Close() + + clients, err := loadClients(db) + if err != nil { + log.Fatal(err) + } + + // Собираем все SID для обработки + var allSIDs []int64 + for _, client := range clients { + ids, err := loadL7IDs(db, client) + if err != nil { + log.Println("binding error:", err) + continue + } + allSIDs = append(allSIDs, ids...) + } + + log.Printf("Total SIDs to process: %d", len(allSIDs)) + log.Printf("Using %d concurrent workers", maxConcurrentWorkers) + + // Создаём канал для заданий и WaitGroup для ожидания + sidChan := make(chan int64, len(allSIDs)) + var wg sync.WaitGroup + + // Запускаем воркеры + for i := 0; i < maxConcurrentWorkers; i++ { + wg.Add(1) + go func(workerID int) { + defer wg.Done() + for sid := range sidChan { + log.Printf("[Worker %d] Processing SID %d", workerID, sid) + if err := processSID(db, sid); err != nil { + log.Printf("[Worker %d] SID %d error: %v", workerID, sid, err) + } + } + }(i) + } + + // Отправляем все SID в канал + for _, sid := range allSIDs { + sidChan <- sid + } + close(sidChan) + + // Ждём завершения всех воркеров + wg.Wait() + + log.Println("Finish sync sp_info") + + // Удаляем из sp_info ресурсы, которых больше нет в apps_settings + log.Println("Start cleanup of removed SIDs") + if err := cleanupRemovedSIDs(db, allSIDs); err != nil { + log.Printf("Cleanup error: %v", err) + } + log.Println("Finish cleanup of removed SIDs") + + // Обновление WAF whitelist + log.Println("Start WAF whitelist update") + if err := updateWAFWhitelist(db); err != nil { + log.Printf("WAF whitelist update error: %v", err) + } + log.Println("Finish WAF whitelist update") +} diff --git a/notify.go b/notify.go new file mode 100644 index 0000000..cb8e73b --- /dev/null +++ b/notify.go @@ -0,0 +1,315 @@ +package main + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "log" + "net" + "net/http" + "net/smtp" + "net/url" + "path/filepath" + "strings" + "time" + + "golang.org/x/net/proxy" +) + +/* ALERT DISPATCHER */ + +// sendAlert отправляет сообщение во все настроенные каналы +func sendAlert(message string) error { + var errs []string + + if err := sendTelegramAlert(message); err != nil { + log.Printf("[Alert] Telegram error: %v", err) + errs = append(errs, "telegram: "+err.Error()) + } + + if mattermostURL != "" && mattermostBotToken != "" && mattermostChannelID != "" { + mmMessage := htmlToMarkdown(message) + if err := sendMattermostAlert(mmMessage); err != nil { + log.Printf("[Alert] Mattermost error: %v", err) + errs = append(errs, "mattermost: "+err.Error()) + } + } + + if smtpHost != "" && smtpFrom != "" && smtpTo != "" { + plainMessage := stripHTMLTags(message) + if err := sendEmailAlert(plainMessage); err != nil { + log.Printf("[Alert] Email error: %v", err) + errs = append(errs, "email: "+err.Error()) + } + } + + if len(errs) > 0 { + return fmt.Errorf("some alerts failed: %s", strings.Join(errs, "; ")) + } + return nil +} + +/* TELEGRAM */ + +// buildTelegramClient пробует HTTP-прокси, затем SOCKS5, затем прямое подключение +func buildTelegramClient() *http.Client { + timeout := 20 * time.Second + + // 1. HTTP прокси + if telegramHTTPProxy != "" { + proxyURL, err := url.Parse(telegramHTTPProxy) + if err == nil { + client := &http.Client{ + Timeout: timeout, + Transport: &http.Transport{Proxy: http.ProxyURL(proxyURL)}, + } + if checkTelegramConnectivity(client) { + log.Println("[Telegram] Using HTTP proxy:", telegramHTTPProxy) + return client + } + log.Println("[Telegram] HTTP proxy unavailable, trying SOCKS5") + } else { + log.Printf("[Telegram] Invalid HTTP proxy URL: %v", err) + } + } + + // 2. SOCKS5 прокси + if telegramSocks5Proxy != "" { + var auth *proxy.Auth + if telegramSocks5User != "" { + auth = &proxy.Auth{ + User: telegramSocks5User, + Password: telegramSocks5Password, + } + } + dialer, err := proxy.SOCKS5("tcp", telegramSocks5Proxy, auth, proxy.Direct) + if err == nil { + transport := &http.Transport{ + DialContext: func(ctx context.Context, network, addr string) (net.Conn, error) { + return dialer.Dial(network, addr) + }, + } + client := &http.Client{Timeout: timeout, Transport: transport} + if checkTelegramConnectivity(client) { + log.Println("[Telegram] Using SOCKS5 proxy:", telegramSocks5Proxy) + return client + } + log.Println("[Telegram] SOCKS5 proxy unavailable, using direct connection") + } else { + log.Printf("[Telegram] Failed to create SOCKS5 dialer: %v", err) + } + } + + // 3. Прямое подключение + log.Println("[Telegram] Using direct connection") + return &http.Client{Timeout: timeout} +} + +// checkTelegramConnectivity проверяет доступность Telegram API через данный клиент +func checkTelegramConnectivity(client *http.Client) bool { + resp, err := client.Get("https://api.telegram.org") + if err != nil { + return false + } + resp.Body.Close() + return true +} + +func sendTelegramAlert(message string) error { + client := buildTelegramClient() + + apiURL := fmt.Sprintf("https://api.telegram.org/bot%s/sendMessage", telegramBotToken) + + payload := map[string]interface{}{ + "chat_id": telegramChatID, + "text": message, + "parse_mode": "HTML", + } + + jsonData, err := json.Marshal(payload) + if err != nil { + return err + } + + resp, err := client.Post(apiURL, "application/json", bytes.NewBuffer(jsonData)) + if err != nil { + return err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("Telegram API returned status %s", resp.Status) + } + + return nil +} + +/* MATTERMOST */ + +func sendMattermostAlert(message string) error { + apiURL := fmt.Sprintf("%s/api/v4/posts", strings.TrimRight(mattermostURL, "/")) + + payload := map[string]string{ + "channel_id": mattermostChannelID, + "message": message, + } + + jsonData, err := json.Marshal(payload) + if err != nil { + return err + } + + req, err := http.NewRequest(http.MethodPost, apiURL, bytes.NewBuffer(jsonData)) + if err != nil { + return err + } + + req.Header.Set("Authorization", "Bearer "+mattermostBotToken) + req.Header.Set("Content-Type", "application/json") + + resp, err := httpClient.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusCreated { + return fmt.Errorf("Mattermost API returned status %s", resp.Status) + } + + return nil +} + +/* EMAIL */ + +func sendEmailAlert(message string) error { + addr := fmt.Sprintf("%s:%d", smtpHost, smtpPort) + + subject := "WAF Monitor Alert" + body := fmt.Sprintf("From: %s\r\nTo: %s\r\nSubject: %s\r\nContent-Type: text/plain; charset=UTF-8\r\n\r\n%s", + smtpFrom, smtpTo, subject, message) + + var auth smtp.Auth + if smtpUser != "" && smtpPassword != "" { + auth = smtp.PlainAuth("", smtpUser, smtpPassword, smtpHost) + } + + to := strings.Split(smtpTo, ",") + for i := range to { + to[i] = strings.TrimSpace(to[i]) + } + + if err := smtp.SendMail(addr, auth, smtpFrom, to, []byte(body)); err != nil { + return fmt.Errorf("smtp.SendMail: %w", err) + } + + return nil +} + +/* HELPERS */ + +// htmlToMarkdown конвертирует базовые HTML теги Telegram в Markdown для Mattermost +func htmlToMarkdown(s string) string { + r := strings.NewReplacer( + "", "**", "", "**", + "", "*", "", "*", + "", "`", "", "`", + "
", "```\n", "
", "\n```", + ) + return r.Replace(s) +} + +// stripHTMLTags убирает HTML теги для plain text email +func stripHTMLTags(s string) string { + r := strings.NewReplacer( + "", "", "", "", + "", "", "", "", + "", "", "", "", + "
", "", "
", "", + ) + return r.Replace(s) +} + +/* TELEGRAM MESSAGE FORMATTERS */ + +func formatTelegramMessage(added, removed []string, ipToSIDs map[string][]int64, sidToDomain map[int64]string) string { + var msg strings.Builder + msg.WriteString("🔄 Обновление PTAF Whitelist DosGate\n\n") + + if len(added) > 0 { + msg.WriteString(fmt.Sprintf("✅ Добавлено origin IP (%d):\n\n", len(added))) + limit := len(added) + if limit > 20 { + limit = 20 + } + for i := 0; i < limit; i++ { + ip := added[i] + if sids, ok := ipToSIDs[ip]; ok && len(sids) > 0 { + sid := sids[0] + domain := sidToDomain[sid] + msg.WriteString(fmt.Sprintf("%s (SID: %d, %s)\n", ip, sid, domain)) + } else { + msg.WriteString(ip + "\n") + } + } + if len(added) > 20 { + msg.WriteString(fmt.Sprintf("... и ещё %d\n", len(added)-20)) + } + msg.WriteString("\n") + } + + if len(removed) > 0 { + msg.WriteString(fmt.Sprintf("❌ Удалено origin IP (%d):\n\n", len(removed))) + limit := len(removed) + if limit > 20 { + limit = 20 + } + for i := 0; i < limit; i++ { + msg.WriteString(removed[i] + "\n") + } + if len(removed) > 20 { + msg.WriteString(fmt.Sprintf("... и ещё %d\n", len(removed)-20)) + } + msg.WriteString("\n") + } + + whitelistPath := filepath.Join(gitRepoPath, whitelistFile) + currentIPs, err := readWhitelistFile(whitelistPath) + totalCount := len(currentIPs) + if err != nil { + totalCount = 0 + } + + msg.WriteString(fmt.Sprintf("📊 Итого в whitelist: %d IP адресов", totalCount)) + + return msg.String() +} + +func formatCommitMessage(added, removed []string) string { + var parts []string + + if len(added) > 0 { + parts = append(parts, fmt.Sprintf("Added %d IP(s)", len(added))) + } + if len(removed) > 0 { + parts = append(parts, fmt.Sprintf("Removed %d IP(s)", len(removed))) + } + + msg := "Auto-update WAF whitelist: " + strings.Join(parts, ", ") + + if len(added)+len(removed) <= 10 { + var details []string + for _, ip := range added { + details = append(details, "+"+ip) + } + for _, ip := range removed { + details = append(details, "-"+ip) + } + if len(details) > 0 { + msg += "\n\n" + strings.Join(details, "\n") + } + } + + return msg +} diff --git a/whitelist.go b/whitelist.go new file mode 100644 index 0000000..eda3daa --- /dev/null +++ b/whitelist.go @@ -0,0 +1,314 @@ +package main + +import ( + "database/sql" + "encoding/json" + "fmt" + "log" + "net" + "os" + "path/filepath" + "sort" + "strings" +) + +/* WAF WHITELIST MANAGEMENT */ + +func updateWAFWhitelist(db *sql.DB) error { + log.Println("[WAF Whitelist] Starting update process") + + if err := syncGitRepo(); err != nil { + return fmt.Errorf("git sync failed: %w", err) + } + + wafNetworks, err := loadWAFNetworks(db) + if err != nil { + return fmt.Errorf("failed to load WAF networks: %w", err) + } + log.Printf("[WAF Whitelist] Loaded %d WAF networks from DB", len(wafNetworks)) + + allOrigins, ipToSIDs, sidToDomain, err := loadAllOrigins(db) + if err != nil { + return fmt.Errorf("failed to load origins: %w", err) + } + log.Printf("[WAF Whitelist] Loaded %d auto origin IPs", len(allOrigins)) + + manualOrigins, manualIPToSIDs, manualSIDToDomain, err := loadManualOrigins(db) + if err != nil { + return fmt.Errorf("failed to load manual origins: %w", err) + } + log.Printf("[WAF Whitelist] Loaded %d manual origin IPs", len(manualOrigins)) + + // Объединяем auto и manual + for ip, sids := range manualIPToSIDs { + ipToSIDs[ip] = append(ipToSIDs[ip], sids...) + } + for sid, domain := range manualSIDToDomain { + sidToDomain[sid] = domain + } + uniqueIPSet := make(map[string]bool) + for _, ip := range allOrigins { + uniqueIPSet[ip] = true + } + for _, ip := range manualOrigins { + uniqueIPSet[ip] = true + } + combined := make([]string, 0, len(uniqueIPSet)) + for ip := range uniqueIPSet { + combined = append(combined, ip) + } + sort.Strings(combined) + log.Printf("[WAF Whitelist] Combined total: %d unique origin IPs", len(combined)) + + filteredOrigins := filterNonWAFOrigins(combined, wafNetworks) + log.Printf("[WAF Whitelist] Filtered to %d non-WAF origin IPs", len(filteredOrigins)) + + whitelistPath := filepath.Join(gitRepoPath, whitelistFile) + currentWhitelist, err := readWhitelistFile(whitelistPath) + if err != nil { + return fmt.Errorf("failed to read whitelist file: %w", err) + } + log.Printf("[WAF Whitelist] Current whitelist contains %d IPs", len(currentWhitelist)) + + toAdd, toRemove := compareWhitelists(currentWhitelist, filteredOrigins) + + if len(toAdd) == 0 && len(toRemove) == 0 { + log.Println("[WAF Whitelist] No changes needed") + return nil + } + + log.Println("[WAF Whitelist] ==================== CHANGES ====================") + if len(toAdd) > 0 { + log.Printf("[WAF Whitelist] IPs TO ADD (%d):", len(toAdd)) + for _, ip := range toAdd { + log.Printf("[WAF Whitelist] + %s", ip) + } + } + if len(toRemove) > 0 { + log.Printf("[WAF Whitelist] IPs TO REMOVE (%d):", len(toRemove)) + for _, ip := range toRemove { + log.Printf("[WAF Whitelist] - %s", ip) + } + } + log.Println("[WAF Whitelist] ===================================================") + + if err := writeWhitelistFile(whitelistPath, filteredOrigins); err != nil { + return fmt.Errorf("failed to write whitelist file: %w", err) + } + log.Println("[WAF Whitelist] File updated successfully") + + commitMsg := formatCommitMessage(toAdd, toRemove) + if err := gitCommitAndPush(commitMsg); err != nil { + return fmt.Errorf("git commit/push failed: %w", err) + } + log.Println("[WAF Whitelist] Changes committed and pushed") + + telegramMsg := formatTelegramMessage(toAdd, toRemove, ipToSIDs, sidToDomain) + if err := sendAlert(telegramMsg); err != nil { + log.Printf("[WAF Whitelist] Failed to send Telegram alert: %v", err) + } else { + log.Println("[WAF Whitelist] Telegram alert sent") + } + + return nil +} + +func loadWAFNetworks(db *sql.DB) ([]*net.IPNet, error) { + rows, err := db.Query(`SELECT unnest(waf_networks) FROM ips`) + if err != nil { + return nil, err + } + defer rows.Close() + + var networks []*net.IPNet + for rows.Next() { + var networkStr string + if err := rows.Scan(&networkStr); err != nil { + return nil, err + } + + _, ipNet, err := net.ParseCIDR(networkStr) + if err != nil { + log.Printf("[WAF Networks] Warning: invalid CIDR %s: %v", networkStr, err) + continue + } + networks = append(networks, ipNet) + log.Printf("[WAF Networks] Loaded network: %s", networkStr) + } + + return networks, nil +} + +func loadAllOrigins(db *sql.DB) ([]string, map[string][]int64, map[int64]string, error) { + rows, err := db.Query(`SELECT sid, domain_name, origins FROM sp_info`) + if err != nil { + return nil, nil, nil, err + } + defer rows.Close() + + uniqueIPs := make(map[string]bool) + ipToSIDs := make(map[string][]int64) + sidToDomain := make(map[int64]string) + + for rows.Next() { + var sid int64 + var domain string + var originsJSON []byte + if err := rows.Scan(&sid, &domain, &originsJSON); err != nil { + return nil, nil, nil, err + } + + sidToDomain[sid] = domain + + var origins []originItem + if err := json.Unmarshal(originsJSON, &origins); err != nil { + log.Printf("[Origins] Warning: failed to parse origins JSON for SID %d: %v", sid, err) + continue + } + + for _, origin := range origins { + uniqueIPs[origin.IP] = true + ipToSIDs[origin.IP] = append(ipToSIDs[origin.IP], sid) + } + } + + result := make([]string, 0, len(uniqueIPs)) + for ip := range uniqueIPs { + result = append(result, ip) + } + sort.Strings(result) + + return result, ipToSIDs, sidToDomain, nil +} + +func loadManualOrigins(db *sql.DB) ([]string, map[string][]int64, map[int64]string, error) { + rows, err := db.Query(`SELECT sid, domain_name, origins FROM manual_info`) + if err != nil { + return nil, nil, nil, err + } + defer rows.Close() + + uniqueIPs := make(map[string]bool) + ipToSIDs := make(map[string][]int64) + sidToDomain := make(map[int64]string) + + for rows.Next() { + var sid int64 + var domain string + var originsJSON []byte + if err := rows.Scan(&sid, &domain, &originsJSON); err != nil { + return nil, nil, nil, err + } + + sidToDomain[sid] = domain + + var origins []originItem + if err := json.Unmarshal(originsJSON, &origins); err != nil { + log.Printf("[Manual Origins] Warning: failed to parse origins JSON for SID %d: %v", sid, err) + continue + } + + for _, origin := range origins { + uniqueIPs[origin.IP] = true + ipToSIDs[origin.IP] = append(ipToSIDs[origin.IP], sid) + } + } + + result := make([]string, 0, len(uniqueIPs)) + for ip := range uniqueIPs { + result = append(result, ip) + } + sort.Strings(result) + + return result, ipToSIDs, sidToDomain, nil +} + +func filterNonWAFOrigins(origins []string, wafNetworks []*net.IPNet) []string { + var filtered []string + + for _, ipStr := range origins { + ip := net.ParseIP(ipStr) + if ip == nil { + log.Printf("[Filter] Warning: invalid IP address: %s", ipStr) + continue + } + + isWAF := false + for _, network := range wafNetworks { + if network.Contains(ip) { + isWAF = true + log.Printf("[Filter] Excluding WAF IP: %s (matches network %s)", ipStr, network.String()) + break + } + } + + if !isWAF { + filtered = append(filtered, ipStr) + } + } + + return filtered +} + +func readWhitelistFile(path string) ([]string, error) { + data, err := os.ReadFile(path) + if err != nil { + if os.IsNotExist(err) { + log.Printf("[Whitelist] File does not exist, will create new one") + return []string{}, nil + } + return nil, err + } + + lines := strings.Split(string(data), "\n") + var ips []string + + for _, line := range lines { + line = strings.TrimSpace(line) + if line == "" || strings.HasPrefix(line, "#") { + continue + } + ips = append(ips, line) + } + + return ips, nil +} + +func compareWhitelists(current, desired []string) (toAdd, toRemove []string) { + currentMap := make(map[string]bool) + desiredMap := make(map[string]bool) + + for _, ip := range current { + currentMap[ip] = true + } + for _, ip := range desired { + desiredMap[ip] = true + } + + for _, ip := range desired { + if !currentMap[ip] { + toAdd = append(toAdd, ip) + } + } + + for _, ip := range current { + if !desiredMap[ip] { + toRemove = append(toRemove, ip) + } + } + + sort.Strings(toAdd) + sort.Strings(toRemove) + + return toAdd, toRemove +} + +func writeWhitelistFile(path string, ips []string) error { + var content strings.Builder + + for _, ip := range ips { + content.WriteString(ip + "\n") + } + + return os.WriteFile(path, []byte(content.String()), 0644) +}