This is an automated email from the ASF dual-hosted git repository.
Caideyipi pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/dev/1.3 by this push:
new 2e4d7d5299a Fix subscription topic authorization bypass (#18418)
(#18435)
2e4d7d5299a is described below
commit 2e4d7d5299aaaf74d042edabe3540181c9ca4eec
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 18 12:19:06 2026 +0800
Fix subscription topic authorization bypass (#18418) (#18435)
(cherry picked from commit bc5b6a2d162fa1f38c0553fe177430629bf55262)
---
.../protocol/thrift/impl/ClientRPCServiceImpl.java | 2 +-
.../agent/SubscriptionReceiverAgent.java | 11 ++++
.../subscription/agent/SubscriptionTopicAgent.java | 55 +++++++++++++++++++
.../receiver/SubscriptionReceiver.java | 2 +
.../receiver/SubscriptionReceiverV1.java | 64 ++++++++++++++++++++--
5 files changed, 127 insertions(+), 7 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java
index c4bfd21a7a4..3987adcaaaf 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java
@@ -2828,7 +2828,7 @@ public class ClientRPCServiceImpl implements
IClientRPCServiceWithHandler {
return getNotLoggedInPipeSubscribeResp();
}
- return SubscriptionAgent.receiver().handle(req);
+ return SubscriptionAgent.receiver().handle(req,
clientSession.getUsername());
} finally {
SESSION_MANAGER.updateIdleTime();
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
index 200728fe203..e2ae494aa12 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
@@ -77,6 +77,16 @@ public class SubscriptionReceiverAgent {
}
public TPipeSubscribeResp handle(final TPipeSubscribeReq req) {
+ return handle(req, null);
+ }
+
+ public TPipeSubscribeResp handle(final TPipeSubscribeReq req, final String
username) {
+ if (username == null) {
+ return new TPipeSubscribeResp(
+ RpcUtils.getStatus(TSStatusCode.NO_PERMISSION),
+ PipeSubscribeResponseVersion.VERSION_1.getVersion(),
+ PipeSubscribeResponseType.ACK.getType());
+ }
if (!SubscriptionConfig.getInstance().getSubscriptionEnabled()) {
return SUBSCRIPTION_NOT_ENABLED_ERROR_RESP;
}
@@ -84,6 +94,7 @@ public class SubscriptionReceiverAgent {
final byte reqVersion = req.getVersion();
if (RECEIVER_CONSTRUCTORS.containsKey(reqVersion)) {
final SubscriptionReceiver receiver = getReceiver(reqVersion);
+ receiver.setAuthenticatedUsername(username);
activeReceivers.add(receiver);
receiver.handleTimeout();
return receiver.handle(req);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionTopicAgent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionTopicAgent.java
index 37cdaa72690..4c178732bc4 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionTopicAgent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionTopicAgent.java
@@ -19,9 +19,18 @@
package org.apache.iotdb.db.subscription.agent;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.auth.entity.PrivilegeType;
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.pipe.datastructure.pattern.IoTDBPipePattern;
+import org.apache.iotdb.commons.pipe.datastructure.pattern.PipePattern;
+import org.apache.iotdb.commons.pipe.datastructure.pattern.PrefixPipePattern;
import org.apache.iotdb.commons.subscription.meta.topic.TopicMeta;
import org.apache.iotdb.commons.subscription.meta.topic.TopicMetaKeeper;
+import org.apache.iotdb.db.auth.AuthorityChecker;
import org.apache.iotdb.mpp.rpc.thrift.TPushTopicMetaRespExceptionMessage;
+import org.apache.iotdb.rpc.RpcUtils;
+import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.rpc.subscription.config.TopicConfig;
import org.apache.iotdb.rpc.subscription.config.TopicConstant;
@@ -30,6 +39,7 @@ import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Set;
import java.util.stream.Collectors;
@@ -187,4 +197,49 @@ public class SubscriptionTopicAgent {
releaseReadLock();
}
}
+
+ /**
+ * Check that the authenticated session can read all data covered by the
requested topics. The
+ * username in ConsumerConfig is client-controlled and therefore must not be
used as the
+ * authorization identity.
+ */
+ public TSStatus checkTopicReadPermissions(
+ final String username, final Iterable<String> topicNames) {
+ if (Objects.isNull(username)) {
+ return RpcUtils.getStatus(TSStatusCode.NO_PERMISSION);
+ }
+
+ acquireReadLock();
+ try {
+ for (final String topicName : topicNames) {
+ final TopicMeta topicMeta = topicMetaKeeper.getTopicMeta(topicName);
+ if (Objects.isNull(topicMeta)) {
+ continue;
+ }
+
+ final TSStatus status = checkTopicReadPermission(username, topicMeta);
+ if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return status;
+ }
+ }
+ return RpcUtils.SUCCESS_STATUS;
+ } finally {
+ releaseReadLock();
+ }
+ }
+
+ private TSStatus checkTopicReadPermission(final String username, final
TopicMeta topicMeta) {
+ final TopicConfig topicConfig = topicMeta.getConfig();
+ final PipePattern pipePattern =
+ topicConfig.getAttribute().containsKey(TopicConstant.PATTERN_KEY)
+ ? new
PrefixPipePattern(topicConfig.getAttribute().get(TopicConstant.PATTERN_KEY))
+ : new IoTDBPipePattern(
+ topicConfig.getStringOrDefault(
+ TopicConstant.PATH_KEY, TopicConstant.PATH_DEFAULT_VALUE));
+ final List<PartialPath> paths = pipePattern.getBaseInclusionPaths();
+ return AuthorityChecker.getTSStatus(
+ AuthorityChecker.checkPatternPermission(username, paths,
PrivilegeType.READ_DATA.ordinal()),
+ paths,
+ PrivilegeType.READ_DATA);
+ }
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
index 36e3c9b74f5..cc7b57eee81 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
@@ -27,6 +27,8 @@ public interface SubscriptionReceiver {
TPipeSubscribeResp handle(TPipeSubscribeReq req);
+ void setAuthenticatedUsername(final String username);
+
PipeSubscribeRequestVersion getVersion();
void handleExit();
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
index 8cb09983689..e15cc333f4f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
@@ -113,6 +113,7 @@ public class SubscriptionReceiverV1 implements
SubscriptionReceiver {
private final ThreadLocal<ConsumerConfig> consumerConfigThreadLocal = new
ThreadLocal<>();
private final ThreadLocal<PollTimer> pollTimerThreadLocal = new
ThreadLocal<>();
+ private volatile String authenticatedUsername;
private volatile ConsumerConfig sharedConsumerConfig;
private volatile boolean consumerInvalidated;
private volatile long lastActivityTimeMs = System.currentTimeMillis();
@@ -167,6 +168,11 @@ public class SubscriptionReceiverV1 implements
SubscriptionReceiver {
return PipeSubscribeRequestVersion.VERSION_1;
}
+ @Override
+ public void setAuthenticatedUsername(final String username) {
+ authenticatedUsername = username;
+ }
+
@Override
public void handleExit() {
final ConsumerConfig consumerConfig = consumerConfigThreadLocal.get();
@@ -183,6 +189,7 @@ public class SubscriptionReceiverV1 implements
SubscriptionReceiver {
consumerConfigThreadLocal.remove();
}
clearSharedConsumerState();
+ authenticatedUsername = null;
}
@Override
@@ -322,17 +329,22 @@ public class SubscriptionReceiverV1 implements
SubscriptionReceiver {
return SUBSCRIPTION_MISSING_CUSTOMER_RESP;
}
- // TODO: do something
+ final Set<String> subscribedTopicNames =
+ SubscriptionAgent.consumer()
+ .getTopicNamesSubscribedByConsumer(
+ consumerConfig.getConsumerGroupId(),
consumerConfig.getConsumerId());
+ final TSStatus readPermissionStatus =
+ SubscriptionAgent.topic()
+ .checkTopicReadPermissions(authenticatedUsername,
subscribedTopicNames);
+ if (readPermissionStatus.getCode() !=
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return
PipeSubscribeHeartbeatResp.toTPipeSubscribeResp(readPermissionStatus);
+ }
LOGGER.info("Subscription: consumer {} heartbeat successfully",
consumerConfig);
// fetch subscribed topics
final Map<String, TopicConfig> topics =
- SubscriptionAgent.topic()
- .getTopicConfigs(
- SubscriptionAgent.consumer()
- .getTopicNamesSubscribedByConsumer(
- consumerConfig.getConsumerGroupId(),
consumerConfig.getConsumerId()));
+ SubscriptionAgent.topic().getTopicConfigs(subscribedTopicNames);
// fetch available endpoints
final Map<Integer, TEndPoint> endPoints = new HashMap<>();
@@ -403,6 +415,11 @@ public class SubscriptionReceiverV1 implements
SubscriptionReceiver {
// subscribe topics
final Set<String> topicNames = req.getTopicNames();
+ final TSStatus readPermissionStatus =
+
SubscriptionAgent.topic().checkTopicReadPermissions(authenticatedUsername,
topicNames);
+ if (readPermissionStatus.getCode() !=
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return
PipeSubscribeSubscribeResp.toTPipeSubscribeResp(readPermissionStatus);
+ }
subscribe(consumerConfig, topicNames);
LOGGER.info("Subscription: consumer {} subscribe {} successfully",
consumerConfig, topicNames);
@@ -494,16 +511,51 @@ public class SubscriptionReceiverV1 implements
SubscriptionReceiver {
if (SubscriptionPollRequestType.isValidatedRequestType(requestType)) {
switch (SubscriptionPollRequestType.valueOf(requestType)) {
case POLL:
+ final Set<String> pollTopicNames = ((PollPayload)
request.getPayload()).getTopicNames();
+ final Set<String> subscribedTopicNames =
+ SubscriptionAgent.consumer()
+ .getTopicNamesSubscribedByConsumer(
+ consumerConfig.getConsumerGroupId(),
consumerConfig.getConsumerId());
+ final Set<String> topicNamesToCheck = new HashSet<>(pollTopicNames);
+ topicNamesToCheck.removeIf(topicName ->
!subscribedTopicNames.contains(topicName));
+ final TSStatus readPermissionStatus =
+ SubscriptionAgent.topic()
+ .checkTopicReadPermissions(authenticatedUsername,
topicNamesToCheck);
+ if (readPermissionStatus.getCode() !=
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return PipeSubscribePollResp.toTPipeSubscribeResp(
+ readPermissionStatus, Collections.emptyList());
+ }
events =
handlePipeSubscribePollRequest(
consumerConfig, (PollPayload) request.getPayload(),
maxBytes);
break;
case POLL_FILE:
+ final String tsFileTopicName =
+ ((PollFilePayload)
request.getPayload()).getCommitContext().getTopicName();
+ final TSStatus tsFileReadPermissionStatus =
+ SubscriptionAgent.topic()
+ .checkTopicReadPermissions(
+ authenticatedUsername,
Collections.singleton(tsFileTopicName));
+ if (tsFileReadPermissionStatus.getCode() !=
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return PipeSubscribePollResp.toTPipeSubscribeResp(
+ tsFileReadPermissionStatus, Collections.emptyList());
+ }
events =
handlePipeSubscribePollTsFileRequest(
consumerConfig, (PollFilePayload) request.getPayload());
break;
case POLL_TABLETS:
+ final String tabletsTopicName =
+ ((PollTabletsPayload)
request.getPayload()).getCommitContext().getTopicName();
+ final TSStatus tabletsReadPermissionStatus =
+ SubscriptionAgent.topic()
+ .checkTopicReadPermissions(
+ authenticatedUsername,
Collections.singleton(tabletsTopicName));
+ if (tabletsReadPermissionStatus.getCode()
+ != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return PipeSubscribePollResp.toTPipeSubscribeResp(
+ tabletsReadPermissionStatus, Collections.emptyList());
+ }
events =
handlePipeSubscribePollTabletsRequest(
consumerConfig, (PollTabletsPayload) request.getPayload());