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) +}