package main import ( "archive/zip" "bytes" "context" "crypto/sha256" "encoding/base64" "encoding/hex" "encoding/json" "errors" "fmt" "go/parser" "go/token" "io" "log" "net" "net/http" "os" "os/exec" "regexp" "sort" "strconv" "strings" "sync" "time" "unicode/utf8" "fission-console/ui" 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" "k8s.io/client-go/rest" "k8s.io/client-go/tools/clientcmd" ) 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"} ) const defaultSATokenPath = "/var/run/secrets/kubernetes.io/serviceaccount/token" 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", } // nsInflightEnsure — состояние in-flight вызова ensureUserNamespace. // Все параллельные горутины ждут close(done), затем читают err. type nsInflightEnsure struct { done chan struct{} err error } type server struct { dyn dynamic.Interface ns string routerURL string http *http.Client saTokenPath string invokeTimeout time.Duration testMode bool // FISSION_TEST_MODE=true — пропускает deck auth, X-Test-Sub задаёт user authUser string authPass string // --- ai/ask feature (удалить блок целиком чтобы выкосить) --- llmURL string // FISSION_LLM_URL llmKey string // FISSION_LLM_KEY // --- end ai/ask feature --- tokenMu sync.Mutex cachedJWT string tokenExpAt time.Time tokenCache sync.Map // nsReconcileCh — сигнал для немедленного запуска NS reconciler. // Буферизирован на 1: несколько сигналов схлопываются в один запуск. nsReconcileCh chan struct{} // ensuredNS — кэш namespace-ов для которых уже отработал ensureUserNamespace. // ensuredNSMu защищает ensuredNS и ensuredNSInFlight. // ensuredNSInFlight — ожидание: если namespace создаётся прямо сейчас, параллельные // запросы ждут завершения (ручной singleflight без внешних зависимостей). ensuredNSMu sync.Mutex ensuredNS map[string]struct{} ensuredNSInFlight map[string]*nsInflightEnsure // nsSemaphore ограничивает параллелизм ensureUserNamespace — не более 3 одновременно. // Без него 10 новых пользователей генерируют 140 K8s API calls одновременно → throttle → 504. nsSemaphore chan struct{} } 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" — пустое = функция не протухает } type langEnvDef struct { Image string BuilderImage string Version int // 0 defaults to 3 } 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}, } type updateCodeRequest struct { Code string `json:"code"` } 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) 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) } authUser := envDefault("FISSION_AUTH_USERNAME", "") authPass := envDefault("FISSION_AUTH_PASSWORD", "") saTokenPath := envDefault("SA_TOKEN_PATH", defaultSATokenPath) s := &server{ dyn: dyn, ns: namespace, routerURL: routerURL, http: &http.Client{Timeout: httpTimeout}, saTokenPath: saTokenPath, invokeTimeout: invokeTimeout, authUser: authUser, authPass: authPass, 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 --- nsReconcileCh: make(chan struct{}, 1), ensuredNS: make(map[string]struct{}), ensuredNSInFlight: make(map[string]*nsInflightEnsure), nsSemaphore: make(chan struct{}, 3), } mux := http.NewServeMux() 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")) }) uiHandler := http.StripPrefix("/console", ui.Handler()) mux.Handle("/console", uiHandler) mux.Handle("/console/", uiHandler) mux.HandleFunc("/api/environments", s.handleList(environmentGVR)) mux.HandleFunc("/api/packages", s.handleList(packageGVR)) mux.HandleFunc("/api/functions", s.handleFunctionsRoot) mux.HandleFunc("/api/functions/", s.handleFunctionsAction) mux.HandleFunc("/api/httptriggers", s.handleList(httpTrigGVR)) mux.HandleFunc("/api/timetriggers", s.handleList(timeTrigGVR)) auth := func(h http.HandlerFunc) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var ns string if s.testMode { // TEST_MODE: deck API не вызывается. // X-Test-Sub или X-Auth-Token задают sub → разные namespace-ы для тестирования. 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 { log.Printf("namespaceFromJWT: %v", err) ns = s.ns } } ctx := context.WithValue(r.Context(), ctxKeyNS{}, ns) // Гарантируем что namespace + RBAC + quota + netpol существуют. // ensureUserNS реализует singleflight + кэш + семафор параллелизма. if ensureErr := s.ensureUserNS(ctx, ns); ensureErr != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", ensureErr)) return } h(w, r.WithContext(ctx)) } } mux.HandleFunc("/console/api/auth", s.handleAuth) mux.HandleFunc("/console/api/environments", auth(s.handleList(environmentGVR))) mux.HandleFunc("/console/api/packages", auth(s.handleList(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(httpTrigGVR))) mux.HandleFunc("/console/api/timetriggers", auth(s.handleList(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 --- httpServer := &http.Server{ Addr: ":" + port, Handler: withSecurityHeaders(withCORS(logRequests(mux))), ReadHeaderTimeout: 10 * time.Second, } log.Printf("fission-console listening on :%s (namespace=%s)", port, namespace) s.startExpiryReaper(envDurationDefault("REAPER_INTERVAL", 5*time.Minute)) s.startNSReconciler(envDurationDefault("NS_RECONCILE_INTERVAL", 2*time.Minute)) log.Fatal(httpServer.ListenAndServe()) } func (s *server) handleFunctionsRoot(w http.ResponseWriter, r *http.Request) { switch r.Method { case http.MethodGet: s.handleList(functionGVR)(w, r) case http.MethodPost: s.handleCreateFunction(w, r) default: http.Error(w, "method not allowed", http.StatusMethodNotAllowed) } } func (s *server) handleFunctionsAction(w http.ResponseWriter, r *http.Request) { 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 { 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 } if len(parts) == 2 && parts[1] == "code" && r.Method == http.MethodPut { s.handleUpdateFunctionCode(w, r, name) return } if len(parts) == 2 && parts[1] == "invoke" && r.Method == http.MethodPost { s.handleInvokeFunction(w, r, name) return } http.NotFound(w, r) } // validFuncName — RFC 1123 subdomain label: lowercase alphanumeric + hyphens, no leading/trailing hyphen, max 63 chars. var validFuncName = regexp.MustCompile(`^[a-z0-9]([a-z0-9-]*[a-z0-9])?$`) // maxCodeSize — максимальный размер кода функции (1 MB). const maxCodeSize = 1 << 20 func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { var req createFunctionRequest ns := s.userNS(r) if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err)) return } // Гарантируем что namespace + RBAC созданы до любых операций с ресурсами. // handleAuth делает это при логине, но в test mode или при прямом вызове API // namespace может отсутствовать — создаём idempotent через кэш+singleflight+семафор. nsCtx, nsCancel := context.WithTimeout(r.Context(), 60*time.Second) defer nsCancel() if err := s.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) // Валидация имени: RFC 1123 label, максимум 57 символов. // Ограничение 57 (не 63): самый длинный суффикс — "-route" (HTTPTrigger) = 6 символов. // 63 - 6 = 57. Fission webhook требует все объекты <= 63 символов. 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 } // Resolve language → environment (lazy creation). // ensureEnvironment создаёт Environment CRD если не существует — Fission увидит и поднимет pool pod. if req.Language != "" { envCtx, envCancel := context.WithTimeout(r.Context(), 15*time.Second) defer envCancel() envName, err := s.ensureEnvironment(envCtx, ns, req.Language) if err != nil { // unsupported language — это клиентская ошибка → 400 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 = defaultEntrypoint(req.Language) } if req.Route == "" { // namespace-prefix route to avoid collisions between users nsShort := s.userNS(r) 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() if _, err := s.dyn.Resource(environmentGVR).Namespace(s.userNS(r)).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) } // Build the package spec: Go uses source archive (builder), others use literal deployment var pkgSpec map[string]any if req.Language == "go" { srcZip, err := s.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 := 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) pkgSpec = map[string]any{ "deployment": map[string]any{ "type": "literal", "literal": literal, }, "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{} fnAnnotations["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(packageGVR).Namespace(s.userNS(r)).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": s.userNS(r)}, "functionName": req.Entrypoint, }, }, }} if _, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Create(ctx, fn, metav1.CreateOptions{}); err != nil { _ = s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).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(httpTrigGVR).Namespace(s.userNS(r)).Create(ctx, httpTrigger, metav1.CreateOptions{}); err != nil { _ = s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Delete(ctx, req.Name, metav1.DeleteOptions{}) _ = s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).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"], }) } func (s *server) buildLangEnvironment(name string, ns string, def langEnvDef) *unstructured.Unstructured { 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), } 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, }} } // ensureEnvironment создаёт Environment CRD для языка lang в namespace ns если не существует. // Вся Fission-специфичная конструкция объекта изолирована в buildLangEnvironment — // при обновлении Fission (изменение схемы CRD) меняем только там. func (s *server) ensureEnvironment(ctx context.Context, ns, lang string) (string, error) { langDef, ok := langEnvMap[lang] if !ok { return "", fmt.Errorf("unsupported language: %q", lang) } envName := "console-" + lang + "-env" _, getErr := s.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) } env := s.buildLangEnvironment(envName, ns, langDef) if _, createErr := s.dyn.Resource(environmentGVR).Namespace(ns).Create(ctx, env, metav1.CreateOptions{}); createErr != nil && !apierrors.IsAlreadyExists(createErr) { return "", fmt.Errorf("create environment %q: %w", envName, createErr) } log.Printf("ensureEnvironment: created %s/%s", ns, envName) return envName, nil } // cleanupEnvironmentIfUnused удаляет Environment CRD если ни одна функция в namespace его не использует. // Fission увидит удаление → уберёт pool Deployment → поды умирают. // Не блокирующий: ошибки логируются, не возвращаются вызывающему коду. func (s *server) cleanupEnvironmentIfUnused(ctx context.Context, ns, envName string) { functions, err := s.dyn.Resource(functionGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) if err != nil { log.Printf("cleanupEnvironmentIfUnused: list functions in %s: %v", ns, err) return } for _, fn := range functions.Items { fnEnv, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") if fnEnv == envName { return // ещё используется хотя бы одной функцией } } // Ни одна функция не ссылается на этот environment — удаляем if delErr := s.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) } // 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 } // startExpiryReaper запускает фоновый goroutine для удаления функций с истёкшим TTL. // Интервал задаётся через env REAPER_INTERVAL (default 5m). // При удалении функции вызывает cleanupEnvironmentIfUnused — поды умирают когда language больше не используется. func (s *server) startExpiryReaper(interval time.Duration) { go func() { ticker := time.NewTicker(interval) defer ticker.Stop() log.Printf("expiryReaper: started, interval=%v", interval) for range ticker.C { s.runExpiryReap() } }() } // runExpiryReap обходит все namespace под управлением fission-console и удаляет протухшие функции. func (s *server) runExpiryReap() { ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) defer cancel() // Находим только наши namespace-ы по метке которую мы ставим при создании nsList, err := s.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 { s.reapExpiredFunctionsInNS(ctx, ns.GetName(), now) } } // reapExpiredFunctionsInNS удаляет протухшие функции в конкретном namespace. // Для каждой удалённой функции вызывает cleanupEnvironmentIfUnused. // Также удаляет orphan packages — пакеты у которых нет соответствующей функции. func (s *server) reapExpiredFunctionsInNS(ctx context.Context, ns string, now time.Time) { functions, err := s.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 := s.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 { _ = s.dyn.Resource(httpTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{}) } } } _ = s.dyn.Resource(functionGVR).Namespace(ns).Delete(ctx, fnName, metav1.DeleteOptions{}) if pkgName != "" { _ = s.dyn.Resource(packageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) } // Убираем pool pod если язык больше не используется if envName != "" { s.cleanupEnvironmentIfUnused(ctx, ns, envName) } } // Сигналим reconciler — он проверит все namespace'ы и почистит список. select { case s.nsReconcileCh <- struct{}{}: default: } // Сканируем orphan packages — пакеты без соответствующей функции // (могут остаться если под упал в середине удаления) packages, pkgListErr := s.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) _ = s.dyn.Resource(packageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) } } } } func (s *server) buildGoSourceZip(code string) ([]byte, error) { var buf bytes.Buffer zw := zip.NewWriter(&buf) fw, err := zw.Create("handler.go") if err != nil { return nil, err } if _, err := fw.Write([]byte(code)); err != nil { return nil, err } 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 } func buildJSDeployZip(code string) ([]byte, error) { var buf bytes.Buffer zw := zip.NewWriter(&buf) // package.json: объявляем ESM тип чтобы Node.js трактовал .js как 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 } // main.js — ESM wrapper + код пользователя инлайн через new Function // new Function безопасно изолирует module/exports от глобального контекста codeJSON, err := json.Marshal(code) if err != nil { return nil, fmt.Errorf("marshal user code: %w", err) } 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)) 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 } // buildScriptZip wraps code into a zip file with the given filename. // Used for PHP, Ruby, Perl where the environment requires a named source file. func buildScriptZip(code, filename string) ([]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([]byte(code)); err != nil { return nil, err } if err := zw.Close(); err != nil { return nil, err } return buf.Bytes(), nil } func defaultEntrypoint(lang string) string { switch lang { case "nodejs": return "main" case "php": return "main.php::handler" case "ruby": return "handler" case "perl": return "handler" case "go": return "Handler" default: return "main.main" } } func (s *server) handleGetFunction(w http.ResponseWriter, r *http.Request, name string) { ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) defer cancel() fn, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).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") code := "" if packageName != "" { pkg, pkgErr := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Get(ctx, packageName, metav1.GetOptions{}) if pkgErr == nil { code = s.extractPackageSourceCode(ctx, pkg) } } route := "" methods := []string{} triggers, trigErr := s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).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": s.userNS(r), "environment": environment, "package": packageName, "entrypoint": entrypoint, "code": code, "route": route, "methods": methods, "raw": fn.Object, }) } func (s *server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request, name string) { var req 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() fn, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).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(packageGVR).Namespace(s.userNS(r)).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 := 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 } if _, err := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Update(ctx, pkg, metav1.UpdateOptions{}); err != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("update package %q: %v", pkgName, err)) return } updatedPkg, err := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Get(ctx, pkgName, metav1.GetOptions{}) if err != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get updated package %q: %v", pkgName, err)) return } 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(functionGVR).Namespace(s.userNS(r)).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()}) } 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() // check function exists in k8s before invoke if _, err2 := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).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 } invokeURL := fmt.Sprintf("%s/fission-function/v2/functions/%s", s.routerURL, name) invokeMethod := http.MethodPost triggers, err := s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).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 := false hasGet := false for _, method := range methods { m := strings.ToUpper(strings.TrimSpace(method)) if m == http.MethodPost { hasPost = true } if m == http.MethodGet { hasGet = true } } if route != "" { if !strings.HasPrefix(route, "/") { route = "/" + route } invokeURL = s.routerURL + route 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 { 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), }) } func (s *server) readSAToken() string { if s.saTokenPath == "" { return "" } data, err := os.ReadFile(s.saTokenPath) if err != nil { return "" } return strings.TrimSpace(string(data)) } 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 { log.Printf("router login failed: %v", err) return s.readSAToken() } defer resp.Body.Close() if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated { respBody, _ := io.ReadAll(resp.Body) log.Printf("router login %d: %s", resp.StatusCode, string(respBody)) return s.readSAToken() } var result struct { AccessToken string `json:"accesstoken"` } if err := json.NewDecoder(resp.Body).Decode(&result); err != nil || result.AccessToken == "" { log.Printf("router login decode error: %v", err) return s.readSAToken() } s.cachedJWT = result.AccessToken s.tokenExpAt = time.Now().Add(100 * time.Second) log.Printf("router JWT obtained, expires in 100s") return s.cachedJWT } // --- namespace isolation helpers --- type ctxKeyNS struct{} // userNS возвращает namespace пользователя из контекста запроса, или fallback (s.userNS(r)). func (s *server) userNS(r *http.Request) string { if ns, ok := r.Context().Value(ctxKeyNS{}).(string); ok && ns != "" { return ns } return s.userNS(r) } // namespaceFromJWT декодирует JWT payload (без верификации подписи), // извлекает claim "sub" и вычисляет namespace: "fission-" + hex(SHA256(sub)[:8]). func namespaceFromJWT(token string) (string, error) { parts := strings.SplitN(token, ".", 3) if len(parts) != 3 { return "", fmt.Errorf("invalid JWT format") } payload := parts[1] // base64url без padding switch len(payload) % 4 { case 2: payload += "==" case 3: payload += "=" } 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 } // ensureUserNamespace создаёт K8s namespace и shared environments если не существуют. // После создания сигналит nsReconciler который синхронизирует FISSION_RESOURCE_NAMESPACES. // // Зачем reconciler, а не прямой патч: // - Прямой патч на горячем пути → rolling restart всех Fission deployments при каждом новом юзере // - Race condition: 100 юзеров одновременно → read-modify-write без мьютекса → namespace'ы теряются // - Reconciler батчит изменения, нет race condition, нет лишних restarts func (s *server) addNSToFission(_ context.Context, _ string) error { return nil } // заменено reconciler'ом // startNSReconciler запускает фоновый reconciler FISSION_RESOURCE_NAMESPACES. // // Проблема при scale: // - addNSToFission на горячем пути (каждый новый юзер) → патч 5 деплоев → rolling restart → cold start для всех // - Без мьютекса: 100 юзеров одновременно → read-modify-write race → namespace'ы теряются // - Ручное kubectl delete ns → namespace остаётся в переменной вечно → executor флудит RBAC ошибками // // Решение: единственный goroutine с debounce-каналом. // - Запускается по таймеру (каждые NS_RECONCILE_INTERVAL) или немедленно через nsReconcileCh // - 1000 юзеров создают namespace'ы одновременно → 1 патч вместо 5000 // - Нет race condition (один goroutine, один writer) // - Автоматически чистит ghost namespace'ы (удалённые kubectl delete ns или reaper'ом) func (s *server) 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: s.reconcileNSList() case <-s.nsReconcileCh: // Немедленный запуск (новый namespace или удаление функции). // Дренируем канал чтобы не запускаться дважды подряд. s.reconcileNSList() drain: for { select { case <-s.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 (s *server) 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 := s.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 namespace'ы убираем из списка (они уже умирают). 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 (источник истины) routerDep, err := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Get(ctx, "router", metav1.GetOptions{}) if err != nil { log.Printf("nsReconciler: get router: %v", err) return } currentVal := "default" containers, _, _ := unstructured.NestedSlice(routerDep.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 != "" { currentVal = v } } } break } // Шаг 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 := s.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) } // ensureUserNS — единая точка входа для гарантии существования пользовательского namespace. // Реализует singleflight + in-memory кэш + семафор параллелизма. // // Singleflight: если один goroutine уже создаёт namespace ns — остальные ждут его результата // вместо того чтобы запускать параллельные K8s API calls (вызывало throttle и 504). // // Кэш: если namespace уже создан в этом запуске процесса — быстрый путь без K8s calls. // // Семафор (3 слота): не более 3 namespace-ов создаются одновременно. // 10 новых пользователей × 14 K8s calls = 140 calls без семафора → throttle → 60s+ → 504. // С семафором: 3 batch-а по 14 calls → ~3 × 5s = 15s total, все укладываются в timeout. func (s *server) ensureUserNS(ctx context.Context, ns string) error { s.ensuredNSMu.Lock() if _, ok := s.ensuredNS[ns]; ok { // Быстрый путь: уже создан в этой жизни процесса. s.ensuredNSMu.Unlock() return nil } if inflight, ok := s.ensuredNSInFlight[ns]; ok { // Кто-то уже создаёт — ждём его результата. s.ensuredNSMu.Unlock() select { case <-inflight.done: return inflight.err case <-ctx.Done(): return ctx.Err() } } // Мы первые для этого namespace. inflight := &nsInflightEnsure{done: make(chan struct{})} s.ensuredNSInFlight[ns] = inflight s.ensuredNSMu.Unlock() // Берём слот семафора — ограничиваем параллелизм. select { case s.nsSemaphore <- struct{}{}: case <-ctx.Done(): s.ensuredNSMu.Lock() delete(s.ensuredNSInFlight, ns) s.ensuredNSMu.Unlock() inflight.err = ctx.Err() close(inflight.done) return ctx.Err() } ensureCtx, ensureCancel := context.WithTimeout(ctx, 60*time.Second) inflight.err = s.ensureUserNamespace(ensureCtx, ns) ensureCancel() <-s.nsSemaphore // освобождаем слот s.ensuredNSMu.Lock() delete(s.ensuredNSInFlight, ns) if inflight.err == nil { s.ensuredNS[ns] = struct{}{} } s.ensuredNSMu.Unlock() close(inflight.done) return inflight.err } func (s *server) ensureUserNamespace(ctx context.Context, ns string) error { // 1. Создать 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 := s.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". // Создаём явно при каждом вызове (idempotent через IsAlreadyExists). 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 := s.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. // // Проблема: Fission компоненты (executor, router, buildermgr и др.) работают в namespace // "fission", но при добавлении нового namespace в FISSION_RESOURCE_NAMESPACES они начинают // туда смотреть (list/watch). По умолчанию у их SA нет прав в чужих namespace-ах → // "forbidden: cannot list environments.fission.io in namespace X". // // Почему 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 — это безопасно. // // Операция idempotent: если RoleBinding уже существует — IsAlreadyExists игнорируется. // Важно: pool pod запускается с serviceAccountName=fission-fetcher в самом user namespace, // а не в system namespace "fission". Поэтому для fetcher/builder нужны ещё локальные bindings. type rbSubject struct { name string namespace string binding string } 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"}, {name: "fission-fetcher", namespace: ns, binding: "fission-fetcher-local-user-ns"}, {name: "fission-builder", namespace: ns, binding: "fission-builder-local-user-ns"}, } fissionSysNS := os.Getenv("FISSION_SYSTEM_NAMESPACE") if fissionSysNS == "" { fissionSysNS = "fission" } 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", // namespace-scoped через RoleBinding, не ClusterRoleBinding }, "subjects": []any{ map[string]any{ "kind": "ServiceAccount", "name": sa.name, "namespace": subjectNS, }, }, }, } _, rbErr := s.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 := s.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. 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 := s.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 (pod → pod в том же ns) // - трафик из fission core namespace (router → function pod) // - трафик из kube-system (kubelet health checks, DNS) // Запрещаем: // - трафик от подов других 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, executor) map[string]any{ "from": []any{ map[string]any{ "namespaceSelector": map[string]any{ "matchLabels": map[string]any{ "kubernetes.io/metadata.name": fissionSysNSForNetpol, }, }, }, }, }, // Разрешаем трафик из kube-system (DNS, health checks) map[string]any{ "from": []any{ map[string]any{ "namespaceSelector": map[string]any{ "matchLabels": map[string]any{ "kubernetes.io/metadata.name": "kube-system", }, }, }, }, }, }, }, }} _, npErr := s.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 s.nsReconcileCh <- struct{}{}: default: // уже есть сигнал в буфере — не блокируем } } // Environments создаются лениво (lazy) в момент создания первой функции на конкретном языке. // Это экономит ресурсы: пул подов поднимается только под те языки что реально используются. return nil } // handleNSStatus возвращает статус инициализации пользовательского namespace. // Используется UI для отображения прогресса при первом входе. 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: NS существует и Active (ensureUserNS уже вызван auth middleware) nsObj, err := s.dyn.Resource(namespaceGVR).Get(ctx, ns, metav1.GetOptions{}) if err == nil { phase, _, _ := unstructured.NestedString(nsObj.Object, "status", "phase") stages[0].Done = phase == "Active" } // Stage 2: NS в FISSION_RESOURCE_NAMESPACES executor deployment if stages[0].Done { execDep, err2 := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Get(ctx, "executor", metav1.GetOptions{}) if err2 == nil { containers, _, _ := unstructured.NestedSlice(execDep.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 { stages[1].Done = true } } } } } break } } } // 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, }) } 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 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 } 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: токен — это sub (email), Deck не вызывается ns = func() string { h32 := sha256.Sum256([]byte(body.Token)) return "fission-" + hex.EncodeToString(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.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}) } func (s *server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, name string) { ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) defer cancel() var pkgName, envName string fn, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).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") triggers, err := s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).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(httpTrigGVR).Namespace(s.userNS(r)).Delete(ctx, trig.GetName(), metav1.DeleteOptions{}) } } } if err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).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(packageGVR).Namespace(s.userNS(r)).Delete(ctx, pkgName, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("delete package %q: %v", pkgName, err)) return } } // Если environment больше не используется ни одной функцией — удаляем его. // Если environment больше не используется ни одной функцией — удаляем его. // Fission увидит удаление Environment CRD и убьёт pool deployment → поды умирают. if envName != "" { cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 15*time.Second) defer cleanupCancel() s.cleanupEnvironmentIfUnused(cleanupCtx, s.userNS(r), envName) } // Сигналим reconciler: он проверит все namespace'ы и уберёт пустые из FISSION_RESOURCE_NAMESPACES. // Не делаем это на горячем пути — reconciler батчит изменения без race condition и rolling restarts. select { case s.nsReconcileCh <- struct{}{}: default: } writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName}) } func buildConfig(kubeconfig string) (*rest.Config, error) { if kubeconfig != "" { cfg, err := clientcmd.BuildConfigFromFlags("", kubeconfig) if err == nil { return cfg, nil } return nil, fmt.Errorf("kubeconfig %s: %w", kubeconfig, err) } cfg, err := rest.InClusterConfig() if err == nil { return cfg, nil } loadingRules := &clientcmd.ClientConfigLoadingRules{} clientCfg := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(loadingRules, &clientcmd.ConfigOverrides{}) return clientCfg.ClientConfig() } 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) } } 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) } 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) } 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}) } 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) }) } 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") if r.Method == http.MethodOptions { w.WriteHeader(http.StatusNoContent) return } next.ServeHTTP(w, r) }) } func withSecurityHeaders(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.Header().Set("X-Content-Type-Options", "nosniff") w.Header().Set("X-Frame-Options", "DENY") w.Header().Set("Referrer-Policy", "strict-origin-when-cross-origin") w.Header().Set("Permissions-Policy", "camera=(), microphone=(), geolocation=()") 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) }) } 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 } 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 } func (s *server) extractPackageSourceCode(ctx context.Context, 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 decodedCode, decErr := decodeLiteralToSource(literal); decErr == nil && strings.TrimSpace(decodedCode) != "" { return decodedCode } } 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, fetchErr := s.fetchPackageArchive(ctx, urlValue) if fetchErr != nil { continue } decodedCode, decErr := decodeArchiveBytesToSource(archiveBytes) if decErr == nil && strings.TrimSpace(decodedCode) != "" { return decodedCode } } return "" } func (s *server) fetchPackageArchive(ctx context.Context, 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, fmt.Errorf("archive request failed: %s", resp.Status) } return io.ReadAll(resp.Body) } func decodeLiteralToSource(literal string) (string, error) { decoded, err := base64.StdEncoding.DecodeString(literal) if err != nil { return "", err } return decodeArchiveBytesToSource(decoded) } func decodeArchiveBytesToSource(decoded []byte) (string, error) { if len(decoded) == 0 { return "", fmt.Errorf("empty payload") } if utf8.Valid(decoded) { return string(decoded), nil } if len(decoded) >= 4 && bytes.Equal(decoded[:4], []byte{'P', 'K', 3, 4}) { if src, zipErr := decodeZipSource(decoded); zipErr == nil { return src, nil } } return "", fmt.Errorf("payload does not contain utf-8 source") } 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 { return "", readErr } if utf8.Valid(content) { return string(content), nil } } } } files := make([]*zip.File, 0, len(reader.File)) for _, file := range reader.File { if file.FileInfo().IsDir() { continue } 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 { continue } if utf8.Valid(content) { return string(content), nil } } return "", fmt.Errorf("zip archive does not contain utf-8 source files") } func readZipFile(file *zip.File) ([]byte, error) { rc, err := file.Open() if err != nil { return nil, err } defer rc.Close() return io.ReadAll(rc) } // goParse проверяет синтаксис Go-кода через go/parser. func goParse(code string) (*token.FileSet, error) { fset := token.NewFileSet() _, err := parser.ParseFile(fset, "code.go", code, parser.AllErrors) if err != nil { return nil, err } return fset, nil } // handleAICheck проверяет синтаксис кода через реальный линтер языка. // POST /console/api/ai/check // Body: {"language":"python","code":"..."} // Response: {"ok":true,"result":"..."} 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 прямо в процессе — без внешних команд _, parseErr := goParse(req.Code) if parseErr == nil { isOK = true result = "✅ Синтаксис корректен." } else { isOK = false result = parseErr.Error() } } else { // Записываем код во временный файл 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 содержит только захардкоженные команды из langsMap 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 // =============================================================================