package main import ( "database/sql" "fmt" "log" "net" "os" "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 в lock файл pid := os.Getpid() f.WriteString(fmt.Sprintf("%d\n", pid)) // Гарантированно удаляем lock при выходе 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, ", ")) } // Глобальная обработка 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=== Программа завершена успешно ===") } // 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 обрабатывает одного клиента. Возвращает true если клиент был обработан. // Роутер по 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 "" } // 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", 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 не найдено") } } // fetchResourcesData REMOVED: логика перенесена в vendor-specific модули // - PTAF: ptaf/ptaf_processor.go → fetchPTAFResourcesData() // - SW: sw/sw_processor.go → fetchSWResourcesData() (TODO) // resolvePortsForClient загружает сохранённые порты или выделяет новые func resolvePortsForClient(clientInfo ClientInfo, resourcesData []ResourceData, portAllocator *PortAllocator) (map[int][]PortMapping, error) { log.Printf("Чтение сохранённых портов для %s", clientInfo.ClientTitle) savedPortMap := make(map[int][]PortMapping) allPortsFound := true 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 break } } 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) } // 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 }