From 6102b277c86711c383989d25298c170c0a69299f Mon Sep 17 00:00:00 2001 From: Naeel Date: Sun, 26 Apr 2026 09:33:36 +0300 Subject: [PATCH] layer1: migrate runtime loops to snapshots step 3 --- .../2026-04-26-namespace-manager-step3.md | 28 +++++++++++++++++++ .../executortype/container/containermgr.go | 2 +- .../executortype/newdeploy/newdeploymgr.go | 4 +-- pkg/executor/executortype/poolmgr/gpm.go | 8 +++--- pkg/storagesvc/archivePruner.go | 2 +- 5 files changed, 36 insertions(+), 8 deletions(-) create mode 100644 doc/thinking/2026-04-26-namespace-manager-step3.md diff --git a/doc/thinking/2026-04-26-namespace-manager-step3.md b/doc/thinking/2026-04-26-namespace-manager-step3.md new file mode 100644 index 00000000..c58ccd59 --- /dev/null +++ b/doc/thinking/2026-04-26-namespace-manager-step3.md @@ -0,0 +1,28 @@ +# 2026-04-26 — NamespaceManager rewrite, step 3 + +## Цель шага + +Срезать еще один слой прямых чтений `DefaultNSResolver().FissionResourceNS` в runtime code path. + +## Почему это отдельный шаг + +После step 1 snapshot API уже существует, но runtime loops в executor и storagesvc все еще +читают общую mutable map напрямую. Это не архитектурный rewrite, а чистый safety refactor: + +- `container.AdoptExistingResources()` +- `newdeploy.AdoptExistingResources()` +- `newdeploy.doIdleObjectReaper()` +- `poolmgr.AdoptExistingResources()` +- `poolmgr.doIdleObjectReaper()` +- `storagesvc.ArchivePruner.getOrphanArchives()` + +## Что меняем + +В этих местах цикл переводится на `DefaultNSResolver().Snapshot()`. + +## Что НЕ меняем + +- не меняем семантику cleanup; +- не меняем behavior watcher-ов; +- не добавляем remove semantics; +- не исправляем router race и buildermgr dedup на этом шаге. \ No newline at end of file diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index 9601dfb4..70acec7a 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -292,7 +292,7 @@ func (caaf *Container) RefreshFuncPods(ctx context.Context, logger *zap.Logger, func (caaf *Container) AdoptExistingResources(ctx context.Context) { wg := &sync.WaitGroup{} - for _, namepsace := range utils.DefaultNSResolver().FissionResourceNS { + for _, namepsace := range utils.DefaultNSResolver().Snapshot() { fnList, err := caaf.fissionClient.CoreV1().Functions(namepsace).List(ctx, metav1.ListOptions{}) if err != nil { caaf.logger.Error("error getting function list", zap.Error(err)) diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 2f41f72b..c325582d 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -312,7 +312,7 @@ func (deploy *NewDeploy) RefreshFuncPods(ctx context.Context, logger *zap.Logger func (deploy *NewDeploy) AdoptExistingResources(ctx context.Context) { wg := &sync.WaitGroup{} - for _, namepsace := range utils.DefaultNSResolver().FissionResourceNS { + for _, namepsace := range utils.DefaultNSResolver().Snapshot() { fnList, err := deploy.fissionClient.CoreV1().Functions(namepsace).List(ctx, metav1.ListOptions{}) if err != nil { deploy.logger.Error("error getting function list", zap.Error(err)) @@ -782,7 +782,7 @@ func (deploy *NewDeploy) idleObjectReaper(ctx context.Context) { func (deploy *NewDeploy) doIdleObjectReaper(ctx context.Context) { envList := make(map[k8sTypes.UID]struct{}) - for _, namespace := range utils.DefaultNSResolver().FissionResourceNS { + for _, namespace := range utils.DefaultNSResolver().Snapshot() { envs, err := deploy.fissionClient.CoreV1().Environments(namespace).List(ctx, metav1.ListOptions{}) if err != nil { deploy.logger.Fatal("failed to get environment list", zap.Error(err), zap.String("namespace", namespace)) diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index be7a991f..dcf98fb8 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -357,7 +357,7 @@ func (gpm *GenericPoolManager) AdoptExistingResources(ctx context.Context) { envMap := make(map[string]fv1.Environment) wg := &sync.WaitGroup{} - for _, namespace := range utils.DefaultNSResolver().FissionResourceNS { + for _, namespace := range utils.DefaultNSResolver().Snapshot() { envs, err := gpm.fissionClient.CoreV1().Environments(namespace).List(ctx, metav1.ListOptions{}) if err != nil { gpm.logger.Error("error getting environment list", zap.Error(err)) @@ -389,7 +389,7 @@ func (gpm *GenericPoolManager) AdoptExistingResources(ctx context.Context) { fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr), } - for _, namespace := range utils.DefaultNSResolver().FissionResourceNS { + for _, namespace := range utils.DefaultNSResolver().Snapshot() { podList, err := gpm.kubernetesClient.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{ LabelSelector: labels.Set(l).AsSelector().String(), }) @@ -624,7 +624,7 @@ func (gpm *GenericPoolManager) idleObjectReaper(ctx context.Context) { func (gpm *GenericPoolManager) doIdleObjectReaper(ctx context.Context) { envList := make(map[k8sTypes.UID]struct{}) - for _, namespace := range utils.DefaultNSResolver().FissionResourceNS { + for _, namespace := range utils.DefaultNSResolver().Snapshot() { envs, err := gpm.fissionClient.CoreV1().Environments(namespace).List(ctx, metav1.ListOptions{}) if err != nil { gpm.logger.Error("failed to get environment list", zap.Error(err), zap.String("namespace", namespace)) @@ -637,7 +637,7 @@ func (gpm *GenericPoolManager) doIdleObjectReaper(ctx context.Context) { } fnList := make(map[k8sTypes.UID]fv1.Function) - for _, namespace := range utils.DefaultNSResolver().FissionResourceNS { + for _, namespace := range utils.DefaultNSResolver().Snapshot() { fns, err := gpm.fissionClient.CoreV1().Functions(namespace).List(ctx, metav1.ListOptions{}) if err != nil { gpm.logger.Error("failed to get environment list", zap.Error(err), zap.String("namespace", namespace)) diff --git a/pkg/storagesvc/archivePruner.go b/pkg/storagesvc/archivePruner.go index eb670467..56a9ad0e 100644 --- a/pkg/storagesvc/archivePruner.go +++ b/pkg/storagesvc/archivePruner.go @@ -92,7 +92,7 @@ func (pruner *ArchivePruner) getOrphanArchives(ctx context.Context) { var archiveID string // get all pkgs from kubernetes - for _, namespace := range utils.DefaultNSResolver().FissionResourceNS { + for _, namespace := range utils.DefaultNSResolver().Snapshot() { pkgList, err := pruner.crdClient.CoreV1().Packages(namespace).List(ctx, metav1.ListOptions{}) if err != nil { pruner.logger.Error("error getting package list from kubernetes", zap.Error(err))