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

448 lines
19 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"
"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.Printf("⚠ Ошибка подготовки Docker образа: %v", err)
log.Printf(" → Продолжаем работу, но новые PTAF-контейнеры могут не запуститься")
}
} 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.
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()
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)
}
if written == 0 {
os.Remove(destPath)
return fmt.Errorf("скачан пустой файл (HTTP %d)", resp.StatusCode)
}
log.Printf(" → Скачано: %.1f МБ (HTTP %d)", float64(written)/1024/1024, resp.StatusCode)
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(" Инициирую удаление контейнеров с текущего хоста (миграция)")
actualCount := getActualContainerCount(clientInfo.ClientTitle)
for containerNum := 1; containerNum <= actualCount; 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, true); 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
actualCount := getActualContainerCount(clientInfo.ClientTitle)
for containerNum := 1; containerNum <= actualCount; 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
}