From 60646af0c83dbcd8225701442241e8a574d2cf02 Mon Sep 17 00:00:00 2001 From: Naeel Date: Sat, 25 Apr 2026 23:53:20 +0300 Subject: [PATCH] refactor: split main.go into packages (internal/api, internal/fission, internal/runtime, cmd/server) --- console/Dockerfile | 2 +- console/cmd/server/main.go | 107 ++++ console/internal/api/ai_check.go | 219 +++++++ console/internal/api/auth.go | 125 ++++ console/internal/api/doc.go | 2 + console/internal/api/handlers.go | 739 ++++++++++++++++++++++++ console/internal/api/middleware.go | 56 ++ console/internal/api/ns_status.go | 120 ++++ console/internal/api/package.go | 156 +++++ console/internal/api/server.go | 295 ++++++++++ console/internal/fission/client.go | 18 + console/internal/fission/environment.go | 117 ++++ console/internal/fission/namespace.go | 606 +++++++++++++++++++ console/internal/model/types.go | 46 ++ console/internal/runtime/doc.go | 4 + console/internal/runtime/entrypoint.go | 30 + console/internal/runtime/go.go | 43 ++ console/internal/runtime/helpers.go | 27 + console/internal/runtime/nodejs.go | 75 +++ console/internal/runtime/script.go | 13 + 20 files changed, 2799 insertions(+), 1 deletion(-) create mode 100644 console/cmd/server/main.go create mode 100644 console/internal/api/ai_check.go create mode 100644 console/internal/api/auth.go create mode 100644 console/internal/api/doc.go create mode 100644 console/internal/api/handlers.go create mode 100644 console/internal/api/middleware.go create mode 100644 console/internal/api/ns_status.go create mode 100644 console/internal/api/package.go create mode 100644 console/internal/api/server.go create mode 100644 console/internal/fission/client.go create mode 100644 console/internal/fission/environment.go create mode 100644 console/internal/fission/namespace.go create mode 100644 console/internal/model/types.go create mode 100644 console/internal/runtime/doc.go create mode 100644 console/internal/runtime/entrypoint.go create mode 100644 console/internal/runtime/go.go create mode 100644 console/internal/runtime/helpers.go create mode 100644 console/internal/runtime/nodejs.go create mode 100644 console/internal/runtime/script.go diff --git a/console/Dockerfile b/console/Dockerfile index 1019e31..78cfe37 100644 --- a/console/Dockerfile +++ b/console/Dockerfile @@ -3,7 +3,7 @@ WORKDIR /build COPY go.mod go.sum ./ RUN go mod download COPY . . -RUN CGO_ENABLED=0 GOOS=linux go build -trimpath -ldflags="-s -w" -o fission-console . +RUN CGO_ENABLED=0 GOOS=linux go build -trimpath -ldflags="-s -w" -o fission-console ./cmd/server/ FROM alpine:3.20 RUN apk add --no-cache ca-certificates nodejs python3 ruby perl php83 diff --git a/console/cmd/server/main.go b/console/cmd/server/main.go new file mode 100644 index 0000000..0c8f227 --- /dev/null +++ b/console/cmd/server/main.go @@ -0,0 +1,107 @@ +// cmd/server/main.go — точка входа fission-console сервера. +// Отвечает только за: парсинг env vars, создание зависимостей, запуск HTTP сервера. +// Вся бизнес-логика находится в internal/api/, internal/fission/, internal/runtime/. +package main + +import ( + "log" + "net/http" + "os" + "strings" + "time" + + "fission-console/internal/api" + + "k8s.io/client-go/dynamic" + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/clientcmd" +) + +func main() { + kubeconfig := strings.TrimSpace(os.Getenv("KUBECONFIG")) + namespace := envDefault("FISSION_NAMESPACE", "default") + routerURL := strings.TrimRight(envDefault("FISSION_ROUTER_URL", "http://router.fission.svc.cluster.local"), "/") + port := envDefault("PORT", "8090") + httpTimeout := envDurationDefault("FISSION_HTTP_TIMEOUT", 30*time.Second) + invokeTimeout := envDurationDefault("FISSION_INVOKE_TIMEOUT", 20*time.Second) + + // Строим kubernetes config: сначала kubeconfig файл, потом in-cluster, потом default + cfg, err := buildConfig(kubeconfig) + if err != nil { + log.Fatalf("build kube config: %v", err) + } + + dyn, err := dynamic.NewForConfig(cfg) + if err != nil { + log.Fatalf("create dynamic client: %v", err) + } + + srv := api.NewServer(api.Config{ + Dyn: dyn, + Namespace: namespace, + RouterURL: routerURL, + HTTPTimeout: httpTimeout, + InvokeTimeout: invokeTimeout, + SATokenPath: envDefault("SA_TOKEN_PATH", "/var/run/secrets/kubernetes.io/serviceaccount/token"), + AuthUser: envDefault("FISSION_AUTH_USERNAME", ""), + AuthPass: envDefault("FISSION_AUTH_PASSWORD", ""), + TestMode: os.Getenv("FISSION_TEST_MODE") == "true", + // --- ai/ask feature --- + LLMUrl: envDefault("FISSION_LLM_URL", "https://api.aillm.ru"), + LLMKey: os.Getenv("FISSION_LLM_KEY"), + // --- end ai/ask feature --- + }) + + // Запускаем фоновые горутины: reaper истёкших функций + reconciler namespace-ов + nsm := srv.NSManager() + nsm.StartExpiryReaper(envDurationDefault("REAPER_INTERVAL", 5*time.Minute)) + nsm.StartNSReconciler(envDurationDefault("NS_RECONCILE_INTERVAL", 2*time.Minute)) + + httpServer := &http.Server{ + Addr: ":" + port, + Handler: srv.Handler(), + // ReadHeaderTimeout защищает от Slowloris атаки (медленная отправка заголовков) + ReadHeaderTimeout: 10 * time.Second, + } + + log.Printf("fission-console listening on :%s (namespace=%s)", port, namespace) + log.Fatal(httpServer.ListenAndServe()) +} + +// buildConfig строит kubernetes rest.Config в порядке приоритета: +// 1. KUBECONFIG env var или флаг — для локальной разработки +// 2. In-cluster config — для запуска внутри pod +// 3. Default kubeconfig — ~/.kube/config +func buildConfig(kubeconfig string) (*rest.Config, error) { + if kubeconfig != "" { + return clientcmd.BuildConfigFromFlags("", kubeconfig) + } + cfg, err := rest.InClusterConfig() + if err == nil { + return cfg, nil + } + // Fallback: пробуем default kubeconfig (empty rules → ~/.kube/config) + loadingRules := &clientcmd.ClientConfigLoadingRules{} + clientCfg := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(loadingRules, &clientcmd.ConfigOverrides{}) + return clientCfg.ClientConfig() +} + +func envDefault(key, fallback string) string { + if v := strings.TrimSpace(os.Getenv(key)); v != "" { + return v + } + return fallback +} + +func envDurationDefault(key string, fallback time.Duration) time.Duration { + raw := strings.TrimSpace(os.Getenv(key)) + if raw == "" { + return fallback + } + d, err := time.ParseDuration(raw) + if err != nil || d <= 0 { + log.Printf("invalid duration for %s=%q, using default %s", key, raw, fallback) + return fallback + } + return d +} diff --git a/console/internal/api/ai_check.go b/console/internal/api/ai_check.go new file mode 100644 index 0000000..8a56941 --- /dev/null +++ b/console/internal/api/ai_check.go @@ -0,0 +1,219 @@ +package api + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "go/parser" + "go/token" + "io" + "net/http" + "os" + "os/exec" + "strings" + "time" +) + +// handleAICheck проверяет синтаксис кода через реальный линтер языка. +// POST /console/api/ai/check +// Body: {"language":"python","code":"..."} +// Response: {"ok":true,"result":"..."} +// +// Почему реальные линтеры, а не LLM: +// LLM часто "исправляет" валидный код и врёт о наличии ошибок. +// node --check / python3 -m py_compile / ruby -c / php -l / perl -c — детерминированы и надёжны. +// Go использует go/parser прямо в процессе — без внешних команд, быстро. +func (s *Server) handleAICheck(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + writeJSONError(w, http.StatusMethodNotAllowed, "method not allowed") + return + } + + var req struct { + Language string `json:"language"` + Code string `json:"code"` + } + if err := json.NewDecoder(io.LimitReader(r.Body, 64*1024)).Decode(&req); err != nil { + writeJSONError(w, http.StatusBadRequest, "invalid JSON: "+err.Error()) + return + } + if strings.TrimSpace(req.Code) == "" { + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(map[string]any{"ok": false, "result": "Код пустой."}) + return + } + + // Конфигурация линтера для каждого языка + type linterCfg struct { + ext string + cmd []string + } + langs := map[string]linterCfg{ + "nodejs": {ext: ".js", cmd: []string{"node", "--check"}}, + "python": {ext: ".py", cmd: []string{"python3", "-m", "py_compile"}}, + "ruby": {ext: ".rb", cmd: []string{"ruby", "-c"}}, + "php": {ext: ".php", cmd: []string{"php", "-l"}}, + "perl": {ext: ".pl", cmd: []string{"perl", "-c"}}, + "go": {ext: ".go", cmd: nil}, // go проверяется через go/parser в процессе + } + + cfg, ok := langs[req.Language] + if !ok { + writeJSONError(w, http.StatusBadRequest, "unsupported language: "+req.Language) + return + } + + var isOK bool + var result string + + if req.Language == "go" { + // Go: используем go/parser прямо в процессе — без внешних команд и temp файлов + fset := token.NewFileSet() + _, parseErr := parser.ParseFile(fset, "code.go", req.Code, parser.AllErrors) + if parseErr == nil { + isOK = true + result = "✅ Синтаксис корректен." + } else { + isOK = false + result = parseErr.Error() + } + } else { + // Записываем код во временный файл с нужным расширением + // (линтеры ruby/php определяют режим проверки по расширению) + tmpf, err := os.CreateTemp("", "fission-lint-*"+cfg.ext) + if err != nil { + writeJSONError(w, http.StatusInternalServerError, "tmp file: "+err.Error()) + return + } + defer os.Remove(tmpf.Name()) // очищаем независимо от результата + + if _, err := tmpf.WriteString(req.Code); err != nil { + tmpf.Close() + writeJSONError(w, http.StatusInternalServerError, "write tmp: "+err.Error()) + return + } + tmpf.Close() + + args := append(cfg.cmd, tmpf.Name()) + ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) + defer cancel() + + //nolint:gosec — cfg.cmd содержит только захардкоженные команды из langs (не user input) + out, err := exec.CommandContext(ctx, args[0], args[1:]...).CombinedOutput() + outStr := strings.TrimSpace(string(out)) + // Убираем путь к tmp-файлу из вывода — пользователю незачем его видеть + outStr = strings.ReplaceAll(outStr, tmpf.Name(), "") + + if err == nil { + isOK = true + result = "✅ Синтаксис корректен." + } else { + isOK = false + if outStr != "" { + result = outStr + } else { + result = "Ошибка синтаксиса (линтер вернул код " + fmt.Sprintf("%v", err) + ")" + } + } + } + + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(map[string]any{ + "ok": isOK, + "result": result, + }) +} + +// ============================================================================= +// ai/ask feature — удалить весь блок до "end ai/ask feature" чтобы выкосить +// ============================================================================= + +// handleAIAsk отвечает на произвольный вопрос пользователя через LLM. +// POST /console/api/ai/ask +// Body: {"question":"..."} +// Response: {"answer":"..."} +func (s *Server) handleAIAsk(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + writeJSONError(w, http.StatusMethodNotAllowed, "method not allowed") + return + } + if s.llmKey == "" { + writeJSONError(w, http.StatusServiceUnavailable, "AI-ассистент не настроен (FISSION_LLM_KEY не задан)") + return + } + + var req struct { + Question string `json:"question"` + } + if err := json.NewDecoder(io.LimitReader(r.Body, 4*1024)).Decode(&req); err != nil { + writeJSONError(w, http.StatusBadRequest, "invalid JSON: "+err.Error()) + return + } + if strings.TrimSpace(req.Question) == "" { + writeJSONError(w, http.StatusBadRequest, "question is required") + return + } + + body, _ := json.Marshal(map[string]any{ + "model": "gpt-oss-120b", + "messages": []map[string]string{ + {"role": "system", "content": "Ты умный ассистент. Отвечай кратко и по делу."}, + {"role": "user", "content": req.Question}, + }, + "max_tokens": 1024, + }) + + ctx, cancel := context.WithTimeout(r.Context(), 30*time.Second) + defer cancel() + + llmReq, err := http.NewRequestWithContext(ctx, http.MethodPost, + strings.TrimRight(s.llmURL, "/")+"/chat/completions", + bytes.NewReader(body), + ) + if err != nil { + writeJSONError(w, http.StatusInternalServerError, err.Error()) + return + } + llmReq.Header.Set("Content-Type", "application/json") + llmReq.Header.Set("Authorization", "Bearer "+s.llmKey) + + resp, err := s.http.Do(llmReq) + if err != nil { + writeJSONError(w, http.StatusBadGateway, "LLM недоступен: "+err.Error()) + return + } + defer resp.Body.Close() + + var llmResp struct { + Choices []struct { + Message struct { + Content string `json:"content"` + } `json:"message"` + } `json:"choices"` + Error *struct { + Message string `json:"message"` + } `json:"error"` + } + if err := json.NewDecoder(io.LimitReader(resp.Body, 64*1024)).Decode(&llmResp); err != nil { + writeJSONError(w, http.StatusBadGateway, "parse error: "+err.Error()) + return + } + if llmResp.Error != nil { + writeJSONError(w, http.StatusBadGateway, llmResp.Error.Message) + return + } + if len(llmResp.Choices) == 0 { + writeJSONError(w, http.StatusBadGateway, "empty response from LLM") + return + } + + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(map[string]any{ + "answer": strings.TrimSpace(llmResp.Choices[0].Message.Content), + }) +} + +// ============================================================================= +// end ai/ask feature +// ============================================================================= diff --git a/console/internal/api/auth.go b/console/internal/api/auth.go new file mode 100644 index 0000000..2186234 --- /dev/null +++ b/console/internal/api/auth.go @@ -0,0 +1,125 @@ +package api + +import ( + "context" + "crypto/sha256" + "encoding/base64" + "encoding/hex" + "encoding/json" + "fmt" + "net/http" + "strings" +) + +// ctxKeyNS — ключ для хранения namespace пользователя в context.Context. +// Использует приватный тип чтобы избежать коллизий с ключами из других пакетов. +type ctxKeyNS struct{} + +// userNS возвращает namespace пользователя из контекста запроса. +// Устанавливается в authMiddleware после успешной аутентификации. +func (s *Server) userNS(r *http.Request) string { + if ns, ok := r.Context().Value(ctxKeyNS{}).(string); ok && ns != "" { + return ns + } + // Fallback: использовать системный namespace (не должно происходить в prod) + return s.ns +} + +// authMiddleware оборачивает handler, добавляя аутентификацию и инициализацию namespace. +// +// В testMode (FISSION_TEST_MODE=true): +// - Deck API не вызывается +// - X-Test-Sub или X-Auth-Token задают sub → разные namespace-ы для тестирования +// +// В production: +// - X-Auth-Token валидируется через Deck API +// - Namespace вычисляется из JWT claim "sub" +func (s *Server) authMiddleware(h http.HandlerFunc) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + var ns string + + if s.testMode { + sub := strings.TrimSpace(r.Header.Get("X-Test-Sub")) + if sub == "" { + sub = strings.TrimSpace(r.Header.Get("X-Auth-Token")) + } + if sub == "" { + writeJSONError(w, http.StatusUnauthorized, "test mode: X-Test-Sub required") + return + } + h32 := sha256.Sum256([]byte(sub)) + ns = "fission-" + hex.EncodeToString(h32[:8]) + } else { + token := strings.TrimSpace(r.Header.Get("X-Auth-Token")) + env := strings.TrimSpace(strings.ToLower(r.Header.Get("X-Auth-Env"))) + if _, ok := deckAPIs[env]; !ok { + env = "test" + } + if token == "" { + writeJSONError(w, http.StatusUnauthorized, "unauthorized") + return + } + if err := s.validateDeckToken(token, env); err != nil { + writeJSONError(w, http.StatusUnauthorized, "unauthorized") + return + } + var err error + ns, err = namespaceFromJWT(token) + if err != nil { + ns = s.ns + } + } + + ctx := context.WithValue(r.Context(), ctxKeyNS{}, ns) + + // Гарантируем что namespace + RBAC + quota + netpol существуют. + // EnsureUserNS реализует singleflight + кэш + семафор параллелизма. + if ensureErr := s.nsManager.EnsureUserNS(ctx, ns); ensureErr != nil { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", ensureErr)) + return + } + + h(w, r.WithContext(ctx)) + } +} + +// namespaceFromJWT декодирует JWT payload (без верификации подписи), +// извлекает claim "sub" и вычисляет namespace: "fission-" + hex(SHA256(sub)[:8]). +// +// Подпись не проверяется — токен уже валидирован через Deck API (validateDeckToken). +// Здесь нам нужен только deterministic namespace name из sub claim. +func namespaceFromJWT(token string) (string, error) { + parts := strings.SplitN(token, ".", 3) + if len(parts) != 3 { + return "", fmt.Errorf("invalid JWT format") + } + payload := parts[1] + + // JWT использует base64url без padding — добавляем если нужно + switch len(payload) % 4 { + case 2: + payload += "==" + case 3: + payload += "=" + } + + // base64url без стандартного padding — пробуем оба варианта + decoded, err := base64.URLEncoding.DecodeString(payload) + if err != nil { + decoded, err = base64.StdEncoding.DecodeString(payload) + if err != nil { + return "", fmt.Errorf("decode JWT payload: %w", err) + } + } + + var claims map[string]any + if err := json.Unmarshal(decoded, &claims); err != nil { + return "", fmt.Errorf("unmarshal JWT claims: %w", err) + } + sub, _ := claims["sub"].(string) + if sub == "" { + return "", fmt.Errorf("JWT missing sub claim") + } + h := sha256.Sum256([]byte(sub)) + return "fission-" + hex.EncodeToString(h[:8]), nil +} diff --git a/console/internal/api/doc.go b/console/internal/api/doc.go new file mode 100644 index 0000000..ae7381a --- /dev/null +++ b/console/internal/api/doc.go @@ -0,0 +1,2 @@ +// Package api содержит HTTP сервер, роутинг, handlers и middleware. +package api diff --git a/console/internal/api/handlers.go b/console/internal/api/handlers.go new file mode 100644 index 0000000..523c219 --- /dev/null +++ b/console/internal/api/handlers.go @@ -0,0 +1,739 @@ +package api + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/base64" + "encoding/json" + "errors" + "fmt" + "io" + "log" + "net" + "net/http" + "regexp" + "strconv" + "strings" + "time" + + "fission-console/internal/fission" + "fission-console/internal/model" + "fission-console/internal/runtime" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" +) + +// validFuncName — RFC 1123 subdomain label: строчные буквы+цифры+дефис, без дефиса в начале/конце. +// Максимум 57 символов (не 63): самый длинный суффикс "-route" (HTTPTrigger) = 6 символов. +// 63 - 6 = 57. Fission webhook требует все связанные объекты <= 63 символов. +var validFuncName = regexp.MustCompile(`^[a-z0-9]([a-z0-9-]*[a-z0-9])?$`) + +// maxCodeSize — максимальный размер кода функции (1 MB). +// Выше — не имеет смысла для inline функции; лучше использовать Package с URL. +const maxCodeSize = 1 << 20 + +// handleFunctionsRoot обрабатывает запросы к /console/api/functions без имени функции. +// GET → список всех функций, POST → создать новую. +func (s *Server) handleFunctionsRoot(w http.ResponseWriter, r *http.Request) { + switch r.Method { + case http.MethodGet: + s.handleList(fission.FunctionGVR)(w, r) + case http.MethodPost: + s.handleCreateFunction(w, r) + default: + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + } +} + +// handleFunctionsAction обрабатывает запросы к /console/api/functions/:name[/action]. +// Парсит имя функции и опциональный sub-path ("code", "invoke"). +func (s *Server) handleFunctionsAction(w http.ResponseWriter, r *http.Request) { + // Убираем оба возможных префикса (legacy /api/ и основной /console/api/) + path := strings.TrimPrefix(r.URL.Path, "/api/functions/") + if path == r.URL.Path { + path = strings.TrimPrefix(r.URL.Path, "/console/api/functions/") + } + path = strings.Trim(path, "/") + if path == "" { + http.NotFound(w, r) + return + } + + parts := strings.Split(path, "/") + name := strings.TrimSpace(parts[0]) + if name == "" { + writeJSONError(w, http.StatusBadRequest, "function name is required") + return + } + + if len(parts) == 1 { + // /functions/:name — CRUD операции с конкретной функцией + switch r.Method { + case http.MethodGet: + s.handleGetFunction(w, r, name) + case http.MethodDelete: + s.handleDeleteFunction(w, r, name) + default: + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + } + return + } + + // /functions/:name/code — обновление кода + if len(parts) == 2 && parts[1] == "code" && r.Method == http.MethodPut { + s.handleUpdateFunctionCode(w, r, name) + return + } + + // /functions/:name/invoke — вызов функции + if len(parts) == 2 && parts[1] == "invoke" && r.Method == http.MethodPost { + s.handleInvokeFunction(w, r, name) + return + } + + http.NotFound(w, r) +} + +// handleCreateFunction создаёт новую функцию: Package + Function + HTTPTrigger. +// +// Порядок создания: Package → Function → HTTPTrigger. +// При ошибке на любом шаге откатываем уже созданные объекты (best-effort). +// TTL парсится ДО создания объектов — невалидный TTL не оставляет мусор. +func (s *Server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { + ns := s.userNS(r) + + var req model.CreateFunctionRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err)) + return + } + + // Гарантируем namespace — на случай прямого вызова API без handleAuth + nsCtx, nsCancel := context.WithTimeout(r.Context(), 60*time.Second) + defer nsCancel() + if err := s.nsManager.EnsureUserNS(nsCtx, ns); err != nil { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", err)) + return + } + + req.Name = strings.TrimSpace(req.Name) + req.Language = strings.TrimSpace(req.Language) + req.Environment = strings.TrimSpace(req.Environment) + req.Code = strings.TrimSpace(req.Code) + req.Entrypoint = strings.TrimSpace(req.Entrypoint) + req.Route = strings.TrimSpace(req.Route) + + // Валидация имени + if req.Name != "" && (!validFuncName.MatchString(req.Name) || len(req.Name) > 57) { + writeJSONError(w, http.StatusBadRequest, "invalid function name: must match ^[a-z0-9]([a-z0-9-]*[a-z0-9])?$ and be <= 57 chars") + return + } + if len(req.Code) > maxCodeSize { + writeJSONError(w, http.StatusBadRequest, "code exceeds 1MB limit") + return + } + + // Lazy создание Environment по языку (если язык указан явно) + if req.Language != "" { + envCtx, envCancel := context.WithTimeout(r.Context(), 15*time.Second) + defer envCancel() + envName, err := fission.EnsureEnvironment(envCtx, s.dyn, ns, req.Language) + if err != nil { + if strings.Contains(err.Error(), "unsupported language") { + writeJSONError(w, http.StatusBadRequest, err.Error()) + } else { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure environment: %v", err)) + } + return + } + req.Environment = envName + } + + if req.Name == "" || req.Environment == "" || req.Code == "" { + writeJSONError(w, http.StatusBadRequest, "name, environment/language and code are required") + return + } + if req.Entrypoint == "" { + req.Entrypoint = runtime.DefaultEntrypoint(req.Language) + } + if req.Route == "" { + // Namespace-prefix route: избегаем коллизий между пользователями + // (разные пользователи могут создать функцию с одинаковым именем) + nsShort := ns + if len(nsShort) > 12 { + nsShort = nsShort[len(nsShort)-12:] + } + req.Route = "/" + nsShort + "/" + req.Name + } + if !strings.HasPrefix(req.Route, "/") { + req.Route = "/" + req.Route + } + req.Methods = normalizeMethods(req.Methods) + + ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) + defer cancel() + + // Проверяем что environment существует (мог быть задан явно без language) + if _, err := s.dyn.Resource(fission.EnvironmentGVR).Namespace(ns).Get(ctx, req.Environment, metav1.GetOptions{}); err != nil { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("environment %q not found: %v", req.Environment, err)) + return + } + + pkgName := req.Name + "-pkg" + triggerName := req.Name + "-route" + + methodValues := make([]any, 0, len(req.Methods)) + for _, method := range req.Methods { + methodValues = append(methodValues, method) + } + + // Строим Package spec в зависимости от языка: + // - Go: source package → builder job компилирует в .so плагин + // - Node.js: deployment zip с ESM wrapper (package.json + main.js) + // - Остальные: deployment literal с кодом напрямую + var pkgSpec map[string]any + if req.Language == "go" { + srcZip, err := runtime.BuildGoSourceZip(req.Code) + if err != nil { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("build go source archive: %v", err)) + return + } + literal := base64.StdEncoding.EncodeToString(srcZip) + pkgSpec = map[string]any{ + "source": map[string]any{ + "type": "literal", + "literal": literal, + }, + "deployment": map[string]any{}, + "environment": map[string]any{"name": req.Environment, "namespace": ns}, + "buildcommand": "build", + } + } else { + var deployBytes []byte + if req.Language == "nodejs" { + zipBytes, zipErr := runtime.BuildJSDeployZip(req.Code) + if zipErr != nil { + writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build nodejs archive: %v", zipErr)) + return + } + deployBytes = zipBytes + } else { + // Python, PHP, Ruby, Perl — код передаётся как есть в deployment.literal + deployBytes = []byte(req.Code) + } + pkgSpec = map[string]any{ + "deployment": map[string]any{"type": "literal", "literal": base64.StdEncoding.EncodeToString(deployBytes)}, + "environment": map[string]any{"name": req.Environment, "namespace": ns}, + "source": map[string]any{}, + } + } + + pkg := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "fission.io/v1", + "kind": "Package", + "metadata": map[string]any{"name": pkgName, "namespace": ns}, + "spec": pkgSpec, + }} + + // Парсим TTL ДО создания K8s ресурсов — невалидный TTL не оставляет мусор + fnAnnotations := map[string]any{ + "fission-console/language": req.Language, + } + if req.TTL != "" { + expiresAt, ttlErr := parseTTL(req.TTL) + if ttlErr != nil { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("invalid ttl %q: %v", req.TTL, ttlErr)) + return + } + fnAnnotations["fission-console/expires-at"] = expiresAt.UTC().Format(time.RFC3339) + } + + if _, err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Create(ctx, pkg, metav1.CreateOptions{}); err != nil { + if apierrors.IsAlreadyExists(err) { + writeJSONError(w, http.StatusConflict, fmt.Sprintf("function %q already exists", req.Name)) + return + } + if apierrors.IsInvalid(err) { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("invalid function spec: %v", err)) + return + } + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create package: %v", err)) + return + } + + fn := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "fission.io/v1", + "kind": "Function", + "metadata": map[string]any{"name": req.Name, "namespace": ns, "annotations": fnAnnotations}, + "spec": map[string]any{ + "environment": map[string]any{"name": req.Environment, "namespace": ns}, + "InvokeStrategy": map[string]any{ + "ExecutionStrategy": map[string]any{"ExecutorType": "poolmgr"}, + "StrategyType": "execution", + }, + "package": map[string]any{ + "packageref": map[string]any{"name": pkgName, "namespace": ns}, + "functionName": req.Entrypoint, + }, + }, + }} + + if _, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Create(ctx, fn, metav1.CreateOptions{}); err != nil { + // Откатываем Package если Function не создалась + _ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) + if apierrors.IsAlreadyExists(err) { + writeJSONError(w, http.StatusConflict, fmt.Sprintf("function %q already exists", req.Name)) + return + } + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create function: %v", err)) + return + } + + httpTrigger := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "fission.io/v1", + "kind": "HTTPTrigger", + "metadata": map[string]any{"name": triggerName, "namespace": ns}, + "spec": map[string]any{ + "relativeurl": req.Route, + "methods": methodValues, + "createingress": true, + "functionref": map[string]any{"type": "name", "name": req.Name}, + }, + }} + + if _, err := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).Create(ctx, httpTrigger, metav1.CreateOptions{}); err != nil { + // Откатываем Function и Package + _ = s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Delete(ctx, req.Name, metav1.DeleteOptions{}) + _ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create httptrigger: %v", err)) + return + } + + writeAnyJSON(w, http.StatusCreated, map[string]any{ + "name": req.Name, + "package": pkgName, + "httptrigger": triggerName, + "route": req.Route, + "expires_at": fnAnnotations["fission-console/expires-at"], + }) +} + +// handleGetFunction возвращает детали функции: код, environment, route, methods. +func (s *Server) handleGetFunction(w http.ResponseWriter, r *http.Request, name string) { + ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) + defer cancel() + ns := s.userNS(r) + + fn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, name, metav1.GetOptions{}) + if err != nil { + status := http.StatusBadGateway + if apierrors.IsNotFound(err) { + status = http.StatusNotFound + } + writeJSONError(w, status, fmt.Sprintf("get function %q: %v", name, err)) + return + } + + packageName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name") + environment, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") + entrypoint, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "functionName") + + // Извлекаем исходный код из Package (пробуем source.literal, потом deployment.literal, потом url) + code := "" + if packageName != "" { + pkg, pkgErr := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Get(ctx, packageName, metav1.GetOptions{}) + if pkgErr == nil { + code = extractPackageSourceCode(ctx, s, pkg) + } + } + + // Ищем HTTPTrigger для получения route и methods + route := "" + methods := []string{} + triggers, trigErr := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if trigErr == nil { + for _, trig := range triggers.Items { + refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name") + if refName != name { + continue + } + route, _, _ = unstructured.NestedString(trig.Object, "spec", "relativeurl") + methods, _, _ = unstructured.NestedStringSlice(trig.Object, "spec", "methods") + break + } + } + + writeAnyJSON(w, http.StatusOK, map[string]any{ + "name": name, + "namespace": ns, + "environment": environment, + "package": packageName, + "entrypoint": entrypoint, + "code": code, + "route": route, + "methods": methods, + "raw": fn.Object, + }) +} + +// handleUpdateFunctionCode обновляет код уже существующей функции. +// Обновляет Package.spec.deployment.literal и синхронизирует resourceVersion в Function. +// +// Почему нужно обновлять resourceVersion в Function: +// Fission executor кэширует Package по resourceVersion. Без обновления в Function +// executor будет использовать старый код до перезапуска pod. +func (s *Server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request, name string) { + var req model.UpdateCodeRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err)) + return + } + req.Code = strings.TrimSpace(req.Code) + if req.Code == "" { + writeJSONError(w, http.StatusBadRequest, "code is required") + return + } + + ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) + defer cancel() + ns := s.userNS(r) + + fn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, name, metav1.GetOptions{}) + if err != nil { + status := http.StatusBadGateway + if apierrors.IsNotFound(err) { + status = http.StatusNotFound + } + writeJSONError(w, status, fmt.Sprintf("get function %q: %v", name, err)) + return + } + + pkgName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name") + if pkgName == "" { + writeJSONError(w, http.StatusBadGateway, "function has no package reference") + return + } + + pkg, err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Get(ctx, pkgName, metav1.GetOptions{}) + if err != nil { + status := http.StatusBadGateway + if apierrors.IsNotFound(err) { + status = http.StatusNotFound + } + writeJSONError(w, status, fmt.Sprintf("get package %q: %v", pkgName, err)) + return + } + + // Определяем язык из аннотации — нужен для правильной упаковки + lang, _, _ := unstructured.NestedString(fn.Object, "metadata", "annotations", "fission-console/language") + var deployBytes []byte + if lang == "nodejs" { + zipBytes, zipErr := runtime.BuildJSDeployZip(req.Code) + if zipErr != nil { + writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build nodejs archive: %v", zipErr)) + return + } + deployBytes = zipBytes + } else { + deployBytes = []byte(req.Code) + } + + literal := base64.StdEncoding.EncodeToString(deployBytes) + if err := unstructured.SetNestedField(pkg.Object, literal, "spec", "deployment", "literal"); err != nil { + writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set package literal: %v", err)) + return + } + + updatedPkg, err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Update(ctx, pkg, metav1.UpdateOptions{}) + if err != nil { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("update package %q: %v", pkgName, err)) + return + } + + // Синхронизируем resourceVersion в Function.spec.package.packageref + // Это триггерит executor перезагрузить код в pool pod + if err := unstructured.SetNestedField(fn.Object, updatedPkg.GetResourceVersion(), "spec", "package", "packageref", "resourceversion"); err != nil { + writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set function package resourceversion: %v", err)) + return + } + if _, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Update(ctx, fn, metav1.UpdateOptions{}); err != nil { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("update function %q package ref: %v", name, err)) + return + } + + writeAnyJSON(w, http.StatusOK, map[string]any{ + "updated": true, + "package": pkgName, + "package_resourceversion": updatedPkg.GetResourceVersion(), + }) +} + +// handleInvokeFunction вызывает функцию через Fission router. +// Определяет реальный URL из HTTPTrigger, выбирает метод (POST/GET). +func (s *Server) handleInvokeFunction(w http.ResponseWriter, r *http.Request, name string) { + bodyBytes, err := io.ReadAll(r.Body) + if err != nil { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("read request body: %v", err)) + return + } + if len(bytes.TrimSpace(bodyBytes)) == 0 { + bodyBytes = []byte("{}") + } + + invokeTimeout := s.invokeTimeout + if invokeTimeout <= 0 { + invokeTimeout = 20 * time.Second + } + + ctx, cancel := context.WithTimeout(r.Context(), invokeTimeout) + defer cancel() + ns := s.userNS(r) + + // Проверяем существование функции до вызова — лучше 404 чем непонятный timeout + if _, err2 := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, name, metav1.GetOptions{}); err2 != nil { + if apierrors.IsNotFound(err2) { + writeJSONError(w, http.StatusNotFound, fmt.Sprintf("function %q not found", name)) + return + } + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function %q: %v", name, err2)) + return + } + + // Ищем HTTPTrigger чтобы получить реальный URL и метод + invokeURL := fmt.Sprintf("%s/fission-function/v2/functions/%s", s.routerURL, name) + invokeMethod := http.MethodPost + triggers, err := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if err == nil { + for _, trig := range triggers.Items { + refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name") + if refName != name { + continue + } + route, _, _ := unstructured.NestedString(trig.Object, "spec", "relativeurl") + methods, _, _ := unstructured.NestedStringSlice(trig.Object, "spec", "methods") + hasPost, hasGet := false, false + for _, m := range methods { + switch strings.ToUpper(strings.TrimSpace(m)) { + case http.MethodPost: + hasPost = true + case http.MethodGet: + hasGet = true + } + } + if route != "" { + if !strings.HasPrefix(route, "/") { + route = "/" + route + } + invokeURL = s.routerURL + route + // Если функция поддерживает только GET — используем GET + if !hasPost && hasGet { + invokeMethod = http.MethodGet + } + break + } + } + } + + start := time.Now() + var invokeBody io.Reader + if invokeMethod == http.MethodPost { + invokeBody = bytes.NewReader(bodyBytes) + } + + req, err := http.NewRequestWithContext(ctx, invokeMethod, invokeURL, invokeBody) + if err != nil { + writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build invoke request: %v", err)) + return + } + if invokeMethod == http.MethodPost { + req.Header.Set("Content-Type", "application/json") + } + if token := s.getRouterToken(); token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + + resp, err := s.http.Do(req) + if err != nil { + // Отличаем timeout от сетевой ошибки — timeout часто означает что функция не специализировалась + if errors.Is(err, context.DeadlineExceeded) { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s: function specialization likely failed (for example, syntax error)", name, invokeTimeout)) + return + } + var netErr net.Error + if errors.As(err, &netErr) && netErr.Timeout() { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s: function specialization likely failed (for example, syntax error)", name, invokeTimeout)) + return + } + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q: %v", name, err)) + return + } + defer resp.Body.Close() + + respBody, _ := io.ReadAll(resp.Body) + writeAnyJSON(w, http.StatusOK, map[string]any{ + "status": resp.StatusCode, + "latency_ms": time.Since(start).Milliseconds(), + "invoke_url": invokeURL, + "response_raw": string(respBody), + }) +} + +// handleDeleteFunction удаляет функцию и связанные объекты: HTTPTrigger, Package. +// После удаления вызывает CleanupEnvironmentIfUnused — убирает environment если язык больше не используется. +func (s *Server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, name string) { + ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) + defer cancel() + ns := s.userNS(r) + + // Получаем Function чтобы знать pkgName и envName для cleanup + var pkgName, envName string + fn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, name, metav1.GetOptions{}) + if err != nil { + if apierrors.IsNotFound(err) { + writeJSONError(w, http.StatusNotFound, fmt.Sprintf("function %q not found", name)) + return + } + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function %q: %v", name, err)) + return + } + pkgName, _, _ = unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name") + envName, _, _ = unstructured.NestedString(fn.Object, "spec", "environment", "name") + + // Удаляем связанные HTTPTrigger-ы + triggers, err := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if err == nil { + for _, trig := range triggers.Items { + refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name") + if refName == name { + _ = s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{}) + } + } + } + + if err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Delete(ctx, name, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("delete function %q: %v", name, err)) + return + } + + if pkgName != "" { + if err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("delete package %q: %v", pkgName, err)) + return + } + } + + // Убираем environment pool pods если язык больше не используется (best-effort) + if envName != "" { + cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cleanupCancel() + fission.CleanupEnvironmentIfUnused(cleanupCtx, s.dyn, ns, envName) + } + + // Сигналим reconciler: батчит изменения без race condition и rolling restarts + select { + case s.nsManager.NSReconcileCh <- struct{}{}: + default: + } + + writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName}) +} + +// handleAuth обрабатывает POST /console/api/auth. +// Валидирует токен, создаёт namespace, возвращает namespace пользователя. +func (s *Server) handleAuth(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + writeJSONError(w, http.StatusMethodNotAllowed, "method not allowed") + return + } + + var body struct { + Token string `json:"token"` + Env string `json:"env"` + } + if err := json.NewDecoder(r.Body).Decode(&body); err != nil || strings.TrimSpace(body.Token) == "" { + writeJSONError(w, http.StatusBadRequest, "token required") + return + } + + env := strings.TrimSpace(strings.ToLower(body.Env)) + if _, ok := deckAPIs[env]; !ok { + env = "test" + } + + var ns string + if s.testMode { + // testMode: токен — это email (sub), Deck не вызывается + if !strings.Contains(body.Token, "@") { + writeJSONError(w, http.StatusUnauthorized, "invalid token") + return + } + h32 := sha256.Sum256([]byte(body.Token)) + ns = "fission-" + fmt.Sprintf("%x", h32[:8]) + } else { + if err := s.validateDeckToken(body.Token, env); err != nil { + writeJSONError(w, http.StatusUnauthorized, "invalid token") + return + } + var err error + ns, err = namespaceFromJWT(body.Token) + if err != nil { + log.Printf("handleAuth: namespaceFromJWT: %v", err) + ns = s.ns + } + } + + ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second) + defer cancel() + if ensureErr := s.nsManager.EnsureUserNS(ctx, ns); ensureErr != nil { + log.Printf("handleAuth: ensureUserNS %s: %v", ns, ensureErr) + } + + w.Header().Set("Content-Type", "application/json; charset=utf-8") + _ = json.NewEncoder(w).Encode(map[string]any{"ok": true, "env": env, "namespace": ns}) +} + +// parseTTL парсит строку TTL и возвращает время истечения. +// Поддерживаемые форматы: Go duration (1h, 30m, 24h) и дни (1d, 7d, 30d). +// Суффикс "d" не поддерживается стандартным time.ParseDuration — обрабатываем отдельно. +func parseTTL(ttl string) (time.Time, error) { + if strings.HasSuffix(ttl, "d") { + days, err := strconv.Atoi(strings.TrimSuffix(ttl, "d")) + if err != nil || days <= 0 { + return time.Time{}, fmt.Errorf("invalid days value: %q", ttl) + } + return time.Now().Add(time.Duration(days) * 24 * time.Hour), nil + } + d, err := time.ParseDuration(ttl) + if err != nil { + return time.Time{}, err + } + if d <= 0 { + return time.Time{}, fmt.Errorf("ttl must be positive") + } + return time.Now().Add(d), nil +} + +// normalizeMethods приводит список HTTP методов к верхнему регистру, убирает дубли. +// Если список пустой или все элементы пустые — возвращает ["GET"]. +func normalizeMethods(in []string) []string { + if len(in) == 0 { + return []string{"GET"} + } + out := make([]string, 0, len(in)) + seen := map[string]bool{} + for _, method := range in { + m := strings.ToUpper(strings.TrimSpace(method)) + if m == "" || seen[m] { + continue + } + seen[m] = true + out = append(out, m) + } + if len(out) == 0 { + return []string{"GET"} + } + return out +} diff --git a/console/internal/api/middleware.go b/console/internal/api/middleware.go new file mode 100644 index 0000000..7eb5489 --- /dev/null +++ b/console/internal/api/middleware.go @@ -0,0 +1,56 @@ +package api + +import ( + "log" + "net/http" +) + +// logRequests логирует метод и путь каждого входящего запроса. +// Middleware-обёртка: оборачивает любой http.Handler. +func logRequests(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + log.Printf("%s %s", r.Method, r.URL.Path) + next.ServeHTTP(w, r) + }) +} + +// withCORS добавляет CORS заголовки для разрешения cross-origin запросов. +// +// Почему Access-Control-Allow-Origin: *: +// Console — одностраничное приложение, может обслуживаться с любого домена. +// Строгий CORS с конкретным Origin потребовал бы конфигурации при деплое. +// Все чувствительные операции защищены X-Auth-Token — не cookie, не ambient credentials. +func withCORS(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Access-Control-Allow-Origin", "*") + w.Header().Set("Access-Control-Allow-Methods", "GET,POST,PUT,PATCH,DELETE,OPTIONS") + w.Header().Set("Access-Control-Allow-Headers", "Content-Type, Authorization, X-Auth-Token, X-Auth-Env") + // Preflight OPTIONS — отвечаем сразу без передачи в handler + if r.Method == http.MethodOptions { + w.WriteHeader(http.StatusNoContent) + return + } + next.ServeHTTP(w, r) + }) +} + +// withSecurityHeaders добавляет стандартные HTTP security заголовки. +// Снижает риск XSS, clickjacking и утечки Referer. +func withSecurityHeaders(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + // nosniff: браузер не должен угадывать MIME тип — только то что в Content-Type + w.Header().Set("X-Content-Type-Options", "nosniff") + // DENY: запрещаем отображение страницы в iframe (clickjacking protection) + w.Header().Set("X-Frame-Options", "DENY") + // strict-origin-when-cross-origin: скрываем полный URL при cross-origin запросах + w.Header().Set("Referrer-Policy", "strict-origin-when-cross-origin") + // Отключаем браузерные API которые не нужны console + w.Header().Set("Permissions-Policy", "camera=(), microphone=(), geolocation=()") + // CSP: разрешаем только ресурсы с того же origin + inline scripts/styles (нужны для embedded UI) + w.Header().Set("Content-Security-Policy", + "default-src 'self'; script-src 'self' 'unsafe-inline'; style-src 'self' 'unsafe-inline'; "+ + "img-src 'self' data:; connect-src 'self'; font-src 'self' data:; object-src 'none'; "+ + "frame-ancestors 'none'; base-uri 'self'; form-action 'self'; upgrade-insecure-requests; block-all-mixed-content") + next.ServeHTTP(w, r) + }) +} diff --git a/console/internal/api/ns_status.go b/console/internal/api/ns_status.go new file mode 100644 index 0000000..d82d9e7 --- /dev/null +++ b/console/internal/api/ns_status.go @@ -0,0 +1,120 @@ +package api + +import ( + "context" + "net/http" + "os" + "strings" + "time" + + "fission-console/internal/fission" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// handleNSStatus возвращает статус инициализации пользовательского namespace. +// Используется UI для отображения прогресса при первом входе нового пользователя. +// +// Три стадии: +// 1. Namespace существует и Active +// 2. Namespace добавлен в FISSION_RESOURCE_NAMESPACES executor deployment +// 3. Executor pod Running + Ready (готов принимать функции) +func (s *Server) handleNSStatus(w http.ResponseWriter, r *http.Request) { + ns := s.userNS(r) + ctx, cancel := context.WithTimeout(r.Context(), 15*time.Second) + defer cancel() + + fissionNS := os.Getenv("FISSION_SYSTEM_NAMESPACE") + if fissionNS == "" { + fissionNS = "fission" + } + + type stageInfo struct { + Name string `json:"name"` + Done bool `json:"done"` + } + stages := []stageInfo{ + {Name: "Создание пространства имён"}, + {Name: "Регистрация в Fission"}, + {Name: "Прогрев окружений"}, + } + + // Stage 1: namespace существует и Active + nsObj, err := s.dyn.Resource(fission.NamespaceGVR).Get(ctx, ns, metav1.GetOptions{}) + if err == nil { + phase, _, _ := unstructured.NestedString(nsObj.Object, "status", "phase") + stages[0].Done = phase == "Active" + } + + // Stage 2: namespace в FISSION_RESOURCE_NAMESPACES executor deployment + if stages[0].Done { + execDep, err2 := s.dyn.Resource(fission.DeploymentGVR).Namespace(fissionNS).Get(ctx, "executor", metav1.GetOptions{}) + if err2 == nil { + stages[1].Done = nsInFissionEnv(execDep, ns) + } + } + + // Stage 3: executor pod Running + Ready + if stages[1].Done { + podGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "pods"} + podList, err3 := s.dyn.Resource(podGVR).Namespace(fissionNS).List(ctx, metav1.ListOptions{ + LabelSelector: "svc=executor", + }) + if err3 == nil { + for _, pod := range podList.Items { + phase, _, _ := unstructured.NestedString(pod.Object, "status", "phase") + if phase != "Running" { + continue + } + conditions, _, _ := unstructured.NestedSlice(pod.Object, "status", "conditions") + for _, c := range conditions { + cond, ok := c.(map[string]any) + if !ok { + continue + } + if cond["type"] == "Ready" && cond["status"] == "True" { + stages[2].Done = true + } + } + } + } + } + + ready := stages[0].Done && stages[1].Done && stages[2].Done + writeAnyJSON(w, http.StatusOK, map[string]any{ + "ready": ready, + "stages": stages, + }) +} + +// nsInFissionEnv проверяет что namespace ns содержится в FISSION_RESOURCE_NAMESPACES +// первого контейнера данного deployment. +func nsInFissionEnv(dep *unstructured.Unstructured, ns string) bool { + containers, _, _ := unstructured.NestedSlice(dep.Object, "spec", "template", "spec", "containers") + for _, c := range containers { + cont, ok := c.(map[string]any) + if !ok { + continue + } + envs, _, _ := unstructured.NestedSlice(cont, "env") + for _, e := range envs { + env, ok := e.(map[string]any) + if !ok { + continue + } + if env["name"] == "FISSION_RESOURCE_NAMESPACES" { + if v, ok := env["value"].(string); ok { + for _, p := range strings.Split(v, ",") { + if strings.TrimSpace(p) == ns { + return true + } + } + } + } + } + break // только первый контейнер + } + return false +} diff --git a/console/internal/api/package.go b/console/internal/api/package.go new file mode 100644 index 0000000..9fad159 --- /dev/null +++ b/console/internal/api/package.go @@ -0,0 +1,156 @@ +package api + +import ( + "archive/zip" + "bytes" + "context" + "encoding/base64" + "io" + "net/http" + "sort" + "strings" + "unicode/utf8" + + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" +) + +// extractPackageSourceCode пытается извлечь исходный код из Fission Package. +// Порядок попыток: +// 1. spec.source.literal (base64) — Go функции (source package) +// 2. spec.deployment.literal (base64) — интерпретируемые языки +// 3. spec.source.url / spec.deployment.url — скачиваем архив по URL +// +// Декодирование: если данные — валидный UTF-8, возвращаем как есть. +// Если это zip — ищем известные файлы (main.py, handler.rb и т.д.). +func extractPackageSourceCode(ctx context.Context, s *Server, pkg *unstructured.Unstructured) string { + literalPaths := [][]string{ + {"spec", "source", "literal"}, + {"spec", "deployment", "literal"}, + } + for _, p := range literalPaths { + literal, found, _ := unstructured.NestedString(pkg.Object, p...) + if !found || strings.TrimSpace(literal) == "" { + continue + } + if code, err := decodeLiteralToSource(literal); err == nil && strings.TrimSpace(code) != "" { + return code + } + } + + // Fallback: скачиваем по URL + urlPaths := [][]string{ + {"spec", "source", "url"}, + {"spec", "deployment", "url"}, + } + for _, p := range urlPaths { + urlValue, found, _ := unstructured.NestedString(pkg.Object, p...) + if !found || strings.TrimSpace(urlValue) == "" { + continue + } + archiveBytes, err := fetchPackageArchive(ctx, s, urlValue) + if err != nil { + continue + } + if code, err := decodeArchiveBytesToSource(archiveBytes); err == nil && strings.TrimSpace(code) != "" { + return code + } + } + + return "" +} + +// fetchPackageArchive скачивает архив функции по URL из Fission storage. +func fetchPackageArchive(ctx context.Context, s *Server, archiveURL string) ([]byte, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, archiveURL, nil) + if err != nil { + return nil, err + } + resp, err := s.http.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return nil, io.ErrUnexpectedEOF + } + return io.ReadAll(resp.Body) +} + +// decodeLiteralToSource декодирует base64-строку из Package.spec.*.literal в исходный код. +func decodeLiteralToSource(literal string) (string, error) { + decoded, err := base64.StdEncoding.DecodeString(literal) + if err != nil { + return "", err + } + return decodeArchiveBytesToSource(decoded) +} + +// decodeArchiveBytesToSource преобразует байты (utf-8 строка или zip архив) в исходный код. +func decodeArchiveBytesToSource(decoded []byte) (string, error) { + if len(decoded) == 0 { + return "", io.ErrUnexpectedEOF + } + // Если байты — валидный UTF-8, возвращаем напрямую (python, ruby, perl, php) + if utf8.Valid(decoded) { + return string(decoded), nil + } + // PK signature: zip архив (nodejs, go source) + if len(decoded) >= 4 && bytes.Equal(decoded[:4], []byte{'P', 'K', 3, 4}) { + if src, zipErr := decodeZipSource(decoded); zipErr == nil { + return src, nil + } + } + return "", io.ErrUnexpectedEOF +} + +// decodeZipSource извлекает исходный UTF-8 файл из zip архива. +// Предпочитает хорошо известные имена файлов (main.py, handler.go и т.д.). +func decodeZipSource(zipBytes []byte) (string, error) { + reader, err := zip.NewReader(bytes.NewReader(zipBytes), int64(len(zipBytes))) + if err != nil { + return "", err + } + + // Сначала ищем по приоритетным именам + preferred := []string{"main.py", "main.js", "main.go", "handler.go", "handler.js", "handler.py"} + for _, name := range preferred { + for _, file := range reader.File { + if strings.EqualFold(file.Name, name) { + content, readErr := readZipFile(file) + if readErr == nil && utf8.Valid(content) { + return string(content), nil + } + } + } + } + + // Fallback: первый UTF-8 файл по алфавиту (не директория) + files := make([]*zip.File, 0, len(reader.File)) + for _, file := range reader.File { + if !file.FileInfo().IsDir() { + files = append(files, file) + } + } + sort.Slice(files, func(i, j int) bool { return files[i].Name < files[j].Name }) + + for _, file := range files { + content, readErr := readZipFile(file) + if readErr == nil && utf8.Valid(content) { + return string(content), nil + } + } + + return "", io.ErrUnexpectedEOF +} + +// readZipFile открывает и читает содержимое файла из zip архива. +func readZipFile(file *zip.File) ([]byte, error) { + rc, err := file.Open() + if err != nil { + return nil, err + } + defer rc.Close() + return io.ReadAll(rc) +} + +// decodeZipSource извлекает исходный UTF-8 файл из zip архива. \ No newline at end of file diff --git a/console/internal/api/server.go b/console/internal/api/server.go new file mode 100644 index 0000000..86e64e9 --- /dev/null +++ b/console/internal/api/server.go @@ -0,0 +1,295 @@ +package api + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "strings" + "sync" + "time" + + "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 *fission.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: fission.NewNSManager(cfg.Dyn), + } +} + +// RegisterRoutes регистрирует все HTTP маршруты на mux. +// Разделено на публичные (без auth) и приватные (с authMiddleware). +func (s *Server) RegisterRoutes(mux *http.ServeMux) { + // Публичные: health checks, UI, auth endpoint + 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) + + // Приватные: требуют 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.handleList(fission.TimeTrigGVR)) + + // Основные /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("/console/api/httptriggers", auth(s.handleList(fission.HTTPTrigGVR))) + mux.HandleFunc("/console/api/timetriggers", auth(s.handleList(fission.TimeTrigGVR))) + mux.HandleFunc("/console/api/ns/status", auth(s.handleNSStatus)) + 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() *fission.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 +} diff --git a/console/internal/fission/client.go b/console/internal/fission/client.go new file mode 100644 index 0000000..388b290 --- /dev/null +++ b/console/internal/fission/client.go @@ -0,0 +1,18 @@ +// Package fission содержит константы GVR и вспомогательные функции +// для работы с Fission CRDs через Kubernetes dynamic client. +package fission + +import "k8s.io/apimachinery/pkg/runtime/schema" + +// GroupVersionResource константы для всех Fission CRD и core K8s ресурсов. +// Вынесены сюда чтобы не дублировать между пакетами и легко обновить при смене версии API. +var ( + EnvironmentGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "environments"} + PackageGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "packages"} + FunctionGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "functions"} + HTTPTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "httptriggers"} + TimeTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "timetriggers"} + + NamespaceGVR = schema.GroupVersionResource{Group: "", Version: "v1", Resource: "namespaces"} + DeploymentGVR = schema.GroupVersionResource{Group: "apps", Version: "v1", Resource: "deployments"} +) diff --git a/console/internal/fission/environment.go b/console/internal/fission/environment.go new file mode 100644 index 0000000..f1965b1 --- /dev/null +++ b/console/internal/fission/environment.go @@ -0,0 +1,117 @@ +package fission + +import ( + "context" + "fmt" + "log" + + "fission-console/internal/model" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/client-go/dynamic" +) + +// EnsureEnvironment создаёт Environment CRD для языка lang в namespace ns если не существует. +// Возвращает имя environment (например "console-python-env"). +// +// Lazy creation: environments создаются только когда пользователь создаёт первую функцию +// на конкретном языке. Это экономит ресурсы — Fission pool pods поднимаются только +// под те языки которые реально используются, а не все сразу при создании namespace. +func EnsureEnvironment(ctx context.Context, dyn dynamic.Interface, ns, lang string) (string, error) { + langDef, ok := model.LangEnvMap[lang] + if !ok { + // Неизвестный язык — клиентская ошибка, возвращаем 400-совместимое сообщение + return "", fmt.Errorf("unsupported language: %q", lang) + } + envName := "console-" + lang + "-env" + + // Проверяем существование — Get быстрее чем Create+IsAlreadyExists + _, getErr := dyn.Resource(EnvironmentGVR).Namespace(ns).Get(ctx, envName, metav1.GetOptions{}) + if getErr == nil { + return envName, nil // уже существует — быстрый путь + } + if !apierrors.IsNotFound(getErr) { + return "", fmt.Errorf("check environment %q: %w", envName, getErr) + } + + // Строим и создаём Environment CRD + env := buildLangEnvironment(envName, ns, langDef) + if _, createErr := dyn.Resource(EnvironmentGVR).Namespace(ns).Create(ctx, env, metav1.CreateOptions{}); createErr != nil { + if apierrors.IsAlreadyExists(createErr) { + // Race condition: другой goroutine создал между Get и Create — это нормально + return envName, nil + } + return "", fmt.Errorf("create environment %q: %w", envName, createErr) + } + log.Printf("ensureEnvironment: created %s/%s", ns, envName) + return envName, nil +} + +// CleanupEnvironmentIfUnused удаляет Environment CRD если ни одна функция в namespace +// его не использует. Вызывается после удаления функции (handleDeleteFunction) +// и при срабатывании reaper-а (runExpiryReap). +// +// Не блокирующий: ошибки логируются, не возвращаются вызывающему — это best-effort cleanup. +// Fission увидит удаление Environment CRD и убьёт pool Deployment → поды умирают. +func CleanupEnvironmentIfUnused(ctx context.Context, dyn dynamic.Interface, ns, envName string) { + functions, err := dyn.Resource(FunctionGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if err != nil { + log.Printf("cleanupEnvironmentIfUnused: list functions in %s: %v", ns, err) + return + } + + // Если хотя бы одна функция ссылается на этот environment — оставляем его + for _, fn := range functions.Items { + fnEnv, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") + if fnEnv == envName { + return // ещё используется + } + } + + // Ни одна функция не ссылается — удаляем + if delErr := dyn.Resource(EnvironmentGVR).Namespace(ns).Delete(ctx, envName, metav1.DeleteOptions{}); delErr != nil && !apierrors.IsNotFound(delErr) { + log.Printf("cleanupEnvironmentIfUnused: delete env %s/%s: %v", ns, envName, delErr) + return + } + log.Printf("cleanupEnvironmentIfUnused: deleted unused env %s/%s", ns, envName) +} + +// buildLangEnvironment строит unstructured.Unstructured объект Environment CRD. +// Вся Fission-специфичная схема изолирована здесь — при обновлении Fission меняем только тут. +func buildLangEnvironment(name, ns string, def model.LangEnvDef) *unstructured.Unstructured { + // Version по умолчанию 3 (V2 protocol с async entrypoint). + // Исключение: perl-env поддерживает только V1 protocol → version=1. + envVersion := int64(3) + if def.Version != 0 { + envVersion = int64(def.Version) + } + + spec := map[string]any{ + "version": envVersion, + "runtime": map[string]any{ + "image": def.Image, + }, + "poolsize": int64(1), + } + + // BuilderImage задан только для Go — остальные языки интерпретируемые, + // им builder не нужен: код передаётся напрямую в deployment.literal. + if def.BuilderImage != "" { + spec["builder"] = map[string]any{ + "image": def.BuilderImage, + "command": "build", + } + } + + return &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "fission.io/v1", + "kind": "Environment", + "metadata": map[string]any{ + "name": name, + "namespace": ns, + }, + "spec": spec, + }} +} diff --git a/console/internal/fission/namespace.go b/console/internal/fission/namespace.go new file mode 100644 index 0000000..6d818f7 --- /dev/null +++ b/console/internal/fission/namespace.go @@ -0,0 +1,606 @@ +package fission + +import ( + "context" + "fmt" + "log" + "os" + "sort" + "strings" + "sync" + "time" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/dynamic" + + "encoding/json" +) + +// NSInflightEnsure — состояние in-flight вызова EnsureUserNamespace. +// Все параллельные горутины ждут close(done), затем читают err. +// Реализует ручной singleflight без внешних зависимостей. +type NSInflightEnsure struct { + done chan struct{} + err error +} + +// NSManager управляет жизненным циклом пользовательских namespace-ов: +// создание, RBAC, quota, network policies, reconciler, reaper. +type NSManager struct { + dyn dynamic.Interface + + // nsReconcileCh — сигнал для немедленного запуска NS reconciler. + // Буферизирован на 1: несколько сигналов схлопываются в один запуск. + NSReconcileCh chan struct{} + + // Singleflight + кэш для ensureUserNamespace: + // ensuredNS — namespace-ы которые уже инициализированы в этом запуске процесса. + // ensuredNSInFlight — текущие in-flight операции (ключ = namespace name). + ensuredNSMu sync.Mutex + ensuredNS map[string]struct{} + ensuredNSInFlight map[string]*NSInflightEnsure + + // nsSemaphore ограничивает параллелизм EnsureUserNamespace — не более 3 одновременно. + // Без него 10 новых пользователей генерируют 140 K8s API calls одновременно → throttle → 504. + // С семафором: 3 batch-а по 14 calls → ~3 × 5s = 15s total, все укладываются в timeout. + nsSemaphore chan struct{} +} + +// NewNSManager создаёт NSManager с инициализированными каналами и структурами. +func NewNSManager(dyn dynamic.Interface) *NSManager { + return &NSManager{ + dyn: dyn, + NSReconcileCh: make(chan struct{}, 1), + ensuredNS: make(map[string]struct{}), + ensuredNSInFlight: make(map[string]*NSInflightEnsure), + nsSemaphore: make(chan struct{}, 3), + } +} + +// EnsureUserNS — единая точка входа для гарантии существования пользовательского namespace. +// Реализует singleflight + in-memory кэш + семафор параллелизма. +// +// Singleflight: если один goroutine уже создаёт namespace ns — остальные ждут его результата +// вместо того чтобы запускать параллельные K8s API calls (вызывало throttle и 504). +// +// Кэш: если namespace уже создан в этом запуске процесса — быстрый путь без K8s calls. +func (m *NSManager) EnsureUserNS(ctx context.Context, ns string) error { + m.ensuredNSMu.Lock() + if _, ok := m.ensuredNS[ns]; ok { + // Быстрый путь: уже создан в этой жизни процесса. + m.ensuredNSMu.Unlock() + return nil + } + if inflight, ok := m.ensuredNSInFlight[ns]; ok { + // Кто-то уже создаёт — ждём его результата (singleflight). + m.ensuredNSMu.Unlock() + select { + case <-inflight.done: + return inflight.err + case <-ctx.Done(): + return ctx.Err() + } + } + // Мы первые для этого namespace. + inflight := &NSInflightEnsure{done: make(chan struct{})} + m.ensuredNSInFlight[ns] = inflight + m.ensuredNSMu.Unlock() + + // Берём слот семафора — ограничиваем параллелизм. + select { + case m.nsSemaphore <- struct{}{}: + case <-ctx.Done(): + m.ensuredNSMu.Lock() + delete(m.ensuredNSInFlight, ns) + m.ensuredNSMu.Unlock() + inflight.err = ctx.Err() + close(inflight.done) + return ctx.Err() + } + + ensureCtx, ensureCancel := context.WithTimeout(ctx, 60*time.Second) + inflight.err = m.ensureUserNamespace(ensureCtx, ns) + ensureCancel() + <-m.nsSemaphore // освобождаем слот + + m.ensuredNSMu.Lock() + delete(m.ensuredNSInFlight, ns) + if inflight.err == nil { + m.ensuredNS[ns] = struct{}{} + } + m.ensuredNSMu.Unlock() + close(inflight.done) + + return inflight.err +} + +// ensureUserNamespace выполняет реальную работу: создаёт namespace + RBAC + quota + netpol. +// Все операции идемпотентны (IsAlreadyExists игнорируется). +func (m *NSManager) ensureUserNamespace(ctx context.Context, ns string) error { + // 1. Создать namespace с меткой managed-by=fission-console. + // Метка используется reconciler-ом для фильтрации наших namespace-ов. + nsObj := &unstructured.Unstructured{ + Object: map[string]any{ + "apiVersion": "v1", + "kind": "Namespace", + "metadata": map[string]any{ + "name": ns, + "labels": map[string]any{ + "managed-by": "fission-console", + }, + }, + }, + } + _, err := m.dyn.Resource(NamespaceGVR).Create(ctx, nsObj, metav1.CreateOptions{}) + newlyCreated := err == nil + if err != nil && !apierrors.IsAlreadyExists(err) { + return fmt.Errorf("create namespace %s: %w", ns, err) + } + + // 1a. ServiceAccounts для Fission в user namespace. + // + // Fission pool pods (fetcher sidecar) запускаются в user namespace и требуют + // `serviceAccountName: fission-fetcher` в том же namespace. Executor (с SERVICEACCOUNT_CHECK_ENABLED=false) + // не создаёт эти SA автоматически → pool pods падают с "serviceaccount not found". + saGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "serviceaccounts"} + for _, saName := range []string{"fission-fetcher", "fission-builder"} { + saObj := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "v1", + "kind": "ServiceAccount", + "metadata": map[string]any{ + "name": saName, + "namespace": ns, + }, + }} + _, saErr := m.dyn.Resource(saGVR).Namespace(ns).Create(ctx, saObj, metav1.CreateOptions{}) + if saErr != nil && !apierrors.IsAlreadyExists(saErr) { + log.Printf("ensureUserNamespace: create SA %s/%s: %v", ns, saName, saErr) + } + } + + // 1b. RoleBindings для Fission SA в user namespace. + // + // Почему cluster-admin, а не admin: + // ClusterRole "admin" не включает custom resource группы (fission.io/*). + // Fission executor при старте пытается создать Role с правами на fission.io/packages, + // и получает "attempting to grant RBAC permissions not currently held" — RBAC escalation + // prevention. ClusterRole "cluster-admin" в контексте RoleBinding (не ClusterRoleBinding) + // даёт полный доступ ТОЛЬКО внутри конкретного namespace — это безопасно. + type rbSubject struct { + name string + namespace string + binding string + } + fissionSysNS := os.Getenv("FISSION_SYSTEM_NAMESPACE") + if fissionSysNS == "" { + fissionSysNS = "fission" + } + fissionSAs := []rbSubject{ + {name: "fission-executor", binding: "fission-executor-user-ns"}, + {name: "fission-router", binding: "fission-router-user-ns"}, + {name: "fission-buildermgr", binding: "fission-buildermgr-user-ns"}, + {name: "fission-kubewatcher", binding: "fission-kubewatcher-user-ns"}, + {name: "fission-timer", binding: "fission-timer-user-ns"}, + {name: "fission-fetcher", binding: "fission-fetcher-system-user-ns"}, + {name: "fission-builder", binding: "fission-builder-system-user-ns"}, + // fetcher/builder локальные в user namespace (там запускаются pool pods) + {name: "fission-fetcher", namespace: ns, binding: "fission-fetcher-local-user-ns"}, + {name: "fission-builder", namespace: ns, binding: "fission-builder-local-user-ns"}, + } + rbGVR := schema.GroupVersionResource{Group: "rbac.authorization.k8s.io", Version: "v1", Resource: "rolebindings"} + for _, sa := range fissionSAs { + subjectNS := sa.namespace + if subjectNS == "" { + subjectNS = fissionSysNS + } + rbObj := &unstructured.Unstructured{ + Object: map[string]any{ + "apiVersion": "rbac.authorization.k8s.io/v1", + "kind": "RoleBinding", + "metadata": map[string]any{ + "name": sa.binding, + "namespace": ns, + }, + "roleRef": map[string]any{ + "apiGroup": "rbac.authorization.k8s.io", + "kind": "ClusterRole", + "name": "cluster-admin", + }, + "subjects": []any{ + map[string]any{ + "kind": "ServiceAccount", + "name": sa.name, + "namespace": subjectNS, + }, + }, + }, + } + _, rbErr := m.dyn.Resource(rbGVR).Namespace(ns).Create(ctx, rbObj, metav1.CreateOptions{}) + if rbErr != nil && !apierrors.IsAlreadyExists(rbErr) { + log.Printf("ensureUserNamespace: create rolebinding %s/%s@%s: %v", ns, sa.name, subjectNS, rbErr) + } + } + + // 1c. ResourceQuota — ограничиваем ресурсы одного пользователя. + // Без этого одна функция может исчерпать CPU/RAM всего кластера. + // Параметры задаются через env vars для гибкой настройки без пересборки. + quotaGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "resourcequotas"} + quotaObj := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "v1", + "kind": "ResourceQuota", + "metadata": map[string]any{"name": "user-quota", "namespace": ns}, + "spec": map[string]any{ + "hard": map[string]any{ + "requests.cpu": envDefault("QUOTA_REQ_CPU", "1"), + "requests.memory": envDefault("QUOTA_REQ_MEM", "1Gi"), + "limits.cpu": envDefault("QUOTA_LIM_CPU", "4"), + "limits.memory": envDefault("QUOTA_LIM_MEM", "4Gi"), + "pods": envDefault("QUOTA_PODS", "30"), + "count/functions.fission.io": envDefault("QUOTA_FUNCTIONS", "20"), + "count/packages.fission.io": envDefault("QUOTA_PACKAGES", "40"), + "count/httptriggers.fission.io": envDefault("QUOTA_HTTPTRIGGERS", "20"), + }, + }, + }} + _, quotaErr := m.dyn.Resource(quotaGVR).Namespace(ns).Create(ctx, quotaObj, metav1.CreateOptions{}) + if quotaErr != nil && !apierrors.IsAlreadyExists(quotaErr) { + log.Printf("ensureUserNamespace: create ResourceQuota %s: %v", ns, quotaErr) + } + + // 1d. LimitRange — дефолтные лимиты на контейнер. + // Поды без явных limits — unbounded; LimitRange автоматически добавляет defaults. + limitRangeGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "limitranges"} + limitRangeObj := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "v1", + "kind": "LimitRange", + "metadata": map[string]any{"name": "user-limits", "namespace": ns}, + "spec": map[string]any{ + "limits": []any{ + map[string]any{ + "type": "Container", + "default": map[string]any{"cpu": envDefault("LIMIT_DEFAULT_CPU", "500m"), "memory": envDefault("LIMIT_DEFAULT_MEM", "256Mi")}, + "defaultRequest": map[string]any{ + "cpu": envDefault("LIMIT_REQ_CPU", "50m"), "memory": envDefault("LIMIT_REQ_MEM", "64Mi"), + }, + "max": map[string]any{"cpu": envDefault("LIMIT_MAX_CPU", "2"), "memory": envDefault("LIMIT_MAX_MEM", "1Gi")}, + }, + }, + }, + }} + _, lrErr := m.dyn.Resource(limitRangeGVR).Namespace(ns).Create(ctx, limitRangeObj, metav1.CreateOptions{}) + if lrErr != nil && !apierrors.IsAlreadyExists(lrErr) { + log.Printf("ensureUserNamespace: create LimitRange %s: %v", ns, lrErr) + } + + // 1e. NetworkPolicy — запрещаем входящий трафик из других user namespace. + // Разрешаем: внутри namespace, из fission core namespace, из kube-system. + // Запрещаем: трафик от подов других user namespace (межтенантная изоляция). + fissionSysNSForNetpol := os.Getenv("FISSION_SYSTEM_NAMESPACE") + if fissionSysNSForNetpol == "" { + fissionSysNSForNetpol = "fission" + } + netpolGVR := schema.GroupVersionResource{Group: "networking.k8s.io", Version: "v1", Resource: "networkpolicies"} + netpolObj := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "networking.k8s.io/v1", + "kind": "NetworkPolicy", + "metadata": map[string]any{"name": "deny-cross-tenant", "namespace": ns}, + "spec": map[string]any{ + "podSelector": map[string]any{}, // применяется ко всем подам namespace + "policyTypes": []any{"Ingress"}, + "ingress": []any{ + // Разрешаем трафик внутри namespace + map[string]any{"from": []any{map[string]any{"podSelector": map[string]any{}}}}, + // Разрешаем трафик из fission core namespace (router → function pod) + map[string]any{"from": []any{map[string]any{"namespaceSelector": map[string]any{ + "matchLabels": map[string]any{"kubernetes.io/metadata.name": fissionSysNSForNetpol}, + }}}}, + // Разрешаем трафик из kube-system (kubelet health checks, DNS) + map[string]any{"from": []any{map[string]any{"namespaceSelector": map[string]any{ + "matchLabels": map[string]any{"kubernetes.io/metadata.name": "kube-system"}, + }}}}, + }, + }, + }} + _, npErr := m.dyn.Resource(netpolGVR).Namespace(ns).Create(ctx, netpolObj, metav1.CreateOptions{}) + if npErr != nil && !apierrors.IsAlreadyExists(npErr) { + log.Printf("ensureUserNamespace: create NetworkPolicy %s: %v", ns, npErr) + } + + // 1f. Сигналим NS reconciler что появился новый namespace. + // Reconciler сам синхронизирует FISSION_RESOURCE_NAMESPACES без race condition и rolling restarts. + if newlyCreated { + select { + case m.NSReconcileCh <- struct{}{}: + default: // уже есть сигнал в буфере — не блокируем + } + } + + // Environments создаются лениво (lazy) в момент создания первой функции на языке. + return nil +} + +// StartNSReconciler запускает фоновый reconciler FISSION_RESOURCE_NAMESPACES. +// +// Проблема при scale: +// - Прямой патч на горячем пути → rolling restart всех Fission deployments при каждом новом юзере +// - Без мьютекса: 100 юзеров одновременно → read-modify-write race → namespace-ы теряются +// - Ручное kubectl delete ns → namespace остаётся в переменной вечно +// +// Решение: единственный goroutine с debounce-каналом. +// - Запускается по таймеру или немедленно через NSReconcileCh +// - 1000 юзеров одновременно → 1 патч вместо 5000 +// - Нет race condition (один goroutine, один writer) +func (m *NSManager) StartNSReconciler(interval time.Duration) { + go func() { + ticker := time.NewTicker(interval) + defer ticker.Stop() + log.Printf("nsReconciler: started, interval=%v", interval) + for { + select { + case <-ticker.C: + m.reconcileNSList() + case <-m.NSReconcileCh: + m.reconcileNSList() + // Дренируем канал чтобы не запускаться дважды подряд после сигнала + drain: + for { + select { + case <-m.NSReconcileCh: + default: + break drain + } + } + } + } + }() +} + +// reconcileNSList синхронизирует FISSION_RESOURCE_NAMESPACES с реально существующими namespace-ами. +// +// Алгоритм: +// 1. Читает все namespace-ы с меткой managed-by=fission-console из k8s (источник истины) +// 2. Читает текущий FISSION_RESOURCE_NAMESPACES из router deployment +// 3. Вычисляет desired = "default" + активные namespace-ы +// 4. Если desired == current → ничего не делает (нет rolling restart!) +// 5. Если разница → один патч всех Fission deployments +func (m *NSManager) reconcileNSList() { + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) + defer cancel() + + fissionNS := os.Getenv("FISSION_SYSTEM_NAMESPACE") + if fissionNS == "" { + fissionNS = "fission" + } + + // Шаг 1: реальные namespace-ы с нашей меткой + nsList, err := m.dyn.Resource(NamespaceGVR).List(ctx, metav1.ListOptions{ + LabelSelector: "managed-by=fission-console", + }) + if err != nil { + log.Printf("nsReconciler: list namespaces: %v", err) + return + } + + // Шаг 2: фильтруем — берём только Active. Terminating убираем (они уже умирают). + desired := map[string]struct{}{"default": {}} + for _, ns := range nsList.Items { + phase, _, _ := unstructured.NestedString(ns.Object, "status", "phase") + if phase == "Active" { + desired[ns.GetName()] = struct{}{} + } + } + + // Шаг 3: текущее значение из router deployment + routerDep, err := m.dyn.Resource(DeploymentGVR).Namespace(fissionNS).Get(ctx, "router", metav1.GetOptions{}) + if err != nil { + log.Printf("nsReconciler: get router: %v", err) + return + } + currentVal := extractFissionNSEnv(routerDep) + + // Шаг 4: сравниваем current с desired + currentSet := map[string]struct{}{} + for _, p := range strings.Split(currentVal, ",") { + if t := strings.TrimSpace(p); t != "" { + currentSet[t] = struct{}{} + } + } + same := len(currentSet) == len(desired) + if same { + for k := range desired { + if _, ok := currentSet[k]; !ok { + same = false + break + } + } + } + if same { + return // ничего менять не нужно — нет патча, нет rolling restart + } + + // Шаг 5: патчим все Fission deployments + parts := make([]string, 0, len(desired)) + for ns := range desired { + parts = append(parts, ns) + } + sort.Strings(parts) + newVal := strings.Join(parts, ",") + + fissionDeployments := []string{"router", "executor", "buildermgr", "kubewatcher", "timer"} + patch := map[string]any{ + "spec": map[string]any{ + "template": map[string]any{ + "spec": map[string]any{ + "containers": []any{ + map[string]any{ + "name": "", + "env": []any{ + map[string]any{ + "name": "FISSION_RESOURCE_NAMESPACES", + "value": newVal, + }, + }, + }, + }, + }, + }, + }, + } + for _, dep := range fissionDeployments { + patch["spec"].(map[string]any)["template"].(map[string]any)["spec"].(map[string]any)["containers"].([]any)[0].(map[string]any)["name"] = dep + patchBytes, _ := json.Marshal(patch) + if _, pErr := m.dyn.Resource(DeploymentGVR).Namespace(fissionNS).Patch( + ctx, dep, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{}); pErr != nil { + log.Printf("nsReconciler: patch %s: %v", dep, pErr) + } + } + log.Printf("nsReconciler: synced FISSION_RESOURCE_NAMESPACES: %s → %s", currentVal, newVal) +} + +// extractFissionNSEnv читает FISSION_RESOURCE_NAMESPACES из первого контейнера deployment. +func extractFissionNSEnv(dep *unstructured.Unstructured) string { + containers, _, _ := unstructured.NestedSlice(dep.Object, "spec", "template", "spec", "containers") + for _, c := range containers { + cont, ok := c.(map[string]any) + if !ok { + continue + } + envs, _, _ := unstructured.NestedSlice(cont, "env") + for _, e := range envs { + env, ok := e.(map[string]any) + if !ok { + continue + } + if env["name"] == "FISSION_RESOURCE_NAMESPACES" { + if v, ok := env["value"].(string); ok && v != "" { + return v + } + } + } + break // смотрим только первый контейнер + } + return "default" +} + +// StartExpiryReaper запускает фоновый goroutine для удаления функций с истёкшим TTL. +func (m *NSManager) StartExpiryReaper(interval time.Duration) { + go func() { + ticker := time.NewTicker(interval) + defer ticker.Stop() + log.Printf("expiryReaper: started, interval=%v", interval) + for range ticker.C { + m.runExpiryReap() + } + }() +} + +// runExpiryReap обходит все управляемые namespace-ы и удаляет протухшие функции. +func (m *NSManager) runExpiryReap() { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + + nsList, err := m.dyn.Resource(NamespaceGVR).List(ctx, metav1.ListOptions{ + LabelSelector: "managed-by=fission-console", + }) + if err != nil { + log.Printf("expiryReaper: list namespaces: %v", err) + return + } + + now := time.Now().UTC() + for _, ns := range nsList.Items { + m.reapExpiredFunctionsInNS(ctx, ns.GetName(), now) + } +} + +// reapExpiredFunctionsInNS удаляет протухшие функции в конкретном namespace. +// Для каждой удалённой функции вызывает CleanupEnvironmentIfUnused. +// Также удаляет orphan packages — пакеты без соответствующей функции. +func (m *NSManager) reapExpiredFunctionsInNS(ctx context.Context, ns string, now time.Time) { + functions, err := m.dyn.Resource(FunctionGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if err != nil { + log.Printf("expiryReaper: list functions in %s: %v", ns, err) + return + } + + // Строим множество активных функций для поиска orphan packages + activeFunctions := make(map[string]struct{}, len(functions.Items)) + for _, fn := range functions.Items { + activeFunctions[fn.GetName()] = struct{}{} + } + + for _, fn := range functions.Items { + expiresAtStr, _, _ := unstructured.NestedString(fn.Object, "metadata", "annotations", "fission-console/expires-at") + if expiresAtStr == "" { + continue // нет TTL — функция живёт вечно + } + expiresAt, parseErr := time.Parse(time.RFC3339, expiresAtStr) + if parseErr != nil { + log.Printf("expiryReaper: parse expires-at for %s/%s: %v", ns, fn.GetName(), parseErr) + continue + } + if now.Before(expiresAt) { + continue // ещё не протухла + } + + fnName := fn.GetName() + envName, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") + pkgName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name") + log.Printf("expiryReaper: deleting expired function %s/%s (expired %s ago)", ns, fnName, now.Sub(expiresAt).Round(time.Second)) + + // Удаляем HTTPTrigger-ы ссылающиеся на эту функцию + triggers, tErr := m.dyn.Resource(HTTPTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if tErr == nil { + for _, trig := range triggers.Items { + refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name") + if refName == fnName { + _ = m.dyn.Resource(HTTPTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{}) + } + } + } + + _ = m.dyn.Resource(FunctionGVR).Namespace(ns).Delete(ctx, fnName, metav1.DeleteOptions{}) + if pkgName != "" { + _ = m.dyn.Resource(PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) + } + if envName != "" { + CleanupEnvironmentIfUnused(ctx, m.dyn, ns, envName) + } + } + + // Сигналим reconciler — он проверит все namespace-ы и почистит список + select { + case m.NSReconcileCh <- struct{}{}: + default: + } + + // Orphan packages — пакеты без соответствующей функции + // (остаются если процесс упал в середине удаления) + packages, pkgListErr := m.dyn.Resource(PackageGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if pkgListErr == nil { + for _, pkg := range packages.Items { + pkgName := pkg.GetName() + // Конвенция именования: {fn-name}-pkg + if !strings.HasSuffix(pkgName, "-pkg") { + continue + } + fnName := strings.TrimSuffix(pkgName, "-pkg") + if _, exists := activeFunctions[fnName]; !exists { + log.Printf("expiryReaper: deleting orphan package %s/%s (no matching function)", ns, pkgName) + _ = m.dyn.Resource(PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) + } + } + } +} + +// envDefault читает env переменную или возвращает fallback. +func envDefault(key, fallback string) string { + if v := strings.TrimSpace(os.Getenv(key)); v != "" { + return v + } + return fallback +} diff --git a/console/internal/model/types.go b/console/internal/model/types.go new file mode 100644 index 0000000..793c19e --- /dev/null +++ b/console/internal/model/types.go @@ -0,0 +1,46 @@ +// Package model содержит общие типы данных (request/response структуры, +// конфигурации языков), которые используются в нескольких пакетах. +// Вынесено отдельно чтобы избежать циклических зависимостей между пакетами. +package model + +// CreateFunctionRequest — тело POST /console/api/functions. +// TTL пустой → функция живёт вечно; "1d", "24h" — протухнет через указанное время. +type CreateFunctionRequest struct { + Name string `json:"name"` + Language string `json:"language"` + Environment string `json:"environment"` + Code string `json:"code"` + Entrypoint string `json:"entrypoint"` + Route string `json:"route"` + Methods []string `json:"methods"` + TTL string `json:"ttl"` // e.g. "24h", "7d" — пустое = функция не протухает +} + +// UpdateCodeRequest — тело PUT /console/api/functions/:name/code. +type UpdateCodeRequest struct { + Code string `json:"code"` +} + +// LangEnvDef описывает Docker-образы для конкретного языка. +// BuilderImage заполнен только для языков которым нужна компиляция (Go). +// Version: версия среды Fission (по умолчанию 3; perl использует 1 — нет поддержки async entrypoint). +type LangEnvDef struct { + Image string + BuilderImage string + Version int // 0 → по умолчанию 3 +} + +// LangEnvMap сопоставляет идентификатор языка (string) с описанием среды выполнения. +// Ключ используется в createFunctionRequest.Language и как суффикс имени Environment. +// +// Почему perl Version=1: +// Fission Environment v3 требует поддержки async entrypoint (V2 protocol). +// ghcr.io/fission/perl-env поддерживает только V1 protocol → version=1. +var LangEnvMap = map[string]LangEnvDef{ + "python": {Image: "ghcr.io/fission/python-env"}, + "nodejs": {Image: "ghcr.io/fission/node-env"}, + "go": {Image: "ghcr.io/fission/go-env", BuilderImage: "naeel/go-builder-fast:v1"}, + "php": {Image: "ghcr.io/fission/php-env"}, + "ruby": {Image: "ghcr.io/fission/ruby-env"}, + "perl": {Image: "ghcr.io/fission/perl-env", Version: 1}, +} diff --git a/console/internal/runtime/doc.go b/console/internal/runtime/doc.go new file mode 100644 index 0000000..96615b3 --- /dev/null +++ b/console/internal/runtime/doc.go @@ -0,0 +1,4 @@ +// Package runtime отвечает за подготовку исходного кода функции +// к загрузке в Fission: упаковка в zip, обёртка ESM для Node.js и т.д. +// Каждый язык — отдельный файл; общая логика zip — в helpers.go. +package runtime diff --git a/console/internal/runtime/entrypoint.go b/console/internal/runtime/entrypoint.go new file mode 100644 index 0000000..b88ef77 --- /dev/null +++ b/console/internal/runtime/entrypoint.go @@ -0,0 +1,30 @@ +package runtime + +// DefaultEntrypoint возвращает точку входа по умолчанию для языка lang. +// +// Формат entrypoint зависит от среды выполнения Fission: +// - python: "module.function" → Fission импортирует модуль и вызывает функцию +// - nodejs: "main" → имя экспортированной default-функции в main.js +// - php: "main.php::handler" → имя файла + "::" + имя функции +// - ruby: "handler" → имя метода в загруженном файле +// - perl: "handler" → имя функции в загруженном файле +// - go: "Handler" → экспортированная Go-функция (с большой буквы) +func DefaultEntrypoint(lang string) string { + switch lang { + case "nodejs": + return "main" + case "php": + // Fission php-env: filename::functionName + return "main.php::handler" + case "ruby": + return "handler" + case "perl": + return "handler" + case "go": + // Go: экспортированная функция (заглавная) — go/plugin требует экспорт + return "Handler" + default: + // python и неизвестные языки: "main.main" = модуль main, функция main + return "main.main" + } +} diff --git a/console/internal/runtime/go.go b/console/internal/runtime/go.go new file mode 100644 index 0000000..16a59e9 --- /dev/null +++ b/console/internal/runtime/go.go @@ -0,0 +1,43 @@ +package runtime + +import ( + "archive/zip" + "bytes" +) + +// BuildGoSourceZip упаковывает исходный Go-код в zip для Fission source package. +// +// Fission Go-среда работает в два этапа: +// 1. Package создаётся с source.literal = zip(handler.go + go.mod) +// 2. go-builder job компилирует его в плагин (.so) — это занимает ~40s первый раз +// +// Важно: source package, а НЕ deployment package — Fission сам компилирует. +// Entrypoint по умолчанию "Handler" (экспортированная функция в handler.go). +func BuildGoSourceZip(code string) ([]byte, error) { + var buf bytes.Buffer + zw := zip.NewWriter(&buf) + + // handler.go — пользовательский код + fw, err := zw.Create("handler.go") + if err != nil { + return nil, err + } + if _, err := fw.Write([]byte(code)); err != nil { + return nil, err + } + + // go.mod — минимальный модуль для сборки; go-builder добавит нужные зависимости + goMod := "module github.com/user/fn\n\ngo 1.23\n" + fw2, err := zw.Create("go.mod") + if err != nil { + return nil, err + } + if _, err := fw2.Write([]byte(goMod)); err != nil { + return nil, err + } + + if err := zw.Close(); err != nil { + return nil, err + } + return buf.Bytes(), nil +} diff --git a/console/internal/runtime/helpers.go b/console/internal/runtime/helpers.go new file mode 100644 index 0000000..d813b6b --- /dev/null +++ b/console/internal/runtime/helpers.go @@ -0,0 +1,27 @@ +package runtime + +import ( + "archive/zip" + "bytes" +) + +// buildZip создаёт zip-архив с одним файлом fileName и содержимым content. +// Вспомогательная функция: используется в buildScriptZip (PHP/Ruby/Perl), +// а также как основа для buildGoSourceZip и buildJSDeployZip. +func buildZip(fileName string, content []byte) ([]byte, error) { + var buf bytes.Buffer + zw := zip.NewWriter(&buf) + + fw, err := zw.Create(fileName) + if err != nil { + return nil, err + } + if _, err := fw.Write(content); err != nil { + return nil, err + } + + if err := zw.Close(); err != nil { + return nil, err + } + return buf.Bytes(), nil +} diff --git a/console/internal/runtime/nodejs.go b/console/internal/runtime/nodejs.go new file mode 100644 index 0000000..6a30c67 --- /dev/null +++ b/console/internal/runtime/nodejs.go @@ -0,0 +1,75 @@ +package runtime + +import ( + "archive/zip" + "bytes" + "encoding/json" + "fmt" +) + +// BuildJSDeployZip упаковывает Node.js код в zip для Fission deployment package. +// +// Почему zip, а не просто literal-код: +// Fission node-env v3 поддерживает ESM модули. Чтобы Node.js трактовал файл +// как ESM, нужен package.json с {"type":"module"} рядом с main.js. +// Без него `import` синтаксис вызовет "require is not defined" или SyntaxError. +// +// Почему new Function: +// Пользовательский код пишет CJS-стиль (module.exports = ...) но мы хотим ESM wrapper. +// new Function() изолирует module/exports от глобального ESM контекста и позволяет +// исполнять CJS-код без изменений. +// +// Результат: zip с двумя файлами: +// - package.json: {"type":"module"} +// - main.js: ESM wrapper + инлайн пользовательский код через new Function +func BuildJSDeployZip(code string) ([]byte, error) { + // Сериализуем пользовательский код в JSON строку чтобы безопасно инлайнить + // в JavaScript-литерал — экранирует кавычки, переводы строк, спецсимволы. + codeJSON, err := json.Marshal(code) + if err != nil { + return nil, fmt.Errorf("marshal user code: %w", err) + } + + // main.js — ESM wrapper: + // 1. new Function создаёт функцию в пустом модульном контексте (нет import/export) + // 2. Передаём ей module и exports как параметры → пользовательский CJS код работает + // 3. Экспортируем default async function для Fission node-env entrypoint + wrapper := fmt.Sprintf(`const __mod = { exports: {} }; +(new Function('module', 'exports', %s))(__mod, __mod.exports); +const _fn = __mod.exports; + +export default async function(ctx) { + const fn = typeof _fn === 'function' ? _fn : (_fn.default || _fn.handler || _fn.main); + if (!fn) throw new Error('no exported function found in user code'); + const result = await fn(ctx); + if (!result) return { status: 200, body: '' }; + if (typeof result.status !== 'undefined') return result; + return { status: 200, ...result }; +} +`, string(codeJSON)) + + var buf bytes.Buffer + zw := zip.NewWriter(&buf) + + // package.json: включает ESM режим для всего архива + pkgfw, err := zw.Create("package.json") + if err != nil { + return nil, err + } + if _, err := pkgfw.Write([]byte(`{"type":"module"}`)); err != nil { + return nil, err + } + + fw, err := zw.Create("main.js") + if err != nil { + return nil, err + } + if _, err := fw.Write([]byte(wrapper)); err != nil { + return nil, err + } + + if err := zw.Close(); err != nil { + return nil, err + } + return buf.Bytes(), nil +} diff --git a/console/internal/runtime/script.go b/console/internal/runtime/script.go new file mode 100644 index 0000000..af26c6a --- /dev/null +++ b/console/internal/runtime/script.go @@ -0,0 +1,13 @@ +package runtime + +// BuildScriptZip создаёт zip-архив с одним файлом fileName и содержимым code. +// Используется для PHP, Ruby, Perl — языков где среда Fission ожидает +// именованный файл (handler.rb, handler.pl, main.php и т.д.) внутри архива. +// +// Почему zip, а не просто literal: +// Fission poolmgr при специализации пода распаковывает deployment archive, +// находит нужный файл по имени и загружает его в среду выполнения. +// Если передать просто байты кода в deployment.literal — среда не знает расширение. +func BuildScriptZip(code, fileName string) ([]byte, error) { + return buildZip(fileName, []byte(code)) +}