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 dc2b9e45c [ISSUE #3292]Modify EventMeshHTTPServer class attributes
public to private
new 5352c59c9 Merge pull request #3293 from mxsm/eventmesh-3292
dc2b9e45c is described below
commit dc2b9e45c1888b1cbc5c09a84009a7f2c14e982c
Author: mxsm <[email protected]>
AuthorDate: Tue Feb 28 00:11:58 2023 +0800
[ISSUE #3292]Modify EventMeshHTTPServer class attributes public to private
---
.../runtime/boot/EventMeshHTTPServer.java | 45 +++++----
.../processor/LocalSubscribeEventProcessor.java | 2 +-
.../processor/RemoteSubscribeEventProcessor.java | 4 +-
.../processor/RemoteUnSubscribeEventProcessor.java | 2 +-
.../http/processor/SubscribeProcessor.java | 2 +-
.../protocol/http/push/AsyncHTTPPushRequest.java | 101 ++++++++++-----------
.../protocol/http/push/HTTPMessageHandler.java | 2 +-
.../runtime/metrics/http/HTTPMetricsServer.java | 6 +-
8 files changed, 78 insertions(+), 86 deletions(-)
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 11d3d61ea..5036f6660 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
@@ -27,7 +27,6 @@ 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;
import org.apache.eventmesh.runtime.core.consumer.SubscriptionManager;
@@ -74,47 +73,45 @@ import lombok.extern.slf4j.Slf4j;
@Slf4j
public class EventMeshHTTPServer extends AbstractHTTPServer {
- private final transient EventMeshServer eventMeshServer;
+ private final EventMeshServer eventMeshServer;
- public transient ServiceState serviceState;
+ private final EventMeshHTTPConfiguration eventMeshHttpConfiguration;
- private final transient EventMeshHTTPConfiguration
eventMeshHttpConfiguration;
-
- private final transient Registry registry;
+ private final Registry registry;
private final Acl acl;
- public final transient EventBus eventBus = new EventBus();
+ public final EventBus eventBus = new EventBus();
- private transient ConsumerManager consumerManager;
+ private ConsumerManager consumerManager;
- private transient SubscriptionManager subscriptionManager;
+ private SubscriptionManager subscriptionManager;
- private transient ProducerManager producerManager;
+ private ProducerManager producerManager;
- private transient HttpRetryer httpRetryer;
+ private HttpRetryer httpRetryer;
- public transient ThreadPoolExecutor batchMsgExecutor;
+ private ThreadPoolExecutor batchMsgExecutor;
- public transient ThreadPoolExecutor sendMsgExecutor;
+ private ThreadPoolExecutor sendMsgExecutor;
- public transient ThreadPoolExecutor remoteMsgExecutor;
+ private ThreadPoolExecutor remoteMsgExecutor;
- public transient ThreadPoolExecutor replyMsgExecutor;
+ private ThreadPoolExecutor replyMsgExecutor;
- public transient ThreadPoolExecutor pushMsgExecutor;
+ private ThreadPoolExecutor pushMsgExecutor;
- public transient ThreadPoolExecutor clientManageExecutor;
+ private ThreadPoolExecutor clientManageExecutor;
- public transient ThreadPoolExecutor adminExecutor;
+ private ThreadPoolExecutor adminExecutor;
- public ThreadPoolExecutor webhookExecutor;
+ private ThreadPoolExecutor webhookExecutor;
private transient RateLimiter msgRateLimiter;
private transient RateLimiter batchRateLimiter;
- public transient HTTPClientPool httpClientPool = new HTTPClientPool(10);
+ private final transient HTTPClientPool httpClientPool = new
HTTPClientPool(10);
public EventMeshHTTPServer(final EventMeshServer eventMeshServer, final
EventMeshHTTPConfiguration eventMeshHttpConfiguration) {
@@ -394,10 +391,6 @@ public class EventMeshHTTPServer extends
AbstractHTTPServer {
return producerManager;
}
- public ServiceState getServiceState() {
- return serviceState;
- }
-
public EventMeshHTTPConfiguration getEventMeshHttpConfiguration() {
return eventMeshHttpConfiguration;
}
@@ -453,4 +446,8 @@ public class EventMeshHTTPServer extends AbstractHTTPServer
{
public Registry getRegistry() {
return registry;
}
+
+ public HTTPClientPool getHttpClientPool() {
+ return httpClientPool;
+ }
}
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 04f98eb8e..b761a2873 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
@@ -159,7 +159,7 @@ public class LocalSubscribeEventProcessor extends
AbstractEventProcessor {
}
// obtain webhook delivery agreement for Abuse Protection
- if
(!WebhookUtil.obtainDeliveryAgreement(eventMeshHTTPServer.httpClientPool.getClient(),
+ if
(!WebhookUtil.obtainDeliveryAgreement(eventMeshHTTPServer.getHttpClientPool().getClient(),
url,
eventMeshHTTPServer.getEventMeshHttpConfiguration().getEventMeshWebhookOrigin()))
{
if (log.isErrorEnabled()) {
log.error("subscriber url {} is not allowed by the target
system", url);
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 b94862c95..05601774e 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
@@ -163,7 +163,7 @@ public class RemoteSubscribeEventProcessor extends
AbstractEventProcessor {
}
// obtain webhook delivery agreement for Abuse Protection
- boolean isWebhookAllowed =
WebhookUtil.obtainDeliveryAgreement(eventMeshHTTPServer.httpClientPool.getClient(),
+ boolean isWebhookAllowed =
WebhookUtil.obtainDeliveryAgreement(eventMeshHTTPServer.getHttpClientPool().getClient(),
url, eventMeshHttpConfiguration.getEventMeshWebhookOrigin());
if (!isWebhookAllowed) {
@@ -192,7 +192,7 @@ public class RemoteSubscribeEventProcessor extends
AbstractEventProcessor {
targetMesh = meshAddress;
}
- CloseableHttpClient closeableHttpClient =
eventMeshHTTPServer.httpClientPool.getClient();
+ CloseableHttpClient closeableHttpClient =
eventMeshHTTPServer.getHttpClientPool().getClient();
String remoteResult = post(closeableHttpClient, targetMesh,
builderRemoteHeaderMap(localAddress), remoteBodyMap,
response -> EntityUtils.toString(response.getEntity(),
Constants.DEFAULT_CHARSET));
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/RemoteUnSubscribeEventProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/RemoteUnSubscribeEventProcessor.java
index 7bbe5f482..5139bee38 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/RemoteUnSubscribeEventProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/RemoteUnSubscribeEventProcessor.java
@@ -161,7 +161,7 @@ public class RemoteUnSubscribeEventProcessor extends
AbstractEventProcessor {
targetMesh = meshAddress;
}
- CloseableHttpClient closeableHttpClient =
eventMeshHTTPServer.httpClientPool.getClient();
+ CloseableHttpClient closeableHttpClient =
eventMeshHTTPServer.getHttpClientPool().getClient();
String remoteResult = post(closeableHttpClient, targetMesh,
builderRemoteHeaderMap(localAddress), remoteBodyMap,
response -> EntityUtils.toString(response.getEntity(),
Constants.DEFAULT_CHARSET));
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 ef99e4f77..e0f816dac 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
@@ -169,7 +169,7 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
}
// obtain webhook delivery agreement for Abuse Protection
- if
(!WebhookUtil.obtainDeliveryAgreement(eventMeshHTTPServer.httpClientPool.getClient(),
+ if
(!WebhookUtil.obtainDeliveryAgreement(eventMeshHTTPServer.getHttpClientPool().getClient(),
url,
eventMeshHTTPServer.getEventMeshHttpConfiguration().getEventMeshWebhookOrigin()))
{
log.error("subscriber url {} is not allowed by the target system",
url);
responseEventMeshCommand = request.createHttpCommandResponse(
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/AsyncHTTPPushRequest.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/AsyncHTTPPushRequest.java
index 1270d162d..ab6621a09 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/AsyncHTTPPushRequest.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/AsyncHTTPPushRequest.java
@@ -44,7 +44,6 @@ import org.apache.http.HttpEntity;
import org.apache.http.HttpResponse;
import org.apache.http.HttpStatus;
import org.apache.http.NameValuePair;
-import org.apache.http.client.ResponseHandler;
import org.apache.http.client.entity.UrlEncodedFormEntity;
import org.apache.http.client.methods.HttpPost;
import org.apache.http.message.BasicNameValuePair;
@@ -200,67 +199,63 @@ public class AsyncHTTPPushRequest extends
AbstractHTTPPushRequest {
}
try {
- eventMeshHTTPServer.httpClientPool.getClient().execute(builder,
new ResponseHandler<Object>() {
- @Override
- public Object handleResponse(HttpResponse response) {
- removeWaitingMap(AsyncHTTPPushRequest.this);
- long cost = System.currentTimeMillis() - lastPushTime;
-
eventMeshHTTPServer.getMetrics().getSummaryMetrics().recordHTTPPushTimeCost(cost);
-
- if
(processResponseStatus(response.getStatusLine().getStatusCode(), response)) {
- // this is successful response, process response
payload
- String res = "";
- try {
- res = EntityUtils.toString(response.getEntity(),
-
Charset.forName(EventMeshConstants.DEFAULT_CHARSET));
- } catch (IOException e) {
+
eventMeshHTTPServer.getHttpClientPool().getClient().execute(builder, response
-> {
+ removeWaitingMap(AsyncHTTPPushRequest.this);
+ long cost = System.currentTimeMillis() - lastPushTime;
+
eventMeshHTTPServer.getMetrics().getSummaryMetrics().recordHTTPPushTimeCost(cost);
+
+ if
(processResponseStatus(response.getStatusLine().getStatusCode(), response)) {
+ // this is successful response, process response payload
+ String res;
+ try {
+ res = EntityUtils.toString(response.getEntity(),
Charset.forName(EventMeshConstants.DEFAULT_CHARSET));
+ } catch (IOException e) {
+ handleMsgContext.finish();
+ return new Object();
+ }
+ ClientRetCode result = processResponseContent(res);
+ if (MESSAGE_LOGGER.isInfoEnabled()) {
+ MESSAGE_LOGGER.info(
+
"message|eventMesh2client|{}|url={}|topic={}|bizSeqNo={}"
+ + "|uniqueId={}|cost={}",
+ result, currPushUrl, handleMsgContext.getTopic(),
+ handleMsgContext.getBizSeqNo(),
handleMsgContext.getUniqueId(), cost);
+ }
+ if (result == ClientRetCode.OK || result ==
ClientRetCode.REMOTE_OK) {
+ complete();
+ if (isComplete()) {
handleMsgContext.finish();
- return new Object();
}
- ClientRetCode result = processResponseContent(res);
- if (MESSAGE_LOGGER.isInfoEnabled()) {
- MESSAGE_LOGGER.info(
-
"message|eventMesh2client|{}|url={}|topic={}|bizSeqNo={}"
- + "|uniqueId={}|cost={}",
- result, currPushUrl,
handleMsgContext.getTopic(),
- handleMsgContext.getBizSeqNo(),
handleMsgContext.getUniqueId(), cost);
- }
- if (result == ClientRetCode.OK || result ==
ClientRetCode.REMOTE_OK) {
- complete();
- if (isComplete()) {
- handleMsgContext.finish();
- }
- } else if (result == ClientRetCode.RETRY) {
- delayRetry();
- if (isComplete()) {
- handleMsgContext.finish();
- }
- } else if (result == ClientRetCode.NOLISTEN) {
- delayRetry();
- if (isComplete()) {
- handleMsgContext.finish();
- }
- } else if (result == ClientRetCode.FAIL) {
- complete();
- if (isComplete()) {
- handleMsgContext.finish();
- }
+ } else if (result == ClientRetCode.RETRY) {
+ delayRetry();
+ if (isComplete()) {
+ handleMsgContext.finish();
}
- } else {
-
eventMeshHTTPServer.getMetrics().getSummaryMetrics().recordHttpPushMsgFailed();
- if (MESSAGE_LOGGER.isInfoEnabled()) {
- MESSAGE_LOGGER.info(
-
"message|eventMesh2client|exception|url={}|topic={}|bizSeqNo={}"
- + "|uniqueId={}|cost={}", currPushUrl,
handleMsgContext.getTopic(),
- handleMsgContext.getBizSeqNo(),
handleMsgContext.getUniqueId(), cost);
+ } else if (result == ClientRetCode.NOLISTEN) {
+ delayRetry();
+ if (isComplete()) {
+ handleMsgContext.finish();
}
-
+ } else if (result == ClientRetCode.FAIL) {
+ complete();
if (isComplete()) {
handleMsgContext.finish();
}
}
- return new Object();
+ } else {
+
eventMeshHTTPServer.getMetrics().getSummaryMetrics().recordHttpPushMsgFailed();
+ if (MESSAGE_LOGGER.isInfoEnabled()) {
+ MESSAGE_LOGGER.info(
+
"message|eventMesh2client|exception|url={}|topic={}|bizSeqNo={}"
+ + "|uniqueId={}|cost={}", currPushUrl,
handleMsgContext.getTopic(),
+ handleMsgContext.getBizSeqNo(),
handleMsgContext.getUniqueId(), cost);
+ }
+
+ if (isComplete()) {
+ handleMsgContext.finish();
+ }
}
+ return new Object();
});
if (MESSAGE_LOGGER.isDebugEnabled()) {
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/HTTPMessageHandler.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/HTTPMessageHandler.java
index 8c86902c2..cad2c394f 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/HTTPMessageHandler.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/HTTPMessageHandler.java
@@ -68,7 +68,7 @@ public class HTTPMessageHandler implements MessageHandler {
public HTTPMessageHandler(EventMeshConsumer eventMeshConsumer) {
this.eventMeshConsumer = eventMeshConsumer;
- this.pushExecutor =
eventMeshConsumer.getEventMeshHTTPServer().pushMsgExecutor;
+ this.pushExecutor =
eventMeshConsumer.getEventMeshHTTPServer().getPushMsgExecutor();
waitingRequests.put(this.eventMeshConsumer.getConsumerGroupConf().getConsumerGroup(),
Sets.newConcurrentHashSet());
SCHEDULER.scheduleAtFixedRate(this::checkTimeout, 0, 1000,
TimeUnit.MILLISECONDS);
}
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/metrics/http/HTTPMetricsServer.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/metrics/http/HTTPMetricsServer.java
index 32655f24d..414f31778 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/metrics/http/HTTPMetricsServer.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/metrics/http/HTTPMetricsServer.java
@@ -48,9 +48,9 @@ public class HTTPMetricsServer {
this.eventMeshHTTPServer = eventMeshHTTPServer;
this.metricsRegistries = metricsRegistries;
this.summaryMetrics = new HttpSummaryMetrics(
- eventMeshHTTPServer.batchMsgExecutor,
- eventMeshHTTPServer.sendMsgExecutor,
- eventMeshHTTPServer.pushMsgExecutor,
+ eventMeshHTTPServer.getBatchMsgExecutor(),
+ eventMeshHTTPServer.getSendMsgExecutor(),
+ eventMeshHTTPServer.getPushMsgExecutor(),
eventMeshHTTPServer.getHttpRetryer().getFailedQueue());
init();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]