package main import ( "database/sql" "fmt" "io" "log" "net" "net/http" "os" "os/exec" "path/filepath" "strings" "time" ) func main() { // ════════════════════════════════════════════════════════════ // Lock-файл для предотвращения параллельного запуска // ════════════════════════════════════════════════════════════ lockFile := "/var/lock/ptaf-main.lock" f, err := os.OpenFile(lockFile, os.O_CREATE|os.O_EXCL|os.O_RDWR, 0600) if err != nil { if os.IsExist(err) { log.Fatal("❌ Программа уже запущена! Lock файл существует: ", lockFile) } log.Fatal("❌ Не удалось создать lock файл: ", err) } pid := os.Getpid() f.WriteString(fmt.Sprintf("%d\n", pid)) defer func() { f.Close() os.Remove(lockFile) log.Printf("✓ Lock файл удален") }() log.Printf("✓ Lock файл создан: %s (PID: %d)", lockFile, pid) // ════════════════════════════════════════════════════════════ log.Println("=== Запуск программы управления конфигурациями PTAF ===") hostname, err := os.Hostname() if err != nil { log.Fatalf("Ошибка получения hostname: %v", err) } log.Printf("Hostname текущего хоста: %s", hostname) tempConfig := getDefaultConfig() db := connectToDB(tempConfig) defer db.Close() config, err := initConfig(db) if err != nil { log.Fatalf("Ошибка загрузки конфигурации: %v", err) } instanceNames, err := getInstancesByHostname(db, hostname) if err != nil { log.Fatalf("Ошибка получения instances для hostname %s: %v", hostname, err) } if len(instanceNames) == 1 { log.Printf("Instance для хоста: %s", instanceNames[0]) } else { 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 for _, instanceName := range instanceNames { count := processInstance(db, config, hostname, instanceName, portAllocator) totalClients += count } if totalClients == 0 { log.Println("\nНет активных клиентов для обработки на текущем хосте") } else { log.Printf("\nВсего обработано клиентов: %d (на %d instance)", totalClients, len(instanceNames)) } // Финальная проверка drain-конфигов Angie log.Println("\n--- Проверка отложенных удалений Angie конфигов ---") processDrainingAngieConfigs() 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) clientInfoList, err := getClientInfoByInstance(db, instanceName) if err != nil { log.Printf("⚠ Ошибка получения клиентов для instance %s: %v", instanceName, err) return 0 } if len(clientInfoList) == 0 { log.Printf("Нет активных клиентов для instance %s", instanceName) return 0 } log.Printf("Найдено клиентов для instance %s: %d", instanceName, len(clientInfoList)) for _, clientInfo := range clientInfoList { processClient(db, config, hostname, instanceName, clientInfo, portAllocator) } log.Printf("\n=== Instance %s обработан ===", instanceName) return len(clientInfoList) } // 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) appsSettingsList, err := getAppsSettingsByClientTitle(db, clientInfo.ClientTitle) if err != nil { log.Printf("Ошибка получения apps_settings для %s: %v", clientInfo.ClientTitle, err) return false } if len(appsSettingsList) == 0 { log.Printf("⚠ Нет apps_settings для клиента %s, пропускаем", clientInfo.ClientTitle) return false } // ════════════════════════════════════════════════════════════ // ТОЧКА ВЕТВЛЕНИЯ ПО WAF VENDOR // ════════════════════════════════════════════════════════════ vendor := determineVendor(appsSettingsList) 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 } } // determineVendor определяет vendor из списка apps_settings. // Все ресурсы клиента должны иметь одинаковый waf_vendor. func determineVendor(appsSettingsList []AppsSettings) string { for _, apps := range appsSettingsList { if apps.WAFVendor.Valid && apps.WAFVendor.String != "" { v := strings.ToLower(strings.TrimSpace(apps.WAFVendor.String)) if v == "ptaf" || v == "sw" { return v } } } return "" } // handleClientMigration обрабатывает миграцию клиента с текущего хоста func handleClientMigration(clientInfo ClientInfo, hostname string) { log.Printf("⚠ Клиент %s должен быть на instance '%s', текущий хост '%s' не принадлежит этому instance", clientInfo.ClientTitle, clientInfo.WAFInstance, hostname) log.Printf(" Инициирую удаление контейнеров с текущего хоста (миграция)") for containerNum := 1; containerNum <= MaxContainersPerClient; containerNum++ { containerFullName := fmt.Sprintf("ptaf_%s_%03d", clientInfo.ClientTitle, containerNum) exists, _ := checkContainerStatus(containerFullName) if exists { log.Printf(" → Помечаю контейнер %s для drain (миграция, таймаут %d сек)", containerFullName, MigrationDrainTimeout) if err := markContainerForDrain(containerFullName, true); err != nil { log.Printf(" ⚠ Ошибка пометки: %v", err) } } } angieConfigPath := fmt.Sprintf("/etc/angie/http.d/angie-ptaf-%s.conf", clientInfo.ClientTitle) if _, err := os.Stat(angieConfigPath); err == nil { log.Printf(" → Помечаю конфиг Angie для удаления (миграция, таймаут %d сек)", MigrationDrainTimeout) if err := markAngieConfigForDrain(clientInfo.ClientTitle); err != nil { log.Printf(" ⚠ Ошибка пометки конфига: %v", err) } else { log.Printf(" ✓ Конфиг Angie помечен для удаления") } } processDrainingContainers(clientInfo.ClientTitle, 0) processDrainingAngieConfigs() log.Printf(" Контейнеры и конфиг будут удалены через %d сек после завершения соединений", MigrationDrainTimeout) } // cleanStaleDrainMarkers очищает устаревшие drain-маркеры если клиент вернулся на хост func cleanStaleDrainMarkers(clientInfo ClientInfo) { log.Printf(" → Проверка и очистка устаревших маркеров drain для %s", clientInfo.ClientTitle) cleanedContainers := 0 cleanedAngieMarker := false for containerNum := 1; containerNum <= MaxContainersPerClient; containerNum++ { containerFullName := fmt.Sprintf("ptaf_%s_%03d", clientInfo.ClientTitle, containerNum) markerFile := fmt.Sprintf("/tmp/ptaf-drain-%s", containerFullName) if _, err := os.Stat(markerFile); err == nil { _, isMigration, found := getContainerDrainInfo(containerFullName) if found && isMigration { if err := os.Remove(markerFile); err == nil { log.Printf(" ✓ Удалён устаревший маркер drain для контейнера %s (миграция отменена)", containerFullName) cleanedContainers++ } } else if found && !isMigration { if containerNum <= clientInfo.ContainersCount { if err := os.Remove(markerFile); err == nil { log.Printf(" ✓ Удалён устаревший маркер drain для контейнера %s (масштабирование отменено)", containerFullName) cleanedContainers++ } } } } } angieMarkerFile := fmt.Sprintf("/tmp/ptaf-angie-drain-%s", clientInfo.ClientTitle) if _, err := os.Stat(angieMarkerFile); err == nil { if err := os.Remove(angieMarkerFile); err == nil { log.Printf(" ✓ Удалён устаревший маркер drain для конфига Angie (миграция отменена)") cleanedAngieMarker = true } } if cleanedContainers > 0 || cleanedAngieMarker { log.Printf(" ✓ Очищено устаревших маркеров: контейнеры=%d, конфиг Angie=%v", cleanedContainers, cleanedAngieMarker) } else { log.Printf(" ✓ Устаревших маркеров drain не найдено") } } // resolvePortsForClient загружает сохранённые порты или выделяет новые. // При увеличении ContainersCount сохраняет порты существующих контейнеров // и доделывает только новые, не трогая уже работающие. func resolvePortsForClient(clientInfo ClientInfo, resourcesData []ResourceData, portAllocator *PortAllocator) (map[int][]PortMapping, error) { log.Printf("Чтение сохранённых портов для %s", clientInfo.ClientTitle) // Инициализируем итоговый 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) 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 } } return resourcePortMap, nil } // checkPortListening проверяет что порт слушает на указанном адресе func checkPortListening(host string, port int) bool { conn, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", host, port), 2*time.Second) if err != nil { return false } conn.Close() return true }