diff --git a/config.go b/config.go index 4462ed4..2d01904 100644 --- a/config.go +++ b/config.go @@ -46,6 +46,7 @@ const ( DefaultPrimaryDBHost = "10.100.10.8" DefaultSecondaryDBHost = "10.100.13.5" DefaultDockerImage = "ptaf-core-nginx-agent:release-4.1.6.409465" + DefaultDockerImageURL = "" // пустой дефолт — не скачивать если не задан в env DefaultWorkerProcesses = 4 DefaultWorkerConnections = 64000 ) @@ -64,6 +65,7 @@ type Config struct { WorkerProcesses int WorkerConnections int DockerImage string + DockerImageURL string // URL для скачивания tar-архива образа если отсутствует локально } // loadEnvFile загружает переменные окружения из файла /etc/magos/magos.env @@ -111,19 +113,20 @@ func getDefaultConfig() Config { SecondaryDBHost: getEnv("SECONDARY_DB_HOST", DefaultSecondaryDBHost), APIToken: getEnv("API_TOKEN", ""), AntibotNetworks: []string{ - "185.66.85.0/24", - "185.66.86.0/24", - "185.35.5.0/24", - "185.35.6.0/24", - "212.67.26.0/24", + "185.66.23.0/24", + "185.66.25.0/24", + "185.35.26.0/24", + "185.35.26.0/24", + "212.67.25.0/24", }, WAFNetworks: []string{ - "109.238.89.0/24", - "89.20.63.0/24", + "109.238.9.0/24", + "89.20.3.0/24", }, WorkerProcesses: getEnvInt("WORKER_PROCESSES", DefaultWorkerProcesses), WorkerConnections: getEnvInt("WORKER_CONNECTIONS", DefaultWorkerConnections), DockerImage: getEnv("DOCKER_IMAGE", DefaultDockerImage), + DockerImageURL: getEnv("DOCKER_IMAGE_URL", DefaultDockerImageURL), } } @@ -153,6 +156,7 @@ func initConfig(db *sql.DB) (Config, error) { SecondaryDBHost: getEnv("SECONDARY_DB_HOST", DefaultSecondaryDBHost), APIToken: getEnv("API_TOKEN", ""), DockerImage: getEnv("DOCKER_IMAGE", DefaultDockerImage), + DockerImageURL: getEnv("DOCKER_IMAGE_URL", DefaultDockerImageURL), WorkerProcesses: getEnvInt("WORKER_PROCESSES", DefaultWorkerProcesses), WorkerConnections: getEnvInt("WORKER_CONNECTIONS", DefaultWorkerConnections), AntibotNetworks: antibotNetworks, diff --git a/db.go b/db.go index 86f077f..91c927c 100644 --- a/db.go +++ b/db.go @@ -153,7 +153,8 @@ func getAppsSettingsByClientTitle(db *sql.DB, clientTitle string) ([]AppsSetting ssl_enabled, custom_angie_ssl, max_fails, - fail_timeout + fail_timeout, + custom_sw_nginx_ssl FROM apps_settings WHERE client_title = $1 ORDER BY l7resourceid ASC @@ -186,6 +187,7 @@ func getAppsSettingsByClientTitle(db *sql.DB, clientTitle string) ([]AppsSetting &a.CustomAngieSSL, &a.MaxFails, &a.FailTimeout, + &a.CustomSWNginxSSL, ) if err != nil { return nil, fmt.Errorf("ошибка сканирования строки: %w", err) diff --git a/main.go b/main.go index 1398d1c..963a3a4 100644 --- a/main.go +++ b/main.go @@ -3,9 +3,12 @@ package main import ( "database/sql" "fmt" + "io" "log" "net" + "net/http" "os" + "os/exec" "path/filepath" "strings" "time" @@ -13,7 +16,7 @@ import ( func main() { // ════════════════════════════════════════════════════════════ - // КРИТИЧЕСКОЕ ИСПРАВЛЕНИЕ: Lock-файл для предотвращения параллельного запуска + // Lock-файл для предотвращения параллельного запуска // ════════════════════════════════════════════════════════════ lockFile := "/var/lock/ptaf-main.lock" @@ -25,11 +28,9 @@ func main() { log.Fatal("❌ Не удалось создать lock файл: ", err) } - // Записываем PID в lock файл pid := os.Getpid() f.WriteString(fmt.Sprintf("%d\n", pid)) - // Гарантированно удаляем lock при выходе defer func() { f.Close() os.Remove(lockFile) @@ -67,10 +68,19 @@ func main() { log.Printf("Instances для хоста (%d): %s", len(instanceNames), strings.Join(instanceNames, ", ")) } + // Проверка наличия Docker образа только если на хосте есть PTAF-клиенты + if hasPTAFClients(db, instanceNames) { + if err := ensureDockerImageExists(config); err != nil { + log.Fatalf("❌ Ошибка подготовки Docker образа: %v", err) + } + } else { + log.Printf("ℹ PTAF-клиентов на хосте нет, проверка Docker образа пропускается") + } + // Глобальная обработка drain-контейнеров и миграций processAllDrainingContainersGlobally(db, hostname) - // Аллокатор портов (заменяет глобальную переменную) + // Аллокатор портов portAllocator := NewPortAllocator() totalClients := 0 @@ -92,6 +102,112 @@ func main() { log.Println("\n=== Программа завершена успешно ===") } +// hasPTAFClients проверяет есть ли хотя бы один клиент с waf_vendor = 'ptaf' +// среди всех клиентов на данном хосте (по списку instances) +func hasPTAFClients(db *sql.DB, instanceNames []string) bool { + for _, instanceName := range instanceNames { + clientInfoList, err := getClientInfoByInstance(db, instanceName) + if err != nil { + continue + } + for _, clientInfo := range clientInfoList { + appsSettingsList, err := getAppsSettingsByClientTitle(db, clientInfo.ClientTitle) + if err != nil { + continue + } + if determineVendor(appsSettingsList) == "ptaf" { + return true + } + } + } + return false +} + +// ensureDockerImageExists проверяет наличие Docker образа локально. +// Если образ отсутствует — скачивает tar-архив по DOCKER_IMAGE_URL и загружает через docker load. +func ensureDockerImageExists(config Config) error { + if config.DockerImage == "" { + log.Printf("⚠ DOCKER_IMAGE не задан, пропускаем проверку образа") + return nil + } + + // Проверяем есть ли образ локально + cmd := exec.Command("docker", "image", "inspect", config.DockerImage) + if err := cmd.Run(); err == nil { + log.Printf("✓ Docker образ уже присутствует: %s", config.DockerImage) + return nil + } + + log.Printf(" Docker образ не найден локально: %s", config.DockerImage) + + if config.DockerImageURL == "" { + log.Printf("⚠ Docker образ %s отсутствует локально, DOCKER_IMAGE_URL не задан — пропускаем", config.DockerImage) + return nil + } + + // Путь для tar-архива: /home/install/ptaf-core-nginx-agent_release-X.X.X.tar + safeName := strings.ReplaceAll(config.DockerImage, ":", "_") + safeName = strings.ReplaceAll(safeName, "/", "_") + tarPath := fmt.Sprintf("/home/install/%s.tar", safeName) + + // Скачиваем архив если ещё не скачан + if _, err := os.Stat(tarPath); os.IsNotExist(err) { + log.Printf(" → Скачивание образа с %s", config.DockerImageURL) + log.Printf(" → Сохранение в %s", tarPath) + if err := downloadFile(tarPath, config.DockerImageURL); err != nil { + return fmt.Errorf("ошибка скачивания образа: %w", err) + } + log.Printf(" ✓ Архив образа скачан: %s", tarPath) + } else { + log.Printf(" → Архив уже существует локально: %s", tarPath) + } + + // Загружаем образ в Docker + log.Printf(" → Загрузка образа в Docker: docker load -i %s", tarPath) + loadCmd := exec.Command("docker", "load", "-i", tarPath) + output, err := loadCmd.CombinedOutput() + if err != nil { + return fmt.Errorf("ошибка docker load: %w\n%s", err, string(output)) + } + log.Printf(" ✓ Образ успешно загружен: %s", config.DockerImage) + + return nil +} + +// downloadFile скачивает файл по URL и сохраняет в destPath. +// Использует общий apiClient с таймаутом 120 сек (для больших tar-архивов может быть мало — +// при необходимости увеличь таймаут в apiClient или создай отдельный клиент). +func downloadFile(destPath, url string) error { + client := &http.Client{ + Timeout: 30 * time.Minute, // образы могут быть большими + } + + resp, err := client.Get(url) + if err != nil { + return err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("HTTP %d при скачивании %s", resp.StatusCode, url) + } + + f, err := os.Create(destPath) + if err != nil { + return fmt.Errorf("ошибка создания файла %s: %w", destPath, err) + } + defer f.Close() + + written, err := io.Copy(f, resp.Body) + if err != nil { + os.Remove(destPath) // удаляем неполный файл + return fmt.Errorf("ошибка записи файла: %w", err) + } + + log.Printf(" → Скачано: %.1f МБ", float64(written)/1024/1024) + return nil +} + // processInstance обрабатывает один instance и возвращает количество обработанных клиентов func processInstance(db *sql.DB, config Config, hostname, instanceName string, portAllocator *PortAllocator) int { log.Printf("\n=== Обработка instance: %s ===", instanceName) @@ -117,8 +233,8 @@ func processInstance(db *sql.DB, config Config, hostname, instanceName string, p return len(clientInfoList) } -// processClient обрабатывает одного клиента. Возвращает true если клиент был обработан. -// Роутер по waf_vendor: направляет клиента в PTAF или SW процессор +// processClient обрабатывает одного клиента. +// Роутер по waf_vendor: направляет клиента в PTAF или SW процессор. func processClient(db *sql.DB, config Config, hostname, instanceName string, clientInfo ClientInfo, portAllocator *PortAllocator) bool { log.Printf("\n--- Обработка клиента: %s (instance: %s) ---", clientInfo.ClientTitle, instanceName) @@ -151,8 +267,8 @@ func processClient(db *sql.DB, config Config, hostname, instanceName string, cli } } -// determineVendor определяет vendor из списка apps_settings -// Все ресурсы клиента должны иметь одинаковый waf_vendor +// determineVendor определяет vendor из списка apps_settings. +// Все ресурсы клиента должны иметь одинаковый waf_vendor. func determineVendor(appsSettingsList []AppsSettings) string { for _, apps := range appsSettingsList { if apps.WAFVendor.Valid && apps.WAFVendor.String != "" { @@ -165,11 +281,6 @@ func determineVendor(appsSettingsList []AppsSettings) string { return "" } -// REMOVED: старая логика processClient перенесена в ptaf/ptaf_processor.go -// Ниже идут вспомогательные функции которые используются обоими vendor - -// processClient_REMOVED_PLACEHOLDER - вся логика обработки теперь в vendor-specific модулях -// Этот блок оставлен для ясности что код был перенесён, не удалён // handleClientMigration обрабатывает миграцию клиента с текущего хоста func handleClientMigration(clientInfo ClientInfo, hostname string) { log.Printf("⚠ Клиент %s должен быть на instance '%s', текущий хост '%s' не принадлежит этому instance", @@ -247,45 +358,81 @@ func cleanStaleDrainMarkers(clientInfo ClientInfo) { } } -// fetchResourcesData REMOVED: логика перенесена в vendor-specific модули -// - PTAF: ptaf/ptaf_processor.go → fetchPTAFResourcesData() -// - SW: sw/sw_processor.go → fetchSWResourcesData() (TODO) - -// resolvePortsForClient загружает сохранённые порты или выделяет новые +// resolvePortsForClient загружает сохранённые порты или выделяет новые. +// При увеличении ContainersCount сохраняет порты существующих контейнеров +// и доделывает только новые, не трогая уже работающие. func resolvePortsForClient(clientInfo ClientInfo, resourcesData []ResourceData, portAllocator *PortAllocator) (map[int][]PortMapping, error) { log.Printf("Чтение сохранённых портов для %s", clientInfo.ClientTitle) - savedPortMap := make(map[int][]PortMapping) - allPortsFound := true + + // Инициализируем итоговый map с нужным количеством слотов + resourcePortMap := make(map[int][]PortMapping) + for _, res := range resourcesData { + resourcePortMap[res.L7ResourceID] = make([]PortMapping, clientInfo.ContainersCount) + } + + // Собираем уже занятые порты из существующих контейнеров + existingPorts := make(map[int]bool) + lastLoadedContainer := 0 for containerNum := 1; containerNum <= clientInfo.ContainersCount; containerNum++ { containerName := fmt.Sprintf("%s-ptaf-agent%03d", clientInfo.ClientTitle, containerNum) containerDir := filepath.Join("/home/install/ptaf", clientInfo.ClientTitle, containerName) - if portMappings, found := loadContainerPorts(containerDir, resourcesData); found { - if containerNum == 1 { - for _, res := range resourcesData { - savedPortMap[res.L7ResourceID] = make([]PortMapping, clientInfo.ContainersCount) - } - } - for idx, res := range resourcesData { - if idx < len(portMappings) { - savedPortMap[res.L7ResourceID][containerNum-1] = portMappings[idx] - } - } - } else { - log.Printf(" → Файл портов для контейнера %d не найден или невалиден", containerNum) - allPortsFound = false + portMappings, found := loadContainerPorts(containerDir, resourcesData) + if !found { + log.Printf(" → Файл портов для контейнера %d не найден — будут выделены новые порты", containerNum) break } + + // Сохраняем порты этого контейнера + for idx, res := range resourcesData { + if idx < len(portMappings) { + resourcePortMap[res.L7ResourceID][containerNum-1] = portMappings[idx] + // Помечаем порты как занятые + for _, dockerPort := range portMappings[idx].HTTPPorts { + existingPorts[dockerPort] = true + portAllocator.markAsUsed(dockerPort) + } + for _, dockerPort := range portMappings[idx].HTTPSPorts { + existingPorts[dockerPort] = true + portAllocator.markAsUsed(dockerPort) + } + } + } + lastLoadedContainer = containerNum + } + + // Все контейнеры загружены из файлов — ничего дополнительно выделять не нужно + if lastLoadedContainer == clientInfo.ContainersCount { + log.Printf(" → Используются сохранённые порты из .ports.json файлов (%d контейнеров)", clientInfo.ContainersCount) + return resourcePortMap, nil + } + + // Нужно выделить порты для контейнеров с lastLoadedContainer+1 по ContainersCount + newContainersCount := clientInfo.ContainersCount - lastLoadedContainer + if lastLoadedContainer == 0 { + // Ни одного файла не найдено — выделяем всё заново + log.Printf(" → Выделение новых портов для %d ресурсов × %d контейнеров", len(resourcesData), clientInfo.ContainersCount) + return portAllocator.allocatePortsForResources(resourcesData, clientInfo.ContainersCount, nil) + } + + log.Printf(" → Сохранены порты для %d существующих контейнеров, выделяю порты для %d новых (контейнеры %d-%d)", + lastLoadedContainer, newContainersCount, lastLoadedContainer+1, clientInfo.ContainersCount) + + newPortMap, err := portAllocator.allocatePortsForResources(resourcesData, newContainersCount, existingPorts) + if err != nil { + return nil, err + } + + // Объединяем: существующие порты + новые + for _, res := range resourcesData { + newMappings := newPortMap[res.L7ResourceID] + for i, mapping := range newMappings { + resourcePortMap[res.L7ResourceID][lastLoadedContainer+i] = mapping + } } - if allPortsFound && len(savedPortMap) > 0 { - log.Printf(" → Используются сохранённые порты из .ports.json файлов (%d контейнеров)", clientInfo.ContainersCount) - return savedPortMap, nil - } - - log.Printf(" → Выделение новых портов для %d ресурсов × %d контейнеров", len(resourcesData), clientInfo.ContainersCount) - return portAllocator.allocatePortsForResources(resourcesData, clientInfo.ContainersCount, nil) + return resourcePortMap, nil } // checkPortListening проверяет что порт слушает на указанном адресе diff --git a/sw_processor.go b/sw_processor.go index 056a2ce..1dd6aa6 100644 --- a/sw_processor.go +++ b/sw_processor.go @@ -496,16 +496,30 @@ func writeSWPublicHTTPServer(f *os.File, config Config, res ResourceData, port i // writeSWPublicHTTPSServer пишет публичный HTTPS server блок func writeSWPublicHTTPSServer(f *os.File, config Config, res ResourceData, port int, domainsStr string, serverDirectives string, locationDirectives string, isFullLocationBlock bool) { + sslEnabled := true + if res.AppsSettings.SSLEnabled.Valid { + sslEnabled = res.AppsSettings.SSLEnabled.Bool + } + + if sslEnabled { + log.Printf(" └─ SW HTTPS порт %d: SSL включен", port) + } else { + log.Printf(" └─ SW HTTPS порт %d: SSL отключен (plain HTTP)", port) + } + f.WriteString("# Proxying HTTPS requests to SolidWall WAF\n") f.WriteString("server {\n") f.WriteString(" include static/server.conf;\n") - f.WriteString(fmt.Sprintf(" listen %d ssl;\n", port)) + if sslEnabled { + f.WriteString(fmt.Sprintf(" listen %d ssl;\n", port)) + } else { + f.WriteString(fmt.Sprintf(" listen %d;\n", port)) + } f.WriteString(fmt.Sprintf(" server_name %s;\n\n", domainsStr)) - f.WriteString(fmt.Sprintf(" ssl_certificate %s;\n", SWSSLCertPath)) - f.WriteString(fmt.Sprintf(" ssl_certificate_key %s;\n", SWSSLKeyPath)) - f.WriteString(" ssl_protocols TLSv1.3 TLSv1.2;\n") - f.WriteString(" ssl_prefer_server_ciphers on;\n\n") + if sslEnabled { + writeSWSSLSettings(f, res.AppsSettings.CustomSWNginxSSL) + } writeSWRealIPConfig(f, config) writeSWStandardSettings(f, res.AppsSettings) @@ -594,16 +608,22 @@ func generateSWCustomServerBlock(f *os.File, config Config, res ResourceData, bl } log.Printf(" └─ SW SERVER_BLOCK: порт %d, server_name %s%s", block.Port, block.ServerName, aliasInfo) + sslEnabled := true + if res.AppsSettings.SSLEnabled.Valid { + sslEnabled = res.AppsSettings.SSLEnabled.Bool + } + f.WriteString("# Custom SERVER_BLOCK for SW\n") f.WriteString("server {\n") f.WriteString(" include static/server.conf;\n") if block.IsHTTPS { - f.WriteString(fmt.Sprintf(" listen %d ssl;\n", block.Port)) - f.WriteString(fmt.Sprintf(" ssl_certificate %s;\n", SWSSLCertPath)) - f.WriteString(fmt.Sprintf(" ssl_certificate_key %s;\n", SWSSLKeyPath)) - f.WriteString(" ssl_protocols TLSv1.3 TLSv1.2;\n") - f.WriteString(" ssl_prefer_server_ciphers on;\n\n") + if sslEnabled { + f.WriteString(fmt.Sprintf(" listen %d ssl;\n", block.Port)) + writeSWSSLSettings(f, res.AppsSettings.CustomSWNginxSSL) + } else { + f.WriteString(fmt.Sprintf(" listen %d;\n", block.Port)) + } } else { f.WriteString(fmt.Sprintf(" listen %d;\n", block.Port)) } @@ -652,6 +672,26 @@ func generateSWCustomServerBlock(f *os.File, config Config, res ResourceData, bl f.WriteString("}\n\n") } +// writeSWSSLSettings записывает SSL настройки для SW nginx +func writeSWSSLSettings(f *os.File, customSSL sql.NullString) { + if customSSL.Valid && strings.TrimSpace(customSSL.String) != "" { + log.Printf(" └─ Использование custom_sw_nginx_ssl (кастомные SSL настройки)") + lines := strings.Split(strings.TrimSpace(customSSL.String), "\n") + for _, line := range lines { + trimmed := strings.TrimSpace(line) + if trimmed != "" { + f.WriteString(" " + trimmed + "\n") + } + } + f.WriteString("\n") + } else { + f.WriteString(fmt.Sprintf(" ssl_certificate %s;\n", SWSSLCertPath)) + f.WriteString(fmt.Sprintf(" ssl_certificate_key %s;\n", SWSSLKeyPath)) + f.WriteString(" ssl_protocols TLSv1.3 TLSv1.2;\n") + f.WriteString(" ssl_prefer_server_ciphers on;\n\n") + } +} + // writeSWRealIPConfig пишет real_ip конфигурацию (antibot networks) func writeSWRealIPConfig(f *os.File, config Config) { f.WriteString(" # antibot networks\n") diff --git a/types.go b/types.go index 326afbc..4cb8201 100644 --- a/types.go +++ b/types.go @@ -36,6 +36,7 @@ type AppsSettings struct { CustomAngieSSL sql.NullString // Кастомные SSL директивы для Angie - полностью заменяют дефолтные если заполнено MaxFails sql.NullInt64 // Количество неудачных попыток к origin (дефолт 5) FailTimeout sql.NullInt64 // Время в секундах после которого origin снова доступен (дефолт 15) + CustomSWNginxSSL sql.NullString // Кастомные SSL директивы для SW nginx - полностью заменяют дефолтные если заполнено } // Структура для таблицы client_info (бывшая nodes)