This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new ed6613bd1 [ISSUE #4463] fix(producer): surface incomplete connection
scans (#4468)
ed6613bd1 is described below
commit ed6613bd11df3ad2c169169955da93687edcec8d
Author: aias00 <[email protected]>
AuthorDate: Mon Sep 21 05:59:21 2026 -0700
[ISSUE #4463] fix(producer): surface incomplete connection scans (#4468)
The topic-wide producer scan in `RocketMQClientProvider` logged a warning
and moved on when a broker's producer table or a group's connection query
failed, and threw 502 only when every group failed. Nothing in the response
said so, so `GET /api/producer/connection?topic=X` returned a short list that
looked complete — during a rolling restart the producers on an unreachable
broker simply vanished from the panel.
The scan now returns a `ProducerConnectionScanResult` carrying the
connections plus `failedBrokers` and `failedProducerGroups`, surfaced on
`ProducerConnectionResultVO` as `complete` alongside both lists, and
`ProducerConnectionSummaryVO` gains an `INCOMPLETE_SCAN` warning that forces
`readiness` to `WARNING`. The Producer page tags each failed broker and group
in the readiness Alert and no longer reports "no connections" for an incomplete
scan. Partial results are kept rather than re [...]
The producer group selector behind `/api/producer/group` still returns a
partial list with no marker. Same defect class, but marking it means changing
that endpoint's contract, so it is left for a follow-up.
Fixes #4463
---
docs/api-spec.md | 10 ++-
.../studio/cluster/client/ClientProvider.java | 6 ++
.../cluster/client/ProducerConnectionResultVO.java | 19 +++++-
.../client/ProducerConnectionScanResult.java | 41 ++++++++++++
.../cluster/client/ProducerConnectionService.java | 9 ++-
.../client/ProducerConnectionSummaryVO.java | 26 ++++++--
.../studio/cluster/client/ProducerController.java | 3 +-
.../provider/apache/RocketMQClientProvider.java | 53 ++++++++++++---
.../client/ProducerConnectionServiceTest.java | 60 +++++++++++------
.../client/ProducerConnectionSummaryVOTest.java | 9 +++
.../cluster/client/ProducerControllerTest.java | 30 +++++++--
.../apache/RocketMQClientProviderTest.java | 78 ++++++++++++++++++++--
web/src/api/producer.test.ts | 17 +++++
web/src/api/producer.ts | 38 ++++++++++-
web/src/i18n/translations.ts | 12 ++++
web/src/pages/studio/Producer.tsx | 21 +++++-
web/src/pages/studio/__tests__/Producer.test.tsx | 45 +++++++++++++
17 files changed, 420 insertions(+), 57 deletions(-)
diff --git a/docs/api-spec.md b/docs/api-spec.md
index a0754b29e..e5588ecba 100644
--- a/docs/api-spec.md
+++ b/docs/api-spec.md
@@ -1622,6 +1622,14 @@ GET /api/producer/connection
|------|------|------|
| `connectionSet` | `ProducerConnection[]` | 匹配的 Producer 连接 |
| `summary` | `ProducerConnectionSummary` | 连接完整性和分布摘要 |
+| `complete` | `boolean` | 是否覆盖全部可发现的 Broker 与 Producer Group |
+| `failedBrokers` | `string[]` | Producer Group 发现失败的 Broker 地址;完整扫描时为空 |
+| `failedProducerGroups` | `string[]` | 连接查询失败的 Producer Group;完整扫描时为空 |
+
+省略 `producerGroup` 的 Topic 聚合查询采用部分成功语义:只要至少一个 Broker 可查询,且并非
+所有已发现的 Producer Group 都查询失败,就返回可用连接;覆盖缺口通过 `complete=false` 和
+失败列表公开。所有 Broker 或所有已发现 Group 都无法查询时仍返回 `502`。显式
+`producerGroup` 查询保持原有语义。
#### ProducerConnection
@@ -1646,7 +1654,7 @@ GET /api/producer/connection
| `languages` | `{ value: string, count: number }[]` | 按客户端语言统计的分布 |
| `versions` | `{ value: string, count: number }[]` | 按客户端版本统计的分布 |
| `duplicateClientIds` | `string[]` | 重复出现的客户端 ID |
-| `warnings` | `string[]` | `NO_CONNECTIONS` / `DUPLICATE_CLIENT_ID` /
`MIXED_CLIENT_VERSION` / `INCOMPLETE_CLIENT_METADATA` |
+| `warnings` | `string[]` | `NO_CONNECTIONS` / `DUPLICATE_CLIENT_ID` /
`MIXED_CLIENT_VERSION` / `INCOMPLETE_CLIENT_METADATA` / `INCOMPLETE_SCAN` |
| `readiness` | `string` | `READY` / `WARNING` / `UNAVAILABLE` |
---
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
index b04736009..bff3853f9 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
@@ -30,4 +30,10 @@ public interface ClientProvider {
List<String> findProducerGroups(String instanceId, String topic, String
query, int limit);
List<ClientConnectionVO> findProducerConnections(String instanceId, String
topic, String producerGroup);
+
+ default ProducerConnectionScanResult scanProducerConnections(
+ String instanceId, String topic, String producerGroup) {
+ return ProducerConnectionScanResult.complete(
+ findProducerConnections(instanceId, topic, producerGroup));
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionResultVO.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionResultVO.java
index 5a77d3fce..0d7cfd698 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionResultVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionResultVO.java
@@ -26,9 +26,24 @@ import java.util.List;
public class ProducerConnectionResultVO {
private List<ProducerConnectionVO> connectionSet;
private ProducerConnectionSummaryVO summary;
+ private boolean complete = true;
+ private List<String> failedBrokers = List.of();
+ private List<String> failedProducerGroups = List.of();
public ProducerConnectionResultVO(List<ProducerConnectionVO>
connectionSet) {
- this.connectionSet = connectionSet == null ? List.of() : connectionSet;
- this.summary = ProducerConnectionSummaryVO.from(this.connectionSet);
+ this(connectionSet, true, List.of(), List.of());
+ }
+
+ public ProducerConnectionResultVO(
+ List<ProducerConnectionVO> connectionSet,
+ boolean complete,
+ List<String> failedBrokers,
+ List<String> failedProducerGroups) {
+ this.connectionSet = connectionSet == null ? List.of() :
List.copyOf(connectionSet);
+ this.complete = complete;
+ this.failedBrokers = failedBrokers == null ? List.of() :
List.copyOf(failedBrokers);
+ this.failedProducerGroups = failedProducerGroups == null
+ ? List.of() : List.copyOf(failedProducerGroups);
+ this.summary = ProducerConnectionSummaryVO.from(this.connectionSet,
complete);
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionScanResult.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionScanResult.java
new file mode 100644
index 000000000..ee8789f03
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionScanResult.java
@@ -0,0 +1,41 @@
+/*
+ * 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
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.cluster.client;
+
+import java.util.List;
+
+/** Producer connection rows plus the coverage gaps encountered while scanning
them. */
+public record ProducerConnectionScanResult(
+ List<ClientConnectionVO> connections,
+ List<String> failedBrokers,
+ List<String> failedProducerGroups) {
+
+ public ProducerConnectionScanResult {
+ connections = connections == null ? List.of() :
List.copyOf(connections);
+ failedBrokers = failedBrokers == null ? List.of() :
List.copyOf(failedBrokers);
+ failedProducerGroups = failedProducerGroups == null
+ ? List.of() : List.copyOf(failedProducerGroups);
+ }
+
+ public static ProducerConnectionScanResult
complete(List<ClientConnectionVO> connections) {
+ return new ProducerConnectionScanResult(connections, List.of(),
List.of());
+ }
+
+ public boolean complete() {
+ return failedBrokers.isEmpty() && failedProducerGroups.isEmpty();
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
index b1f8ca293..5514e34b8 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
@@ -33,16 +33,21 @@ public class ProducerConnectionService {
private final ClientProvider clientProvider;
- public List<ProducerConnectionVO> listConnections(String instanceId,
String topic, String producerGroup) {
+ public ProducerConnectionResultVO listConnections(
+ String instanceId, String topic, String producerGroup) {
log.info("Listing producer connections, instanceId={}, topic={},
producerGroup={}",
instanceId, topic, producerGroup);
String normalizedInstanceId = requireFilter(instanceId, "instanceId");
String normalizedTopic = requireFilter(topic, "topic");
String normalizedProducerGroup =
normalizeOptionalFilter(producerGroup);
- return clientProvider.findProducerConnections(normalizedInstanceId,
normalizedTopic, normalizedProducerGroup)
+ ProducerConnectionScanResult scan =
clientProvider.scanProducerConnections(
+ normalizedInstanceId, normalizedTopic,
normalizedProducerGroup);
+ List<ProducerConnectionVO> connections = scan.connections()
.stream()
.map(this::toProducerConnection)
.toList();
+ return new ProducerConnectionResultVO(
+ connections, scan.complete(), scan.failedBrokers(),
scan.failedProducerGroups());
}
public List<String> listProducerGroups(String instanceId, String topic,
String query, Integer limit) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVO.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVO.java
index 75362d5e7..0b16e98b0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVO.java
@@ -36,6 +36,7 @@ public class ProducerConnectionSummaryVO {
public static final String DUPLICATE_CLIENT_ID = "DUPLICATE_CLIENT_ID";
public static final String MIXED_CLIENT_VERSION = "MIXED_CLIENT_VERSION";
public static final String INCOMPLETE_CLIENT_METADATA =
"INCOMPLETE_CLIENT_METADATA";
+ public static final String INCOMPLETE_SCAN = "INCOMPLETE_SCAN";
private int totalConnections;
private int uniqueClientCount;
@@ -49,6 +50,11 @@ public class ProducerConnectionSummaryVO {
private String readiness = READY;
public static ProducerConnectionSummaryVO from(List<ProducerConnectionVO>
connections) {
+ return from(connections, true);
+ }
+
+ public static ProducerConnectionSummaryVO from(
+ List<ProducerConnectionVO> connections, boolean complete) {
List<ProducerConnectionVO> safeConnections = connections == null
? List.of()
: connections.stream().filter(Objects::nonNull).toList();
@@ -61,8 +67,8 @@ public class ProducerConnectionSummaryVO {
summary.uniqueLanguageCount = summary.languages.size();
summary.uniqueVersionCount = summary.versions.size();
summary.duplicateClientIds = duplicateClientIds(safeConnections);
- summary.warnings = warnings(summary, safeConnections);
- summary.readiness = readiness(summary);
+ summary.warnings = warnings(summary, safeConnections, complete);
+ summary.readiness = readiness(summary, complete);
return summary;
}
@@ -105,11 +111,13 @@ public class ProducerConnectionSummaryVO {
}
private static List<String> warnings(
- ProducerConnectionSummaryVO summary, List<ProducerConnectionVO>
connections) {
+ ProducerConnectionSummaryVO summary,
+ List<ProducerConnectionVO> connections,
+ boolean complete) {
+ List<String> warnings = new java.util.ArrayList<>();
if (summary.totalConnections == 0) {
- return List.of(NO_CONNECTIONS);
+ warnings.add(NO_CONNECTIONS);
}
- List<String> warnings = new java.util.ArrayList<>();
if (!summary.duplicateClientIds.isEmpty()) {
warnings.add(DUPLICATE_CLIENT_ID);
}
@@ -119,10 +127,16 @@ public class ProducerConnectionSummaryVO {
if
(connections.stream().anyMatch(ProducerConnectionSummaryVO::hasIncompleteMetadata))
{
warnings.add(INCOMPLETE_CLIENT_METADATA);
}
+ if (!complete) {
+ warnings.add(INCOMPLETE_SCAN);
+ }
return List.copyOf(warnings);
}
- private static String readiness(ProducerConnectionSummaryVO summary) {
+ private static String readiness(ProducerConnectionSummaryVO summary,
boolean complete) {
+ if (!complete) {
+ return WARNING;
+ }
if (summary.totalConnections == 0) {
return UNAVAILABLE;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
index 61a6ba544..c5ff41307 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
@@ -49,8 +49,7 @@ public class ProducerController {
@RequestParam(required = false) String producerGroup) {
requireParameter(instanceId, "instanceId");
requireParameter(topic, "topic");
- return new ProducerConnectionResultVO(
- producerConnectionService.listConnections(instanceId, topic,
producerGroup));
+ return producerConnectionService.listConnections(instanceId, topic,
producerGroup);
}
private void requireParameter(String value, String name) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
index 63b7cdcc5..9978ee2a5 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
@@ -30,6 +30,7 @@ import
org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import org.apache.rocketmq.studio.cluster.client.ClientConnectionVO;
import org.apache.rocketmq.studio.cluster.client.ClientProvider;
+import org.apache.rocketmq.studio.cluster.client.ProducerConnectionScanResult;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
@@ -94,8 +95,14 @@ public class RocketMQClientProvider implements
ClientProvider {
@Override
public List<ClientConnectionVO> findProducerConnections(String instanceId,
String topic, String producerGroup) {
+ return scanProducerConnections(instanceId, topic,
producerGroup).connections();
+ }
+
+ @Override
+ public ProducerConnectionScanResult scanProducerConnections(
+ String instanceId, String topic, String producerGroup) {
return runtimeAdminClientResolver.execute(instanceId,
- adminExt -> findProducerConnections(adminExt, topic,
producerGroup));
+ adminExt -> scanProducerConnections(adminExt, topic,
producerGroup));
}
@Override
@@ -105,12 +112,17 @@ public class RocketMQClientProvider implements
ClientProvider {
}
private List<String> findProducerGroups(MQAdminExt adminExt, String topic,
String query, int limit) {
+ return scanProducerGroups(adminExt, query, limit).groups();
+ }
+
+ private ProducerGroupScanResult scanProducerGroups(MQAdminExt adminExt,
String query, int limit) {
BrokerTopology topology = discoverBrokerTopology(adminExt, null,
"producer group selector");
if (topology.brokerAddresses().isEmpty()) {
- return List.of();
+ return new ProducerGroupScanResult(List.of(), List.of());
}
String normalizedQuery = query == null ? null :
query.toLowerCase(Locale.ROOT);
LinkedHashSet<String> groups = new LinkedHashSet<>();
+ List<String> failedBrokers = new ArrayList<>();
int successfulBrokers = 0;
for (String brokerAddress : topology.brokerAddresses()) {
try {
@@ -118,44 +130,56 @@ public class RocketMQClientProvider implements
ClientProvider {
successfulBrokers++;
collectProducerGroups(groups, producerTable, normalizedQuery);
} catch (Exception e) {
+ failedBrokers.add(brokerAddress);
log.warn("Failed to fetch producer groups from broker={},
skipping", brokerAddress, e);
}
}
if (successfulBrokers == 0) {
throw new BusinessException(502, "Failed to query producer groups
from all brokers");
}
- return groups.stream()
+ List<String> matchedGroups = groups.stream()
.sorted(Comparator.naturalOrder())
.limit(limit)
.toList();
+ return new ProducerGroupScanResult(matchedGroups, failedBrokers);
}
- private List<ClientConnectionVO> findProducerConnections(MQAdminExt
adminExt, String topic, String producerGroup) {
+ private ProducerConnectionScanResult scanProducerConnections(
+ MQAdminExt adminExt, String topic, String producerGroup) {
if (producerGroup == null || producerGroup.isBlank()) {
return findProducerConnectionsForActiveGroups(adminExt, topic);
}
- return findProducerConnectionsForGroup(adminExt, topic, producerGroup);
+ return ProducerConnectionScanResult.complete(
+ findProducerConnectionsForGroup(adminExt, topic,
producerGroup));
}
- private List<ClientConnectionVO>
findProducerConnectionsForActiveGroups(MQAdminExt adminExt, String topic) {
- List<String> producerGroups = findProducerGroups(adminExt, topic,
null, Integer.MAX_VALUE);
+ private ProducerConnectionScanResult
findProducerConnectionsForActiveGroups(
+ MQAdminExt adminExt, String topic) {
+ ProducerGroupScanResult groupScan = scanProducerGroups(adminExt, null,
Integer.MAX_VALUE);
+ List<String> producerGroups = groupScan.groups();
if (producerGroups.isEmpty()) {
- return List.of();
+ return new ProducerConnectionScanResult(
+ List.of(), groupScan.failedBrokers(), List.of());
}
List<ClientConnectionVO> connections = new ArrayList<>();
+ List<String> failedProducerGroups = new ArrayList<>();
int successfulGroupQueries = 0;
for (String producerGroup : producerGroups) {
try {
connections.addAll(findProducerConnectionsForGroup(adminExt,
topic, producerGroup));
successfulGroupQueries++;
} catch (BusinessException e) {
- log.warn("Failed to query producer connections for group={},
skipping", producerGroup, e);
+ failedProducerGroups.add(producerGroup);
+ log.warn("Failed to query producer connections for group={},
skipping: {}",
+ producerGroup, rootMessage(e));
}
}
if (successfulGroupQueries == 0) {
- throw new BusinessException(502, "Failed to query producer
connections from all groups");
+ throw new BusinessException(502,
+ "Failed to query producer connections from all groups");
}
- return connections;
+ return new ProducerConnectionScanResult(
+ connections, groupScan.failedBrokers(), failedProducerGroups);
}
private List<ClientConnectionVO> findProducerConnectionsForGroup(
@@ -460,4 +484,11 @@ public class RocketMQClientProvider implements
ClientProvider {
return clusterByAddress.get(brokerAddress);
}
}
+
+ private record ProducerGroupScanResult(List<String> groups, List<String>
failedBrokers) {
+ private ProducerGroupScanResult {
+ groups = List.copyOf(groups);
+ failedBrokers = List.copyOf(failedBrokers);
+ }
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
index 47f569207..992ca2a91 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
@@ -53,20 +53,22 @@ class ProducerConnectionServiceTest {
.language(ClientLanguage.Java)
.version("5.1.0")
.build();
- when(clientProvider.findProducerConnections("instance-1",
"order-topic", "pg-order"))
- .thenReturn(List.of(producer));
+ when(clientProvider.scanProducerConnections("instance-1",
"order-topic", "pg-order"))
+
.thenReturn(ProducerConnectionScanResult.complete(List.of(producer)));
- List<ProducerConnectionVO> result =
producerConnectionService.listConnections(
+ ProducerConnectionResultVO result =
producerConnectionService.listConnections(
"instance-1", "order-topic", "pg-order");
- assertThat(result).hasSize(1);
- assertThat(result.get(0).getClientId()).isEqualTo("producer-1");
- assertThat(result.get(0).getClientAddr()).isEqualTo("10.0.0.1:38888");
- assertThat(result.get(0).getTopic()).isEqualTo("order-topic");
- assertThat(result.get(0).getProducerGroup()).isEqualTo("pg-order");
- assertThat(result.get(0).getLanguage()).isEqualTo("Java");
- assertThat(result.get(0).getVersionDesc()).isEqualTo("5.1.0");
- verify(clientProvider).findProducerConnections("instance-1",
"order-topic", "pg-order");
+
assertThat(result.getConnectionSet()).singleElement().satisfies(connection -> {
+ assertThat(connection.getClientId()).isEqualTo("producer-1");
+ assertThat(connection.getClientAddr()).isEqualTo("10.0.0.1:38888");
+ assertThat(connection.getTopic()).isEqualTo("order-topic");
+ assertThat(connection.getProducerGroup()).isEqualTo("pg-order");
+ assertThat(connection.getLanguage()).isEqualTo("Java");
+ assertThat(connection.getVersionDesc()).isEqualTo("5.1.0");
+ });
+ assertThat(result.isComplete()).isTrue();
+ verify(clientProvider).scanProducerConnections("instance-1",
"order-topic", "pg-order");
}
@Test
@@ -80,25 +82,41 @@ class ProducerConnectionServiceTest {
@Test
void listConnectionsShouldAllowMissingProducerGroupForAllGroupScan() {
- when(clientProvider.findProducerConnections("instance-1",
"order-topic", null))
- .thenReturn(List.of());
+ when(clientProvider.scanProducerConnections("instance-1",
"order-topic", null))
+ .thenReturn(ProducerConnectionScanResult.complete(List.of()));
- List<ProducerConnectionVO> result =
+ ProducerConnectionResultVO result =
producerConnectionService.listConnections("instance-1",
"order-topic", " ");
- assertThat(result).isEmpty();
- verify(clientProvider).findProducerConnections("instance-1",
"order-topic", null);
+ assertThat(result.getConnectionSet()).isEmpty();
+ verify(clientProvider).scanProducerConnections("instance-1",
"order-topic", null);
}
@Test
void listConnectionsShouldTrimRequiredValues() {
- when(clientProvider.findProducerConnections("instance-1",
"order-topic", "pg-order"))
- .thenReturn(List.of());
+ when(clientProvider.scanProducerConnections("instance-1",
"order-topic", "pg-order"))
+ .thenReturn(ProducerConnectionScanResult.complete(List.of()));
- List<ProducerConnectionVO> result =
producerConnectionService.listConnections(
+ ProducerConnectionResultVO result =
producerConnectionService.listConnections(
" instance-1 ", " order-topic ", " pg-order ");
- assertThat(result).isEmpty();
- verify(clientProvider).findProducerConnections("instance-1",
"order-topic", "pg-order");
+ assertThat(result.getConnectionSet()).isEmpty();
+ verify(clientProvider).scanProducerConnections("instance-1",
"order-topic", "pg-order");
+ }
+
+ @Test
+ void listConnectionsShouldPreservePartialScanMetadataTest() {
+ when(clientProvider.scanProducerConnections("instance-1",
"order-topic", null))
+ .thenReturn(new ProducerConnectionScanResult(
+ List.of(), List.of("broker-a:10911"),
List.of("pg-orders")));
+
+ ProducerConnectionResultVO result =
+ producerConnectionService.listConnections("instance-1",
"order-topic", null);
+
+ assertThat(result.isComplete()).isFalse();
+
assertThat(result.getFailedBrokers()).containsExactly("broker-a:10911");
+
assertThat(result.getFailedProducerGroups()).containsExactly("pg-orders");
+ assertThat(result.getSummary().getWarnings())
+ .contains(ProducerConnectionSummaryVO.INCOMPLETE_SCAN);
}
@Test
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVOTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVOTest.java
index 14065d5cb..b0059bcd8 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVOTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVOTest.java
@@ -79,6 +79,15 @@ class ProducerConnectionSummaryVOTest {
.containsExactly("UNKNOWN");
}
+ @Test
+ void fromShouldWarnWhenTheProducerScanIsIncompleteTest() {
+ ProducerConnectionSummaryVO summary =
ProducerConnectionSummaryVO.from(List.of(
+ connection("producer-a", "10.0.0.1:38888", "Java", "5.1.0")),
false);
+
+
assertThat(summary.getReadiness()).isEqualTo(ProducerConnectionSummaryVO.WARNING);
+
assertThat(summary.getWarnings()).containsExactly(ProducerConnectionSummaryVO.INCOMPLETE_SCAN);
+ }
+
private ProducerConnectionVO connection(String clientId, String address,
String language, String version) {
return ProducerConnectionVO.builder()
.clientId(clientId)
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
index 0b4688639..c70f26a6f 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
@@ -76,7 +76,7 @@ class ProducerControllerTest extends WebMvcAuthTestSupport {
.versionDesc("5.1.0")
.build();
when(producerConnectionService.listConnections("instance-1",
"order-topic", "pg-order"))
- .thenReturn(List.of(connection));
+ .thenReturn(new
ProducerConnectionResultVO(List.of(connection)));
mockMvc.perform(get("/api/producer/connection")
.param("instanceId", "instance-1")
@@ -93,7 +93,10 @@ class ProducerControllerTest extends WebMvcAuthTestSupport {
.andExpect(jsonPath("$.summary.totalConnections").value(1))
.andExpect(jsonPath("$.summary.uniqueClientCount").value(1))
.andExpect(jsonPath("$.summary.uniqueAddressCount").value(1))
- .andExpect(jsonPath("$.summary.readiness").value("READY"));
+ .andExpect(jsonPath("$.summary.readiness").value("READY"))
+ .andExpect(jsonPath("$.complete").value(true))
+ .andExpect(jsonPath("$.failedBrokers").isEmpty())
+ .andExpect(jsonPath("$.failedProducerGroups").isEmpty());
verify(producerConnectionService).listConnections("instance-1",
"order-topic", "pg-order");
}
@@ -113,7 +116,7 @@ class ProducerControllerTest extends WebMvcAuthTestSupport {
@Test
void listConnectionsShouldAllowMissingProducerGroup() throws Exception {
when(producerConnectionService.listConnections("instance-1",
"order-topic", null))
- .thenReturn(List.of());
+ .thenReturn(new ProducerConnectionResultVO(List.of()));
mockMvc.perform(get("/api/producer/connection")
.param("instanceId", "instance-1")
@@ -121,11 +124,30 @@ class ProducerControllerTest extends
WebMvcAuthTestSupport {
.andExpect(status().isOk())
.andExpect(jsonPath("$.connectionSet").isArray())
.andExpect(jsonPath("$.summary.totalConnections").value(0))
-
.andExpect(jsonPath("$.summary.readiness").value("UNAVAILABLE"));
+
.andExpect(jsonPath("$.summary.readiness").value("UNAVAILABLE"))
+ .andExpect(jsonPath("$.complete").value(true));
verify(producerConnectionService).listConnections("instance-1",
"order-topic", null);
}
+ @Test
+ void listConnectionsShouldExposePartialScanMetadataTest() throws Exception
{
+ when(producerConnectionService.listConnections("instance-1",
"order-topic", null))
+ .thenReturn(new ProducerConnectionResultVO(
+ List.of(), false, List.of("broker-a:10911"),
List.of("pg-orders")));
+
+ mockMvc.perform(get("/api/producer/connection")
+ .param("instanceId", "instance-1")
+ .param("topic", "order-topic"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.complete").value(false))
+
.andExpect(jsonPath("$.failedBrokers[0]").value("broker-a:10911"))
+
.andExpect(jsonPath("$.failedProducerGroups[0]").value("pg-orders"))
+ .andExpect(jsonPath("$.summary.readiness").value("WARNING"))
+
.andExpect(jsonPath("$.summary.warnings[0]").value("NO_CONNECTIONS"))
+
.andExpect(jsonPath("$.summary.warnings[1]").value("INCOMPLETE_SCAN"));
+ }
+
@Test
void listConnectionsShouldRejectBlankParameters() throws Exception {
mockMvc.perform(get("/api/producer/connection")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
index a732b6ebd..2bf3a429b 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
@@ -29,6 +29,7 @@ import
org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
import org.apache.rocketmq.studio.cluster.client.ClientConnectionVO;
+import org.apache.rocketmq.studio.cluster.client.ProducerConnectionScanResult;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
@@ -262,6 +263,21 @@ class RocketMQClientProviderTest {
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
}
+ @Test
+ void producerGroupSelectorKeepsBestEffortResultsWhenOneBrokerFailsTest()
throws Exception {
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo(
+ "127.0.0.1:10911", "127.0.0.2:10911"));
+ when(adminExt.getAllProducerInfo("127.0.0.1:10911"))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+ when(adminExt.getAllProducerInfo("127.0.0.2:10911"))
+ .thenReturn(new ProducerTableInfo(Map.of(
+ "pg-payment", List.of(producerInfo("producer-payment",
"10.0.0.2:1000")))));
+
+ List<String> groups = provider.findProducerGroups("instance-a",
"TopicA", "pg", 20);
+
+ assertThat(groups).containsExactly("pg-payment");
+ }
+
@Test
void exactProducerQueryPassesNonBlankGroupToAdminApi() throws Exception {
ProducerConnection producerConnection = new ProducerConnection();
@@ -313,26 +329,80 @@ class RocketMQClientProviderTest {
}
@Test
- void producerQueryWithoutGroupReturnsPartialResultsWhenOneGroupFails()
throws Exception {
+ void
producerQueryWithoutGroupReturnsPartialResultWhenOneGroupQueryFailsTest()
throws Exception {
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo("127.0.0.1:10911"));
when(adminExt.getAllProducerInfo("127.0.0.1:10911"))
.thenReturn(new ProducerTableInfo(Map.of(
"pg-order", List.of(producerInfo("producer-order",
"10.0.0.1:1000")),
"pg-payment", List.of(producerInfo("producer-payment",
"10.0.0.2:1000")))));
+ when(adminExt.examineProducerConnectionInfo("pg-order", "TopicA"))
+ .thenThrow(new IllegalStateException("broker unavailable"));
ProducerConnection paymentConnection = new ProducerConnection();
paymentConnection.setConnectionSet(new HashSet<>(List.of(
connection("producer-payment", "10.0.0.2:1000"))));
- when(adminExt.examineProducerConnectionInfo("pg-order", "TopicA"))
+ when(adminExt.examineProducerConnectionInfo("pg-payment", "TopicA"))
+ .thenReturn(paymentConnection);
+
+ ProducerConnectionScanResult result =
+ provider.scanProducerConnections("instance-a", "TopicA", null);
+
+ assertThat(result.connections()).singleElement().satisfies(connection
->
+
assertThat(connection.getProducerGroup()).isEqualTo("pg-payment"));
+ assertThat(result.complete()).isFalse();
+ assertThat(result.failedBrokers()).isEmpty();
+ assertThat(result.failedProducerGroups()).containsExactly("pg-order");
+ }
+
+ @Test
+ void
producerQueryWithoutGroupReturnsPartialResultWhenOneBrokerGroupDiscoveryFailsTest()
throws Exception {
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo(
+ "127.0.0.1:10911", "127.0.0.2:10911"));
+ when(adminExt.getAllProducerInfo("127.0.0.1:10911"))
.thenThrow(new IllegalStateException("broker unavailable"));
+ when(adminExt.getAllProducerInfo("127.0.0.2:10911"))
+ .thenReturn(new ProducerTableInfo(Map.of(
+ "pg-payment", List.of(producerInfo("producer-payment",
"10.0.0.2:1000")))));
+ ProducerConnection paymentConnection = new ProducerConnection();
+ paymentConnection.setConnectionSet(new HashSet<>(List.of(
+ connection("producer-payment", "10.0.0.2:1000"))));
when(adminExt.examineProducerConnectionInfo("pg-payment", "TopicA"))
.thenReturn(paymentConnection);
- List<ClientConnectionVO> connections =
provider.findProducerConnections("instance-a", "TopicA", null);
+ ProducerConnectionScanResult result =
+ provider.scanProducerConnections("instance-a", "TopicA", null);
- assertThat(connections).singleElement().satisfies(connection -> {
+ assertThat(result.connections()).singleElement().satisfies(connection
->
+
assertThat(connection.getProducerGroup()).isEqualTo("pg-payment"));
+ assertThat(result.complete()).isFalse();
+ assertThat(result.failedBrokers()).containsExactly("127.0.0.1:10911");
+ assertThat(result.failedProducerGroups()).isEmpty();
+ }
+
+ @Test
+ void
producerQueryWithoutGroupTreatsOfflineGroupAsCompleteEmptyResultTest() throws
Exception {
+
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo("127.0.0.1:10911"));
+ when(adminExt.getAllProducerInfo("127.0.0.1:10911"))
+ .thenReturn(new ProducerTableInfo(Map.of(
+ "pg-offline", List.of(producerInfo("offline-producer",
"10.0.0.1:1000")),
+ "pg-payment", List.of(producerInfo("producer-payment",
"10.0.0.2:1000")))));
+ when(adminExt.examineProducerConnectionInfo("pg-offline", "TopicA"))
+ .thenThrow(new MQClientException("Not found the producer group
connection", null));
+ ProducerConnection paymentConnection = new ProducerConnection();
+ paymentConnection.setConnectionSet(new HashSet<>(List.of(
+ connection("producer-payment", "10.0.0.2:1000"))));
+ when(adminExt.examineProducerConnectionInfo("pg-payment", "TopicA"))
+ .thenReturn(paymentConnection);
+
+ ProducerConnectionScanResult result =
+ provider.scanProducerConnections("instance-a", "TopicA", null);
+
+ assertThat(result.connections()).singleElement().satisfies(connection
-> {
assertThat(connection.getClientId()).isEqualTo("producer-payment");
assertThat(connection.getProducerGroup()).isEqualTo("pg-payment");
});
+ assertThat(result.complete()).isTrue();
+ assertThat(result.failedBrokers()).isEmpty();
+ assertThat(result.failedProducerGroups()).isEmpty();
}
@Test
diff --git a/web/src/api/producer.test.ts b/web/src/api/producer.test.ts
index 9878c6744..cd79af4bd 100644
--- a/web/src/api/producer.test.ts
+++ b/web/src/api/producer.test.ts
@@ -149,6 +149,23 @@ describe('Producer API', () => {
expect(result.summary.totalConnections).toBe(1);
});
+ it('preserves partial producer scan metadata', async () => {
+ mock.onGet('/producer/connection').reply(200, {
+ connectionSet: [],
+ complete: false,
+ failedBrokers: ['broker-a:10911'],
+ failedProducerGroups: ['pg-orders'],
+ });
+
+ const result = await queryProducerConnection('instance-1', 'order-events');
+
+ expect(result.complete).toBe(false);
+ expect(result.failedBrokers).toEqual(['broker-a:10911']);
+ expect(result.failedProducerGroups).toEqual(['pg-orders']);
+ expect(result.summary.readiness).toBe('WARNING');
+ expect(result.summary.warnings).toEqual(['NO_CONNECTIONS',
'INCOMPLETE_SCAN']);
+ });
+
it('handles empty producer connections', async () => {
mock.onGet('/producer/connection').reply(200, { connectionSet: [] });
diff --git a/web/src/api/producer.ts b/web/src/api/producer.ts
index c62ac1f12..0c1215f06 100644
--- a/web/src/api/producer.ts
+++ b/web/src/api/producer.ts
@@ -30,7 +30,11 @@ export interface ProducerConnection {
export type ProducerReadiness = 'READY' | 'WARNING' | 'UNAVAILABLE';
export type ProducerConnectionWarning =
- 'NO_CONNECTIONS' | 'DUPLICATE_CLIENT_ID' | 'MIXED_CLIENT_VERSION' |
'INCOMPLETE_CLIENT_METADATA';
+ | 'NO_CONNECTIONS'
+ | 'DUPLICATE_CLIENT_ID'
+ | 'MIXED_CLIENT_VERSION'
+ | 'INCOMPLETE_CLIENT_METADATA'
+ | 'INCOMPLETE_SCAN';
export interface ProducerConnectionSummaryItem {
value: string;
@@ -53,6 +57,9 @@ export interface ProducerConnectionSummary {
export interface ProducerConnectionResult {
connectionSet: ProducerConnection[];
summary: ProducerConnectionSummary;
+ complete: boolean;
+ failedBrokers: string[];
+ failedProducerGroups: string[];
}
interface TopicRecord {
@@ -67,6 +74,9 @@ interface TopicListResponse {
interface ProducerConnectionResponse {
connectionSet?: ProducerConnection[];
summary?: ProducerConnectionSummary;
+ complete?: boolean;
+ failedBrokers?: string[];
+ failedProducerGroups?: string[];
}
// ─── API ────────────────────────────────────────────────────────
@@ -102,6 +112,7 @@ const distribution = (
export function buildProducerConnectionSummary(
connections: ProducerConnection[],
+ complete = true,
): ProducerConnectionSummary {
const duplicateClientIds = [
...connections.reduce((counts, connection) => {
@@ -134,6 +145,7 @@ export function buildProducerConnectionSummary(
warnings.push('INCOMPLETE_CLIENT_METADATA');
}
}
+ if (!complete) warnings.push('INCOMPLETE_SCAN');
return {
totalConnections: connections.length,
@@ -145,7 +157,13 @@ export function buildProducerConnectionSummary(
versions,
duplicateClientIds,
warnings,
- readiness: connections.length === 0 ? 'UNAVAILABLE' : warnings.length > 0
? 'WARNING' : 'READY',
+ readiness: !complete
+ ? 'WARNING'
+ : connections.length === 0
+ ? 'UNAVAILABLE'
+ : warnings.length > 0
+ ? 'WARNING'
+ : 'READY',
};
}
@@ -186,8 +204,22 @@ export async function queryProducerConnection(
params: { instanceId, topic, producerGroup },
});
const connectionSet = res.data?.connectionSet ?? [];
+ const complete = res.data?.complete ?? true;
+ const backendSummary = res.data?.summary;
+ const summary = backendSummary ??
buildProducerConnectionSummary(connectionSet, complete);
+ const normalizedSummary =
+ !complete && !summary.warnings.includes('INCOMPLETE_SCAN')
+ ? {
+ ...summary,
+ warnings: [...summary.warnings, 'INCOMPLETE_SCAN' as const],
+ readiness: 'WARNING' as const,
+ }
+ : summary;
return {
connectionSet,
- summary: res.data?.summary ??
buildProducerConnectionSummary(connectionSet),
+ summary: normalizedSummary,
+ complete,
+ failedBrokers: res.data?.failedBrokers ?? [],
+ failedProducerGroups: res.data?.failedProducerGroups ?? [],
};
}
diff --git a/web/src/i18n/translations.ts b/web/src/i18n/translations.ts
index aeb4aff90..3e05f0687 100644
--- a/web/src/i18n/translations.ts
+++ b/web/src/i18n/translations.ts
@@ -2873,6 +2873,18 @@ const translations: Record<string, Record<Lang, string>>
= {
zh: '连接元数据不完整',
en: 'Incomplete connection metadata',
},
+ 'producer.warningIncompleteScan': {
+ zh: '扫描结果不完整',
+ en: 'Incomplete scan',
+ },
+ 'producer.failedBroker': {
+ zh: 'Broker 失败:{name}',
+ en: 'Broker failed: {name}',
+ },
+ 'producer.failedGroup': {
+ zh: '生产者组失败:{name}',
+ en: 'Producer group failed: {name}',
+ },
// ─── Namespace ───
'ns.title': { zh: '命名空间管理', en: 'Namespace Management' },
diff --git a/web/src/pages/studio/Producer.tsx
b/web/src/pages/studio/Producer.tsx
index 0fb20a820..dc85973ce 100644
--- a/web/src/pages/studio/Producer.tsx
+++ b/web/src/pages/studio/Producer.tsx
@@ -80,6 +80,8 @@ const ProducerPage = () => {
const [connectionSummary, setConnectionSummary] =
useState<ProducerConnectionSummary | null>(
null,
);
+ const [failedBrokers, setFailedBrokers] = useState<string[]>([]);
+ const [failedProducerGroups, setFailedProducerGroups] =
useState<string[]>([]);
const [instances, setInstances] = useState<Instance[]>([]);
const [selectedInstanceId, setSelectedInstanceId] = useState<string |
undefined>(undefined);
const [loading, setLoading] = useState(false);
@@ -126,6 +128,8 @@ const ProducerPage = () => {
queryInFlightRef.current = null;
setConnectionList([]);
setConnectionSummary(null);
+ setFailedBrokers([]);
+ setFailedProducerGroups([]);
setLoading(false);
};
@@ -205,6 +209,8 @@ const ProducerPage = () => {
queryInFlightRef.current = requestId;
setConnectionList([]);
setConnectionSummary(null);
+ setFailedBrokers([]);
+ setFailedProducerGroups([]);
setLoading(true);
try {
const result = await queryProducerConnection(
@@ -216,7 +222,9 @@ const ProducerPage = () => {
const connections = result.connectionSet;
setConnectionList(connections);
setConnectionSummary(result.summary);
- if (connections.length === 0) {
+ setFailedBrokers(result.failedBrokers);
+ setFailedProducerGroups(result.failedProducerGroups);
+ if (connections.length === 0 && result.complete) {
message.info(t('producer.noConnections'));
}
} catch {
@@ -267,6 +275,7 @@ const ProducerPage = () => {
DUPLICATE_CLIENT_ID: t('producer.warningDuplicateClientId'),
MIXED_CLIENT_VERSION: t('producer.warningMixedVersion'),
INCOMPLETE_CLIENT_METADATA: t('producer.warningIncompleteMetadata'),
+ INCOMPLETE_SCAN: t('producer.warningIncompleteScan'),
};
const renderDistribution = (items: ProducerConnectionSummary['languages']) =>
@@ -400,6 +409,16 @@ const ProducerPage = () => {
{warningLabel[warning] ?? warning}
</Tag>
))}
+ {failedBrokers.map((broker) => (
+ <Tag key={`broker:${broker}`} color="error">
+ {t('producer.failedBroker', { name: broker })}
+ </Tag>
+ ))}
+ {failedProducerGroups.map((group) => (
+ <Tag key={`group:${group}`} color="error">
+ {t('producer.failedGroup', { name: group })}
+ </Tag>
+ ))}
</Flex>
}
style={{ marginBottom: 12 }}
diff --git a/web/src/pages/studio/__tests__/Producer.test.tsx
b/web/src/pages/studio/__tests__/Producer.test.tsx
index b15e05305..b54ddb2cf 100644
--- a/web/src/pages/studio/__tests__/Producer.test.tsx
+++ b/web/src/pages/studio/__tests__/Producer.test.tsx
@@ -66,6 +66,9 @@ const renderWithProviders = (ui: React.ReactElement) => {
const producerResult = (connectionSet: ProducerConnection[]):
ProducerConnectionResult => ({
connectionSet,
+ complete: true,
+ failedBrokers: [],
+ failedProducerGroups: [],
summary: {
totalConnections: connectionSet.length,
uniqueClientCount: new Set(connectionSet.map((connection) =>
connection.clientId)).size,
@@ -267,6 +270,9 @@ describe('ProducerPage', () => {
versionDesc: '5.2.0',
},
],
+ complete: true,
+ failedBrokers: [],
+ failedProducerGroups: [],
summary: {
totalConnections: 2,
uniqueClientCount: 1,
@@ -299,6 +305,42 @@ describe('ProducerPage', () => {
expect(screen.getByText('JAVA: 2')).toBeInTheDocument();
});
+ it('surfaces partial topic-wide producer scans with failed targets', async
() => {
+ vi.mocked(queryProducerConnection).mockResolvedValue({
+ connectionSet: [],
+ complete: false,
+ failedBrokers: ['broker-a:10911'],
+ failedProducerGroups: ['pg-orders'],
+ summary: {
+ totalConnections: 0,
+ uniqueClientCount: 0,
+ uniqueAddressCount: 0,
+ uniqueLanguageCount: 0,
+ uniqueVersionCount: 0,
+ languages: [],
+ versions: [],
+ duplicateClientIds: [],
+ warnings: ['NO_CONNECTIONS', 'INCOMPLETE_SCAN'],
+ readiness: 'WARNING',
+ },
+ });
+ const user = userEvent.setup();
+ renderWithProviders(<ProducerPage />);
+
+ await waitFor(() => expect(fetchTopicList).toHaveBeenCalledTimes(1));
+ const [, topicSelect] = screen.getAllByRole('combobox');
+ fireEvent.mouseDown(topicSelect.parentElement!);
+ await user.click(
+ await screen.findByText('order-events', { selector:
'.ant-select-item-option-content' }),
+ );
+ await user.click(screen.getByRole('button', { name: /搜索/ }));
+
+ expect(await screen.findByText('扫描结果不完整')).toBeInTheDocument();
+ expect(screen.getByText('Broker 失败:broker-a:10911')).toBeInTheDocument();
+ expect(screen.getByText('生产者组失败:pg-orders')).toBeInTheDocument();
+ expect(screen.queryByText('暂无生产者连接')).not.toBeInTheDocument();
+ });
+
it('exports the current producer connection diagnostics as CSV', async () =>
{
const createObjectURL = vi.fn((blob: Blob | MediaSource) => {
expect(blob).toBeInstanceOf(Blob);
@@ -326,6 +368,9 @@ describe('ProducerPage', () => {
versionDesc: '5.1.0',
},
],
+ complete: true,
+ failedBrokers: [],
+ failedProducerGroups: [],
summary: {
totalConnections: 1,
uniqueClientCount: 1,