api_sync/main.go

157 lines
3.8 KiB
Go
Executable file

package main
import (
"log"
"os"
"sync"
_ "github.com/lib/pq" // PostgreSQL driver
)
func main() {
// Режим --cd: только проверка дублей и ошибок API
if len(os.Args) > 1 && os.Args[1] == "--cd" {
log.Println("Running duplicate domain/alias check only")
db, err := dbConnect()
if err != nil {
log.Fatal(err)
}
defer db.Close()
// Проверяем ошибки API по всем SID из apps_settings
var allSIDs []int64
clients, err := loadClients(db)
if err == nil {
for _, client := range clients {
ids, err := loadL7IDs(db, client)
if err == nil {
allSIDs = append(allSIDs, ids...)
}
}
}
// Проверяем доступность API параллельно
sidChan := make(chan int64, len(allSIDs))
var wg sync.WaitGroup
var mu sync.Mutex
apiErrors := make(map[int64]string)
for i := 0; i < maxConcurrentWorkers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for sid := range sidChan {
if err := checkSIDAvailability(sid); err != nil {
mu.Lock()
apiErrors[sid] = err.Error()
mu.Unlock()
}
}
}()
}
for _, sid := range allSIDs {
sidChan <- sid
}
close(sidChan)
wg.Wait()
sendAPIErrorsAlert(apiErrors)
if err := checkDuplicateDomains(db); err != nil {
log.Printf("Duplicate check error: %v", err)
}
return
}
log.Println("Start sync sp_info")
db, err := dbConnect()
if err != nil {
log.Fatal(err)
}
defer db.Close()
clients, err := loadClients(db)
if err != nil {
log.Fatal(err)
}
// Собираем все SID для обработки
var allSIDs []int64
for _, client := range clients {
ids, err := loadL7IDs(db, client)
if err != nil {
log.Println("binding error:", err)
continue
}
allSIDs = append(allSIDs, ids...)
}
log.Printf("Total SIDs to process: %d", len(allSIDs))
log.Printf("Using %d concurrent workers", maxConcurrentWorkers)
// Создаём канал для заданий и WaitGroup для ожидания
sidChan := make(chan int64, len(allSIDs))
var wg sync.WaitGroup
var mu sync.Mutex
apiErrors := make(map[int64]string)
// Запускаем воркеры
for i := 0; i < maxConcurrentWorkers; i++ {
wg.Add(1)
go func(workerID int) {
defer wg.Done()
for sid := range sidChan {
log.Printf("[Worker %d] Processing SID %d", workerID, sid)
if err := processSID(db, sid); err != nil {
log.Printf("[Worker %d] SID %d error: %v", workerID, sid, err)
mu.Lock()
apiErrors[sid] = err.Error()
mu.Unlock()
}
}
}(i)
}
// Отправляем все SID в канал
for _, sid := range allSIDs {
sidChan <- sid
}
close(sidChan)
// Ждём завершения всех воркеров
wg.Wait()
log.Println("Finish sync sp_info")
// Отправляем сводный алерт по ошибкам API
if checkDuplicates {
sendAPIErrorsAlert(apiErrors)
}
// Удаляем из sp_info ресурсы, которых больше нет в apps_settings
log.Println("Start cleanup of removed SIDs")
if err := cleanupRemovedSIDs(db, allSIDs); err != nil {
log.Printf("Cleanup error: %v", err)
}
log.Println("Finish cleanup of removed SIDs")
// Проверка дублирующихся domain_name и aliases между разными SID
if checkDuplicates {
log.Println("Start duplicate domain/alias check")
if err := checkDuplicateDomains(db); err != nil {
log.Printf("Duplicate check error: %v", err)
}
log.Println("Finish duplicate domain/alias check")
} else {
log.Println("Duplicate domain/alias check is disabled")
}
// Обновление WAF whitelist
log.Println("Start WAF whitelist update")
if err := updateWAFWhitelist(db); err != nil {
log.Printf("WAF whitelist update error: %v", err)
}
log.Println("Finish WAF whitelist update")
}