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") }