Files
fission-console/console/internal/api/function_archive.go
T
“Naeel” f42fa366de refactor: split handlers.go into focused modules
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.
2026-05-08 09:00:25 +04:00

343 lines
13 KiB
Go

// Package api — создание и обновление функций из zip-архива (multipart/form-data).
//
// Этот файл отвечает за два сценария:
// 1. handleCreateFunctionFromArchive — создание новой функции из загруженного .zip файла.
// 2. handleUpdateFunctionArchive — обновление существующей функции новым .zip файлом.
//
// Архив загружается пользователем через форму с полем "archive".
// Содержимое архива передаётся в storagesvc → S3 без модификации.
// В отличие от function_code.go, здесь нет трансформации кода — архив идёт как есть.
//
// Связанная аннотация: fission-console/source-type = "archive"
// позволяет UI определить режим редактирования при открытии функции.
package api
import (
"context"
"fmt"
"io"
"net/http"
"strconv"
"strings"
"time"
"fission-console/internal/fission"
"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"
)
// handleCreateFunctionFromArchive создаёт функцию из загруженного zip-архива (multipart/form-data).
// Поля формы: name, language (или environment), entrypoint, route, methods, timeout, ttl.
// Файловое поле: archive (.zip).
func (s *Server) handleCreateFunctionFromArchive(w http.ResponseWriter, r *http.Request, ns string) {
if err := r.ParseMultipartForm(maxArchiveUploadSize); err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("parse multipart form: %v", err))
return
}
name := strings.TrimSpace(r.FormValue("name"))
if name == "" || (!validFuncName.MatchString(name) || len(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
}
lang := strings.TrimSpace(r.FormValue("language"))
envName := strings.TrimSpace(r.FormValue("environment"))
f, fhCreate, err := r.FormFile("archive")
if err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("archive file required: %v", err))
return
}
defer f.Close()
archiveBytes, err := io.ReadAll(io.LimitReader(f, maxArchiveUploadSize))
if err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("read archive: %v", err))
return
}
archiveFilenameCreate := ""
if fhCreate != nil {
archiveFilenameCreate = fhCreate.Filename
}
// Проверяем magic bytes: zip должен начинаться с PK (0x50 0x4B)
if len(archiveBytes) < 4 || archiveBytes[0] != 0x50 || archiveBytes[1] != 0x4B {
writeJSONError(w, http.StatusBadRequest, "загруженный файл не является zip-архивом (ожидается .zip)")
return
}
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
}
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
defer cancel()
// Определяем environment: по языку или явно
if lang != "" {
envCtx, envCancel := context.WithTimeout(r.Context(), 15*time.Second)
defer envCancel()
resolved, envErr := fission.EnsureEnvironment(envCtx, s.dyn, ns, lang)
if envErr != nil {
if strings.Contains(envErr.Error(), "unsupported language") {
writeJSONError(w, http.StatusBadRequest, envErr.Error())
} else {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure environment: %v", envErr))
}
return
}
envName = resolved
}
if envName == "" {
writeJSONError(w, http.StatusBadRequest, "language or environment is required")
return
}
if _, err := s.dyn.Resource(fission.EnvironmentGVR).Namespace(ns).Get(ctx, envName, metav1.GetOptions{}); err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("environment %q not found: %v", envName, err))
return
}
entrypoint := strings.TrimSpace(r.FormValue("entrypoint"))
if entrypoint == "" {
entrypoint = runtime.DefaultEntrypoint(lang)
}
route := strings.TrimSpace(r.FormValue("route"))
if route == "" {
nsShort := ns
if len(nsShort) > 12 {
nsShort = nsShort[len(nsShort)-12:]
}
route = "/" + nsShort + "/" + name
}
if !strings.HasPrefix(route, "/") {
route = "/" + route
}
methods := normalizeMethods(strings.Split(r.FormValue("methods"), ","))
timeout := normalizeFunctionTimeout(0)
if tv := r.FormValue("timeout"); tv != "" {
if n, err := strconv.ParseInt(tv, 10, 64); err == nil {
timeout = normalizeFunctionTimeout(n)
}
}
// Загружаем архив в storagesvc
deploySpec, uploadErr := s.buildDeploySpec(ctx, archiveBytes)
if uploadErr != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("upload archive: %v", uploadErr))
return
}
pkgName := name + "-pkg"
triggerName := name + "-route"
pkg := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "fission.io/v1",
"kind": "Package",
"metadata": map[string]any{"name": pkgName, "namespace": ns},
"spec": map[string]any{
"deployment": deploySpec,
"environment": map[string]any{"name": envName, "namespace": ns},
"source": map[string]any{},
},
}}
now := time.Now().UTC()
fnAnnotations := map[string]any{
"fission-console/language": lang,
fissionSourceTypeAnnotation: "archive",
functionCreatedAtAnnotation: now.Format(time.RFC3339),
functionUpdatedAtAnnotation: now.Format(time.RFC3339),
}
if archiveFilenameCreate != "" {
fnAnnotations["fission-console/archive-filename"] = archiveFilenameCreate
}
if ttl := r.FormValue("ttl"); ttl != "" {
if expiresAt, ttlErr := parseTTL(ttl); ttlErr == nil {
fnAnnotations["fission-console/expires-at"] = expiresAt.UTC().Format(time.RFC3339)
}
}
methodValues := make([]any, 0, len(methods))
for _, m := range methods {
methodValues = append(methodValues, m)
}
fn := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "fission.io/v1",
"kind": "Function",
"metadata": map[string]any{"name": name, "namespace": ns, "annotations": fnAnnotations},
"spec": map[string]any{
"environment": map[string]any{"name": envName, "namespace": ns},
"package": map[string]any{"packageref": map[string]any{"name": pkgName, "namespace": ns}},
"InvokeStrategy": map[string]any{
"ExecutionStrategy": map[string]any{"ExecutorType": "poolmgr"},
"StrategyType": "execution",
},
"functionTimeout": timeout,
},
}}
if entrypoint != "" {
_ = unstructured.SetNestedField(fn.Object, entrypoint, "spec", "package", "functionName")
}
trigger := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "fission.io/v1",
"kind": "HTTPTrigger",
"metadata": map[string]any{"name": triggerName, "namespace": ns},
"spec": map[string]any{
"functionref": map[string]any{"name": name, "type": "name"},
"relativeurl": route,
"methods": methodValues,
},
}}
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", name))
return
}
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create package: %v", err))
return
}
if _, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Create(ctx, fn, metav1.CreateOptions{}); err != nil {
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{})
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create function: %v", err))
return
}
if _, err := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).Create(ctx, trigger, metav1.CreateOptions{}); err != nil {
_ = s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Delete(ctx, name, metav1.DeleteOptions{})
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{})
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create trigger: %v", err))
return
}
writeAnyJSON(w, http.StatusCreated, map[string]any{
"name": name,
"namespace": ns,
"environment": envName,
"route": route,
"source_type": "archive",
})
}
// handleUpdateFunctionArchive обновляет функцию из загруженного zip-архива (multipart/form-data).
// Поля формы: timeout (optional), entrypoint (optional). Файловое поле: archive (.zip).
// Создаёт новый Package (новое имя) — чтобы executor сбросил кэш function service.
func (s *Server) handleUpdateFunctionArchive(w http.ResponseWriter, r *http.Request, name string) {
if err := r.ParseMultipartForm(maxArchiveUploadSize); err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("parse multipart form: %v", err))
return
}
f, _, err := r.FormFile("archive")
if err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("archive file required: %v", err))
return
}
defer f.Close()
archiveBytes, err := io.ReadAll(io.LimitReader(f, maxArchiveUploadSize))
if err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("read archive: %v", err))
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")
envName, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name")
newPkgName := name + "-pkg-" + strconv.FormatInt(time.Now().UnixMilli(), 36)
deploySpec, uploadErr := s.buildDeploySpec(ctx, archiveBytes)
if uploadErr != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("upload archive: %v", 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
}
// Обновляем timeout если задан
timeout := normalizeFunctionTimeout(0)
if tv := r.FormValue("timeout"); tv != "" {
if n, err := strconv.ParseInt(tv, 10, 64); err == nil {
timeout = normalizeFunctionTimeout(n)
}
}
_ = unstructured.SetNestedField(fn.Object, timeout, "spec", "functionTimeout")
// Обновляем entrypoint если передан
if ep := strings.TrimSpace(r.FormValue("entrypoint")); ep != "" {
_ = unstructured.SetNestedField(fn.Object, ep, "spec", "package", "functionName")
}
fnAnnotations := fn.GetAnnotations()
if fnAnnotations == nil {
fnAnnotations = map[string]string{}
}
fnAnnotations[fissionSourceTypeAnnotation] = "archive"
fnAnnotations[functionUpdatedAtAnnotation] = time.Now().UTC().Format(time.RFC3339)
if fh, fhErr := r.MultipartForm.File["archive"]; fhErr == false || len(fh) > 0 {
if files := r.MultipartForm.File["archive"]; len(files) > 0 && files[0].Filename != "" {
fnAnnotations["fission-console/archive-filename"] = files[0].Filename
}
}
fn.SetAnnotations(fnAnnotations)
if err := unstructured.SetNestedField(fn.Object, map[string]any{
"name": newPkgName,
"namespace": ns,
"resourceversion": createdPkg.GetResourceVersion(),
}, "spec", "package", "packageref"); err != nil {
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{})
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set packageref: %v", err))
return
}
if _, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Update(ctx, fn, metav1.UpdateOptions{}); err != nil {
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{})
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("update function: %v", err))
return
}
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,
"source_type": "archive",
})
}