diff --git a/pkg/executor/multitenant/namespace_subscriber.go b/pkg/executor/multitenant/namespace_subscriber.go index 087f7612..d57d023f 100644 --- a/pkg/executor/multitenant/namespace_subscriber.go +++ b/pkg/executor/multitenant/namespace_subscriber.go @@ -21,15 +21,16 @@ func NewNamespaceSubscriber( return utils.NamespaceSubscriberFuncs{ SubscriberName: "executor", AddFunc: func(ctx context.Context, record utils.NamespaceRecord) error { - registerNamespace(ctx, logger, kubernetesClient, record.Name, executorTypes, mgr) - return nil + return registerNamespace(ctx, logger, kubernetesClient, record.Name, executorTypes, mgr) }, RemoveFunc: func(ctx context.Context, record utils.NamespaceRecord) error { return deregisterNamespace(ctx, logger, record.Name, executorTypes) }, ResyncFunc: func(ctx context.Context, record utils.NamespaceRecord) error { - registerNamespace(ctx, logger, kubernetesClient, record.Name, executorTypes, mgr) - return nil + // Reconciler calls this for namespaces in NamespacePhaseFailed. + // registerNamespace is idempotent: SA creation is a no-op if SA exists, + // executor type AddNamespace guards against duplicate informer creation. + return registerNamespace(ctx, logger, kubernetesClient, record.Name, executorTypes, mgr) }, } } diff --git a/pkg/executor/multitenant/ns_watcher.go b/pkg/executor/multitenant/ns_watcher.go index 70e68a82..b1544ede 100644 --- a/pkg/executor/multitenant/ns_watcher.go +++ b/pkg/executor/multitenant/ns_watcher.go @@ -62,6 +62,7 @@ package multitenant import ( "context" "errors" + "fmt" "go.uber.org/zap" "k8s.io/client-go/kubernetes" @@ -100,6 +101,10 @@ func StartNSWatcher( // Each executor type uses its own internal state for deduplication instead of the // global resolver, so all executor types receive the AddNamespace call regardless of // iteration order. +// +// Returns an error if SA provisioning or any executor-type initialization fails. +// The error is propagated to the NamespaceSubscriber so the NamespaceManager can +// mark the namespace as NamespacePhaseFailed and the reconciler will retry automatically. func registerNamespace( ctx context.Context, logger *zap.Logger, @@ -107,14 +112,21 @@ func registerNamespace( ns string, executorTypes map[fv1.ExecutorType]executortype.ExecutorType, mgr manager.Interface, -) { +) error { // Update the global resolver once here. Each executor type must NOT call // DefaultNSResolver().AddNamespace() for dedup — they have their own checks. utils.DefaultNSResolver().AddNamespace(ns) // Ensure fission-fetcher SA exists in the new namespace so pool pods can start. - utils.EnsureNamespaceSA(ctx, kubernetesClient, logger, ns) - registerExecutorTypes(ctx, logger, ns, executorTypes, mgr) + // A failure here means function pods will crash (no SA to pull fetcher image) — + // propagate so the reconciler retries until the API is available again. + if err := utils.EnsureNamespaceSA(ctx, kubernetesClient, logger, ns); err != nil { + return fmt.Errorf("EnsureNamespaceSA: %w", err) + } + if err := registerExecutorTypes(ctx, logger, ns, executorTypes, mgr); err != nil { + return fmt.Errorf("registerExecutorTypes: %w", err) + } logger.Info("multitenant.NSWatcher: registered namespace", zap.String("namespace", ns)) + return nil } func registerExecutorTypes( diff --git a/pkg/utils/namespace_manager.go b/pkg/utils/namespace_manager.go index 73a59bfd..435bbbc3 100644 --- a/pkg/utils/namespace_manager.go +++ b/pkg/utils/namespace_manager.go @@ -76,7 +76,8 @@ type NamespaceManager interface { 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) + // logger is used to report retry attempts and outcomes; pass zap.NewNop() to silence. + RunReconciler(ctx context.Context, logger *zap.Logger) } type inMemoryNamespaceManager struct { @@ -148,7 +149,7 @@ func RunManagedNamespaceWatcher(ctx context.Context, logger *zap.Logger, kubeCli // 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) + manager.RunReconciler(ctx, logger) logger.Info(config.Component + ": namespace reconciler stopped") }) LogNamespaceManagerSummary(logger, config.Component+": started namespace watcher", manager.Summary()) @@ -612,9 +613,14 @@ func (m *inMemoryNamespaceManager) Remove(name string) bool { } // 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) { +// via DispatchResync. This ensures transient k8s API errors (e.g. temporary 503 on +// EnsureNamespaceSA or executor-type AddNamespace) do not permanently strand a namespace. +// logger receives one log line per retry attempt and per outcome. +// Exits when ctx is cancelled. +func (m *inMemoryNamespaceManager) RunReconciler(ctx context.Context, logger *zap.Logger) { + if logger == nil { + logger = zap.NewNop() + } ticker := time.NewTicker(30 * time.Second) defer ticker.Stop() for { @@ -630,10 +636,23 @@ func (m *inMemoryNamespaceManager) RunReconciler(ctx context.Context) { } } m.mu.RUnlock() + if len(failedNS) == 0 { + continue + } + logger.Info("namespace reconciler: retrying failed namespaces", + zap.Int("count", len(failedNS)), + zap.Strings("namespaces", failedNS), + ) for _, ns := range failedNS { - if _, _, err := m.DispatchResync(ctx, ns); err != nil { - // Still failing — will retry on next tick. - _ = err + logger.Info("namespace reconciler: dispatching resync", zap.String("namespace", ns)) + _, _, err := m.DispatchResync(ctx, ns) + if err != nil { + logger.Error("namespace reconciler: resync still failing, will retry", + zap.String("namespace", ns), + zap.Error(err), + ) + } else { + logger.Info("namespace reconciler: resync succeeded", zap.String("namespace", ns)) } } } diff --git a/pkg/utils/serviceaccount.go b/pkg/utils/serviceaccount.go index f4937e75..4caba9bd 100644 --- a/pkg/utils/serviceaccount.go +++ b/pkg/utils/serviceaccount.go @@ -121,7 +121,8 @@ func (sa *ServiceAccount) runSACheck(ctx context.Context) { for _, baseNS := range sa.nsResolver.Snapshot() { for _, permission := range sa.permissions { targetNS := sa.resolveSANamespace(baseNS, permission.saName) - setupSAAndRoleBindings(ctx, sa.kubernetesClient, sa.logger, targetNS, permission) + // Errors are already logged inside setupSAAndRoleBindings; periodic loop ignores them. + _ = setupSAAndRoleBindings(ctx, sa.kubernetesClient, sa.logger, targetNS, permission) } } } @@ -134,14 +135,14 @@ func (sa *ServiceAccount) resolveSANamespace(baseNS, saName string) string { return sa.nsResolver.GetFunctionNS(baseNS) } -func setupSAAndRoleBindings(ctx context.Context, client kubernetes.Interface, logger *zap.Logger, namespace string, ps *ServiceAccountPermissions) { +func setupSAAndRoleBindings(ctx context.Context, client kubernetes.Interface, logger *zap.Logger, namespace string, ps *ServiceAccountPermissions) error { SAObj, err := createGetSA(ctx, client, ps.saName, namespace) if err != nil { logger.Error("error while creating or getting service account", zap.String("sa_name", ps.saName), zap.String("namespace", namespace), zap.Error(err)) - return + return err } var rules []rbac.PolicyRule @@ -177,14 +178,15 @@ func setupSAAndRoleBindings(ctx context.Context, client kubernetes.Interface, lo role, err := setupRoles(ctx, client, logger, SAObj, rules, suffix) if err != nil { logger.Error("error while creating roles", zap.Error(err)) - return + return err } _, err = setupRoleBinding(ctx, client, logger, SAObj, role, suffix) if err != nil { logger.Error("error while creating role bindings", zap.Error(err)) - return + return err } } + return nil } func setupRoles(ctx context.Context, client kubernetes.Interface, logger *zap.Logger, sa *v1.ServiceAccount, rules []rbac.PolicyRule, suffix string) (*rbac.Role, error) { @@ -314,8 +316,10 @@ func getSAInterval() time.Duration { } // EnsureNamespaceSA creates the fission-fetcher ServiceAccount and its Role/RoleBinding +// in the given namespace. Returns an error if SA or RoleBinding creation fails so callers +// can propagate it to the namespace lifecycle manager and trigger a reconcile retry. // in the given namespace if they do not already exist. Safe to call repeatedly. // Used by the multi-tenant NS watcher to provision per-namespace SA on NS registration. -func EnsureNamespaceSA(ctx context.Context, client kubernetes.Interface, logger *zap.Logger, ns string) { - setupSAAndRoleBindings(ctx, client, logger, ns, fetcherCheck) +func EnsureNamespaceSA(ctx context.Context, client kubernetes.Interface, logger *zap.Logger, ns string) error { + return setupSAAndRoleBindings(ctx, client, logger, ns, fetcherCheck) }