package api import ( "bytes" "context" "encoding/base64" "encoding/json" "errors" "fmt" "io" "log" "net" "net/http" "net/url" "regexp" "strconv" "strings" "time" "fission-console/internal/auth" "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 // defaultFunctionInvokeTimeout совпадает с дефолтом Fission для spec.functionTimeout. const defaultFunctionInvokeTimeout = 60 * time.Second func buildDeployArchive(lang, code string) ([]byte, error) { switch lang { case "nodejs": return runtime.BuildJSDeployZip(code) case "php": return runtime.BuildScriptZip(code, "main.php") case "ruby": return runtime.BuildScriptZip(code, "handler.rb") default: return []byte(code), nil } } // deleteFromStoragesvc удаляет архив из S3 через storagesvc по URL из Package spec. // URL имеет формат: http://storagesvc.../v1/archive?id=fission/UUID // Best-effort: ошибка логируется, но не прерывает операцию удаления. func (s *Server) deleteFromStoragesvc(ctx context.Context, archiveURL string) { if s.storagesvcURL == "" || archiveURL == "" { return } // archiveURL = "http://storagesvc.../v1/archive?id=fission/UUID" // Строим DELETE URL к storagesvc, сохраняя query-параметр id parsed, err := url.Parse(archiveURL) if err != nil { log.Printf("deleteFromStoragesvc: parse url %q: %v", archiveURL, err) return } deleteURL := strings.TrimRight(s.storagesvcURL, "/") + "/v1/archive?" + parsed.RawQuery req, err := http.NewRequestWithContext(ctx, http.MethodDelete, deleteURL, nil) if err != nil { log.Printf("deleteFromStoragesvc: build request: %v", err) return } resp, err := s.http.Do(req) if err != nil { log.Printf("deleteFromStoragesvc: %v", err) return } defer resp.Body.Close() if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusNoContent { body, _ := io.ReadAll(resp.Body) log.Printf("deleteFromStoragesvc: status %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) return } log.Printf("deleteFromStoragesvc: deleted %s", parsed.Query().Get("id")) } // uploadToStoragesvc загружает байты в Fission storagesvc и возвращает URL для package archive. // Если storagesvcURL не задан — возвращает пустую строку (fallback на literal). func (s *Server) uploadToStoragesvc(ctx context.Context, data []byte) (string, error) { if s.storagesvcURL == "" { log.Printf("uploadToStoragesvc: storagesvcURL is empty, skip upload") return "", nil } log.Printf("uploadToStoragesvc: uploading %d bytes to %s", len(data), s.storagesvcURL) uploadURL := strings.TrimRight(s.storagesvcURL, "/") + "/v1/archive" body := &bytes.Reader{} // multipart/form-data с полем uploadfile var buf bytes.Buffer boundary := fmt.Sprintf("fission%d", time.Now().UnixNano()) buf.WriteString("--" + boundary + "\r\n") buf.WriteString(fmt.Sprintf("Content-Disposition: form-data; name=\"uploadfile\"; filename=\"archive.zip\"\r\n")) buf.WriteString("Content-Type: application/octet-stream\r\n\r\n") buf.Write(data) buf.WriteString("\r\n--" + boundary + "--\r\n") body = bytes.NewReader(buf.Bytes()) req, err := http.NewRequestWithContext(ctx, http.MethodPost, uploadURL, body) if err != nil { return "", fmt.Errorf("build storagesvc upload request: %w", err) } req.Header.Set("Content-Type", "multipart/form-data; boundary="+boundary) req.Header.Set("X-File-Size", fmt.Sprintf("%d", len(data))) resp, err := s.http.Do(req) if err != nil { return "", fmt.Errorf("storagesvc upload: %w", err) } defer resp.Body.Close() respBody, _ := io.ReadAll(resp.Body) if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated { return "", fmt.Errorf("storagesvc upload status %d: %s", resp.StatusCode, strings.TrimSpace(string(respBody))) } var result struct { ID string `json:"id"` } if err := json.Unmarshal(respBody, &result); err != nil || result.ID == "" { return "", fmt.Errorf("storagesvc upload: bad response: %s", string(respBody)) } archiveURL := strings.TrimRight(s.storagesvcURL, "/") + "/v1/archive?id=" + result.ID return archiveURL, nil } // buildDeploySpec строит spec.deployment для Fission Package. // Если storagesvcURL задан — загружает архив в S3 через storagesvc и возвращает type: url. // Иначе — возвращает type: literal с base64-кодом. func (s *Server) buildDeploySpec(ctx context.Context, data []byte) (map[string]any, error) { archiveURL, err := s.uploadToStoragesvc(ctx, data) if err != nil { log.Printf("storagesvc upload failed, falling back to literal: %v", err) // fallback — сохраняем как literal return map[string]any{"type": "literal", "literal": base64.StdEncoding.EncodeToString(data)}, nil } if archiveURL == "" { return map[string]any{"type": "literal", "literal": base64.StdEncoding.EncodeToString(data)}, nil } return map[string]any{"type": "url", "url": archiveURL}, nil } func (s *Server) resolveInvokeTimeout(fn *unstructured.Unstructured) time.Duration { if fn != nil { seconds, found, err := unstructured.NestedInt64(fn.Object, "spec", "functionTimeout") if err == nil && found && seconds > 0 { return time.Duration(seconds) * time.Second } } if s.invokeTimeout > 0 { return s.invokeTimeout } return defaultFunctionInvokeTimeout } func normalizeFunctionTimeout(seconds int64) int64 { if seconds <= 0 { return int64(defaultFunctionInvokeTimeout / time.Second) } return seconds } func normalizeRoute(route string) string { route = strings.TrimSpace(route) if route == "" || route == "/" { return "/" } if !strings.HasPrefix(route, "/") { return "/" + route } return route } func routeAllowsMethod(methods []string, method string) bool { if len(methods) == 0 { return true } for _, candidate := range methods { if strings.EqualFold(strings.TrimSpace(candidate), method) { return true } } return false } func appendUniqueMethods(dst []string, src []string) []string { for _, method := range src { method = strings.ToUpper(strings.TrimSpace(method)) if method == "" { continue } seen := false for _, existing := range dst { if existing == method { seen = true break } } if !seen { dst = append(dst, method) } } return dst } func shouldForwardRequestBody(method string) bool { switch method { case http.MethodGet, http.MethodHead: return false default: return true } } func copyProxyRequestHeaders(dst, src http.Header) { for key, values := range src { switch http.CanonicalHeaderKey(key) { case "Authorization", "X-Auth-Token", "X-Auth-Env", "Host", "Content-Length": continue } for _, value := range values { dst.Add(key, value) } } } func copyProxyResponseHeaders(dst, src http.Header) { for key, values := range src { for _, value := range values { dst.Add(key, value) } } } func doRequestWithContextTimeout(client *http.Client, req *http.Request) (*http.Response, error) { if client == nil { return http.DefaultClient.Do(req) } invokeClient := *client // Для invoke реальный лимит должен определяться context timeout функции, // а не общим HTTP timeout console. invokeClient.Timeout = 0 return invokeClient.Do(req) } // 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) req.Timeout = normalizeFunctionTimeout(req.Timeout) 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 archive с кодом (S3 или literal fallback) 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 } srcSpec, srcErr := s.buildDeploySpec(ctx, srcZip) if srcErr != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("upload go source: %v", srcErr)) return } pkgSpec = map[string]any{ "source": srcSpec, "deployment": map[string]any{}, "environment": map[string]any{"name": req.Environment, "namespace": ns}, "buildcommand": "build", } } else { deployBytes, archiveErr := buildDeployArchive(req.Language, req.Code) if archiveErr != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build %s archive: %v", req.Language, archiveErr)) return } deploySpec, uploadErr := s.buildDeploySpec(ctx, deployBytes) if uploadErr != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("upload %s archive: %v", req.Language, uploadErr)) return } pkgSpec = map[string]any{ "deployment": deploySpec, "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, } now := time.Now().UTC() fnAnnotations[functionCreatedAtAnnotation] = now.Format(time.RFC3339) fnAnnotations[functionUpdatedAtAnnotation] = now.Format(time.RFC3339) 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}, "functionTimeout": req.Timeout, "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") functionTimeout, foundTimeout, _ := unstructured.NestedInt64(fn.Object, "spec", "functionTimeout") if !foundTimeout || functionTimeout <= 0 { functionTimeout = int64(defaultFunctionInvokeTimeout / time.Second) } // Извлекаем исходный код из 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, "timeout": functionTimeout, "created_at": functionTimestampResponse(fn)["created_at"], "updated_at": functionTimestampResponse(fn)["updated_at"], "code": code, "route": route, "methods": methods, "raw": fn.Object, }) } // handleUpdateFunctionCode обновляет код уже существующей функции. // Создаёт НОВЫЙ Package (вместо обновления старого) чтобы executor сбросил кэш: // executor кэширует function service по functionUid и не видит изменений в том же Package. // Новое имя пакета гарантирует cache miss в executor. 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 } oldPkgName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name") // Определяем язык из аннотации — нужен для правильной упаковки lang, _, _ := unstructured.NestedString(fn.Object, "metadata", "annotations", "fission-console/language") deployBytes, archiveErr := buildDeployArchive(lang, req.Code) if archiveErr != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build %s archive: %v", lang, archiveErr)) return } // Создаём новый Package с уникальным именем. // Это единственный способ сбросить кэш executor: он кэширует по functionUid и // не замечает изменений в существующем Package. envName, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") createdAt := func() time.Time { ann := fn.GetAnnotations() if ann != nil { if v := strings.TrimSpace(ann[functionCreatedAtAnnotation]); v != "" { if ts, err := parseRFC3339(v); err == nil { return ts.UTC() } } } if ts := fn.GetCreationTimestamp(); !ts.IsZero() { return ts.UTC() } return time.Time{} }() now := time.Now().UTC() newPkgName := name + "-pkg-" + strconv.FormatInt(time.Now().UnixMilli(), 36) deploySpec, uploadErr := s.buildDeploySpec(ctx, deployBytes) if uploadErr != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("upload %s archive: %v", lang, uploadErr)) return } newPkg := &unstructured.Unstructured{Object: map[string]any{ "apiVersion": "fission.io/v1", "kind": "Package", "metadata": map[string]any{"name": newPkgName, "namespace": ns}, "spec": map[string]any{ "deployment": deploySpec, "environment": map[string]any{"name": envName, "namespace": ns}, "source": map[string]any{}, }, }} createdPkg, err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Create(ctx, newPkg, metav1.CreateOptions{}) if err != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create new package: %v", err)) return } // Обновляем Function на новый Package if err := unstructured.SetNestedField(fn.Object, normalizeFunctionTimeout(req.Timeout), "spec", "functionTimeout"); err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set function timeout: %v", err)) _ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{}) return } ensureFunctionTimestamps(fn, now) if createdAt.IsZero() { createdAt = now } fnAnnotations := fn.GetAnnotations() if fnAnnotations == nil { fnAnnotations = map[string]string{} } fnAnnotations[functionCreatedAtAnnotation] = createdAt.UTC().Format(time.RFC3339) fnAnnotations[functionUpdatedAtAnnotation] = now.Format(time.RFC3339) fn.SetAnnotations(fnAnnotations) if err := unstructured.SetNestedField(fn.Object, map[string]any{ "name": newPkgName, "namespace": ns, "resourceversion": createdPkg.GetResourceVersion(), }, "spec", "package", "packageref"); err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set function packageref: %v", err)) _ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{}) 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: %v", name, err)) _ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{}) return } // Удаляем старый Package (best effort) if oldPkgName != "" && oldPkgName != newPkgName { _ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, oldPkgName, metav1.DeleteOptions{}) } writeAnyJSON(w, http.StatusOK, map[string]any{ "updated": true, "package": newPkgName, }) } // 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("{}") } ns := s.userNS(r) lookupCtx, lookupCancel := context.WithTimeout(r.Context(), 10*time.Second) defer lookupCancel() // Проверяем существование функции до вызова — лучше 404 чем непонятный timeout fn, err2 := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(lookupCtx, name, metav1.GetOptions{}) if 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 } invokeTimeout := s.resolveInvokeTimeout(fn) ctx, cancel := context.WithTimeout(r.Context(), invokeTimeout) defer cancel() // Ищем HTTPTrigger чтобы получить реальный URL и метод invokeURL := buildInternalInvokeURL(s.routerURL, ns, 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 := doRequestWithContextTimeout(s.http, req) if err != nil { // Отличаем timeout от сетевой ошибки. if errors.Is(err, context.DeadlineExceeded) { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s", 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", 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 buildInternalInvokeURL(routerURL, namespace, functionName string) string { if namespace == "default" || namespace == "" { return fmt.Sprintf("%s/fission-function/%s", routerURL, functionName) } return fmt.Sprintf("%s/fission-function/%s/%s", routerURL, namespace, functionName) } // handleFissionFunctionGateway принимает внутренние invoke-запросы timer/router // и проксирует их через console в upstream router с корректным router JWT. func (s *Server) handleFissionFunctionGateway(w http.ResponseWriter, r *http.Request) { rawPath := strings.Trim(strings.TrimPrefix(r.URL.Path, "/fission-function"), "/") if rawPath == "" { http.NotFound(w, r) return } parts := strings.Split(rawPath, "/") namespace := s.ns functionName := "" remainingPath := "" if len(parts) == 1 { functionName = strings.TrimSpace(parts[0]) } else { namespace = strings.TrimSpace(parts[0]) functionName = strings.TrimSpace(parts[1]) if len(parts) > 2 { remainingPath = "/" + strings.Join(parts[2:], "/") } } if namespace == "" || functionName == "" { writeJSONError(w, http.StatusBadRequest, "namespace and function name are required") return } s.invokeInternalFunction(w, r, namespace, functionName, remainingPath) } func (s *Server) invokeInternalFunction(w http.ResponseWriter, r *http.Request, namespace, functionName, extraPath 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("{}") } lookupCtx, lookupCancel := context.WithTimeout(r.Context(), 10*time.Second) defer lookupCancel() fn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(namespace).Get(lookupCtx, functionName, metav1.GetOptions{}) if err != nil { if apierrors.IsNotFound(err) { writeJSONError(w, http.StatusNotFound, fmt.Sprintf("function %q not found", functionName)) return } writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function %q: %v", functionName, err)) return } invokeTimeout := s.resolveInvokeTimeout(fn) ctx, cancel := context.WithTimeout(r.Context(), invokeTimeout) defer cancel() invokeURL := buildInternalInvokeURL(s.routerURL, namespace, functionName) + extraPath if r.URL.RawQuery != "" { invokeURL += "?" + r.URL.RawQuery } var invokeBody io.Reader if shouldForwardRequestBody(r.Method) { invokeBody = bytes.NewReader(bodyBytes) } req, err := http.NewRequestWithContext(ctx, r.Method, invokeURL, invokeBody) if err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build invoke request: %v", err)) return } if shouldForwardRequestBody(r.Method) { req.Header.Set("Content-Type", "application/json") } copyProxyRequestHeaders(req.Header, r.Header) if token := s.getRouterToken(); token != "" { req.Header.Set("Authorization", "Bearer "+token) } start := time.Now() resp, err := doRequestWithContextTimeout(s.http, req) if err != nil { if errors.Is(err, context.DeadlineExceeded) { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s", functionName, invokeTimeout)) return } var netErr net.Error if errors.As(err, &netErr) && netErr.Timeout() { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s", functionName, invokeTimeout)) return } writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q: %v", functionName, 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), }) } // handleInvokeRoute даёт пользователю прямой HTTP gateway к своей функции по route. // Внешний контракт: /fn/ + Authorization: Bearer . func (s *Server) handleInvokeRoute(w http.ResponseWriter, r *http.Request) { route := normalizeRoute(strings.TrimPrefix(r.URL.Path, "/fn")) if route == "/" { writeJSONError(w, http.StatusBadRequest, "route is required") return } ns := s.userNS(r) lookupCtx, lookupCancel := context.WithTimeout(r.Context(), 10*time.Second) defer lookupCancel() triggers, err := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).List(lookupCtx, metav1.ListOptions{}) if err != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("list httptriggers: %v", err)) return } matchedFunction := "" allowedMethods := make([]string, 0, 4) for _, trig := range triggers.Items { trigRoute, _, _ := unstructured.NestedString(trig.Object, "spec", "relativeurl") if normalizeRoute(trigRoute) != route { continue } methods, _, _ := unstructured.NestedStringSlice(trig.Object, "spec", "methods") allowedMethods = appendUniqueMethods(allowedMethods, methods) if !routeAllowsMethod(methods, r.Method) { continue } matchedFunction, _, _ = unstructured.NestedString(trig.Object, "spec", "functionref", "name") if matchedFunction != "" { break } } if matchedFunction == "" { if len(allowedMethods) > 0 { w.Header().Set("Allow", strings.Join(allowedMethods, ", ")) writeJSONError(w, http.StatusMethodNotAllowed, fmt.Sprintf("route %q does not allow method %s", route, r.Method)) return } writeJSONError(w, http.StatusNotFound, fmt.Sprintf("route %q not found", route)) return } fn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(lookupCtx, matchedFunction, metav1.GetOptions{}) if err != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function %q: %v", matchedFunction, err)) return } bodyBytes, err := io.ReadAll(r.Body) if err != nil { writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("read request body: %v", err)) return } invokeTimeout := s.resolveInvokeTimeout(fn) ctx, cancel := context.WithTimeout(r.Context(), invokeTimeout) defer cancel() invokeURL := s.routerURL + route if r.URL.RawQuery != "" { invokeURL += "?" + r.URL.RawQuery } var invokeBody io.Reader if shouldForwardRequestBody(r.Method) { invokeBody = bytes.NewReader(bodyBytes) } req, err := http.NewRequestWithContext(ctx, r.Method, invokeURL, invokeBody) if err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build invoke request: %v", err)) return } copyProxyRequestHeaders(req.Header, r.Header) if token := s.getRouterToken(); token != "" { req.Header.Set("Authorization", "Bearer "+token) } resp, err := doRequestWithContextTimeout(s.http, req) if err != nil { if errors.Is(err, context.DeadlineExceeded) { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke route %q timeout after %s", route, invokeTimeout)) return } var netErr net.Error if errors.As(err, &netErr) && netErr.Timeout() { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke route %q timeout after %s", route, invokeTimeout)) return } writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke route %q: %v", route, err)) return } defer resp.Body.Close() copyProxyResponseHeaders(w.Header(), resp.Header) w.WriteHeader(resp.StatusCode) _, _ = io.Copy(w, resp.Body) } // handleDeleteFunction удаляет функцию и связанные объекты: HTTPTrigger, TimeTrigger, 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{}) } } } // Удаляем связанные TimeTrigger-ы if triggers, err := s.dyn.Resource(fission.TimeTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{}); err == nil { for _, trig := range triggers.Items { refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name") if refName == name { _ = s.dyn.Resource(fission.TimeTrigGVR).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 != "" { // Получаем URL архива из Package spec.deployment перед удалением, чтобы потом очистить S3 var archiveURL string if pkg, pkgErr := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Get(ctx, pkgName, metav1.GetOptions{}); pkgErr == nil { deployType, _, _ := unstructured.NestedString(pkg.Object, "spec", "deployment", "type") if deployType == "url" { archiveURL, _, _ = unstructured.NestedString(pkg.Object, "spec", "deployment", "url") } } 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 } // Удаляем архив из S3 после успешного удаления Package (best-effort) if archiveURL != "" { go s.deleteFromStoragesvc(context.Background(), archiveURL) } } // Убираем 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 NS удалён — за FISSION_RESOURCE_NAMESPACES теперь отвечает Layer 1 NSWatcher) writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName}) } // handleAuth обрабатывает POST /console/api/auth. // Валидирует токен через authenticator, создаёт namespace, возвращает namespace + email. 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 := normalizeEnv(body.Env) identity, err := s.authenticator.Authenticate(r.Context(), body.Token, env) if err != nil { writeJSONError(w, http.StatusUnauthorized, "invalid token") return } ns := auth.NamespaceForSub(identity.Sub) 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, "email": identity.Email}) } // 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 }