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) // Auth: не требует токена — сам проверяет и возвращает namespace mux.HandleFunc("/console/api/auth", s.handleAuth) // Cron: публичные маршруты для метрик cronHandler := ui.CronHandler() mux.Handle("/cron", cronHandler) mux.Handle("/cron/", cronHandler) mux.Handle("/console/cron", cronHandler) mux.Handle("/console/cron/", cronHandler) mux.HandleFunc("/cron/api/metrics", s.handleCronMetrics) mux.HandleFunc("/console/cron/api/metrics", s.handleCronMetrics) // Приватные: требуют 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) // Основные /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 }