317 lines
12 KiB
Go
317 lines
12 KiB
Go
package api
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"io"
|
||
"net/http"
|
||
"os"
|
||
"strings"
|
||
"sync"
|
||
"time"
|
||
|
||
"fission-console/internal/cloud"
|
||
"fission-console/internal/fission"
|
||
"fission-console/ui"
|
||
|
||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||
"k8s.io/client-go/dynamic"
|
||
)
|
||
|
||
// deckAPIs — карта окружений Deck API.
|
||
// Ключ используется в X-Auth-Env заголовке для выбора нужного сервера.
|
||
var deckAPIs = map[string]string{
|
||
"prod": "https://deck-api.ngcloud.ru/api/v1",
|
||
"dev": "https://deck-api-dev.ngcloud.ru/api/v1",
|
||
"test": "https://deck-api-test.ngcloud.ru/api/v1",
|
||
}
|
||
|
||
// defaultSATokenPath — путь к service account токену внутри pod-а.
|
||
// Используется для авторизации запросов от console к Fission router.
|
||
const defaultSATokenPath = "/var/run/secrets/kubernetes.io/serviceaccount/token"
|
||
|
||
// Server — основная структура HTTP сервера.
|
||
// Содержит все зависимости: kubernetes client, конфиги, кэши токенов.
|
||
type Server struct {
|
||
dyn dynamic.Interface
|
||
ns string // системный namespace (fallback, обычно "fission")
|
||
routerURL string
|
||
http *http.Client
|
||
|
||
saTokenPath string
|
||
invokeTimeout time.Duration
|
||
testMode bool // FISSION_TEST_MODE=true — пропускает deck auth
|
||
|
||
authUser string
|
||
authPass string
|
||
|
||
// --- ai/ask feature (удалить блок целиком чтобы выкосить) ---
|
||
llmURL string // FISSION_LLM_URL
|
||
llmKey string // FISSION_LLM_KEY
|
||
// --- end ai/ask feature ---
|
||
|
||
// tokenMu защищает кэш JWT токена для аутентификации в Fission router.
|
||
tokenMu sync.Mutex
|
||
cachedJWT string
|
||
tokenExpAt time.Time
|
||
|
||
// tokenCache кэширует результаты валидации Deck токенов.
|
||
// Ключ: "env:token", значение: time.Time — когда кэш истекает.
|
||
tokenCache sync.Map
|
||
|
||
// nsManager управляет жизненным циклом пользовательских namespace-ов.
|
||
nsManager *cloud.NSManager
|
||
}
|
||
|
||
// Config содержит все параметры для создания Server.
|
||
type Config struct {
|
||
Dyn dynamic.Interface
|
||
Namespace string
|
||
RouterURL string
|
||
HTTPTimeout time.Duration
|
||
InvokeTimeout time.Duration
|
||
SATokenPath string
|
||
AuthUser string
|
||
AuthPass string
|
||
TestMode bool
|
||
LLMUrl string
|
||
LLMKey string
|
||
}
|
||
|
||
// NewServer создаёт и настраивает HTTP Server со всеми зависимостями.
|
||
func NewServer(cfg Config) *Server {
|
||
return &Server{
|
||
dyn: cfg.Dyn,
|
||
ns: cfg.Namespace,
|
||
routerURL: cfg.RouterURL,
|
||
http: &http.Client{Timeout: cfg.HTTPTimeout},
|
||
saTokenPath: cfg.SATokenPath,
|
||
invokeTimeout: cfg.InvokeTimeout,
|
||
authUser: cfg.AuthUser,
|
||
authPass: cfg.AuthPass,
|
||
testMode: cfg.TestMode,
|
||
llmURL: cfg.LLMUrl,
|
||
llmKey: cfg.LLMKey,
|
||
nsManager: cloud.NewNSManager(cfg.Dyn),
|
||
}
|
||
}
|
||
|
||
// RegisterRoutes регистрирует все HTTP маршруты на mux.
|
||
// Разделено на публичные (без auth) и приватные (с authMiddleware).
|
||
func (s *Server) RegisterRoutes(mux *http.ServeMux) {
|
||
// Публичные: health checks, UI, auth endpoint
|
||
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
|
||
if r.URL.Path != "/" {
|
||
http.NotFound(w, r)
|
||
return
|
||
}
|
||
http.Redirect(w, r, "/console/", http.StatusTemporaryRedirect)
|
||
})
|
||
mux.HandleFunc("/health", func(w http.ResponseWriter, _ *http.Request) {
|
||
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
|
||
_, _ = w.Write([]byte("ok\n"))
|
||
})
|
||
mux.HandleFunc("/console/health", func(w http.ResponseWriter, _ *http.Request) {
|
||
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
|
||
_, _ = w.Write([]byte("ok\n"))
|
||
})
|
||
|
||
// UI: embedded single-page app
|
||
uiHandler := http.StripPrefix("/console", ui.Handler())
|
||
mux.Handle("/console", uiHandler)
|
||
mux.Handle("/console/", uiHandler)
|
||
cronHandler := ui.CronHandler()
|
||
mux.Handle("/cron", cronHandler)
|
||
mux.Handle("/cron/", cronHandler)
|
||
mux.Handle("/console/cron", cronHandler)
|
||
mux.Handle("/console/cron/", cronHandler)
|
||
|
||
// Auth: не требует токена — сам проверяет и возвращает namespace
|
||
mux.HandleFunc("/console/api/auth", s.handleAuth)
|
||
|
||
// Приватные: требуют X-Auth-Token / X-Test-Sub
|
||
auth := s.authMiddleware
|
||
|
||
// Legacy /api/* (без /console префикса) — для обратной совместимости
|
||
mux.HandleFunc("/api/environments", s.handleList(fission.EnvironmentGVR))
|
||
mux.HandleFunc("/api/packages", s.handleList(fission.PackageGVR))
|
||
mux.HandleFunc("/api/functions", s.handleFunctionsRoot)
|
||
mux.HandleFunc("/api/functions/", s.handleFunctionsAction)
|
||
mux.HandleFunc("/api/httptriggers", s.handleList(fission.HTTPTrigGVR))
|
||
mux.HandleFunc("/api/timetriggers", s.handleTimeTriggersRoot)
|
||
mux.HandleFunc("/api/timetriggers/", s.handleTimeTriggersAction)
|
||
mux.HandleFunc("/cron/api/metrics", s.handleCronMetrics)
|
||
mux.HandleFunc("/console/cron/api/metrics", s.handleCronMetrics)
|
||
|
||
// Основные /console/api/* маршруты
|
||
mux.HandleFunc("/console/api/environments", auth(s.handleList(fission.EnvironmentGVR)))
|
||
mux.HandleFunc("/console/api/packages", auth(s.handleList(fission.PackageGVR)))
|
||
mux.HandleFunc("/console/api/functions", auth(s.handleFunctionsRoot))
|
||
mux.HandleFunc("/console/api/functions/", auth(s.handleFunctionsAction))
|
||
mux.HandleFunc("/fission-function", s.handleFissionFunctionGateway)
|
||
mux.HandleFunc("/fission-function/", s.handleFissionFunctionGateway)
|
||
mux.HandleFunc("/fn/", auth(s.handleInvokeRoute))
|
||
mux.HandleFunc("/console/api/httptriggers", auth(s.handleList(fission.HTTPTrigGVR)))
|
||
mux.HandleFunc("/console/api/timetriggers", auth(s.handleTimeTriggersRoot))
|
||
mux.HandleFunc("/console/api/timetriggers/", auth(s.handleTimeTriggersAction))
|
||
mux.HandleFunc("/console/api/ns/status", auth(s.handleNSStatus))
|
||
mux.HandleFunc("/console/api/ns/debug", auth(s.handleNSDebug))
|
||
mux.HandleFunc("/console/api/ai/check", auth(s.handleAICheck))
|
||
// --- ai/ask feature (удалить строку чтобы выкосить роут) ---
|
||
mux.HandleFunc("/console/api/ai/ask", auth(s.handleAIAsk))
|
||
// --- end ai/ask feature ---
|
||
}
|
||
|
||
// Handler возвращает http.Handler со всеми middleware: CORS, security headers, логирование.
|
||
func (s *Server) Handler() http.Handler {
|
||
mux := http.NewServeMux()
|
||
s.RegisterRoutes(mux)
|
||
return withSecurityHeaders(withCORS(logRequests(mux)))
|
||
}
|
||
|
||
// NSManager возвращает указатель на NSManager для запуска фоновых горутин из main.
|
||
func (s *Server) NSManager() *cloud.NSManager {
|
||
return s.nsManager
|
||
}
|
||
|
||
// handleList возвращает handler который читает список ресурсов из namespace пользователя.
|
||
// Общий для environments, packages, functions, httptriggers, timetriggers.
|
||
func (s *Server) handleList(gvr schema.GroupVersionResource) http.HandlerFunc {
|
||
return func(w http.ResponseWriter, r *http.Request) {
|
||
if r.Method != http.MethodGet {
|
||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||
return
|
||
}
|
||
|
||
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
|
||
defer cancel()
|
||
|
||
list, err := s.dyn.Resource(gvr).Namespace(s.userNS(r)).List(ctx, metav1.ListOptions{})
|
||
if err != nil {
|
||
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("list %s: %v", gvr.Resource, err))
|
||
return
|
||
}
|
||
|
||
writeJSON(w, http.StatusOK, list.Items)
|
||
}
|
||
}
|
||
|
||
// writeJSON сериализует список unstructured объектов в JSON response.
|
||
func writeJSON(w http.ResponseWriter, status int, data []unstructured.Unstructured) {
|
||
w.Header().Set("Content-Type", "application/json; charset=utf-8")
|
||
w.WriteHeader(status)
|
||
_ = json.NewEncoder(w).Encode(data)
|
||
}
|
||
|
||
// writeAnyJSON сериализует любое значение в JSON response.
|
||
func writeAnyJSON(w http.ResponseWriter, status int, data any) {
|
||
w.Header().Set("Content-Type", "application/json; charset=utf-8")
|
||
w.WriteHeader(status)
|
||
_ = json.NewEncoder(w).Encode(data)
|
||
}
|
||
|
||
// writeJSONError возвращает JSON с ключом "error" и заданным HTTP статусом.
|
||
func writeJSONError(w http.ResponseWriter, status int, msg string) {
|
||
w.Header().Set("Content-Type", "application/json; charset=utf-8")
|
||
w.WriteHeader(status)
|
||
_ = json.NewEncoder(w).Encode(map[string]any{"error": msg})
|
||
}
|
||
|
||
// readSAToken читает service account JWT токен из файла.
|
||
// Используется как fallback аутентификация для запросов к Fission router
|
||
// когда authUser/authPass не заданы.
|
||
func (s *Server) readSAToken() string {
|
||
if s.saTokenPath == "" {
|
||
return ""
|
||
}
|
||
data, err := os.ReadFile(s.saTokenPath)
|
||
if err != nil {
|
||
return ""
|
||
}
|
||
return strings.TrimSpace(string(data))
|
||
}
|
||
|
||
// getRouterToken возвращает Bearer токен для запросов к Fission router.
|
||
// Если authUser+authPass заданы — получает JWT через login endpoint (с кэшем).
|
||
// Иначе — использует SA токен из pod filesystem.
|
||
func (s *Server) getRouterToken() string {
|
||
if s.authUser == "" || s.authPass == "" {
|
||
return s.readSAToken()
|
||
}
|
||
|
||
s.tokenMu.Lock()
|
||
defer s.tokenMu.Unlock()
|
||
|
||
if s.cachedJWT != "" && time.Now().Before(s.tokenExpAt) {
|
||
return s.cachedJWT // кэш ещё актуален
|
||
}
|
||
|
||
loginURL := s.routerURL + "/auth/login"
|
||
body, _ := json.Marshal(map[string]string{"username": s.authUser, "password": s.authPass})
|
||
resp, err := s.http.Post(loginURL, "application/json", bytes.NewReader(body))
|
||
if err != nil {
|
||
return s.readSAToken()
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
|
||
return s.readSAToken()
|
||
}
|
||
|
||
var result struct {
|
||
AccessToken string `json:"accesstoken"`
|
||
}
|
||
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil || result.AccessToken == "" {
|
||
return s.readSAToken()
|
||
}
|
||
|
||
s.cachedJWT = result.AccessToken
|
||
// Fission JWT живёт ~120s; кэшируем на 100s чтобы не использовать просроченный
|
||
s.tokenExpAt = time.Now().Add(100 * time.Second)
|
||
return s.cachedJWT
|
||
}
|
||
|
||
// validateDeckToken проверяет токен через Deck API с кэшированием результата на 5 минут.
|
||
func (s *Server) validateDeckToken(token, env string) error {
|
||
cacheKey := env + ":" + token
|
||
if v, ok := s.tokenCache.Load(cacheKey); ok {
|
||
if time.Now().Before(v.(time.Time)) {
|
||
return nil // кэш актуален — токен валиден
|
||
}
|
||
s.tokenCache.Delete(cacheKey)
|
||
}
|
||
|
||
apiBase, ok := deckAPIs[env]
|
||
if !ok {
|
||
return fmt.Errorf("unknown deck env: %s", env)
|
||
}
|
||
|
||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||
defer cancel()
|
||
|
||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, apiBase+"/index.cfm/instances", nil)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
req.Header.Set("Authorization", "Bearer "+token)
|
||
|
||
resp, err := s.http.Do(req)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer resp.Body.Close()
|
||
_, _ = io.ReadAll(resp.Body)
|
||
|
||
if resp.StatusCode == http.StatusUnauthorized {
|
||
return fmt.Errorf("invalid token")
|
||
}
|
||
|
||
// Кэшируем успешный результат
|
||
s.tokenCache.Store(cacheKey, time.Now().Add(5*time.Minute))
|
||
return nil
|
||
}
|