package main import ( "database/sql" "fmt" "log" "os" "path/filepath" "strings" "time" ) var _ = sql.ErrNoRows // ensure import // getActualContainerCount возвращает реальное количество контейнеров клиента — // максимум из: существующих директорий на диске и drain-маркеров в /tmp. // Используется вместо захардкоженного MaxContainersPerClient. func getActualContainerCount(clientTitle string) int { max := 0 // Считаем по директориям на диске baseDir := filepath.Join("/home/install/ptaf", clientTitle) for num := 1; ; num++ { containerName := fmt.Sprintf("%s-ptaf-agent%03d", clientTitle, num) containerDir := filepath.Join(baseDir, containerName) if _, err := os.Stat(containerDir); err != nil { break } max = num } // Считаем по drain-маркерам в /tmp (контейнеры могли быть уже удалены с диска) for num := max + 1; ; num++ { markerFile := fmt.Sprintf("/tmp/ptaf-drain-ptaf_%s_%03d", clientTitle, num) if _, err := os.Stat(markerFile); err != nil { break } max = num } return max } func processDrainingAngieConfigs() { currentTime := time.Now().Unix() timeout := int64(MigrationDrainTimeout) hostname, err := os.Hostname() if err != nil { log.Printf("⚠ Ошибка получения hostname: %v", err) return } tempConfig := getDefaultConfig() db := connectToDB(tempConfig) defer db.Close() matches, err := filepath.Glob("/tmp/ptaf-angie-drain-*") if err != nil { log.Printf("⚠ Ошибка поиска маркеров Angie: %v", err) return } for _, markerFile := range matches { parts := strings.Split(filepath.Base(markerFile), "-") if len(parts) < 4 { continue } clientTitle := strings.Join(parts[3:], "-") drainTimestamp, found := getAngieConfigDrainTimestamp(clientTitle) if !found { continue } elapsedTime := currentTime - drainTimestamp remainingTime := timeout - elapsedTime angieConfigPath := fmt.Sprintf("/etc/angie/http.d/angie-ptaf-%s.conf", clientTitle) if elapsedTime >= timeout { log.Printf(" Конфиг Angie для %s завершил drain (прошло %d сек), проверяю актуальность...", clientTitle, elapsedTime) clientInfo, err := getClientInfoByClientTitle(db, clientTitle) if err == nil && clientInfo.WAFInstance != "" { belongsNow, checkErr := checkHostBelongsToInstance(db, hostname, clientInfo.WAFInstance) if checkErr == nil && belongsNow { log.Printf(" ⚠ Клиент %s снова принадлежит этому хосту, ОТМЕНЯЮ удаление конфига", clientTitle) os.Remove(markerFile) continue } } log.Printf(" → Клиент %s не принадлежит этому хосту, удаляю конфиг", clientTitle) if _, err := os.Stat(angieConfigPath); err == nil { if err := os.Remove(angieConfigPath); err != nil { log.Printf(" ⚠ Ошибка удаления конфига Angie: %v", err) } else { log.Printf(" ✓ Конфиг Angie удалён: %s", angieConfigPath) os.Remove(markerFile) reloadAngie() } } else { log.Printf(" Конфиг Angie для %s уже не существует, удаляю маркер", clientTitle) os.Remove(markerFile) } } else { log.Printf(" ○ Конфиг Angie для %s в процессе drain [миграция] (осталось %d сек)", clientTitle, remainingTime) } } } func processAllDrainingContainersGlobally(db *sql.DB, hostname string) { log.Println("\n--- Глобальная проверка клиентов и drain-контейнеров ---") entries, err := os.ReadDir("/home/install/ptaf") if err != nil { log.Printf("⚠ Не удалось прочитать директорию: %v", err) return } processedClients := 0 markedForDrain := 0 for _, entry := range entries { if !entry.IsDir() { continue } clientTitle := entry.Name() log.Printf(" → Проверка клиента: %s", clientTitle) clientInfo, err := getClientInfoByClientTitle(db, clientTitle) if err != nil { log.Printf(" ⚠ Ошибка получения данных клиента из БД: %v", err) checkAndProcessExistingDrainMarkers(clientTitle, 0, &processedClients) continue } shouldBeHere := false skipDrainCheck := false if clientInfo.WAFInstance != "" { belongsToInstance, err := checkHostBelongsToInstance(db, hostname, clientInfo.WAFInstance) if err != nil { log.Printf(" ⚠ Ошибка проверки принадлежности хоста: %v", err) } else { shouldBeHere = belongsToInstance } } else { log.Printf(" ℹ waf_instance не указан (NULL), ресурс может быть на любом хосте") shouldBeHere = true skipDrainCheck = true } if !shouldBeHere && !skipDrainCheck { log.Printf(" ⚠ Клиент %s не должен быть на этом хосте (waf_instance='%s')", clientTitle, clientInfo.WAFInstance) hasExistingMarkers := false actualCount := getActualContainerCount(clientTitle) for containerNum := 1; containerNum <= actualCount; containerNum++ { if _, err := os.Stat(fmt.Sprintf("/tmp/ptaf-drain-ptaf_%s_%03d", clientTitle, containerNum)); err == nil { hasExistingMarkers = true break } } if !hasExistingMarkers { log.Printf(" → Помечаю контейнеры для миграции") for containerNum := 1; containerNum <= actualCount; containerNum++ { containerFullName := fmt.Sprintf("ptaf_%s_%03d", clientTitle, containerNum) exists, _ := checkContainerStatus(containerFullName) if exists { if err := markContainerForDrain(containerFullName, true); err != nil { log.Printf(" ⚠ Ошибка пометки контейнера %s: %v", containerFullName, err) } else { log.Printf(" ✓ Контейнер %s помечен для drain (миграция)", containerFullName) markedForDrain++ } } } angieConfigPath := fmt.Sprintf("/etc/angie/http.d/angie-ptaf-%s.conf", clientTitle) if _, err := os.Stat(angieConfigPath); err == nil { if err := markAngieConfigForDrain(clientTitle); err != nil { log.Printf(" ⚠ Ошибка пометки Angie конфига: %v", err) } else { log.Printf(" ✓ Angie конфиг помечен для удаления") } } } processDrainingContainers(clientTitle, 0) processedClients++ } else { log.Printf(" ✓ Клиент должен быть здесь, проверяю устаревшие drain-маркеры") checkAndProcessExistingDrainMarkers(clientTitle, clientInfo.ContainersCount, &processedClients) } } if markedForDrain > 0 { log.Printf("\n✓ Помечено для миграции: %d контейнеров", markedForDrain) } if processedClients > 0 { log.Printf("✓ Обработано клиентов: %d", processedClients) } else { log.Println(" Нет клиентов требующих обработки") } processDrainingAngieConfigs() } func checkAndProcessExistingDrainMarkers(clientTitle string, containersCount int, processedCount *int) { hasAnyDrainMarkers := false actualCount := getActualContainerCount(clientTitle) for containerNum := 1; containerNum <= actualCount; containerNum++ { if _, err := os.Stat(fmt.Sprintf("/tmp/ptaf-drain-ptaf_%s_%03d", clientTitle, containerNum)); err == nil { hasAnyDrainMarkers = true break } } if !hasAnyDrainMarkers { if _, err := os.Stat(fmt.Sprintf("/tmp/ptaf-angie-drain-%s", clientTitle)); err == nil { hasAnyDrainMarkers = true } } if hasAnyDrainMarkers { log.Printf(" → Найдены drain-маркеры, обрабатываю") processDrainingContainers(clientTitle, containersCount) *processedCount++ } } func processDrainingContainers(clientTitle string, maxContainerNum int) { baseDir := filepath.Join("/home/install/ptaf", clientTitle) currentTime := time.Now().Unix() hostname, err := os.Hostname() if err != nil { log.Printf("⚠ Ошибка получения hostname: %v", err) return } tempConfig := getDefaultConfig() db := connectToDB(tempConfig) defer db.Close() actualCount := getActualContainerCount(clientTitle) for containerNum := maxContainerNum + 1; containerNum <= actualCount; containerNum++ { containerName := fmt.Sprintf("%s-ptaf-agent%03d", clientTitle, containerNum) containerFullName := fmt.Sprintf("ptaf_%s_%03d", clientTitle, containerNum) containerDir := filepath.Join(baseDir, containerName) composeFile := filepath.Join(containerDir, "docker-compose.yml") exists, _ := checkContainerStatus(containerFullName) if !exists { markerFile := fmt.Sprintf("/tmp/ptaf-drain-%s", containerFullName) if _, err := os.Stat(markerFile); err == nil { log.Printf(" Удаляю файл-маркер для несуществующего контейнера %s", containerFullName) os.Remove(markerFile) } continue } drainTimestamp, isMigration, isDraining := getContainerDrainInfo(containerFullName) if isDraining { var timeout int64 timeoutName := "масштабирование" if isMigration { timeout = int64(MigrationDrainTimeout) timeoutName = "миграция" } else { timeout = int64(ContainerDrainTimeout) } elapsedTime := currentTime - drainTimestamp remainingTime := timeout - elapsedTime if elapsedTime >= timeout { if isMigration { log.Printf(" Контейнер %s завершил drain [миграция] (прошло %d сек), проверяю актуальность...", containerFullName, elapsedTime) clientInfo, err := getClientInfoByClientTitle(db, clientTitle) if err == nil && clientInfo.WAFInstance != "" { belongsNow, checkErr := checkHostBelongsToInstance(db, hostname, clientInfo.WAFInstance) if checkErr == nil && belongsNow { log.Printf(" ⚠ Клиент %s снова принадлежит этому хосту, ОТМЕНЯЮ удаление контейнера %s", clientTitle, containerFullName) os.Remove(fmt.Sprintf("/tmp/ptaf-drain-%s", containerFullName)) continue } } log.Printf(" → Клиент %s не принадлежит этому хосту, удаляю контейнер %s", clientTitle, containerFullName) } else { log.Printf(" Контейнер %s завершил drain [%s] (прошло %d сек), удаляем", containerFullName, timeoutName, elapsedTime) } if err := stopAndRemoveContainer(containerFullName, composeFile, containerDir); err != nil { log.Printf(" ⚠ Ошибка удаления контейнера %s: %v", containerFullName, err) continue } if err := removeContainerFiles(clientTitle, containerNum); err != nil { log.Printf(" ⚠ Ошибка удаления файлов контейнера %s: %v", containerFullName, err) } else { if err := cleanupClientDirectoryIfEmpty(clientTitle); err != nil { log.Printf(" ⚠ Ошибка очистки директорий клиента: %v", err) } } } else { log.Printf(" ○ Контейнер %s в процессе drain [%s] (осталось %d сек)", containerFullName, timeoutName, remainingTime) } } else { log.Printf(" Контейнер %s больше не нужен, помечаем для drain", containerFullName) if err := markContainerForDrain(containerFullName, false); err != nil { log.Printf(" ⚠ Ошибка пометки контейнера для drain: %v", err) } else { log.Printf(" ✓ Контейнер %s помечен для drain, удаление через %d сек", containerFullName, ContainerDrainTimeout) } } } }