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 e54e1e2f6 refactor(server): bound topic-consumer scan, fail-closed
SSRF guard (#2398)
e54e1e2f6 is described below
commit e54e1e2f6ce70a32669992d0feb2ed791717f8bd
Author: zhaohai <[email protected]>
AuthorDate: Fri Aug 21 17:11:56 2026 +0800
refactor(server): bound topic-consumer scan, fail-closed SSRF guard (#2398)
- B1: getTopicConsumers no longer requests Integer.MAX_VALUE groups from the
broker (one examineConsumeStats + one examineConsumerConnectionInfo RPC
per
group = unbounded N+1). Cap the non-paginated scan at
TOPIC_CONSUMER_SCAN_LIMIT
(1000); topics with more consumer groups should use the paginated
/consumers/page endpoint instead.
- B2: SettingsService.isAllowedDataSourceHost now fails closed on
UnknownHostException (returns false) instead of true, matching
UrlHostGuard.check and removing a blind reachability oracle for
unresolvable
hosts. Loopback/private rejection policy is unchanged, so existing tests
(including the direct areAllowedDataSourceAddresses tests) keep passing.
- B12: the two catch (Exception ignored) blocks in getTopicConsumersPage now
log the skipped group at debug level instead of swallowing silently.
- B11: LlmConfigService.getConfig dropped a redundant double negation
(!!StringUtils.hasText) in favor of StringUtils.hasText.
Note: TencentClientFactory @PreDestroy (originally planned) was NOT added —
the Tencent SDK AbstractClient exposes no close/shutdown, so there is
nothing
to release at shutdown.
---
.../org/apache/rocketmq/studio/ops/ai/LlmConfigService.java | 2 +-
.../studio/provider/apache/RocketMQMetadataProvider.java | 10 +++++++++-
.../org/apache/rocketmq/studio/settings/SettingsService.java | 6 ++++--
3 files changed, 14 insertions(+), 4 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
index fe6107e47..d7218f599 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
@@ -97,7 +97,7 @@ public class LlmConfigService {
? copy(overrides)
: fromGeneralSettings(settingsService.getGeneralSettings());
String token = envToken();
- if (!!StringUtils.hasText(token)) {
+ if (StringUtils.hasText(token)) {
config.setApiKey(token.trim());
config.setEnabled(true);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
index 06f5cc28a..9ed5a8d01 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
@@ -100,6 +100,12 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
"RMQ_SYS_TRACE_TOPIC", "RMQ_SYS_TRANS_OP_HALF_TOPIC"
);
+ // Upper bound on groups scanned by the non-paginated getTopicConsumers.
Without this the
+ // call used Integer.MAX_VALUE, fanning out one examineConsumeStats + one
+ // examineConsumerConnectionInfo RPC per group (an unbounded N+1 against
the broker).
+ // Topics with more groups should use the paginated /consumers/page
endpoint instead.
+ private static final int TOPIC_CONSUMER_SCAN_LIMIT = 1000;
+
private final MqAdminExtFactory adminFactory;
private final RocketMQProperties properties;
private final RmqTopicMapper topicMapper;
@@ -415,7 +421,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
@Override
public List<TopicConsumerVO> getTopicConsumers(String instanceId, String
name) {
- return getTopicConsumersPage(instanceId, name, 1,
Integer.MAX_VALUE).getItems();
+ return getTopicConsumersPage(instanceId, name, 1,
TOPIC_CONSUMER_SCAN_LIMIT).getItems();
}
@Override
@@ -479,6 +485,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
}
} catch (Exception ignored) {
// group may be offline
+ log.debug("Skipping consumer connection info for group
{}: {}", group, ignored.getMessage());
}
consumers.add(TopicConsumerVO.builder()
@@ -490,6 +497,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
.build());
} catch (Exception ignored) {
// stats unavailable for this group, still list it below
without numbers
+ log.debug("Skipping consume stats for group {}: {}",
group, ignored.getMessage());
consumers.add(TopicConsumerVO.builder()
.group(group)
.consumeType(ConsumeType.CLUSTERING)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/settings/SettingsService.java
b/server/src/main/java/org/apache/rocketmq/studio/settings/SettingsService.java
index ecb125dcc..27c40e085 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/settings/SettingsService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/settings/SettingsService.java
@@ -370,8 +370,10 @@ public class SettingsService {
try {
return
areAllowedDataSourceAddresses(InetAddress.getAllByName(normalized));
} catch (UnknownHostException exception) {
- // Unresolvable host: let the connection attempt surface the real
connectivity error.
- return true;
+ // Fail closed: an unresolvable host must not be handed to the
connection
+ // layer (this used to return true, creating a blind reachability
oracle that
+ // differed from UrlHostGuard.check, which also fails closed).
+ return false;
}
}