From 80571bf47f78b71d62f303cbea4cde13cc3760bc Mon Sep 17 00:00:00 2001 From: Magnus Root Date: Wed, 3 Jun 2026 16:07:28 +0300 Subject: [PATCH] Added new docs and sid blocks --- .gitignore | 0 ARCHITECTURE.md | 48 +++++++++++++++--- DATABASE_SCHEMA.md | 45 +++++++++++++++-- LICENSE | 0 README.md | 0 api.go | 0 config.go | 0 containers_part1.go | 0 containers_part2.go | 0 db.go | 8 ++- drain.go | 0 go.mod | 0 go.sum | 0 main.go | 44 ++++++++++++----- nginx.go | 0 ports.go | 0 ptaf_processor.go | 117 ++++++++++++++++++++++++++++++++++++++++++++ sw_processor.go | 0 types.go | 2 + utils.go | 40 +++++++++++++++ 20 files changed, 280 insertions(+), 24 deletions(-) mode change 100644 => 100755 .gitignore mode change 100644 => 100755 ARCHITECTURE.md mode change 100644 => 100755 DATABASE_SCHEMA.md mode change 100644 => 100755 LICENSE mode change 100644 => 100755 README.md mode change 100644 => 100755 api.go mode change 100644 => 100755 config.go mode change 100644 => 100755 containers_part1.go mode change 100644 => 100755 containers_part2.go mode change 100644 => 100755 db.go mode change 100644 => 100755 drain.go mode change 100644 => 100755 go.mod mode change 100644 => 100755 go.sum mode change 100644 => 100755 main.go mode change 100644 => 100755 nginx.go mode change 100644 => 100755 ports.go mode change 100644 => 100755 ptaf_processor.go mode change 100644 => 100755 sw_processor.go mode change 100644 => 100755 types.go mode change 100644 => 100755 utils.go diff --git a/.gitignore b/.gitignore old mode 100644 new mode 100755 diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md old mode 100644 new mode 100755 index 3197660..a525c12 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -117,6 +117,8 @@ waf-host-02 | instance-c - sp_antiddos # Использовать антиддос сети в set_real_ip (дефолт: true) - real_ip_networks # Кастомные сети для set_real_ip если sp_antiddos = false (inet[]) - worker_processes # Количество worker процессов nginx внутри контейнера (дефолт: 4) +- one_container # Все SID в один контейнер (дефолт: true). При false — блочный режим +- sid_block # Блоки SID для разбивки по контейнерам: {sid1,sid2}{sid3} (TEXT) ``` **Важно:** `waf_instance` определяет на каком хосте должен работать клиент @@ -260,24 +262,46 @@ main() → initConfig() → connectToDB() 4. fetchPTAFResourcesData — получение данных ресурсов из API/кеша -5. resolvePortsForClient — загрузка/выделение портов +5. one_container = false И sid_block задан? + ├─ Да → processPTAFClientBlocks (блочный режим, см. ниже) + └─ Нет → стандартный режим: + +6. resolvePortsForClient — загрузка/выделение портов ├─ .ports.json существует → переиспользуем порты ├─ Scale up → сохраняем старые порты, добавляем новые └─ Нет файла → выделяем все заново -6. generateAngieConfigsWithoutReload — генерация Angie конфига +7. generateAngieConfigsWithoutReload — генерация Angie конфига -7. setupContainers — генерация конфигов и управление контейнерами +8. setupContainers — генерация конфигов и управление контейнерами ├─ Запущен недавно, порты не слушают → пересоздаём ├─ Compose изменился → пересоздаём ├─ Конфиг изменился → nginx -s reload └─ Новый → запускаем (с ожиданием освобождения портов) -8. Проверка портов → reload Angie если нужно +9. Проверка портов → reload Angie если нужно -9. processDrainingContainers — drain лишних контейнеров +10. processDrainingContainers — drain лишних контейнеров ``` +### 4. Блочный режим (one_container = false) + +Используется когда разные группы SID должны обслуживаться разными наборами контейнеров. + +``` +sid_block = '{10111,10231}{10112}{10113}' + +Блок 1 → SID 10111, 10231 → контейнеры {clientTitle}_block1_AZ3_hn01_a001..a{N} +Блок 2 → SID 10112 → контейнеры {clientTitle}_block2_AZ3_hn01_a001..a{N} +Блок 3 → SID 10113 → контейнеры {clientTitle}_block3_AZ3_hn01_a001..a{N} +``` + +- `containers_count` применяется к каждому блоку независимо +- Каждый блок получает свои Angie конфиги и nginx конфиги только со своими SID +- SID не указанные ни в одном блоке не обрабатываются +- Имя контейнера: `{clientTitle}_block{N}_{AZx}_{hn/n}{NN}_a{001}` +- Каждый блок обрабатывается через стандартный `resolvePortsForClient` → `generateAngieConfigsWithoutReload` → `setupContainers` + ### 4. Обработка SW-клиента ``` @@ -786,7 +810,17 @@ UPDATE client_info SET shm_size = 8 WHERE client_title = 'CLIENT001'; UPDATE apps_settings SET balancing_method = 'least_conn' WHERE l7resourceid = 12345; ``` -### Настроить ptaf_fallback +### Включить блочный режим (разбивка SID по контейнерам) + +```sql +-- Три блока SID, каждый в своём наборе контейнеров +UPDATE client_info +SET one_container = false, + sid_block = '{10111,10231}{10112}{10113}' +WHERE client_title = 'CLIENT001'; +``` + +Формат `sid_block`: `{sid1,sid2}{sid3}{sid4,sid5,sid6}` — каждая пара `{}` это один блок. ```sql -- Пропускать трафик если WAF недоступен @@ -809,4 +843,4 @@ UPDATE client_info SET sp_antiddos = false, real_ip_networks = '{10.0.0.0/8, 192.168.1.0/24}' WHERE client_title = 'CLIENT001'; -``` +``` \ No newline at end of file diff --git a/DATABASE_SCHEMA.md b/DATABASE_SCHEMA.md old mode 100644 new mode 100755 index d7afa3b..3030bce --- a/DATABASE_SCHEMA.md +++ b/DATABASE_SCHEMA.md @@ -102,6 +102,8 @@ WHERE hostname = $1 AND instance = $2; | `sp_antiddos` | BOOLEAN | Использовать антиддос сети в set_real_ip | `true` | | `real_ip_networks` | INET[] | Кастомные сети для set_real_ip (если sp_antiddos=false) | `{10.0.0.0/8}` | | `worker_processes` | INTEGER | Количество worker процессов nginx | `4` | +| `one_container` | BOOLEAN | Все SID в один контейнер | `true` | +| `sid_block` | TEXT | Блоки SID для разбивки: `{sid1,sid2}{sid3}` | `NULL` | ### Используемые запросы @@ -121,7 +123,9 @@ SELECT ptaf_fallback_code, sp_antiddos, real_ip_networks, - worker_processes + worker_processes, + one_container, + sid_block FROM client_info WHERE waf_instance = $1; ``` @@ -144,7 +148,9 @@ SELECT ptaf_fallback_code, sp_antiddos, real_ip_networks, - worker_processes + worker_processes, + one_container, + sid_block FROM client_info WHERE client_title = $1; ``` @@ -177,6 +183,12 @@ ALTER TABLE client_info ALTER TABLE client_info ADD COLUMN IF NOT EXISTS worker_processes INTEGER DEFAULT 4; + +ALTER TABLE client_info + ADD COLUMN IF NOT EXISTS one_container BOOLEAN DEFAULT TRUE; + +ALTER TABLE client_info + ADD COLUMN IF NOT EXISTS sid_block TEXT DEFAULT NULL; ``` ### Ключевые поля @@ -220,7 +232,29 @@ UPDATE client_info SET ptaf_fallback_code = '503' WHERE client_title = 'CLIENT00 UPDATE client_info SET worker_processes = 8 WHERE client_title = 'CLIENT001'; ``` -#### sp_antiddos / real_ip_networks +#### one_container / sid_block + +`one_container = true` (дефолт) — все SID клиента обслуживаются одним набором контейнеров. Стандартный режим. + +`one_container = false` + заполненный `sid_block` — блочный режим. Каждый блок `{}` в `sid_block` создаёт отдельный набор контейнеров только со своими SID. + +**Формат `sid_block`:** `{sid1,sid2}{sid3}{sid4,sid5}` — каждая пара фигурных скобок это один блок. + +```sql +-- Три блока: первый с двумя SID, второй и третий с одним +UPDATE client_info +SET one_container = false, + sid_block = '{10111,10231}{10112}{10113}' +WHERE client_title = 'CLIENT001'; +``` + +**Именование контейнеров в блочном режиме:** +- Блок 1 → `{clientTitle}_block1_AZ3_hn01_a001` +- Блок 2 → `{clientTitle}_block2_AZ3_hn01_a001` + +**containers_count** применяется к каждому блоку — если `containers_count=2` и 3 блока → 6 контейнеров суммарно. + +**SID не указанные ни в одном блоке не обрабатываются.** `sp_antiddos = true` (дефолт) — в Angie используются антиддос сети из таблицы `ips` для `set_real_ip_from`. @@ -597,6 +631,9 @@ SELECT client_title FROM client_info WHERE debug = true; -- Клиенты с индивидуальным образом SELECT client_title, docker_image FROM client_info WHERE docker_image IS NOT NULL; +-- Клиенты в блочном режиме +SELECT client_title, sid_block FROM client_info WHERE one_container = false; + -- Клиенты с нестандартным количеством worker процессов SELECT client_title, worker_processes FROM client_info WHERE worker_processes != 4; @@ -628,4 +665,4 @@ CREATE INDEX idx_client_info_waf_instance ON client_info(waf_instance); CREATE INDEX idx_client_info_client_title ON client_info(client_title); CREATE INDEX idx_apps_settings_client_title ON apps_settings(client_title); CREATE INDEX idx_apps_settings_l7resourceid ON apps_settings(l7resourceid); -``` +``` \ No newline at end of file diff --git a/LICENSE b/LICENSE old mode 100644 new mode 100755 diff --git a/README.md b/README.md old mode 100644 new mode 100755 diff --git a/api.go b/api.go old mode 100644 new mode 100755 diff --git a/config.go b/config.go old mode 100644 new mode 100755 diff --git a/containers_part1.go b/containers_part1.go old mode 100644 new mode 100755 diff --git a/containers_part2.go b/containers_part2.go old mode 100644 new mode 100755 diff --git a/db.go b/db.go old mode 100644 new mode 100755 index 259f651..f8bed52 --- a/db.go +++ b/db.go @@ -99,7 +99,7 @@ func getClientInfoByInstance(db *sql.DB, instance string) ([]ClientInfo, error) query := ` SELECT containers_count, ptaf_config, client_title, fluent_bit_port, waf_instance, docker_image, docker_image_download, debug, shm_size, ptaf_fallback_code, - sp_antiddos, real_ip_networks, worker_processes + sp_antiddos, real_ip_networks, worker_processes, one_container, sid_block FROM client_info WHERE waf_instance = $1 ` @@ -131,6 +131,8 @@ func getClientInfoByInstance(db *sql.DB, instance string) ([]ClientInfo, error) &ci.SpAntiddos, pq.Array(&ci.RealIPNetworks), &ci.WorkerProcesses, + &ci.OneContainer, + &ci.SidBlock, ) if err != nil { return nil, fmt.Errorf("ошибка сканирования строки: %w", err) @@ -217,7 +219,7 @@ func getClientInfoByClientTitle(db *sql.DB, clientTitle string) (*ClientInfo, er query := ` SELECT containers_count, ptaf_config, client_title, fluent_bit_port, waf_instance, docker_image, docker_image_download, debug, shm_size, ptaf_fallback_code, - sp_antiddos, real_ip_networks, worker_processes + sp_antiddos, real_ip_networks, worker_processes, one_container, sid_block FROM client_info WHERE client_title = $1 ` @@ -238,6 +240,8 @@ func getClientInfoByClientTitle(db *sql.DB, clientTitle string) (*ClientInfo, er &clientInfo.SpAntiddos, pq.Array(&clientInfo.RealIPNetworks), &clientInfo.WorkerProcesses, + &clientInfo.OneContainer, + &clientInfo.SidBlock, ) if err != nil { if strings.Contains(err.Error(), "does not exist") { diff --git a/drain.go b/drain.go old mode 100644 new mode 100755 diff --git a/go.mod b/go.mod old mode 100644 new mode 100755 diff --git a/go.sum b/go.sum old mode 100644 new mode 100755 diff --git a/main.go b/main.go old mode 100644 new mode 100755 index fabf315..f4007b0 --- a/main.go +++ b/main.go @@ -1,6 +1,7 @@ package main import ( + "context" "database/sql" "fmt" "io" @@ -26,6 +27,22 @@ func updateLockStatus(status string) { lockStatusFile.WriteString(fmt.Sprintf("%d\n%s\n", os.Getpid(), status)) } +// checkDockerAvailable проверяет доступность Docker daemon с таймаутом +func checkDockerAvailable() error { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + cmd := exec.CommandContext(ctx, "docker", "info", "--format", "{{.ServerVersion}}") + output, err := cmd.Output() + if ctx.Err() == context.DeadlineExceeded { + return fmt.Errorf("таймаут подключения к Docker daemon (10 сек) — возможно Docker не запущен") + } + if err != nil { + return fmt.Errorf("ошибка подключения к Docker daemon: %w", err) + } + log.Printf("✓ Docker daemon доступен, версия: %s", strings.TrimSpace(string(output))) + return nil +} + func main() { // ════════════════════════════════════════════════════════════ // Lock-файл для предотвращения параллельного запуска @@ -110,6 +127,11 @@ func main() { log.Fatalf("Ошибка загрузки конфигурации: %v", err) } + // Проверяем доступность Docker daemon + if err := checkDockerAvailable(); err != nil { + log.Fatalf("❌ Docker daemon недоступен: %v", err) + } + instanceNames, err := getInstancesByHostname(db, hostname) if err != nil { log.Fatalf("Ошибка получения instances для hostname %s: %v", hostname, err) @@ -336,14 +358,14 @@ func processClient(db *sql.DB, config Config, hostname, instanceName string, cli log.Printf("WAF Vendor: %s", vendor) switch vendor { - case "ptaf": - return processPTAFClient(db, config, hostname, instanceName, clientInfo, appsSettingsList, portAllocator) - case "sw": - return processSWClient(db, config, hostname, instanceName, clientInfo, appsSettingsList, portAllocator) - default: - log.Printf("⚠ Неизвестный или не указан waf_vendor для клиента %s, пропускаем", clientInfo.ClientTitle) - log.Printf(" Ожидается: 'ptaf' или 'sw' в поле waf_vendor таблицы apps_settings") - return false + case "ptaf": + return processPTAFClient(db, config, hostname, instanceName, clientInfo, appsSettingsList, portAllocator) + case "sw": + return processSWClient(db, config, hostname, instanceName, clientInfo, appsSettingsList, portAllocator) + default: + log.Printf("⚠ Неизвестный или не указан waf_vendor для клиента %s, пропускаем", clientInfo.ClientTitle) + log.Printf(" Ожидается: 'ptaf' или 'sw' в поле waf_vendor таблицы apps_settings") + return false } } @@ -364,7 +386,7 @@ func determineVendor(appsSettingsList []AppsSettings) string { // handleClientMigration обрабатывает миграцию клиента с текущего хоста func handleClientMigration(clientInfo ClientInfo, hostname string) { log.Printf("⚠ Клиент %s должен быть на instance '%s', текущий хост '%s' не принадлежит этому instance", - clientInfo.ClientTitle, clientInfo.WAFInstance, hostname) + clientInfo.ClientTitle, clientInfo.WAFInstance, hostname) log.Printf(" Инициирую удаление контейнеров с текущего хоста (миграция)") actualCount := getActualContainerCount(clientInfo.ClientTitle) @@ -373,7 +395,7 @@ func handleClientMigration(clientInfo ClientInfo, hostname string) { exists, _ := checkContainerStatus(containerFullName) if exists { log.Printf(" → Помечаю контейнер %s для drain (миграция, таймаут %d сек)", - containerFullName, MigrationDrainTimeout) + containerFullName, MigrationDrainTimeout) if err := markContainerForDrain(containerFullName, true); err != nil { log.Printf(" ⚠ Ошибка пометки: %v", err) } @@ -499,7 +521,7 @@ func resolvePortsForClient(clientInfo ClientInfo, resourcesData []ResourceData, } log.Printf(" → Сохранены порты для %d существующих контейнеров, выделяю порты для %d новых (контейнеры %d-%d)", - lastLoadedContainer, newContainersCount, lastLoadedContainer+1, clientInfo.ContainersCount) + lastLoadedContainer, newContainersCount, lastLoadedContainer+1, clientInfo.ContainersCount) newPortMap, err := portAllocator.allocatePortsForResources(resourcesData, newContainersCount, existingPorts) if err != nil { diff --git a/nginx.go b/nginx.go old mode 100644 new mode 100755 diff --git a/ports.go b/ports.go old mode 100644 new mode 100755 diff --git a/ptaf_processor.go b/ptaf_processor.go old mode 100644 new mode 100755 index 1759418..2149318 --- a/ptaf_processor.go +++ b/ptaf_processor.go @@ -81,6 +81,11 @@ func processPTAFClient( return false } + // Если one_container = false — разбиваем SID по блокам и обрабатываем каждый блок отдельно + if !clientInfo.OneContainer && clientInfo.SidBlock.Valid && clientInfo.SidBlock.String != "" { + return processPTAFClientBlocks(db, config, hostname, instanceName, clientInfo, resourcesData, portAllocator) + } + // Выделение портов resourcePortMap, err := resolvePortsForClient(clientInfo, resourcesData, portAllocator) if err != nil { @@ -289,3 +294,115 @@ func checkPortListeningPTAF(host string, port int) bool { conn.Close() return true } + +// processPTAFClientBlocks обрабатывает клиента с разбивкой SID по блокам контейнеров +// Используется когда one_container = false +func processPTAFClientBlocks( + db *sql.DB, + config Config, + hostname, instanceName string, + clientInfo ClientInfo, + resourcesData []ResourceData, + portAllocator *PortAllocator, +) bool { + blocks, err := parseSidBlocks(clientInfo.SidBlock.String) + if err != nil { + log.Printf("⚠ Ошибка парсинга sid_block для %s: %v", clientInfo.ClientTitle, err) + return false + } + if len(blocks) == 0 { + log.Printf("⚠ sid_block пустой для %s", clientInfo.ClientTitle) + return false + } + + log.Printf(" → Режим блоков: %d блоков SID", len(blocks)) + + // Строим map SID → ResourceData для быстрого поиска + sidToResource := make(map[int]ResourceData) + for _, res := range resourcesData { + sidToResource[res.L7ResourceID] = res + } + + success := true + for blockIdx, sidList := range blocks { + blockNum := blockIdx + 1 + blockName := fmt.Sprintf("block%d", blockNum) + log.Printf(" → Обработка блока %d: SID %v", blockNum, sidList) + + // Собираем ресурсы для этого блока + var blockResources []ResourceData + for _, sid := range sidList { + if res, ok := sidToResource[sid]; ok { + blockResources = append(blockResources, res) + } else { + log.Printf(" ⚠ SID %d не найден в данных ресурсов", sid) + } + } + if len(blockResources) == 0 { + log.Printf(" ⚠ Нет данных для блока %d, пропускаем", blockNum) + continue + } + + // Создаём clientInfo для этого блока с модифицированным именем + blockClientInfo := clientInfo + blockClientInfo.ClientTitle = fmt.Sprintf("%s_%s", clientInfo.ClientTitle, blockName) + + // Выделяем порты для блока + blockPortMap, err := resolvePortsForClient(blockClientInfo, blockResources, portAllocator) + if err != nil { + log.Printf(" ⚠ Ошибка выделения портов для блока %d: %v", blockNum, err) + success = false + continue + } + + // Генерируем конфиги Angie для блока + angieChanged, err := generateAngieConfigsWithoutReload(config, blockClientInfo.ClientTitle, blockClientInfo, blockResources, blockPortMap, blockClientInfo.Debug) + if err != nil { + log.Printf(" ⚠ Ошибка генерации конфигов Angie для блока %d: %v", blockNum, err) + success = false + continue + } + + // Настраиваем контейнеры для блока + err = setupContainers(config, blockClientInfo.ClientTitle, &blockClientInfo, blockResources, blockPortMap, hostname) + if err != nil { + log.Printf(" ⚠ Ошибка настройки контейнеров для блока %d: %v", blockNum, err) + success = false + continue + } + + // Reload Angie если нужно + if angieChanged { + allPortsReady := true + for _, res := range blockResources { + containerPorts, ok := blockPortMap[res.L7ResourceID] + if !ok { + continue + } + for _, ports := range containerPorts { + for _, dockerPort := range ports.HTTPSPorts { + if !checkPortListening("127.0.0.1", dockerPort) { + allPortsReady = false + break + } + } + if !allPortsReady { + break + } + } + if !allPortsReady { + break + } + } + if allPortsReady { + reloadAngie() + } else { + log.Printf(" ⚠ Не все порты готовы для блока %d, пропускаем reload Angie", blockNum) + } + } + + log.Printf(" ✓ Блок %d обработан успешно", blockNum) + } + + return success +} diff --git a/sw_processor.go b/sw_processor.go old mode 100644 new mode 100755 diff --git a/types.go b/types.go old mode 100644 new mode 100755 index 8c60d9b..f720514 --- a/types.go +++ b/types.go @@ -56,6 +56,8 @@ type ClientInfo struct { SpAntiddos bool // Использовать антиддос сети для set_real_ip (дефолт true) RealIPNetworks []string // Кастомные сети для set_real_ip если SpAntiddos = false WorkerProcesses int // Количество worker процессов nginx (дефолт 4) + OneContainer bool // Все SID в один контейнер (дефолт true) + SidBlock sql.NullString // Блоки SID для разбивки по контейнерам: {sid1,sid2}{sid3} } // Структуры для API ответов diff --git a/utils.go b/utils.go old mode 100644 new mode 100755 index e9fc850..cfe5357 --- a/utils.go +++ b/utils.go @@ -384,3 +384,43 @@ func buildContainerFullName(clientTitle string, containerNum int, hostname strin // prefix = "AZ3_hn01", итог: "PSB_CFA_AZ3_hn01_a001" return fmt.Sprintf("%s_%s_a%03d", clientTitle, prefix, containerNum) } + +// parseSidBlocks парсит строку вида {10111,10231}{10112} в срезы блоков SID +// Возвращает [][]int где каждый элемент — блок SID для одного набора контейнеров +func parseSidBlocks(sidBlock string) ([][]int, error) { + sidBlock = strings.TrimSpace(sidBlock) + if sidBlock == "" { + return nil, fmt.Errorf("пустая строка sid_block") + } + + var blocks [][]int + i := 0 + for i < len(sidBlock) { + if sidBlock[i] != '{' { + return nil, fmt.Errorf("ожидался '{' на позиции %d", i) + } + end := strings.Index(sidBlock[i:], "}") + if end < 0 { + return nil, fmt.Errorf("не найдена закрывающая '}' начиная с позиции %d", i) + } + inner := sidBlock[i+1 : i+end] + parts := strings.Split(inner, ",") + var block []int + for _, p := range parts { + p = strings.TrimSpace(p) + if p == "" { + continue + } + var sid int + if _, err := fmt.Sscanf(p, "%d", &sid); err != nil { + return nil, fmt.Errorf("невалидный SID '%s': %w", p, err) + } + block = append(block, sid) + } + if len(block) > 0 { + blocks = append(blocks, block) + } + i += end + 1 + } + return blocks, nil +}