// Package api — CRUD операции с функциями: чтение, удаление, обновление таймаута, логи, env vars. // // Этот файл содержит операции, не связанные с заменой кода/архива: // - handleGetFunction — GET /functions/:name (детали: код, route, environment, source_type) // - handleDeleteFunction — DELETE /functions/:name (каскадное удаление: триггеры, Package, S3) // - handleUpdateFunctionTimeout — PUT /functions/:name/timeout (только таймаут, без замены кода) // - handleGetFunctionLogs — GET /functions/:name/logs (логи пода через Kubernetes API) // - handleGetFunctionEnvVars — GET /functions/:name/envvars (переменные окружения из CRD) // - handlePutFunctionEnvVars — PUT /functions/:name/envvars (обновить env vars в CRD) // // Операции с кодом и архивом — в function_code.go и function_archive.go соответственно. // Вызов функции — в function_invoke.go. package api import ( "context" "encoding/json" "fmt" "io" "net/http" "strings" "time" "fission-console/internal/billing" "fission-console/internal/fission" corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" k8stypes "k8s.io/apimachinery/pkg/types" ) // 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 } } // Читаем source-type аннотацию (code / archive) sourceType := "code" archiveFilename := "" if ann := fn.GetAnnotations(); ann != nil { if v := ann[fissionSourceTypeAnnotation]; v != "" { sourceType = v } if v := ann["fission-console/archive-filename"]; v != "" { archiveFilename = v } } 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, "source_type": sourceType, "archive_filename": archiveFilename, "route": route, "methods": methods, "env_vars": extractEnvVars(fn), "raw": fn.Object, }) } // handleUpdateFunctionTimeout обновляет только spec.functionTimeout функции (без замены кода/архива). func (s *Server) handleUpdateFunctionTimeout(w http.ResponseWriter, r *http.Request, name string) { ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) defer cancel() ns := s.userNS(r) var req struct { Timeout int64 `json:"timeout"` Entrypoint string `json:"entrypoint"` } if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err)) return } timeout := normalizeFunctionTimeout(req.Timeout) fn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, name, metav1.GetOptions{}) if err != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function %q: %v", name, err)) return } if err := unstructured.SetNestedField(fn.Object, timeout, "spec", "functionTimeout"); err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set timeout: %v", err)) return } if req.Entrypoint != "" { _ = unstructured.SetNestedField(fn.Object, req.Entrypoint, "spec", "package", "functionName") } now := time.Now().UTC().Format(time.RFC3339) ann := fn.GetAnnotations() if ann == nil { ann = map[string]string{} } ann[functionUpdatedAtAnnotation] = now fn.SetAnnotations(ann) 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)) return } writeAnyJSON(w, http.StatusOK, map[string]any{"updated": true, "timeout": timeout}) } // 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) s.billing.RecordInvocation(billing.Invocation{ Namespace: ns, FunctionName: name, TriggerType: billing.TriggerEvent, StartedAt: time.Now(), StatusCode: http.StatusOK, RecordedBy: "console", EventType: "delete", }) writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName}) } // handleGetFunctionLogs возвращает логи пода функции (последние 100 строк). // Ищет под по лейблу functionName= в namespace пользователя. func (s *Server) handleGetFunctionLogs(w http.ResponseWriter, r *http.Request, name string) { ns := s.userNS(r) ctx := r.Context() if s.kube == nil { writeJSONError(w, http.StatusServiceUnavailable, "kubernetes client not available") return } labelSelector := "functionName=" + name pods, err := s.kube.CoreV1().Pods(ns).List(ctx, metav1.ListOptions{ LabelSelector: labelSelector, }) if err != nil { writeJSONError(w, http.StatusInternalServerError, "list pods: "+err.Error()) return } if len(pods.Items) == 0 { writeAnyJSON(w, http.StatusOK, map[string]any{ "logs": "(нет запущенных подов для функции " + name + ")", }) return } var allLogs strings.Builder tailLines := int64(100) for _, pod := range pods.Items { containerName := "" if len(pod.Spec.Containers) > 0 { containerName = pod.Spec.Containers[0].Name } req := s.kube.CoreV1().Pods(ns).GetLogs(pod.Name, &corev1.PodLogOptions{ Container: containerName, TailLines: &tailLines, }) rc, err := req.Stream(ctx) if err != nil { allLogs.WriteString("[" + pod.Name + ": ошибка чтения логов: " + err.Error() + "]\n") continue } data, _ := io.ReadAll(rc) rc.Close() if allLogs.Len() > 0 { allLogs.WriteString("\n--- " + pod.Name + " ---\n") } else { allLogs.WriteString("--- " + pod.Name + " ---\n") } allLogs.Write(data) } writeAnyJSON(w, http.StatusOK, map[string]any{ "logs": allLogs.String(), }) } // handleGetFunctionEnvVars возвращает переменные окружения функции из .spec.runtime.container.env func (s *Server) handleGetFunctionEnvVars(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 } envVars := extractEnvVars(fn) writeAnyJSON(w, http.StatusOK, map[string]any{"env_vars": envVars}) } // handlePutFunctionEnvVars обновляет переменные окружения функции. // // Логика переключения ExecutorType: // - Если env vars непустые → ExecutorType: newdeploy + spec.podspec.containers[0].env // (newdeploy создаёт dedicated Deployment, Kubernetes ставит env vars на уровне ОС) // - Если env vars пустые → ExecutorType: poolmgr, podspec удаляется // (poolmgr использует warm pool, быстрый cold start) // // Это единственный универсальный способ передать env vars в pod для всех языков // (Python, Go, Ruby, PHP, Node.js) без изменений в env-серверах. func (s *Server) handlePutFunctionEnvVars(w http.ResponseWriter, r *http.Request, name string) { ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) defer cancel() ns := s.userNS(r) var req struct { EnvVars []map[string]string `json:"env_vars"` // [{name: "KEY", value: "VAL"}, ...] } if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err)) return } // Валидация: имена переменных for _, ev := range req.EnvVars { k := ev["name"] if k == "" { writeJSONError(w, http.StatusBadRequest, "env var name cannot be empty") return } } 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 } // Сериализуем в JSON для аннотации (для UI) envJSON, err := json.Marshal(req.EnvVars) if err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("marshal env vars: %v", err)) return } // Обновляем аннотации now := time.Now().UTC().Format(time.RFC3339) ann := fn.GetAnnotations() if ann == nil { ann = map[string]string{} } ann[functionUpdatedAtAnnotation] = now ann["fission-console/env-vars"] = string(envJSON) fn.SetAnnotations(ann) if len(req.EnvVars) > 0 { // Есть env vars → newdeploy + podspec с env vars envName, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") // Строим список env vars для Kubernetes envList := make([]any, 0, len(req.EnvVars)) for _, ev := range req.EnvVars { envList = append(envList, map[string]any{ "name": ev["name"], "value": ev["value"], }) } // Устанавливаем podspec.containers[0] с env vars // Имя контейнера = имя environment (стандарт Fission) if err := unstructured.SetNestedSlice(fn.Object, []any{ map[string]any{ "name": envName, "env": envList, }, }, "spec", "podspec", "containers"); err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set podspec: %v", err)) return } // Переключаем на newdeploy (только он поддерживает podspec env) if err := unstructured.SetNestedField(fn.Object, map[string]any{ "ExecutionStrategy": map[string]any{ "ExecutorType": "newdeploy", "MinScale": int64(0), "MaxScale": int64(1), "SpecializationTimeout": int64(120), }, "StrategyType": "execution", }, "spec", "InvokeStrategy"); err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set invoke strategy: %v", err)) return } } else { // Нет env vars → poolmgr, убираем podspec unstructured.RemoveNestedField(fn.Object, "spec", "podspec") if err := unstructured.SetNestedField(fn.Object, map[string]any{ "ExecutionStrategy": map[string]any{ "ExecutorType": "poolmgr", "SpecializationTimeout": int64(120), }, "StrategyType": "execution", }, "spec", "InvokeStrategy"); err != nil { writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set invoke strategy: %v", err)) 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)) return } // Fission newdeploy НЕ применяет fn.Spec.PodSpec при построении Deployment. // Поэтому напрямую патчим существующий Deployment через Kubernetes API. if len(req.EnvVars) > 0 && s.kube != nil { envName, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") if envName != "" { if err := s.patchDeploymentEnvVars(ctx, ns, name, envName, req.EnvVars); err != nil { // Не фатальная ошибка — CRD обновлён, Deployment будет обновлён позже _ = err // warn only } } } executor := "poolmgr" if len(req.EnvVars) > 0 { executor = "newdeploy" } writeAnyJSON(w, http.StatusOK, map[string]any{"updated": true, "count": len(req.EnvVars), "executor": executor}) } // extractEnvVars читает аннотацию fission-console/env-vars (JSON) из Function CRD // Возвращает [{name, value}, ...] func extractEnvVars(fn *unstructured.Unstructured) []map[string]string { ann := fn.GetAnnotations() if ann == nil { return []map[string]string{} } raw := ann["fission-console/env-vars"] if raw == "" { return []map[string]string{} } var result []map[string]string if err := json.Unmarshal([]byte(raw), &result); err != nil { return []map[string]string{} } return result } // patchDeploymentEnvVars находит Deployment newdeploy для функции и патчит его env vars. // Fission не применяет fn.Spec.PodSpec при построении Deployment, поэтому патчим напрямую. // Поиск по labels: functionName=, functionNamespace=, executorType=newdeploy func (s *Server) patchDeploymentEnvVars(ctx context.Context, ns, fnName, envContainerName string, envVars []map[string]string) error { selector := fmt.Sprintf("functionName=%s,functionNamespace=%s,executorType=newdeploy", fnName, ns) deplList, err := s.kube.AppsV1().Deployments(ns).List(ctx, metav1.ListOptions{LabelSelector: selector}) if err != nil { return fmt.Errorf("list deployments: %w", err) } if len(deplList.Items) == 0 { return nil // Deployment ещё не создан Fission — ничего страшного } // Строим env vars для patch (StrategicMergePatch мержит по "name") envItems := make([]map[string]string, 0, len(envVars)) for _, ev := range envVars { envItems = append(envItems, map[string]string{"name": ev["name"], "value": ev["value"]}) } patch := map[string]any{ "spec": map[string]any{ "template": map[string]any{ "spec": map[string]any{ "containers": []any{ map[string]any{ "name": envContainerName, "env": envItems, }, }, }, }, }, } patchBytes, err := json.Marshal(patch) if err != nil { return fmt.Errorf("marshal patch: %w", err) } for _, depl := range deplList.Items { if _, err := s.kube.AppsV1().Deployments(ns).Patch( ctx, depl.Name, k8stypes.StrategicMergePatchType, patchBytes, metav1.PatchOptions{}, ); err != nil { return fmt.Errorf("patch deployment %s: %w", depl.Name, err) } } return nil }