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]