299 lines
13 KiB
Go
299 lines
13 KiB
Go
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
|
||
}
|