auspex/portcheck.go
2026-06-09 17:12:18 +03:00

539 lines
16 KiB
Go
Executable file
Raw Blame History

This file contains invisible Unicode characters

This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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 (
"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")
}
if err := mattermostSend(message); err != nil {
fmt.Printf("❌ Mattermost: Ошибка отправки: %v\n", err)
} else {
fmt.Println("✓ Уведомление отправлено в Mattermost")
}
}