spec.runtime.container.env срезается CRD-схемой Fission. Хранение перенесено в аннотацию fission-console/env-vars (JSON-массив). Версия v1.3.73
386 lines
15 KiB
Go
386 lines
15 KiB
Go
// 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/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"
|
|
)
|
|
|
|
// 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)
|
|
|
|
writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName})
|
|
}
|
|
|
|
// handleGetFunctionLogs возвращает логи пода функции (последние 100 строк).
|
|
// Ищет под по лейблу functionName=<name> в 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 обновляет переменные окружения функции в аннотации fission-console/env-vars (JSON)
|
|
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 для аннотации
|
|
envJSON, err := json.Marshal(req.EnvVars)
|
|
if err != nil {
|
|
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("marshal env vars: %v", err))
|
|
return
|
|
}
|
|
|
|
// Обновляем аннотацию updated-at
|
|
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 _, 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, "count": len(req.EnvVars)})
|
|
}
|
|
|
|
// 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
|
|
}
|