feat: lazy env creation + TTL functions + expiry reaper
This commit is contained in:
+195
-39
@@ -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})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user