From 923ad8267ca7c3a3ecdbf7faadc172dea7deb46b Mon Sep 17 00:00:00 2001 From: Magnus Root Date: Fri, 20 Mar 2026 18:21:53 +0300 Subject: [PATCH] Added change balance and agent version to db --- containers_part2.go | 7 +++- db.go | 12 +++++-- main.go | 78 ++++++++++++++++++++++++++++++--------------- nginx.go | 21 ++++++++++++ sw_processor.go | 5 +++ types.go | 13 +++++--- 6 files changed, 102 insertions(+), 34 deletions(-) diff --git a/containers_part2.go b/containers_part2.go index 55f1901..775e985 100644 --- a/containers_part2.go +++ b/containers_part2.go @@ -598,10 +598,15 @@ func generateDockerCompose(filePath string, config Config, clientTitle string, c sort.Slice(allHTTPPorts, func(i, j int) bool { return allHTTPPorts[i].DockerPort < allHTTPPorts[j].DockerPort }) sort.Slice(allHTTPSPorts, func(i, j int) bool { return allHTTPSPorts[i].DockerPort < allHTTPSPorts[j].DockerPort }) + dockerImage := config.DockerImage + if clientInfo.DockerImage.Valid && clientInfo.DockerImage.String != "" { + dockerImage = clientInfo.DockerImage.String + } + data := ComposeData{ ClientTitle: clientTitle, ContainerNum: containerNum, - DockerImage: config.DockerImage, + DockerImage: dockerImage, PTAFConfig: "", WorkerProcesses: config.WorkerProcesses, WorkerConnections: config.WorkerConnections, diff --git a/db.go b/db.go index 91c927c..ccc2474 100644 --- a/db.go +++ b/db.go @@ -97,7 +97,8 @@ func checkHostBelongsToInstance(db *sql.DB, hostname, instance string) (bool, er // getClientInfoByInstance получает всех клиентов для заданного waf_instance func getClientInfoByInstance(db *sql.DB, instance string) ([]ClientInfo, error) { query := ` - SELECT containers_count, ptaf_config, client_title, fluent_bit_port, waf_instance + SELECT containers_count, ptaf_config, client_title, fluent_bit_port, waf_instance, + docker_image, docker_image_download FROM client_info WHERE waf_instance = $1 ` @@ -121,6 +122,8 @@ func getClientInfoByInstance(db *sql.DB, instance string) ([]ClientInfo, error) &ci.ClientTitle, &ci.FluentBitPort, &wafInstance, + &ci.DockerImage, + &ci.DockerImageDownload, ) if err != nil { return nil, fmt.Errorf("ошибка сканирования строки: %w", err) @@ -154,6 +157,7 @@ func getAppsSettingsByClientTitle(db *sql.DB, clientTitle string) ([]AppsSetting custom_angie_ssl, max_fails, fail_timeout, + balancing_method, custom_sw_nginx_ssl FROM apps_settings WHERE client_title = $1 @@ -187,6 +191,7 @@ func getAppsSettingsByClientTitle(db *sql.DB, clientTitle string) ([]AppsSetting &a.CustomAngieSSL, &a.MaxFails, &a.FailTimeout, + &a.BalancingMethod, &a.CustomSWNginxSSL, ) if err != nil { @@ -201,7 +206,8 @@ func getAppsSettingsByClientTitle(db *sql.DB, clientTitle string) ([]AppsSetting // getClientInfoByClientTitle получает информацию о клиенте из таблицы client_info func getClientInfoByClientTitle(db *sql.DB, clientTitle string) (*ClientInfo, error) { query := ` - SELECT containers_count, ptaf_config, client_title, fluent_bit_port, waf_instance + SELECT containers_count, ptaf_config, client_title, fluent_bit_port, waf_instance, + docker_image, docker_image_download FROM client_info WHERE client_title = $1 ` @@ -214,6 +220,8 @@ func getClientInfoByClientTitle(db *sql.DB, clientTitle string) (*ClientInfo, er &clientInfo.ClientTitle, &clientInfo.FluentBitPort, &wafInstance, + &clientInfo.DockerImage, + &clientInfo.DockerImageDownload, ) if err != nil { if strings.Contains(err.Error(), "does not exist") { diff --git a/main.go b/main.go index b143d5e..befcc12 100644 --- a/main.go +++ b/main.go @@ -68,14 +68,24 @@ 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.Printf("⚠ Ошибка подготовки Docker образа: %v", err) - log.Printf(" → Продолжаем работу, но новые PTAF-контейнеры могут не запуститься") - } + // Проверка наличия Docker образов для каждого PTAF-клиента отдельно + ptafClients := getPTAFClients(db, instanceNames) + if len(ptafClients) == 0 { + log.Printf("ℹ PTAF-клиентов на хосте нет, проверка Docker образов пропускается") } else { - log.Printf("ℹ PTAF-клиентов на хосте нет, проверка Docker образа пропускается") + checkedImages := make(map[string]bool) // не проверяем один образ дважды + for _, clientInfo := range ptafClients { + image, url := getClientDockerImage(clientInfo, config) + if image == "" || checkedImages[image] { + continue + } + checkedImages[image] = true + log.Printf(" → Проверка Docker образа для клиента %s: %s", clientInfo.ClientTitle, image) + if err := ensureDockerImageExists(image, url); err != nil { + log.Printf("⚠ Ошибка подготовки Docker образа для %s: %v", clientInfo.ClientTitle, err) + log.Printf(" → Продолжаем работу, но контейнеры %s могут не запуститься", clientInfo.ClientTitle) + } + } } // Глобальная обработка drain-контейнеров и миграций @@ -103,9 +113,9 @@ func main() { log.Println("\n=== Программа завершена успешно ===") } -// hasPTAFClients проверяет есть ли хотя бы один клиент с waf_vendor = 'ptaf' -// среди всех клиентов на данном хосте (по списку instances) -func hasPTAFClients(db *sql.DB, instanceNames []string) bool { +// getPTAFClients возвращает список PTAF-клиентов на хосте +func getPTAFClients(db *sql.DB, instanceNames []string) []ClientInfo { + var ptafClients []ClientInfo for _, instanceName := range instanceNames { clientInfoList, err := getClientInfoByInstance(db, instanceName) if err != nil { @@ -117,45 +127,61 @@ func hasPTAFClients(db *sql.DB, instanceNames []string) bool { continue } if determineVendor(appsSettingsList) == "ptaf" { - return true + ptafClients = append(ptafClients, clientInfo) } } } - return false + return ptafClients +} + +// getClientDockerImage возвращает образ и URL для клиента. +// Если у клиента заданы свои значения — использует их, иначе fallback на config. +func getClientDockerImage(clientInfo ClientInfo, config Config) (image string, url string) { + if clientInfo.DockerImage.Valid && clientInfo.DockerImage.String != "" { + image = clientInfo.DockerImage.String + } else { + image = config.DockerImage + } + if clientInfo.DockerImageDownload.Valid && clientInfo.DockerImageDownload.String != "" { + url = clientInfo.DockerImageDownload.String + } else { + url = config.DockerImageURL + } + return } // ensureDockerImageExists проверяет наличие Docker образа локально. -// Если образ отсутствует — скачивает tar-архив по DOCKER_IMAGE_URL и загружает через docker load. -func ensureDockerImageExists(config Config) error { - if config.DockerImage == "" { - log.Printf("⚠ DOCKER_IMAGE не задан, пропускаем проверку образа") +// Если образ отсутствует — скачивает tar-архив по url и загружает через docker load. +func ensureDockerImageExists(image, url string) error { + if image == "" { + log.Printf("⚠ Docker образ не задан, пропускаем проверку") return nil } // Проверяем есть ли образ локально - cmd := exec.Command("docker", "image", "inspect", config.DockerImage) + cmd := exec.Command("docker", "image", "inspect", image) if err := cmd.Run(); err == nil { - log.Printf("✓ Docker образ уже присутствует: %s", config.DockerImage) + log.Printf("✓ Docker образ уже присутствует: %s", image) return nil } - log.Printf(" Docker образ не найден локально: %s", config.DockerImage) + log.Printf(" Docker образ не найден локально: %s", image) - if config.DockerImageURL == "" { - log.Printf("⚠ Docker образ %s отсутствует локально, DOCKER_IMAGE_URL не задан — пропускаем", config.DockerImage) + if url == "" { + log.Printf("⚠ Docker образ %s отсутствует локально, URL для скачивания не задан — пропускаем", image) return nil } - // Путь для tar-архива: /home/install/ptaf-core-nginx-agent_release-X.X.X.tar - safeName := strings.ReplaceAll(config.DockerImage, ":", "_") + // Путь для tar-архива: /home/install/{image_name}.tar + safeName := strings.ReplaceAll(image, ":", "_") 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", url) log.Printf(" → Сохранение в %s", tarPath) - if err := downloadFile(tarPath, config.DockerImageURL); err != nil { + if err := downloadFile(tarPath, url); err != nil { return fmt.Errorf("ошибка скачивания образа: %w", err) } log.Printf(" ✓ Архив образа скачан: %s", tarPath) @@ -170,7 +196,7 @@ func ensureDockerImageExists(config Config) error { if err != nil { return fmt.Errorf("ошибка docker load: %w\n%s", err, string(output)) } - log.Printf(" ✓ Образ успешно загружен: %s", config.DockerImage) + log.Printf(" ✓ Образ успешно загружен: %s", image) return nil } diff --git a/nginx.go b/nginx.go index 76a9cb3..534959e 100644 --- a/nginx.go +++ b/nginx.go @@ -153,6 +153,9 @@ func writeAngieResourceConfig(f *os.File, config Config, res ResourceData, conta for _, customPort := range httpsPortsToListen { upstreamName := fmt.Sprintf("secure%d_%d_%s", res.L7ResourceID, customPort, res.ServerName) f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName)) + if res.AppsSettings.BalancingMethod.Valid && res.AppsSettings.BalancingMethod.String != "" { + f.WriteString(fmt.Sprintf(" %s;\n", strings.TrimSpace(res.AppsSettings.BalancingMethod.String))) + } maxFails := getMaxFails(res.AppsSettings) failTimeout := getFailTimeout(res.AppsSettings) for _, ports := range containerPorts { @@ -168,6 +171,9 @@ func writeAngieResourceConfig(f *os.File, config Config, res ResourceData, conta for _, customPort := range httpPortsToListen { upstreamName := fmt.Sprintf("unsecure%d_%d_%s", res.L7ResourceID, customPort, res.ServerName) f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName)) + if res.AppsSettings.BalancingMethod.Valid && res.AppsSettings.BalancingMethod.String != "" { + f.WriteString(fmt.Sprintf(" %s;\n", strings.TrimSpace(res.AppsSettings.BalancingMethod.String))) + } maxFails := getMaxFails(res.AppsSettings) failTimeout := getFailTimeout(res.AppsSettings) for _, ports := range containerPorts { @@ -192,9 +198,16 @@ func writeAngieResourceConfig(f *os.File, config Config, res ResourceData, conta var upstreamName string maxFails := getMaxFails(res.AppsSettings) failTimeout := getFailTimeout(res.AppsSettings) + balancingMethod := "" + if res.AppsSettings.BalancingMethod.Valid && res.AppsSettings.BalancingMethod.String != "" { + balancingMethod = strings.TrimSpace(res.AppsSettings.BalancingMethod.String) + } if block.IsHTTPS { upstreamName = fmt.Sprintf("secure%d_%d_%s", res.L7ResourceID, block.Port, block.ServerName) f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName)) + if balancingMethod != "" { + f.WriteString(fmt.Sprintf(" %s;\n", balancingMethod)) + } // Определяем input-порт для поиска docker-порта: точное совпадение или фоллбэк на [0] resolvedHTTPSPort := block.Port foundHTTPS := false @@ -217,6 +230,9 @@ func writeAngieResourceConfig(f *os.File, config Config, res ResourceData, conta } else { upstreamName = fmt.Sprintf("unsecure%d_%d_%s", res.L7ResourceID, block.Port, block.ServerName) f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName)) + if balancingMethod != "" { + f.WriteString(fmt.Sprintf(" %s;\n", balancingMethod)) + } resolvedHTTPPort := block.Port foundHTTP := false for _, p := range httpPortsToListen { @@ -519,6 +535,11 @@ func generateAngieCustomServerBlock(f *os.File, config Config, res ResourceData, // writeUpstreamServers записывает серверы для upstream блока nginx (с origins) func writeUpstreamServers(f *os.File, origins []OriginItem, outputPort int, settings AppsSettings) { + // Метод балансировки (если задан) + if settings.BalancingMethod.Valid && settings.BalancingMethod.String != "" { + f.WriteString(fmt.Sprintf(" %s;\n", strings.TrimSpace(settings.BalancingMethod.String))) + } + allBackup := true for _, origin := range origins { if strings.ToLower(origin.Mode) != "backup" { diff --git a/sw_processor.go b/sw_processor.go index 8512fee..40a5bcd 100644 --- a/sw_processor.go +++ b/sw_processor.go @@ -400,6 +400,11 @@ func writeSWConfigContent(f *os.File, config Config, res ResourceData) error { // writeSWUpstreamServers записывает серверы в upstream блок func writeSWUpstreamServers(f *os.File, origins []OriginItem, port int, settings AppsSettings) { + // Метод балансировки (если задан) + if settings.BalancingMethod.Valid && settings.BalancingMethod.String != "" { + f.WriteString(fmt.Sprintf(" %s;\n", strings.TrimSpace(settings.BalancingMethod.String))) + } + allBackup := true for _, origin := range origins { if strings.ToLower(origin.Mode) != "backup" { diff --git a/types.go b/types.go index 4cb8201..d0a882a 100644 --- a/types.go +++ b/types.go @@ -36,16 +36,19 @@ type AppsSettings struct { CustomAngieSSL sql.NullString // Кастомные SSL директивы для Angie - полностью заменяют дефолтные если заполнено MaxFails sql.NullInt64 // Количество неудачных попыток к origin (дефолт 5) FailTimeout sql.NullInt64 // Время в секундах после которого origin снова доступен (дефолт 15) + BalancingMethod sql.NullString // Метод балансировки: least_conn, ip_hash, random и т.д. CustomSWNginxSSL sql.NullString // Кастомные SSL директивы для SW nginx - полностью заменяют дефолтные если заполнено } // Структура для таблицы client_info (бывшая nodes) type ClientInfo struct { - ContainersCount int - PTAFConfig sql.NullString - ClientTitle string - FluentBitPort sql.NullInt64 - WAFInstance string // Instance для связи с apps_settings + ContainersCount int + PTAFConfig sql.NullString + ClientTitle string + FluentBitPort sql.NullInt64 + WAFInstance string // Instance для связи с apps_settings + DockerImage sql.NullString // Образ Docker для контейнеров клиента + DockerImageDownload sql.NullString // URL для скачивания образа } // Структуры для API ответов