From 4c319eeed71be92adccf2fabf2f074544f445e6b Mon Sep 17 00:00:00 2001 From: fxiang1 Date: Thu, 24 Sep 2026 17:49:45 -0400 Subject: [PATCH 1/2] Performance fixes Signed-off-by: fxiang1 --- backend/cmd/console/main.go | 4 +- backend/internal/aggregate/argo.go | 1 + backend/internal/aggregate/engine.go | 5 +- backend/internal/aggregate/engine_test.go | 14 +++-- backend/internal/aggregate/handler_test.go | 3 ++ backend/internal/aggregate/rbac.go | 55 ++++++++++++++++++-- backend/internal/aggregate/transform_test.go | 3 ++ backend/internal/auth/auth.go | 36 +++++++++++-- backend/internal/events/rbac/access.go | 53 ++++++++++++++++++- backend/internal/informers/factory.go | 6 +-- backend/internal/informers/factory_test.go | 26 +++++++++ backend/internal/informers/specs.go | 1 + backend/internal/informers/specs_test.go | 12 ++--- frontend/src/components/LoadData.tsx | 3 -- frontend/src/components/LoadEventsData.tsx | 6 +++ 15 files changed, 200 insertions(+), 28 deletions(-) diff --git a/backend/cmd/console/main.go b/backend/cmd/console/main.go index c3f9eee80c..b3ef559079 100644 --- a/backend/cmd/console/main.go +++ b/backend/cmd/console/main.go @@ -95,7 +95,9 @@ func run() error { if err = rbacevents.StartInformer(ctx, kube, store); err != nil { return err } - rbacHandler := rbacevents.NewHandler(store, rbacevents.NewAPIAuth(restCfg), rbacevents.NewSSARAccess(restCfg)) + rbacSSAR := rbacevents.NewSSARAccess(restCfg) + rbacSSAR.StartCleanup(ctx) + rbacHandler := rbacevents.NewHandler(store, rbacevents.NewAPIAuth(restCfg), rbacSSAR) infCfg := informers.RESTConfig(restCfg) infDyn, err := dynamic.NewForConfig(infCfg) diff --git a/backend/internal/aggregate/argo.go b/backend/internal/aggregate/argo.go index 38c1507913..dbf3819141 100644 --- a/backend/internal/aggregate/argo.go +++ b/backend/internal/aggregate/argo.go @@ -259,6 +259,7 @@ func placementFromAppSet(obj map[string]any) string { } func (e *Engine) createArgoStatusMap(search searchapi.ResultBucket, clusters []Cluster) map[string]StatusMap { + e.appStatusByName = map[string]map[string]AppHealthSync{} out := map[string]StatusMap{} ids := map[string]*statusIDs{} sorted := make([]string, 0, len(clusters)) diff --git a/backend/internal/aggregate/engine.go b/backend/internal/aggregate/engine.go index 70ff2b8007..e0153e9250 100644 --- a/backend/internal/aggregate/engine.go +++ b/backend/internal/aggregate/engine.go @@ -177,9 +177,8 @@ func (e *Engine) searchLoop(ctx context.Context) { } func (e *Engine) applications() []App { - e.mu.Lock() - defer e.mu.Unlock() - e.withListCache(e.rebuildSubscriptionLocked) + e.mu.RLock() + defer e.mu.RUnlock() items := getApplicationsHelper(e.cache, cacheKeys) if items == nil { return []App{} diff --git a/backend/internal/aggregate/engine_test.go b/backend/internal/aggregate/engine_test.go index f38cc959aa..70236ad811 100644 --- a/backend/internal/aggregate/engine_test.go +++ b/backend/internal/aggregate/engine_test.go @@ -119,7 +119,7 @@ func (c *countingLister) count(key string) int { return c.n[key] } -func TestApplicationsRebuildsOnlySubscriptions(t *testing.T) { +func TestApplicationsReadOnly(t *testing.T) { cl := &countingLister{inner: MapLister{ "app.k8s.io/v1beta1|Application": { uObj("app.k8s.io/v1beta1", "Application", "sub-app", "ns", nil), @@ -130,15 +130,21 @@ func TestApplicationsRebuildsOnlySubscriptions(t *testing.T) { "cluster.open-cluster-management.io/v1|ManagedCluster": {localCluster()}, }} e := NewEngine(cl, nil, nil) + e.mu.Lock() + e.rebuildLocalLocked() + e.mu.Unlock() + cl.mu.Lock() + cl.n = map[string]int{} + cl.mu.Unlock() e.cache[cacheLocalArgo].Resources = []App{ {Object: map[string]any{"metadata": map[string]any{"name": "cached-argo"}}}, } apps := e.applications() if cl.count("argoproj.io/v1alpha1|Application") != 0 { - t.Fatalf("listed local argo %d", cl.count("argoproj.io/v1alpha1|Application")) + t.Fatalf("applications() should not list anything, got argo %d", cl.count("argoproj.io/v1alpha1|Application")) } - if cl.count("app.k8s.io/v1beta1|Application") != 1 { - t.Fatalf("listed subscription apps %d", cl.count("app.k8s.io/v1beta1|Application")) + if cl.count("app.k8s.io/v1beta1|Application") != 0 { + t.Fatalf("applications() should not list anything, got subscription %d", cl.count("app.k8s.io/v1beta1|Application")) } found := false for _, a := range apps { diff --git a/backend/internal/aggregate/handler_test.go b/backend/internal/aggregate/handler_test.go index 9cb85f3c93..3236dcb0a9 100644 --- a/backend/internal/aggregate/handler_test.go +++ b/backend/internal/aggregate/handler_test.go @@ -25,6 +25,9 @@ func testHandler(t *testing.T, lister Lister) *Handler { eng := NewEngine(lister, nil, nil) zero := 0 eng.PreLimit = &zero + eng.mu.Lock() + eng.rebuildLocalLocked() + eng.mu.Unlock() h := NewHandler(eng, nil, AllowAll{}) h.Authn = testAuthOK return h diff --git a/backend/internal/aggregate/rbac.go b/backend/internal/aggregate/rbac.go index a779e89c6f..ffb27c350a 100644 --- a/backend/internal/aggregate/rbac.go +++ b/backend/internal/aggregate/rbac.go @@ -6,6 +6,7 @@ import ( "context" "crypto/sha256" "encoding/hex" + "errors" "sort" "sync" "time" @@ -55,8 +56,11 @@ type cacheEntry struct { } type tokenState struct { - last time.Time - entries map[ssarKey]cacheEntry + last time.Time + entries map[ssarKey]cacheEntry + client kubernetes.Interface + clientErr error + clientWait chan struct{} } // SSARAccess ports Node getAuthorizedResources / canAccess. @@ -148,6 +152,51 @@ func (a *SSARAccess) canAccessRemote(ctx context.Context, token string, clusters return false, nil } +func (a *SSARAccess) clientFor(token, th string) (kubernetes.Interface, error) { + a.mu.Lock() + st := a.byToken[th] + if st == nil { + st = &tokenState{entries: map[ssarKey]cacheEntry{}} + a.byToken[th] = st + } + if st.client != nil || st.clientErr != nil { + c, err := st.client, st.clientErr + a.mu.Unlock() + return c, err + } + if st.clientWait != nil { + wait := st.clientWait + a.mu.Unlock() + <-wait + a.mu.Lock() + st = a.byToken[th] + if st == nil { + a.mu.Unlock() + return nil, errors.New("token state evicted during client creation") + } + c, err := st.client, st.clientErr + a.mu.Unlock() + return c, err + } + st.clientWait = make(chan struct{}) + a.mu.Unlock() + + client, err := a.newClient(token) + + a.mu.Lock() + st = a.byToken[th] + if st == nil { + st = &tokenState{entries: map[ssarKey]cacheEntry{}} + a.byToken[th] = st + } + st.client = client + st.clientErr = err + close(st.clientWait) + st.clientWait = nil + a.mu.Unlock() + return client, err +} + func (a *SSARAccess) ssar(ctx context.Context, token string, obj map[string]any, verb, name, namespace string) (bool, error) { kind := kindOf(obj) group := apiGroup(apiVersionOf(obj)) @@ -165,7 +214,7 @@ func (a *SSARAccess) ssar(ctx context.Context, token string, obj map[string]any, } a.mu.Unlock() - client, err := a.newClient(token) + client, err := a.clientFor(token, th) if err != nil { return false, err } diff --git a/backend/internal/aggregate/transform_test.go b/backend/internal/aggregate/transform_test.go index 07bb88092f..2e74178199 100644 --- a/backend/internal/aggregate/transform_test.go +++ b/backend/internal/aggregate/transform_test.go @@ -93,6 +93,9 @@ func TestPaginationPerPageAllAndBreakpoint(t *testing.T) { eng := NewEngine(lister, nil, nil) limit := 500 eng.PreLimit = &limit + eng.mu.Lock() + eng.rebuildLocalLocked() + eng.mu.Unlock() h2 := NewHandler(eng, nil, AllowAll{}) h2.Authn = testAuthOK resp2 := postAggregate(t, h2, "/aggregate/applications", RequestListView{Page: 1, PerPage: 10, Search: "zzz"}) diff --git a/backend/internal/auth/auth.go b/backend/internal/auth/auth.go index 96e6e56f07..d14d5720b5 100644 --- a/backend/internal/auth/auth.go +++ b/backend/internal/auth/auth.go @@ -12,6 +12,7 @@ import ( "os" "path/filepath" "strings" + "sync" authv1 "k8s.io/api/authentication/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -158,22 +159,49 @@ func UserRESTConfig(base *rest.Config, userToken string) *rest.Config { return c } +// tokenValidationClient caches the HTTP client used by ValidateUserTokenStatus +// so that a new transport (with its TLS session state and connection pool) is not +// allocated on every request. +var ( + tokenClientMu sync.Mutex + tokenClientHost string + tokenClient *http.Client +) + +func tokenValidationClient(base *rest.Config) (*http.Client, string, error) { + tokenClientMu.Lock() + defer tokenClientMu.Unlock() + host := strings.TrimRight(base.Host, "/") + if tokenClient != nil && tokenClientHost == host { + return tokenClient, host, nil + } + cfg := rest.CopyConfig(base) + cfg.BearerToken = "" + cfg.BearerTokenFile = "" + c, err := rest.HTTPClientFor(cfg) + if err != nil { + return nil, "", err + } + tokenClient = c + tokenClientHost = host + return c, host, nil +} + // ValidateUserTokenStatus probes GET /api with the user token and returns the HTTP status. func ValidateUserTokenStatus(ctx context.Context, base *rest.Config, token string) (int, error) { if base == nil { return 0, errors.New("rest config is required") } - cfg := UserRESTConfig(base, token) - httpClient, err := rest.HTTPClientFor(cfg) + client, host, err := tokenValidationClient(base) if err != nil { return 0, err } - host := strings.TrimRight(cfg.Host, "/") req, err := http.NewRequestWithContext(ctx, http.MethodGet, host+"/api", nil) if err != nil { return 0, err } - resp, err := httpClient.Do(req) + req.Header.Set("Authorization", bearerSchemePrefix+token) + resp, err := client.Do(req) if err != nil { return 0, err } diff --git a/backend/internal/events/rbac/access.go b/backend/internal/events/rbac/access.go index e2e4fd4e18..1006e63dd0 100644 --- a/backend/internal/events/rbac/access.go +++ b/backend/internal/events/rbac/access.go @@ -4,6 +4,7 @@ package rbac import ( "context" + "sort" "sync" "time" @@ -16,7 +17,11 @@ import ( "github.com/stolostron/console/backend/internal/auth" ) -const accessCacheTTL = 60 * time.Second +const ( + accessCacheTTL = 60 * time.Second + accessCleanupEvery = 90 * time.Second + accessCacheMaxSize = 5000 +) // AccessChecker decides whether a user token may see a ClusterRole. type AccessChecker interface { @@ -109,3 +114,49 @@ func (a *SSARAccess) ssar(ctx context.Context, userToken, verb, name string) (bo a.mu.Unlock() return allowed, nil } + +// StartCleanup expires SSAR cache entries periodically. +func (a *SSARAccess) StartCleanup(ctx context.Context) { + if a == nil { + return + } + go func() { + tick := time.NewTicker(accessCleanupEvery) + defer tick.Stop() + for { + select { + case <-ctx.Done(): + return + case <-tick.C: + a.cleanup(time.Now()) + } + } + }() +} + +func (a *SSARAccess) cleanup(now time.Time) { + a.mu.Lock() + defer a.mu.Unlock() + for k, e := range a.cache { + if !e.expiry.After(now) { + delete(a.cache, k) + } + } + if len(a.cache) <= accessCacheMaxSize { + return + } + // Evict oldest entries when the cache exceeds the size limit. + type pair struct { + key cacheKey + expiry time.Time + } + all := make([]pair, 0, len(a.cache)) + for k, e := range a.cache { + all = append(all, pair{k, e.expiry}) + } + sort.Slice(all, func(i, j int) bool { return all[i].expiry.Before(all[j].expiry) }) + extra := len(all) - accessCacheMaxSize + for i := 0; i < extra; i++ { + delete(a.cache, all[i].key) + } +} diff --git a/backend/internal/informers/factory.go b/backend/internal/informers/factory.go index 9e19d688f3..85b15859b7 100644 --- a/backend/internal/informers/factory.go +++ b/backend/internal/informers/factory.go @@ -105,9 +105,9 @@ func (c *InformerCache) runSpec(ctx context.Context, dyn dynamic.Interface, mapp } else { applog.Logger().Warn("informer GVR resolve failed; retrying", "kind", st.spec.Kind, "apiVersion", st.spec.APIVersion, "error", err) - } - if inv, ok := mapper.(CacheInvalidator); ok { - inv.Invalidate() + if inv, ok := mapper.(CacheInvalidator); ok { + inv.Invalidate() + } } if !waitRetry(ctx) { return diff --git a/backend/internal/informers/factory_test.go b/backend/internal/informers/factory_test.go index 5fee8b0ccf..5fb5b4ac72 100644 --- a/backend/internal/informers/factory_test.go +++ b/backend/internal/informers/factory_test.go @@ -374,6 +374,32 @@ func TestStaleDiscoveryCacheInvalidatedOnRetry(t *testing.T) { t.Fatal("expected mapper.Invalidate() to be called on stale discovery error") } +type unavailableMapper struct { + invalidated atomic.Bool +} + +func (m *unavailableMapper) ServerResourcesForGroupVersion(string) (*metav1.APIResourceList, error) { + return nil, apierrors.NewNotFound(schema.GroupResource{Group: "tower.ansible.com", Resource: "ansiblejobs"}, "") +} + +func (m *unavailableMapper) Invalidate() { + m.invalidated.Store(true) +} + +func TestUnavailableCRDDoesNotInvalidateCache(t *testing.T) { + mapper := &unavailableMapper{} + + ctx, cancel := context.WithCancel(context.Background()) + _ = StartSpecs(ctx, nil, mapper, []WatchSpec{watch("AnsibleJob", "tower.ansible.com/v1alpha1")}) + + time.Sleep(200 * time.Millisecond) + cancel() + + if mapper.invalidated.Load() { + t.Fatal("Invalidate() should not be called for unavailable CRDs") + } +} + func TestStartConcurrencyLimitsLists(t *testing.T) { orig := startConcurrency startConcurrency = 2 diff --git a/backend/internal/informers/specs.go b/backend/internal/informers/specs.go index 384477da6a..1174056f5c 100644 --- a/backend/internal/informers/specs.go +++ b/backend/internal/informers/specs.go @@ -138,6 +138,7 @@ func DefaultWatchSpecs() []WatchSpec { watch("Secret", "v1").fields("metadata.name", "auto-import-secret"), watch("Secret", "v1").labels("argocd.argoproj.io/secret-type", "repository"), watch("PolicyReport", "wgpolicyk8s.io/v1alpha2"), + watch("ClusterRole", "rbac.authorization.k8s.io/v1").labels("rbac.open-cluster-management.io/filter", "vm-clusterroles"), watch("HostedCluster", "hypershift.openshift.io/v1beta1"), watch("NodePool", "hypershift.openshift.io/v1beta1"), watch("AgentMachine", "capi-provider.agent-install.openshift.io/v1alpha1"), diff --git a/backend/internal/informers/specs_test.go b/backend/internal/informers/specs_test.go index d1a26e5846..4599196306 100644 --- a/backend/internal/informers/specs_test.go +++ b/backend/internal/informers/specs_test.go @@ -8,8 +8,8 @@ import ( func TestDefaultWatchSpecsCount(t *testing.T) { specs := DefaultWatchSpecs() - if len(specs) != 68 { - t.Fatalf("got %d specs, want 68", len(specs)) + if len(specs) != 69 { + t.Fatalf("got %d specs, want 69", len(specs)) } var polled, cacheOnly, withSel int for _, s := range specs { @@ -29,8 +29,8 @@ func TestDefaultWatchSpecsCount(t *testing.T) { if cacheOnly != 2 { t.Fatalf("cacheOnly=%d want 2 (Authentication, MultiClusterHub)", cacheOnly) } - if withSel != 12 { - t.Fatalf("selector specs=%d want 12", withSel) + if withSel != 13 { + t.Fatalf("selector specs=%d want 13", withSel) } } @@ -92,8 +92,8 @@ func TestDefaultWatchSpecsShouldForwardCount(t *testing.T) { skip++ } } - if forward != 64 { - t.Fatalf("forward=%d want 64", forward) + if forward != 65 { + t.Fatalf("forward=%d want 65", forward) } if skip != 4 { t.Fatalf("skip=%d want 4 (2 polled + 2 cacheOnly)", skip) diff --git a/frontend/src/components/LoadData.tsx b/frontend/src/components/LoadData.tsx index c7bac777dd..94ad44741c 100644 --- a/frontend/src/components/LoadData.tsx +++ b/frontend/src/components/LoadData.tsx @@ -1,17 +1,14 @@ /* Copyright Contributors to the Open Cluster Management project */ import { ReactNode } from 'react' import { LoadEventsData } from './LoadEventsData' -import { LoadRbacData } from './LoadRbacData' /** * Composition root for backend event streams. - * One business domain → one GET /events/ → one LoadXxxData → one line here. */ export function LoadData(props: Readonly<{ children?: ReactNode }>) { return ( <> - {props.children} ) diff --git a/frontend/src/components/LoadEventsData.tsx b/frontend/src/components/LoadEventsData.tsx index 1345932ae6..2310f67736 100644 --- a/frontend/src/components/LoadEventsData.tsx +++ b/frontend/src/components/LoadEventsData.tsx @@ -108,6 +108,8 @@ import { SubscriptionOperatorKind, ClusterExtensionApiVersion, ClusterExtensionKind, + ClusterRoleKind, + RbacApiVersion, SubscriptionReportApiVersion, SubscriptionReportKind, UserApiVersion, @@ -181,6 +183,7 @@ import { subscriptionReportsState, subscriptionsState, usersState, + vmClusterRolesState, WatchEvent, } from '../atoms' import { applyWatchEventsToCache, groupWatchEventsByKind } from '../hooks/applyWatchEventsToCache' @@ -258,6 +261,7 @@ export function LoadEventsData() { const setSubscriptionReportsState = useSetRecoilState(subscriptionReportsState) const setSubscriptionsState = useSetRecoilState(subscriptionsState) const setUsers = useSetRecoilState(usersState) + const setVMClusterRoles = useSetRecoilState(vmClusterRolesState) const { setters, mappers, caches } = useMemo(() => { const setters: Record>> = {} @@ -343,6 +347,7 @@ export function LoadEventsData() { addSetter(PolicyAutomationApiVersion, PolicyAutomationKind, setPolicyAutomationState) addSetter(PolicyReportApiVersion, PolicyReportKind, setPolicyReports) addSetter(PolicySetApiVersion, PolicySetKind, setPolicySetsState) + addSetter(RbacApiVersion, ClusterRoleKind, setVMClusterRoles) addSetter(SearchOperatorApiVersion, SearchOperatorKind, setSearchOperator) addSetter(SecretApiVersion, SecretKind, setSecrets) addSetter(ServiceApiVersion, ServiceKind, setServices) @@ -413,6 +418,7 @@ export function LoadEventsData() { setSubscriptionReportsState, setSubscriptionsState, setUsers, + setVMClusterRoles, ]) const applyWatchEvents = useCallback( From d78b35876a67d17a7609e0ac679e7829eadda0fb Mon Sep 17 00:00:00 2001 From: Enrique Mingorance Cano Date: Fri, 25 Sep 2026 03:10:40 +0200 Subject: [PATCH 2/2] frontend changes from #46 reverted back Signed-off-by: Enrique Mingorance Cano --- frontend/src/components/LoadData.tsx | 833 +++++++++++++++++- .../src/components/LoadDataAbstract.test.tsx | 127 --- frontend/src/components/LoadDataAbstract.tsx | 186 ---- frontend/src/components/LoadEventsData.tsx | 720 --------------- .../src/components/LoadPluginData.test.tsx | 10 +- frontend/src/components/LoadRbacData.test.tsx | 142 --- frontend/src/components/LoadRbacData.tsx | 17 - .../src/hooks/applyWatchEventsToCache.test.ts | 66 -- frontend/src/hooks/applyWatchEventsToCache.ts | 37 - .../src/hooks/useWatchEventStream.test.ts | 140 --- frontend/src/hooks/useWatchEventStream.ts | 113 --- frontend/src/lib/test-event-source.ts | 48 - .../src/resources/utils/resource-request.ts | 2 +- frontend/webpack.config.ts | 1 - 14 files changed, 824 insertions(+), 1618 deletions(-) delete mode 100644 frontend/src/components/LoadDataAbstract.test.tsx delete mode 100644 frontend/src/components/LoadDataAbstract.tsx delete mode 100644 frontend/src/components/LoadEventsData.tsx delete mode 100644 frontend/src/components/LoadRbacData.test.tsx delete mode 100644 frontend/src/components/LoadRbacData.tsx delete mode 100644 frontend/src/hooks/applyWatchEventsToCache.test.ts delete mode 100644 frontend/src/hooks/applyWatchEventsToCache.ts delete mode 100644 frontend/src/hooks/useWatchEventStream.test.ts delete mode 100644 frontend/src/hooks/useWatchEventStream.ts delete mode 100644 frontend/src/lib/test-event-source.ts diff --git a/frontend/src/components/LoadData.tsx b/frontend/src/components/LoadData.tsx index 94ad44741c..61781b7378 100644 --- a/frontend/src/components/LoadData.tsx +++ b/frontend/src/components/LoadData.tsx @@ -1,15 +1,824 @@ /* Copyright Contributors to the Open Cluster Management project */ -import { ReactNode } from 'react' -import { LoadEventsData } from './LoadEventsData' - -/** - * Composition root for backend event streams. - */ -export function LoadData(props: Readonly<{ children?: ReactNode }>) { - return ( - <> - - {props.children} - +import get from 'lodash/get' +import { Fragment, ReactNode, useCallback, useContext, useEffect, useMemo, useRef, useState } from 'react' +// eslint-disable-next-line @typescript-eslint/no-restricted-imports +import { SetterOrUpdater, useRecoilValue, useSetRecoilState } from 'recoil' +import { tokenExpired } from '../logout' +import { + AgentClusterInstallApiVersion, + AgentClusterInstallKind, + AgentKind, + AgentKindVersion, + AgentMachineApiVersion, + AgentMachineKind, + AgentServiceConfigKind, + AgentServiceConfigKindVersion, + AnsibleJobApiVersion, + AnsibleJobKind, + AnsibleWorkflowKind, + ApplicationApiVersion, + ApplicationKind, + BareMetalHostApiVersion, + BareMetalHostKind, + CertificateSigningRequestApiVersion, + CertificateSigningRequestKind, + ChannelApiVersion, + ChannelKind, + ClusterClaimApiVersion, + ClusterClaimKind, + ClusterCuratorApiVersion, + ClusterCuratorKind, + ClusterDeploymentApiVersion, + ClusterDeploymentKind, + ClusterImageSetApiVersion, + ClusterImageSetKind, + ClusterManagementAddOnApiVersion, + ClusterManagementAddOnKind, + ClusterPoolApiVersion, + ClusterPoolKind, + ClusterProvisionApiVersion, + ClusterProvisionKind, + ClusterRoleKind, + ClusterVersionApiVersion, + ClusterVersionKind, + ConfigMapApiVersion, + ConfigMapKind, + DiscoveredClusterApiVersion, + DiscoveredClusterKind, + DiscoveryConfigApiVersion, + DiscoveryConfigKind, + GitOpsClusterApiVersion, + GitOpsClusterKind, + GroupKind, + HelmReleaseApiVersion, + HelmReleaseKind, + HostedClusterApiVersion, + HostedClusterKind, + InfraEnvApiVersion, + InfraEnvKind, + InfrastructureApiVersion, + InfrastructureKind, + IResource, + MachinePoolApiVersion, + MachinePoolKind, + ManagedClusterAddOnApiVersion, + ManagedClusterAddOnKind, + ManagedClusterApiVersion, + ManagedClusterInfoApiVersion, + ManagedClusterInfoKind, + ManagedClusterKind, + ManagedClusterSetApiVersion, + ManagedClusterSetBindingApiVersion, + ManagedClusterSetBindingKind, + ManagedClusterSetKind, + MulticlusterRoleAssignmentApiVersion, + MulticlusterRoleAssignmentKind, + MultiClusterEngineApiVersion, + MultiClusterEngineKind, + NamespaceApiVersion, + NamespaceKind, + NMStateConfigApiVersion, + NMStateConfigKind, + NodePoolApiVersion, + NodePoolKind, + PlacementApiVersionAlpha, + PlacementBindingApiVersion, + PlacementBindingKind, + PlacementDecisionApiVersion, + PlacementDecisionKind, + PlacementKind, + PolicyApiVersion, + PolicyAutomationApiVersion, + PolicyAutomationKind, + PolicyKind, + PolicyReportApiVersion, + PolicyReportKind, + PolicySetApiVersion, + PolicySetKind, + RbacApiVersion, + SearchOperatorApiVersion, + SearchOperatorKind, + SecretApiVersion, + SecretKind, + StorageClassApiVersion, + StorageClassKind, + SubmarinerConfigApiVersion, + SubmarinerConfigKind, + SubscriptionApiVersion, + SubscriptionKind, + SubscriptionOperatorApiVersion, + SubscriptionOperatorKind, + ClusterExtensionApiVersion, + ClusterExtensionKind, + SubscriptionReportApiVersion, + SubscriptionReportKind, + UserApiVersion, + UserKind, + ServiceApiVersion, + ServiceKind, +} from '../resources' +import { getBackendUrl, getRequest } from '../resources/utils' +// eslint-disable-next-line @typescript-eslint/no-restricted-imports +import { + agentClusterInstallsState, + agentMachinesState, + agentServiceConfigsState, + agentsState, + ansibleJobState, + ansibleWorkflowState, + applicationsState, + argoCDsState, + bareMetalHostsState, + certificateSigningRequestsState, + channelsState, + claimMappingsState, + clusterClaimsState, + clusterCuratorsState, + clusterDeploymentsState, + clusterImageSetsState, + clusterManagementAddonsState, + clusterPoolsState, + clusterProvisionsState, + clusterVersionState, + configMapsState, + discoveredClusterState, + discoveryConfigState, + gitOpsClustersState, + groupsState, + helmReleaseState, + hostedClustersState, + infraEnvironmentsState, + infrastructuresState, + isDirectAuthenticationEnabledState, + isFineGrainedRbacEnabledState, + isGlobalHubState, + isHubSelfManagedState, + localHubNameState, + machinePoolsState, + managedClusterAddonsState, + managedClusterInfosState, + managedClusterSetBindingsState, + managedClusterSetsState, + managedClustersState, + multiClusterEnginesState, + multiclusterRoleAssignmentState, + namespacesState, + nmStateConfigsState, + nodePoolsState, + placementBindingsState, + placementDecisionsState, + placementsState, + policiesState, + policyAutomationState, + policyreportState, + policySetsState, + searchOperatorState, + secretsState, + ServerSideEventData, + settingsState, + useEventStreamIdleGracePeriod, + useEventStreamIdleTimeout, + servicesState, + storageClassState, + submarinerConfigsState, + subscriptionOperatorsState, + clusterExtensionsState, + subscriptionReportsState, + subscriptionsState, + usersState, + vmClusterRolesState, + WatchEvent, +} from '../atoms' +import { PluginDataContext } from '../lib/PluginDataContext' +import { useQuery } from '../lib/useQuery' +import { MultiClusterHubComponent } from '../resources/multi-cluster-hub-component' +import { ClaimMappings } from '~/resources/authentication' +import { usePageActivity } from '../lib/usePageActivity' + +export function LoadData(props: { children?: ReactNode }) { + const { loadCompleted, setLoadStarted, setLoadCompleted, setIsStreamIdle, setIsReconnecting, mounted } = + useContext(PluginDataContext) + const [eventsLoaded, setEventsLoaded] = useState(false) + const idleTimeoutMs = useEventStreamIdleTimeout() + const gracePeriodMs = useEventStreamIdleGracePeriod() + const { isActive } = usePageActivity(idleTimeoutMs, mounted) + const wasActiveRef = useRef(true) + const isReconnectingRef = useRef(false) + const streamStoppedRef = useRef(false) + const graceTimerRef = useRef>() + const eventSourceRef = useRef() + const processIntervalRef = useRef>() + const [restartKey, setRestartKey] = useState(0) + + const setAgentClusterInstalls = useSetRecoilState(agentClusterInstallsState) + const setAgentMachinesState = useSetRecoilState(agentMachinesState) + const setAgents = useSetRecoilState(agentsState) + const setAgentServiceConfigs = useSetRecoilState(agentServiceConfigsState) + const setAnsibleJobs = useSetRecoilState(ansibleJobState) + const setAnsibleWorkflows = useSetRecoilState(ansibleWorkflowState) + const setApplicationsState = useSetRecoilState(applicationsState) + const setArgoCDsState = useSetRecoilState(argoCDsState) + const setBareMetalHosts = useSetRecoilState(bareMetalHostsState) + const setCertificateSigningRequests = useSetRecoilState(certificateSigningRequestsState) + const setChannelsState = useSetRecoilState(channelsState) + const setClusterClaims = useSetRecoilState(clusterClaimsState) + const setClusterCurators = useSetRecoilState(clusterCuratorsState) + const setClusterDeployments = useSetRecoilState(clusterDeploymentsState) + const setClusterImageSets = useSetRecoilState(clusterImageSetsState) + const setClusterManagementAddons = useSetRecoilState(clusterManagementAddonsState) + const setClusterPools = useSetRecoilState(clusterPoolsState) + const setClusterProvisions = useSetRecoilState(clusterProvisionsState) + const setVMClusterRoles = useSetRecoilState(vmClusterRolesState) + const setClusterVerions = useSetRecoilState(clusterVersionState) + const setConfigMaps = useSetRecoilState(configMapsState) + const setDiscoveredClusters = useSetRecoilState(discoveredClusterState) + const setDiscoveryConfigs = useSetRecoilState(discoveryConfigState) + const setGitOpsClustersState = useSetRecoilState(gitOpsClustersState) + const setGroups = useSetRecoilState(groupsState) + const setHelmReleases = useSetRecoilState(helmReleaseState) + const setHostedClustersState = useSetRecoilState(hostedClustersState) + const setInfraEnvironments = useSetRecoilState(infraEnvironmentsState) + const setInfrastructure = useSetRecoilState(infrastructuresState) + const setClaimMappings = useSetRecoilState(claimMappingsState) + const setIsDirectAuthenticationEnabled = useSetRecoilState(isDirectAuthenticationEnabledState) + const setIsFineGrainedRbacEnabled = useSetRecoilState(isFineGrainedRbacEnabledState) + const setIsGlobalHub = useSetRecoilState(isGlobalHubState) + const setIsHubSelfManaged = useSetRecoilState(isHubSelfManagedState) + const setlocalHubName = useSetRecoilState(localHubNameState) + const setMachinePools = useSetRecoilState(machinePoolsState) + const setManagedClusterAddons = useSetRecoilState(managedClusterAddonsState) + const setManagedClusterInfos = useSetRecoilState(managedClusterInfosState) + const setManagedClusterSetBindings = useSetRecoilState(managedClusterSetBindingsState) + const setManagedClusterSets = useSetRecoilState(managedClusterSetsState) + const setManagedClusters = useSetRecoilState(managedClustersState) + const setMultiClusterEngines = useSetRecoilState(multiClusterEnginesState) + const setMulticlusterRoleAssignments = useSetRecoilState(multiclusterRoleAssignmentState) + const setNamespaces = useSetRecoilState(namespacesState) + const setNMStateConfigs = useSetRecoilState(nmStateConfigsState) + const setNodePoolsState = useSetRecoilState(nodePoolsState) + const setPlacementBindingsState = useSetRecoilState(placementBindingsState) + const setPlacementDecisionsState = useSetRecoilState(placementDecisionsState) + const setPlacementsState = useSetRecoilState(placementsState) + const setPoliciesState = useSetRecoilState(policiesState) + const setPolicyAutomationState = useSetRecoilState(policyAutomationState) + const setPolicyReports = useSetRecoilState(policyreportState) + const setPolicySetsState = useSetRecoilState(policySetsState) + const setSearchOperator = useSetRecoilState(searchOperatorState) + const setSecrets = useSetRecoilState(secretsState) + const setSettings = useSetRecoilState(settingsState) + const setServices = useSetRecoilState(servicesState) + const setStorageClassState = useSetRecoilState(storageClassState) + const setSubmarinerConfigs = useSetRecoilState(submarinerConfigsState) + const setSubscriptionOperatorsState = useSetRecoilState(subscriptionOperatorsState) + const setClusterExtensionsState = useSetRecoilState(clusterExtensionsState) + const setSubscriptionReportsState = useSetRecoilState(subscriptionReportsState) + const setSubscriptionsState = useSetRecoilState(subscriptionsState) + const setUsers = useSetRecoilState(usersState) + + const { setters, mappers, caches } = useMemo(() => { + const setters: Record>> = {} + + const mappers: Record< + string, + Record< + string, + { + setter: SetterOrUpdater> + mcaches: Record>> + keyBy: string[] + } + > + > = {} + const caches: Record>> = {} + const mcaches: Record>> = {} + function addSetter(apiVersion: string, kind: string, setter: SetterOrUpdater) { + const groupVersion = apiVersion.split('/')[0] + if (!setters[groupVersion]) setters[groupVersion] = {} + setters[groupVersion][kind] = setter + if (!caches[groupVersion]) caches[groupVersion] = {} + caches[groupVersion][kind] = {} + } + function addMapper( + apiVersion: string, + kind: string, + setter: SetterOrUpdater>, + keyBy: string[] + ) { + const groupVersion = apiVersion.split('/')[0] + if (!mappers[groupVersion]) mappers[groupVersion] = {} + if (!mcaches[groupVersion]) mcaches[groupVersion] = {} + mcaches[groupVersion][kind] = {} + mappers[groupVersion][kind] = { setter, mcaches, keyBy } + } + + // mappers (key=>[values]) + addMapper(ManagedClusterAddOnApiVersion, ManagedClusterAddOnKind, setManagedClusterAddons, ['metadata.namespace']) + + // setters + addSetter('argoproj.io/v1alpha1', 'ArgoCD', setArgoCDsState) + addSetter(AgentClusterInstallApiVersion, AgentClusterInstallKind, setAgentClusterInstalls) + addSetter(AgentKindVersion, AgentKind, setAgents) + addSetter(AgentMachineApiVersion, AgentMachineKind, setAgentMachinesState) + addSetter(AgentServiceConfigKindVersion, AgentServiceConfigKind, setAgentServiceConfigs) + addSetter(AnsibleJobApiVersion, AnsibleJobKind, setAnsibleJobs) + addSetter(AnsibleJobApiVersion, AnsibleWorkflowKind, setAnsibleWorkflows) + addSetter(ApplicationApiVersion, ApplicationKind, setApplicationsState) + addSetter(BareMetalHostApiVersion, BareMetalHostKind, setBareMetalHosts) + addSetter(CertificateSigningRequestApiVersion, CertificateSigningRequestKind, setCertificateSigningRequests) + addSetter(ChannelApiVersion, ChannelKind, setChannelsState) + addSetter(ClusterClaimApiVersion, ClusterClaimKind, setClusterClaims) + addSetter(ClusterCuratorApiVersion, ClusterCuratorKind, setClusterCurators) + addSetter(ClusterDeploymentApiVersion, ClusterDeploymentKind, setClusterDeployments) + addSetter(ClusterImageSetApiVersion, ClusterImageSetKind, setClusterImageSets) + addSetter(ClusterManagementAddOnApiVersion, ClusterManagementAddOnKind, setClusterManagementAddons) + addSetter(ClusterPoolApiVersion, ClusterPoolKind, setClusterPools) + addSetter(ClusterProvisionApiVersion, ClusterProvisionKind, setClusterProvisions) + addSetter(ClusterVersionApiVersion, ClusterVersionKind, setClusterVerions) + addSetter(ConfigMapApiVersion, ConfigMapKind, setConfigMaps) + addSetter(DiscoveredClusterApiVersion, DiscoveredClusterKind, setDiscoveredClusters) + addSetter(DiscoveryConfigApiVersion, DiscoveryConfigKind, setDiscoveryConfigs) + addSetter(GitOpsClusterApiVersion, GitOpsClusterKind, setGitOpsClustersState) + addSetter(HelmReleaseApiVersion, HelmReleaseKind, setHelmReleases) + addSetter(HostedClusterApiVersion, HostedClusterKind, setHostedClustersState) + addSetter(InfraEnvApiVersion, InfraEnvKind, setInfraEnvironments) + addSetter(InfrastructureApiVersion, InfrastructureKind, setInfrastructure) + addSetter(MachinePoolApiVersion, MachinePoolKind, setMachinePools) + addSetter(ManagedClusterApiVersion, ManagedClusterKind, setManagedClusters) + addSetter(ManagedClusterInfoApiVersion, ManagedClusterInfoKind, setManagedClusterInfos) + addSetter(ManagedClusterSetApiVersion, ManagedClusterSetKind, setManagedClusterSets) + addSetter(ManagedClusterSetBindingApiVersion, ManagedClusterSetBindingKind, setManagedClusterSetBindings) + addSetter(MulticlusterRoleAssignmentApiVersion, MulticlusterRoleAssignmentKind, setMulticlusterRoleAssignments) + addSetter(MultiClusterEngineApiVersion, MultiClusterEngineKind, setMultiClusterEngines) + addSetter(NamespaceApiVersion, NamespaceKind, setNamespaces) + addSetter(NMStateConfigApiVersion, NMStateConfigKind, setNMStateConfigs) + addSetter(NodePoolApiVersion, NodePoolKind, setNodePoolsState) + addSetter(PlacementApiVersionAlpha, PlacementKind, setPlacementsState) + addSetter(PlacementBindingApiVersion, PlacementBindingKind, setPlacementBindingsState) + addSetter(PlacementDecisionApiVersion, PlacementDecisionKind, setPlacementDecisionsState) + addSetter(PolicyApiVersion, PolicyKind, setPoliciesState) + addSetter(PolicyAutomationApiVersion, PolicyAutomationKind, setPolicyAutomationState) + addSetter(PolicyReportApiVersion, PolicyReportKind, setPolicyReports) + addSetter(PolicySetApiVersion, PolicySetKind, setPolicySetsState) + addSetter(RbacApiVersion, ClusterRoleKind, setVMClusterRoles) + addSetter(SearchOperatorApiVersion, SearchOperatorKind, setSearchOperator) + addSetter(SecretApiVersion, SecretKind, setSecrets) + addSetter(ServiceApiVersion, ServiceKind, setServices) + addSetter(StorageClassApiVersion, StorageClassKind, setStorageClassState) + addSetter(SubmarinerConfigApiVersion, SubmarinerConfigKind, setSubmarinerConfigs) + addSetter(SubscriptionApiVersion, SubscriptionKind, setSubscriptionsState) + addSetter(SubscriptionOperatorApiVersion, SubscriptionOperatorKind, setSubscriptionOperatorsState) + addSetter(ClusterExtensionApiVersion, ClusterExtensionKind, setClusterExtensionsState) + addSetter(SubscriptionReportApiVersion, SubscriptionReportKind, setSubscriptionReportsState) + addSetter(UserApiVersion, GroupKind, setGroups) + addSetter(UserApiVersion, UserKind, setUsers) + + return { setters, mappers, caches } + }, [ + setAgentClusterInstalls, + setAgentMachinesState, + setAgents, + setAgentServiceConfigs, + setAnsibleJobs, + setAnsibleWorkflows, + setApplicationsState, + setArgoCDsState, + setBareMetalHosts, + setCertificateSigningRequests, + setChannelsState, + setClusterClaims, + setClusterCurators, + setClusterDeployments, + setClusterImageSets, + setClusterManagementAddons, + setClusterPools, + setClusterProvisions, + setVMClusterRoles, + setClusterVerions, + setConfigMaps, + setDiscoveredClusters, + setDiscoveryConfigs, + setGitOpsClustersState, + setGroups, + setHelmReleases, + setHostedClustersState, + setInfraEnvironments, + setInfrastructure, + setMachinePools, + setManagedClusterAddons, + setManagedClusterInfos, + setManagedClusterSetBindings, + setManagedClusterSets, + setManagedClusters, + setMultiClusterEngines, + setMulticlusterRoleAssignments, + setNamespaces, + setNMStateConfigs, + setNodePoolsState, + setPlacementBindingsState, + setPlacementDecisionsState, + setPlacementsState, + setPoliciesState, + setPolicyAutomationState, + setPolicyReports, + setPolicySetsState, + setSearchOperator, + setSecrets, + setServices, + setStorageClassState, + setSubmarinerConfigs, + setSubscriptionOperatorsState, + setClusterExtensionsState, + setSubscriptionReportsState, + setSubscriptionsState, + setUsers, + ]) + + const stopStream = useCallback(() => { + streamStoppedRef.current = true + eventSourceRef.current?.close() + eventSourceRef.current = undefined + if (processIntervalRef.current) { + clearInterval(processIntervalRef.current) + processIntervalRef.current = undefined + } + }, []) + + useEffect(() => { + if (!isActive && wasActiveRef.current) { + wasActiveRef.current = false + setIsStreamIdle(true) + if (gracePeriodMs <= 0) { + // No grace period: stop stream immediately + stopStream() + } else { + // Start grace timer (stream keeps running during grace period) + graceTimerRef.current = setTimeout(stopStream, gracePeriodMs) + } + } else if (isActive && !wasActiveRef.current) { + wasActiveRef.current = true + if (graceTimerRef.current) { + clearTimeout(graceTimerRef.current) + graceTimerRef.current = undefined + } + if (streamStoppedRef.current) { + // stopped → reconnecting: stream was killed, need full reload + streamStoppedRef.current = false + setIsStreamIdle(false) + isReconnectingRef.current = true + setIsReconnecting(true) + resetCaches(caches) + resetMapperCaches(mappers) + setEventsLoaded(false) + setRestartKey((k) => k + 1) + } else { + // idle → active: returned during grace period, just hide overlay + setIsStreamIdle(false) + } + } + }, [isActive, gracePeriodMs, caches, mappers, stopStream, setIsStreamIdle, setIsReconnecting]) + + useEffect(() => { + const eventQueue: WatchEvent[] = [] + + function processEventQueue() { + if (eventQueue.length === 0) return + + const resourceTypeMap = eventQueue?.reduce( + (resourceTypeMap, eventData) => { + const apiVersion = eventData.object.apiVersion + const groupVersion = apiVersion.split('/')[0] + const kind = eventData.object.kind + if (!resourceTypeMap[groupVersion]) resourceTypeMap[groupVersion] = {} + if (!resourceTypeMap[groupVersion][kind]) resourceTypeMap[groupVersion][kind] = [] + resourceTypeMap[groupVersion][kind].push(eventData) + return resourceTypeMap + }, + {} as Record> + ) + eventQueue.length = 0 + + for (const groupVersion in resourceTypeMap) { + for (const kind in resourceTypeMap[groupVersion]) { + const watchEvents = resourceTypeMap[groupVersion]?.[kind] + if (watchEvents) { + const setter = setters[groupVersion]?.[kind] + if (setter) { + updateSetterCache(caches, groupVersion, kind, watchEvents) + if (!isReconnectingRef.current) { + setter(Object.values(caches[groupVersion]?.[kind])) + } + } else { + const mapper = mappers[groupVersion]?.[kind] + if (mapper) { + updateMapperCache(mapper, groupVersion, kind, watchEvents) + if (!isReconnectingRef.current) { + mapper.setter({ ...mapper.mcaches[groupVersion]?.[kind] }) + } + } + } + } + } + } + } + + function flushCachesToRecoil() { + for (const groupVersion in setters) { + for (const kind in setters[groupVersion]) { + setters[groupVersion][kind](Object.values(caches[groupVersion]?.[kind])) + } + } + for (const groupVersion in mappers) { + for (const kind in mappers[groupVersion]) { + const { setter, mcaches } = mappers[groupVersion][kind] + setter({ ...mcaches[groupVersion]?.[kind] }) + } + } + } + + function processMessage(event: MessageEvent) { + if (event.data) { + try { + const data = JSON.parse(event.data) as ServerSideEventData + switch (data.type) { + case 'ADDED': + case 'MODIFIED': + case 'DELETED': + eventQueue.push(data) + break + case 'START': + eventQueue.length = 0 + break + // instead of waiting for entire backend data to load + // data is broken up into packets with list resources first + // tables show skeleton until firs packet is received + // then list grows as subsequent packets packets are received + case 'EOP': // END OF A PACKET + processEventQueue() + if (!isReconnectingRef.current) { + setLoadStarted(true) + } + break + case 'LOADED': + processEventQueue() + if (isReconnectingRef.current) { + flushCachesToRecoil() + isReconnectingRef.current = false + setIsReconnecting(false) + } + setEventsLoaded(true) + break + case 'SETTINGS': + setSettings(data.settings) + break + } + } catch (err) { + console.error(err) + } + } + } + + let evtSource: EventSource | undefined + function startWatch() { + evtSource = new EventSource(`${getBackendUrl()}/events`, { withCredentials: true }) + eventSourceRef.current = evtSource + evtSource.onmessage = processMessage + evtSource.onerror = function () { + console.log('EventSource', 'error', 'readyState', evtSource?.readyState) + if (streamStoppedRef.current) return + switch (evtSource?.readyState) { + case EventSource.CLOSED: + setTimeout(() => { + startWatch() + }, 1000) + break + } + } + } + startWatch() + + const timeout = setInterval(processEventQueue, 500) + processIntervalRef.current = timeout + return () => { + clearInterval(timeout) + if (evtSource) evtSource.close() + eventSourceRef.current = undefined + processIntervalRef.current = undefined + } + }, [caches, mappers, restartKey, setIsReconnecting, setLoadStarted, setSettings, setters]) + + const { + data: globalHubRes, + loading: globalHubLoading, + startPolling: globalHubStartPoll, + stopPolling: globalHubStopPoll, + } = useQuery( + globalHubQueryFn, + [ + { + isGlobalHub: false, + localHubName: 'local-cluster', + isHubSelfManaged: undefined, + authentication: { isDirectAuthenticationEnabled: false }, + }, + ], + { + pollInterval: 30, + } ) + + // Start all Polls for Global values here + useEffect(() => { + globalHubStartPoll() + return () => { + // Stop polls on dismount + globalHubStopPoll() + } + }, [globalHubStartPoll, globalHubStopPoll]) + + // Update global value setters when data has finished + const isGlobalHub = useRecoilValue(isGlobalHubState) + if (globalHubRes && !globalHubLoading && !isGlobalHub) { + setIsGlobalHub(globalHubRes[0]?.isGlobalHub) + setlocalHubName(globalHubRes[0]?.localHubName) + setIsHubSelfManaged(globalHubRes[0]?.isHubSelfManaged) + setIsDirectAuthenticationEnabled(globalHubRes[0]?.authentication?.isDirectAuthenticationEnabled ?? false) + setClaimMappings(globalHubRes[0]?.authentication?.claimMappings) + } + + const { + data: mchResponse, + loading: mchLoading, + startPolling: startMCHPoll, + stopPolling: stopMCHPoll, + } = useQuery(mchQueryFn, [], { + pollInterval: 30, + }) + + // Start all Polls for MCH resource + useEffect(() => { + startMCHPoll() + return () => { + // Stop polls on dismount + stopMCHPoll() + } + }, [startMCHPoll, stopMCHPoll]) + + // Update fine-grained RBAC state from mch response + const isFineGrainedRbacEnabled = useRecoilValue(isFineGrainedRbacEnabledState) + if (mchResponse && !mchLoading && !isFineGrainedRbacEnabled) { + setIsFineGrainedRbacEnabled(mchResponse?.find((e) => e?.name === 'fine-grained-rbac')?.enabled ?? false) + } + + // If all data not loaded (!loaded) & events data is loaded (eventsLoaded) && global hub value is loaded (!globalHubLoading) -> set loaded to true + if (!loadCompleted && eventsLoaded && !globalHubLoading) { + setLoadCompleted(true) + } + + useEffect(() => { + function checkLoggedIn() { + fetch(`${getBackendUrl()}/authenticated`, { + credentials: 'include', + headers: { accept: 'application/json' }, + }) + .then((res) => { + switch (res.status) { + case 200: + break + default: + /* istanbul ignore if */ + if (process.env.NODE_ENV === 'development' && res.status === 504) { + window.location.reload() + } else { + tokenExpired() + } + break + } + }) + .catch(() => { + tokenExpired() + }) + .finally(() => { + setTimeout(checkLoggedIn, 30 * 1000) + }) + } + + if (process.env.MODE !== 'plugin') { + checkLoggedIn() + } + }, []) + + const children = useMemo(() => {props.children}, [props.children]) + + return children +} + +function resetCaches(caches: Record>>) { + for (const groupVersion in caches) { + for (const kind in caches[groupVersion]) { + caches[groupVersion][kind] = {} + } + } +} + +function resetMapperCaches( + mappers: Record< + string, + Record< + string, + { + setter: SetterOrUpdater> + mcaches: Record>> + keyBy: string[] + } + > + > +) { + for (const groupVersion in mappers) { + for (const kind in mappers[groupVersion]) { + const { mcaches } = mappers[groupVersion][kind] + for (const gv in mcaches) { + for (const k in mcaches[gv]) { + mcaches[gv][k] = {} + } + } + } + } +} + +function updateSetterCache( + caches: Record>>, + groupVersion: string, + kind: string, + watchEvents: WatchEvent[] +) { + const cache = caches[groupVersion]?.[kind] + for (const watchEvent of watchEvents) { + const key = `${watchEvent.object.metadata.namespace}/${watchEvent.object.metadata.name}` + switch (watchEvent.type) { + case 'ADDED': + case 'MODIFIED': + cache[key] = watchEvent.object + break + case 'DELETED': + delete cache[key] + break + } + } +} + +function updateMapperCache( + mapper: { + setter: SetterOrUpdater> + mcaches: Record>> + keyBy: string[] + }, + groupVersion: string, + kind: string, + watchEvents: WatchEvent[] +) { + const { mcaches, keyBy } = mapper + const map = mcaches[groupVersion]?.[kind] + for (const watchEvent of watchEvents) { + const key = keyBy + .reduce((keys, partKey) => { + keys.push(get(watchEvent.object, partKey)) + return keys + }, [] as string[]) + .join('/') + map[key] = [...(map[key] || [])] + const arr = map[key] + const index = arr.findIndex( + (resource) => + resource.metadata?.name === watchEvent.object.metadata.name && + resource.metadata?.namespace === watchEvent.object.metadata.namespace + ) + switch (watchEvent.type) { + case 'ADDED': + case 'MODIFIED': + if (index !== -1) arr[index] = watchEvent.object + else arr.push(watchEvent.object) + break + case 'DELETED': + if (index !== -1) arr.splice(index, 1) + break + } + } +} + +// Query for GlobalHub check and name +const globalHubQueryFn = () => { + return getRequest<{ + isGlobalHub: boolean + localHubName: string + isHubSelfManaged: boolean | undefined + authentication: { + isDirectAuthenticationEnabled: boolean + claimMappings?: ClaimMappings + } + }>(getBackendUrl() + '/hub') +} + +// Query for GlobalHub check and name +const mchQueryFn = () => { + return getRequest(getBackendUrl() + '/multiclusterhub/components') } diff --git a/frontend/src/components/LoadDataAbstract.test.tsx b/frontend/src/components/LoadDataAbstract.test.tsx deleted file mode 100644 index 801c5390ae..0000000000 --- a/frontend/src/components/LoadDataAbstract.test.tsx +++ /dev/null @@ -1,127 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ - -import { act, render, waitFor } from '@testing-library/react' -import { createElement, ReactElement, type ComponentProps } from 'react' -import { MutableSnapshot, RecoilRoot } from 'recoil' -import { settingsState } from '../atoms' -import { PluginDataContext, defaultContext, PluginData } from '../lib/PluginDataContext' -import { installFakeEventSource } from '../lib/test-event-source' -import { LoadDataAbstract } from './LoadDataAbstract' - -let mockIsActive = true -jest.mock('../lib/usePageActivity', () => ({ - usePageActivity: () => ({ isActive: mockIsActive, deadline: null, pageInUse: true }), -})) - -jest.mock('../resources/utils', () => ({ - getBackendUrl: () => '', -})) - -function createTestContext(overrides: Partial = {}): PluginData { - return { - ...defaultContext, - loadStarted: true, - loadCompleted: true, - startLoading: true, - mounted: true, - ...overrides, - } -} - -function Wrapper({ ctx, children }: { ctx: PluginData; children: ReactElement }) { - return createElement( - PluginDataContext.Provider, - { value: ctx }, - createElement( - RecoilRoot, - { - initializeState: (snapshot: MutableSnapshot) => { - snapshot.set(settingsState, { EVENT_STREAM_IDLE_TIMEOUT: '1', EVENT_STREAM_IDLE_GRACE_PERIOD: '0' }) - }, - } as ComponentProps, - children - ) - ) -} - -describe('LoadDataAbstract', () => { - let fake: ReturnType - - beforeEach(() => { - mockIsActive = true - fake = installFakeEventSource() - }) - - afterEach(() => { - fake.restore() - }) - - it('does not drive overlay flags by default', async () => { - const setIsStreamIdle = jest.fn() - const setIsReconnecting = jest.fn() - const ctx = createTestContext({ setIsStreamIdle, setIsReconnecting }) - const { rerender } = render( - - - - ) - await waitFor(() => expect(fake.sources).toHaveLength(1)) - - mockIsActive = false - rerender( - - - - ) - - expect(setIsStreamIdle).not.toHaveBeenCalled() - expect(setIsReconnecting).not.toHaveBeenCalled() - expect(fake.sources[0].close).toHaveBeenCalled() - }) - - it('drives overlay flags when driveAppLifecycle is set', async () => { - const setIsStreamIdle = jest.fn() - const ctx = createTestContext({ setIsStreamIdle }) - const { rerender } = render( - - - - ) - await waitFor(() => expect(fake.sources).toHaveLength(1)) - - mockIsActive = false - rerender( - - - - ) - - expect(setIsStreamIdle).toHaveBeenCalledWith(true) - }) - - it('applies resources[] watch events into the Recoil setter', async () => { - const setState = jest.fn() - const ctx = createTestContext() - render( - - - - ) - await waitFor(() => expect(fake.sources).toHaveLength(1)) - - const object = { - kind: 'ClusterRole', - apiVersion: 'rbac.authorization.k8s.io/v1', - metadata: { name: 'kubevirt.io:admin', uid: 'uid-1' }, - } - act(() => { - fake.sources[0].emit({ type: 'ADDED', object }) - fake.sources[0].emit({ type: 'EOP' }) - }) - - expect(setState).toHaveBeenCalledWith([object]) - }) -}) diff --git a/frontend/src/components/LoadDataAbstract.tsx b/frontend/src/components/LoadDataAbstract.tsx deleted file mode 100644 index 543f467e68..0000000000 --- a/frontend/src/components/LoadDataAbstract.tsx +++ /dev/null @@ -1,186 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ -import { ReactNode, useCallback, useContext, useEffect, useMemo, useRef, useState } from 'react' -// eslint-disable-next-line @typescript-eslint/no-restricted-imports -import { SetterOrUpdater } from 'recoil' -// eslint-disable-next-line @typescript-eslint/no-restricted-imports -import { useEventStreamIdleGracePeriod, useEventStreamIdleTimeout, WatchEvent } from '../atoms' -import { applyWatchEventsToCache, groupWatchEventsByKind } from '../hooks/applyWatchEventsToCache' -import { useWatchEventStream } from '../hooks/useWatchEventStream' -import { PluginDataContext } from '../lib/PluginDataContext' -import { usePageActivity } from '../lib/usePageActivity' -import type { IResource } from '../resources' - -export interface StreamResource { - apiVersion: string - kind: string - setState: SetterOrUpdater -} - -export interface LoadedContext { - isReconnecting: boolean -} - -export interface LoadDataAbstractProps { - path: string - /** Recoil atom contract for simple single-kind (or few-kind) streams. */ - resources?: StreamResource[] - /** Escape hatch for streams that need custom caches (mappers, reconnect flush). */ - applyWatchEvents?: (events: WatchEvent[]) => void - reset?: () => void - onSettings?: (settings: Record) => void - onEndOfPacket?: () => void - onLoaded?: (ctx: LoadedContext) => void - /** When true, drive PluginDataContext idle/reconnect overlay. Default false. */ - driveAppLifecycle?: boolean - children?: ReactNode -} - -function resourceCacheKey(apiVersion: string, kind: string): string { - return `${apiVersion.split('/')[0]}/${kind}` -} - -export function LoadDataAbstract(props: LoadDataAbstractProps) { - const { mounted, setIsStreamIdle, setIsReconnecting } = useContext(PluginDataContext) - const idleTimeoutMs = useEventStreamIdleTimeout() - const gracePeriodMs = useEventStreamIdleGracePeriod() - const { isActive } = usePageActivity(idleTimeoutMs, mounted) - const wasActiveRef = useRef(true) - const isReconnectingRef = useRef(false) - const streamStoppedRef = useRef(false) - const graceTimerRef = useRef>() - const eventSourceRef = useRef() - const processIntervalRef = useRef>() - const [restartKey, setRestartKey] = useState(0) - const cachesRef = useRef>>({}) - - const resourcesRef = useRef(props.resources) - resourcesRef.current = props.resources - const applyWatchEventsRef = useRef(props.applyWatchEvents) - applyWatchEventsRef.current = props.applyWatchEvents - const resetRef = useRef(props.reset) - resetRef.current = props.reset - const onLoadedRef = useRef(props.onLoaded) - onLoadedRef.current = props.onLoaded - const driveAppLifecycle = props.driveAppLifecycle ?? false - - const applyFromResources = useCallback((events: WatchEvent[]) => { - const resources = resourcesRef.current - if (!resources?.length) return - const grouped = groupWatchEventsByKind(events) - for (const resource of resources) { - const groupVersion = resource.apiVersion.split('/')[0] - const watchEvents = grouped[groupVersion]?.[resource.kind] - if (!watchEvents) continue - const cacheKey = resourceCacheKey(resource.apiVersion, resource.kind) - if (!cachesRef.current[cacheKey]) cachesRef.current[cacheKey] = {} - applyWatchEventsToCache(cachesRef.current[cacheKey], watchEvents) - resource.setState(Object.values(cachesRef.current[cacheKey])) - } - }, []) - - const resetFromResources = useCallback(() => { - const resources = resourcesRef.current - if (!resources?.length) return - for (const resource of resources) { - const cacheKey = resourceCacheKey(resource.apiVersion, resource.kind) - cachesRef.current[cacheKey] = {} - resource.setState([]) - } - }, []) - - const applyWatchEvents = useCallback( - (events: WatchEvent[]) => { - if (applyWatchEventsRef.current) { - applyWatchEventsRef.current(events) - return - } - applyFromResources(events) - }, - [applyFromResources] - ) - - const handleReset = useCallback(() => { - if (resetRef.current) { - resetRef.current() - return - } - resetFromResources() - }, [resetFromResources]) - - const stopStream = useCallback(() => { - streamStoppedRef.current = true - eventSourceRef.current?.close() - eventSourceRef.current = undefined - if (processIntervalRef.current) { - clearInterval(processIntervalRef.current) - processIntervalRef.current = undefined - } - }, []) - - const onBecameInactive = useCallback(() => { - wasActiveRef.current = false - if (driveAppLifecycle) { - setIsStreamIdle(true) - } - if (gracePeriodMs <= 0) { - stopStream() - return - } - graceTimerRef.current = setTimeout(stopStream, gracePeriodMs) - }, [driveAppLifecycle, gracePeriodMs, setIsStreamIdle, stopStream]) - - const onBecameActive = useCallback(() => { - wasActiveRef.current = true - if (graceTimerRef.current) { - clearTimeout(graceTimerRef.current) - graceTimerRef.current = undefined - } - if (!streamStoppedRef.current) { - if (driveAppLifecycle) { - setIsStreamIdle(false) - } - return - } - streamStoppedRef.current = false - isReconnectingRef.current = true - if (driveAppLifecycle) { - setIsStreamIdle(false) - setIsReconnecting(true) - } - handleReset() - setRestartKey((k) => k + 1) - }, [driveAppLifecycle, handleReset, setIsReconnecting, setIsStreamIdle]) - - useEffect(() => { - if (!isActive && wasActiveRef.current) { - onBecameInactive() - } else if (isActive && !wasActiveRef.current) { - onBecameActive() - } - }, [isActive, onBecameActive, onBecameInactive]) - - const onLoaded = useCallback(() => { - const isReconnecting = isReconnectingRef.current - onLoadedRef.current?.({ isReconnecting }) - if (isReconnecting) { - isReconnectingRef.current = false - if (driveAppLifecycle) { - setIsReconnecting(false) - } - } - }, [driveAppLifecycle, setIsReconnecting]) - - useWatchEventStream({ - path: props.path, - restartKey, - streamStoppedRef, - eventSourceRef, - processIntervalRef, - applyWatchEvents, - onSettings: props.onSettings, - onEndOfPacket: props.onEndOfPacket, - onLoaded, - }) - - return useMemo(() => (props.children ? <>{props.children} : null), [props.children]) -} diff --git a/frontend/src/components/LoadEventsData.tsx b/frontend/src/components/LoadEventsData.tsx deleted file mode 100644 index 2310f67736..0000000000 --- a/frontend/src/components/LoadEventsData.tsx +++ /dev/null @@ -1,720 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ -import get from 'lodash/get' -import { useCallback, useContext, useEffect, useMemo, useRef, useState } from 'react' -// eslint-disable-next-line @typescript-eslint/no-restricted-imports -import { SetterOrUpdater, useRecoilValue, useSetRecoilState } from 'recoil' -import { tokenExpired } from '../logout' -import { - AgentClusterInstallApiVersion, - AgentClusterInstallKind, - AgentKind, - AgentKindVersion, - AgentMachineApiVersion, - AgentMachineKind, - AgentServiceConfigKind, - AgentServiceConfigKindVersion, - AnsibleJobApiVersion, - AnsibleJobKind, - AnsibleWorkflowKind, - ApplicationApiVersion, - ApplicationKind, - BareMetalHostApiVersion, - BareMetalHostKind, - CertificateSigningRequestApiVersion, - CertificateSigningRequestKind, - ChannelApiVersion, - ChannelKind, - ClusterClaimApiVersion, - ClusterClaimKind, - ClusterCuratorApiVersion, - ClusterCuratorKind, - ClusterDeploymentApiVersion, - ClusterDeploymentKind, - ClusterImageSetApiVersion, - ClusterImageSetKind, - ClusterManagementAddOnApiVersion, - ClusterManagementAddOnKind, - ClusterPoolApiVersion, - ClusterPoolKind, - ClusterProvisionApiVersion, - ClusterProvisionKind, - ClusterVersionApiVersion, - ClusterVersionKind, - ConfigMapApiVersion, - ConfigMapKind, - DiscoveredClusterApiVersion, - DiscoveredClusterKind, - DiscoveryConfigApiVersion, - DiscoveryConfigKind, - GitOpsClusterApiVersion, - GitOpsClusterKind, - GroupKind, - HelmReleaseApiVersion, - HelmReleaseKind, - HostedClusterApiVersion, - HostedClusterKind, - InfraEnvApiVersion, - InfraEnvKind, - InfrastructureApiVersion, - InfrastructureKind, - IResource, - MachinePoolApiVersion, - MachinePoolKind, - ManagedClusterAddOnApiVersion, - ManagedClusterAddOnKind, - ManagedClusterApiVersion, - ManagedClusterInfoApiVersion, - ManagedClusterInfoKind, - ManagedClusterKind, - ManagedClusterSetApiVersion, - ManagedClusterSetBindingApiVersion, - ManagedClusterSetBindingKind, - ManagedClusterSetKind, - MulticlusterRoleAssignmentApiVersion, - MulticlusterRoleAssignmentKind, - MultiClusterEngineApiVersion, - MultiClusterEngineKind, - NamespaceApiVersion, - NamespaceKind, - NMStateConfigApiVersion, - NMStateConfigKind, - NodePoolApiVersion, - NodePoolKind, - PlacementApiVersionAlpha, - PlacementBindingApiVersion, - PlacementBindingKind, - PlacementDecisionApiVersion, - PlacementDecisionKind, - PlacementKind, - PolicyApiVersion, - PolicyAutomationApiVersion, - PolicyAutomationKind, - PolicyKind, - PolicyReportApiVersion, - PolicyReportKind, - PolicySetApiVersion, - PolicySetKind, - SearchOperatorApiVersion, - SearchOperatorKind, - SecretApiVersion, - SecretKind, - StorageClassApiVersion, - StorageClassKind, - SubmarinerConfigApiVersion, - SubmarinerConfigKind, - SubscriptionApiVersion, - SubscriptionKind, - SubscriptionOperatorApiVersion, - SubscriptionOperatorKind, - ClusterExtensionApiVersion, - ClusterExtensionKind, - ClusterRoleKind, - RbacApiVersion, - SubscriptionReportApiVersion, - SubscriptionReportKind, - UserApiVersion, - UserKind, - ServiceApiVersion, - ServiceKind, -} from '../resources' -import { getBackendUrl, getRequest } from '../resources/utils' -// eslint-disable-next-line @typescript-eslint/no-restricted-imports -import { - agentClusterInstallsState, - agentMachinesState, - agentServiceConfigsState, - agentsState, - ansibleJobState, - ansibleWorkflowState, - applicationsState, - argoCDsState, - bareMetalHostsState, - certificateSigningRequestsState, - channelsState, - claimMappingsState, - clusterClaimsState, - clusterCuratorsState, - clusterDeploymentsState, - clusterImageSetsState, - clusterManagementAddonsState, - clusterPoolsState, - clusterProvisionsState, - clusterVersionState, - configMapsState, - discoveredClusterState, - discoveryConfigState, - gitOpsClustersState, - groupsState, - helmReleaseState, - hostedClustersState, - infraEnvironmentsState, - infrastructuresState, - isDirectAuthenticationEnabledState, - isFineGrainedRbacEnabledState, - isGlobalHubState, - isHubSelfManagedState, - localHubNameState, - machinePoolsState, - managedClusterAddonsState, - managedClusterInfosState, - managedClusterSetBindingsState, - managedClusterSetsState, - managedClustersState, - multiClusterEnginesState, - multiclusterRoleAssignmentState, - namespacesState, - nmStateConfigsState, - nodePoolsState, - placementBindingsState, - placementDecisionsState, - placementsState, - policiesState, - policyAutomationState, - policyreportState, - policySetsState, - searchOperatorState, - secretsState, - settingsState, - servicesState, - storageClassState, - submarinerConfigsState, - subscriptionOperatorsState, - clusterExtensionsState, - subscriptionReportsState, - subscriptionsState, - usersState, - vmClusterRolesState, - WatchEvent, -} from '../atoms' -import { applyWatchEventsToCache, groupWatchEventsByKind } from '../hooks/applyWatchEventsToCache' -import { PluginDataContext } from '../lib/PluginDataContext' -import { useQuery } from '../lib/useQuery' -import { MultiClusterHubComponent } from '../resources/multi-cluster-hub-component' -import { ClaimMappings } from '~/resources/authentication' -import { LoadDataAbstract } from './LoadDataAbstract' - -export function LoadEventsData() { - const { loadCompleted, setLoadStarted, setLoadCompleted } = useContext(PluginDataContext) - const [eventsLoaded, setEventsLoaded] = useState(false) - const isReconnectingRef = useRef(false) - - const setAgentClusterInstalls = useSetRecoilState(agentClusterInstallsState) - const setAgentMachinesState = useSetRecoilState(agentMachinesState) - const setAgents = useSetRecoilState(agentsState) - const setAgentServiceConfigs = useSetRecoilState(agentServiceConfigsState) - const setAnsibleJobs = useSetRecoilState(ansibleJobState) - const setAnsibleWorkflows = useSetRecoilState(ansibleWorkflowState) - const setApplicationsState = useSetRecoilState(applicationsState) - const setArgoCDsState = useSetRecoilState(argoCDsState) - const setBareMetalHosts = useSetRecoilState(bareMetalHostsState) - const setCertificateSigningRequests = useSetRecoilState(certificateSigningRequestsState) - const setChannelsState = useSetRecoilState(channelsState) - const setClusterClaims = useSetRecoilState(clusterClaimsState) - const setClusterCurators = useSetRecoilState(clusterCuratorsState) - const setClusterDeployments = useSetRecoilState(clusterDeploymentsState) - const setClusterImageSets = useSetRecoilState(clusterImageSetsState) - const setClusterManagementAddons = useSetRecoilState(clusterManagementAddonsState) - const setClusterPools = useSetRecoilState(clusterPoolsState) - const setClusterProvisions = useSetRecoilState(clusterProvisionsState) - const setClusterVerions = useSetRecoilState(clusterVersionState) - const setConfigMaps = useSetRecoilState(configMapsState) - const setDiscoveredClusters = useSetRecoilState(discoveredClusterState) - const setDiscoveryConfigs = useSetRecoilState(discoveryConfigState) - const setGitOpsClustersState = useSetRecoilState(gitOpsClustersState) - const setGroups = useSetRecoilState(groupsState) - const setHelmReleases = useSetRecoilState(helmReleaseState) - const setHostedClustersState = useSetRecoilState(hostedClustersState) - const setInfraEnvironments = useSetRecoilState(infraEnvironmentsState) - const setInfrastructure = useSetRecoilState(infrastructuresState) - const setClaimMappings = useSetRecoilState(claimMappingsState) - const setIsDirectAuthenticationEnabled = useSetRecoilState(isDirectAuthenticationEnabledState) - const setIsFineGrainedRbacEnabled = useSetRecoilState(isFineGrainedRbacEnabledState) - const setIsGlobalHub = useSetRecoilState(isGlobalHubState) - const setIsHubSelfManaged = useSetRecoilState(isHubSelfManagedState) - const setlocalHubName = useSetRecoilState(localHubNameState) - const setMachinePools = useSetRecoilState(machinePoolsState) - const setManagedClusterAddons = useSetRecoilState(managedClusterAddonsState) - const setManagedClusterInfos = useSetRecoilState(managedClusterInfosState) - const setManagedClusterSetBindings = useSetRecoilState(managedClusterSetBindingsState) - const setManagedClusterSets = useSetRecoilState(managedClusterSetsState) - const setManagedClusters = useSetRecoilState(managedClustersState) - const setMultiClusterEngines = useSetRecoilState(multiClusterEnginesState) - const setMulticlusterRoleAssignments = useSetRecoilState(multiclusterRoleAssignmentState) - const setNamespaces = useSetRecoilState(namespacesState) - const setNMStateConfigs = useSetRecoilState(nmStateConfigsState) - const setNodePoolsState = useSetRecoilState(nodePoolsState) - const setPlacementBindingsState = useSetRecoilState(placementBindingsState) - const setPlacementDecisionsState = useSetRecoilState(placementDecisionsState) - const setPlacementsState = useSetRecoilState(placementsState) - const setPoliciesState = useSetRecoilState(policiesState) - const setPolicyAutomationState = useSetRecoilState(policyAutomationState) - const setPolicyReports = useSetRecoilState(policyreportState) - const setPolicySetsState = useSetRecoilState(policySetsState) - const setSearchOperator = useSetRecoilState(searchOperatorState) - const setSecrets = useSetRecoilState(secretsState) - const setSettings = useSetRecoilState(settingsState) - const setServices = useSetRecoilState(servicesState) - const setStorageClassState = useSetRecoilState(storageClassState) - const setSubmarinerConfigs = useSetRecoilState(submarinerConfigsState) - const setSubscriptionOperatorsState = useSetRecoilState(subscriptionOperatorsState) - const setClusterExtensionsState = useSetRecoilState(clusterExtensionsState) - const setSubscriptionReportsState = useSetRecoilState(subscriptionReportsState) - const setSubscriptionsState = useSetRecoilState(subscriptionsState) - const setUsers = useSetRecoilState(usersState) - const setVMClusterRoles = useSetRecoilState(vmClusterRolesState) - - const { setters, mappers, caches } = useMemo(() => { - const setters: Record>> = {} - - const mappers: Record< - string, - Record< - string, - { - setter: SetterOrUpdater> - mcaches: Record>> - keyBy: string[] - } - > - > = {} - const caches: Record>> = {} - const mcaches: Record>> = {} - function addSetter(apiVersion: string, kind: string, setter: SetterOrUpdater) { - const groupVersion = apiVersion.split('/')[0] - if (!setters[groupVersion]) setters[groupVersion] = {} - setters[groupVersion][kind] = setter - if (!caches[groupVersion]) caches[groupVersion] = {} - caches[groupVersion][kind] = {} - } - function addMapper( - apiVersion: string, - kind: string, - setter: SetterOrUpdater>, - keyBy: string[] - ) { - const groupVersion = apiVersion.split('/')[0] - if (!mappers[groupVersion]) mappers[groupVersion] = {} - if (!mcaches[groupVersion]) mcaches[groupVersion] = {} - mcaches[groupVersion][kind] = {} - mappers[groupVersion][kind] = { setter, mcaches, keyBy } - } - - // mappers (key=>[values]) - addMapper(ManagedClusterAddOnApiVersion, ManagedClusterAddOnKind, setManagedClusterAddons, ['metadata.namespace']) - - // setters - addSetter('argoproj.io/v1alpha1', 'ArgoCD', setArgoCDsState) - addSetter(AgentClusterInstallApiVersion, AgentClusterInstallKind, setAgentClusterInstalls) - addSetter(AgentKindVersion, AgentKind, setAgents) - addSetter(AgentMachineApiVersion, AgentMachineKind, setAgentMachinesState) - addSetter(AgentServiceConfigKindVersion, AgentServiceConfigKind, setAgentServiceConfigs) - addSetter(AnsibleJobApiVersion, AnsibleJobKind, setAnsibleJobs) - addSetter(AnsibleJobApiVersion, AnsibleWorkflowKind, setAnsibleWorkflows) - addSetter(ApplicationApiVersion, ApplicationKind, setApplicationsState) - addSetter(BareMetalHostApiVersion, BareMetalHostKind, setBareMetalHosts) - addSetter(CertificateSigningRequestApiVersion, CertificateSigningRequestKind, setCertificateSigningRequests) - addSetter(ChannelApiVersion, ChannelKind, setChannelsState) - addSetter(ClusterClaimApiVersion, ClusterClaimKind, setClusterClaims) - addSetter(ClusterCuratorApiVersion, ClusterCuratorKind, setClusterCurators) - addSetter(ClusterDeploymentApiVersion, ClusterDeploymentKind, setClusterDeployments) - addSetter(ClusterImageSetApiVersion, ClusterImageSetKind, setClusterImageSets) - addSetter(ClusterManagementAddOnApiVersion, ClusterManagementAddOnKind, setClusterManagementAddons) - addSetter(ClusterPoolApiVersion, ClusterPoolKind, setClusterPools) - addSetter(ClusterProvisionApiVersion, ClusterProvisionKind, setClusterProvisions) - addSetter(ClusterVersionApiVersion, ClusterVersionKind, setClusterVerions) - addSetter(ConfigMapApiVersion, ConfigMapKind, setConfigMaps) - addSetter(DiscoveredClusterApiVersion, DiscoveredClusterKind, setDiscoveredClusters) - addSetter(DiscoveryConfigApiVersion, DiscoveryConfigKind, setDiscoveryConfigs) - addSetter(GitOpsClusterApiVersion, GitOpsClusterKind, setGitOpsClustersState) - addSetter(HelmReleaseApiVersion, HelmReleaseKind, setHelmReleases) - addSetter(HostedClusterApiVersion, HostedClusterKind, setHostedClustersState) - addSetter(InfraEnvApiVersion, InfraEnvKind, setInfraEnvironments) - addSetter(InfrastructureApiVersion, InfrastructureKind, setInfrastructure) - addSetter(MachinePoolApiVersion, MachinePoolKind, setMachinePools) - addSetter(ManagedClusterApiVersion, ManagedClusterKind, setManagedClusters) - addSetter(ManagedClusterInfoApiVersion, ManagedClusterInfoKind, setManagedClusterInfos) - addSetter(ManagedClusterSetApiVersion, ManagedClusterSetKind, setManagedClusterSets) - addSetter(ManagedClusterSetBindingApiVersion, ManagedClusterSetBindingKind, setManagedClusterSetBindings) - addSetter(MulticlusterRoleAssignmentApiVersion, MulticlusterRoleAssignmentKind, setMulticlusterRoleAssignments) - addSetter(MultiClusterEngineApiVersion, MultiClusterEngineKind, setMultiClusterEngines) - addSetter(NamespaceApiVersion, NamespaceKind, setNamespaces) - addSetter(NMStateConfigApiVersion, NMStateConfigKind, setNMStateConfigs) - addSetter(NodePoolApiVersion, NodePoolKind, setNodePoolsState) - addSetter(PlacementApiVersionAlpha, PlacementKind, setPlacementsState) - addSetter(PlacementBindingApiVersion, PlacementBindingKind, setPlacementBindingsState) - addSetter(PlacementDecisionApiVersion, PlacementDecisionKind, setPlacementDecisionsState) - addSetter(PolicyApiVersion, PolicyKind, setPoliciesState) - addSetter(PolicyAutomationApiVersion, PolicyAutomationKind, setPolicyAutomationState) - addSetter(PolicyReportApiVersion, PolicyReportKind, setPolicyReports) - addSetter(PolicySetApiVersion, PolicySetKind, setPolicySetsState) - addSetter(RbacApiVersion, ClusterRoleKind, setVMClusterRoles) - addSetter(SearchOperatorApiVersion, SearchOperatorKind, setSearchOperator) - addSetter(SecretApiVersion, SecretKind, setSecrets) - addSetter(ServiceApiVersion, ServiceKind, setServices) - addSetter(StorageClassApiVersion, StorageClassKind, setStorageClassState) - addSetter(SubmarinerConfigApiVersion, SubmarinerConfigKind, setSubmarinerConfigs) - addSetter(SubscriptionApiVersion, SubscriptionKind, setSubscriptionsState) - addSetter(SubscriptionOperatorApiVersion, SubscriptionOperatorKind, setSubscriptionOperatorsState) - addSetter(ClusterExtensionApiVersion, ClusterExtensionKind, setClusterExtensionsState) - addSetter(SubscriptionReportApiVersion, SubscriptionReportKind, setSubscriptionReportsState) - addSetter(UserApiVersion, GroupKind, setGroups) - addSetter(UserApiVersion, UserKind, setUsers) - - return { setters, mappers, caches } - }, [ - setAgentClusterInstalls, - setAgentMachinesState, - setAgents, - setAgentServiceConfigs, - setAnsibleJobs, - setAnsibleWorkflows, - setApplicationsState, - setArgoCDsState, - setBareMetalHosts, - setCertificateSigningRequests, - setChannelsState, - setClusterClaims, - setClusterCurators, - setClusterDeployments, - setClusterImageSets, - setClusterManagementAddons, - setClusterPools, - setClusterProvisions, - setClusterVerions, - setConfigMaps, - setDiscoveredClusters, - setDiscoveryConfigs, - setGitOpsClustersState, - setGroups, - setHelmReleases, - setHostedClustersState, - setInfraEnvironments, - setInfrastructure, - setMachinePools, - setManagedClusterAddons, - setManagedClusterInfos, - setManagedClusterSetBindings, - setManagedClusterSets, - setManagedClusters, - setMultiClusterEngines, - setMulticlusterRoleAssignments, - setNamespaces, - setNMStateConfigs, - setNodePoolsState, - setPlacementBindingsState, - setPlacementDecisionsState, - setPlacementsState, - setPoliciesState, - setPolicyAutomationState, - setPolicyReports, - setPolicySetsState, - setSearchOperator, - setSecrets, - setServices, - setStorageClassState, - setSubmarinerConfigs, - setSubscriptionOperatorsState, - setClusterExtensionsState, - setSubscriptionReportsState, - setSubscriptionsState, - setUsers, - setVMClusterRoles, - ]) - - const applyWatchEvents = useCallback( - (watchEvents: WatchEvent[]) => { - const resourceTypeMap = groupWatchEventsByKind(watchEvents) - const skipAtomUpdate = isReconnectingRef.current - for (const groupVersion in resourceTypeMap) { - for (const kind in resourceTypeMap[groupVersion]) { - const kindEvents = resourceTypeMap[groupVersion]?.[kind] - if (!kindEvents) continue - applyWatchEventsForKind(setters, caches, mappers, groupVersion, kind, kindEvents, skipAtomUpdate) - } - } - }, - [caches, mappers, setters] - ) - - const reset = useCallback(() => { - isReconnectingRef.current = true - resetCaches(caches) - resetMapperCaches(mappers) - setEventsLoaded(false) - }, [caches, mappers]) - - const flushCachesToRecoil = useCallback(() => { - for (const groupVersion in setters) { - for (const kind in setters[groupVersion]) { - setters[groupVersion][kind](Object.values(caches[groupVersion]?.[kind])) - } - } - for (const groupVersion in mappers) { - for (const kind in mappers[groupVersion]) { - const { setter, mcaches } = mappers[groupVersion][kind] - setter({ ...mcaches[groupVersion]?.[kind] }) - } - } - }, [caches, mappers, setters]) - - const onEndOfPacket = useCallback(() => { - if (!isReconnectingRef.current) { - setLoadStarted(true) - } - }, [setLoadStarted]) - - const onLoaded = useCallback( - ({ isReconnecting }: { isReconnecting: boolean }) => { - if (isReconnecting) { - flushCachesToRecoil() - isReconnectingRef.current = false - } - setEventsLoaded(true) - }, - [flushCachesToRecoil] - ) - - const onSettings = useCallback( - (settings: Record) => { - setSettings(settings) - }, - [setSettings] - ) - - const { - data: globalHubRes, - loading: globalHubLoading, - startPolling: globalHubStartPoll, - stopPolling: globalHubStopPoll, - } = useQuery( - globalHubQueryFn, - [ - { - isGlobalHub: false, - localHubName: 'local-cluster', - isHubSelfManaged: undefined, - authentication: { isDirectAuthenticationEnabled: false }, - }, - ], - { - pollInterval: 30, - } - ) - - // Start all Polls for Global values here - useEffect(() => { - globalHubStartPoll() - return () => { - // Stop polls on dismount - globalHubStopPoll() - } - }, [globalHubStartPoll, globalHubStopPoll]) - - // Update global value setters when data has finished - const isGlobalHub = useRecoilValue(isGlobalHubState) - if (globalHubRes && !globalHubLoading && !isGlobalHub) { - setIsGlobalHub(globalHubRes[0]?.isGlobalHub) - setlocalHubName(globalHubRes[0]?.localHubName) - setIsHubSelfManaged(globalHubRes[0]?.isHubSelfManaged) - setIsDirectAuthenticationEnabled(globalHubRes[0]?.authentication?.isDirectAuthenticationEnabled ?? false) - setClaimMappings(globalHubRes[0]?.authentication?.claimMappings) - } - - const { - data: mchResponse, - loading: mchLoading, - startPolling: startMCHPoll, - stopPolling: stopMCHPoll, - } = useQuery(mchQueryFn, [], { - pollInterval: 30, - }) - - // Start all Polls for MCH resource - useEffect(() => { - startMCHPoll() - return () => { - // Stop polls on dismount - stopMCHPoll() - } - }, [startMCHPoll, stopMCHPoll]) - - // Update fine-grained RBAC state from mch response - const isFineGrainedRbacEnabled = useRecoilValue(isFineGrainedRbacEnabledState) - if (mchResponse && !mchLoading && !isFineGrainedRbacEnabled) { - setIsFineGrainedRbacEnabled(mchResponse?.find((e) => e?.name === 'fine-grained-rbac')?.enabled ?? false) - } - - // If all data not loaded (!loaded) & events data is loaded (eventsLoaded) && global hub value is loaded (!globalHubLoading) -> set loaded to true - if (!loadCompleted && eventsLoaded && !globalHubLoading) { - setLoadCompleted(true) - } - - useEffect(() => { - function checkLoggedIn() { - fetch(`${getBackendUrl()}/authenticated`, { - credentials: 'include', - headers: { accept: 'application/json' }, - }) - .then((res) => { - if (res.status === 200) { - return - } - /* istanbul ignore if */ - if (process.env.NODE_ENV === 'development' && res.status === 504) { - window.location.reload() - } else { - tokenExpired() - } - }) - .catch(() => { - tokenExpired() - }) - .finally(() => { - setTimeout(checkLoggedIn, 30 * 1000) - }) - } - - if (process.env.MODE !== 'plugin') { - checkLoggedIn() - } - }, []) - - return ( - - ) -} - -function applyWatchEventsForKind( - setters: Record>>, - caches: Record>>, - mappers: Record< - string, - Record< - string, - { - setter: SetterOrUpdater> - mcaches: Record>> - keyBy: string[] - } - > - >, - groupVersion: string, - kind: string, - kindEvents: WatchEvent[], - skipAtomUpdate: boolean -) { - const setter = setters[groupVersion]?.[kind] - if (setter) { - const cache = caches[groupVersion]?.[kind] - if (!cache) return - applyWatchEventsToCache(cache, kindEvents) - if (!skipAtomUpdate) { - setter(Object.values(cache)) - } - return - } - const mapper = mappers[groupVersion]?.[kind] - if (!mapper) return - updateMapperCache(mapper, groupVersion, kind, kindEvents) - if (!skipAtomUpdate) { - mapper.setter({ ...mapper.mcaches[groupVersion]?.[kind] }) - } -} - -function resetCaches(caches: Record>>) { - for (const groupVersion in caches) { - for (const kind in caches[groupVersion]) { - caches[groupVersion][kind] = {} - } - } -} - -function resetMapperCaches( - mappers: Record< - string, - Record< - string, - { - setter: SetterOrUpdater> - mcaches: Record>> - keyBy: string[] - } - > - > -) { - for (const groupVersion in mappers) { - for (const kind in mappers[groupVersion]) { - const { mcaches } = mappers[groupVersion][kind] - for (const gv in mcaches) { - for (const k in mcaches[gv]) { - mcaches[gv][k] = {} - } - } - } - } -} - -function updateMapperCache( - mapper: { - setter: SetterOrUpdater> - mcaches: Record>> - keyBy: string[] - }, - groupVersion: string, - kind: string, - watchEvents: WatchEvent[] -) { - const { mcaches, keyBy } = mapper - const map = mcaches[groupVersion]?.[kind] - for (const watchEvent of watchEvents) { - const key = keyBy - .reduce((keys, partKey) => { - keys.push(get(watchEvent.object, partKey)) - return keys - }, [] as string[]) - .join('/') - map[key] = [...(map[key] || [])] - const arr = map[key] - const index = arr.findIndex( - (resource) => - resource.metadata?.name === watchEvent.object.metadata.name && - resource.metadata?.namespace === watchEvent.object.metadata.namespace - ) - switch (watchEvent.type) { - case 'ADDED': - case 'MODIFIED': - if (index !== -1) arr[index] = watchEvent.object - else arr.push(watchEvent.object) - break - case 'DELETED': - if (index !== -1) arr.splice(index, 1) - break - } - } -} - -// Query for GlobalHub check and name -const globalHubQueryFn = () => { - return getRequest<{ - isGlobalHub: boolean - localHubName: string - isHubSelfManaged: boolean | undefined - authentication: { - isDirectAuthenticationEnabled: boolean - claimMappings?: ClaimMappings - } - }>(getBackendUrl() + '/hub') -} - -// Query for GlobalHub check and name -const mchQueryFn = () => { - return getRequest(getBackendUrl() + '/multiclusterhub/components') -} diff --git a/frontend/src/components/LoadPluginData.test.tsx b/frontend/src/components/LoadPluginData.test.tsx index a79bdc0cbf..dec397513f 100644 --- a/frontend/src/components/LoadPluginData.test.tsx +++ b/frontend/src/components/LoadPluginData.test.tsx @@ -2,7 +2,6 @@ import { render, screen } from '@testing-library/react' import { MemoryRouter } from 'react-router' -import { NavigationPath } from '../NavigationPath' import { defaultContext, PluginData, PluginDataContext } from '../lib/PluginDataContext' import { PluginContext, defaultPlugin } from '../lib/PluginContext' import { LoadPluginData } from './LoadPluginData' @@ -13,10 +12,10 @@ jest.mock('../lib/acm-i18next', () => ({ }), })) -function renderWithContext(contextOverrides: Partial, children = 'Page Content', initialPath = '/') { +function renderWithContext(contextOverrides: Partial, children = 'Page Content') { const ctx: PluginData = { ...defaultContext, ...contextOverrides } return render( - + {children} @@ -33,11 +32,6 @@ describe('LoadPluginData', () => { expect(screen.getByText('Loading')).toBeInTheDocument() }) - it('fast-loads /multicloud/home when loadStarted without waiting for loadCompleted', () => { - renderWithContext({ loadCompleted: false, loadStarted: true }, 'Page Content', NavigationPath.home) - expect(screen.getByText('Page Content')).toBeInTheDocument() - }) - it('shows children when loadCompleted is true', () => { renderWithContext({ loadCompleted: true }) expect(screen.getByText('Page Content')).toBeInTheDocument() diff --git a/frontend/src/components/LoadRbacData.test.tsx b/frontend/src/components/LoadRbacData.test.tsx deleted file mode 100644 index a679bfd00c..0000000000 --- a/frontend/src/components/LoadRbacData.test.tsx +++ /dev/null @@ -1,142 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ - -import { act, render, waitFor } from '@testing-library/react' -import { createElement, ReactElement, type ComponentProps } from 'react' -import { MutableSnapshot, RecoilRoot, useRecoilValue } from 'recoil' -import { settingsState, vmClusterRolesState } from '../atoms' -import { PluginDataContext, defaultContext, PluginData } from '../lib/PluginDataContext' -import { installFakeEventSource } from '../lib/test-event-source' -import { ClusterRole } from '../resources' -import { LoadRbacData } from './LoadRbacData' - -let mockIsActive = true -jest.mock('../lib/usePageActivity', () => ({ - usePageActivity: () => ({ isActive: mockIsActive, deadline: null, pageInUse: true }), -})) - -jest.mock('../resources/utils', () => ({ - getBackendUrl: () => '', -})) - -function createTestContext(overrides: Partial = {}): PluginData { - return { - ...defaultContext, - loadStarted: true, - loadCompleted: true, - startLoading: true, - mounted: true, - ...overrides, - } -} - -function RolesProbe() { - const roles = useRecoilValue(vmClusterRolesState) - return createElement('div', { id: 'roles' }, String(roles.length)) -} - -function Wrapper({ ctx, children }: { ctx: PluginData; children: ReactElement }) { - return createElement( - PluginDataContext.Provider, - { value: ctx }, - createElement( - RecoilRoot, - { - initializeState: (snapshot: MutableSnapshot) => { - snapshot.set(settingsState, { EVENT_STREAM_IDLE_TIMEOUT: '1', EVENT_STREAM_IDLE_GRACE_PERIOD: '0' }) - }, - } as ComponentProps, - children - ) - ) -} - -const sampleRole: ClusterRole = { - apiVersion: 'rbac.authorization.k8s.io/v1', - kind: 'ClusterRole', - metadata: { name: 'kubevirt.io:admin', uid: 'uid-1' }, - rules: [], -} - -describe('LoadRbacData', () => { - let fake: ReturnType - - beforeEach(() => { - mockIsActive = true - fake = installFakeEventSource() - }) - - afterEach(() => { - fake.restore() - }) - - it('opens /events/rbac with credentials and applies ADDED into the atom', async () => { - const ctx = createTestContext() - render( - - <> - - - - - ) - - await waitFor(() => expect(fake.sources).toHaveLength(1)) - expect(fake.sources[0].url).toBe('/events/rbac') - expect(fake.sources[0].withCredentials).toBe(true) - - act(() => { - fake.sources[0].emit({ type: 'START' }) - fake.sources[0].emit({ type: 'ADDED', object: sampleRole }) - fake.sources[0].emit({ type: 'EOP' }) - fake.sources[0].emit({ type: 'LOADED' }) - }) - - await waitFor(() => { - expect(document.getElementById('roles')?.textContent).toBe('1') - }) - }) - - it('removes DELETED roles from the atom', async () => { - const ctx = createTestContext() - render( - - <> - - - - - ) - await waitFor(() => expect(fake.sources).toHaveLength(1)) - act(() => { - fake.sources[0].emit({ type: 'ADDED', object: sampleRole }) - fake.sources[0].emit({ type: 'EOP' }) - }) - await waitFor(() => expect(document.getElementById('roles')?.textContent).toBe('1')) - act(() => { - fake.sources[0].emit({ type: 'DELETED', object: sampleRole }) - fake.sources[0].emit({ type: 'EOP' }) - }) - await waitFor(() => expect(document.getElementById('roles')?.textContent).toBe('0')) - }) - - it('closes the stream when idle with no grace period and does not drive overlay flags', async () => { - const setIsStreamIdle = jest.fn() - const setIsReconnecting = jest.fn() - const ctx = createTestContext({ setIsStreamIdle, setIsReconnecting }) - const { rerender } = render( - - - - ) - await waitFor(() => expect(fake.sources).toHaveLength(1)) - mockIsActive = false - rerender( - - - - ) - expect(fake.sources[0].close).toHaveBeenCalled() - expect(setIsStreamIdle).not.toHaveBeenCalled() - expect(setIsReconnecting).not.toHaveBeenCalled() - }) -}) diff --git a/frontend/src/components/LoadRbacData.tsx b/frontend/src/components/LoadRbacData.tsx deleted file mode 100644 index 7a8ec650af..0000000000 --- a/frontend/src/components/LoadRbacData.tsx +++ /dev/null @@ -1,17 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ -import { useMemo } from 'react' -// eslint-disable-next-line @typescript-eslint/no-restricted-imports -import { useSetRecoilState } from 'recoil' -// eslint-disable-next-line @typescript-eslint/no-restricted-imports -import { vmClusterRolesState } from '../atoms' -import { ClusterRoleKind, RbacApiVersion } from '../resources' -import { LoadDataAbstract } from './LoadDataAbstract' - -export function LoadRbacData() { - const setVMClusterRoles = useSetRecoilState(vmClusterRolesState) - const resources = useMemo( - () => [{ apiVersion: RbacApiVersion, kind: ClusterRoleKind, setState: setVMClusterRoles }], - [setVMClusterRoles] - ) - return -} diff --git a/frontend/src/hooks/applyWatchEventsToCache.test.ts b/frontend/src/hooks/applyWatchEventsToCache.test.ts deleted file mode 100644 index c675bc4f1f..0000000000 --- a/frontend/src/hooks/applyWatchEventsToCache.test.ts +++ /dev/null @@ -1,66 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ -import type { WatchEvent } from '../atoms' -import type { IResource } from '../resources' -import { applyWatchEventsToCache, groupWatchEventsByKind, resourceKey } from './applyWatchEventsToCache' - -const added: WatchEvent = { - type: 'ADDED', - object: { - kind: 'ClusterRole', - apiVersion: 'rbac.authorization.k8s.io/v1', - metadata: { name: 'admin', namespace: undefined as unknown as string, resourceVersion: '1' }, - }, -} - -const modified: WatchEvent = { - type: 'MODIFIED', - object: { - kind: 'ClusterRole', - apiVersion: 'rbac.authorization.k8s.io/v1', - metadata: { name: 'admin', namespace: undefined as unknown as string, resourceVersion: '2' }, - }, -} - -const deleted: WatchEvent = { - type: 'DELETED', - object: { - kind: 'ClusterRole', - apiVersion: 'rbac.authorization.k8s.io/v1', - metadata: { name: 'admin', namespace: undefined as unknown as string, resourceVersion: '2' }, - }, -} - -const pod: WatchEvent = { - type: 'ADDED', - object: { - kind: 'Pod', - apiVersion: 'v1', - metadata: { name: 'nginx', namespace: 'default', resourceVersion: '1' }, - }, -} - -describe('resourceKey', () => { - it('joins namespace and name', () => { - expect(resourceKey(pod.object)).toBe('default/nginx') - }) -}) - -describe('applyWatchEventsToCache', () => { - it('adds, updates, and deletes by namespace/name', () => { - const cache: Record = {} - applyWatchEventsToCache(cache, [added]) - expect(Object.keys(cache)).toEqual(['undefined/admin']) - applyWatchEventsToCache(cache, [modified]) - expect(cache['undefined/admin'].metadata?.resourceVersion).toBe('2') - applyWatchEventsToCache(cache, [deleted]) - expect(cache).toEqual({}) - }) -}) - -describe('groupWatchEventsByKind', () => { - it('groups by apiVersion group and kind', () => { - const grouped = groupWatchEventsByKind([added, pod]) - expect(grouped['rbac.authorization.k8s.io'].ClusterRole).toEqual([added]) - expect(grouped.v1.Pod).toEqual([pod]) - }) -}) diff --git a/frontend/src/hooks/applyWatchEventsToCache.ts b/frontend/src/hooks/applyWatchEventsToCache.ts deleted file mode 100644 index 4db48bf5f3..0000000000 --- a/frontend/src/hooks/applyWatchEventsToCache.ts +++ /dev/null @@ -1,37 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ -// eslint-disable-next-line @typescript-eslint/no-restricted-imports -import type { WatchEvent } from '../atoms' -import type { IResource } from '../resources' - -export function resourceKey(object: WatchEvent['object']): string { - return `${object.metadata.namespace}/${object.metadata.name}` -} - -export function applyWatchEventsToCache(cache: Record, watchEvents: WatchEvent[]): void { - for (const watchEvent of watchEvents) { - const key = resourceKey(watchEvent.object) - switch (watchEvent.type) { - case 'ADDED': - case 'MODIFIED': - cache[key] = watchEvent.object - break - case 'DELETED': - delete cache[key] - break - } - } -} - -export function groupWatchEventsByKind(watchEvents: WatchEvent[]): Record> { - return watchEvents.reduce( - (resourceTypeMap, eventData) => { - const groupVersion = eventData.object.apiVersion.split('/')[0] - const kind = eventData.object.kind - if (!resourceTypeMap[groupVersion]) resourceTypeMap[groupVersion] = {} - if (!resourceTypeMap[groupVersion][kind]) resourceTypeMap[groupVersion][kind] = [] - resourceTypeMap[groupVersion][kind].push(eventData) - return resourceTypeMap - }, - {} as Record> - ) -} diff --git a/frontend/src/hooks/useWatchEventStream.test.ts b/frontend/src/hooks/useWatchEventStream.test.ts deleted file mode 100644 index 7dc5d067f9..0000000000 --- a/frontend/src/hooks/useWatchEventStream.test.ts +++ /dev/null @@ -1,140 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ -import { act, renderHook } from '@testing-library/react-hooks' -import { useRef } from 'react' -import { installFakeEventSource } from '../lib/test-event-source' -import { useWatchEventStream } from './useWatchEventStream' - -jest.mock('../resources/utils', () => ({ - getBackendUrl: () => '', -})) - -describe('useWatchEventStream', () => { - let fake: ReturnType - - beforeEach(() => { - fake = installFakeEventSource() - jest.useFakeTimers() - }) - - afterEach(() => { - jest.useRealTimers() - fake.restore() - }) - - function renderStream(path = '/events/rbac', streamStopped = false) { - const applyWatchEvents = jest.fn() - const onEndOfPacket = jest.fn() - const onLoaded = jest.fn() - const onSettings = jest.fn() - const { result, rerender, unmount } = renderHook( - (props: { path: string; restartKey: number }) => { - const streamStoppedRef = useRef(streamStopped) - const eventSourceRef = useRef() - const processIntervalRef = useRef>() - streamStoppedRef.current = streamStopped - useWatchEventStream({ - path: props.path, - restartKey: props.restartKey, - streamStoppedRef, - eventSourceRef, - processIntervalRef, - applyWatchEvents, - onEndOfPacket, - onLoaded, - onSettings, - }) - return { eventSourceRef, streamStoppedRef } - }, - { initialProps: { path, restartKey: 0 } } - ) - return { applyWatchEvents, onEndOfPacket, onLoaded, onSettings, result, rerender, unmount } - } - - it('opens the path with credentials', () => { - renderStream('/events/rbac') - expect(fake.sources).toHaveLength(1) - expect(fake.sources[0].url).toBe('/events/rbac') - expect(fake.sources[0].withCredentials).toBe(true) - }) - - it('applies ADDED on EOP and calls onEndOfPacket', () => { - const { applyWatchEvents, onEndOfPacket } = renderStream() - const object = { - kind: 'ClusterRole', - apiVersion: 'rbac.authorization.k8s.io/v1', - metadata: { name: 'admin', namespace: '', resourceVersion: '1' }, - } - act(() => { - fake.sources[0].emit({ type: 'START' }) - fake.sources[0].emit({ type: 'ADDED', object }) - fake.sources[0].emit({ type: 'EOP' }) - }) - expect(applyWatchEvents).toHaveBeenCalledWith([expect.objectContaining({ type: 'ADDED', object })]) - expect(onEndOfPacket).toHaveBeenCalledTimes(1) - }) - - it('drops queued events on START', () => { - const { applyWatchEvents } = renderStream() - const object = { - kind: 'ClusterRole', - apiVersion: 'rbac.authorization.k8s.io/v1', - metadata: { name: 'admin', namespace: '', resourceVersion: '1' }, - } - act(() => { - fake.sources[0].emit({ type: 'ADDED', object }) - fake.sources[0].emit({ type: 'START' }) - fake.sources[0].emit({ type: 'EOP' }) - }) - expect(applyWatchEvents).not.toHaveBeenCalled() - }) - - it('calls onLoaded after flushing the queue', () => { - const { applyWatchEvents, onLoaded } = renderStream() - act(() => { - fake.sources[0].emit({ type: 'LOADED' }) - }) - expect(applyWatchEvents).not.toHaveBeenCalled() - expect(onLoaded).toHaveBeenCalledTimes(1) - }) - - it('forwards SETTINGS', () => { - const { onSettings } = renderStream() - act(() => { - fake.sources[0].emit({ type: 'SETTINGS', settings: { FOO: 'bar' } }) - }) - expect(onSettings).toHaveBeenCalledWith({ FOO: 'bar' }) - }) - - it('reconnects after CLOSED unless the stream was stopped', () => { - renderStream('/events', false) - act(() => { - fake.sources[0].triggerError() - }) - expect(fake.sources).toHaveLength(1) - act(() => { - jest.advanceTimersByTime(1000) - }) - expect(fake.sources).toHaveLength(2) - }) - - it('does not reconnect when the stream was stopped for idle', () => { - renderStream('/events', true) - act(() => { - fake.sources[0].triggerError() - jest.advanceTimersByTime(1000) - }) - expect(fake.sources).toHaveLength(1) - }) - - it('does not schedule a second reconnect while one is pending', () => { - renderStream('/events', false) - act(() => { - fake.sources[0].triggerError() - fake.sources[0].triggerError() - }) - act(() => { - jest.advanceTimersByTime(1000) - }) - expect(fake.sources).toHaveLength(2) - }) -}) diff --git a/frontend/src/hooks/useWatchEventStream.ts b/frontend/src/hooks/useWatchEventStream.ts deleted file mode 100644 index 03e4db7ccb..0000000000 --- a/frontend/src/hooks/useWatchEventStream.ts +++ /dev/null @@ -1,113 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ -import { MutableRefObject, useEffect, useRef } from 'react' -// eslint-disable-next-line @typescript-eslint/no-restricted-imports -import { ServerSideEventData, THROTTLE_EVENTS_DELAY, WatchEvent } from '../atoms' -import { getBackendUrl } from '../resources/utils' - -export interface WatchEventStreamHandlers { - applyWatchEvents: (events: WatchEvent[]) => void - onSettings?: (settings: Record) => void - onEndOfPacket?: () => void - onLoaded?: () => void -} - -export interface UseWatchEventStreamOptions extends WatchEventStreamHandlers { - path: string - restartKey: number - streamStoppedRef: MutableRefObject - eventSourceRef: MutableRefObject - processIntervalRef: MutableRefObject | undefined> -} - -export function useWatchEventStream({ - path, - restartKey, - streamStoppedRef, - eventSourceRef, - processIntervalRef, - applyWatchEvents, - onSettings, - onEndOfPacket, - onLoaded, -}: UseWatchEventStreamOptions): void { - const applyWatchEventsRef = useRef(applyWatchEvents) - applyWatchEventsRef.current = applyWatchEvents - const onSettingsRef = useRef(onSettings) - onSettingsRef.current = onSettings - const onEndOfPacketRef = useRef(onEndOfPacket) - onEndOfPacketRef.current = onEndOfPacket - const onLoadedRef = useRef(onLoaded) - onLoadedRef.current = onLoaded - - useEffect(() => { - const eventQueue: WatchEvent[] = [] - - function processEventQueue() { - if (eventQueue.length === 0) return - const watchEvents = eventQueue.splice(0) - applyWatchEventsRef.current(watchEvents) - } - - function processMessage(event: MessageEvent) { - if (!event.data) return - try { - const data = JSON.parse(event.data) as ServerSideEventData - switch (data.type) { - case 'ADDED': - case 'MODIFIED': - case 'DELETED': - eventQueue.push(data) - break - case 'START': - eventQueue.length = 0 - break - case 'EOP': - processEventQueue() - onEndOfPacketRef.current?.() - break - case 'LOADED': - processEventQueue() - onLoadedRef.current?.() - break - case 'SETTINGS': - onSettingsRef.current?.(data.settings) - break - } - } catch (err) { - console.error(err) - } - } - - let evtSource: EventSource | undefined - let reconnectTimer: ReturnType | undefined - - function startWatch() { - evtSource = new EventSource(`${getBackendUrl()}${path}`, { withCredentials: true }) - eventSourceRef.current = evtSource - evtSource.onmessage = processMessage - evtSource.onerror = function () { - console.log('EventSource', 'error', 'readyState', evtSource?.readyState) - if (streamStoppedRef.current) return - if (evtSource?.readyState === EventSource.CLOSED) { - if (reconnectTimer) return - reconnectTimer = setTimeout(() => { - reconnectTimer = undefined - if (streamStoppedRef.current) return - startWatch() - }, 1000) - } - } - } - startWatch() - - const timeout = setInterval(processEventQueue, THROTTLE_EVENTS_DELAY) - processIntervalRef.current = timeout - return () => { - clearInterval(timeout) - if (reconnectTimer) clearTimeout(reconnectTimer) - if (evtSource) evtSource.close() - eventSourceRef.current = undefined - processIntervalRef.current = undefined - } - }, [eventSourceRef, path, processIntervalRef, restartKey, streamStoppedRef]) -} diff --git a/frontend/src/lib/test-event-source.ts b/frontend/src/lib/test-event-source.ts deleted file mode 100644 index 309be17bc9..0000000000 --- a/frontend/src/lib/test-event-source.ts +++ /dev/null @@ -1,48 +0,0 @@ -/* Copyright Contributors to the Open Cluster Management project */ - -export type FakeEventSource = { - url: string - withCredentials: boolean - readyState: number - onmessage: ((ev: MessageEvent) => void) | null - onerror: (() => void) | null - close: jest.Mock - emit: (data: unknown) => void - triggerError: (readyState?: number) => void -} - -export function installFakeEventSource(): { sources: FakeEventSource[]; restore: () => void } { - const sources: FakeEventSource[] = [] - const OriginalEventSource = global.EventSource - global.EventSource = class { - static readonly CONNECTING = 0 - static readonly OPEN = 1 - static readonly CLOSED = 2 - url: string - withCredentials: boolean - readyState = 1 - onmessage: ((ev: MessageEvent) => void) | null = null - onerror: (() => void) | null = null - close = jest.fn() - constructor(url: string | URL, init?: EventSourceInit) { - this.url = url.toString() - this.withCredentials = !!init?.withCredentials - const self = this as unknown as FakeEventSource - self.emit = (data: unknown) => { - this.onmessage?.({ data: JSON.stringify(data) } as MessageEvent) - } - self.triggerError = (readyState = EventSource.CLOSED) => { - this.readyState = readyState - this.onerror?.() - } - sources.push(self) - } - } as unknown as typeof EventSource - - return { - sources, - restore: () => { - global.EventSource = OriginalEventSource - }, - } -} diff --git a/frontend/src/resources/utils/resource-request.ts b/frontend/src/resources/utils/resource-request.ts index d8ae5d2fd8..bdc37f73ca 100644 --- a/frontend/src/resources/utils/resource-request.ts +++ b/frontend/src/resources/utils/resource-request.ts @@ -12,7 +12,7 @@ import { getResourceApiPath, getResourceName, getResourceNameApiPath, IResource, import { Status, StatusKind } from '../status' import { AnsibleTowerInventory, AnsibleTowerInventoryList } from '../ansible-inventory' -// must match ansibletower.Paths in backend/internal/ansibletower +// must match ansiblePaths in backend/src/routes/ansibletower.ts const ansibleControllerPaths = ['/api/v2/job_templates/', '/api/v2/workflow_job_templates/'] // Ansible Automation Platform Operator v2.5 and later only supports the Gateway URL. // For Gateway URLs, use the following path prefixes: diff --git a/frontend/webpack.config.ts b/frontend/webpack.config.ts index 7e0468de10..23362be177 100644 --- a/frontend/webpack.config.ts +++ b/frontend/webpack.config.ts @@ -179,7 +179,6 @@ module.exports = function (env: any, argv: { hot?: boolean; mode: string | undef '/multicloud/configure', '/multicloud/console-links', '/multicloud/events', - '/multicloud/events/rbac', '/multicloud/hub', '/multicloud/upgrade-risks-prediction', '/multicloud/login',