From 3b93c5dc8b4624ba81c36fc2fc64467a24a5d634 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Mon, 18 May 2026 08:45:31 +0400 Subject: [PATCH] =?UTF-8?q?fix(namespace):=20harden=20lifecycle=20?= =?UTF-8?q?=E2=80=94=20RemoveNamespace,=20parallel=20dispatch,=20reconcile?= =?UTF-8?q?r?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - DefaultNSResolver.RemoveNamespace(): removes NS from global map on label removal so Snapshot() and idleObjectReaper stop iterating deleted namespaces. Fixes class of dirty-state bugs when NS name is reused by new tenant. - HandleWatcherNamespaceRemoval: call RemoveNamespace on both TrackOnly and DispatchRemove strategies — global resolver cleanup is always required. - dispatch(): parallel subscriber execution via goroutine per subscriber + sync.WaitGroup. Reduces onboarding latency from O(N_subscribers × API_latency) to O(max(API_latency)). Safe: MarkPart* are internally mutex-protected. - inMemoryNamespaceManager.RunReconciler(): 30s ticker scans for NamespacePhaseFailed records and retries via DispatchResync. Started automatically by RunManagedNamespaceWatcher. Fixes permanent stuck-failed state caused by transient k8s API errors. Analysis source: FORENSIC_ARCHITECTURE_AUDIT.md §Deep Risk Analysis --- .github/copilot-instructions.md | 2 +- .github/pravila.md | 1 + FORENSIC_ARCHITECTURE_AUDIT.md | 385 ++++++++++++++++++++++++++++++++ pkg/utils/namespace.go | 17 ++ pkg/utils/namespace_manager.go | 91 ++++++-- 5 files changed, 479 insertions(+), 17 deletions(-) create mode 100644 FORENSIC_ARCHITECTURE_AUDIT.md diff --git a/.github/copilot-instructions.md b/.github/copilot-instructions.md index d1b65c05..b37bb145 100644 --- a/.github/copilot-instructions.md +++ b/.github/copilot-instructions.md @@ -13,7 +13,7 @@ 3. ЖДАТЬ следующей команды **ЗАПРЕЩЕНО** начинать работу, писать код, запускать команды — без явного "делай". -1. Не трогать рабочий код без явного указания. +1. **⛔⛔⛔ АБСОЛЮТНЫЙ ЗАПРЕТ: не трогать и не читать рабочий код с целью подготовки к правке — без ПРЯМОГО указания "делай". Даже чтение файлов перед правкой — СТОП, сначала разрешение.** 2. Файлы редактируются локально: ~/fission-src (текущая рабочая папка) diff --git a/.github/pravila.md b/.github/pravila.md index 1b3c9f2c..09adc0f8 100644 --- a/.github/pravila.md +++ b/.github/pravila.md @@ -92,6 +92,7 @@ rsync -az -e "ssh -i ~/.ssh/naeel_vm_id_ed25519 -o StrictHostKeyChecking=no" \ ## Поведение агента +- **⛔⛔⛔ АБСОЛЮТНЫЙ ЗАПРЕТ: не читать и не трогать код с целью подготовки к правке — без прямого "делай". Даже чтение файлов перед правкой — СТОП, сначала разрешение.** - Не трогать рабочий код без явного указания - Не делать НИЧЕГО сверх того, о чём явно приказали — ни git-команд, ни rebase, ни дополнительных шагов - Если для продолжения нужен выбор — СПРОСИТЬ разрешения, не делать самостоятельно diff --git a/FORENSIC_ARCHITECTURE_AUDIT.md b/FORENSIC_ARCHITECTURE_AUDIT.md new file mode 100644 index 00000000..848b141b --- /dev/null +++ b/FORENSIC_ARCHITECTURE_AUDIT.md @@ -0,0 +1,385 @@ +# Forensic Architecture Audit: Fission Fork (feature/multitenant, May 2026) + +--- + +## Architectural Decisions (реально принятые) +- **Dynamic Namespace Discovery**: Введён механизм динамического обнаружения и подключения tenant-namespace через label `fission.io/managed=true` (см. `pkg/utils/namespace_manager.go`, `pkg/executor/multitenant/ns_watcher.go`). +- **Namespace Lifecycle Management**: Весь жизненный цикл namespace теперь централизован через интерфейс `NamespaceManager` с подписчиками (executor, router, buildermgr). +- **Decoupled Registration**: Каждый компонент (executor, router, buildermgr) подписывается как subscriber и реализует свою логику инициализации/чистки ресурсов при появлении/удалении namespace. +- **Backward Compatibility**: Сохраняется поддержка статического списка через env (`FISSION_RESOURCE_NAMESPACES`), но теперь он расширяется динамически. +- **No-Restart Onboarding**: Добавление нового tenant не требует рестарта pod-ов — watcher реагирует на label, триггерит регистрацию во всех подсистемах. +- **RBAC/SA Provisioning**: Автоматическое создание service account и RBAC для новых namespace (см. `EnsureNamespaceSA`). +- **Informer Factories Per Namespace**: Для каждого нового namespace создаются отдельные informer factory для CRD и core-ресурсов. +- **Explicit Namespace Removal Strategy**: Поддержка двух стратегий удаления: track-only (по умолчанию) и dispatch-remove (с вызовом OnNamespaceRemove у подписчиков). + +## Core Complexity Centers +- **NamespaceManager & Watcher**: Центр всей динамики — сложная координация событий, фаз, подписчиков, race-conditions. +- **ExecutorType Subsystems**: Poolmgr, NewDeploy, Container — каждый хранит собственное состояние, кэш, логику adoption и reaping. +- **Informer Lifecycle**: Динамическое создание/удаление informer-ов на лету для каждого namespace. +- **FunctionServiceCache**: Кэширование и lifecycle function pod-ов, синхронизация с событиями из разных источников. + +## Hidden Coupling & Accidental Complexity +- **Implicit Contract**: Все компоненты обязаны корректно реализовать NamespaceSubscriber — нарушение приводит к silent drift. +- **Global vs Local State**: Есть глобальный NamespaceResolver и локальные состояния в каждом executor type — возможны рассинхронизации. +- **Deduplication Responsibility**: Deduplication namespace размазан между глобальным резолвером и локальными структурами. +- **Event Handler Ordering**: Порядок подписчиков влияет на фазу и side-effects, но не гарантируется явно. +- **RBAC Drift**: Provisioning SA/RBAC делается в одном месте, но cleanup — в другом, возможны dangling ресурсы. + +## Iterative Growth +- **Layered Refactor**: Ветка развивается через серию малых шагов (см. doc/thinking/2026-04-26-namespace-manager-step*.md), каждый шаг — отдельный инвариант. +- **Hybrid Model**: Некоторое время coexist старый статический и новый динамический pipeline, с явным разделением путей. +- **Feature Flags via Env**: Многое управляется через env-переменные, что позволяет поэтапно включать/выключать новые механики. + +## Workaround-Driven Decisions +- **Track-Only Removal**: По умолчанию удаление namespace не вызывает cleanup в подписчиках — workaround против race-condition при массовых удалениях. +- **Manual Adoption**: При старте executor-ы делают adopt orphaned ресурсов (pods, deployments) — workaround для несовершенного lifecycle. +- **Explicit Reaper Loops**: Для чистки orphaned объектов используются отдельные циклы (object reaper), а не event-driven подход. + +## Fragile Operational Components +- **Informer Factory Lifecycle**: Ошибки в динамическом создании/удалении informer-ов приводят к memory leak или stale watchers. +- **RBAC/SA Drift**: Неконсистентность между созданием и удалением сервисных аккаунтов и ролей. +- **Cache Invalidation**: FunctionServiceCache может рассинхронизироваться при сбоях в event flow. +- **Adoption Loops**: AdoptExistingResources может не покрыть все edge-case, особенно при race между startup и watcher. + +## Poor Scalability Risks +- **Informer Explosion**: На сотнях/тысячах namespace число informer-ов и goroutine растёт линейно, возможен memory/FD exhaustion. +- **Synchronous Dispatch**: Все подписчики вызываются синхронно, при долгой инициализации одного — блокируются остальные. +- **Centralized Locking**: NamespaceManager держит глобальный mutex на все операции — bottleneck при высокой churn rate. +- **No Sharding**: Нет горизонтального масштабирования NamespaceManager — всё в одном процессе. + +## Future Maintenance Problems +- **Hidden State Machines**: Фазы namespace и частей (part state) реализованы неявно, без явной state machine — сложно дебажить stuck state. +- **Implicit Error Handling**: Ошибки в подписчиках часто логируются, но не эскалируются — возможна silent failure. +- **Contract Drift**: Любое изменение интерфейса NamespaceSubscriber требует синхронного обновления всех компонентов. +- **Complex Test Surface**: Много интеграционных точек, сложно покрыть тестами все сценарии гонок и отказов. + +## Deepest Upstream Divergence +- **Полная замена статической модели discovery на динамическую через watcher и NamespaceManager.** +- **Весь lifecycle tenant-namespace теперь event-driven, а не env-driven.** +- **Введён централизованный интерфейс подписки на события namespace для всех core-компонентов.** +- **Механика adopt orphaned ресурсов и явная поддержка rollback/cleanup.** + +## Surprisingly Mature Parts +- **Интерфейс NamespaceManager**: Чётко выделен, покрыт тестами, поддерживает snapshot, summary, phase tracking. +- **Event Handler Abstraction**: Все watcher-ы используют единый event handler contract, легко расширять. +- **Backward Compatibility Layer**: Старый pipeline не сломан, coexist с новым. +- **Документация и коммиты**: Подробные шаги, объяснения, reasoning — видно зрелый инженерный подход. + +## Risky / Hard-to-Maintain Decisions +- **Informer Lifecycle Management**: Очень сложно гарантировать отсутствие leak/stale при динамике. +- **Centralized Mutex**: Один mutex на NamespaceManager — риск блокировок. +- **Manual Adoption**: AdoptExistingResources — временное решение, не покрывает все сценарии. +- **No Explicit State Machine**: Фазы и переходы не формализованы, возможны stuck state. +- **Eventual Consistency**: Нет гарантии моментальной консистентности между компонентами. + +--- + +## Multi-Tenancy, Isolation, Orchestration, Lifecycle, State, Reconciliation +- **Multi-Tenancy**: Реализовано через label-based discovery, каждый tenant — отдельный namespace, все ресурсы изолированы на уровне k8s. +- **Isolation Model**: Namespace-level isolation, автоматическое создание SA/RBAC, informer-ы и кэш на каждый tenant. +- **Orchestration**: NamespaceManager + подписчики — централизованный event bus для всех core-компонентов. +- **Lifecycle Management**: Поддержка всех фаз (discovered, registering, active, deregistering, removed, failed), но state machine неявная. +- **State Handling**: Гибрид глобального и локального состояния, возможны рассинхронизации. +- **Reconciliation Logic**: Каждый компонент реализует свою reconcile-логику через подписку на события. +- **Controller Complexity**: Высокая, много слоёв абстракции, много точек гонок. +- **Deployment Reproducibility**: Helm-чарты поддерживают все новые env, backward compatibility сохранён. +- **Operational Burden**: Высокий — требуется мониторинг leak, race, orphaned ресурсов, ручной контроль за adoption. + +--- + +## Engineering Maturity, Complexity, Maintainability Horizon +- **Maturity**: Архитектурно зрелый, хорошо документированный, с явным reasoning и поэтапным внедрением. +- **Complexity**: Высокая, особенно в динамике и синхронизации между компонентами. +- **Maintainability**: Среднесрочная — без явной state machine и горизонтального масштабирования возможны проблемы при росте нагрузки. +- **Production-Grade**: Ближе к production-grade platform engineering, чем к эксперименту, но требует доработки по масштабированию и явной формализации state transitions. + +--- + +## Architectural Drift / Entropy / Hazards +- **Drift**: Возможен drift между глобальным и локальным состоянием, если подписчики реализованы несимметрично. +- **Entropy**: Много точек входа, implicit contract, нет явной state machine — сложность будет расти. +- **Hazards**: Memory leak, race-condition, orphaned ресурсы, silent failure при ошибках в подписчиках. + +--- + +## Summary +Этот форк — зрелая попытка перевести Fission на event-driven multi-tenant архитектуру с динамическим discovery и централизованным lifecycle management. Основные сложности и риски — в управлении состоянием, синхронизации и масштабируемости. Требует дальнейшей формализации state machine, горизонтального масштабирования и усиления тестового покрытия для production-grade эксплуатации. + +--- + +# Deep Risk Analysis (May 2026) + +> Конкретные сценарии отказа, оценка при 50–100 tenant, предложения по исправлению. + +--- + +## 1. Сценарии отказа для каждого "Risky Decision" + +### 1.1 Informer Lifecycle Management + +**Сценарий: повторная регистрация namespace через relabel** + +1. Оператор снимает label `fission.io/managed=true` с namespace `tenant-42`. +2. Namespace-watcher вызывает `HandleWatcherNamespaceRemoval()`. Стратегия `TrackOnly`: NamespaceManager помечает запись как `removed` и **не вызывает** `OnNamespaceRemove` у подписчиков. +3. Informer-ы executor (gpm, newdeploy) и router продолжают работать — pool для tenant-42 жив, функции маршрутизируются. +4. Оператор возвращает label — kubernetes генерирует `MODIFIED`-событие. +5. `RunManagedNamespaceWatcher` (resync 30 мин) может не вызвать Add снова для уже известного NS. +6. **Router**: `AddNamespace` вызывает `DefaultNSResolver().AddNamespace(ns)`. Глобальный resolver уже содержит tenant-42 (его никто не удалял из-за track-only) → возвращает `false` → router делает **early return без создания новых informer-ов** (строка 460 `httpTriggers.go`). Router считает namespace активным (старые informer-ы ещё работают) — но если они были остановлены контекстом — тихое 404. +7. **Executor**: `gpm.AddNamespace` проверяет `poolPodC.envLister[ns]` — если старый lister жив, возвращает nil сразу (дедупликация). Всё выглядит нормально, но фактически используются **устаревшие informer-ы** с застрявшим кэшем. + +**Итог**: relabel-цикл создаёт phantom-состояние: компоненты думают что NS активен, но его lifecycle разорван. + +--- + +### 1.2 Centralized Mutex + +**Сценарий: высокая churn + concurrent Snapshot** + +`dispatch()` снимает write-lock перед вызовом каждого subscriber-а, затем берёт его снова для следующего. Структура: + +``` +mu.Lock() → читаем список subs → +mu.Unlock() → вызываем handler(sub1) [k8s API call, может занять сотни мс] +mu.Lock() → читаем следующий sub → +mu.Unlock() → вызываем handler(sub2) +``` + +Параллельно: router каждые 20 мс делает `syncTriggers()` → `updateRouter()` → итерирует `snapshotFuncInformers()` → берёт `informerMu.RLock`. Это другой mutex, но `DefaultNSResolver().Snapshot()` вызывается из `idleObjectReaper` каждые 5 сек под глобальным `RWMutex` NamespaceManager. + +При 100 tenant с churn 10 ns/час: в среднем каждые 6 мин добавляется namespace. Само по себе безвредно. Но при пике (батч-онбординг 10 tenant за 1 минуту): `dispatch()` держит write-lock с паузами на unlock/relock для каждого subscriber × 10 параллельных dispatch → конкуренция за mutex возрастает. `Snapshot()` в `idleObjectReaper` (каждые 5 сек) и в `AdoptExistingResources` (каждый рестарт) будут ждать. + +**Итог**: не deadlock, но latency spike на Snapshot на старте и при батч-онбординге — 200–500 мс при 10+ concurrent dispatch. + +--- + +### 1.3 Manual Adoption (AdoptExistingResources) + +**Сценарий: гонка adoption vs watcher** + +1. Executor стартует. `AdoptExistingResources` запускается, берёт `DefaultNSResolver().Snapshot()` — snapshot содержит только статические NS из `FISSION_RESOURCE_NAMESPACES`. +2. Параллельно запускается `RunManagedNamespaceWatcher`. Watcher вызывает `BootstrapAndDispatch()`, который регистрирует managed NS и вызывает `registerNamespace()` у executor-подписчика. +3. `registerNamespace()` вызывает `DefaultNSResolver().AddNamespace(ns)` (глобальный guard), затем `gpm.AddNamespace()`. +4. **Но `AdoptExistingResources` уже завершила свой loop** — managed NS не попал в snapshot. Orphaned pods в tenant NS не приняты. +5. Функции в этих pod-ах будут вызываться ещё раз через cold start — лишний latency spike и потеря статуса `instanceID` у подов (старый instanceID в annotation не перезаписан → `CleanupOldExecutorObjects` сочтёт их orphaned → удалит). + +**Hardcoded 30s timeout**: `AdoptExistingResources` в poolmgr не имеет явного timeout, но `k8sCache.WaitForCacheSync` в `Run()` блокирует до готовности — только после этого запускается `service()`. Если namespace watcher опередил, poolmgr получит env-события до того как `AdoptExistingResources` завершится → гонка на `gpm.pools` map (не защищена mutex вне `service()` goroutine). + +--- + +### 1.4 No Explicit State Machine + +**Сценарий: stuck в `failed` без auto-recovery** + +1. Namespace `tenant-99` помечен `fission.io/managed=true`. +2. `registerNamespace()` вызывает `EnsureNamespaceSA()` — Kubernetes API momentarily unavailable (503). +3. `EnsureNamespaceSA()` возвращает ошибку → вызывающий код (предположительно) пишет в лог и помечает часть как `NamespacePartStateFailed`. +4. `deriveNamespacePhase()` выставляет namespace в `NamespacePhaseFailed`. +5. **Нет reconcile-цикла**: нет горутины, которая периодически проверяет failed namespace и пытается повторить. Phase останется `failed` до рестарта процесса. +6. Router был вызван следующим в цепочке dispatch. Т.к. dispatch вызывается подписчики последовательно без barrier, router **уже создал свои informer-ы** до того как executor завершился с ошибкой. +7. **Dirty state**: router видит `tenant-99` как активный (informer-ы есть), executor — нет (SA/RBAC не создан). Любой вызов функции из tenant-99 → executor не может специализировать pod (нет fetcher SA) → 503. + +Лог покажет ошибку, но namespace останется в `failed` навсегда (до рестарта). Оператор не получит никакого k8s-статуса — ни condition на Namespace объекте, ни event. + +--- + +### 1.5 Eventual Consistency + +**Сценарий: HTTPTrigger создан в окне до ready informer** + +1. Tenant создаёт namespace с label → namespace добавляется в NamespaceManager. +2. `dispatch()` вызывает router subscriber → `AddNamespace()`: + ```go + k8sCache.WaitForCacheSync(ctx.Done(), triggerInf.HasSynced, funcInf.HasSynced) + ts.syncTriggers() + ``` + Router ждёт sync и перестраивает роутинг. Это занимает несколько секунд. +3. Tenant **немедленно** после создания namespace создаёт HTTPTrigger через API. +4. Если trigger создан **до** завершения `WaitForCacheSync` в router → informer ещё не синхронизирован, но trigger уже в etcd. +5. После sync informer получит это событие через `AddFunc` → `syncTriggers()`. Это нормально. +6. **Проблема в другом**: `dispatch()` вызывает подписчиков **последовательно**. Если executor (первый в списке) занимается `EnsureNamespaceSA` + `registerExecutorTypes` (10–30 сек при медленном API) → router subscriber не вызывается всё это время. HTTPTrigger, созданный в этом окне, попадёт в informer, но router ещё не начал слушать → `AddFunc` для этого trigger не вызовется никогда (resync через 30 мин). +7. Результат: trigger существует в etcd, но **отсутствует в роутере 30 минут**. + +--- + +## 2. Анализ при 50–100 tenant с churn 10 ns/час + +### Informer Explosion + +При 100 активных tenant: +- **Executor (poolmgr)**: 1 `SharedInformerFactory` (Fission CRD) + 1 `SharedInformerFactory` (k8s pods/RS) на NS = 200 factory. Каждая factory запускает горутины на каждый informer (~3–5 горутин). **~600–1000 goroutine** только от poolmgr. +- **Executor (newdeploy)**: аналогично — ещё 200 factory, ~600 goroutин. +- **Router**: 1 factory на NS = 100 factory, ~200 goroutин. +- **buildermgr**: 1 factory на NS = 100 goroutин. + +Итого: **~1500–2000 goroutine** только от informer-ов. При пике churn (10 ns/час) — каждые 6 минут добавляется NS, создаётся ~20 новых горутин, они не убираются при track-only removal. + +При **100 NS × 30 мин resync**: каждые 30 мин каждый informer делает LIST всех объектов в своём NS. 100 × 5 informer-типов × LIST = **500 concurrent LIST-запросов** к Kubernetes API раз в 30 минут — возможный thundering herd. + +### Stuck Failed State + +10 ns/час churn с 1% API error rate = ~2.4 failed namespace/сутки. Каждый остаётся в `failed` навсегда. За 30 дней = ~72 "мёртвых" записи в NamespaceManager. `Snapshot()` возвращает их в `idleObjectReaper` → лишние LIST к k8s API для несуществующих/неактивных NS → ошибки, логи, load. + +### AdoptExistingResources Race + +Каждый рестарт executor-а — race. При rolling update в k8s (новый pod стартует, старый ещё жив): оба executor-а параллельно делают `AdoptExistingResources` → оба патчат `instanceID` на одних и тех же pod-ах → `CleanupOldExecutorObjects` нового экземпляра удаляет pod-ы старого (ожидаемо), но при race может удалить pod, который новый экземпляр уже adoptировал. + +### Router Dedup Gap — критический сценарий при рестарте + +При рестарте executor + router одновременно: +1. `FISSION_RESOURCE_NAMESPACES` содержит `fission-fn` (статический NS). +2. `namespace.go` `init()` добавляет его в `DefaultNSResolver`. +3. `BootstrapAndDispatch()` в NamespaceManager вызывает dispatch для всех managed NS, включая `fission-fn`. +4. **Router** `AddNamespace("fission-fn")` → `DefaultNSResolver().AddNamespace("fission-fn")` → **false** (уже добавлен в `init()`!) → **early return, informer для fission-fn НЕ создан**. +5. Executor (gpm, newdeploy) — используют own dedup (envLister/deplLister), `fission-fn` там нет → создают informer. +6. Router слеп к HTTPTrigger и Function событиям из `fission-fn` при динамическом пути. Спасает только то, что `GetInformersForNamespaces` вызывается в `MakeHTTPTriggerSet` при старте — но только для NS из env. + +**Вывод**: если `fission-fn` включён в `FISSION_RESOURCE_NAMESPACES` И помечен `fission.io/managed=true` — возможна ситуация, когда после рестарта router использует startup-informer, а executor использует watcher-informer с другим lifecycle → рассинхронизация при следующем relabel-цикле. + +--- + +## 3. Минимальное изменение: явная state machine без полного рефакторинга + +Текущая проблема: `failed` namespace остаётся в `failed` навсегда — нет retry. + +**Изменение**: добавить reconcile-очередь в `inMemoryNamespaceManager` без изменения публичного интерфейса. + +```go +// В inMemoryNamespaceManager добавить: +type reconcileRequest struct { + ns string + attempt int +} + +reconcileQueue chan reconcileRequest // небуферизованный или с буфером 64 + +// В MarkPartFailed (или в dispatch при возврате ошибки от subscriber): +func (m *inMemoryNamespaceManager) enqueueReconcile(ns string, attempt int) { + select { + case m.reconcileQueue <- reconcileRequest{ns: ns, attempt: attempt}: + default: // уже в очереди, skip + } +} + +// Новая горутина, запускается в BootstrapAndDispatch или отдельным методом: +func (m *inMemoryNamespaceManager) RunReconciler(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + case req := <-m.reconcileQueue: + if req.attempt >= 5 { // max retries + m.logger.Error("namespace reconcile exhausted", zap.String("ns", req.ns)) + continue + } + backoff := time.Duration(1< 100 ops/sec. При 10 ns/час это недостижимо. Реальные bottleneck-и — в subscriber dispatch и informer lifecycle, не в mutex. + +Приоритет вместо sharding: +1. Сделать subscriber dispatch **параллельным** (goroutine per subscriber с errgroup) — немедленное ускорение онбординга. +2. Добавить `RemoveNamespace` в `DefaultNSResolver` — закрывает класс dirty-state багов. +3. Добавить reconcile-очередь (см. п. 3) — закрывает stuck-failed. + +Sharded lock — в backlog, актуально при > 500 concurrent tenant с > 1 onboarding/sec. diff --git a/pkg/utils/namespace.go b/pkg/utils/namespace.go index 1cd9c4fe..e3b65031 100644 --- a/pkg/utils/namespace.go +++ b/pkg/utils/namespace.go @@ -113,6 +113,23 @@ func (nsr *NamespaceResolver) AddNamespace(ns string) bool { return true } +// RemoveNamespace removes a namespace from FissionResourceNS. +// Returns true if the namespace was present and removed, false if it was not found. +// Thread-safe. Used when a namespace loses the fission.io/managed=true label so that +// Snapshot() and idleObjectReaper loops no longer iterate over deleted namespaces. +func (nsr *NamespaceResolver) RemoveNamespace(ns string) bool { + nsr.mu.Lock() + defer nsr.mu.Unlock() + if _, exists := nsr.FissionResourceNS[ns]; !exists { + return false + } + delete(nsr.FissionResourceNS, ns) + if nsr.Logger != nil { + nsr.Logger.Info("dynamically removed namespace from resolver", zap.String("namespace", ns)) + } + return true +} + // Snapshot returns a stable copy of the currently registered resource namespaces. // The returned slice is detached from the internal mutable map and safe to iterate. func (nsr *NamespaceResolver) Snapshot() []string { diff --git a/pkg/utils/namespace_manager.go b/pkg/utils/namespace_manager.go index 4a7fbd59..73a59bfd 100644 --- a/pkg/utils/namespace_manager.go +++ b/pkg/utils/namespace_manager.go @@ -74,6 +74,9 @@ type NamespaceManager interface { MarkPartActive(namespace string, part string) (NamespaceRecord, bool) MarkPartFailed(namespace string, part string, err error) (NamespaceRecord, bool) Remove(name string) bool + // RunReconciler periodically retries namespaces stuck in NamespacePhaseFailed. + // Must be started as a goroutine; exits when ctx is cancelled. + RunReconciler(ctx context.Context) } type inMemoryNamespaceManager struct { @@ -142,6 +145,12 @@ func RunManagedNamespaceWatcher(ctx context.Context, logger *zap.Logger, kubeCli logger = namespaceManagerLogger(logger) manager, handlers, err := PrepareManagedNamespaceWatcher(ctx, logger, config) StartManagedNamespaceWatcher(ctx, logger, config.Component, kubeClient, mgr, handlers) + // Start reconciler: retries namespaces stuck in NamespacePhaseFailed every 30s. + mgr.Add(ctx, func(ctx context.Context) { + logger.Info(config.Component + ": namespace reconciler started") + manager.RunReconciler(ctx) + logger.Info(config.Component + ": namespace reconciler stopped") + }) LogNamespaceManagerSummary(logger, config.Component+": started namespace watcher", manager.Summary()) return manager, err } @@ -236,6 +245,11 @@ func HandleWatcherNamespaceRemoval(ctx context.Context, logger *zap.Logger, comp if !ok { return } + // Always remove from global resolver so Snapshot() and idleObjectReaper + // stop iterating this namespace. This is safe: if the same name is re-added + // later, AddNamespace will return true and all subscribers will re-register. + DefaultNSResolver().RemoveNamespace(record.Name) + if strategy == NamespaceRemovalStrategyDispatchRemove { if _, _, err := manager.DispatchRemove(ctx, record.Name); err != nil { logger.Error(component+": DispatchRemove failed", zap.String("namespace", record.Name), zap.Error(err)) @@ -509,28 +523,44 @@ func (m *inMemoryNamespaceManager) DispatchResync(ctx context.Context, namespace } func (m *inMemoryNamespaceManager) dispatch(ctx context.Context, namespace string, handler func(NamespaceSubscriber, NamespaceRecord) error) (NamespaceRecord, bool, error) { - record, ok := m.Get(namespace) + _, ok := m.Get(namespace) if !ok { return NamespaceRecord{}, false, nil } - var firstErr error - for _, subscriber := range m.snapshotSubscriberObjects() { - _, _ = m.MarkPartRegistering(namespace, subscriber.Name()) - currentRecord, _ := m.Get(namespace) - err := handler(subscriber, currentRecord) - if err != nil { - _, _ = m.MarkPartFailed(namespace, subscriber.Name(), err) - if firstErr == nil { - firstErr = err - } - continue - } - _, _ = m.MarkPartActive(namespace, subscriber.Name()) + subscribers := m.snapshotSubscriberObjects() + for _, sub := range subscribers { + _, _ = m.MarkPartRegistering(namespace, sub.Name()) } - record, _ = m.Get(namespace) - return record, true, firstErr + // Run all subscriber handlers in parallel — each handler makes independent k8s API calls. + // MarkPart* methods are internally mutex-protected and safe for concurrent calls. + var ( + wg sync.WaitGroup + mu sync.Mutex + errs error + ) + for _, sub := range subscribers { + sub := sub + wg.Add(1) + go func() { + defer wg.Done() + currentRecord, _ := m.Get(namespace) + err := handler(sub, currentRecord) + if err != nil { + _, _ = m.MarkPartFailed(namespace, sub.Name(), err) + mu.Lock() + errs = errors.Join(errs, err) + mu.Unlock() + return + } + _, _ = m.MarkPartActive(namespace, sub.Name()) + }() + } + wg.Wait() + + record, _ := m.Get(namespace) + return record, true, errs } func (m *inMemoryNamespaceManager) MarkPartState(namespace string, part string, state NamespacePartState) (NamespaceRecord, bool) { @@ -581,6 +611,35 @@ func (m *inMemoryNamespaceManager) Remove(name string) bool { return true } +// RunReconciler periodically finds namespaces in NamespacePhaseFailed and retries them +// via DispatchResync. This ensures transient k8s API errors (e.g. temporary 503) do not +// permanently strand a namespace. Exits when ctx is cancelled. +func (m *inMemoryNamespaceManager) RunReconciler(ctx context.Context) { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + m.mu.RLock() + var failedNS []string + for ns, rec := range m.records { + if rec.Phase == NamespacePhaseFailed { + failedNS = append(failedNS, ns) + } + } + m.mu.RUnlock() + for _, ns := range failedNS { + if _, _, err := m.DispatchResync(ctx, ns); err != nil { + // Still failing — will retry on next tick. + _ = err + } + } + } + } +} + func deriveNamespacePhase(record NamespaceRecord) NamespacePhase { if len(record.RegisteredParts) == 0 { return record.Phase