355 lines
11 KiB
Go
355 lines
11 KiB
Go
package configserver
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/tls"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"mime/multipart"
|
|
"net/http"
|
|
"net/http/cookiejar"
|
|
"net/url"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// maxResponseBytes — жёсткий предохранитель на размер тела ответа (защита от
|
|
// OOM). Доменные ответы (compose/env/schema) заметно меньше.
|
|
const maxResponseBytes = 16 << 20 // 16 MiB
|
|
|
|
// Session — авторизованная HTTP-сессия к одному config_server.
|
|
//
|
|
// Авторизация устроена как Flask-сессия: login кладёт подписанную cookie, она
|
|
// хранится в cookiejar и автоматически прикладывается к запросам. Сессия
|
|
// потокобезопасна: login выполняется один раз (single-flight), при 401 cookie
|
|
// пересоздаётся и запрос повторяется ровно один раз.
|
|
type Session struct {
|
|
baseURL string
|
|
username string
|
|
password string
|
|
label string // alias (+ "(ro)") для сообщений об ошибках
|
|
|
|
client *http.Client
|
|
|
|
mu sync.Mutex
|
|
loggedIn bool
|
|
}
|
|
|
|
// response — низкоуровневый ответ (тело уже прочитано и ограничено).
|
|
type response struct {
|
|
status int
|
|
body []byte
|
|
ctype string
|
|
}
|
|
|
|
// NewSession создаёт сессию для сервера. readonly=true выбирает отдельную
|
|
// read-only учётную запись, если она задана (иначе — основная).
|
|
func NewSession(srv *Server, readonly bool) (*Session, error) {
|
|
user, pass := srv.Username, srv.Password
|
|
label := srv.Alias
|
|
if readonly && srv.ReadonlyUsername != "" {
|
|
user, pass = srv.ReadonlyUsername, srv.ReadonlyPassword
|
|
label += "(ro)"
|
|
}
|
|
|
|
jar, err := cookiejar.New(nil)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("cookie jar: %w", err)
|
|
}
|
|
// #nosec G402 — InsecureSkipVerify включается осознанно оператором для
|
|
// стендов с самоподписанными сертификатами (флаг в конфиге).
|
|
transport := &http.Transport{
|
|
TLSClientConfig: &tls.Config{InsecureSkipVerify: srv.InsecureSkipVerify},
|
|
}
|
|
return &Session{
|
|
baseURL: srv.BaseURL,
|
|
username: user,
|
|
password: pass,
|
|
label: label,
|
|
client: &http.Client{Jar: jar, Transport: transport},
|
|
}, nil
|
|
}
|
|
|
|
// BaseURL возвращает адрес сервера сессии.
|
|
func (s *Session) BaseURL() string { return s.baseURL }
|
|
|
|
// ensureLogin выполняет login, если сессия ещё не авторизована.
|
|
func (s *Session) ensureLogin(ctx context.Context) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if s.loggedIn {
|
|
return nil
|
|
}
|
|
return s.loginLocked(ctx)
|
|
}
|
|
|
|
// loginLocked — login под уже взятым мьютексом (single-flight).
|
|
func (s *Session) loginLocked(ctx context.Context) error {
|
|
if s.username == "" || s.password == "" {
|
|
return &APIError{Status: http.StatusUnauthorized,
|
|
Message: fmt.Sprintf("для сервера %s не заданы username/password", s.label)}
|
|
}
|
|
resp, err := s.once(ctx, http.MethodPost, "/api/login", nil,
|
|
map[string]string{"username": s.username, "password": s.password})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := classify(resp); err != nil {
|
|
return err
|
|
}
|
|
s.loggedIn = true
|
|
return nil
|
|
}
|
|
|
|
// invalidate сбрасывает признак авторизации (cookie могла истечь).
|
|
func (s *Session) invalidate() {
|
|
s.mu.Lock()
|
|
s.loggedIn = false
|
|
s.mu.Unlock()
|
|
}
|
|
|
|
// getJSON выполняет авторизованный GET и разбирает JSON в out.
|
|
func (s *Session) getJSON(ctx context.Context, path string, query url.Values, out any) error {
|
|
resp, err := s.do(ctx, http.MethodGet, path, query, nil, true)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := classify(resp); err != nil {
|
|
return err
|
|
}
|
|
return decodeJSON(resp, out)
|
|
}
|
|
|
|
// postJSON выполняет авторизованный POST с JSON-телом и разбирает ответ.
|
|
func (s *Session) postJSON(ctx context.Context, path string, body, out any) error {
|
|
resp, err := s.do(ctx, http.MethodPost, path, nil, body, true)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := classify(resp); err != nil {
|
|
return err
|
|
}
|
|
return decodeJSON(resp, out)
|
|
}
|
|
|
|
// getPublic выполняет GET без авторизации (например, /health).
|
|
func (s *Session) getPublic(ctx context.Context, path string) (*response, error) {
|
|
return s.do(ctx, http.MethodGet, path, nil, nil, false)
|
|
}
|
|
|
|
// do выполняет запрос с автоматическим повтором: один релогin при 401 и
|
|
// (только для GET) один повтор при сетевом сбое или 5xx. Мутации НЕ
|
|
// повторяются, чтобы не выполнить действие дважды.
|
|
func (s *Session) do(ctx context.Context, method, path string, query url.Values, body any, authenticated bool) (*response, error) {
|
|
if authenticated {
|
|
if err := s.ensureLogin(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
attempts := 1
|
|
if method == http.MethodGet {
|
|
attempts = 2
|
|
}
|
|
|
|
var resp *response
|
|
for attempt := 1; attempt <= attempts; attempt++ {
|
|
r, err := s.once(ctx, method, path, query, body)
|
|
if err != nil {
|
|
if attempt < attempts && ctx.Err() == nil {
|
|
if err := sleepBackoff(ctx, attempt); err != nil {
|
|
return nil, err
|
|
}
|
|
continue
|
|
}
|
|
return nil, err
|
|
}
|
|
resp = r
|
|
|
|
if authenticated && resp.status == http.StatusUnauthorized {
|
|
// Cookie-сессия истекла: перелогиниваемся и повторяем ровно один раз.
|
|
s.invalidate()
|
|
if lerr := s.ensureLogin(ctx); lerr != nil {
|
|
return resp, nil
|
|
}
|
|
if r2, err2 := s.once(ctx, method, path, query, body); err2 == nil {
|
|
return r2, nil
|
|
}
|
|
return resp, nil
|
|
}
|
|
if attempt < attempts && resp.status >= 500 {
|
|
if err := sleepBackoff(ctx, attempt); err != nil {
|
|
return nil, err
|
|
}
|
|
continue
|
|
}
|
|
return resp, nil
|
|
}
|
|
return resp, nil
|
|
}
|
|
|
|
// once выполняет ровно один HTTP-запрос.
|
|
func (s *Session) once(ctx context.Context, method, path string, query url.Values, body any) (*response, error) {
|
|
u := s.baseURL + path
|
|
if len(query) > 0 {
|
|
u += "?" + query.Encode()
|
|
}
|
|
|
|
var reader io.Reader
|
|
if body != nil {
|
|
raw, err := json.Marshal(body)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("кодирование тела запроса: %w", err)
|
|
}
|
|
reader = bytes.NewReader(raw)
|
|
}
|
|
|
|
req, err := http.NewRequestWithContext(ctx, method, u, reader)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("сборка запроса: %w", err)
|
|
}
|
|
req.Header.Set("Accept", "application/json")
|
|
if body != nil {
|
|
req.Header.Set("Content-Type", "application/json")
|
|
}
|
|
|
|
res, err := s.client.Do(req)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("%s %s: %w", method, path, err)
|
|
}
|
|
defer func() { _ = res.Body.Close() }()
|
|
|
|
data, err := io.ReadAll(io.LimitReader(res.Body, maxResponseBytes))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("чтение ответа %s: %w", path, err)
|
|
}
|
|
return &response{status: res.StatusCode, body: data, ctype: res.Header.Get("Content-Type")}, nil
|
|
}
|
|
|
|
// postMultipart загружает файлы как multipart/form-data (update_front). Файлы
|
|
// передаются потоково из памяти вызывающего; имена — относительные пути.
|
|
func (s *Session) postMultipart(ctx context.Context, path string, files []UploadFile) (json.RawMessage, error) {
|
|
if err := s.ensureLogin(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
var buf bytes.Buffer
|
|
writer := multipart.NewWriter(&buf)
|
|
for _, f := range files {
|
|
part, err := writer.CreateFormFile("files", f.Name)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("multipart: %w", err)
|
|
}
|
|
if _, err := part.Write(f.Data); err != nil {
|
|
return nil, fmt.Errorf("multipart: %w", err)
|
|
}
|
|
}
|
|
if err := writer.Close(); err != nil {
|
|
return nil, fmt.Errorf("multipart: %w", err)
|
|
}
|
|
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, s.baseURL+path, &buf)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("сборка запроса: %w", err)
|
|
}
|
|
req.Header.Set("Content-Type", writer.FormDataContentType())
|
|
req.Header.Set("Accept", "application/json")
|
|
|
|
res, err := s.client.Do(req)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("POST %s: %w", path, err)
|
|
}
|
|
defer func() { _ = res.Body.Close() }()
|
|
|
|
data, err := io.ReadAll(io.LimitReader(res.Body, maxResponseBytes))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("чтение ответа %s: %w", path, err)
|
|
}
|
|
resp := &response{status: res.StatusCode, body: data, ctype: res.Header.Get("Content-Type")}
|
|
if err := classify(resp); err != nil {
|
|
return nil, err
|
|
}
|
|
return json.RawMessage(data), nil
|
|
}
|
|
|
|
// UploadFile — файл для multipart-загрузки фронта.
|
|
type UploadFile struct {
|
|
Name string
|
|
Data []byte
|
|
}
|
|
|
|
// classify превращает HTTP-статус ≥400 в доменную APIError с текстом сервера.
|
|
// Текст берём из JSON-поля "error"; HTML/мусорные тела заменяем кратким
|
|
// сообщением, чтобы не засорять контекст модели.
|
|
func classify(resp *response) error {
|
|
if resp.status < 400 {
|
|
return nil
|
|
}
|
|
var payload struct {
|
|
Error string `json:"error"`
|
|
}
|
|
_ = json.Unmarshal(resp.body, &payload)
|
|
msg := strings.TrimSpace(payload.Error)
|
|
if msg == "" {
|
|
msg = defaultHTTPMessage(resp.status)
|
|
}
|
|
return &APIError{Status: resp.status, Message: msg}
|
|
}
|
|
|
|
// defaultHTTPMessage — краткое человекочитаемое описание статуса.
|
|
func defaultHTTPMessage(status int) string {
|
|
switch status {
|
|
case 400:
|
|
return "неверный запрос (400)"
|
|
case 401:
|
|
return "не авторизован (401)"
|
|
case 404:
|
|
return "не найдено (404)"
|
|
case 409:
|
|
return "конфликт (409)"
|
|
case 500:
|
|
return "внутренняя ошибка сервера (500)"
|
|
default:
|
|
return fmt.Sprintf("HTTP %d", status)
|
|
}
|
|
}
|
|
|
|
func decodeJSON(resp *response, out any) error {
|
|
if out == nil || len(resp.body) == 0 {
|
|
return nil
|
|
}
|
|
if err := json.Unmarshal(resp.body, out); err != nil {
|
|
return fmt.Errorf("разбор ответа сервера: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// sleepBackoff делает паузу перед повтором, уважая отмену контекста.
|
|
func sleepBackoff(ctx context.Context, attempt int) error {
|
|
d := time.Duration(200*attempt) * time.Millisecond
|
|
timer := time.NewTimer(d)
|
|
defer timer.Stop()
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case <-timer.C:
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// IsDomainError сообщает, относится ли ошибка к доменным (config_server вернул
|
|
// 4xx/5xx, ошибка валидации модуля или политики вроде отсутствия front_roots),
|
|
// а не к инфраструктурным (сеть). Доменные — recoverable: текст видит модель.
|
|
func IsDomainError(err error) bool {
|
|
var apiErr *APIError
|
|
var valErr *ValidationError
|
|
if errors.As(err, &apiErr) || errors.As(err, &valErr) {
|
|
return true
|
|
}
|
|
return errors.Is(err, ErrNoFrontRoots)
|
|
}
|