From 9ebd36eed301d7d2719eebe9a65052d373780fd2 Mon Sep 17 00:00:00 2001 From: Naeel Date: Sun, 19 Apr 2026 15:05:12 +0300 Subject: [PATCH] feat: lazy env creation + TTL functions + expiry reaper --- console/main.go | 234 ++++++++++++++++++++++++++++++++++++++++-------- 1 file changed, 195 insertions(+), 39 deletions(-) diff --git a/console/main.go b/console/main.go index 2856d67..54d2ccf 100644 --- a/console/main.go +++ b/console/main.go @@ -16,6 +16,7 @@ import ( "net/http" "os" "sort" + "strconv" "strings" "sync" "time" @@ -76,6 +77,7 @@ type createFunctionRequest struct { Entrypoint string `json:"entrypoint"` Route string `json:"route"` Methods []string `json:"methods"` + TTL string `json:"ttl"` // e.g. "24h", "7d" — пустое = функция не протухает } type langEnvDef struct { @@ -190,6 +192,7 @@ func main() { } log.Printf("fission-console listening on :%s (namespace=%s)", port, namespace) + s.startExpiryReaper(envDurationDefault("REAPER_INTERVAL", 5*time.Minute)) log.Fatal(httpServer.ListenAndServe()) } @@ -262,25 +265,14 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { req.Entrypoint = strings.TrimSpace(req.Entrypoint) req.Route = strings.TrimSpace(req.Route) - // Resolve language → environment (auto-create if needed) + // Resolve language → environment (lazy creation). + // ensureEnvironment создаёт Environment CRD если не существует — Fission увидит и поднимет pool pod. if req.Language != "" { - langDef, ok := langEnvMap[req.Language] - if !ok { - writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("unsupported language: %q", req.Language)) - return - } - envName := "console-" + req.Language + "-env" - ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) - defer cancel() - _, err := s.dyn.Resource(environmentGVR).Namespace(s.userNS(r)).Get(ctx, envName, metav1.GetOptions{}) - if apierrors.IsNotFound(err) { - env := s.buildLangEnvironment(envName, s.userNS(r), langDef) - if _, err := s.dyn.Resource(environmentGVR).Namespace(s.userNS(r)).Create(ctx, env, metav1.CreateOptions{}); err != nil { - writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create environment %q: %v", envName, err)) - return - } - } else if err != nil { - writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("check environment: %v", err)) + envCtx, envCancel := context.WithTimeout(r.Context(), 15*time.Second) + defer envCancel() + envName, err := s.ensureEnvironment(envCtx, ns, req.Language) + if err != nil { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure environment: %v", err)) return } req.Environment = envName @@ -367,12 +359,25 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { return } + // Парсим TTL — если указан, записываем аннотацию на Function CRD. + // reaper периодически читает эту аннотацию и удаляет протухшие функции + чистит environment если он больше не используется. + fnAnnotations := map[string]any{} + 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) + } + fn := &unstructured.Unstructured{Object: map[string]any{ "apiVersion": "fission.io/v1", "kind": "Function", "metadata": map[string]any{ - "name": req.Name, - "namespace": ns, + "name": req.Name, + "namespace": ns, + "annotations": fnAnnotations, }, "spec": map[string]any{ "environment": map[string]any{ @@ -422,10 +427,11 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { } writeAnyJSON(w, http.StatusCreated, map[string]any{ - "name": req.Name, - "package": pkgName, + "name": req.Name, + "package": pkgName, "httptrigger": triggerName, - "route": req.Route, + "route": req.Route, + "expires_at": fnAnnotations["fission-console/expires-at"], }) } @@ -454,6 +460,160 @@ func (s *server) buildLangEnvironment(name string, ns string, def langEnvDef) *u }} } +// ensureEnvironment создаёт Environment CRD для языка lang в namespace ns если не существует. +// Вся Fission-специфичная конструкция объекта изолирована в buildLangEnvironment — +// при обновлении Fission (изменение схемы CRD) меняем только там. +func (s *server) ensureEnvironment(ctx context.Context, ns, lang string) (string, error) { + langDef, ok := langEnvMap[lang] + if !ok { + return "", fmt.Errorf("unsupported language: %q", lang) + } + envName := "console-" + lang + "-env" + _, getErr := s.dyn.Resource(environmentGVR).Namespace(ns).Get(ctx, envName, metav1.GetOptions{}) + if getErr == nil { + return envName, nil // уже существует + } + if !apierrors.IsNotFound(getErr) { + return "", fmt.Errorf("check environment %q: %w", envName, getErr) + } + env := s.buildLangEnvironment(envName, ns, langDef) + if _, createErr := s.dyn.Resource(environmentGVR).Namespace(ns).Create(ctx, env, metav1.CreateOptions{}); createErr != nil && !apierrors.IsAlreadyExists(createErr) { + return "", fmt.Errorf("create environment %q: %w", envName, createErr) + } + log.Printf("ensureEnvironment: created %s/%s", ns, envName) + return envName, nil +} + +// cleanupEnvironmentIfUnused удаляет Environment CRD если ни одна функция в namespace его не использует. +// Fission увидит удаление → уберёт pool Deployment → поды умирают. +// Не блокирующий: ошибки логируются, не возвращаются вызывающему коду. +func (s *server) cleanupEnvironmentIfUnused(ctx context.Context, ns, envName string) { + functions, err := s.dyn.Resource(functionGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if err != nil { + log.Printf("cleanupEnvironmentIfUnused: list functions in %s: %v", ns, err) + return + } + for _, fn := range functions.Items { + fnEnv, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") + if fnEnv == envName { + return // ещё используется хотя бы одной функцией + } + } + // Ни одна функция не ссылается на этот environment — удаляем + if delErr := s.dyn.Resource(environmentGVR).Namespace(ns).Delete(ctx, envName, metav1.DeleteOptions{}); delErr != nil && !apierrors.IsNotFound(delErr) { + log.Printf("cleanupEnvironmentIfUnused: delete env %s/%s: %v", ns, envName, delErr) + return + } + log.Printf("cleanupEnvironmentIfUnused: deleted unused env %s/%s", ns, envName) +} + +// parseTTL парсит строку TTL и возвращает время истечения. +// Поддерживаемые форматы: Go duration (1h, 30m, 24h) и дни (1d, 7d, 30d). +// Суффикс "d" не поддерживается стандартным time.ParseDuration — обрабатываем отдельно. +func parseTTL(ttl string) (time.Time, error) { + if strings.HasSuffix(ttl, "d") { + days, err := strconv.Atoi(strings.TrimSuffix(ttl, "d")) + if err != nil || days <= 0 { + return time.Time{}, fmt.Errorf("invalid days value: %q", ttl) + } + return time.Now().Add(time.Duration(days) * 24 * time.Hour), nil + } + d, err := time.ParseDuration(ttl) + if err != nil { + return time.Time{}, err + } + if d <= 0 { + return time.Time{}, fmt.Errorf("ttl must be positive") + } + return time.Now().Add(d), nil +} + +// startExpiryReaper запускает фоновый goroutine для удаления функций с истёкшим TTL. +// Интервал задаётся через env REAPER_INTERVAL (default 5m). +// При удалении функции вызывает cleanupEnvironmentIfUnused — поды умирают когда language больше не используется. +func (s *server) startExpiryReaper(interval time.Duration) { + go func() { + ticker := time.NewTicker(interval) + defer ticker.Stop() + log.Printf("expiryReaper: started, interval=%v", interval) + for range ticker.C { + s.runExpiryReap() + } + }() +} + +// runExpiryReap обходит все namespace под управлением fission-console и удаляет протухшие функции. +func (s *server) runExpiryReap() { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + + // Находим только наши namespace-ы по метке которую мы ставим при создании + nsList, err := s.dyn.Resource(namespaceGVR).List(ctx, metav1.ListOptions{ + LabelSelector: "managed-by=fission-console", + }) + if err != nil { + log.Printf("expiryReaper: list namespaces: %v", err) + return + } + + now := time.Now().UTC() + for _, ns := range nsList.Items { + s.reapExpiredFunctionsInNS(ctx, ns.GetName(), now) + } +} + +// reapExpiredFunctionsInNS удаляет протухшие функции в конкретном namespace. +// Для каждой удалённой функции вызывает cleanupEnvironmentIfUnused. +func (s *server) reapExpiredFunctionsInNS(ctx context.Context, ns string, now time.Time) { + functions, err := s.dyn.Resource(functionGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if err != nil { + log.Printf("expiryReaper: list functions in %s: %v", ns, err) + return + } + + for _, fn := range functions.Items { + expiresAtStr, _, _ := unstructured.NestedString(fn.Object, "metadata", "annotations", "fission-console/expires-at") + if expiresAtStr == "" { + continue // нет TTL — функция живёт вечно + } + expiresAt, parseErr := time.Parse(time.RFC3339, expiresAtStr) + if parseErr != nil { + log.Printf("expiryReaper: parse expires-at for %s/%s: %v", ns, fn.GetName(), parseErr) + continue + } + if now.Before(expiresAt) { + continue // ещё не протухла + } + + // Функция протухла — удаляем всё + fnName := fn.GetName() + envName, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name") + pkgName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name") + log.Printf("expiryReaper: deleting expired function %s/%s (expired %s ago)", ns, fnName, now.Sub(expiresAt).Round(time.Second)) + + // Удаляем HTTPTrigger-ы ссылающиеся на эту функцию + triggers, tErr := s.dyn.Resource(httpTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) + if tErr == nil { + for _, trig := range triggers.Items { + refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name") + if refName == fnName { + _ = s.dyn.Resource(httpTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{}) + } + } + } + + _ = s.dyn.Resource(functionGVR).Namespace(ns).Delete(ctx, fnName, metav1.DeleteOptions{}) + if pkgName != "" { + _ = s.dyn.Resource(packageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) + } + + // Убираем pool pod если язык больше не используется + if envName != "" { + s.cleanupEnvironmentIfUnused(ctx, ns, envName) + } + } +} + func (s *server) buildGoSourceZip(code string) ([]byte, error) { var buf bytes.Buffer zw := zip.NewWriter(&buf) @@ -983,21 +1143,8 @@ func (s *server) ensureUserNamespace(ctx context.Context, ns string) error { } } - // 2. Создать shared environments для всех языков - for lang, def := range langEnvMap { - envName := lang + "-env" - _, getErr := s.dyn.Resource(environmentGVR).Namespace(ns).Get(ctx, envName, metav1.GetOptions{}) - if getErr == nil { - continue // уже существует - } - if !apierrors.IsNotFound(getErr) { - continue // другая ошибка — пропускаем - } - env := s.buildLangEnvironment(envName, ns, def) - if _, createErr := s.dyn.Resource(environmentGVR).Namespace(ns).Create(ctx, env, metav1.CreateOptions{}); createErr != nil && !apierrors.IsAlreadyExists(createErr) { - log.Printf("ensureUserNamespace: create env %s/%s: %v", ns, envName, createErr) - } - } + // Environments создаются лениво (lazy) в момент создания первой функции на конкретном языке. + // Это экономит ресурсы: пул подов поднимается только под те языки что реально используются. return nil } @@ -1072,10 +1219,11 @@ func (s *server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, na ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) defer cancel() - var pkgName string + var pkgName, envName string fn, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Get(ctx, name, metav1.GetOptions{}) if err == nil { pkgName, _, _ = unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name") + envName, _, _ = unstructured.NestedString(fn.Object, "spec", "environment", "name") } triggers, err := s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).List(ctx, metav1.ListOptions{}) @@ -1100,6 +1248,14 @@ func (s *server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, na } } + // Если environment больше не используется ни одной функцией — удаляем его. + // Fission увидит удаление Environment CRD и убьёт pool deployment → поды умирают. + if envName != "" { + cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cleanupCancel() + s.cleanupEnvironmentIfUnused(cleanupCtx, s.userNS(r), envName) + } + writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName}) }