magos/drain.go
2026-03-16 17:43:14 +03:00

310 lines
12 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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()
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, isMigration, found := getAngieConfigDrainInfo(clientTitle)
if !found {
continue
}
var timeout int64
timeoutName := "масштабирование"
if isMigration {
timeout = int64(MigrationDrainTimeout)
timeoutName = "миграция"
} else {
timeout = int64(ContainerDrainTimeout)
}
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 [%s] (прошло %d сек), проверяю актуальность...", clientTitle, timeoutName, elapsedTime)
if isMigration {
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)
} else {
log.Printf(" → Удаляю конфиг Angie для %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 [%s] (осталось %d сек)", clientTitle, timeoutName, 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, true); 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)
}
}
}
}