package api import ( "bytes" "context" "encoding/json" "errors" "fmt" "io" "log" "net" "net/http" "regexp" "strconv" "strings" "time" "fission-console/internal/fission" "fission-console/internal/model" "fission-console/internal/runtime" apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" ) // validFuncName — RFC 1123 subdomain label: строчные буквы+цифры+дефис, без дефиса в начале/конце. // Максимум 57 символов (не 63): самый длинный суффикс "-route" (HTTPTrigger) = 6 символов. // 63 - 6 = 57. Fission webhook требует все связанные объекты <= 63 символов. var validFuncName = regexp.MustCompile(`^[a-z0-9]([a-z0-9-]*[a-z0-9])?$`) // maxCodeSize — максимальный размер кода функции (1 MB). // Выше — не имеет смысла для inline функции; лучше использовать Package с URL. const maxCodeSize = 1 << 20 // defaultFunctionInvokeTimeout совпадает с дефолтом Fission для spec.functionTimeout. const defaultFunctionInvokeTimeout = 60 * time.Second 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 } result, err := fission.CreateFunction(ctx, s.dyn, ns, req) if err != nil { status := http.StatusBadGateway if apierrors.IsAlreadyExists(err) { status = http.StatusConflict } else if apierrors.IsInvalid(err) { status = http.StatusBadRequest } writeJSONError(w, status, err.Error()) return } writeAnyJSON(w, http.StatusCreated, result) } // 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.http, 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, }) } 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 } if req.Timeout <= 0 { req.Timeout = int64(defaultFunctionInvokeTimeout / time.Second) } ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) defer cancel() ns := s.userNS(r) result, err := fission.UpdateFunctionCode(ctx, s.dyn, ns, name, req) if err != nil { status := http.StatusBadGateway if apierrors.IsNotFound(err) { status = http.StatusNotFound } writeJSONError(w, status, err.Error()) return } writeAnyJSON(w, http.StatusOK, result) } // 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) } 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) if err := fission.DeleteFunction(ctx, s.dyn, ns, name); err != nil { status := http.StatusBadGateway if apierrors.IsNotFound(err) { status = http.StatusNotFound } writeJSONError(w, status, err.Error()) return } writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name}) } // handleAuth обрабатывает POST /console/api/auth. // Валидирует токен, создаёт namespace, возвращает namespace пользователя. // 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"].