Added database_schema

This commit is contained in:
Magnus Root 2026-04-17 11:07:49 +03:00
parent faf236ad0a
commit cca899e659
15 changed files with 2120 additions and 0 deletions

46
.gitignore vendored Normal file
View file

@ -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

149
ARCHITECTURE.md Normal file
View file

@ -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 <whitelist_file>`, `git commit -m <message>`, `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 (стандартная библиотека) |

69
DATABASE_SCHEMA.md Normal file
View file

@ -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"}` |

230
README.md Normal file
View file

@ -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, и только если оба недоступны — прямое подключение.

86
api.go Normal file
View file

@ -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"`
}

148
compare.go Normal file
View file

@ -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("<b>📡 Origins:</b>")
if len(added) > 0 {
result.WriteString("\n\n<b>✅ Добавлено:</b>\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<b>❌ Удалено:</b>\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<b>🔄 Изменено:</b>\n")
result.WriteString(strings.Join(modified, "\n"))
}
// Добавляем текущее состояние после всех изменений
result.WriteString("\n\n<b>📋 Текущее состояние:</b>\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("<b>🌐 Aliases:</b>")
if len(added) > 0 {
result.WriteString("\n\n<b>✅ Добавлено:</b>\n\n")
for _, alias := range added {
result.WriteString(alias + "\n")
}
}
if len(removed) > 0 {
result.WriteString("\n<b>❌ Удалено:</b>\n\n")
for _, alias := range removed {
result.WriteString(alias + "\n")
}
}
return result.String(), true
}

125
config.go Normal file
View file

@ -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
}

440
database.go Normal file
View file

@ -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("<b>🔄 Домен изменён:</b>\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("<b>WAF Enabled:</b>\n <b>Было:</b> %d\n <b>Стало:</b> %d", oldWafEnabled, wafEnabledSP))
}
if oldWafVendor != wafVendor {
hasWAFChanges = true
oldDisplay := oldWafVendor
if oldDisplay == "" {
oldDisplay = "(не указан)"
}
newDisplay := wafVendor
if newDisplay == "" {
newDisplay = "(не указан)"
}
wafChangeDetails = append(wafChangeDetails,
fmt.Sprintf("<b>WAF Provider:</b>\n <b>Было:</b> %s\n <b>Стало:</b> %s", oldDisplay, newDisplay))
}
if oldInstanceSP != instanceSP {
hasWAFChanges = true
oldDisplay := oldInstanceSP
if oldDisplay == "" {
oldDisplay = "(не указан)"
}
newDisplay := instanceSP
if newDisplay == "" {
newDisplay = "(не указан)"
}
wafChangeDetails = append(wafChangeDetails,
fmt.Sprintf("<b>WAF Instance:</b>\n <b>Было:</b> %s\n <b>Стало:</b> %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<b>⚠️ ВНИМАНИЕ: Необходимо внести изменения в кабинете PT AF!</b>"
}
message := fmt.Sprintf(
"<b>🔔 Обновление WAF Info</b>\n\n"+
"<b>SID:</b> %d\n"+
"<b>Домен:</b> %s\n\n"+
"<b>Изменения:</b>\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(
"<b>⚙️ Изменения настроек инстанса на стороне SP</b>\n\n"+
"<b>SID:</b> %d\n"+
"<b>Домен:</b> %s\n\n"+
"<b>Изменения:</b>\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 "<b>⚠️ Проверьте, не нужно ли внести изменения в PT AF!</b>"
case "sw":
return "<b>⚠️ Проверьте, не нужно ли внести изменения в SW!</b>"
case "wmx":
return "<b>⚠️ Проверьте, не нужно ли внести изменения в WMX!</b>"
default:
return "<b>⚠️ Проверьте, не нужно ли внести изменения в WAF!</b>"
}
}
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(
"<b>🗑 Ресурс удалён из мониторинга</b>\n\n"+
"<b>SID:</b> %d\n"+
"<b>Домен:</b> %s\n\n"+
"Ресурс отсутствует в <code>apps_settings</code> и был удалён из <code>sp_info</code>.\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
}

24
dns.go Normal file
View file

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

78
git.go Normal file
View file

@ -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
}

10
go.mod Normal file
View file

@ -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

6
go.sum Normal file
View file

@ -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=

80
main.go Normal file
View file

@ -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")
}

315
notify.go Normal file
View file

@ -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(
"<b>", "**", "</b>", "**",
"<i>", "*", "</i>", "*",
"<code>", "`", "</code>", "`",
"<pre>", "```\n", "</pre>", "\n```",
)
return r.Replace(s)
}
// stripHTMLTags убирает HTML теги для plain text email
func stripHTMLTags(s string) string {
r := strings.NewReplacer(
"<b>", "", "</b>", "",
"<i>", "", "</i>", "",
"<code>", "", "</code>", "",
"<pre>", "", "</pre>", "",
)
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("<b>🔄 Обновление PTAF Whitelist DosGate</b>\n\n")
if len(added) > 0 {
msg.WriteString(fmt.Sprintf("<b>✅ Добавлено origin IP (%d):</b>\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("<b>❌ Удалено origin IP (%d):</b>\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("<b>📊 Итого в whitelist:</b> %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
}

314
whitelist.go Normal file
View file

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