This is an automated email from the ASF dual-hosted git repository.
mikexue pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git
The following commit(s) were added to refs/heads/develop by this push:
new 19af7cf [ISSUE #490] Support service invocation (#508)
19af7cf is described below
commit 19af7cf134bedf2614efd396dc7abda51394d02d
Author: mike_xwm <[email protected]>
AuthorDate: Wed Aug 25 19:45:57 2021 +0800
[ISSUE #490] Support service invocation (#508)
* [ISSUE #490]Support service invocation
* remove unused code
---
.../eventmesh/api/producer/MeshMQProducer.java | 4 +-
.../connector/rocketmq/producer/ProducerImpl.java | 44 ++++++++++++++++++++
.../rocketmq/producer/RocketMQProducerImpl.java | 15 ++++---
.../connector/rocketmq/utils/OMSUtil.java | 8 +++-
.../runtime/core/plugin/MQProducerWrapper.java | 8 +---
.../http/processor/SendSyncMessageProcessor.java | 48 +---------------------
.../protocol/http/producer/EventMeshProducer.java | 8 +---
.../tcp/client/group/ClientGroupWrapper.java | 4 +-
.../tcp/client/session/send/SessionSender.java | 4 +-
9 files changed, 68 insertions(+), 75 deletions(-)
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/MeshMQProducer.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/MeshMQProducer.java
index 93fd3d3..ef0d2f5 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/MeshMQProducer.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/MeshMQProducer.java
@@ -33,9 +33,7 @@ public interface MeshMQProducer extends Producer {
void send(Message message, SendCallback sendCallback) throws Exception;
- void request(Message message, SendCallback sendCallback, RRCallback
rrCallback, long timeout) throws Exception;
-
- Message request(Message message, long timeout) throws Exception;
+ void request(Message message, RRCallback rrCallback, long timeout) throws
Exception;
boolean reply(final Message message, final SendCallback sendCallback)
throws Exception;
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/ProducerImpl.java
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/ProducerImpl.java
index f8d4302..c34a991 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/ProducerImpl.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/ProducerImpl.java
@@ -20,6 +20,7 @@ package org.apache.eventmesh.connector.rocketmq.producer;
import java.util.Properties;
import java.util.concurrent.ExecutorService;
+import com.alibaba.fastjson.JSONObject;
import io.openmessaging.api.Message;
import io.openmessaging.api.MessageBuilder;
import io.openmessaging.api.OnExceptionContext;
@@ -28,14 +29,26 @@ import io.openmessaging.api.SendCallback;
import io.openmessaging.api.SendResult;
import io.openmessaging.api.exception.OMSRuntimeException;
+import org.apache.eventmesh.api.RRCallback;
import org.apache.eventmesh.connector.rocketmq.utils.OMSUtil;
+import org.apache.rocketmq.client.exception.MQBrokerException;
+import org.apache.rocketmq.client.exception.MQClientException;
+import org.apache.rocketmq.client.producer.RequestCallback;
+import org.apache.rocketmq.client.producer.RequestResponseFuture;
+import org.apache.rocketmq.client.utils.MessageUtil;
import org.apache.rocketmq.common.message.MessageClientIDSetter;
+import org.apache.rocketmq.common.message.MessageExt;
+import org.apache.rocketmq.remoting.exception.RemotingException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
public class ProducerImpl extends AbstractOMSProducer implements Producer {
public static final int eventMeshServerAsyncAccumulationThreshold = 1000;
+ private final Logger logger = LoggerFactory.getLogger(ProducerImpl.class);
+
public ProducerImpl(final Properties properties) {
super(properties);
}
@@ -100,6 +113,37 @@ public class ProducerImpl extends AbstractOMSProducer
implements Producer {
}
}
+ public void request(Message message, RRCallback rrCallback, long timeout)
+ throws InterruptedException, RemotingException, MQClientException,
MQBrokerException {
+
+
this.checkProducerServiceState(this.rocketmqProducer.getDefaultMQProducerImpl());
+ org.apache.rocketmq.common.message.Message msgRMQ =
OMSUtil.msgConvert(message);
+ rocketmqProducer.request(msgRMQ, rrCallbackConvert(message,
rrCallback), timeout);
+ }
+
+ private RequestCallback rrCallbackConvert(final Message message, final
RRCallback rrCallback) {
+ return new RequestCallback() {
+ @Override
+ public void onSuccess(org.apache.rocketmq.common.message.Message
message) {
+ Message openMessage = OMSUtil.msgConvert((MessageExt) message);
+ rrCallback.onSuccess(openMessage);
+ }
+
+ @Override
+ public void onException(Throwable e) {
+ String topic = message.getTopic();
+ String msgId = message.getMsgID();
+ OMSRuntimeException onsEx =
ProducerImpl.this.checkProducerException(topic, msgId, e);
+ OnExceptionContext context = new OnExceptionContext();
+ context.setTopic(topic);
+ context.setMessageId(msgId);
+ context.setException(onsEx);
+ rrCallback.onException(e);
+
+ }
+ };
+ }
+
private org.apache.rocketmq.client.producer.SendCallback
sendCallbackConvert(final Message message, final SendCallback sendCallback) {
org.apache.rocketmq.client.producer.SendCallback rmqSendCallback = new
org.apache.rocketmq.client.producer.SendCallback() {
@Override
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/RocketMQProducerImpl.java
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/RocketMQProducerImpl.java
index 6769a0e..eba5208 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/RocketMQProducerImpl.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/RocketMQProducerImpl.java
@@ -36,6 +36,8 @@ import
org.apache.eventmesh.connector.rocketmq.config.ClientConfiguration;
import org.apache.eventmesh.connector.rocketmq.config.ConfigurationWrapper;
import org.apache.rocketmq.client.exception.MQBrokerException;
import org.apache.rocketmq.client.exception.MQClientException;
+import org.apache.rocketmq.common.MixAll;
+import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.remoting.exception.RemotingException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -98,19 +100,16 @@ public class RocketMQProducerImpl implements
MeshMQProducer {
}
@Override
- public void request(Message message, SendCallback sendCallback, RRCallback
rrCallback, long timeout)
+ public void request(Message message, RRCallback rrCallback, long timeout)
throws InterruptedException, RemotingException, MQClientException,
MQBrokerException {
- throw new UnsupportedOperationException("not support request-reply
mode when eventstore=rocketmq");
- }
-
- @Override
- public Message request(Message message, long timeout) throws
InterruptedException, RemotingException, MQClientException, MQBrokerException {
- throw new UnsupportedOperationException("not support request-reply
mode when eventstore=rocketmq");
+ producer.request(message, rrCallback, timeout);
}
@Override
public boolean reply(final Message message, final SendCallback
sendCallback) throws Exception {
- throw new UnsupportedOperationException("not support request-reply
mode when eventstore=rocketmq");
+ message.putSystemProperties(MessageConst.PROPERTY_MESSAGE_TYPE,
MixAll.REPLY_MESSAGE_FLAG);
+ producer.sendAsync(message, sendCallback);
+ return true;
}
@Override
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/utils/OMSUtil.java
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/utils/OMSUtil.java
index 906be5f..00bcb85 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/utils/OMSUtil.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/utils/OMSUtil.java
@@ -136,9 +136,13 @@ public class OMSUtil {
}
}
- systemProperties.put(Constants.PROPERTY_MESSAGE_MESSAGE_ID,
rmqMsg.getMsgId());
+ if (rmqMsg.getMsgId() != null){
+ systemProperties.put(Constants.PROPERTY_MESSAGE_MESSAGE_ID,
rmqMsg.getMsgId());
+ }
- systemProperties.put(Constants.PROPERTY_MESSAGE_DESTINATION,
rmqMsg.getTopic());
+ if (rmqMsg.getTopic() != null){
+ systemProperties.put(Constants.PROPERTY_MESSAGE_DESTINATION,
rmqMsg.getTopic());
+ }
// omsMsg.putSysHeaders(BuiltinKeys.SEARCH_KEYS, rmqMsg.getKeys());
systemProperties.put(Constants.PROPERTY_MESSAGE_BORN_HOST,
String.valueOf(rmqMsg.getBornHost()));
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 f979dc4..1c519af 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
@@ -80,13 +80,9 @@ public class MQProducerWrapper extends MQWrapper {
meshMQProducer.send(message, sendCallback);
}
- public void request(Message message, SendCallback sendCallback, RRCallback
rrCallback, long timeout)
+ public void request(Message message, RRCallback rrCallback, long timeout)
throws Exception {
- meshMQProducer.request(message, sendCallback, rrCallback, timeout);
- }
-
- public Message request(Message message, long timeout) throws Exception {
- return meshMQProducer.request(message, timeout);
+ meshMQProducer.request(message, rrCallback, timeout);
}
public boolean reply(final Message message, final SendCallback
sendCallback) throws Exception {
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 136961a..edb1a58 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
@@ -205,53 +205,7 @@ public class SendSyncMessageProcessor implements
HttpRequestProcessor {
.setProp(sendMessageRequestBody.getExtFields());
try {
- eventMeshProducer.request(sendMessageContext, new SendCallback() {
- @Override
- public void onSuccess(SendResult sendResult) {
- long endTime = System.currentTimeMillis();
-
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgCost(endTime -
startTime);
-
messageLogger.info("message|eventMesh2mq|REQ|SYNC|send2MQCost={}ms|topic={}|bizSeqNo={}|uniqueId={}",
- endTime - startTime,
- sendMessageRequestBody.getTopic(),
- sendMessageRequestBody.getBizSeqNo(),
- sendMessageRequestBody.getUniqueId());
- }
-
- @Override
- public void onException(OnExceptionContext context) {
- HttpCommand err =
asyncContext.getRequest().createHttpCommandResponse(
- sendMessageResponseHeader,
-
SendMessageResponseBody.buildBody(EventMeshRetCode.EVENTMESH_SEND_SYNC_MSG_ERR.getRetCode(),
-
EventMeshRetCode.EVENTMESH_SEND_SYNC_MSG_ERR.getErrMsg() +
EventMeshUtil.stackTrace(context.getException(), 2)));
- asyncContext.onComplete(err, handler);
- long endTime = System.currentTimeMillis();
-
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgFailed();
-
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgCost(endTime -
startTime);
-
messageLogger.error("message|eventMesh2mq|REQ|SYNC|send2MQCost={}ms|topic={}|bizSeqNo={}|uniqueId={}",
- endTime - startTime,
- sendMessageRequestBody.getTopic(),
- sendMessageRequestBody.getBizSeqNo(),
- sendMessageRequestBody.getUniqueId(),
context.getException());
- }
-// }
-//
-// @Override
-// public void onException(Throwable e) {
-// HttpCommand err =
asyncContext.getRequest().createHttpCommandResponse(
-// sendMessageResponseHeader,
-//
SendMessageResponseBody.buildBody(EventMeshRetCode.EVENTMESH_SEND_SYNC_MSG_ERR.getRetCode(),
-//
EventMeshRetCode.EVENTMESH_SEND_SYNC_MSG_ERR.getErrMsg() +
EventMeshUtil.stackTrace(e, 2)));
-// asyncContext.onComplete(err, handler);
-// long endTime = System.currentTimeMillis();
-//
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgFailed();
-//
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgCost(endTime -
startTime);
-//
messageLogger.error("message|eventMesh2mq|REQ|SYNC|send2MQCost={}ms|topic={}|bizSeqNo={}|uniqueId={}",
-// endTime - startTime,
-// sendMessageRequestBody.getTopic(),
-// sendMessageRequestBody.getBizSeqNo(),
-// sendMessageRequestBody.getUniqueId(), e);
-// }
- }, new RRCallback() {
+ eventMeshProducer.request(sendMessageContext, new RRCallback() {
@Override
public void onSuccess(Message omsMsg) {
omsMsg.getUserProperties().put(Constants.PROPERTY_MESSAGE_BORN_TIMESTAMP,
omsMsg.getSystemProperties("BORN_TIMESTAMP"));
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 fe32180..a058c9b 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
@@ -55,13 +55,9 @@ public class EventMeshProducer {
mqProducerWrapper.send(sendMsgContext.getMsg(), sendCallback);
}
- public void request(SendMessageContext sendMsgContext, SendCallback
sendCallback, RRCallback rrCallback, long timeout)
+ public void request(SendMessageContext sendMsgContext, RRCallback
rrCallback, long timeout)
throws Exception {
- mqProducerWrapper.request(sendMsgContext.getMsg(), sendCallback,
rrCallback, timeout);
- }
-
- public Message request(SendMessageContext sendMessageContext, long
timeout) throws Exception {
- return mqProducerWrapper.request(sendMessageContext.getMsg(), timeout);
+ mqProducerWrapper.request(sendMsgContext.getMsg(), rrCallback,
timeout);
}
public boolean reply(final SendMessageContext sendMsgContext, final
SendCallback sendCallback) throws Exception {
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 946cc8f..311d6eb 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
@@ -145,9 +145,9 @@ public class ClientGroupWrapper {
return true;
}
- public void request(UpStreamMsgContext upStreamMsgContext, SendCallback
sendCallback, RRCallback rrCallback, long timeout)
+ public void request(UpStreamMsgContext upStreamMsgContext, RRCallback
rrCallback, long timeout)
throws Exception {
- mqProducerWrapper.request(upStreamMsgContext.getMsg(), sendCallback,
rrCallback, timeout);
+ mqProducerWrapper.request(upStreamMsgContext.getMsg(), rrCallback,
timeout);
}
public boolean reply(UpStreamMsgContext upStreamMsgContext) throws
Exception {
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
index 33bfd0b..6378fa2 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
@@ -82,13 +82,15 @@ public class SessionSender {
if (Command.REQUEST_TO_SERVER == cmd) {
long ttl =
msg.getSystemProperties(EventMeshConstants.PROPERTY_MESSAGE_TTL) != null ?
Long.parseLong(msg.getSystemProperties(EventMeshConstants.PROPERTY_MESSAGE_TTL))
: EventMeshConstants.DEFAULT_TIMEOUT_IN_MILLISECONDS;
upStreamMsgContext = new
UpStreamMsgContext(header.getSeq(), session, msg);
-
session.getClientGroupWrapper().get().request(upStreamMsgContext, sendCallback,
initSyncRRCallback(header, startTime, taskExecuteTime), ttl);
+
session.getClientGroupWrapper().get().request(upStreamMsgContext,
initSyncRRCallback(header, startTime, taskExecuteTime), ttl);
+ upstreamBuff.release();
} else if (Command.RESPONSE_TO_SERVER == cmd) {
String cluster =
msg.getUserProperties(EventMeshConstants.PROPERTY_MESSAGE_CLUSTER);
if (!StringUtils.isEmpty(cluster)) {
String replyTopic = EventMeshConstants.RR_REPLY_TOPIC;
replyTopic = cluster + "-" + replyTopic;
msg.getSystemProperties().put(Constants.PROPERTY_MESSAGE_DESTINATION,
replyTopic);
+ msg.setTopic(replyTopic);
}
// //for rocketmq support
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]