Added change balance and agent version to db

This commit is contained in:
Magnus Root 2026-03-20 18:21:53 +03:00
parent 44285ad992
commit 923ad8267c
6 changed files with 102 additions and 34 deletions

View file

@ -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(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 }) 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{ data := ComposeData{
ClientTitle: clientTitle, ClientTitle: clientTitle,
ContainerNum: containerNum, ContainerNum: containerNum,
DockerImage: config.DockerImage, DockerImage: dockerImage,
PTAFConfig: "", PTAFConfig: "",
WorkerProcesses: config.WorkerProcesses, WorkerProcesses: config.WorkerProcesses,
WorkerConnections: config.WorkerConnections, WorkerConnections: config.WorkerConnections,

12
db.go
View file

@ -97,7 +97,8 @@ func checkHostBelongsToInstance(db *sql.DB, hostname, instance string) (bool, er
// getClientInfoByInstance получает всех клиентов для заданного waf_instance // getClientInfoByInstance получает всех клиентов для заданного waf_instance
func getClientInfoByInstance(db *sql.DB, instance string) ([]ClientInfo, error) { func getClientInfoByInstance(db *sql.DB, instance string) ([]ClientInfo, error) {
query := ` 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 FROM client_info
WHERE waf_instance = $1 WHERE waf_instance = $1
` `
@ -121,6 +122,8 @@ func getClientInfoByInstance(db *sql.DB, instance string) ([]ClientInfo, error)
&ci.ClientTitle, &ci.ClientTitle,
&ci.FluentBitPort, &ci.FluentBitPort,
&wafInstance, &wafInstance,
&ci.DockerImage,
&ci.DockerImageDownload,
) )
if err != nil { if err != nil {
return nil, fmt.Errorf("ошибка сканирования строки: %w", err) return nil, fmt.Errorf("ошибка сканирования строки: %w", err)
@ -154,6 +157,7 @@ func getAppsSettingsByClientTitle(db *sql.DB, clientTitle string) ([]AppsSetting
custom_angie_ssl, custom_angie_ssl,
max_fails, max_fails,
fail_timeout, fail_timeout,
balancing_method,
custom_sw_nginx_ssl custom_sw_nginx_ssl
FROM apps_settings FROM apps_settings
WHERE client_title = $1 WHERE client_title = $1
@ -187,6 +191,7 @@ func getAppsSettingsByClientTitle(db *sql.DB, clientTitle string) ([]AppsSetting
&a.CustomAngieSSL, &a.CustomAngieSSL,
&a.MaxFails, &a.MaxFails,
&a.FailTimeout, &a.FailTimeout,
&a.BalancingMethod,
&a.CustomSWNginxSSL, &a.CustomSWNginxSSL,
) )
if err != nil { if err != nil {
@ -201,7 +206,8 @@ func getAppsSettingsByClientTitle(db *sql.DB, clientTitle string) ([]AppsSetting
// getClientInfoByClientTitle получает информацию о клиенте из таблицы client_info // getClientInfoByClientTitle получает информацию о клиенте из таблицы client_info
func getClientInfoByClientTitle(db *sql.DB, clientTitle string) (*ClientInfo, error) { func getClientInfoByClientTitle(db *sql.DB, clientTitle string) (*ClientInfo, error) {
query := ` 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 FROM client_info
WHERE client_title = $1 WHERE client_title = $1
` `
@ -214,6 +220,8 @@ func getClientInfoByClientTitle(db *sql.DB, clientTitle string) (*ClientInfo, er
&clientInfo.ClientTitle, &clientInfo.ClientTitle,
&clientInfo.FluentBitPort, &clientInfo.FluentBitPort,
&wafInstance, &wafInstance,
&clientInfo.DockerImage,
&clientInfo.DockerImageDownload,
) )
if err != nil { if err != nil {
if strings.Contains(err.Error(), "does not exist") { if strings.Contains(err.Error(), "does not exist") {

78
main.go
View file

@ -68,14 +68,24 @@ func main() {
log.Printf("Instances для хоста (%d): %s", len(instanceNames), strings.Join(instanceNames, ", ")) log.Printf("Instances для хоста (%d): %s", len(instanceNames), strings.Join(instanceNames, ", "))
} }
// Проверка наличия Docker образа только если на хосте есть PTAF-клиенты // Проверка наличия Docker образов для каждого PTAF-клиента отдельно
if hasPTAFClients(db, instanceNames) { ptafClients := getPTAFClients(db, instanceNames)
if err := ensureDockerImageExists(config); err != nil { if len(ptafClients) == 0 {
log.Printf("⚠ Ошибка подготовки Docker образа: %v", err) log.Printf(" PTAF-клиентов на хосте нет, проверка Docker образов пропускается")
log.Printf(" → Продолжаем работу, но новые PTAF-контейнеры могут не запуститься")
}
} else { } 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-контейнеров и миграций // Глобальная обработка drain-контейнеров и миграций
@ -103,9 +113,9 @@ func main() {
log.Println("\n=== Программа завершена успешно ===") log.Println("\n=== Программа завершена успешно ===")
} }
// hasPTAFClients проверяет есть ли хотя бы один клиент с waf_vendor = 'ptaf' // getPTAFClients возвращает список PTAF-клиентов на хосте
// среди всех клиентов на данном хосте (по списку instances) func getPTAFClients(db *sql.DB, instanceNames []string) []ClientInfo {
func hasPTAFClients(db *sql.DB, instanceNames []string) bool { var ptafClients []ClientInfo
for _, instanceName := range instanceNames { for _, instanceName := range instanceNames {
clientInfoList, err := getClientInfoByInstance(db, instanceName) clientInfoList, err := getClientInfoByInstance(db, instanceName)
if err != nil { if err != nil {
@ -117,45 +127,61 @@ func hasPTAFClients(db *sql.DB, instanceNames []string) bool {
continue continue
} }
if determineVendor(appsSettingsList) == "ptaf" { 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 образа локально. // ensureDockerImageExists проверяет наличие Docker образа локально.
// Если образ отсутствует — скачивает tar-архив по DOCKER_IMAGE_URL и загружает через docker load. // Если образ отсутствует — скачивает tar-архив по url и загружает через docker load.
func ensureDockerImageExists(config Config) error { func ensureDockerImageExists(image, url string) error {
if config.DockerImage == "" { if image == "" {
log.Printf("⚠ DOCKER_IMAGE не задан, пропускаем проверку образа") log.Printf("⚠ Docker образ не задан, пропускаем проверку")
return nil return nil
} }
// Проверяем есть ли образ локально // Проверяем есть ли образ локально
cmd := exec.Command("docker", "image", "inspect", config.DockerImage) cmd := exec.Command("docker", "image", "inspect", image)
if err := cmd.Run(); err == nil { if err := cmd.Run(); err == nil {
log.Printf("✓ Docker образ уже присутствует: %s", config.DockerImage) log.Printf("✓ Docker образ уже присутствует: %s", image)
return nil return nil
} }
log.Printf(" Docker образ не найден локально: %s", config.DockerImage) log.Printf(" Docker образ не найден локально: %s", image)
if config.DockerImageURL == "" { if url == "" {
log.Printf("⚠ Docker образ %s отсутствует локально, DOCKER_IMAGE_URL не задан — пропускаем", config.DockerImage) log.Printf("⚠ Docker образ %s отсутствует локально, URL для скачивания не задан — пропускаем", image)
return nil return nil
} }
// Путь для tar-архива: /home/install/ptaf-core-nginx-agent_release-X.X.X.tar // Путь для tar-архива: /home/install/{image_name}.tar
safeName := strings.ReplaceAll(config.DockerImage, ":", "_") safeName := strings.ReplaceAll(image, ":", "_")
safeName = strings.ReplaceAll(safeName, "/", "_") safeName = strings.ReplaceAll(safeName, "/", "_")
tarPath := fmt.Sprintf("/home/install/%s.tar", safeName) tarPath := fmt.Sprintf("/home/install/%s.tar", safeName)
// Скачиваем архив если ещё не скачан // Скачиваем архив если ещё не скачан
if _, err := os.Stat(tarPath); os.IsNotExist(err) { if _, err := os.Stat(tarPath); os.IsNotExist(err) {
log.Printf(" → Скачивание образа с %s", config.DockerImageURL) log.Printf(" → Скачивание образа с %s", url)
log.Printf(" → Сохранение в %s", tarPath) log.Printf(" → Сохранение в %s", tarPath)
if err := downloadFile(tarPath, config.DockerImageURL); err != nil { if err := downloadFile(tarPath, url); err != nil {
return fmt.Errorf("ошибка скачивания образа: %w", err) return fmt.Errorf("ошибка скачивания образа: %w", err)
} }
log.Printf(" ✓ Архив образа скачан: %s", tarPath) log.Printf(" ✓ Архив образа скачан: %s", tarPath)
@ -170,7 +196,7 @@ func ensureDockerImageExists(config Config) error {
if err != nil { if err != nil {
return fmt.Errorf("ошибка docker load: %w\n%s", err, string(output)) return fmt.Errorf("ошибка docker load: %w\n%s", err, string(output))
} }
log.Printf(" ✓ Образ успешно загружен: %s", config.DockerImage) log.Printf(" ✓ Образ успешно загружен: %s", image)
return nil return nil
} }

View file

@ -153,6 +153,9 @@ func writeAngieResourceConfig(f *os.File, config Config, res ResourceData, conta
for _, customPort := range httpsPortsToListen { for _, customPort := range httpsPortsToListen {
upstreamName := fmt.Sprintf("secure%d_%d_%s", res.L7ResourceID, customPort, res.ServerName) upstreamName := fmt.Sprintf("secure%d_%d_%s", res.L7ResourceID, customPort, res.ServerName)
f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName)) 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) maxFails := getMaxFails(res.AppsSettings)
failTimeout := getFailTimeout(res.AppsSettings) failTimeout := getFailTimeout(res.AppsSettings)
for _, ports := range containerPorts { for _, ports := range containerPorts {
@ -168,6 +171,9 @@ func writeAngieResourceConfig(f *os.File, config Config, res ResourceData, conta
for _, customPort := range httpPortsToListen { for _, customPort := range httpPortsToListen {
upstreamName := fmt.Sprintf("unsecure%d_%d_%s", res.L7ResourceID, customPort, res.ServerName) upstreamName := fmt.Sprintf("unsecure%d_%d_%s", res.L7ResourceID, customPort, res.ServerName)
f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName)) 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) maxFails := getMaxFails(res.AppsSettings)
failTimeout := getFailTimeout(res.AppsSettings) failTimeout := getFailTimeout(res.AppsSettings)
for _, ports := range containerPorts { for _, ports := range containerPorts {
@ -192,9 +198,16 @@ func writeAngieResourceConfig(f *os.File, config Config, res ResourceData, conta
var upstreamName string var upstreamName string
maxFails := getMaxFails(res.AppsSettings) maxFails := getMaxFails(res.AppsSettings)
failTimeout := getFailTimeout(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 { if block.IsHTTPS {
upstreamName = fmt.Sprintf("secure%d_%d_%s", res.L7ResourceID, block.Port, block.ServerName) upstreamName = fmt.Sprintf("secure%d_%d_%s", res.L7ResourceID, block.Port, block.ServerName)
f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName)) f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName))
if balancingMethod != "" {
f.WriteString(fmt.Sprintf(" %s;\n", balancingMethod))
}
// Определяем input-порт для поиска docker-порта: точное совпадение или фоллбэк на [0] // Определяем input-порт для поиска docker-порта: точное совпадение или фоллбэк на [0]
resolvedHTTPSPort := block.Port resolvedHTTPSPort := block.Port
foundHTTPS := false foundHTTPS := false
@ -217,6 +230,9 @@ func writeAngieResourceConfig(f *os.File, config Config, res ResourceData, conta
} else { } else {
upstreamName = fmt.Sprintf("unsecure%d_%d_%s", res.L7ResourceID, block.Port, block.ServerName) upstreamName = fmt.Sprintf("unsecure%d_%d_%s", res.L7ResourceID, block.Port, block.ServerName)
f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName)) f.WriteString(fmt.Sprintf("upstream %s {\n", upstreamName))
if balancingMethod != "" {
f.WriteString(fmt.Sprintf(" %s;\n", balancingMethod))
}
resolvedHTTPPort := block.Port resolvedHTTPPort := block.Port
foundHTTP := false foundHTTP := false
for _, p := range httpPortsToListen { for _, p := range httpPortsToListen {
@ -519,6 +535,11 @@ func generateAngieCustomServerBlock(f *os.File, config Config, res ResourceData,
// writeUpstreamServers записывает серверы для upstream блока nginx (с origins) // writeUpstreamServers записывает серверы для upstream блока nginx (с origins)
func writeUpstreamServers(f *os.File, origins []OriginItem, outputPort int, settings AppsSettings) { 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 allBackup := true
for _, origin := range origins { for _, origin := range origins {
if strings.ToLower(origin.Mode) != "backup" { if strings.ToLower(origin.Mode) != "backup" {

View file

@ -400,6 +400,11 @@ func writeSWConfigContent(f *os.File, config Config, res ResourceData) error {
// writeSWUpstreamServers записывает серверы в upstream блок // writeSWUpstreamServers записывает серверы в upstream блок
func writeSWUpstreamServers(f *os.File, origins []OriginItem, port int, settings AppsSettings) { 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 allBackup := true
for _, origin := range origins { for _, origin := range origins {
if strings.ToLower(origin.Mode) != "backup" { if strings.ToLower(origin.Mode) != "backup" {

View file

@ -36,6 +36,7 @@ type AppsSettings struct {
CustomAngieSSL sql.NullString // Кастомные SSL директивы для Angie - полностью заменяют дефолтные если заполнено CustomAngieSSL sql.NullString // Кастомные SSL директивы для Angie - полностью заменяют дефолтные если заполнено
MaxFails sql.NullInt64 // Количество неудачных попыток к origin (дефолт 5) MaxFails sql.NullInt64 // Количество неудачных попыток к origin (дефолт 5)
FailTimeout sql.NullInt64 // Время в секундах после которого origin снова доступен (дефолт 15) FailTimeout sql.NullInt64 // Время в секундах после которого origin снова доступен (дефолт 15)
BalancingMethod sql.NullString // Метод балансировки: least_conn, ip_hash, random и т.д.
CustomSWNginxSSL sql.NullString // Кастомные SSL директивы для SW nginx - полностью заменяют дефолтные если заполнено CustomSWNginxSSL sql.NullString // Кастомные SSL директивы для SW nginx - полностью заменяют дефолтные если заполнено
} }
@ -46,6 +47,8 @@ type ClientInfo struct {
ClientTitle string ClientTitle string
FluentBitPort sql.NullInt64 FluentBitPort sql.NullInt64
WAFInstance string // Instance для связи с apps_settings WAFInstance string // Instance для связи с apps_settings
DockerImage sql.NullString // Образ Docker для контейнеров клиента
DockerImageDownload sql.NullString // URL для скачивания образа
} }
// Структуры для API ответов // Структуры для API ответов