package fission import ( "context" "fmt" "log" "os" "sort" "strings" "sync" "time" apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/dynamic" "encoding/json" ) // NSInflightEnsure — состояние in-flight вызова EnsureUserNamespace. // Все параллельные горутины ждут close(done), затем читают err. // Реализует ручной singleflight без внешних зависимостей. type NSInflightEnsure struct { done chan struct{} err error } // NSManager управляет жизненным циклом пользовательских namespace-ов: // создание, RBAC, quota, network policies, reconciler, reaper. type NSManager struct { dyn dynamic.Interface // nsReconcileCh — сигнал для немедленного запуска NS reconciler. // Буферизирован на 1: несколько сигналов схлопываются в один запуск. NSReconcileCh chan struct{} // Singleflight + кэш для ensureUserNamespace: // ensuredNS — namespace-ы которые уже инициализированы в этом запуске процесса. // ensuredNSInFlight — текущие in-flight операции (ключ = namespace name). ensuredNSMu sync.Mutex ensuredNS map[string]struct{} ensuredNSInFlight map[string]*NSInflightEnsure // nsSemaphore ограничивает параллелизм EnsureUserNamespace — не более 3 одновременно. // Без него 10 новых пользователей генерируют 140 K8s API calls одновременно → throttle → 504. // С семафором: 3 batch-а по 14 calls → ~3 × 5s = 15s total, все укладываются в timeout. nsSemaphore chan struct{} } // NewNSManager создаёт NSManager с инициализированными каналами и структурами. func NewNSManager(dyn dynamic.Interface) *NSManager { return &NSManager{ dyn: dyn, NSReconcileCh: make(chan struct{}, 1), ensuredNS: make(map[string]struct{}), ensuredNSInFlight: make(map[string]*NSInflightEnsure), nsSemaphore: make(chan struct{}, 3), } } // EnsureUserNS — единая точка входа для гарантии существования пользовательского namespace. // Реализует singleflight + in-memory кэш + семафор параллелизма. // // Singleflight: если один goroutine уже создаёт namespace ns — остальные ждут его результата // вместо того чтобы запускать параллельные K8s API calls (вызывало throttle и 504). // // Кэш: если namespace уже создан в этом запуске процесса — быстрый путь без K8s calls. func (m *NSManager) EnsureUserNS(ctx context.Context, ns string) error { m.ensuredNSMu.Lock() if _, ok := m.ensuredNS[ns]; ok { // Быстрый путь: уже создан в этой жизни процесса. m.ensuredNSMu.Unlock() return nil } if inflight, ok := m.ensuredNSInFlight[ns]; ok { // Кто-то уже создаёт — ждём его результата (singleflight). m.ensuredNSMu.Unlock() select { case <-inflight.done: return inflight.err case <-ctx.Done(): return ctx.Err() } } // Мы первые для этого namespace. inflight := &NSInflightEnsure{done: make(chan struct{})} m.ensuredNSInFlight[ns] = inflight m.ensuredNSMu.Unlock() // Берём слот семафора — ограничиваем параллелизм. select { case m.nsSemaphore <- struct{}{}: case <-ctx.Done(): m.ensuredNSMu.Lock() delete(m.ensuredNSInFlight, ns) m.ensuredNSMu.Unlock() inflight.err = ctx.Err() close(inflight.done) return ctx.Err() } ensureCtx, ensureCancel := context.WithTimeout(ctx, 60*time.Second) inflight.err = m.ensureUserNamespace(ensureCtx, ns) ensureCancel() <-m.nsSemaphore // освобождаем слот m.ensuredNSMu.Lock() delete(m.ensuredNSInFlight, ns) if inflight.err == nil { m.ensuredNS[ns] = struct{}{} } m.ensuredNSMu.Unlock() close(inflight.done) return inflight.err } // ensureUserNamespace выполняет реальную работу: создаёт namespace + RBAC + quota + netpol. // Все операции идемпотентны (IsAlreadyExists игнорируется). func (m *NSManager) ensureUserNamespace(ctx context.Context, ns string) error { // 1. Создать namespace с меткой managed-by=fission-console. // Метка используется reconciler-ом для фильтрации наших namespace-ов. nsObj := &unstructured.Unstructured{ Object: map[string]any{ "apiVersion": "v1", "kind": "Namespace", "metadata": map[string]any{ "name": ns, "labels": map[string]any{ "managed-by": "fission-console", }, }, }, } _, err := m.dyn.Resource(NamespaceGVR).Create(ctx, nsObj, metav1.CreateOptions{}) newlyCreated := err == nil if err != nil && !apierrors.IsAlreadyExists(err) { return fmt.Errorf("create namespace %s: %w", ns, err) } // 1a. ServiceAccounts для Fission в user namespace. // // Fission pool pods (fetcher sidecar) запускаются в user namespace и требуют // `serviceAccountName: fission-fetcher` в том же namespace. Executor (с SERVICEACCOUNT_CHECK_ENABLED=false) // не создаёт эти SA автоматически → pool pods падают с "serviceaccount not found". saGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "serviceaccounts"} for _, saName := range []string{"fission-fetcher", "fission-builder"} { saObj := &unstructured.Unstructured{Object: map[string]any{ "apiVersion": "v1", "kind": "ServiceAccount", "metadata": map[string]any{ "name": saName, "namespace": ns, }, }} _, saErr := m.dyn.Resource(saGVR).Namespace(ns).Create(ctx, saObj, metav1.CreateOptions{}) if saErr != nil && !apierrors.IsAlreadyExists(saErr) { log.Printf("ensureUserNamespace: create SA %s/%s: %v", ns, saName, saErr) } } // 1b. RoleBindings для Fission SA в user namespace. // // Почему cluster-admin, а не admin: // ClusterRole "admin" не включает custom resource группы (fission.io/*). // Fission executor при старте пытается создать Role с правами на fission.io/packages, // и получает "attempting to grant RBAC permissions not currently held" — RBAC escalation // prevention. ClusterRole "cluster-admin" в контексте RoleBinding (не ClusterRoleBinding) // даёт полный доступ ТОЛЬКО внутри конкретного namespace — это безопасно. type rbSubject struct { name string namespace string binding string } fissionSysNS := os.Getenv("FISSION_SYSTEM_NAMESPACE") if fissionSysNS == "" { fissionSysNS = "fission" } fissionSAs := []rbSubject{ {name: "fission-executor", binding: "fission-executor-user-ns"}, {name: "fission-router", binding: "fission-router-user-ns"}, {name: "fission-buildermgr", binding: "fission-buildermgr-user-ns"}, {name: "fission-kubewatcher", binding: "fission-kubewatcher-user-ns"}, {name: "fission-timer", binding: "fission-timer-user-ns"}, {name: "fission-fetcher", binding: "fission-fetcher-system-user-ns"}, {name: "fission-builder", binding: "fission-builder-system-user-ns"}, // fetcher/builder локальные в user namespace (там запускаются pool pods) {name: "fission-fetcher", namespace: ns, binding: "fission-fetcher-local-user-ns"}, {name: "fission-builder", namespace: ns, binding: "fission-builder-local-user-ns"}, } rbGVR := schema.GroupVersionResource{Group: "rbac.authorization.k8s.io", Version: "v1", Resource: "rolebindings"} for _, sa := range fissionSAs { subjectNS := sa.namespace if subjectNS == "" { subjectNS = fissionSysNS } rbObj := &unstructured.Unstructured{ Object: map[string]any{ "apiVersion": "rbac.authorization.k8s.io/v1", "kind": "RoleBinding", "metadata": map[string]any{ "name": sa.binding, "namespace": ns, }, "roleRef": map[string]any{ "apiGroup": "rbac.authorization.k8s.io", "kind": "ClusterRole", "name": "cluster-admin", }, "subjects": []any{ map[string]any{ "kind": "ServiceAccount", "name": sa.name, "namespace": subjectNS, }, }, }, } _, rbErr := m.dyn.Resource(rbGVR).Namespace(ns).Create(ctx, rbObj, metav1.CreateOptions{}) if rbErr != nil && !apierrors.IsAlreadyExists(rbErr) { log.Printf("ensureUserNamespace: create rolebinding %s/%s@%s: %v", ns, sa.name, subjectNS, rbErr) } } // 1c. ResourceQuota — ограничиваем ресурсы одного пользователя. // Без этого одна функция может исчерпать CPU/RAM всего кластера. // Параметры задаются через env vars для гибкой настройки без пересборки. quotaGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "resourcequotas"} quotaObj := &unstructured.Unstructured{Object: map[string]any{ "apiVersion": "v1", "kind": "ResourceQuota", "metadata": map[string]any{"name": "user-quota", "namespace": ns}, "spec": map[string]any{ "hard": map[string]any{ "requests.cpu": envDefault("QUOTA_REQ_CPU", "1"), "requests.memory": envDefault("QUOTA_REQ_MEM", "1Gi"), "limits.cpu": envDefault("QUOTA_LIM_CPU", "4"), "limits.memory": envDefault("QUOTA_LIM_MEM", "4Gi"), "pods": envDefault("QUOTA_PODS", "30"), "count/functions.fission.io": envDefault("QUOTA_FUNCTIONS", "20"), "count/packages.fission.io": envDefault("QUOTA_PACKAGES", "40"), "count/httptriggers.fission.io": envDefault("QUOTA_HTTPTRIGGERS", "20"), }, }, }} _, quotaErr := m.dyn.Resource(quotaGVR).Namespace(ns).Create(ctx, quotaObj, metav1.CreateOptions{}) if quotaErr != nil && !apierrors.IsAlreadyExists(quotaErr) { log.Printf("ensureUserNamespace: create ResourceQuota %s: %v", ns, quotaErr) } // 1d. LimitRange — дефолтные лимиты на контейнер. // Поды без явных limits — unbounded; LimitRange автоматически добавляет defaults. limitRangeGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "limitranges"} limitRangeObj := &unstructured.Unstructured{Object: map[string]any{ "apiVersion": "v1", "kind": "LimitRange", "metadata": map[string]any{"name": "user-limits", "namespace": ns}, "spec": map[string]any{ "limits": []any{ map[string]any{ "type": "Container", "default": map[string]any{"cpu": envDefault("LIMIT_DEFAULT_CPU", "500m"), "memory": envDefault("LIMIT_DEFAULT_MEM", "256Mi")}, "defaultRequest": map[string]any{ "cpu": envDefault("LIMIT_REQ_CPU", "50m"), "memory": envDefault("LIMIT_REQ_MEM", "64Mi"), }, "max": map[string]any{"cpu": envDefault("LIMIT_MAX_CPU", "2"), "memory": envDefault("LIMIT_MAX_MEM", "1Gi")}, }, }, }, }} _, lrErr := m.dyn.Resource(limitRangeGVR).Namespace(ns).Create(ctx, limitRangeObj, metav1.CreateOptions{}) if lrErr != nil && !apierrors.IsAlreadyExists(lrErr) { log.Printf("ensureUserNamespace: create LimitRange %s: %v", ns, lrErr) } // 1e. NetworkPolicy — запрещаем входящий трафик из других user namespace. // Разрешаем: внутри namespace, из fission core namespace, из kube-system. // Запрещаем: трафик от подов других user namespace (межтенантная изоляция). fissionSysNSForNetpol := os.Getenv("FISSION_SYSTEM_NAMESPACE") if fissionSysNSForNetpol == "" { fissionSysNSForNetpol = "fission" } netpolGVR := schema.GroupVersionResource{Group: "networking.k8s.io", Version: "v1", Resource: "networkpolicies"} netpolObj := &unstructured.Unstructured{Object: map[string]any{ "apiVersion": "networking.k8s.io/v1", "kind": "NetworkPolicy", "metadata": map[string]any{"name": "deny-cross-tenant", "namespace": ns}, "spec": map[string]any{ "podSelector": map[string]any{}, // применяется ко всем подам namespace "policyTypes": []any{"Ingress"}, "ingress": []any{ // Разрешаем трафик внутри namespace map[string]any{"from": []any{map[string]any{"podSelector": map[string]any{}}}}, // Разрешаем трафик из fission core namespace (router → function pod) map[string]any{"from": []any{map[string]any{"namespaceSelector": map[string]any{ "matchLabels": map[string]any{"kubernetes.io/metadata.name": fissionSysNSForNetpol}, }}}}, // Разрешаем трафик из kube-system (kubelet health checks, DNS) map[string]any{"from": []any{map[string]any{"namespaceSelector": map[string]any{ "matchLabels": map[string]any{"kubernetes.io/metadata.name": "kube-system"}, }}}}, }, }, }} _, npErr := m.dyn.Resource(netpolGVR).Namespace(ns).Create(ctx, netpolObj, metav1.CreateOptions{}) if npErr != nil && !apierrors.IsAlreadyExists(npErr) { log.Printf("ensureUserNamespace: create NetworkPolicy %s: %v", ns, npErr) } // 1f. Сигналим NS reconciler что появился новый namespace. // Reconciler сам синхронизирует FISSION_RESOURCE_NAMESPACES без race condition и rolling restarts. if newlyCreated { select { case m.NSReconcileCh <- struct{}{}: default: // уже есть сигнал в буфере — не блокируем } } // Environments создаются лениво (lazy) в момент создания первой функции на языке. return nil } // StartNSReconciler запускает фоновый reconciler FISSION_RESOURCE_NAMESPACES. // // Проблема при scale: // - Прямой патч на горячем пути → rolling restart всех Fission deployments при каждом новом юзере // - Без мьютекса: 100 юзеров одновременно → read-modify-write race → namespace-ы теряются // - Ручное kubectl delete ns → namespace остаётся в переменной вечно // // Решение: единственный goroutine с debounce-каналом. // - Запускается по таймеру или немедленно через NSReconcileCh // - 1000 юзеров одновременно → 1 патч вместо 5000 // - Нет race condition (один goroutine, один writer) func (m *NSManager) StartNSReconciler(interval time.Duration) { go func() { ticker := time.NewTicker(interval) defer ticker.Stop() log.Printf("nsReconciler: started, interval=%v", interval) for { select { case <-ticker.C: m.reconcileNSList() case <-m.NSReconcileCh: m.reconcileNSList() // Дренируем канал чтобы не запускаться дважды подряд после сигнала drain: for { select { case <-m.NSReconcileCh: default: break drain } } } } }() } // reconcileNSList синхронизирует FISSION_RESOURCE_NAMESPACES с реально существующими namespace-ами. // // Алгоритм: // 1. Читает все namespace-ы с меткой managed-by=fission-console из k8s (источник истины) // 2. Читает текущий FISSION_RESOURCE_NAMESPACES из router deployment // 3. Вычисляет desired = "default" + активные namespace-ы // 4. Если desired == current → ничего не делает (нет rolling restart!) // 5. Если разница → один патч всех Fission deployments func (m *NSManager) reconcileNSList() { ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) defer cancel() fissionNS := os.Getenv("FISSION_SYSTEM_NAMESPACE") if fissionNS == "" { fissionNS = "fission" } // Шаг 1: реальные namespace-ы с нашей меткой nsList, err := m.dyn.Resource(NamespaceGVR).List(ctx, metav1.ListOptions{ LabelSelector: "managed-by=fission-console", }) if err != nil { log.Printf("nsReconciler: list namespaces: %v", err) return } // Шаг 2: фильтруем — берём только Active. Terminating убираем (они уже умирают). desired := map[string]struct{}{"default": {}} for _, ns := range nsList.Items { phase, _, _ := unstructured.NestedString(ns.Object, "status", "phase") if phase == "Active" { desired[ns.GetName()] = struct{}{} } } // Шаг 3: текущее значение из router deployment routerDep, err := m.dyn.Resource(DeploymentGVR).Namespace(fissionNS).Get(ctx, "router", metav1.GetOptions{}) if err != nil { log.Printf("nsReconciler: get router: %v", err) return } currentVal := extractFissionNSEnv(routerDep) // Шаг 4: сравниваем current с desired currentSet := map[string]struct{}{} for _, p := range strings.Split(currentVal, ",") { if t := strings.TrimSpace(p); t != "" { currentSet[t] = struct{}{} } } same := len(currentSet) == len(desired) if same { for k := range desired { if _, ok := currentSet[k]; !ok { same = false break } } } if same { return // ничего менять не нужно — нет патча, нет rolling restart } // Шаг 5: патчим все Fission deployments parts := make([]string, 0, len(desired)) for ns := range desired { parts = append(parts, ns) } sort.Strings(parts) newVal := strings.Join(parts, ",") fissionDeployments := []string{"router", "executor", "buildermgr", "kubewatcher", "timer"} patch := map[string]any{ "spec": map[string]any{ "template": map[string]any{ "spec": map[string]any{ "containers": []any{ map[string]any{ "name": "", "env": []any{ map[string]any{ "name": "FISSION_RESOURCE_NAMESPACES", "value": newVal, }, }, }, }, }, }, }, } for _, dep := range fissionDeployments { patch["spec"].(map[string]any)["template"].(map[string]any)["spec"].(map[string]any)["containers"].([]any)[0].(map[string]any)["name"] = dep patchBytes, _ := json.Marshal(patch) if _, pErr := m.dyn.Resource(DeploymentGVR).Namespace(fissionNS).Patch( ctx, dep, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{}); pErr != nil { log.Printf("nsReconciler: patch %s: %v", dep, pErr) } } log.Printf("nsReconciler: synced FISSION_RESOURCE_NAMESPACES: %s → %s", currentVal, newVal) } // extractFissionNSEnv читает FISSION_RESOURCE_NAMESPACES из первого контейнера deployment. func extractFissionNSEnv(dep *unstructured.Unstructured) string { containers, _, _ := unstructured.NestedSlice(dep.Object, "spec", "template", "spec", "containers") for _, c := range containers { cont, ok := c.(map[string]any) if !ok { continue } envs, _, _ := unstructured.NestedSlice(cont, "env") for _, e := range envs { env, ok := e.(map[string]any) if !ok { continue } if env["name"] == "FISSION_RESOURCE_NAMESPACES" { if v, ok := env["value"].(string); ok && v != "" { return v } } } break // смотрим только первый контейнер } return "default" } // StartExpiryReaper запускает фоновый goroutine для удаления функций с истёкшим TTL. func (m *NSManager) StartExpiryReaper(interval time.Duration) { go func() { ticker := time.NewTicker(interval) defer ticker.Stop() log.Printf("expiryReaper: started, interval=%v", interval) for range ticker.C { m.runExpiryReap() } }() } // runExpiryReap обходит все управляемые namespace-ы и удаляет протухшие функции. func (m *NSManager) runExpiryReap() { ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) defer cancel() nsList, err := m.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 { m.reapExpiredFunctionsInNS(ctx, ns.GetName(), now) } } // reapExpiredFunctionsInNS удаляет протухшие функции в конкретном namespace. // Для каждой удалённой функции вызывает CleanupEnvironmentIfUnused. // Также удаляет orphan packages — пакеты без соответствующей функции. func (m *NSManager) reapExpiredFunctionsInNS(ctx context.Context, ns string, now time.Time) { functions, err := m.dyn.Resource(FunctionGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) if err != nil { log.Printf("expiryReaper: list functions in %s: %v", ns, err) return } // Строим множество активных функций для поиска orphan packages activeFunctions := make(map[string]struct{}, len(functions.Items)) for _, fn := range functions.Items { activeFunctions[fn.GetName()] = struct{}{} } 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 := m.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 { _ = m.dyn.Resource(HTTPTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{}) } } } _ = m.dyn.Resource(FunctionGVR).Namespace(ns).Delete(ctx, fnName, metav1.DeleteOptions{}) if pkgName != "" { _ = m.dyn.Resource(PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) } if envName != "" { CleanupEnvironmentIfUnused(ctx, m.dyn, ns, envName) } } // Сигналим reconciler — он проверит все namespace-ы и почистит список select { case m.NSReconcileCh <- struct{}{}: default: } // Orphan packages — пакеты без соответствующей функции // (остаются если процесс упал в середине удаления) packages, pkgListErr := m.dyn.Resource(PackageGVR).Namespace(ns).List(ctx, metav1.ListOptions{}) if pkgListErr == nil { for _, pkg := range packages.Items { pkgName := pkg.GetName() // Конвенция именования: {fn-name}-pkg if !strings.HasSuffix(pkgName, "-pkg") { continue } fnName := strings.TrimSuffix(pkgName, "-pkg") if _, exists := activeFunctions[fnName]; !exists { log.Printf("expiryReaper: deleting orphan package %s/%s (no matching function)", ns, pkgName) _ = m.dyn.Resource(PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) } } } } // envDefault читает env переменную или возвращает fallback. func envDefault(key, fallback string) string { if v := strings.TrimSpace(os.Getenv(key)); v != "" { return v } return fallback }