api_sync/database.go

615 lines
16 KiB
Go
Executable file
Raw Blame History

This file contains ambiguous Unicode characters

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 (
"database/sql"
"encoding/json"
"fmt"
"log"
"sort"
"strings"
"time"
)
/* DATABASE */
func dbConnect() (*sql.DB, error) {
dsn := fmt.Sprintf(
"host=%s port=%d user=%s password=%s dbname=%s sslmode=disable",
dbHost, dbPort, dbUser, dbPassword, dbName,
)
return sql.Open("postgres", dsn)
}
func loadClients(db *sql.DB) ([]string, error) {
rows, err := db.Query(`SELECT client_title FROM client_info`)
if err != nil {
return nil, err
}
defer rows.Close()
var result []string
for rows.Next() {
var c string
if err := rows.Scan(&c); err != nil {
return nil, err
}
result = append(result, c)
}
return result, nil
}
func loadL7IDs(db *sql.DB, client string) ([]int64, error) {
rows, err := db.Query(
`SELECT l7resourceid FROM apps_settings WHERE client_title = $1 AND mode = 'auto'`,
client,
)
if err != nil {
return nil, err
}
defer rows.Close()
var ids []int64
for rows.Next() {
var id int64
if err := rows.Scan(&id); err != nil {
return nil, err
}
ids = append(ids, id)
}
return ids, nil
}
func processSID(db *sql.DB, sid int64) error {
// WHOIS
var whois whoisResponse
if err := apiGet(whoisURL+fmt.Sprint(sid)+"/global", &whois); err != nil {
return fmt.Errorf("whois: %w", err)
}
// ORIGINS
var originsResp apiList[originItem]
if err := apiGet(originURL+fmt.Sprint(sid), &originsResp); err != nil {
return fmt.Errorf("origins: %w", err)
}
// ALIASES
var aliasResp apiList[aliasItem]
if err := apiGet(aliasURL+fmt.Sprint(sid), &aliasResp); err != nil {
return fmt.Errorf("aliases: %w", err)
}
originsJSON, _ := json.Marshal(originsResp.Data.Result.Items)
var aliases []string
var latestModifiedAt int64 = 0
for _, a := range aliasResp.Data.Result.Items {
aliases = append(aliases, a.Domain)
if a.ModifiedAt > latestModifiedAt {
latestModifiedAt = a.ModifiedAt
}
}
aliasesJSON, _ := json.Marshal(aliases)
// Проверяем AntiDDOS
antiddosEnabled := false
protectedIP := whois.Data.Result.ProtectedIp
var resolvedIP string
if protectedIP != "" && whois.Data.Result.L7ResourceName != "" {
// Резолвим домен
resolvedIPAddr, err := resolveDomain(whois.Data.Result.L7ResourceName)
if err != nil {
log.Printf("[AntiDDOS Check] Failed to resolve %s: %v", whois.Data.Result.L7ResourceName, err)
} else {
resolvedIP = resolvedIPAddr
log.Printf("[AntiDDOS Check] Domain: %s, Resolved: %s, Protected: %s",
whois.Data.Result.L7ResourceName, resolvedIP, protectedIP)
if resolvedIP == protectedIP {
antiddosEnabled = true
log.Printf("[AntiDDOS Check] ✅ AntiDDOS is ENABLED for %s", whois.Data.Result.L7ResourceName)
} else {
log.Printf("[AntiDDOS Check] ❌ AntiDDOS is DISABLED for %s", whois.Data.Result.L7ResourceName)
}
}
}
// Извлекаем WAF данные из API
wafEnabledSP := whois.Data.Result.WafEnabled
wafVendor := whois.Data.Result.WafProvider
instanceSP := whois.Data.Result.WafInstance
log.Printf("[WAF Info] SID: %d, WAF Enabled: %d, Vendor: %s, Instance: %s",
sid, wafEnabledSP, wafVendor, instanceSP)
return upsertSPInfo(
db,
sid,
whois.Data.Result.L7ResourceName,
originsResp.Data.Result.Items,
aliases,
originsJSON,
aliasesJSON,
latestModifiedAt,
protectedIP,
resolvedIP,
antiddosEnabled,
wafEnabledSP,
wafVendor,
instanceSP,
)
}
func upsertSPInfo(
db *sql.DB,
sid int64,
domain string,
newOrigins []originItem,
newAliases []string,
originsJSON []byte,
aliasesJSON []byte,
aliasesModifiedAt int64,
protectedIP string,
resolvedIP string,
antiddosEnabled bool,
wafEnabledSP int,
wafVendor string,
instanceSP string,
) error {
// Получаем client_title из apps_settings и waf_provider из client_info
var clientTitle string
db.QueryRow(`SELECT client_title FROM apps_settings WHERE l7resourceid = $1 LIMIT 1`, sid).Scan(&clientTitle)
if clientTitle == "" {
clientTitle = "—"
}
var wafProviderClient string
if clientTitle != "—" {
db.QueryRow(`SELECT waf_provider FROM client_info WHERE client_title = $1`, clientTitle).Scan(&wafProviderClient)
}
if wafProviderClient == "" {
wafProviderClient = "—"
}
// Проверяем существующие данные
var existingDomain string
var existingOriginsJSON, existingAliasesJSON []byte
var existingWafEnabled sql.NullInt64
var existingWafVendor, existingInstanceSP sql.NullString
err := db.QueryRow(`
SELECT domain_name, origins, aliases, waf_enabled_sp, waf_vendor, instance_sp
FROM sp_info
WHERE sid = $1
`, sid).Scan(&existingDomain, &existingOriginsJSON, &existingAliasesJSON,
&existingWafEnabled, &existingWafVendor, &existingInstanceSP)
isNewRecord := err == sql.ErrNoRows
if err != nil && !isNewRecord {
return err
}
hasChanges := false
var changeDetails []string
needsPTAFUpdate := false
hasWAFChanges := false
var wafChangeDetails []string
if !isNewRecord {
var existingOrigins []originItem
var existingAliases []string
json.Unmarshal(existingOriginsJSON, &existingOrigins)
json.Unmarshal(existingAliasesJSON, &existingAliases)
if existingDomain != domain {
hasChanges = true
needsPTAFUpdate = true
changeDetails = append(changeDetails,
fmt.Sprintf("Домен изменён: %s -> %s", existingDomain, domain))
}
originChanges := compareOrigins(existingOrigins, newOrigins)
if originChanges != "" {
hasChanges = true
changeDetails = append(changeDetails, originChanges)
}
aliasChanges, aliasesChanged := compareAliases(existingAliases, newAliases)
if aliasesChanged {
hasChanges = true
needsPTAFUpdate = true
changeDetails = append(changeDetails, aliasChanges)
}
// Проверяем изменения WAF настроек
oldWafEnabled := 0
if existingWafEnabled.Valid {
oldWafEnabled = int(existingWafEnabled.Int64)
}
oldWafVendor := ""
if existingWafVendor.Valid {
oldWafVendor = existingWafVendor.String
}
oldInstanceSP := ""
if existingInstanceSP.Valid {
oldInstanceSP = existingInstanceSP.String
}
if oldWafEnabled != wafEnabledSP {
hasWAFChanges = true
wafChangeDetails = append(wafChangeDetails,
fmt.Sprintf("WAF Enabled: %d -> %d", oldWafEnabled, wafEnabledSP))
}
if oldWafVendor != wafVendor {
hasWAFChanges = true
oldDisplay := oldWafVendor
if oldDisplay == "" {
oldDisplay = "(не указан)"
}
newDisplay := wafVendor
if newDisplay == "" {
newDisplay = "(не указан)"
}
wafChangeDetails = append(wafChangeDetails,
fmt.Sprintf("WAF Provider: %s -> %s", oldDisplay, newDisplay))
}
if oldInstanceSP != instanceSP {
hasWAFChanges = true
oldDisplay := oldInstanceSP
if oldDisplay == "" {
oldDisplay = "(не указан)"
}
newDisplay := instanceSP
if newDisplay == "" {
newDisplay = "(не указан)"
}
wafChangeDetails = append(wafChangeDetails,
fmt.Sprintf("WAF Instance: %s -> %s", oldDisplay, newDisplay))
}
}
var aliasesModifiedTime *time.Time
if aliasesModifiedAt > 0 {
t := time.Unix(aliasesModifiedAt, 0)
aliasesModifiedTime = &t
}
_, err = db.Exec(`
INSERT INTO sp_info (sid, domain_name, origins, aliases, aliaces_in_sp_modified, protected_ip, resolved_ip, antiddos_enable, waf_enabled_sp, waf_vendor, instance_sp)
VALUES ($1, $2, $3::jsonb, $4::jsonb, $5, $6, $7, $8, $9, $10, $11)
ON CONFLICT (sid) DO UPDATE SET
domain_name = EXCLUDED.domain_name,
origins = EXCLUDED.origins,
aliases = EXCLUDED.aliases,
aliaces_in_sp_modified = EXCLUDED.aliaces_in_sp_modified,
protected_ip = EXCLUDED.protected_ip,
resolved_ip = EXCLUDED.resolved_ip,
antiddos_enable = EXCLUDED.antiddos_enable,
waf_enabled_sp = EXCLUDED.waf_enabled_sp,
waf_vendor = EXCLUDED.waf_vendor,
instance_sp = EXCLUDED.instance_sp,
updated_at = now()
`,
sid,
domain,
originsJSON,
aliasesJSON,
aliasesModifiedTime,
protectedIP,
resolvedIP,
antiddosEnabled,
wafEnabledSP,
wafVendor,
instanceSP,
)
if err != nil {
return err
}
// Отправляем алерт для обычных изменений
if hasChanges {
ptafWarning := ""
if needsPTAFUpdate {
ptafWarning = "\n<b>⚠️ ВНИМАНИЕ: Необходимо внести изменения в кабинете " + wafVendorName(wafVendor) + "!</b>"
}
now := time.Now().In(time.FixedZone("UTC+3", 3*60*60))
message := fmt.Sprintf(
"🟡 <b>Обновление WAF Info</b>\n\n"+
"<b>SID:</b> %d <b>Домен:</b> %s <b>TENANT:</b> %s <b>WAF provider:</b> %s\n\n"+
"<b>Изменения:</b>\n%s%s\n\n"+
"<b>Время проверки:</b> %s",
sid,
domain,
clientTitle,
wafProviderClient,
strings.Join(changeDetails, "\n"),
ptafWarning,
now.Format("2006-01-02 15:04:05")+" UTC+3",
)
if err := sendAlert(message); err != nil {
log.Printf("Failed to send Telegram alert for SID %d: %v", sid, err)
} else {
log.Printf("Alert sent for SID %d", sid)
}
}
// Отправляем отдельный алерт для изменений WAF настроек
if hasWAFChanges {
now := time.Now().In(time.FixedZone("UTC+3", 3*60*60))
wafMessage := fmt.Sprintf(
"⚪ <b>Изменения настроек инстанса на стороне SP</b>\n\n"+
"<b>SID:</b> %d <b>Домен:</b> %s <b>TENANT:</b> %s\n\n"+
"<b>Изменения:</b>\n%s\n\n"+
"<b>Время проверки:</b> %s",
sid,
domain,
clientTitle,
strings.Join(wafChangeDetails, "\n"),
now.Format("2006-01-02 15:04:05")+" UTC+3",
)
if err := sendAlert(wafMessage); err != nil {
log.Printf("Failed to send WAF settings alert for SID %d: %v", sid, err)
} else {
log.Printf("WAF settings alert sent for SID %d", sid)
}
}
return nil
}
func checkSIDAvailability(sid int64) error {
var whois whoisResponse
if err := apiGet(whoisURL+fmt.Sprint(sid)+"/global", &whois); err != nil {
return fmt.Errorf("whois: %w", err)
}
var originsResp apiList[originItem]
if err := apiGet(originURL+fmt.Sprint(sid), &originsResp); err != nil {
return fmt.Errorf("origins: %w", err)
}
var aliasResp apiList[aliasItem]
if err := apiGet(aliasURL+fmt.Sprint(sid), &aliasResp); err != nil {
return fmt.Errorf("aliases: %w", err)
}
return nil
}
func sendAPIErrorsAlert(errors map[int64]string) {
if len(errors) == 0 {
return
}
// Сортируем SID для стабильного вывода
sids := make([]int64, 0, len(errors))
for sid := range errors {
sids = append(sids, sid)
}
sort.Slice(sids, func(i, j int) bool { return sids[i] < sids[j] })
now := time.Now().In(time.FixedZone("UTC+3", 3*60*60))
var msg strings.Builder
msg.WriteString(fmt.Sprintf("🔴 <b>Ошибки запросов к API ServicePipe (%d)</b>\n\n", len(errors)))
for _, sid := range sids {
msg.WriteString(fmt.Sprintf("<b>SID %d:</b> %s\n", sid, errors[sid]))
}
msg.WriteString(fmt.Sprintf("\n<b>Время проверки:</b> %s", now.Format("2006-01-02 15:04:05")+" UTC+3"))
if err := sendAlert(msg.String()); err != nil {
log.Printf("[API Errors Alert] Failed to send alert: %v", err)
}
}
func wafVendorName(vendor string) string {
switch strings.ToLower(vendor) {
case "ptaf":
return "PT AF"
case "sw":
return "SW"
case "wmx":
return "WMX"
default:
return "WAF"
}
}
func wafVendorWarning(vendor string) string {
switch strings.ToLower(vendor) {
case "ptaf":
return "<b>⚠️ Проверьте, не нужно ли внести изменения в PT AF!</b>"
case "sw":
return "<b>⚠️ Проверьте, не нужно ли внести изменения в SW!</b>"
case "wmx":
return "<b>⚠️ Проверьте, не нужно ли внести изменения в WMX!</b>"
default:
return "<b>⚠️ Проверьте, не нужно ли внести изменения в WAF!</b>"
}
}
func cleanupRemovedSIDs(db *sql.DB, activeSIDs []int64) error {
if len(activeSIDs) == 0 {
log.Println("[Cleanup] No active SIDs provided, skipping cleanup to avoid accidental data loss")
return nil
}
rows, err := db.Query(`
SELECT s.sid, s.domain_name, a.waf_vendor
FROM sp_info s
LEFT JOIN apps_settings a ON a.l7resourceid = s.sid
`)
if err != nil {
return fmt.Errorf("failed to query sp_info: %w", err)
}
defer rows.Close()
activeSet := make(map[int64]bool, len(activeSIDs))
for _, sid := range activeSIDs {
activeSet[sid] = true
}
type staleRecord struct {
sid int64
domain string
wafVendor string
}
var stale []staleRecord
for rows.Next() {
var sid int64
var domain string
var wafVendor sql.NullString
if err := rows.Scan(&sid, &domain, &wafVendor); err != nil {
return err
}
if !activeSet[sid] {
vendor := ""
if wafVendor.Valid {
vendor = wafVendor.String
}
stale = append(stale, staleRecord{sid, domain, vendor})
}
}
if len(stale) == 0 {
log.Println("[Cleanup] No stale SIDs found in sp_info")
return nil
}
log.Printf("[Cleanup] Found %d stale SID(s) to remove", len(stale))
for _, rec := range stale {
log.Printf("[Cleanup] Removing SID %d (%s) from sp_info", rec.sid, rec.domain)
_, err := db.Exec(`DELETE FROM sp_info WHERE sid = $1`, rec.sid)
if err != nil {
log.Printf("[Cleanup] Failed to delete SID %d: %v", rec.sid, err)
continue
}
wafWarning := wafVendorWarning(rec.wafVendor)
message := fmt.Sprintf(
"<b>🗑 Ресурс удалён из мониторинга</b>\n\n"+
"<b>SID:</b> %d\n"+
"<b>Домен:</b> %s\n\n"+
"Ресурс отсутствует в <code>apps_settings</code> и был удалён из <code>sp_info</code>.\n"+
"%s",
rec.sid,
rec.domain,
wafWarning,
)
if err := sendAlert(message); err != nil {
log.Printf("[Cleanup] Failed to send Telegram alert for SID %d: %v", rec.sid, err)
} else {
log.Printf("[Cleanup] Alert sent for removed SID %d", rec.sid)
}
}
return nil
}
func checkDuplicateDomains(db *sql.DB) error {
rows, err := db.Query(`SELECT sid, domain_name, aliases FROM sp_info`)
if err != nil {
return fmt.Errorf("failed to query sp_info: %w", err)
}
defer rows.Close()
// Единая карта: домен/алиас -> список SID где встречается
valueToSIDs := make(map[string][]int64)
for rows.Next() {
var sid int64
var domain string
var aliasesJSON []byte
if err := rows.Scan(&sid, &domain, &aliasesJSON); err != nil {
return err
}
if domain != "" {
valueToSIDs[domain] = append(valueToSIDs[domain], sid)
}
var aliases []string
if err := json.Unmarshal(aliasesJSON, &aliases); err == nil {
for _, alias := range aliases {
// Добавляем только если этот SID ещё не учтён для данного значения
alreadyAdded := false
for _, s := range valueToSIDs[alias] {
if s == sid {
alreadyAdded = true
break
}
}
if !alreadyAdded {
valueToSIDs[alias] = append(valueToSIDs[alias], sid)
}
}
}
}
// Собираем дубли — значения встречающиеся у более чем одного SID
type duplicate struct {
value string
sids []int64
}
var dups []duplicate
for value, sids := range valueToSIDs {
if len(sids) > 1 {
dups = append(dups, duplicate{value, sids})
}
}
if len(dups) == 0 {
log.Println("[Duplicate Check] No duplicates found")
return nil
}
// Сортируем для стабильного вывода
sort.Slice(dups, func(i, j int) bool {
return dups[i].value < dups[j].value
})
log.Printf("[Duplicate Check] Found %d duplicate(s)", len(dups))
for _, d := range dups {
var sidStrs []string
for _, sid := range d.sids {
sidStrs = append(sidStrs, fmt.Sprintf("%d", sid))
}
log.Printf("[Duplicate Check] %s → SID: %s", d.value, strings.Join(sidStrs, ", "))
}
var msg strings.Builder
msg.WriteString("🔴 <b>Обнаружены дублирующиеся домены/алиасы</b>\n\n")
for _, d := range dups {
var sidStrs []string
for _, sid := range d.sids {
sidStrs = append(sidStrs, fmt.Sprintf("%d", sid))
}
msg.WriteString(fmt.Sprintf(" %s → SID: %s\n", d.value, strings.Join(sidStrs, ", ")))
}
now := time.Now().In(time.FixedZone("UTC+3", 3*60*60))
msg.WriteString(fmt.Sprintf("\n<b>Время проверки:</b> %s", now.Format("2006-01-02 15:04:05")+" UTC+3"))
if err := sendAlert(msg.String()); err != nil {
log.Printf("[Duplicate Check] Failed to send alert: %v", err)
}
return nil
}