534 lines
16 KiB
Go
534 lines
16 KiB
Go
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")
|
||
}
|
||
}
|