magos/api.go
2026-04-01 12:52:09 +03:00

383 lines
12 KiB
Go
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"
"io"
"log"
"net/http"
"strings"
"time"
)
// apiClient — переиспользуемый HTTP-клиент для всех API запросов.
// http.Client безопасен для конкурентного использования и переиспользует соединения.
var apiClient = &http.Client{
Timeout: 120 * time.Second,
}
// makeAPIRequest выполняет GET запрос к API с авторизацией и повторными попытками
func makeAPIRequest(url, token string) ([]byte, error) {
maxRetries := 3
var lastErr error
for attempt := 1; attempt <= maxRetries; attempt++ {
body, err := doAPIRequest(url, token)
if err == nil {
if attempt > 1 {
log.Printf(" ✓ Запрос успешен после %d попыток", attempt)
}
return body, nil
}
lastErr = err
if attempt < maxRetries {
waitTime := time.Duration(attempt*attempt) * time.Second // 1s, 4s, 9s
log.Printf(" ⚠ Попытка %d/%d не удалась: %v, повтор через %v", attempt, maxRetries, err, waitTime)
time.Sleep(waitTime)
}
}
return nil, fmt.Errorf("запрос не удался после %d попыток: %w", maxRetries, lastErr)
}
// doAPIRequest выполняет одиночный API запрос (без retry)
func doAPIRequest(url, token string) ([]byte, error) {
req, err := http.NewRequest("GET", url, nil)
if err != nil {
return nil, fmt.Errorf("ошибка создания запроса: %w", err)
}
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("Content-Type", "application/json")
resp, err := apiClient.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
return nil, fmt.Errorf("ошибка чтения ответа: %w", err)
}
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("API вернул код %d, URL: %s, ответ: %s", resp.StatusCode, url, string(body))
}
return body, nil
}
// getResourceDataFromCache пытается получить данные из таблицы sp_info
func getResourceDataFromCache(db *sql.DB, config Config, l7ResourceID int) (serverName string, aliases []string, origins []OriginItem, found bool, err error) {
query := `
SELECT
domain_name,
aliases,
origins,
updated_at,
EXTRACT(EPOCH FROM (NOW() - updated_at)) / 60 AS age_minutes
FROM sp_info
WHERE sid = $1
`
var domainName sql.NullString
var aliasesJSON sql.NullString
var originsJSON sql.NullString
var updatedAt sql.NullTime
var ageMinutes sql.NullFloat64
err = db.QueryRow(query, l7ResourceID).Scan(&domainName, &aliasesJSON, &originsJSON, &updatedAt, &ageMinutes)
if err != nil {
if err == sql.ErrNoRows {
return "", nil, nil, false, nil // Данных нет в кеше
}
return "", nil, nil, false, fmt.Errorf("ошибка чтения из sp_info: %w", err)
}
// Проверяем свежесть данных (не старее 10 минут)
if !ageMinutes.Valid {
log.Printf(" ⚠ updated_at не заполнен в кеше, обращаюсь к API")
return "", nil, nil, false, nil
}
age := ageMinutes.Float64
if age < 0 {
age = 0
}
if age > 10 {
log.Printf(" ⚠ Данные в кеше устарели (возраст: %.1f мин), обращаюсь к API", age)
return "", nil, nil, false, nil
}
log.Printf(" ✓ Данные из кеша актуальны (возраст: %.1f мин)", age)
if domainName.Valid {
serverName = domainName.String
}
aliases, err = parseAliasesJSON(aliasesJSON)
if err != nil {
return "", nil, nil, false, fmt.Errorf("ошибка парсинга aliases из sp_info: %w", err)
}
origins, err = parseOriginsJSON(originsJSON, config.WAFNetworks)
if err != nil {
return "", nil, nil, false, fmt.Errorf("ошибка парсинга origins из sp_info: %w", err)
}
return serverName, aliases, origins, true, nil
}
// getResourceDataFromManualInfo получает данные из таблицы manual_info (ручное управление)
// В отличие от sp_info, не проверяет timestamp - данные всегда актуальны
func getResourceDataFromManualInfo(db *sql.DB, config Config, l7ResourceID int) (serverName string, aliases []string, origins []OriginItem, err error) {
query := `
SELECT
domain_name,
aliases,
origins
FROM manual_info
WHERE sid = $1
`
var domainName sql.NullString
var aliasesJSON sql.NullString
var originsJSON sql.NullString
err = db.QueryRow(query, l7ResourceID).Scan(&domainName, &aliasesJSON, &originsJSON)
if err != nil {
if err == sql.ErrNoRows {
return "", nil, nil, fmt.Errorf("данные не найдены в manual_info для sid=%d", l7ResourceID)
}
return "", nil, nil, fmt.Errorf("ошибка чтения из manual_info: %w", err)
}
if domainName.Valid {
serverName = domainName.String
}
aliases, err = parseAliasesJSON(aliasesJSON)
if err != nil {
return "", nil, nil, fmt.Errorf("ошибка парсинга aliases из manual_info: %w", err)
}
origins, err = parseOriginsJSON(originsJSON, config.WAFNetworks)
if err != nil {
return "", nil, nil, fmt.Errorf("ошибка парсинга origins из manual_info: %w", err)
}
return serverName, aliases, origins, nil
}
// parseAliasesJSON парсит JSON массив алиасов из sql.NullString
func parseAliasesJSON(aliasesJSON sql.NullString) ([]string, error) {
if !aliasesJSON.Valid || aliasesJSON.String == "" || aliasesJSON.String == "null" {
return nil, nil
}
var aliases []string
if err := json.Unmarshal([]byte(aliasesJSON.String), &aliases); err != nil {
return nil, err
}
return aliases, nil
}
// parseOriginsJSON парсит JSON массив origins из sql.NullString, фильтруя WAF IP
func parseOriginsJSON(originsJSON sql.NullString, wafNetworks []string) ([]OriginItem, error) {
if !originsJSON.Valid || originsJSON.String == "" {
return nil, nil
}
var originsData []struct {
IP string `json:"ip"`
Mode string `json:"mode"`
Weight int `json:"weight"`
}
if err := json.Unmarshal([]byte(originsJSON.String), &originsData); err != nil {
return nil, err
}
var origins []OriginItem
for _, o := range originsData {
if !isIPInWAFNetworks(o.IP, wafNetworks) {
origins = append(origins, OriginItem{
IP: o.IP,
Weight: o.Weight,
Mode: o.Mode,
})
}
}
return origins, nil
}
// getServerName получает имя сервера (l7ResourceName) через API
func getServerName(config Config, l7ResourceID int) (string, error) {
url := fmt.Sprintf("https://api.servicepipe.ru/api/v1/l7/resource/%d/global", l7ResourceID)
body, err := makeAPIRequest(url, config.APIToken)
if err != nil {
return "", err
}
var response ResourceResponse
if err := json.Unmarshal(body, &response); err != nil {
return "", fmt.Errorf("ошибка парсинга JSON: %w", err)
}
return response.Data.Result.L7ResourceName, nil
}
// getAliases получает список aliases через API
func getAliases(config Config, l7ResourceID int) ([]string, error) {
url := fmt.Sprintf("https://api.servicepipe.ru/api/v1/l7/alias/global?limit=1000&l7ResourceId=%d", l7ResourceID)
body, err := makeAPIRequest(url, config.APIToken)
if err != nil {
return nil, err
}
var response AliasResponse
if err := json.Unmarshal(body, &response); err != nil {
return nil, fmt.Errorf("ошибка парсинга JSON: %w", err)
}
var aliases []string
for _, item := range response.Data.Result.Items {
aliases = append(aliases, item.Domain)
}
return aliases, nil
}
// getOrigins получает список origins через API, исключая WAF IP
func getOrigins(config Config, l7ResourceID int) ([]OriginItem, error) {
url := fmt.Sprintf("https://api.servicepipe.ru/api/v1/l7/origin/global?limit=1000&l7ResourceId=%d", l7ResourceID)
body, err := makeAPIRequest(url, config.APIToken)
if err != nil {
return nil, err
}
var response OriginResponse
if err := json.Unmarshal(body, &response); err != nil {
return nil, fmt.Errorf("ошибка парсинга JSON: %w", err)
}
var origins []OriginItem
for _, item := range response.Data.Result.Items {
if !isIPInWAFNetworks(item.IP, config.WAFNetworks) {
origins = append(origins, OriginItem{
IP: item.IP,
Weight: item.Weight,
Mode: item.Mode,
})
}
}
return origins, nil
}
// getInputHTTPPorts возвращает список HTTP INPUT портов из AppsSettings или дефолтные
func getInputHTTPPorts(settings AppsSettings) []int {
if len(settings.CustomInputHTTPPorts) > 0 {
return int64ArrayToIntSlice(settings.CustomInputHTTPPorts)
}
return []int{80}
}
// getInputHTTPSPorts возвращает список HTTPS INPUT портов из AppsSettings или дефолтные
func getInputHTTPSPorts(settings AppsSettings) []int {
if len(settings.CustomInputHTTPSPorts) > 0 {
return int64ArrayToIntSlice(settings.CustomInputHTTPSPorts)
}
return []int{443}
}
// getOutputHTTPPorts возвращает список HTTP OUTPUT портов из AppsSettings или дефолтные
func getOutputHTTPPorts(settings AppsSettings) []int {
if len(settings.CustomOutputHTTPPorts) > 0 {
return int64ArrayToIntSlice(settings.CustomOutputHTTPPorts)
}
return []int{80}
}
// getOutputHTTPSPorts возвращает список HTTPS OUTPUT портов из AppsSettings или дефолтные
func getOutputHTTPSPorts(settings AppsSettings) []int {
if len(settings.CustomOutputHTTPSPorts) > 0 {
return int64ArrayToIntSlice(settings.CustomOutputHTTPSPorts)
}
return []int{443}
}
// int64ArrayToIntSlice конвертирует pq.Int64Array в []int
func int64ArrayToIntSlice(arr []int64) []int {
result := make([]int, len(arr))
for i, v := range arr {
result[i] = int(v)
}
return result
}
// writeCustomUpstreamOrDefault записывает кастомные или дефолтные upstream-директивы
func writeCustomUpstreamOrDefault(w stringWriter, customUpstream sql.NullString, sni sql.NullInt64) {
hasCustomUpstream := customUpstream.Valid && customUpstream.String != ""
if hasCustomUpstream {
w.WriteString("\n # Custom upstream directives (replace defaults)\n")
lines := strings.Split(strings.TrimSpace(customUpstream.String), "\n")
for _, line := range lines {
trimmed := strings.TrimSpace(line)
if trimmed != "" {
w.WriteString(" " + trimmed + "\n")
}
}
} else {
// Если нет custom - используем дефолтные keepalive директивы (только если sni != 1)
if !sni.Valid || sni.Int64 != 1 {
w.WriteString("\n # Default upstream directives\n")
w.WriteString(" keepalive 60;\n")
w.WriteString(" keepalive_timeout 70s;\n")
}
}
}
// stringWriter — интерфейс для записи строк (совпадает с *os.File и *strings.Builder)
type stringWriter interface {
WriteString(s string) (int, error)
}
// getProxyNextUpstream формирует строку proxy_next_upstream из кодов в settings.
// Дефолт: error timeout invalid_header http_500 http_502 http_503 http_504
func getProxyNextUpstream(settings AppsSettings) string {
codes := "500,502,503,504"
if settings.ProxyNextUpstreamCodes.Valid && strings.TrimSpace(settings.ProxyNextUpstreamCodes.String) != "" {
codes = strings.TrimSpace(settings.ProxyNextUpstreamCodes.String)
}
base := "error timeout invalid_header"
parts := strings.Split(codes, ",")
for _, part := range parts {
code := strings.TrimSpace(part)
if code != "" {
base += " http_" + code
}
}
return base
}
// getMaxFails возвращает max_fails из AppsSettings или дефолтное значение 5
func getMaxFails(settings AppsSettings) int {
if settings.MaxFails.Valid && settings.MaxFails.Int64 >= 0 {
return int(settings.MaxFails.Int64)
}
return 5 // дефолт
}
// getFailTimeout возвращает fail_timeout из AppsSettings или дефолтное значение 15
func getFailTimeout(settings AppsSettings) int {
if settings.FailTimeout.Valid && settings.FailTimeout.Int64 >= 0 {
return int(settings.FailTimeout.Int64)
}
return 15 // дефолт
}