247 lines
9.6 KiB
Go
247 lines
9.6 KiB
Go
package cloud
|
||
|
||
import (
|
||
"context"
|
||
"log"
|
||
"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/client-go/dynamic"
|
||
|
||
"fission-console/internal/fission"
|
||
)
|
||
|
||
// NSManager управляет полным lifecycle пользовательских namespace-ов:
|
||
// - Оркестрирует создание NS через fission.SetupFissionNamespace (Layer 1 adapter)
|
||
// - Накладывает облачные политики: ResourceQuota, LimitRange, NetworkPolicy
|
||
// - Реализует singleflight + кэш + семафор для защиты от concurrent создания
|
||
// - Запускает ExpiryReaper для удаления функций с истёкшим TTL
|
||
//
|
||
// NSManager не занимается синхронизацией FISSION_RESOURCE_NAMESPACES —
|
||
// это теперь задача Layer 1 (multitenant.NSWatcher в executor).
|
||
type NSManager struct {
|
||
dyn dynamic.Interface
|
||
|
||
// Singleflight + кэш: ensuredNS содержит namespace-ы уже инициализированные в этом процессе.
|
||
// ensuredNSInFlight — текущие in-flight операции (ключ = namespace).
|
||
ensuredNSMu sync.Mutex
|
||
ensuredNS map[string]struct{}
|
||
ensuredNSInFlight map[string]*nsInflightEnsure
|
||
|
||
// nsSemaphore ограничивает параллелизм EnsureUserNS — не более 3 одновременно.
|
||
// Без него 10 новых пользователей → 140 K8s API calls одновременно → throttle → 504.
|
||
nsSemaphore chan struct{}
|
||
}
|
||
|
||
// nsInflightEnsure — состояние in-flight вызова EnsureUserNS.
|
||
// Все конкурентные goroutine ждут close(done), затем читают err (singleflight без зависимостей).
|
||
type nsInflightEnsure struct {
|
||
done chan struct{}
|
||
err error
|
||
}
|
||
|
||
// NewNSManager создаёт NSManager.
|
||
func NewNSManager(dyn dynamic.Interface) *NSManager {
|
||
return &NSManager{
|
||
dyn: dyn,
|
||
ensuredNS: make(map[string]struct{}),
|
||
ensuredNSInFlight: make(map[string]*nsInflightEnsure),
|
||
nsSemaphore: make(chan struct{}, 3),
|
||
}
|
||
}
|
||
|
||
// EnsureUserNS — единая точка входа для создания namespace при первом запросе пользователя.
|
||
// Гарантирует что namespace + Fission RBAC + ResourceQuota + NetworkPolicy существуют.
|
||
//
|
||
// Паттерны:
|
||
// - Singleflight: параллельные вызовы для одного NS объединяются — один делает работу, остальные ждут
|
||
// - In-memory кэш: повторные вызовы для уже созданного NS — быстрый путь без K8s calls
|
||
// - Семафор (3): ограничивает параллельное создание новых NS
|
||
func (m *NSManager) EnsureUserNS(ctx context.Context, ns string) error {
|
||
m.ensuredNSMu.Lock()
|
||
if _, ok := m.ensuredNS[ns]; ok {
|
||
m.ensuredNSMu.Unlock()
|
||
if m.namespaceExists(ctx, ns) {
|
||
return nil // быстрый путь: уже создан и реально существует
|
||
}
|
||
|
||
m.ensuredNSMu.Lock()
|
||
delete(m.ensuredNS, ns)
|
||
m.ensuredNSMu.Unlock()
|
||
}
|
||
if inflight, ok := m.ensuredNSInFlight[ns]; ok {
|
||
m.ensuredNSMu.Unlock()
|
||
select {
|
||
case <-inflight.done:
|
||
return inflight.err
|
||
case <-ctx.Done():
|
||
return ctx.Err()
|
||
}
|
||
}
|
||
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, cancel := context.WithTimeout(ctx, 60*time.Second)
|
||
inflight.err = m.provisionNamespace(ensureCtx, ns)
|
||
cancel()
|
||
<-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
|
||
}
|
||
|
||
func (m *NSManager) namespaceExists(ctx context.Context, ns string) bool {
|
||
_, err := m.dyn.Resource(fission.NamespaceGVR).Get(ctx, ns, metav1.GetOptions{})
|
||
return err == nil
|
||
}
|
||
|
||
// provisionNamespace выполняет полный lifecycle setup нового namespace:
|
||
// 1. fission.SetupFissionNamespace — создаёт NS с меткой fission.io/managed=true
|
||
// + Fission ServiceAccounts + Fission RoleBindings (Fission-специфичная часть)
|
||
// 2. cloud.ApplyResourceQuota — ResourceQuota + LimitRange (облачная политика)
|
||
// 3. cloud.ApplyNetworkPolicy — NetworkPolicy изоляция (облачная политика)
|
||
//
|
||
// Все операции идемпотентны (IsAlreadyExists игнорируется).
|
||
// Ошибки quota/network логируются но не блокируют — NS будет работать и без них.
|
||
func (m *NSManager) provisionNamespace(ctx context.Context, ns string) error {
|
||
// Шаг 1: Fission-специфичный setup (Layer 1 adapter).
|
||
// Создаёт NS с label fission.io/managed=true → executor auto-discovers через NSWatcher.
|
||
if err := fission.SetupFissionNamespace(ctx, m.dyn, ns); err != nil {
|
||
return err
|
||
}
|
||
|
||
// Шаг 2: облачные политики — best-effort (не блокируют запуск функций).
|
||
ApplyResourceQuota(ctx, m.dyn, ns)
|
||
ApplyNetworkPolicy(ctx, m.dyn, ns)
|
||
|
||
return nil
|
||
}
|
||
|
||
// StartExpiryReaper запускает фоновый goroutine для удаления функций с истёкшим TTL.
|
||
// Interval — как часто проверять. Рекомендуемое значение: 5 минут.
|
||
func (m *NSManager) StartExpiryReaper(interval time.Duration) {
|
||
go func() {
|
||
ticker := time.NewTicker(interval)
|
||
defer ticker.Stop()
|
||
log.Printf("cloud.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(fission.NamespaceGVR).List(ctx, metav1.ListOptions{
|
||
LabelSelector: "managed-by=fission-console",
|
||
})
|
||
if err != nil {
|
||
log.Printf("cloud.ExpiryReaper: list namespaces: %v", err)
|
||
return
|
||
}
|
||
|
||
now := time.Now().UTC()
|
||
for _, ns := range nsList.Items {
|
||
m.reapExpiredFunctionsInNS(ctx, ns.GetName(), now)
|
||
}
|
||
}
|
||
|
||
// reapExpiredFunctionsInNS удаляет протухшие функции в конкретном namespace.
|
||
// Для каждой удалённой функции вызывает fission.CleanupEnvironmentIfUnused.
|
||
// Также удаляет orphan packages — пакеты без соответствующей функции.
|
||
func (m *NSManager) reapExpiredFunctionsInNS(ctx context.Context, ns string, now time.Time) {
|
||
functions, err := m.dyn.Resource(fission.FunctionGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
|
||
if err != nil {
|
||
log.Printf("cloud.ExpiryReaper: list functions in %s: %v", ns, err)
|
||
return
|
||
}
|
||
|
||
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
|
||
}
|
||
expiresAt, parseErr := time.Parse(time.RFC3339, expiresAtStr)
|
||
if parseErr != nil {
|
||
log.Printf("cloud.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("cloud.ExpiryReaper: deleting expired function %s/%s (expired %s ago)", ns, fnName, now.Sub(expiresAt).Round(time.Second))
|
||
|
||
triggers, tErr := m.dyn.Resource(fission.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(fission.HTTPTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{})
|
||
}
|
||
}
|
||
}
|
||
|
||
_ = m.dyn.Resource(fission.FunctionGVR).Namespace(ns).Delete(ctx, fnName, metav1.DeleteOptions{})
|
||
if pkgName != "" {
|
||
_ = m.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{})
|
||
}
|
||
if envName != "" {
|
||
fission.CleanupEnvironmentIfUnused(ctx, m.dyn, ns, envName)
|
||
}
|
||
}
|
||
|
||
// Orphan packages — пакеты без соответствующей функции
|
||
packages, pkgListErr := m.dyn.Resource(fission.PackageGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
|
||
if pkgListErr == nil {
|
||
for _, pkg := range packages.Items {
|
||
pkgName := pkg.GetName()
|
||
if !strings.HasSuffix(pkgName, "-pkg") {
|
||
continue
|
||
}
|
||
fnName := strings.TrimSuffix(pkgName, "-pkg")
|
||
if _, exists := activeFunctions[fnName]; !exists {
|
||
log.Printf("cloud.ExpiryReaper: deleting orphan package %s/%s", ns, pkgName)
|
||
_ = m.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{})
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
// IsAlreadyExists — вспомогательная обёртка для читаемости.
|
||
var _ = apierrors.IsAlreadyExists
|