package main import ( "bufio" "database/sql" "encoding/json" "fmt" "net" "os" "path/filepath" "regexp" "sort" "strings" "time" _ "github.com/lib/pq" ) // ── Типы ───────────────────────────────────────────────────────────────────── type pcEndpoint struct { Host string Port string } func (e pcEndpoint) String() string { return e.Host + ":" + e.Port } type pcEndpointStatus struct { Endpoint string `json:"endpoint"` Available bool `json:"available"` Error string `json:"error,omitempty"` ConfigFile string `json:"config_file"` } type pcStateFile struct { Timestamp string `json:"timestamp"` Status map[string]pcEndpointStatus `json:"status"` } // ── База данных ────────────────────────────────────────────────────────────── // pcConnectDB подключается к БД с failover func pcConnectDB() (*sql.DB, error) { // Пробуем primary primaryDSN := fmt.Sprintf("host=%s port=%d user=%s password=%s dbname=%s sslmode=disable", cfgPrimaryDBHost, cfgDBPort, cfgDBUser, cfgDBPassword, cfgDBName) db, err := sql.Open("postgres", primaryDSN) if err == nil { if err := db.Ping(); err == nil { return db, nil } db.Close() } // Пробуем secondary secondaryDSN := fmt.Sprintf("host=%s port=%d user=%s password=%s dbname=%s sslmode=disable", cfgSecondaryDBHost, cfgDBPort, cfgDBUser, cfgDBPassword, cfgDBName) db, err = sql.Open("postgres", secondaryDSN) if err != nil { return nil, fmt.Errorf("не удалось открыть соединение с secondary: %v", err) } if err := db.Ping(); err != nil { db.Close() return nil, fmt.Errorf("secondary БД недоступна: %v", err) } return db, nil } // pcExtractL7ResourceID извлекает l7resourceid из имени файла // Например: "ITAR-TASS_17822.conf" -> "17822" func pcExtractL7ResourceID(filename string) string { // Убираем путь, оставляем только имя файла baseName := filepath.Base(filename) // Убираем расширение nameWithoutExt := strings.TrimSuffix(baseName, filepath.Ext(baseName)) // Ищем цифры после последнего подчеркивания parts := strings.Split(nameWithoutExt, "_") if len(parts) >= 2 { // Берем последнюю часть (она должна содержать цифры) return parts[len(parts)-1] } return "" } // pcGetCheckPortsStatus проверяет нужно ли проверять порты для данного l7resourceid func pcGetCheckPortsStatus(db *sql.DB, l7resourceid string) (bool, error) { if l7resourceid == "" { return false, nil } query := `SELECT COALESCE(check_ports, false) FROM apps_settings WHERE l7resourceid = $1` var checkPorts bool err := db.QueryRow(query, l7resourceid).Scan(&checkPorts) if err != nil { if err == sql.ErrNoRows { // Если записи нет в БД, по умолчанию не проверяем return false, nil } return false, fmt.Errorf("ошибка запроса к БД: %v", err) } return checkPorts, nil } // pcLoadCheckPortsMap загружает все статусы check_ports для оптимизации func pcLoadCheckPortsMap(db *sql.DB) (map[string]bool, error) { query := `SELECT l7resourceid, COALESCE(check_ports, false) FROM apps_settings` rows, err := db.Query(query) if err != nil { return nil, fmt.Errorf("ошибка запроса к БД: %v", err) } defer rows.Close() checkMap := make(map[string]bool) for rows.Next() { var l7resourceid string var checkPorts bool if err := rows.Scan(&l7resourceid, &checkPorts); err != nil { continue } checkMap[l7resourceid] = checkPorts } return checkMap, rows.Err() } // ── Точка входа ────────────────────────────────────────────────────────────── func runPortCheck() error { hostname, _ := os.Hostname() if hostname == "" { hostname = "unknown" } fmt.Printf("=== Port Checker запущен на %s ===\n\n", hostname) // Подключаемся к БД fmt.Println("📊 Подключение к базе данных...") db, err := pcConnectDB() if err != nil { fmt.Printf("⚠️ Не удалось подключиться к БД: %v\n", err) fmt.Println("⚠️ Продолжаем без фильтрации по check_ports (проверяем все порты)\n") db = nil } else { fmt.Println("✅ Соединение с БД установлено\n") defer db.Close() } // Загружаем карту check_ports для оптимизации var checkPortsMap map[string]bool if db != nil { checkPortsMap, err = pcLoadCheckPortsMap(db) if err != nil { fmt.Printf("⚠️ Ошибка загрузки check_ports: %v\n", err) fmt.Println("⚠️ Продолжаем без фильтрации\n") checkPortsMap = nil } else { fmt.Printf("📋 Загружено статусов check_ports: %d\n\n", len(checkPortsMap)) } } files, err := filepath.Glob(cfgNginxConfigGlob) if err != nil { return fmt.Errorf("ошибка при поиске конфигурационных файлов: %v", err) } if len(files) == 0 { fmt.Printf("Конфигурационные файлы не найдены по пути: %s\n", cfgNginxConfigGlob) return nil } previousState := pcLoadState() currentState := make(map[string]pcEndpointStatus) fileEndpoints := make(map[string][]pcEndpoint) skippedConfigs := 0 for _, file := range files { // Извлекаем l7resourceid из имени файла l7resourceid := pcExtractL7ResourceID(file) // Проверяем нужно ли проверять этот конфиг if checkPortsMap != nil && l7resourceid != "" { checkPorts, exists := checkPortsMap[l7resourceid] if exists && !checkPorts { // check_ports = false, пропускаем skippedConfigs++ continue } // Если записи нет в БД или check_ports = true, проверяем } eps := pcParseNginxConfig(file) if len(eps) > 0 { fileEndpoints[file] = eps } } if skippedConfigs > 0 { fmt.Printf("ℹ️ Пропущено конфигов (check_ports = false): %d\n", skippedConfigs) } if len(fileEndpoints) == 0 { fmt.Println("Серверы не найдены в конфигурационных файлах") return nil } // Уникальные endpoints globalEps := make(map[pcEndpoint]string) for file, eps := range fileEndpoints { baseName := filepath.Base(file) cfgName := strings.TrimSuffix(baseName, filepath.Ext(baseName)) for _, ep := range eps { if _, exists := globalEps[ep]; !exists { globalEps[ep] = cfgName } } } fmt.Printf("Найдено уникальных endpoints: %d\n\n", len(globalEps)) configGroups := make(map[string][]pcEndpoint) for ep, cfgName := range globalEps { configGroups[cfgName] = append(configGroups[cfgName], ep) } sortedConfigs := make([]string, 0, len(configGroups)) for cfg := range configGroups { sortedConfigs = append(sortedConfigs, cfg) } sort.Strings(sortedConfigs) // Параллельная проверка type checkResult struct { status pcEndpointStatus configName string } results := make(chan checkResult, len(globalEps)) checkCount := 0 for _, cfgName := range sortedConfigs { eps := configGroups[cfgName] sort.Slice(eps, func(i, j int) bool { if eps[i].Host != eps[j].Host { return eps[i].Host < eps[j].Host } return eps[i].Port < eps[j].Port }) for _, ep := range eps { checkCount++ go func(e pcEndpoint, c string) { results <- checkResult{status: pcCheckPort(e, c), configName: c} }(ep, cfgName) } } statusByConfig := make(map[string][]pcEndpointStatus) for i := 0; i < checkCount; i++ { r := <-results currentState[r.status.Endpoint] = r.status statusByConfig[r.configName] = append(statusByConfig[r.configName], r.status) } close(results) // Вывод результатов for _, cfgName := range sortedConfigs { statuses := statusByConfig[cfgName] sort.Slice(statuses, func(i, j int) bool { return statuses[i].Endpoint < statuses[j].Endpoint }) fmt.Printf("=== %s ===\n", cfgName) for _, s := range statuses { pcPrintStatus(s) } fmt.Println() } // Сравнение с предыдущим if len(previousState.Status) > 0 { changes := pcCompareStates(previousState.Status, currentState) if len(changes) > 0 { fmt.Println("=== Изменения с последней проверки ===") for _, ch := range changes { fmt.Println(ch) } fmt.Println() if isTelegramConfigured() { pcSendNotification(changes) } else { fmt.Println("ℹ️ Telegram уведомления отключены (не настроены данные)") } } } pcSaveState(currentState) fmt.Println("=== Проверка завершена ===") return nil } // ── Парсинг nginx ──────────────────────────────────────────────────────────── func pcParseNginxConfig(filename string) []pcEndpoint { file, err := os.Open(filename) if err != nil { return nil } defer file.Close() var endpoints []pcEndpoint inUpstream := false scanner := bufio.NewScanner(file) re := regexp.MustCompile(`^\s*server\s+([0-9.]+):(\d+)`) for scanner.Scan() { line := strings.TrimSpace(scanner.Text()) if strings.HasPrefix(line, "upstream") { inUpstream = true continue } if inUpstream && line == "}" { inUpstream = false continue } if inUpstream { m := re.FindStringSubmatch(line) if len(m) == 3 { endpoints = append(endpoints, pcEndpoint{Host: m[1], Port: m[2]}) } } } return endpoints } // ── Проверка порта ─────────────────────────────────────────────────────────── func pcCheckPort(ep pcEndpoint, cfgName string) pcEndpointStatus { address := ep.String() timeout := time.Duration(cfgCheckTimeoutSec) * time.Second delay := time.Duration(cfgRetryDelaySec) * time.Second var lastErr error for attempt := 1; attempt <= cfgMaxRetries; attempt++ { conn, err := net.DialTimeout("tcp", address, timeout) if err == nil { conn.Close() return pcEndpointStatus{Endpoint: address, Available: true, ConfigFile: cfgName} } lastErr = err if attempt < cfgMaxRetries { time.Sleep(delay) } } return pcEndpointStatus{Endpoint: address, Available: false, Error: lastErr.Error(), ConfigFile: cfgName} } func pcPrintStatus(s pcEndpointStatus) { if s.Available { fmt.Printf("%s - ДОСТУПЕН\n", s.Endpoint) } else { msg := s.Error switch { case strings.Contains(msg, "i/o timeout"): msg = "timeout" case strings.Contains(msg, "connection refused"): msg = "connection refused" case strings.Contains(msg, "no route to host"): msg = "no route to host" } fmt.Printf("%s - НЕДОСТУПЕН (%s)\n", s.Endpoint, msg) } } // ── Состояние ──────────────────────────────────────────────────────────────── func pcLoadState() pcStateFile { data, err := os.ReadFile(cfgStateFilePath) if err != nil { return pcStateFile{Status: make(map[string]pcEndpointStatus)} } var state pcStateFile if json.Unmarshal(data, &state) != nil { return pcStateFile{Status: make(map[string]pcEndpointStatus)} } return state } func pcSaveState(status map[string]pcEndpointStatus) { state := pcStateFile{ Timestamp: time.Now().Format("2006-01-02 15:04:05"), Status: status, } data, err := json.MarshalIndent(state, "", " ") if err != nil { fmt.Printf("Ошибка при сохранении состояния: %v\n", err) return } if err := os.WriteFile(cfgStateFilePath, data, 0644); err != nil { fmt.Printf("Ошибка при записи файла состояния: %v\n", err) } } // ── Сравнение состояний ────────────────────────────────────────────────────── func pcCompareStates(prev, curr map[string]pcEndpointStatus) []string { changesByConfig := make(map[string][]string) for ep, cs := range curr { if ps, exists := prev[ep]; exists { if ps.Available != cs.Available { cfg := cs.ConfigFile if cs.Available { changesByConfig[cfg] = append(changesByConfig[cfg], fmt.Sprintf("%s: НЕДОСТУПЕН → ДОСТУПЕН", ep)) } else { changesByConfig[cfg] = append(changesByConfig[cfg], fmt.Sprintf("%s: ДОСТУПЕН → НЕДОСТУПЕН (%s)", ep, pcShortError(cs.Error))) } } } else { cfg := cs.ConfigFile if cs.Available { changesByConfig[cfg] = append(changesByConfig[cfg], fmt.Sprintf("%s: НОВЫЙ (доступен)", ep)) } else { changesByConfig[cfg] = append(changesByConfig[cfg], fmt.Sprintf("%s: НОВЫЙ (недоступен - %s)", ep, pcShortError(cs.Error))) } } } for ep, ps := range prev { if _, exists := curr[ep]; !exists { changesByConfig[ps.ConfigFile] = append(changesByConfig[ps.ConfigFile], fmt.Sprintf("%s: УДАЛЕН из конфигурации", ep)) } } var result []string configs := make([]string, 0, len(changesByConfig)) for c := range changesByConfig { configs = append(configs, c) } sort.Strings(configs) for _, c := range configs { chs := changesByConfig[c] sort.Strings(chs) result = append(result, fmt.Sprintf("[%s]", c)) result = append(result, chs...) result = append(result, "") } if len(result) > 0 && result[len(result)-1] == "" { result = result[:len(result)-1] } return result } func pcShortError(full string) string { switch { case strings.Contains(full, "i/o timeout"): return "timeout" case strings.Contains(full, "connection refused"): return "connection refused" case strings.Contains(full, "no route to host"): return "no route" case strings.Contains(full, "network is unreachable"): return "network unreachable" } if len(full) > 50 { return full[:50] + "..." } return full } // ── Telegram ───────────────────────────────────────────────────────────────── func pcSendNotification(changes []string) { hostname, _ := os.Hostname() if hostname == "" { hostname = "unknown" } message := fmt.Sprintf("🔔 Изменения в статусе портов на %s\n\n", hostname) for _, line := range changes { if strings.HasPrefix(line, "[") && strings.HasSuffix(line, "]") { cfgName := strings.TrimPrefix(strings.TrimSuffix(line, "]"), "[") message += fmt.Sprintf("\n--- %s ---\n", cfgName) continue } if line == "" { continue } switch { case strings.Contains(line, "НЕДОСТУПЕН → ДОСТУПЕН"): message += "✅ " + line + "\n" case strings.Contains(line, "ДОСТУПЕН → НЕДОСТУПЕН"): message += "❌ " + line + "\n" case strings.Contains(line, "НОВЫЙ"): message += "🆕 " + line + "\n" case strings.Contains(line, "УДАЛЕН"): message += "🗑️ " + line + "\n" default: message += line + "\n" } } message += fmt.Sprintf("\n⏰ Время проверки: %s", time.Now().Format("2006-01-02 15:04:05")) if err := telegramSendPlainText(message); err != nil { fmt.Printf("❌ Telegram: Ошибка отправки: %v\n", err) } else { fmt.Println("✓ Уведомление отправлено в Telegram") } if err := emailSendHTML("Auspex: Изменения в статусе портов", message); err != nil { fmt.Printf("❌ Email: Ошибка отправки: %v\n", err) } else { fmt.Println("✓ Уведомление отправлено на Email") } }