handlers.go (267 lines) — routing, constants, handleAuth, utils storagesvc.go — S3 upload/delete via storagesvc function_code.go — create/update function from inline code function_archive.go — create/update function from zip archive function_crud.go — GET, DELETE, timeout update, pod logs function_invoke.go — invoke via router, /fn route gateway, cron gateway No logic changes. go build OK.
411 lines
16 KiB
Go
411 lines
16 KiB
Go
// Package api — создание и обновление функций из исходного кода (inline code).
|
||
//
|
||
// Этот файл отвечает за два сценария:
|
||
// 1. handleCreateFunction — создание новой функции из кода (JSON body).
|
||
// 2. handleUpdateFunctionCode — обновление существующей функции: новый код → новый Package.
|
||
//
|
||
// Поддерживаемые языки: python, nodejs, php, ruby, go.
|
||
// Для Go создаётся source package (builder компилирует .so плагин).
|
||
// Для остальных языков — deployment archive (zip загружается в storagesvc или как literal).
|
||
//
|
||
// Почему новый Package при обновлении:
|
||
// Fission executor кэширует function service по functionUid и не видит изменений
|
||
// в существующем Package. Новое имя пакета гарантирует cache miss в executor.
|
||
package api
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"net/http"
|
||
"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"
|
||
)
|
||
|
||
// buildDeployArchive упаковывает исходный код в байты для deployment Package.
|
||
// Для nodejs — ESM-обёртка (package.json + main.js).
|
||
// Для php/ruby — zip с одним файлом скрипта.
|
||
// Для остальных (python) — raw bytes кода.
|
||
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
|
||
}
|
||
}
|
||
|
||
// 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)
|
||
|
||
// Поддерживаем два формата: JSON (код) и multipart/form-data (архив).
|
||
isArchiveUpload := strings.HasPrefix(r.Header.Get("Content-Type"), "multipart/form-data")
|
||
if isArchiveUpload {
|
||
s.handleCreateFunctionFromArchive(w, r, ns)
|
||
return
|
||
}
|
||
|
||
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,
|
||
fissionSourceTypeAnnotation: "code",
|
||
}
|
||
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"],
|
||
})
|
||
}
|
||
|
||
// 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,
|
||
})
|
||
}
|