This is an automated email from the ASF dual-hosted git repository.

mikexue pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git


The following commit(s) were added to refs/heads/master by this push:
     new bdca426a2 [ISSUE #3382]update eventmesh-runtime module for using 
storage api
     new d8e889c71 Merge pull request #3428 from mxsm/eventmesh-3382
bdca426a2 is described below

commit bdca426a2c1122f904d62934c66ec3047b2ecce1
Author: mxsm <[email protected]>
AuthorDate: Fri Mar 10 21:29:34 2023 +0800

    [ISSUE #3382]update eventmesh-runtime module for using storage api
---
 .../common/config/CommonConfiguration.java         |  3 ++
 .../common/config/CommonConfigurationTest.java     |  4 ++-
 .../src/test/resources/configuration.properties    |  1 +
 eventmesh-runtime/build.gradle                     |  2 +-
 eventmesh-runtime/conf/eventmesh.properties        |  3 ++
 .../admin/controller/ClientManageController.java   |  4 +--
 .../eventmesh/runtime/boot/EventMeshServer.java    | 10 +++----
 .../runtime/core/plugin/MQAdminWrapper.java        |  6 ++--
 .../runtime/core/plugin/MQConsumerWrapper.java     |  6 ++--
 .../runtime/core/plugin/MQProducerWrapper.java     |  6 ++--
 .../protocol/grpc/consumer/EventMeshConsumer.java  |  4 +--
 .../protocol/grpc/producer/EventMeshProducer.java  |  2 +-
 .../protocol/http/consumer/EventMeshConsumer.java  |  6 ++--
 .../protocol/http/producer/EventMeshProducer.java  |  2 +-
 .../tcp/client/group/ClientGroupWrapper.java       |  9 ++----
 .../StorageResource.java}                          | 35 ++++++++++------------
 .../controller/ClientManageControllerTest.java     |  2 +-
 .../runtime/boot/EventMeshServerTest.java          |  1 +
 .../EventMeshGrpcConfigurationTest.java            |  3 +-
 .../EventMeshHTTPConfigurationTest.java            |  3 +-
 .../EventMeshTCPConfigurationTest.java             |  3 +-
 .../src/test/resources/configuration.properties    |  1 +
 22 files changed, 60 insertions(+), 56 deletions(-)

diff --git 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/config/CommonConfiguration.java
 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/config/CommonConfiguration.java
index 97239115b..61e4d9ac6 100644
--- 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/config/CommonConfiguration.java
+++ 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/config/CommonConfiguration.java
@@ -72,6 +72,9 @@ public class CommonConfiguration {
     @ConfigFiled(field = "connector.plugin.type", notEmpty = true)
     private String eventMeshConnectorPluginType = "rocketmq";
 
+    @ConfigFiled(field = "storage.plugin.type", notEmpty = true)
+    private String eventMeshStoragePluginType = "rocketmq";
+
     @ConfigFiled(field = "security.validation.type.token", notEmpty = true)
     private boolean eventMeshSecurityValidateTypeToken = false;
 
diff --git 
a/eventmesh-common/src/test/java/org/apache/eventmesh/common/config/CommonConfigurationTest.java
 
b/eventmesh-common/src/test/java/org/apache/eventmesh/common/config/CommonConfigurationTest.java
index 755ac1d5e..acc9327bb 100644
--- 
a/eventmesh-common/src/test/java/org/apache/eventmesh/common/config/CommonConfigurationTest.java
+++ 
b/eventmesh-common/src/test/java/org/apache/eventmesh/common/config/CommonConfigurationTest.java
@@ -45,7 +45,9 @@ public class CommonConfigurationTest {
         Assert.assertEquals("cluster-succeed!!!", 
config.getEventMeshCluster());
         Assert.assertEquals("name-succeed!!!", config.getEventMeshName());
         Assert.assertEquals("816", config.getSysID());
-        Assert.assertEquals("connector-succeed!!!", 
config.getEventMeshConnectorPluginType());
+        //Assert.assertEquals("connector-succeed!!!", 
config.getEventMeshConnectorPluginType());
+        Assert.assertEquals("storage-succeed!!!", 
config.getEventMeshStoragePluginType());
+        Assert.assertEquals("storage-succeed!!!", 
config.getEventMeshStoragePluginType());
         Assert.assertEquals("security-succeed!!!", 
config.getEventMeshSecurityPluginType());
         Assert.assertEquals("registry-succeed!!!", 
config.getEventMeshRegistryPluginType());
         Assert.assertEquals("trace-succeed!!!", 
config.getEventMeshTracePluginType());
diff --git a/eventmesh-common/src/test/resources/configuration.properties 
b/eventmesh-common/src/test/resources/configuration.properties
index 47bbb1c36..74b5cc025 100644
--- a/eventmesh-common/src/test/resources/configuration.properties
+++ b/eventmesh-common/src/test/resources/configuration.properties
@@ -21,6 +21,7 @@ eventMesh.server.cluster=cluster-succeed!!!
 eventMesh.server.name=name-succeed!!!
 eventMesh.server.hostIp=hostIp-succeed!!!
 eventMesh.connector.plugin.type=connector-succeed!!!
+eventMesh.storage.plugin.type=storage-succeed!!!
 eventMesh.security.plugin.type=security-succeed!!!
 eventMesh.registry.plugin.type=registry-succeed!!!
 eventMesh.trace.plugin=trace-succeed!!!
diff --git a/eventmesh-runtime/build.gradle b/eventmesh-runtime/build.gradle
index 8d9a8c217..c6102a100 100644
--- a/eventmesh-runtime/build.gradle
+++ b/eventmesh-runtime/build.gradle
@@ -37,7 +37,7 @@ dependencies {
 
     implementation project(":eventmesh-common")
     implementation project(":eventmesh-spi")
-    implementation 
project(":eventmesh-connector-plugin:eventmesh-connector-api")
+    implementation project(":eventmesh-storage:eventmesh-storage-api")
     implementation 
project(":eventmesh-connector-plugin:eventmesh-connector-standalone")
     implementation project(":eventmesh-security-plugin:eventmesh-security-api")
     implementation project(":eventmesh-security-plugin:eventmesh-security-acl")
diff --git a/eventmesh-runtime/conf/eventmesh.properties 
b/eventmesh-runtime/conf/eventmesh.properties
index bdcc87e99..b195e9414 100644
--- a/eventmesh-runtime/conf/eventmesh.properties
+++ b/eventmesh-runtime/conf/eventmesh.properties
@@ -74,6 +74,9 @@ eventMesh.server.blacklist.ipv6=::/128,::1/128,ff00::/8
 #connector plugin
 eventMesh.connector.plugin.type=standalone
 
+#storage plugin
+eventMesh.storage.plugin.type=standalone
+
 #security plugin
 eventMesh.server.security.enabled=false
 eventMesh.security.plugin.type=security
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/admin/controller/ClientManageController.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/admin/controller/ClientManageController.java
index 737a00405..1614fddfa 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/admin/controller/ClientManageController.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/admin/controller/ClientManageController.java
@@ -126,8 +126,8 @@ public class ClientManageController {
             eventMeshHTTPServer.getEventMeshHttpConfiguration(),
             eventMeshGrpcServer.getEventMeshGrpcConfiguration(), 
httpHandlerManager);
         new MetricsHandler(eventMeshHTTPServer, eventMeshTCPServer, 
httpHandlerManager);
-        new 
TopicHandler(eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshConnectorPluginType(),
 httpHandlerManager);
-        new 
EventHandler(eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshConnectorPluginType(),
 httpHandlerManager);
+        new 
TopicHandler(eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshStoragePluginType(),
 httpHandlerManager);
+        new 
EventHandler(eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshStoragePluginType(),
 httpHandlerManager);
         new RegistryHandler(eventMeshRegistry, httpHandlerManager);
 
         if 
(Objects.nonNull(adminWebHookConfigOperationManage.getWebHookConfigOperation()))
 {
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshServer.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshServer.java
index 5d9fe44db..dfbad0a38 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshServer.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshServer.java
@@ -24,10 +24,10 @@ import 
org.apache.eventmesh.common.utils.ConfigurationContextUtil;
 import org.apache.eventmesh.runtime.acl.Acl;
 import org.apache.eventmesh.runtime.admin.controller.ClientManageController;
 import org.apache.eventmesh.runtime.common.ServiceState;
-import org.apache.eventmesh.runtime.connector.ConnectorResource;
 import org.apache.eventmesh.runtime.constants.EventMeshConstants;
 import 
org.apache.eventmesh.runtime.core.protocol.http.producer.ProducerTopicManager;
 import org.apache.eventmesh.runtime.registry.Registry;
+import org.apache.eventmesh.runtime.storage.StorageResource;
 import org.apache.eventmesh.runtime.trace.Trace;
 
 import java.util.List;
@@ -46,7 +46,7 @@ public class EventMeshServer {
 
     private static Trace trace;
 
-    private final ConnectorResource connectorResource;
+    private final StorageResource storageResource;
 
     private ServiceState serviceState;
 
@@ -69,7 +69,7 @@ public class EventMeshServer {
         this.registry = 
Registry.getInstance(this.configuration.getEventMeshRegistryPluginType());
 
         trace = 
Trace.getInstance(this.configuration.getEventMeshTracePluginType(), 
this.configuration.isEventMeshServerTraceEnable());
-        this.connectorResource = 
ConnectorResource.getInstance(this.configuration.getEventMeshConnectorPluginType());
+        this.storageResource = 
StorageResource.getInstance(this.configuration.getEventMeshStoragePluginType());
 
         final List<String> provideServerProtocols = 
configuration.getEventMeshProvideServerProtocols();
         for (final String provideServerProtocol : provideServerProtocols) {
@@ -86,7 +86,7 @@ public class EventMeshServer {
     }
 
     public void init() throws Exception {
-        connectorResource.init();
+        storageResource.init();
         if (configuration.isEventMeshServerSecurityEnable()) {
             acl.init();
         }
@@ -177,7 +177,7 @@ public class EventMeshServer {
             registry.shutdown();
         }
 
-        connectorResource.release();
+        storageResource.release();
 
         if (configuration != null && 
configuration.isEventMeshServerSecurityEnable()) {
             acl.shutdown();
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQAdminWrapper.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQAdminWrapper.java
index 1d7a1459b..af74fadf0 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQAdminWrapper.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQAdminWrapper.java
@@ -21,7 +21,7 @@ package org.apache.eventmesh.runtime.core.plugin;
 
 import org.apache.eventmesh.api.admin.Admin;
 import org.apache.eventmesh.api.admin.TopicProperties;
-import org.apache.eventmesh.api.factory.ConnectorPluginFactory;
+import org.apache.eventmesh.api.factory.StoragePluginFactory;
 
 import java.util.List;
 import java.util.Properties;
@@ -36,8 +36,8 @@ public class MQAdminWrapper extends MQWrapper {
 
     protected Admin meshMQAdmin;
 
-    public MQAdminWrapper(String connectorPluginType) {
-        this.meshMQAdmin = 
ConnectorPluginFactory.getMeshMQAdmin(connectorPluginType);
+    public MQAdminWrapper(String storagePluginType) {
+        this.meshMQAdmin = 
StoragePluginFactory.getMeshMQAdmin(storagePluginType);
         if (meshMQAdmin == null) {
             log.error("can't load the meshMQAdmin plugin, please check.");
             throw new RuntimeException("doesn't load the meshMQAdmin plugin, 
please check.");
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQConsumerWrapper.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQConsumerWrapper.java
index 27bb56bd7..420a6b27c 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQConsumerWrapper.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQConsumerWrapper.java
@@ -20,7 +20,7 @@ package org.apache.eventmesh.runtime.core.plugin;
 import org.apache.eventmesh.api.AbstractContext;
 import org.apache.eventmesh.api.EventListener;
 import org.apache.eventmesh.api.consumer.Consumer;
-import org.apache.eventmesh.api.factory.ConnectorPluginFactory;
+import org.apache.eventmesh.api.factory.StoragePluginFactory;
 
 import java.util.List;
 import java.util.Properties;
@@ -35,8 +35,8 @@ public class MQConsumerWrapper extends MQWrapper {
 
     protected Consumer meshMQPushConsumer;
 
-    public MQConsumerWrapper(String connectorPluginType) {
-        this.meshMQPushConsumer = 
ConnectorPluginFactory.getMeshMQPushConsumer(connectorPluginType);
+    public MQConsumerWrapper(String storagePluginType) {
+        this.meshMQPushConsumer = 
StoragePluginFactory.getMeshMQPushConsumer(storagePluginType);
         if (meshMQPushConsumer == null) {
             log.error("can't load the meshMQPushConsumer plugin, please 
check.");
             throw new RuntimeException("doesn't load the meshMQPushConsumer 
plugin, please check.");
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQProducerWrapper.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQProducerWrapper.java
index 210af96ec..7b64e8281 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQProducerWrapper.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/plugin/MQProducerWrapper.java
@@ -19,7 +19,7 @@ package org.apache.eventmesh.runtime.core.plugin;
 
 import org.apache.eventmesh.api.RequestReplyCallback;
 import org.apache.eventmesh.api.SendCallback;
-import org.apache.eventmesh.api.factory.ConnectorPluginFactory;
+import org.apache.eventmesh.api.factory.StoragePluginFactory;
 import org.apache.eventmesh.api.producer.Producer;
 
 import java.util.Properties;
@@ -34,8 +34,8 @@ public class MQProducerWrapper extends MQWrapper {
 
     protected Producer meshMQProducer;
 
-    public MQProducerWrapper(String connectorPluginType) {
-        this.meshMQProducer = 
ConnectorPluginFactory.getMeshMQProducer(connectorPluginType);
+    public MQProducerWrapper(String storagePluginType) {
+        this.meshMQProducer = 
StoragePluginFactory.getMeshMQProducer(storagePluginType);
         if (meshMQProducer == null) {
             log.error("can't load the meshMQProducer plugin, please check.");
             throw new RuntimeException("doesn't load the meshMQProducer 
plugin, please check.");
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/EventMeshConsumer.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/EventMeshConsumer.java
index 1a1f81e7e..f5e7b0037 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/EventMeshConsumer.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/EventMeshConsumer.java
@@ -88,8 +88,8 @@ public class EventMeshConsumer {
         this.eventMeshGrpcConfiguration = 
eventMeshGrpcServer.getEventMeshGrpcConfiguration();
         this.consumerGroup = consumerGroup;
         this.messageHandler = new MessageHandler(consumerGroup, 
eventMeshGrpcServer.getPushMsgExecutor());
-        this.persistentMqConsumer = new 
MQConsumerWrapper(eventMeshGrpcConfiguration.getEventMeshConnectorPluginType());
-        this.broadcastMqConsumer = new 
MQConsumerWrapper(eventMeshGrpcConfiguration.getEventMeshConnectorPluginType());
+        this.persistentMqConsumer = new 
MQConsumerWrapper(eventMeshGrpcConfiguration.getEventMeshStoragePluginType());
+        this.broadcastMqConsumer = new 
MQConsumerWrapper(eventMeshGrpcConfiguration.getEventMeshStoragePluginType());
     }
 
     /**
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/producer/EventMeshProducer.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/producer/EventMeshProducer.java
index 12bc3844d..a5abe27d8 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/producer/EventMeshProducer.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/producer/EventMeshProducer.java
@@ -65,7 +65,7 @@ public class EventMeshProducer {
         //TODO for defibus
         keyValue.put(EventMeshConstants.EVENT_MESH_IDC, 
eventMeshGrpcConfiguration.getEventMeshIDC());
         mqProducerWrapper = new MQProducerWrapper(
-            eventMeshGrpcConfiguration.getEventMeshConnectorPluginType());
+            eventMeshGrpcConfiguration.getEventMeshStoragePluginType());
         mqProducerWrapper.init(keyValue);
         serviceState = ServiceState.INITED;
         log.info("EventMeshProducer [{}] inited...........", 
producerGroupConfig.getGroupName());
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/consumer/EventMeshConsumer.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/consumer/EventMeshConsumer.java
index ec6a61384..cb81a1d89 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/consumer/EventMeshConsumer.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/consumer/EventMeshConsumer.java
@@ -86,10 +86,8 @@ public class EventMeshConsumer {
     public EventMeshConsumer(EventMeshHTTPServer eventMeshHTTPServer, 
ConsumerGroupConf consumerGroupConf) {
         this.eventMeshHTTPServer = eventMeshHTTPServer;
         this.consumerGroupConf = consumerGroupConf;
-        this.persistentMqConsumer = new MQConsumerWrapper(
-            
eventMeshHTTPServer.getEventMeshHttpConfiguration().getEventMeshConnectorPluginType());
-        this.broadcastMqConsumer = new MQConsumerWrapper(
-            
eventMeshHTTPServer.getEventMeshHttpConfiguration().getEventMeshConnectorPluginType());
+        this.persistentMqConsumer = new 
MQConsumerWrapper(eventMeshHTTPServer.getEventMeshHttpConfiguration().getEventMeshStoragePluginType());
+        this.broadcastMqConsumer = new 
MQConsumerWrapper(eventMeshHTTPServer.getEventMeshHttpConfiguration().getEventMeshStoragePluginType());
     }
 
     private MessageHandler httpMessageHandler;
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/producer/EventMeshProducer.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/producer/EventMeshProducer.java
index 1ebf12815..7e4c47d1a 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/producer/EventMeshProducer.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/producer/EventMeshProducer.java
@@ -90,7 +90,7 @@ public class EventMeshProducer {
 
         //TODO for defibus
         keyValue.put("eventMeshIDC", 
eventMeshHttpConfiguration.getEventMeshIDC());
-        mqProducerWrapper = new 
MQProducerWrapper(eventMeshHttpConfiguration.getEventMeshConnectorPluginType());
+        mqProducerWrapper = new 
MQProducerWrapper(eventMeshHttpConfiguration.getEventMeshStoragePluginType());
         mqProducerWrapper.init(keyValue);
         log.info("EventMeshProducer [{}] inited.............", 
producerGroupConfig.getGroupName());
     }
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientGroupWrapper.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientGroupWrapper.java
index b7c511352..0860ba9c7 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientGroupWrapper.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientGroupWrapper.java
@@ -124,12 +124,9 @@ public class ClientGroupWrapper {
         this.eventMeshTcpMonitor =
             
Preconditions.checkNotNull(eventMeshTCPServer.getEventMeshTcpMonitor());
         this.downstreamDispatchStrategy = downstreamDispatchStrategy;
-        this.persistentMsgConsumer = new MQConsumerWrapper(
-            
eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshConnectorPluginType());
-        this.broadCastMsgConsumer = new MQConsumerWrapper(
-            
eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshConnectorPluginType());
-        this.mqProducerWrapper = new MQProducerWrapper(
-            
eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshConnectorPluginType());
+        this.persistentMsgConsumer = new 
MQConsumerWrapper(eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshStoragePluginType());
+        this.broadCastMsgConsumer = new 
MQConsumerWrapper(eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshStoragePluginType());
+        this.mqProducerWrapper = new 
MQProducerWrapper(eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshStoragePluginType());
     }
 
     public ConcurrentHashMap<String, Set<Session>> 
getTopic2sessionInGroupMapping() {
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/connector/ConnectorResource.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/storage/StorageResource.java
similarity index 61%
rename from 
eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/connector/ConnectorResource.java
rename to 
eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/storage/StorageResource.java
index 21ce82ccb..5dd8720e6 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/connector/ConnectorResource.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/storage/StorageResource.java
@@ -15,59 +15,54 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.runtime.connector;
+package org.apache.eventmesh.runtime.storage;
 
-import org.apache.eventmesh.api.connector.ConnectorResourceService;
+import org.apache.eventmesh.api.storage.StorageResourceService;
 import org.apache.eventmesh.spi.EventMeshExtensionFactory;
 
-
 import java.util.HashMap;
 import java.util.Map;
 import java.util.concurrent.atomic.AtomicBoolean;
 
-
 import lombok.extern.slf4j.Slf4j;
 
 @Slf4j
-public class ConnectorResource {
+public class StorageResource {
 
-    private static final Map<String, ConnectorResource> 
CONNECTOR_RESOURCE_CACHE = new HashMap<>(16);
+    private static final Map<String, StorageResource> STORAGE_RESOURCE_CACHE = 
new HashMap<>(16);
 
-    private ConnectorResourceService connectorResourceService;
+    private StorageResourceService storageResourceService;
 
     private final AtomicBoolean inited = new AtomicBoolean(false);
 
     private final AtomicBoolean released = new AtomicBoolean(false);
 
-    private ConnectorResource() {
+    private StorageResource() {
 
     }
 
-    public static ConnectorResource getInstance(String 
connectorResourcePluginType) {
-        return CONNECTOR_RESOURCE_CACHE.computeIfAbsent(
-            connectorResourcePluginType,
-            ConnectorResource::connectorResourceBuilder
-        );
+    public static StorageResource getInstance(String 
storageResourcePluginType) {
+        return 
STORAGE_RESOURCE_CACHE.computeIfAbsent(storageResourcePluginType, 
StorageResource::connectorResourceBuilder);
     }
 
-    private static ConnectorResource connectorResourceBuilder(String 
connectorResourcePluginType) {
-        ConnectorResourceService connectorResourceServiceExt = 
EventMeshExtensionFactory.getExtension(ConnectorResourceService.class,
+    private static StorageResource connectorResourceBuilder(String 
connectorResourcePluginType) {
+        StorageResourceService connectorResourceServiceExt = 
EventMeshExtensionFactory.getExtension(StorageResourceService.class,
             connectorResourcePluginType);
         if (connectorResourceServiceExt == null) {
             String errorMsg = "can't load the connectorResourceService plugin, 
please check.";
             log.error(errorMsg);
             throw new RuntimeException(errorMsg);
         }
-        ConnectorResource connectorResource = new ConnectorResource();
-        connectorResource.connectorResourceService = 
connectorResourceServiceExt;
-        return connectorResource;
+        StorageResource storageResource = new StorageResource();
+        storageResource.storageResourceService = connectorResourceServiceExt;
+        return storageResource;
     }
 
     public void init() throws Exception {
         if (!inited.compareAndSet(false, true)) {
             return;
         }
-        connectorResourceService.init();
+        storageResourceService.init();
     }
 
     public void release() throws Exception {
@@ -75,6 +70,6 @@ public class ConnectorResource {
             return;
         }
         inited.compareAndSet(true, false);
-        connectorResourceService.release();
+        storageResourceService.release();
     }
 }
diff --git 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/admin/controller/ClientManageControllerTest.java
 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/admin/controller/ClientManageControllerTest.java
index beb1a92e2..ab8738a2c 100644
--- 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/admin/controller/ClientManageControllerTest.java
+++ 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/admin/controller/ClientManageControllerTest.java
@@ -82,7 +82,7 @@ public class ClientManageControllerTest {
             eventMeshHTTPServer, eventMeshGrpcServer, registry);
         
controller.setAdminWebHookConfigOperationManage(adminWebHookConfigOperationManage);
 
-        
eventMeshTCPServer.getEventMeshTCPConfiguration().setEventMeshConnectorPluginType("standalone");
+        
eventMeshTCPServer.getEventMeshTCPConfiguration().setEventMeshStoragePluginType("standalone");
 
         try (MockedStatic<HttpServer> dummyStatic = 
Mockito.mockStatic(HttpServer.class)) {
             HttpServer server = mock(HttpServer.class);
diff --git 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/boot/EventMeshServerTest.java
 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/boot/EventMeshServerTest.java
index fb35c499c..5f6b57b1a 100644
--- 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/boot/EventMeshServerTest.java
+++ 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/boot/EventMeshServerTest.java
@@ -94,6 +94,7 @@ public class EventMeshServerTest {
         Assert.assertEquals("name-succeed!!!", config.getEventMeshName());
         Assert.assertEquals("816", config.getSysID());
         Assert.assertEquals("connector-succeed!!!", 
config.getEventMeshConnectorPluginType());
+        Assert.assertEquals("storage-succeed!!!", 
config.getEventMeshStoragePluginType());
         Assert.assertEquals("security-succeed!!!", 
config.getEventMeshSecurityPluginType());
         Assert.assertEquals("registry-succeed!!!", 
config.getEventMeshRegistryPluginType());
         Assert.assertEquals("trace-succeed!!!", 
config.getEventMeshTracePluginType());
diff --git 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshGrpcConfigurationTest.java
 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshGrpcConfigurationTest.java
index 5f107ab3b..75afe7d2e 100644
--- 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshGrpcConfigurationTest.java
+++ 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshGrpcConfigurationTest.java
@@ -73,6 +73,7 @@ public class EventMeshGrpcConfigurationTest {
         Assert.assertEquals("name-succeed!!!", config.getEventMeshName());
         Assert.assertEquals("816", config.getSysID());
         Assert.assertEquals("connector-succeed!!!", 
config.getEventMeshConnectorPluginType());
+        Assert.assertEquals("storage-succeed!!!", 
config.getEventMeshStoragePluginType());
         Assert.assertEquals("security-succeed!!!", 
config.getEventMeshSecurityPluginType());
         Assert.assertEquals("registry-succeed!!!", 
config.getEventMeshRegistryPluginType());
         Assert.assertEquals("trace-succeed!!!", 
config.getEventMeshTracePluginType());
@@ -90,4 +91,4 @@ public class EventMeshGrpcConfigurationTest {
 
         Assert.assertEquals("eventmesh.idc-succeed!!!", 
config.getEventMeshWebhookOrigin());
     }
-}
\ No newline at end of file
+}
diff --git 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshHTTPConfigurationTest.java
 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshHTTPConfigurationTest.java
index 72c0b5e14..8faa4f1a1 100644
--- 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshHTTPConfigurationTest.java
+++ 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshHTTPConfigurationTest.java
@@ -88,6 +88,7 @@ public class EventMeshHTTPConfigurationTest {
         Assert.assertEquals("name-succeed!!!", config.getEventMeshName());
         Assert.assertEquals("816", config.getSysID());
         Assert.assertEquals("connector-succeed!!!", 
config.getEventMeshConnectorPluginType());
+        Assert.assertEquals("storage-succeed!!!", 
config.getEventMeshStoragePluginType());
         Assert.assertEquals("security-succeed!!!", 
config.getEventMeshSecurityPluginType());
         Assert.assertEquals("registry-succeed!!!", 
config.getEventMeshRegistryPluginType());
         Assert.assertEquals("trace-succeed!!!", 
config.getEventMeshTracePluginType());
@@ -105,4 +106,4 @@ public class EventMeshHTTPConfigurationTest {
 
         Assert.assertEquals("eventmesh.idc-succeed!!!", 
config.getEventMeshWebhookOrigin());
     }
-}
\ No newline at end of file
+}
diff --git 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshTCPConfigurationTest.java
 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshTCPConfigurationTest.java
index 1f3163221..fa8ce064a 100644
--- 
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshTCPConfigurationTest.java
+++ 
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/configuration/EventMeshTCPConfigurationTest.java
@@ -75,6 +75,7 @@ public class EventMeshTCPConfigurationTest {
         Assert.assertEquals("name-succeed!!!", config.getEventMeshName());
         Assert.assertEquals("816", config.getSysID());
         Assert.assertEquals("connector-succeed!!!", 
config.getEventMeshConnectorPluginType());
+        Assert.assertEquals("storage-succeed!!!", 
config.getEventMeshStoragePluginType());
         Assert.assertEquals("security-succeed!!!", 
config.getEventMeshSecurityPluginType());
         Assert.assertEquals("registry-succeed!!!", 
config.getEventMeshRegistryPluginType());
         Assert.assertEquals("trace-succeed!!!", 
config.getEventMeshTracePluginType());
@@ -92,4 +93,4 @@ public class EventMeshTCPConfigurationTest {
 
         Assert.assertEquals("eventmesh.idc-succeed!!!", 
config.getEventMeshWebhookOrigin());
     }
-}
\ No newline at end of file
+}
diff --git a/eventmesh-runtime/src/test/resources/configuration.properties 
b/eventmesh-runtime/src/test/resources/configuration.properties
index 190fd2fc5..ef3fe9266 100644
--- a/eventmesh-runtime/src/test/resources/configuration.properties
+++ b/eventmesh-runtime/src/test/resources/configuration.properties
@@ -23,6 +23,7 @@ eventMesh.server.cluster=cluster-succeed!!!
 eventMesh.server.name=name-succeed!!!
 eventMesh.server.hostIp=hostIp-succeed!!!
 eventMesh.connector.plugin.type=connector-succeed!!!
+eventMesh.storage.plugin.type=storage-succeed!!!
 eventMesh.security.plugin.type=security-succeed!!!
 eventMesh.registry.plugin.type=registry-succeed!!!
 eventMesh.trace.plugin=trace-succeed!!!


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to