This is an automated email from the ASF dual-hosted git repository.
jonyang 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 d0f05e24b [ISSUE #3118]Refactor Acl (#3119)
d0f05e24b is described below
commit d0f05e24b74d4a33d502954b5b3c5fd099eaaa0f
Author: mxsm <[email protected]>
AuthorDate: Thu Feb 16 09:04:05 2023 +0800
[ISSUE #3118]Refactor Acl (#3119)
* [ISSUE #3118]Refactor Acl
* polish code
* add init,start,stop flag to avoid call multiple times
---
.../eventmesh/common/config/ConfigService.java | 6 +-
.../java/org/apache/eventmesh/runtime/acl/Acl.java | 96 +++++++++------
.../runtime/boot/EventMeshGrpcBootstrap.java | 9 +-
.../runtime/boot/EventMeshGrpcServer.java | 23 +++-
.../runtime/boot/EventMeshHTTPServer.java | 8 ++
.../runtime/boot/EventMeshHttpBootstrap.java | 5 +-
.../eventmesh/runtime/boot/EventMeshServer.java | 32 +++--
.../eventmesh/runtime/boot/EventMeshStartup.java | 1 +
.../eventmesh/runtime/boot/EventMeshTCPServer.java | 13 +-
.../runtime/boot/EventMeshTcpBootstrap.java | 8 +-
.../processor/BatchPublishMessageProcessor.java | 7 +-
.../grpc/processor/HeartbeatProcessor.java | 7 +-
.../grpc/processor/ReplyMessageProcessor.java | 7 +-
.../grpc/processor/RequestMessageProcessor.java | 7 +-
.../grpc/processor/SendAsyncMessageProcessor.java | 7 +-
.../grpc/processor/SubscribeProcessor.java | 7 +-
.../grpc/processor/SubscribeStreamProcessor.java | 11 +-
.../http/processor/BatchSendMessageProcessor.java | 7 +-
.../processor/BatchSendMessageV2Processor.java | 5 +-
.../http/processor/HeartBeatProcessor.java | 5 +-
.../processor/LocalSubscribeEventProcessor.java | 5 +-
.../processor/RemoteSubscribeEventProcessor.java | 9 +-
.../http/processor/SendAsyncEventProcessor.java | 9 +-
.../http/processor/SendAsyncMessageProcessor.java | 11 +-
.../processor/SendAsyncRemoteEventProcessor.java | 5 +-
.../http/processor/SendSyncMessageProcessor.java | 5 +-
.../http/processor/SubscribeProcessor.java | 20 ++--
.../protocol/tcp/client/task/HeartBeatTask.java | 5 +-
.../core/protocol/tcp/client/task/HelloTask.java | 5 +-
.../tcp/client/task/MessageTransferTask.java | 131 ++++++++++-----------
.../protocol/tcp/client/task/SubscribeTask.java | 12 +-
31 files changed, 300 insertions(+), 188 deletions(-)
diff --git
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/config/ConfigService.java
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/config/ConfigService.java
index f960ee138..0cc9ea94f 100644
---
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/config/ConfigService.java
+++
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/config/ConfigService.java
@@ -53,7 +53,7 @@ public class ConfigService {
return INSTANCE;
}
- public ConfigService() {
+ private ConfigService() {
}
public ConfigService setConfigPath(String configPath) {
@@ -164,7 +164,7 @@ public class ConfigService {
}
Object configObject = this.getConfig(configInfo);
-
+
try {
field.setAccessible(true);
field.set(object, configObject);
@@ -178,4 +178,4 @@ public class ConfigService {
configMonitorService.monitor(configInfo);
}
}
-}
\ No newline at end of file
+}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/acl/Acl.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/acl/Acl.java
index c37f3ed94..c7facea48 100644
--- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/acl/Acl.java
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/acl/Acl.java
@@ -25,46 +25,86 @@ import org.apache.eventmesh.spi.EventMeshExtensionFactory;
import org.apache.commons.lang3.StringUtils;
+import java.util.HashMap;
+import java.util.Map;
import java.util.Properties;
-import java.util.ServiceLoader;
+import java.util.concurrent.atomic.AtomicBoolean;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+
+import lombok.extern.slf4j.Slf4j;
+
+@Slf4j
public class Acl {
- private static final Logger logger = LoggerFactory.getLogger(Acl.class);
- private static AclService aclService;
- public void init(String aclPluginType) throws AclException {
- aclService = EventMeshExtensionFactory.getExtension(AclService.class,
aclPluginType);
- if (aclService == null) {
- logger.error("can't load the aclService plugin, please check.");
+ private static final Map<String, Acl> ACL_CACHE = new HashMap<>(16);
+
+ private AclService aclService;
+
+ private final AtomicBoolean inited = new AtomicBoolean(false);
+
+ private final AtomicBoolean started = new AtomicBoolean(false);
+
+ private final AtomicBoolean shutdown = new AtomicBoolean(false);
+
+ private Acl() {
+
+ }
+
+ public static Acl getInstance(String aclPluginType) {
+ return ACL_CACHE.computeIfAbsent(aclPluginType, key ->
aclBuilder(key));
+ }
+
+ private static Acl aclBuilder(String aclPluginType) {
+ AclService aclServiceExt =
EventMeshExtensionFactory.getExtension(AclService.class, aclPluginType);
+ if (aclServiceExt == null) {
+ log.error("can't load the aclService plugin, please check.");
throw new RuntimeException("doesn't load the aclService plugin,
please check.");
}
+ Acl acl = new Acl();
+ acl.aclService = aclServiceExt;
+
+ return acl;
+ }
+
+ public void init() throws AclException {
+ if (!inited.compareAndSet(false, true)) {
+ return;
+ }
aclService.init();
}
public void start() throws AclException {
+ if (!started.compareAndSet(false, true)) {
+ return;
+ }
aclService.start();
}
public void shutdown() throws AclException {
+ inited.compareAndSet(true, false);
+ started.compareAndSet(true, false);
+ if (!shutdown.compareAndSet(false, true)) {
+ return;
+ }
aclService.shutdown();
}
- public static void doAclCheckInTcpConnect(String remoteAddr, UserAgent
userAgent, int requestCode) throws AclException {
+ public void doAclCheckInTcpConnect(String remoteAddr, UserAgent userAgent,
int requestCode) throws AclException {
aclService.doAclCheckInConnect(buildTcpAclProperties(remoteAddr,
userAgent, null, requestCode));
}
- public static void doAclCheckInTcpHeartbeat(String remoteAddr, UserAgent
userAgent, int requestCode) throws AclException {
+ public void doAclCheckInTcpHeartbeat(String remoteAddr, UserAgent
userAgent, int requestCode) throws AclException {
aclService.doAclCheckInHeartbeat(buildTcpAclProperties(remoteAddr,
userAgent, null, requestCode));
}
- public static void doAclCheckInTcpSend(String remoteAddr, UserAgent
userAgent, String topic, int requestCode) throws AclException {
+ public void doAclCheckInTcpSend(String remoteAddr, UserAgent userAgent,
String topic, int requestCode) throws AclException {
aclService.doAclCheckInSend(buildTcpAclProperties(remoteAddr,
userAgent, topic, requestCode));
}
- public static void doAclCheckInTcpReceive(String remoteAddr, UserAgent
userAgent, String topic, int requestCode) throws AclException {
+ public void doAclCheckInTcpReceive(String remoteAddr, UserAgent userAgent,
String topic, int requestCode) throws AclException {
aclService.doAclCheckInReceive(buildTcpAclProperties(remoteAddr,
userAgent, topic, requestCode));
}
@@ -81,33 +121,32 @@ public class Acl {
return aclProperties;
}
- public static void doAclCheckInHttpSend(String remoteAddr, String user,
String pass, String subsystem, String topic,
- int requestCode) throws
AclException {
+ public void doAclCheckInHttpSend(String remoteAddr, String user, String
pass, String subsystem, String topic, int requestCode)
+ throws AclException {
aclService.doAclCheckInSend(buildHttpAclProperties(remoteAddr, user,
pass, subsystem, topic, requestCode));
}
- public static void doAclCheckInHttpSend(String remoteAddr, String user,
String pass, String subsystem, String topic,
- String requestURI) throws
AclException {
+ public void doAclCheckInHttpSend(String remoteAddr, String user, String
pass, String subsystem, String topic,
+ String requestURI) throws AclException {
aclService.doAclCheckInSend(buildHttpAclProperties(remoteAddr, user,
pass, subsystem, topic, requestURI));
}
- public static void doAclCheckInHttpReceive(String remoteAddr, String user,
String pass, String subsystem, String topic,
- int requestCode) throws
AclException {
+ public void doAclCheckInHttpReceive(String remoteAddr, String user, String
pass, String subsystem, String topic,
+ int requestCode) throws AclException {
aclService.doAclCheckInReceive(buildHttpAclProperties(remoteAddr,
user, pass, subsystem, topic, requestCode));
}
- public static void doAclCheckInHttpReceive(String remoteAddr, String user,
String pass, String subsystem, String topic,
- String requestURI) throws
AclException {
+ public void doAclCheckInHttpReceive(String remoteAddr, String user, String
pass, String subsystem, String topic,
+ String requestURI) throws AclException {
aclService.doAclCheckInReceive(buildHttpAclProperties(remoteAddr,
user, pass, subsystem, topic, requestURI));
}
- public static void doAclCheckInHttpHeartbeat(String remoteAddr, String
user, String pass, String subsystem, String topic,
- int requestCode) throws
AclException {
+ public void doAclCheckInHttpHeartbeat(String remoteAddr, String user,
String pass, String subsystem, String topic,
+ int requestCode) throws AclException {
aclService.doAclCheckInHeartbeat(buildHttpAclProperties(remoteAddr,
user, pass, subsystem, topic, requestCode));
}
- private static Properties buildHttpAclProperties(String remoteAddr, String
user, String pass, String subsystem,
- String topic, int
requestCode) {
+ private Properties buildHttpAclProperties(String remoteAddr, String user,
String pass, String subsystem, String topic, int requestCode) {
Properties aclProperties = new Properties();
aclProperties.put(AclPropertyKeys.CLIENT_IP, remoteAddr);
aclProperties.put(AclPropertyKeys.USER, user);
@@ -120,8 +159,7 @@ public class Acl {
return aclProperties;
}
- private static Properties buildHttpAclProperties(String remoteAddr, String
user, String pass, String subsystem,
- String topic, String
requestURI) {
+ private Properties buildHttpAclProperties(String remoteAddr, String user,
String pass, String subsystem, String topic, String requestURI) {
Properties aclProperties = new Properties();
aclProperties.put(AclPropertyKeys.CLIENT_IP, remoteAddr);
aclProperties.put(AclPropertyKeys.USER, user);
@@ -133,12 +171,4 @@ public class Acl {
}
return aclProperties;
}
-
- private AclService getSpiAclService() {
- ServiceLoader<AclService> serviceLoader =
ServiceLoader.load(AclService.class);
- if (serviceLoader.iterator().hasNext()) {
- return serviceLoader.iterator().next();
- }
- return null;
- }
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshGrpcBootstrap.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshGrpcBootstrap.java
index 51a6ed5ee..1654c8312 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshGrpcBootstrap.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshGrpcBootstrap.java
@@ -28,11 +28,10 @@ public class EventMeshGrpcBootstrap implements
EventMeshBootstrap {
private EventMeshGrpcServer eventMeshGrpcServer;
- private final Registry registry;
-
- public EventMeshGrpcBootstrap(Registry registry) {
- this.registry = registry;
+ private final EventMeshServer eventMeshServer;
+ public EventMeshGrpcBootstrap(final EventMeshServer eventMeshServer) {
+ this.eventMeshServer = eventMeshServer;
ConfigService configService = ConfigService.getInstance();
this.eventMeshGrpcConfiguration =
configService.buildConfigInstance(EventMeshGrpcConfiguration.class);
@@ -43,7 +42,7 @@ public class EventMeshGrpcBootstrap implements
EventMeshBootstrap {
public void init() throws Exception {
// server init
if (eventMeshGrpcConfiguration != null) {
- eventMeshGrpcServer = new
EventMeshGrpcServer(eventMeshGrpcConfiguration, registry);
+ eventMeshGrpcServer = new
EventMeshGrpcServer(this.eventMeshServer, this.eventMeshGrpcConfiguration);
eventMeshGrpcServer.init();
}
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshGrpcServer.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshGrpcServer.java
index 7aa4439a3..5494d9f72 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshGrpcServer.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshGrpcServer.java
@@ -25,6 +25,7 @@ import
org.apache.eventmesh.common.utils.ConfigurationContextUtil;
import org.apache.eventmesh.common.utils.IPUtils;
import org.apache.eventmesh.metrics.api.MetricsPluginFactory;
import org.apache.eventmesh.metrics.api.MetricsRegistry;
+import org.apache.eventmesh.runtime.acl.Acl;
import org.apache.eventmesh.runtime.configuration.EventMeshGrpcConfiguration;
import org.apache.eventmesh.runtime.constants.EventMeshConstants;
import
org.apache.eventmesh.runtime.core.protocol.grpc.consumer.ConsumerManager;
@@ -87,13 +88,19 @@ public class EventMeshGrpcServer {
private RateLimiter msgRateLimiter;
- private Registry registry;
+ private final Registry registry;
+
+ private final Acl acl;
+
+ private final EventMeshServer eventMeshServer;
private EventMeshGrpcMonitor eventMeshGrpcMonitor;
- public EventMeshGrpcServer(EventMeshGrpcConfiguration
eventMeshGrpcConfiguration, Registry registry) {
+ public EventMeshGrpcServer(final EventMeshServer eventMeshServer, final
EventMeshGrpcConfiguration eventMeshGrpcConfiguration) {
+ this.eventMeshServer = eventMeshServer;
this.eventMeshGrpcConfiguration = eventMeshGrpcConfiguration;
- this.registry = registry;
+ this.registry = eventMeshServer.getRegistry();
+ this.acl = eventMeshServer.getAcl();
}
public void init() throws Exception {
@@ -307,4 +314,12 @@ public class EventMeshGrpcServer {
itr.remove();
}
}
-}
\ No newline at end of file
+
+ public Registry getRegistry() {
+ return registry;
+ }
+
+ public Acl getAcl() {
+ return acl;
+ }
+}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshHTTPServer.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshHTTPServer.java
index 5a08d6334..7bf624f32 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshHTTPServer.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshHTTPServer.java
@@ -26,6 +26,7 @@ import
org.apache.eventmesh.common.utils.ConfigurationContextUtil;
import org.apache.eventmesh.common.utils.IPUtils;
import org.apache.eventmesh.metrics.api.MetricsPluginFactory;
import org.apache.eventmesh.metrics.api.MetricsRegistry;
+import org.apache.eventmesh.runtime.acl.Acl;
import org.apache.eventmesh.runtime.common.ServiceState;
import org.apache.eventmesh.runtime.configuration.EventMeshHTTPConfiguration;
import org.apache.eventmesh.runtime.constants.EventMeshConstants;
@@ -80,6 +81,8 @@ public class EventMeshHTTPServer extends AbstractHTTPServer {
private final transient Registry registry;
+ private final Acl acl;
+
public final transient EventBus eventBus = new EventBus();
private transient ConsumerManager consumerManager;
@@ -119,6 +122,7 @@ public class EventMeshHTTPServer extends AbstractHTTPServer
{
this.eventMeshServer = eventMeshServer;
this.eventMeshHttpConfiguration = eventMeshHttpConfiguration;
this.registry = eventMeshServer.getRegistry();
+ this.acl = eventMeshServer.getAcl();
}
@@ -433,4 +437,8 @@ public class EventMeshHTTPServer extends AbstractHTTPServer
{
public HttpRetryer getHttpRetryer() {
return httpRetryer;
}
+
+ public Acl getAcl() {
+ return acl;
+ }
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshHttpBootstrap.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshHttpBootstrap.java
index d478400b7..48e1a2ea4 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshHttpBootstrap.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshHttpBootstrap.java
@@ -30,11 +30,8 @@ public class EventMeshHttpBootstrap implements
EventMeshBootstrap {
private final EventMeshServer eventMeshServer;
- private final Registry registry;
-
- public EventMeshHttpBootstrap(EventMeshServer eventMeshServer, Registry
registry) {
+ public EventMeshHttpBootstrap(final EventMeshServer eventMeshServer) {
this.eventMeshServer = eventMeshServer;
- this.registry = registry;
ConfigService configService = ConfigService.getInstance();
this.eventMeshHttpConfiguration =
configService.buildConfigInstance(EventMeshHTTPConfiguration.class);
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 e8eddb296..f6d0d026c 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
@@ -39,7 +39,7 @@ public class EventMeshServer {
public static final Logger LOGGER =
LoggerFactory.getLogger(EventMeshServer.class);
- private final Acl acl;
+ private Acl acl;
private Registry registry;
@@ -57,36 +57,38 @@ public class EventMeshServer {
private static final String SERVER_STATE_MSG = "server state:{}";
- public EventMeshServer() throws Exception {
- ConfigService configService = ConfigService.getInstance();
+ private static final ConfigService configService =
ConfigService.getInstance();
+
+ public EventMeshServer() {
+
this.configuration =
configService.buildConfigInstance(CommonConfiguration.class);
- this.acl = new Acl();
+ this.acl =
Acl.getInstance(this.configuration.getEventMeshSecurityPluginType());
this.registry = new Registry();
- trace = new Trace(configuration.isEventMeshServerTraceEnable());
+
+ trace = new Trace(this.configuration.isEventMeshServerTraceEnable());
this.connectorResource = new ConnectorResource();
final List<String> provideServerProtocols =
configuration.getEventMeshProvideServerProtocols();
for (final String provideServerProtocol : provideServerProtocols) {
if (ConfigurationContextUtil.HTTP.equals(provideServerProtocol)) {
- BOOTSTRAP_LIST.add(new EventMeshHttpBootstrap(this, registry));
+ BOOTSTRAP_LIST.add(new EventMeshHttpBootstrap(this));
}
if (ConfigurationContextUtil.TCP.equals(provideServerProtocol)) {
- BOOTSTRAP_LIST.add(new EventMeshTcpBootstrap(this, registry));
+ BOOTSTRAP_LIST.add(new EventMeshTcpBootstrap(this));
}
if (ConfigurationContextUtil.GRPC.equals(provideServerProtocol)) {
- BOOTSTRAP_LIST.add(new EventMeshGrpcBootstrap(registry));
+ BOOTSTRAP_LIST.add(new EventMeshGrpcBootstrap(this));
}
}
-
- init();
}
- private void init() throws Exception {
+ public void init() throws Exception {
if (Objects.nonNull(configuration)) {
+
connectorResource.init(configuration.getEventMeshConnectorPluginType());
if (configuration.isEventMeshServerSecurityEnable()) {
- acl.init(configuration.getEventMeshSecurityPluginType());
+ acl.init();
}
if (configuration.isEventMeshServerRegistryEnable()) {
registry.init(configuration.getEventMeshRegistryPluginType());
@@ -94,6 +96,7 @@ public class EventMeshServer {
if (configuration.isEventMeshServerTraceEnable()) {
trace.init(configuration.getEventMeshTracePluginType());
}
+
}
EventMeshTCPServer eventMeshTCPServer = null;
@@ -154,7 +157,6 @@ public class EventMeshServer {
clientManageController.start();
}
-
serviceState = ServiceState.RUNNING;
if (LOGGER.isInfoEnabled()) {
LOGGER.info(SERVER_STATE_MSG, serviceState);
@@ -209,4 +211,8 @@ public class EventMeshServer {
public void setRegistry(final Registry registry) {
this.registry = registry;
}
+
+ public Acl getAcl() {
+ return acl;
+ }
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshStartup.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshStartup.java
index 986726c92..ca56c052e 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshStartup.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshStartup.java
@@ -36,6 +36,7 @@ public class EventMeshStartup {
.setRootConfig(EventMeshConstants.EVENTMESH_CONF_FILE);
EventMeshServer server = new EventMeshServer();
+ server.init();
server.start();
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
try {
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTCPServer.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTCPServer.java
index 91a9fb84f..da08004e8 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTCPServer.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTCPServer.java
@@ -28,6 +28,7 @@ import org.apache.eventmesh.common.utils.IPUtils;
import org.apache.eventmesh.common.utils.ThreadUtils;
import org.apache.eventmesh.metrics.api.MetricsPluginFactory;
import org.apache.eventmesh.metrics.api.MetricsRegistry;
+import org.apache.eventmesh.runtime.acl.Acl;
import org.apache.eventmesh.runtime.configuration.EventMeshTCPConfiguration;
import org.apache.eventmesh.runtime.constants.EventMeshConstants;
import
org.apache.eventmesh.runtime.core.protocol.tcp.client.EventMeshTcpConnectionHandler;
@@ -90,6 +91,8 @@ public class EventMeshTCPServer extends
AbstractRemotingServer {
private final transient Registry registry;
+ private final Acl acl;
+
private transient EventMeshRebalanceService eventMeshRebalanceService;
private transient AdminWebHookConfigOperationManager
adminWebHookConfigOperationManage;
@@ -129,12 +132,12 @@ public class EventMeshTCPServer extends
AbstractRemotingServer {
}
- public EventMeshTCPServer(final EventMeshServer eventMeshServer,
- final EventMeshTCPConfiguration
eventMeshTCPConfiguration, final Registry registry) {
+ public EventMeshTCPServer(final EventMeshServer eventMeshServer, final
EventMeshTCPConfiguration eventMeshTCPConfiguration) {
super();
this.eventMeshServer = eventMeshServer;
this.eventMeshTCPConfiguration = eventMeshTCPConfiguration;
- this.registry = registry;
+ this.registry = eventMeshServer.getRegistry();
+ this.acl = eventMeshServer.getAcl();
}
private void startServer() {
@@ -409,4 +412,8 @@ public class EventMeshTCPServer extends
AbstractRemotingServer {
public void
setAdminWebHookConfigOperationManage(AdminWebHookConfigOperationManager
adminWebHookConfigOperationManage) {
this.adminWebHookConfigOperationManage =
adminWebHookConfigOperationManage;
}
+
+ public Acl getAcl() {
+ return acl;
+ }
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTcpBootstrap.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTcpBootstrap.java
index eecbd1cf2..ee38f7fc6 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTcpBootstrap.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTcpBootstrap.java
@@ -30,11 +30,9 @@ public class EventMeshTcpBootstrap implements
EventMeshBootstrap {
private final EventMeshServer eventMeshServer;
- private final Registry registry;
-
- public EventMeshTcpBootstrap(EventMeshServer eventMeshServer, Registry
registry) {
+ public EventMeshTcpBootstrap(EventMeshServer eventMeshServer) {
this.eventMeshServer = eventMeshServer;
- this.registry = registry;
+
ConfigService configService = ConfigService.getInstance();
this.eventMeshTcpConfiguration =
configService.buildConfigInstance(EventMeshTCPConfiguration.class);
@@ -46,7 +44,7 @@ public class EventMeshTcpBootstrap implements
EventMeshBootstrap {
public void init() throws Exception {
// server init
if (eventMeshTcpConfiguration != null) {
- eventMeshTcpServer = new EventMeshTCPServer(eventMeshServer,
eventMeshTcpConfiguration, registry);
+ eventMeshTcpServer = new EventMeshTCPServer(eventMeshServer,
eventMeshTcpConfiguration);
eventMeshTcpServer.init();
}
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/BatchPublishMessageProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/BatchPublishMessageProcessor.java
index becb88f90..840fc8633 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/BatchPublishMessageProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/BatchPublishMessageProcessor.java
@@ -56,8 +56,11 @@ public class BatchPublishMessageProcessor {
private final EventMeshGrpcServer eventMeshGrpcServer;
- public BatchPublishMessageProcessor(EventMeshGrpcServer
eventMeshGrpcServer) {
+ private final Acl acl;
+
+ public BatchPublishMessageProcessor(final EventMeshGrpcServer
eventMeshGrpcServer) {
this.eventMeshGrpcServer = eventMeshGrpcServer;
+ this.acl = eventMeshGrpcServer.getAcl();
}
public void process(BatchMessage message, EventEmitter<Response> emitter)
throws Exception {
@@ -133,7 +136,7 @@ public class BatchPublishMessageProcessor {
String pass = requestHeader.getPassword();
String subsystem = requestHeader.getSys();
String topic = message.getTopic();
- Acl.doAclCheckInHttpSend(remoteAdd, user, pass, subsystem, topic,
RequestCode.MSG_SEND_ASYNC.getRequestCode());
+ this.acl.doAclCheckInHttpSend(remoteAdd, user, pass, subsystem,
topic, RequestCode.MSG_SEND_ASYNC.getRequestCode());
}
}
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/HeartbeatProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/HeartbeatProcessor.java
index 26de2cb8e..8219f2b40 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/HeartbeatProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/HeartbeatProcessor.java
@@ -42,8 +42,11 @@ public class HeartbeatProcessor {
private final EventMeshGrpcServer eventMeshGrpcServer;
- public HeartbeatProcessor(EventMeshGrpcServer eventMeshGrpcServer) {
+ private final Acl acl;
+
+ public HeartbeatProcessor(final EventMeshGrpcServer eventMeshGrpcServer) {
this.eventMeshGrpcServer = eventMeshGrpcServer;
+ this.acl = eventMeshGrpcServer.getAcl();
}
public void process(Heartbeat heartbeat, EventEmitter<Response> emitter)
throws Exception {
@@ -110,7 +113,7 @@ public class HeartbeatProcessor {
String sys = header.getSys();
int requestCode = RequestCode.HEARTBEAT.getRequestCode();
for (Heartbeat.HeartbeatItem item :
heartbeat.getHeartbeatItemsList()) {
- Acl.doAclCheckInHttpHeartbeat(remoteAdd, user, pass, sys,
item.getTopic(), requestCode);
+ this.acl.doAclCheckInHttpHeartbeat(remoteAdd, user, pass, sys,
item.getTopic(), requestCode);
}
}
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/ReplyMessageProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/ReplyMessageProcessor.java
index b7b8f62ec..aa34154b4 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/ReplyMessageProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/ReplyMessageProcessor.java
@@ -54,8 +54,11 @@ public class ReplyMessageProcessor {
private final EventMeshGrpcServer eventMeshGrpcServer;
- public ReplyMessageProcessor(EventMeshGrpcServer eventMeshGrpcServer) {
+ private final Acl acl;
+
+ public ReplyMessageProcessor(final EventMeshGrpcServer
eventMeshGrpcServer) {
this.eventMeshGrpcServer = eventMeshGrpcServer;
+ this.acl = eventMeshGrpcServer.getAcl();
}
public void process(SimpleMessage message, EventEmitter<SimpleMessage>
emitter) throws Exception {
@@ -134,7 +137,7 @@ public class ReplyMessageProcessor {
String pass = requestHeader.getPassword();
String subsystem = requestHeader.getSys();
String topic = message.getTopic();
- Acl.doAclCheckInHttpSend(remoteAdd, user, pass, subsystem, topic,
RequestCode.REPLY_MESSAGE.getRequestCode());
+ this.acl.doAclCheckInHttpSend(remoteAdd, user, pass, subsystem,
topic, RequestCode.REPLY_MESSAGE.getRequestCode());
}
}
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/RequestMessageProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/RequestMessageProcessor.java
index 9da3e1556..3f23c8d1d 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/RequestMessageProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/RequestMessageProcessor.java
@@ -52,8 +52,11 @@ public class RequestMessageProcessor {
private final EventMeshGrpcServer eventMeshGrpcServer;
- public RequestMessageProcessor(EventMeshGrpcServer eventMeshGrpcServer) {
+ private final Acl acl;
+
+ public RequestMessageProcessor(final EventMeshGrpcServer
eventMeshGrpcServer) {
this.eventMeshGrpcServer = eventMeshGrpcServer;
+ this.acl = eventMeshGrpcServer.getAcl();
}
public void process(SimpleMessage message, EventEmitter<SimpleMessage>
emitter) throws Exception {
@@ -145,7 +148,7 @@ public class RequestMessageProcessor {
String pass = requestHeader.getPassword();
String subsystem = requestHeader.getSys();
String topic = message.getTopic();
- Acl.doAclCheckInHttpSend(remoteAdd, user, pass, subsystem, topic,
RequestCode.MSG_SEND_ASYNC.getRequestCode());
+ this.acl.doAclCheckInHttpSend(remoteAdd, user, pass, subsystem,
topic, RequestCode.MSG_SEND_ASYNC.getRequestCode());
}
}
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SendAsyncMessageProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SendAsyncMessageProcessor.java
index 934f8c55b..c24d7e0ba 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SendAsyncMessageProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SendAsyncMessageProcessor.java
@@ -55,8 +55,11 @@ public class SendAsyncMessageProcessor {
private final EventMeshGrpcServer eventMeshGrpcServer;
- public SendAsyncMessageProcessor(EventMeshGrpcServer eventMeshGrpcServer) {
+ private final Acl acl;
+
+ public SendAsyncMessageProcessor(final EventMeshGrpcServer
eventMeshGrpcServer) {
this.eventMeshGrpcServer = eventMeshGrpcServer;
+ this.acl = eventMeshGrpcServer.getAcl();
}
public void process(SimpleMessage message, EventEmitter<Response> emitter)
throws Exception {
@@ -135,7 +138,7 @@ public class SendAsyncMessageProcessor {
String pass = requestHeader.getPassword();
String subsystem = requestHeader.getSys();
String topic = message.getTopic();
- Acl.doAclCheckInHttpSend(remoteAdd, user, pass, subsystem, topic,
RequestCode.MSG_SEND_ASYNC.getRequestCode());
+ this.acl.doAclCheckInHttpSend(remoteAdd, user, pass, subsystem,
topic, RequestCode.MSG_SEND_ASYNC.getRequestCode());
}
}
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SubscribeProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SubscribeProcessor.java
index 7c9290dae..75d1cfc02 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SubscribeProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SubscribeProcessor.java
@@ -47,8 +47,11 @@ public class SubscribeProcessor {
private final transient GrpcType grpcType = GrpcType.WEBHOOK;
+ private final Acl acl;
+
public SubscribeProcessor(final EventMeshGrpcServer eventMeshGrpcServer) {
this.eventMeshGrpcServer = eventMeshGrpcServer;
+ this.acl = eventMeshGrpcServer.getAcl();
}
public void process(final Subscription subscription, final
EventEmitter<Response> emitter) throws Exception {
@@ -131,9 +134,9 @@ public class SubscribeProcessor {
final RequestHeader header = subscription.getHeader();
if
(eventMeshGrpcServer.getEventMeshGrpcConfiguration().isEventMeshServerSecurityEnable())
{
for (final Subscription.SubscriptionItem item :
subscription.getSubscriptionItemsList()) {
- Acl.doAclCheckInHttpReceive(header.getIp(),
header.getUsername(), header.getPassword(),
+ this.acl.doAclCheckInHttpReceive(header.getIp(),
header.getUsername(), header.getPassword(),
header.getSys(), item.getTopic(),
RequestCode.SUBSCRIBE.getRequestCode());
}
}
}
-}
\ No newline at end of file
+}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SubscribeStreamProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SubscribeStreamProcessor.java
index fcaf78027..e96218532 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SubscribeStreamProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/processor/SubscribeStreamProcessor.java
@@ -49,8 +49,12 @@ public class SubscribeStreamProcessor {
private final GrpcType grpcType = GrpcType.STREAM;
- public SubscribeStreamProcessor(EventMeshGrpcServer eventMeshGrpcServer) {
+ private final Acl acl;
+
+ public SubscribeStreamProcessor(final EventMeshGrpcServer
eventMeshGrpcServer) {
this.eventMeshGrpcServer = eventMeshGrpcServer;
+ this.acl = eventMeshGrpcServer.getAcl();
+
}
public void process(Subscription subscription, EventEmitter<SimpleMessage>
emitter) throws Exception {
@@ -133,9 +137,8 @@ public class SubscribeStreamProcessor {
String pass = header.getPassword();
String subsystem = header.getSys();
for (Subscription.SubscriptionItem item :
subscription.getSubscriptionItemsList()) {
- Acl.doAclCheckInHttpReceive(remoteAdd, user, pass, subsystem,
item.getTopic(),
- RequestCode.SUBSCRIBE.getRequestCode());
+ this.acl.doAclCheckInHttpReceive(remoteAdd, user, pass,
subsystem, item.getTopic(), RequestCode.SUBSCRIBE.getRequestCode());
}
}
}
-}
\ No newline at end of file
+}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageProcessor.java
index 45a5a1c76..e2bdf7d36 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageProcessor.java
@@ -70,8 +70,11 @@ public class BatchSendMessageProcessor implements
HttpRequestProcessor {
private EventMeshHTTPServer eventMeshHTTPServer;
- public BatchSendMessageProcessor(EventMeshHTTPServer eventMeshHTTPServer) {
+ private final Acl acl;
+
+ public BatchSendMessageProcessor(final EventMeshHTTPServer
eventMeshHTTPServer) {
this.eventMeshHTTPServer = eventMeshHTTPServer;
+ this.acl = eventMeshHTTPServer.getAcl();
}
public Logger batchMessageLogger = LoggerFactory.getLogger("batchMessage");
@@ -232,7 +235,7 @@ public class BatchSendMessageProcessor implements
HttpRequestProcessor {
//do acl check
if
(eventMeshHTTPServer.getEventMeshHttpConfiguration().isEventMeshServerSecurityEnable())
{
try {
- Acl.doAclCheckInHttpSend(remoteAddr, user, pass,
subsystem, cloudEvent.getSubject(), requestCode);
+ this.acl.doAclCheckInHttpSend(remoteAddr, user, pass,
subsystem, cloudEvent.getSubject(), requestCode);
} catch (Exception e) {
//String errorMsg = String.format("CLIENT HAS NO
PERMISSION,send failed, topic:%s, subsys:%s, realIp:%s", topic, subsys, realIp);
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageV2Processor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageV2Processor.java
index cdcb9ebce..9da514d63 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageV2Processor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageV2Processor.java
@@ -64,8 +64,11 @@ public class BatchSendMessageV2Processor implements
HttpRequestProcessor {
private final EventMeshHTTPServer eventMeshHTTPServer;
+ private final Acl acl;
+
public BatchSendMessageV2Processor(EventMeshHTTPServer
eventMeshHTTPServer) {
this.eventMeshHTTPServer = eventMeshHTTPServer;
+ this.acl = eventMeshHTTPServer.getAcl();
}
public Logger batchMessageLogger = LoggerFactory.getLogger("batchMessage");
@@ -178,7 +181,7 @@ public class BatchSendMessageV2Processor implements
HttpRequestProcessor {
String pass = getExtension(event,
ProtocolKey.ClientInstanceKey.PASSWD);
String subsystem = getExtension(event,
ProtocolKey.ClientInstanceKey.SYS);
try {
- Acl.doAclCheckInHttpSend(remoteAddr, user, pass, subsystem,
topic, requestCode);
+ this.acl.doAclCheckInHttpSend(remoteAddr, user, pass,
subsystem, topic, requestCode);
} catch (Exception e) {
responseEventMeshCommand =
asyncContext.getRequest().createHttpCommandResponse(
sendMessageBatchV2ResponseHeader,
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/HeartBeatProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/HeartBeatProcessor.java
index eac0f42ac..ef43a8932 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/HeartBeatProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/HeartBeatProcessor.java
@@ -56,8 +56,11 @@ public class HeartBeatProcessor implements
HttpRequestProcessor {
private final transient EventMeshHTTPServer eventMeshHTTPServer;
+ private final Acl acl;
+
public HeartBeatProcessor(final EventMeshHTTPServer eventMeshHTTPServer) {
this.eventMeshHTTPServer = eventMeshHTTPServer;
+ this.acl = eventMeshHTTPServer.getAcl();
}
@Override
@@ -127,7 +130,7 @@ public class HeartBeatProcessor implements
HttpRequestProcessor {
//do acl check
if
(eventMeshHTTPServer.getEventMeshHttpConfiguration().isEventMeshServerSecurityEnable())
{
try {
- Acl.doAclCheckInHttpHeartbeat(
+ this.acl.doAclCheckInHttpHeartbeat(
RemotingHelper.parseChannelRemoteAddr(ctx.channel()),
heartbeatRequestHeader.getUsername(),
heartbeatRequestHeader.getPasswd(),
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/LocalSubscribeEventProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/LocalSubscribeEventProcessor.java
index 729acff6e..d9976ef29 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/LocalSubscribeEventProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/LocalSubscribeEventProcessor.java
@@ -53,8 +53,11 @@ import lombok.extern.slf4j.Slf4j;
@Slf4j
public class LocalSubscribeEventProcessor extends AbstractEventProcessor {
+ private final Acl acl;
+
public LocalSubscribeEventProcessor(final EventMeshHTTPServer
eventMeshHTTPServer) {
super(eventMeshHTTPServer);
+ this.acl = eventMeshHTTPServer.getAcl();
}
@Override
@@ -115,7 +118,7 @@ public class LocalSubscribeEventProcessor extends
AbstractEventProcessor {
if
(eventMeshHTTPServer.getEventMeshHttpConfiguration().isEventMeshServerSecurityEnable())
{
for (final SubscriptionItem item : subscriptionList) {
try {
-
Acl.doAclCheckInHttpReceive(RemotingHelper.parseChannelRemoteAddr(channel),
+
this.acl.doAclCheckInHttpReceive(RemotingHelper.parseChannelRemoteAddr(channel),
sysHeaderMap.get(ProtocolKey.ClientInstanceKey.USERNAME).toString(),
sysHeaderMap.get(ProtocolKey.ClientInstanceKey.PASSWD).toString(),
sysHeaderMap.get(ProtocolKey.ClientInstanceKey.SYS).toString(),
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/RemoteSubscribeEventProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/RemoteSubscribeEventProcessor.java
index 976ddb557..72683e660 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/RemoteSubscribeEventProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/RemoteSubscribeEventProcessor.java
@@ -61,9 +61,11 @@ public class RemoteSubscribeEventProcessor extends
AbstractEventProcessor {
public Logger aclLogger = LoggerFactory.getLogger(EventMeshConstants.ACL);
+ private final Acl acl;
public RemoteSubscribeEventProcessor(EventMeshHTTPServer
eventMeshHTTPServer) {
super(eventMeshHTTPServer);
+ this.acl = eventMeshHTTPServer.getAcl();
}
@Override
@@ -136,13 +138,10 @@ public class RemoteSubscribeEventProcessor extends
AbstractEventProcessor {
String subsystem =
sysHeaderMap.get(ProtocolKey.ClientInstanceKey.SYS).toString();
for (SubscriptionItem item : subscriptionList) {
try {
- Acl.doAclCheckInHttpReceive(remoteAddr, user, pass,
subsystem, item.getTopic(),
- requestWrapper.getRequestURI());
+ this.acl.doAclCheckInHttpReceive(remoteAddr, user, pass,
subsystem, item.getTopic(), requestWrapper.getRequestURI());
} catch (Exception e) {
aclLogger.warn("CLIENT HAS NO
PERMISSION,SubscribeProcessor subscribe failed", e);
-
handlerSpecific.sendErrorResponse(EventMeshRetCode.EVENTMESH_ACL_ERR,
responseHeaderMap,
- responseBodyMap, null);
-
+
handlerSpecific.sendErrorResponse(EventMeshRetCode.EVENTMESH_ACL_ERR,
responseHeaderMap, responseBodyMap, null);
return;
}
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncEventProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncEventProcessor.java
index b45cea715..6d1adc82d 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncEventProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncEventProcessor.java
@@ -61,11 +61,16 @@ import lombok.extern.slf4j.Slf4j;
@Slf4j
@EventMeshTrace(isEnable = true)
-@RequiredArgsConstructor
public class SendAsyncEventProcessor implements AsyncHttpProcessor {
private final transient EventMeshHTTPServer eventMeshHTTPServer;
+ private final Acl acl;
+
+ public SendAsyncEventProcessor(EventMeshHTTPServer eventMeshHTTPServer) {
+ this.eventMeshHTTPServer = eventMeshHTTPServer;
+ this.acl = eventMeshHTTPServer.getAcl();
+ }
@Override
public void handler(final HandlerService.HandlerSpecific handlerSpecific,
final HttpRequest httpRequest) throws Exception {
@@ -170,7 +175,7 @@ public class SendAsyncEventProcessor implements
AsyncHttpProcessor {
final String subsystem =
Objects.requireNonNull(event.getExtension(ProtocolKey.ClientInstanceKey.SYS)).toString();
final String requestURI = requestWrapper.getRequestURI();
try {
- Acl.doAclCheckInHttpSend(remoteAddr, user, pass, subsystem,
topic, requestURI);
+ this.acl.doAclCheckInHttpSend(remoteAddr, user, pass,
subsystem, topic, requestURI);
} catch (Exception e) {
handlerSpecific.sendErrorResponse(EventMeshRetCode.EVENTMESH_ACL_ERR,
responseHeaderMap,
responseBodyMap,
EventMeshUtil.getCloudEventExtensionMap(SpecVersion.V1.toString(), event));
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncMessageProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncMessageProcessor.java
index fc73c4baa..36d3e1cb5 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncMessageProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncMessageProcessor.java
@@ -62,7 +62,7 @@ import io.opentelemetry.api.trace.Span;
import lombok.RequiredArgsConstructor;
-@RequiredArgsConstructor
+
public class SendAsyncMessageProcessor implements HttpRequestProcessor {
public Logger messageLogger =
LoggerFactory.getLogger(EventMeshConstants.MESSAGE);
@@ -75,6 +75,13 @@ public class SendAsyncMessageProcessor implements
HttpRequestProcessor {
private final EventMeshHTTPServer eventMeshHTTPServer;
+ private final Acl acl;
+
+ public SendAsyncMessageProcessor(final EventMeshHTTPServer
eventMeshHTTPServer) {
+ this.eventMeshHTTPServer = eventMeshHTTPServer;
+ this.acl = eventMeshHTTPServer.getAcl();
+ }
+
@Override
public void processRequest(ChannelHandlerContext ctx,
AsyncContext<HttpCommand> asyncContext) throws Exception {
@@ -171,7 +178,7 @@ public class SendAsyncMessageProcessor implements
HttpRequestProcessor {
String subsystem =
Objects.requireNonNull(event.getExtension(ProtocolKey.ClientInstanceKey.SYS)).toString();
int requestCode = Integer.parseInt(request.getRequestCode());
try {
- Acl.doAclCheckInHttpSend(remoteAddr, user, pass, subsystem,
topic, requestCode);
+ this.acl.doAclCheckInHttpSend(remoteAddr, user, pass,
subsystem, topic, requestCode);
} catch (Exception e) {
responseEventMeshCommand = request.createHttpCommandResponse(
sendMessageResponseHeader,
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncRemoteEventProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncRemoteEventProcessor.java
index af4f25e42..930da401d 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncRemoteEventProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncRemoteEventProcessor.java
@@ -70,8 +70,11 @@ public class SendAsyncRemoteEventProcessor implements
AsyncHttpProcessor {
private final transient EventMeshHTTPServer eventMeshHTTPServer;
+ private final Acl acl;
+
public SendAsyncRemoteEventProcessor(final EventMeshHTTPServer
eventMeshHTTPServer) {
this.eventMeshHTTPServer = eventMeshHTTPServer;
+ this.acl = eventMeshHTTPServer.getAcl();
}
@Override
@@ -211,7 +214,7 @@ public class SendAsyncRemoteEventProcessor implements
AsyncHttpProcessor {
//do acl check
if
(eventMeshHTTPServer.getEventMeshHttpConfiguration().isEventMeshServerSecurityEnable())
{
try {
-
Acl.doAclCheckInHttpSend(RemotingHelper.parseChannelRemoteAddr(ctx.channel()),
+
this.acl.doAclCheckInHttpSend(RemotingHelper.parseChannelRemoteAddr(ctx.channel()),
getExtension(event,
ProtocolKey.ClientInstanceKey.USERNAME),
getExtension(event,
ProtocolKey.ClientInstanceKey.PASSWD),
getExtension(event, ProtocolKey.ClientInstanceKey.SYS),
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendSyncMessageProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendSyncMessageProcessor.java
index 7ad277a3b..093e6c207 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendSyncMessageProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendSyncMessageProcessor.java
@@ -58,8 +58,11 @@ import lombok.extern.slf4j.Slf4j;
public class SendSyncMessageProcessor implements HttpRequestProcessor {
private transient EventMeshHTTPServer eventMeshHTTPServer;
+ private final Acl acl;
+
public SendSyncMessageProcessor(final EventMeshHTTPServer
eventMeshHTTPServer) {
this.eventMeshHTTPServer = eventMeshHTTPServer;
+ this.acl = eventMeshHTTPServer.getAcl();
}
@Override
@@ -161,7 +164,7 @@ public class SendSyncMessageProcessor implements
HttpRequestProcessor {
final int requestCode =
Integer.parseInt(asyncContext.getRequest().getRequestCode());
try {
- Acl.doAclCheckInHttpSend(remoteAddr, user, pass, sys, topic,
requestCode);
+ this.acl.doAclCheckInHttpSend(remoteAddr, user, pass, sys,
topic, requestCode);
} catch (Exception e) {
responseEventMeshCommand =
asyncContext.getRequest().createHttpCommandResponse(
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
index fd085b507..5dbf9098d 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
@@ -56,8 +56,11 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
private final transient EventMeshHTTPServer eventMeshHTTPServer;
+ private final Acl acl;
+
public SubscribeProcessor(final EventMeshHTTPServer eventMeshHTTPServer) {
this.eventMeshHTTPServer = eventMeshHTTPServer;
+ this.acl = eventMeshHTTPServer.getAcl();
}
@Override
@@ -118,18 +121,15 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
if
(eventMeshHTTPServer.getEventMeshHttpConfiguration().isEventMeshServerSecurityEnable())
{
for (final SubscriptionItem item : subTopicList) {
try {
-
Acl.doAclCheckInHttpReceive(RemotingHelper.parseChannelRemoteAddr(ctx.channel()),
- subscribeRequestHeader.getUsername(),
- subscribeRequestHeader.getPasswd(),
- subscribeRequestHeader.getSys(), item.getTopic(),
- requestCode);
+
this.acl.doAclCheckInHttpReceive(RemotingHelper.parseChannelRemoteAddr(ctx.channel()),
+ subscribeRequestHeader.getUsername(),
+ subscribeRequestHeader.getPasswd(),
+ subscribeRequestHeader.getSys(), item.getTopic(),
+ requestCode);
} catch (Exception e) {
- responseEventMeshCommand =
asyncContext.getRequest().createHttpCommandResponse(
- subscribeResponseHeader,
- SendMessageResponseBody
-
.buildBody(EventMeshRetCode.EVENTMESH_ACL_ERR.getRetCode(),
- e.getMessage()));
+ responseEventMeshCommand =
asyncContext.getRequest().createHttpCommandResponse(subscribeResponseHeader,
+
SendMessageResponseBody.buildBody(EventMeshRetCode.EVENTMESH_ACL_ERR.getRetCode(),
e.getMessage()));
asyncContext.onComplete(responseEventMeshCommand);
if (LOGGER.isWarnEnabled()) {
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/HeartBeatTask.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/HeartBeatTask.java
index 461b23ea4..7d2b8fce9 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/HeartBeatTask.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/HeartBeatTask.java
@@ -39,8 +39,11 @@ public class HeartBeatTask extends AbstractTask {
private static final Logger LOGGER =
LoggerFactory.getLogger(HeartBeatTask.class);
+ private final Acl acl;
+
public HeartBeatTask(Package pkg, ChannelHandlerContext ctx, long
startTime, EventMeshTCPServer eventMeshTCPServer) {
super(pkg, ctx, startTime, eventMeshTCPServer);
+ this.acl = eventMeshTCPServer.getAcl();
}
@Override
@@ -51,7 +54,7 @@ public class HeartBeatTask extends AbstractTask {
//do acl check in heartbeat
if
(eventMeshTCPServer.getEventMeshTCPConfiguration().isEventMeshServerSecurityEnable())
{
String remoteAddr =
RemotingHelper.parseChannelRemoteAddr(ctx.channel());
- Acl.doAclCheckInTcpHeartbeat(remoteAddr, session.getClient(),
HEARTBEAT_REQUEST.getValue());
+ this.acl.doAclCheckInTcpHeartbeat(remoteAddr,
session.getClient(), HEARTBEAT_REQUEST.getValue());
}
if (session != null) {
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/HelloTask.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/HelloTask.java
index 6f5e4a9f1..1ea200801 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/HelloTask.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/HelloTask.java
@@ -49,8 +49,11 @@ public class HelloTask extends AbstractTask {
private static final Logger MESSAGE_LOGGER =
LoggerFactory.getLogger("message");
+ private final Acl acl;
+
public HelloTask(Package pkg, ChannelHandlerContext ctx, long startTime,
EventMeshTCPServer eventMeshTCPServer) {
super(pkg, ctx, startTime, eventMeshTCPServer);
+ this.acl = eventMeshTCPServer.getAcl();
}
@Override
@@ -63,7 +66,7 @@ public class HelloTask extends AbstractTask {
//do acl check in connect
if
(eventMeshTCPServer.getEventMeshTCPConfiguration().isEventMeshServerSecurityEnable())
{
String remoteAddr =
RemotingHelper.parseChannelRemoteAddr(ctx.channel());
- Acl.doAclCheckInTcpConnect(remoteAddr, user,
HELLO_REQUEST.getValue());
+ this.acl.doAclCheckInTcpConnect(remoteAddr, user,
HELLO_REQUEST.getValue());
}
if (eventMeshTCPServer.getEventMeshServer().getServiceState() !=
ServiceState.RUNNING) {
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/MessageTransferTask.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/MessageTransferTask.java
index 0c13db839..30f82ff02 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/MessageTransferTask.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/MessageTransferTask.java
@@ -67,9 +67,11 @@ public class MessageTransferTask extends AbstractTask {
private static final int TRY_PERMIT_TIME_OUT = 5;
- public MessageTransferTask(Package pkg, ChannelHandlerContext ctx, long
startTime,
- EventMeshTCPServer eventMeshTCPServer) {
+ private final Acl acl;
+
+ public MessageTransferTask(Package pkg, ChannelHandlerContext ctx, long
startTime, EventMeshTCPServer eventMeshTCPServer) {
super(pkg, ctx, startTime, eventMeshTCPServer);
+ this.acl = eventMeshTCPServer.getAcl();
}
@Override
@@ -79,11 +81,11 @@ public class MessageTransferTask extends AbstractTask {
try {
if
(eventMeshTCPServer.getEventMeshTCPConfiguration().isEventMeshServerTraceEnable()
- && RESPONSE_TO_SERVER != cmd) {
+ && RESPONSE_TO_SERVER != cmd) {
//attach the span to the server context
Span span =
TraceUtils.prepareServerSpan(pkg.getHeader().getProperties(),
-
EventMeshTraceConstants.TRACE_UPSTREAM_EVENTMESH_SERVER_SPAN,
- startTime, TimeUnit.MILLISECONDS, true);
+
EventMeshTraceConstants.TRACE_UPSTREAM_EVENTMESH_SERVER_SPAN,
+ startTime, TimeUnit.MILLISECONDS, true);
Context context = Context.current().with(SpanKey.SERVER_KEY,
span);
//put the context in channel
ctx.channel().attr(AttributeKeys.SERVER_CONTEXT).set(context);
@@ -92,7 +94,6 @@ public class MessageTransferTask extends AbstractTask {
LOGGER.warn("upload trace fail in
MessageTransferTask[server-span-start]", ex);
}
-
Command replyCmd = getReplyCmd(cmd);
Package msg = new Package();
@@ -102,11 +103,11 @@ public class MessageTransferTask extends AbstractTask {
try {
String protocolType = "eventmeshmessage";
if (pkg.getHeader().getProperties() != null
- && pkg.getHeader().getProperty(Constants.PROTOCOL_TYPE) !=
null) {
+ && pkg.getHeader().getProperty(Constants.PROTOCOL_TYPE) !=
null) {
protocolType = (String)
pkg.getHeader().getProperty(Constants.PROTOCOL_TYPE);
}
ProtocolAdaptor<ProtocolTransportObject> protocolAdaptor =
- ProtocolPluginFactory.getProtocolAdaptor(protocolType);
+ ProtocolPluginFactory.getProtocolAdaptor(protocolType);
event = protocolAdaptor.toCloudEvent(pkg);
if (event == null) {
@@ -122,32 +123,26 @@ public class MessageTransferTask extends AbstractTask {
//do acl check in sending msg
if
(eventMeshTCPServer.getEventMeshTCPConfiguration().isEventMeshServerSecurityEnable())
{
String remoteAddr =
RemotingHelper.parseChannelRemoteAddr(ctx.channel());
- Acl.doAclCheckInTcpSend(remoteAddr, session.getClient(),
event.getSubject(),
- cmd.getValue());
+ this.acl.doAclCheckInTcpSend(remoteAddr, session.getClient(),
event.getSubject(), cmd.getValue());
}
if (!eventMeshTCPServer.getRateLimiter()
- .tryAcquire(TRY_PERMIT_TIME_OUT, TimeUnit.MILLISECONDS)) {
+ .tryAcquire(TRY_PERMIT_TIME_OUT, TimeUnit.MILLISECONDS)) {
- msg.setHeader(new Header(replyCmd, OPStatus.FAIL.getCode(),
- "Tps overload, global flow control",
- pkg.getHeader().getSeq()));
+ msg.setHeader(new Header(replyCmd, OPStatus.FAIL.getCode(),
"Tps overload, global flow control", pkg.getHeader().getSeq()));
ctx.writeAndFlush(msg).addListener(
- new ChannelFutureListener() {
- @Override
- public void operationComplete(ChannelFuture
future) throws Exception {
- Utils.logSucceedMessageFlow(msg,
session.getClient(), startTime,
- taskExecuteTime);
- }
+ new ChannelFutureListener() {
+ @Override
+ public void operationComplete(ChannelFuture future)
throws Exception {
+ Utils.logSucceedMessageFlow(msg,
session.getClient(), startTime,
+ taskExecuteTime);
}
+ }
);
- TraceUtils.finishSpanWithException(ctx, event, "Tps overload,
global flow control",
- null);
+ TraceUtils.finishSpanWithException(ctx, event, "Tps overload,
global flow control", null);
- LOGGER.warn(
- "======Tps overload, global flow control, rate:{}!
PLEASE CHECK!========",
- eventMeshTCPServer.getRateLimiter().getRate());
+ LOGGER.warn("======Tps overload, global flow control, rate:{}!
PLEASE CHECK!========", eventMeshTCPServer.getRateLimiter().getRate());
return;
}
@@ -156,29 +151,29 @@ public class MessageTransferTask extends AbstractTask {
event = addTimestamp(event, cmd, sendTime);
sendStatus = session
- .upstreamMsg(pkg.getHeader(), event,
- createSendCallback(replyCmd, taskExecuteTime,
event),
- startTime, taskExecuteTime);
+ .upstreamMsg(pkg.getHeader(), event,
+ createSendCallback(replyCmd, taskExecuteTime, event),
+ startTime, taskExecuteTime);
if (StringUtils.equals(EventMeshTcpSendStatus.SUCCESS.name(),
- sendStatus.getSendStatus().name())) {
+ sendStatus.getSendStatus().name())) {
MESSAGE_LOGGER.info("pkg|eventMesh2mq|cmd={}|Msg={}|user={}|wait={}ms|cost={}ms",
- cmd, event,
- session.getClient(), taskExecuteTime - startTime,
sendTime - startTime);
+ cmd, event,
+ session.getClient(), taskExecuteTime - startTime,
sendTime - startTime);
} else {
throw new Exception(sendStatus.getDetail());
}
}
} catch (Exception e) {
LOGGER.error("MessageTransferTask failed|cmd={}|event={}|user={}",
cmd, event,
- session.getClient(),
- e);
+ session.getClient(),
+ e);
if (cmd != RESPONSE_TO_SERVER) {
msg.setHeader(
- new Header(replyCmd, OPStatus.FAIL.getCode(),
e.toString(),
- pkg.getHeader()
- .getSeq()));
+ new Header(replyCmd, OPStatus.FAIL.getCode(), e.toString(),
+ pkg.getHeader()
+ .getSeq()));
Utils.writeAndFlush(msg, startTime, taskExecuteTime,
session.getContext(), session);
if (event != null) {
@@ -191,22 +186,22 @@ public class MessageTransferTask extends AbstractTask {
private CloudEvent addTimestamp(CloudEvent event, Command cmd, long
sendTime) {
if (cmd == RESPONSE_TO_SERVER) {
event = CloudEventBuilder.from(event)
-
.withExtension(EventMeshConstants.RSP_C2EVENTMESH_TIMESTAMP,
- String.valueOf(startTime))
-
.withExtension(EventMeshConstants.RSP_EVENTMESH2MQ_TIMESTAMP,
- String.valueOf(sendTime))
- .withExtension(EventMeshConstants.RSP_SEND_EVENTMESH_IP,
-
eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshServerIp())
- .build();
+ .withExtension(EventMeshConstants.RSP_C2EVENTMESH_TIMESTAMP,
+ String.valueOf(startTime))
+ .withExtension(EventMeshConstants.RSP_EVENTMESH2MQ_TIMESTAMP,
+ String.valueOf(sendTime))
+ .withExtension(EventMeshConstants.RSP_SEND_EVENTMESH_IP,
+
eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshServerIp())
+ .build();
} else {
event = CloudEventBuilder.from(event)
-
.withExtension(EventMeshConstants.REQ_C2EVENTMESH_TIMESTAMP,
- String.valueOf(startTime))
-
.withExtension(EventMeshConstants.REQ_EVENTMESH2MQ_TIMESTAMP,
- String.valueOf(sendTime))
- .withExtension(EventMeshConstants.REQ_SEND_EVENTMESH_IP,
-
eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshServerIp())
- .build();
+ .withExtension(EventMeshConstants.REQ_C2EVENTMESH_TIMESTAMP,
+ String.valueOf(startTime))
+ .withExtension(EventMeshConstants.REQ_EVENTMESH2MQ_TIMESTAMP,
+ String.valueOf(sendTime))
+ .withExtension(EventMeshConstants.REQ_SEND_EVENTMESH_IP,
+
eventMeshTCPServer.getEventMeshTCPConfiguration().getEventMeshServerIp())
+ .build();
}
return event;
}
@@ -225,7 +220,7 @@ public class MessageTransferTask extends AbstractTask {
}
protected SendCallback createSendCallback(Command replyCmd, long
taskExecuteTime,
- CloudEvent event) {
+ CloudEvent event) {
final long createTime = System.currentTimeMillis();
Package msg = new Package();
@@ -234,16 +229,16 @@ public class MessageTransferTask extends AbstractTask {
public void onSuccess(SendResult sendResult) {
session.getSender().getUpstreamBuff().release();
MESSAGE_LOGGER.info("upstreamMsg message
success|user={}|callback cost={}",
- session.getClient(),
- System.currentTimeMillis() - createTime);
+ session.getClient(),
+ System.currentTimeMillis() - createTime);
if (replyCmd == Command.BROADCAST_MESSAGE_TO_SERVER_ACK
- || replyCmd == Command.ASYNC_MESSAGE_TO_SERVER_ACK) {
+ || replyCmd == Command.ASYNC_MESSAGE_TO_SERVER_ACK) {
msg.setHeader(
- new Header(replyCmd, OPStatus.SUCCESS.getCode(),
OPStatus.SUCCESS.getDesc(),
- pkg.getHeader().getSeq()));
+ new Header(replyCmd, OPStatus.SUCCESS.getCode(),
OPStatus.SUCCESS.getDesc(),
+ pkg.getHeader().getSeq()));
msg.setBody(event);
Utils.writeAndFlush(msg, startTime, taskExecuteTime,
session.getContext(),
- session);
+ session);
//async request need finish span when callback, rr request
will finish span when rrCallback
TraceUtils.finishSpan(ctx, event);
@@ -256,21 +251,21 @@ public class MessageTransferTask extends AbstractTask {
// retry
UpStreamMsgContext upStreamMsgContext = new UpStreamMsgContext(
- session, event, pkg.getHeader(), startTime,
taskExecuteTime);
+ session, event, pkg.getHeader(), startTime,
taskExecuteTime);
upStreamMsgContext.delay(10000);
Objects.requireNonNull(
-
session.getClientGroupWrapper().get()).getEventMeshTcpRetryer()
- .pushRetry(upStreamMsgContext);
+
session.getClientGroupWrapper().get()).getEventMeshTcpRetryer()
+ .pushRetry(upStreamMsgContext);
session.getSender().failMsgCount.incrementAndGet();
MESSAGE_LOGGER
- .error("upstreamMsg mq message error|user={}|callback
cost={}, errMsg={}",
- session.getClient(),
- (System.currentTimeMillis() - createTime),
- new Exception(context.getException()));
+ .error("upstreamMsg mq message error|user={}|callback
cost={}, errMsg={}",
+ session.getClient(),
+ (System.currentTimeMillis() - createTime),
+ new Exception(context.getException()));
msg.setHeader(
- new Header(replyCmd, OPStatus.FAIL.getCode(),
context.getException().toString(),
- pkg.getHeader().getSeq()));
+ new Header(replyCmd, OPStatus.FAIL.getCode(),
context.getException().toString(),
+ pkg.getHeader().getSeq()));
msg.setBody(event);
Utils.writeAndFlush(msg, startTime, taskExecuteTime,
session.getContext(), session);
@@ -278,8 +273,8 @@ public class MessageTransferTask extends AbstractTask {
if (replyCmd != RESPONSE_TO_SERVER) {
//upload trace
TraceUtils.finishSpanWithException(ctx, event,
- "upload trace fail in
MessageTransferTask.createSendCallback.onException",
- context.getException());
+ "upload trace fail in
MessageTransferTask.createSendCallback.onException",
+ context.getException());
}
}
};
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/SubscribeTask.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/SubscribeTask.java
index e3fbc927b..a6dd62b42 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/SubscribeTask.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/SubscribeTask.java
@@ -41,8 +41,11 @@ public class SubscribeTask extends AbstractTask {
private static final Logger LOGGER =
LoggerFactory.getLogger(SubscribeTask.class);
+ private final Acl acl;
+
public SubscribeTask(final Package pkg, final ChannelHandlerContext ctx,
long startTime, final EventMeshTCPServer eventMeshTCPServer) {
super(pkg, ctx, startTime, eventMeshTCPServer);
+ this.acl = eventMeshTCPServer.getAcl();
}
@Override
@@ -60,8 +63,7 @@ public class SubscribeTask extends AbstractTask {
subscriptionInfo.getTopicList().forEach(item -> {
//do acl check for receive msg
if (eventMeshServerSecurityEnable) {
- Acl.doAclCheckInTcpReceive(remoteAddr,
session.getClient(), item.getTopic(),
- Command.SUBSCRIBE_REQUEST.getValue());
+ this.acl.doAclCheckInTcpReceive(remoteAddr,
session.getClient(), item.getTopic(), Command.SUBSCRIBE_REQUEST.getValue());
}
subscriptionItems.add(item);
@@ -73,12 +75,10 @@ public class SubscribeTask extends AbstractTask {
LOGGER.info("SubscribeTask succeed|user={}|topics={}",
session.getClient(), subscriptionItems);
}
}
- msg.setHeader(new Header(Command.SUBSCRIBE_RESPONSE,
OPStatus.SUCCESS.getCode(), OPStatus.SUCCESS.getDesc(),
- pkg.getHeader().getSeq()));
+ msg.setHeader(new Header(Command.SUBSCRIBE_RESPONSE,
OPStatus.SUCCESS.getCode(), OPStatus.SUCCESS.getDesc(),
pkg.getHeader().getSeq()));
} catch (Exception e) {
LOGGER.error("SubscribeTask failed|user={}|errMsg={}",
session.getClient(), e);
- msg.setHeader(new Header(Command.SUBSCRIBE_RESPONSE,
OPStatus.FAIL.getCode(), e.toString(), pkg.getHeader()
- .getSeq()));
+ msg.setHeader(new Header(Command.SUBSCRIBE_RESPONSE,
OPStatus.FAIL.getCode(), e.toString(), pkg.getHeader().getSeq()));
} finally {
Utils.writeAndFlush(msg, startTime, taskExecuteTime,
session.getContext(), session);
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]