This is an automated email from the ASF dual-hosted git repository. ycycse pushed a commit to branch FillAINode in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 98f87013316a0ae1b685f6978d5c0f65c50522a2 Author: ycycse <[email protected]> AuthorDate: Tue Sep 24 00:58:25 2024 +0800 Add logic of plan deserialization in configNode --- .../consensus/request/ConfigPhysicalPlan.java | 40 ++++++++++++++++++++++ .../read/ainode/GetAINodeConfigurationPlan.java | 21 +++++++++++- .../request/read/model/GetModelInfoPlan.java | 22 +++++++++++- .../request/read/model/ShowModelPlan.java | 24 +++++++++++++ 4 files changed, 105 insertions(+), 2 deletions(-) diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java index 7bfa5c50c69..9d020aa52d7 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java @@ -20,7 +20,13 @@ package org.apache.iotdb.confignode.consensus.request; import org.apache.iotdb.commons.exception.runtime.SerializationRunTimeException; +import org.apache.iotdb.confignode.consensus.request.read.ainode.GetAINodeConfigurationPlan; +import org.apache.iotdb.confignode.consensus.request.read.model.GetModelInfoPlan; +import org.apache.iotdb.confignode.consensus.request.read.model.ShowModelPlan; import org.apache.iotdb.confignode.consensus.request.read.subscription.ShowTopicPlan; +import org.apache.iotdb.confignode.consensus.request.write.ainode.RegisterAINodePlan; +import org.apache.iotdb.confignode.consensus.request.write.ainode.RemoveAINodePlan; +import org.apache.iotdb.confignode.consensus.request.write.ainode.UpdateAINodePlan; import org.apache.iotdb.confignode.consensus.request.write.auth.AuthorPlan; import org.apache.iotdb.confignode.consensus.request.write.confignode.ApplyConfigNodePlan; import org.apache.iotdb.confignode.consensus.request.write.confignode.RemoveConfigNodePlan; @@ -43,6 +49,10 @@ import org.apache.iotdb.confignode.consensus.request.write.datanode.RemoveDataNo import org.apache.iotdb.confignode.consensus.request.write.datanode.UpdateDataNodePlan; import org.apache.iotdb.confignode.consensus.request.write.function.CreateFunctionPlan; import org.apache.iotdb.confignode.consensus.request.write.function.DropFunctionPlan; +import org.apache.iotdb.confignode.consensus.request.write.model.CreateModelPlan; +import org.apache.iotdb.confignode.consensus.request.write.model.DropModelInNodePlan; +import org.apache.iotdb.confignode.consensus.request.write.model.DropModelPlan; +import org.apache.iotdb.confignode.consensus.request.write.model.UpdateModelInfoPlan; import org.apache.iotdb.confignode.consensus.request.write.partition.AddRegionLocationPlan; import org.apache.iotdb.confignode.consensus.request.write.partition.CreateDataPartitionPlan; import org.apache.iotdb.confignode.consensus.request.write.partition.CreateSchemaPartitionPlan; @@ -167,6 +177,18 @@ public abstract class ConfigPhysicalPlan implements IConsensusRequest { case RemoveDataNode: plan = new RemoveDataNodePlan(); break; + case RegisterAINode: + plan = new RegisterAINodePlan(); + break; + case RemoveAINode: + plan = new RemoveAINodePlan(); + break; + case GetAINodeConfiguration: + plan = new GetAINodeConfigurationPlan(); + break; + case UpdateAINodeConfiguration: + plan = new UpdateAINodePlan(); + break; case CreateDatabase: plan = new DatabaseSchemaPlan(ConfigPhysicalPlanType.CreateDatabase); break; @@ -420,6 +442,24 @@ public abstract class ConfigPhysicalPlan implements IConsensusRequest { case UPDATE_CQ_LAST_EXEC_TIME: plan = new UpdateCQLastExecTimePlan(); break; + case CreateModel: + plan = new CreateModelPlan(); + break; + case UpdateModelInfo: + plan = new UpdateModelInfoPlan(); + break; + case DropModel: + plan = new DropModelPlan(); + break; + case ShowModel: + plan = new ShowModelPlan(); + break; + case DropModelInNode: + plan = new DropModelInNodePlan(); + break; + case GetModelInfo: + plan = new GetModelInfoPlan(); + break; case CreatePipePlugin: plan = new CreatePipePluginPlan(); break; diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/ainode/GetAINodeConfigurationPlan.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/ainode/GetAINodeConfigurationPlan.java index b080cf77c5d..7222a8f53f8 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/ainode/GetAINodeConfigurationPlan.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/ainode/GetAINodeConfigurationPlan.java @@ -22,10 +22,18 @@ package org.apache.iotdb.confignode.consensus.request.read.ainode; import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlanType; import org.apache.iotdb.confignode.consensus.request.read.ConfigPhysicalReadPlan; +import java.io.DataOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; + public class GetAINodeConfigurationPlan extends ConfigPhysicalReadPlan { // if aiNodeId is set to -1, return all AINode configurations. - private final int aiNodeId; + private int aiNodeId; + + public GetAINodeConfigurationPlan() { + super(ConfigPhysicalPlanType.GetAINodeConfiguration); + } public GetAINodeConfigurationPlan(final int aiNodeId) { super(ConfigPhysicalPlanType.GetAINodeConfiguration); @@ -36,6 +44,17 @@ public class GetAINodeConfigurationPlan extends ConfigPhysicalReadPlan { return aiNodeId; } + @Override + protected void serializeImpl(DataOutputStream stream) throws IOException { + stream.writeShort(getType().getPlanType()); + stream.writeInt(aiNodeId); + } + + @Override + protected void deserializeImpl(ByteBuffer buffer) throws IOException { + this.aiNodeId = buffer.getInt(); + } + @Override public boolean equals(final Object o) { if (this == o) { diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/model/GetModelInfoPlan.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/model/GetModelInfoPlan.java index dfec065f206..9c33c267882 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/model/GetModelInfoPlan.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/model/GetModelInfoPlan.java @@ -23,11 +23,20 @@ import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlanType; import org.apache.iotdb.confignode.consensus.request.read.ConfigPhysicalReadPlan; import org.apache.iotdb.confignode.rpc.thrift.TGetModelInfoReq; +import org.apache.tsfile.utils.ReadWriteIOUtils; + +import java.io.DataOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; import java.util.Objects; public class GetModelInfoPlan extends ConfigPhysicalReadPlan { - private final String modelId; + private String modelId; + + public GetModelInfoPlan() { + super(ConfigPhysicalPlanType.GetModelInfo); + } public GetModelInfoPlan(final TGetModelInfoReq getModelInfoReq) { super(ConfigPhysicalPlanType.GetModelInfo); @@ -38,6 +47,17 @@ public class GetModelInfoPlan extends ConfigPhysicalReadPlan { return modelId; } + @Override + protected void serializeImpl(DataOutputStream stream) throws IOException { + stream.writeShort(getType().getPlanType()); + ReadWriteIOUtils.write(modelId, stream); + } + + @Override + protected void deserializeImpl(ByteBuffer buffer) throws IOException { + this.modelId = ReadWriteIOUtils.readString(buffer); + } + @Override public boolean equals(final Object o) { if (this == o) { diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/model/ShowModelPlan.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/model/ShowModelPlan.java index c3d0ff79a37..df924c97f5b 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/model/ShowModelPlan.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/model/ShowModelPlan.java @@ -23,12 +23,21 @@ import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlanType; import org.apache.iotdb.confignode.consensus.request.read.ConfigPhysicalReadPlan; import org.apache.iotdb.confignode.rpc.thrift.TShowModelReq; +import org.apache.tsfile.utils.ReadWriteIOUtils; + +import java.io.DataOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; import java.util.Objects; public class ShowModelPlan extends ConfigPhysicalReadPlan { private String modelName; + public ShowModelPlan() { + super(ConfigPhysicalPlanType.ShowModel); + } + public ShowModelPlan(final TShowModelReq showModelReq) { super(ConfigPhysicalPlanType.ShowModel); if (showModelReq.isSetModelId()) { @@ -44,6 +53,21 @@ public class ShowModelPlan extends ConfigPhysicalReadPlan { return modelName; } + @Override + protected void serializeImpl(DataOutputStream stream) throws IOException { + stream.writeShort(getType().getPlanType()); + ReadWriteIOUtils.write(modelName != null, stream); + ReadWriteIOUtils.write(modelName, stream); + } + + @Override + protected void deserializeImpl(ByteBuffer buffer) throws IOException { + boolean isSetModelId = ReadWriteIOUtils.readBool(buffer); + if (isSetModelId) { + this.modelName = ReadWriteIOUtils.readString(buffer); + } + } + @Override public boolean equals(final Object o) { if (this == o) {
