magos/main.go

550 lines
23 KiB
Go
Executable file
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 (
"context"
"database/sql"
"fmt"
"io"
"log"
"net"
"net/http"
"os"
"os/exec"
"path/filepath"
"strings"
"time"
)
var lockStatusFile *os.File
// updateLockStatus обновляет статус в lock-файле для диагностики зависаний
func updateLockStatus(status string) {
if lockStatusFile == nil {
return
}
lockStatusFile.Seek(0, 0)
lockStatusFile.Truncate(0)
lockStatusFile.WriteString(fmt.Sprintf("%d\n%s\n", os.Getpid(), status))
}
// checkDockerAvailable проверяет доступность Docker daemon с таймаутом
func checkDockerAvailable() error {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
cmd := exec.CommandContext(ctx, "docker", "info", "--format", "{{.ServerVersion}}")
output, err := cmd.Output()
if ctx.Err() == context.DeadlineExceeded {
return fmt.Errorf("таймаут подключения к Docker daemon (10 сек) — возможно Docker не запущен")
}
if err != nil {
return fmt.Errorf("ошибка подключения к Docker daemon: %w", err)
}
log.Printf("✓ Docker daemon доступен, версия: %s", strings.TrimSpace(string(output)))
return nil
}
func main() {
// ════════════════════════════════════════════════════════════
// Lock-файл для предотвращения параллельного запуска
// ════════════════════════════════════════════════════════════
lockFile := "/var/lock/ptaf-main.lock"
const maxLockAge = 15 * time.Minute
f, err := os.OpenFile(lockFile, os.O_CREATE|os.O_EXCL|os.O_RDWR, 0600)
if err != nil {
if os.IsExist(err) {
data, readErr := os.ReadFile(lockFile)
if readErr == nil {
lines := strings.SplitN(strings.TrimSpace(string(data)), "\n", 2)
var existingPID int
if _, parseErr := fmt.Sscanf(lines[0], "%d", &existingPID); parseErr == nil && existingPID > 0 {
// Проверяем существует ли процесс
if _, procErr := os.Stat(fmt.Sprintf("/proc/%d", existingPID)); os.IsNotExist(procErr) {
log.Printf("⚠ Lock файл найден но процесс PID %d не существует — удаляем устаревший lock", existingPID)
os.Remove(lockFile)
f, err = os.OpenFile(lockFile, os.O_CREATE|os.O_EXCL|os.O_RDWR, 0600)
if err != nil {
log.Fatal("❌ Не удалось создать lock файл после очистки: ", err)
}
} else {
// Процесс существует — проверяем возраст и показываем статус
status := "unknown"
if len(lines) > 1 && lines[1] != "" {
status = lines[1]
}
if info, statErr := os.Stat(lockFile); statErr == nil {
age := time.Since(info.ModTime())
if age > maxLockAge {
log.Printf("⚠ Программа уже запущена (PID %d) но работает более %.0f минут — возможно зависла", existingPID, age.Minutes())
log.Printf(" Текущий статус: %s", status)
log.Printf(" Если утилита зависла, выполните: kill %d && rm %s", existingPID, lockFile)
} else {
log.Printf(" Текущий статус: %s", status)
}
}
log.Fatalf("❌ Программа уже запущена! PID: %d, статус: %s, Lock файл: %s", existingPID, status, lockFile)
}
} else {
log.Fatal("❌ Программа уже запущена! Lock файл существует: ", lockFile)
}
} else {
log.Fatal("❌ Программа уже запущена! Lock файл существует: ", lockFile)
}
} else {
log.Fatal("❌ Не удалось создать lock файл: ", err)
}
}
pid := os.Getpid()
f.WriteString(fmt.Sprintf("%d\n", pid))
lockStatusFile = f
defer func() {
lockStatusFile = nil
f.Close()
os.Remove(lockFile)
log.Printf("✓ Lock файл удален")
}()
log.Printf("✓ Lock файл создан: %s (PID: %d)", lockFile, pid)
updateLockStatus("init")
// ════════════════════════════════════════════════════════════
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)
}
// Проверяем доступность Docker daemon
if err := checkDockerAvailable(); err != nil {
log.Fatalf("❌ Docker daemon недоступен: %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-клиента отдельно
ptafClients := getPTAFClients(db, instanceNames)
if len(ptafClients) == 0 {
log.Printf(" PTAF-клиентов на хосте нет, проверка Docker образов пропускается")
} else {
checkedImages := make(map[string]bool) // не проверяем один образ дважды
for _, clientInfo := range ptafClients {
image, url := getClientDockerImage(clientInfo, config)
if image == "" || checkedImages[image] {
continue
}
checkedImages[image] = true
log.Printf(" → Проверка Docker образа для клиента %s: %s", clientInfo.ClientTitle, image)
if err := ensureDockerImageExists(image, url); err != nil {
log.Printf("⚠ Ошибка подготовки Docker образа для %s: %v", clientInfo.ClientTitle, err)
log.Printf(" → Продолжаем работу, но контейнеры %s могут не запуститься", clientInfo.ClientTitle)
}
}
}
// Глобальная обработка 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=== Программа завершена успешно ===")
}
// getPTAFClients возвращает список PTAF-клиентов на хосте
func getPTAFClients(db *sql.DB, instanceNames []string) []ClientInfo {
var ptafClients []ClientInfo
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" {
ptafClients = append(ptafClients, clientInfo)
}
}
}
return ptafClients
}
// getClientDockerImage возвращает образ и URL для клиента.
// Если у клиента заданы свои значения — использует их, иначе fallback на config.
func getClientDockerImage(clientInfo ClientInfo, config Config) (image string, url string) {
if clientInfo.DockerImage.Valid && clientInfo.DockerImage.String != "" {
image = clientInfo.DockerImage.String
} else {
image = config.DockerImage
}
if clientInfo.DockerImageDownload.Valid && clientInfo.DockerImageDownload.String != "" {
url = clientInfo.DockerImageDownload.String
} else {
url = config.DockerImageURL
}
return
}
// ensureDockerImageExists проверяет наличие Docker образа локально.
// Если образ отсутствует — скачивает tar-архив по url и загружает через docker load.
func ensureDockerImageExists(image, url string) error {
if image == "" {
log.Printf("⚠ Docker образ не задан, пропускаем проверку")
return nil
}
// Проверяем есть ли образ локально
cmd := exec.Command("docker", "image", "inspect", image)
if err := cmd.Run(); err == nil {
log.Printf("✓ Docker образ уже присутствует: %s", image)
return nil
}
log.Printf(" Docker образ не найден локально: %s", image)
if url == "" {
log.Printf("⚠ Docker образ %s отсутствует локально, URL для скачивания не задан — пропускаем", image)
return nil
}
// Путь для tar-архива: /home/install/{image_name}.tar
safeName := strings.ReplaceAll(image, ":", "_")
safeName = strings.ReplaceAll(safeName, "/", "_")
tarPath := fmt.Sprintf("/home/install/%s.tar", safeName)
// Скачиваем архив если ещё не скачан
if _, err := os.Stat(tarPath); os.IsNotExist(err) {
log.Printf(" → Скачивание образа с %s", url)
log.Printf(" → Сохранение в %s", tarPath)
if err := downloadFile(tarPath, url); 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", image)
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)
updateLockStatus(fmt.Sprintf("processing client: %s", clientInfo.ClientTitle))
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 := buildContainerFullName(clientInfo.ClientTitle, containerNum, hostname)
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, hostname string) {
log.Printf(" → Проверка и очистка устаревших маркеров drain для %s", clientInfo.ClientTitle)
cleanedContainers := 0
cleanedAngieMarker := false
actualCount := getActualContainerCount(clientInfo.ClientTitle)
for containerNum := 1; containerNum <= actualCount; containerNum++ {
containerFullName := buildContainerFullName(clientInfo.ClientTitle, containerNum, hostname)
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
}