610 lines
21 KiB
Go
610 lines
21 KiB
Go
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/<route> + Authorization: Bearer <user-token>.
|
||
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"].
|
||
|