This is an automated email from the ASF dual-hosted git repository. Caideyipi pushed a commit to branch fix/config-pipe-plan-authorization in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 9a798b16b402f1143727ba002b0eef3ddf805ee4 Author: Caideyipi <[email protected]> AuthorDate: Tue Aug 18 16:42:36 2026 +0800 Fix fail-open authorization for pipe config plans --- .../iotdb/pipe/it/single/IoTDBPipeReceiverIT.java | 19 +++++++++++++++++++ .../receiver/protocol/IoTDBConfigNodeReceiver.java | 15 +++++++++++++-- 2 files changed, 32 insertions(+), 2 deletions(-) diff --git a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipeReceiverIT.java b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipeReceiverIT.java index 5ca62ebb15f..136943f00b0 100644 --- a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipeReceiverIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipeReceiverIT.java @@ -24,12 +24,15 @@ import org.apache.iotdb.commons.client.property.ThriftClientProperty; import org.apache.iotdb.commons.conf.CommonDescriptor; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.path.MeasurementPath; +import org.apache.iotdb.commons.pipe.agent.plugin.meta.PipePluginMeta; import org.apache.iotdb.commons.pipe.sink.client.IoTDBSyncClient; import org.apache.iotdb.commons.pipe.sink.payload.thrift.common.PipeTransferHandshakeConstant; import org.apache.iotdb.commons.pipe.sink.payload.thrift.request.IoTDBSinkRequestVersion; import org.apache.iotdb.commons.pipe.sink.payload.thrift.request.PipeRequestType; +import org.apache.iotdb.confignode.consensus.request.write.pipe.plugin.CreatePipePluginPlan; import org.apache.iotdb.confignode.manager.pipe.sink.payload.PipeTransferConfigNodeHandshakeV1Req; import org.apache.iotdb.confignode.manager.pipe.sink.payload.PipeTransferConfigNodeHandshakeV2Req; +import org.apache.iotdb.confignode.manager.pipe.sink.payload.PipeTransferConfigPlanReq; import org.apache.iotdb.db.pipe.sink.payload.evolvable.request.PipeTransferDataNodeHandshakeV1Req; import org.apache.iotdb.db.pipe.sink.payload.evolvable.request.PipeTransferDataNodeHandshakeV2Req; import org.apache.iotdb.db.pipe.sink.payload.evolvable.request.PipeTransferTabletRawReq; @@ -56,6 +59,7 @@ import org.apache.iotdb.service.rpc.thrift.TSyncTransportMetaInfo; import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.external.commons.io.FileUtils; import org.apache.tsfile.file.metadata.enums.TSEncoding; +import org.apache.tsfile.utils.Binary; import org.apache.tsfile.utils.ReadWriteIOUtils; import org.apache.tsfile.write.record.Tablet; import org.apache.tsfile.write.schema.MeasurementSchema; @@ -256,6 +260,21 @@ public class IoTDBPipeReceiverIT { Assert.assertNotEquals( TSStatusCode.NOT_LOGIN.getStatusCode(), client.pipeTransfer(buildEmptyConfigPlanReq()).getStatus().getCode()); + Assert.assertEquals( + TSStatusCode.NO_PERMISSION.getStatusCode(), + client + .pipeTransfer( + PipeTransferConfigPlanReq.toTPipeTransferReq( + new CreatePipePluginPlan( + new PipePluginMeta( + "receiver-security-test-plugin", + "attacker.Plugin", + false, + "attacker.jar", + "deadbeef"), + new Binary(new byte[] {1})))) + .getStatus() + .getCode()); } } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/receiver/protocol/IoTDBConfigNodeReceiver.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/receiver/protocol/IoTDBConfigNodeReceiver.java index 281867d0ec5..3afeaebddae 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/receiver/protocol/IoTDBConfigNodeReceiver.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/receiver/protocol/IoTDBConfigNodeReceiver.java @@ -88,6 +88,7 @@ import org.apache.iotdb.confignode.consensus.request.write.table.view.SetViewCom import org.apache.iotdb.confignode.consensus.request.write.table.view.SetViewPropertiesPlan; import org.apache.iotdb.confignode.consensus.request.write.template.CommitSetSchemaTemplatePlan; import org.apache.iotdb.confignode.consensus.request.write.template.CreateSchemaTemplatePlan; +import org.apache.iotdb.confignode.consensus.request.write.template.DropSchemaTemplatePlan; import org.apache.iotdb.confignode.consensus.request.write.template.ExtendSchemaTemplatePlan; import org.apache.iotdb.confignode.consensus.request.write.trigger.DeleteTriggerInTablePlan; import org.apache.iotdb.confignode.consensus.request.write.trigger.UpdateTriggerStateInTablePlan; @@ -413,6 +414,12 @@ public class IoTDBConfigNodeReceiver extends IoTDBFileReceiver { case CreateSchemaTemplate: templateName = ((CreateSchemaTemplatePlan) plan).getTemplate().getName(); return checkGlobalStatus(userEntity, PrivilegeType.SYSTEM, templateName, true); + case DropSchemaTemplate: + return checkGlobalStatus( + userEntity, + PrivilegeType.SYSTEM, + ((DropSchemaTemplatePlan) plan).getTemplateName(), + true); case CommitSetSchemaTemplate: templateName = ((CommitSetSchemaTemplatePlan) plan).getName(); return checkGlobalStatus(userEntity, PrivilegeType.SYSTEM, templateName, true); @@ -747,7 +754,7 @@ public class IoTDBConfigNodeReceiver extends IoTDBFileReceiver { return checkGlobalStatus( userEntity, PrivilegeType.MANAGE_ROLE, ((AuthorPlan) plan).getRoleName(), true); default: - return StatusUtils.OK; + return RpcUtils.getStatus(TSStatusCode.NO_PERMISSION); } } @@ -1204,10 +1211,14 @@ public class IoTDBConfigNodeReceiver extends IoTDBFileReceiver { .getPermissionManager() .operatePermission((AuthorPlan) plan, shouldMarkAsPipeRequest.get()); case CreateSchemaTemplate: - default: + case DropSchemaTemplate: + // Only explicitly supported config-region pipe plans may be written to consensus. New plan + // types must be added to an explicit case after their authorization is implemented. return configManager .getConsensusManager() .write(shouldMarkAsPipeRequest.get() ? new PipeEnrichedPlan(plan) : plan); + default: + return RpcUtils.getStatus(TSStatusCode.NO_PERMISSION); } }
