- handleCreateFunction: deps → аннотация fission-console/deps - handleUpdateFunctionCode: deps → аннотация fission-console/deps (или удаление если пусто) - GET /functions/:name: возвращает deps из аннотации - openEdit: заполняет e-deps.value = fn.deps (вместо пустого поля)
530 lines
20 KiB
Go
530 lines
20 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/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 := ""
|
||
deps := ""
|
||
if ann := fn.GetAnnotations(); ann != nil {
|
||
if v := ann[fissionSourceTypeAnnotation]; v != "" {
|
||
sourceType = v
|
||
}
|
||
if v := ann["fission-console/archive-filename"]; v != "" {
|
||
archiveFilename = v
|
||
}
|
||
if v := ann["fission-console/deps"]; v != "" {
|
||
deps = 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,
|
||
"deps": deps,
|
||
"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=<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 обновляет переменные окружения функции.
|
||
//
|
||
// Логика переключения 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=<name>, functionNamespace=<ns>, 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
|
||
}
|