RockteMQ-AI commented on code in PR #3133: URL: https://github.com/apache/rocketmq-dashboard/pull/3133#discussion_r3935597946
########## server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/KubernetesNameServerDiscoveryService.java: ########## @@ -0,0 +1,323 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0. + */ +package org.apache.rocketmq.studio.cluster.nameserver; + +import io.fabric8.kubernetes.api.model.Container; +import io.fabric8.kubernetes.api.model.ContainerPort; +import io.fabric8.kubernetes.api.model.Pod; +import io.fabric8.kubernetes.api.model.Service; +import io.fabric8.kubernetes.api.model.ServicePort; +import io.fabric8.kubernetes.api.model.discovery.v1.Endpoint; +import io.fabric8.kubernetes.api.model.discovery.v1.EndpointPort; +import io.fabric8.kubernetes.api.model.discovery.v1.EndpointSlice; +import io.fabric8.kubernetes.client.KubernetesClientException; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.rocketmq.studio.common.exception.BusinessException; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.util.StringUtils; + +import java.time.LocalDateTime; +import java.util.Collection; +import java.util.Comparator; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Objects; +import java.util.SortedSet; +import java.util.TreeSet; +import java.util.function.Predicate; +import java.util.stream.Stream; + +@Slf4j [email protected] +@RequiredArgsConstructor +@EnableConfigurationProperties(KubernetesNameServerDiscoveryProperties.class) +public class KubernetesNameServerDiscoveryService { + + private static final String SERVICE_NAME_LABEL = "kubernetes.io/service-name"; + + private final KubernetesNameServerDiscoveryClientFactory clientFactory; + private final KubernetesNameServerDiscoveryProperties properties; + + public KubernetesNameServerDiscoveryVO discover(DiscoverKubernetesNameServersDTO command) { + if (command == null || !StringUtils.hasText(command.getNamespace())) { + throw new BusinessException(400, "namespace is required"); + } + if (!properties.isEnabled()) { + throw new BusinessException(503, "Kubernetes NameServer discovery is disabled"); + } + String namespace = command.getNamespace().trim(); + validateConfiguration(); + + try (KubernetesNameServerDiscoveryClient client = clientFactory.create()) { + List<Service> services = client.listServices(namespace); + List<EndpointSlice> endpointSlices = client.listEndpointSlices(namespace); + List<KubernetesNameServerCandidateVO> candidates = servicePortCandidates( + services, endpointSlices, namespace); + if (candidates.isEmpty()) { + candidates = serviceHintCandidates(services, endpointSlices, namespace); + } + if (candidates.isEmpty()) { + candidates = endpointSliceCandidates(endpointSlices, namespace); + } + if (candidates.isEmpty() && properties.isPodFallbackEnabled()) { + candidates = podCandidates(client.listPodsByComponent(namespace), namespace, "POD_LABEL", + pod -> true); + } + if (candidates.isEmpty() && properties.isPodFallbackEnabled()) { + candidates = podCandidates(client.listRocketMqPods(namespace), namespace, "POD_LABEL", + this::hasNameServerRoleEvidence); + } + if (candidates.isEmpty() && properties.isPodFallbackEnabled()) { + candidates = podCandidates(client.listPods(namespace), namespace, "POD_IMAGE", + pod -> hasImageHint(pod) && hasNameServerRoleEvidence(pod)); + } + return KubernetesNameServerDiscoveryVO.builder() + .namespace(namespace) + .observedAt(LocalDateTime.now()) + .candidates(limitAndDeduplicate(candidates)) + .build(); + } catch (BusinessException exception) { + throw exception; + } catch (KubernetesClientException exception) { + int status = exception.getCode() == 403 ? 403 : 503; + String message = status == 403 + ? "Kubernetes RBAC denied NameServer discovery in namespace " + namespace + : "Kubernetes API is unavailable for NameServer discovery"; + throw new BusinessException(status, message); + } catch (RuntimeException exception) { + log.warn("Kubernetes NameServer discovery failed in namespace {}: {}", namespace, Review Comment: The fallback chain (SERVICE_PORT -> SERVICE_HINT -> ENDPOINT_SLICE -> POD_LABEL -> POD_IMAGE) is well-designed. One minor suggestion: when the POD_IMAGE fallback fires, consider adding a log line at INFO level so operators can see that the less-reliable path was used. ########## server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/KubernetesNameServerDiscoveryService.java: ########## @@ -0,0 +1,323 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0. + */ +package org.apache.rocketmq.studio.cluster.nameserver; + +import io.fabric8.kubernetes.api.model.Container; +import io.fabric8.kubernetes.api.model.ContainerPort; +import io.fabric8.kubernetes.api.model.Pod; +import io.fabric8.kubernetes.api.model.Service; +import io.fabric8.kubernetes.api.model.ServicePort; +import io.fabric8.kubernetes.api.model.discovery.v1.Endpoint; +import io.fabric8.kubernetes.api.model.discovery.v1.EndpointPort; +import io.fabric8.kubernetes.api.model.discovery.v1.EndpointSlice; +import io.fabric8.kubernetes.client.KubernetesClientException; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.rocketmq.studio.common.exception.BusinessException; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.util.StringUtils; + +import java.time.LocalDateTime; +import java.util.Collection; +import java.util.Comparator; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Objects; +import java.util.SortedSet; +import java.util.TreeSet; +import java.util.function.Predicate; +import java.util.stream.Stream; + +@Slf4j [email protected] +@RequiredArgsConstructor +@EnableConfigurationProperties(KubernetesNameServerDiscoveryProperties.class) +public class KubernetesNameServerDiscoveryService { + + private static final String SERVICE_NAME_LABEL = "kubernetes.io/service-name"; + + private final KubernetesNameServerDiscoveryClientFactory clientFactory; + private final KubernetesNameServerDiscoveryProperties properties; + + public KubernetesNameServerDiscoveryVO discover(DiscoverKubernetesNameServersDTO command) { + if (command == null || !StringUtils.hasText(command.getNamespace())) { + throw new BusinessException(400, "namespace is required"); + } + if (!properties.isEnabled()) { + throw new BusinessException(503, "Kubernetes NameServer discovery is disabled"); + } + String namespace = command.getNamespace().trim(); + validateConfiguration(); + + try (KubernetesNameServerDiscoveryClient client = clientFactory.create()) { + List<Service> services = client.listServices(namespace); + List<EndpointSlice> endpointSlices = client.listEndpointSlices(namespace); + List<KubernetesNameServerCandidateVO> candidates = servicePortCandidates( + services, endpointSlices, namespace); + if (candidates.isEmpty()) { + candidates = serviceHintCandidates(services, endpointSlices, namespace); + } + if (candidates.isEmpty()) { + candidates = endpointSliceCandidates(endpointSlices, namespace); + } Review Comment: The discover method creates a KubernetesClient via clientFactory.create() and relies on try-with-resources to close it. Consider logging a warning if close() throws, since a failed close on the HTTP transport could leak connections in long-running Studio deployments. ########## web/src/pages/cluster/index.tsx: ########## @@ -293,6 +305,68 @@ const ClusterPage = () => { [nsCreateForm], ); + const openKubernetesDiscovery = useCallback(() => { + kubernetesDiscoveryRequestRef.current += 1; + setKubernetesDiscoveryCandidates([]); + setKubernetesDiscoverySearched(false); + setKubernetesDiscoveryLoading(false); + kubernetesDiscoveryForm.resetFields(); + setKubernetesDiscoveryOpen(true); + }, [kubernetesDiscoveryForm]); + + const closeKubernetesDiscovery = useCallback(() => { + kubernetesDiscoveryRequestRef.current += 1; + setKubernetesDiscoveryOpen(false); + setKubernetesDiscoveryLoading(false); + }, []); + + const handleKubernetesDiscovery = useCallback(async () => { + let namespace: string; + try { + ({ namespace } = await kubernetesDiscoveryForm.validateFields()); + } catch { + return; + } + const requestId = ++kubernetesDiscoveryRequestRef.current; + setKubernetesDiscoveryLoading(true); + try { + const result = await discoverKubernetesNameServers(namespace.trim()); + if (requestId !== kubernetesDiscoveryRequestRef.current) return; + setKubernetesDiscoveryCandidates(result.candidates); + setKubernetesDiscoverySearched(true); + if (result.candidates.length === 0) { + message.info(t('cluster.k8sDiscoveryEmpty')); + } + } catch (error) { Review Comment: Good use of kubernetesDiscoveryRequestRef to handle stale responses. The pattern correctly invalidates in-flight requests when the modal is closed/reopened. Consider extracting this request-ID pattern into a custom hook if reused elsewhere. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
