This is an automated email from the ASF dual-hosted git repository.
mikexue pushed a commit to branch cloudevents
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git
The following commit(s) were added to refs/heads/cloudevents by this push:
new 5f61c24 Fix standalone connector interface, fix example (#608)
5f61c24 is described below
commit 5f61c248ed4c78eeb4602ebd89b5c4c473aa9269
Author: Wenjun Ruan <[email protected]>
AuthorDate: Wed Nov 24 10:52:59 2021 +0800
Fix standalone connector interface, fix example (#608)
---
.../apache/eventmesh/api/consumer/Consumer.java | 2 -
.../apache/eventmesh/api/producer/Producer.java | 6 +-
.../eventmesh-connector-standalone/build.gradle | 4 +-
.../standalone/MessagingAccessPointImpl.java | 91 -----------
.../standalone/broker/StandaloneBroker.java | 9 +-
.../standalone/broker/model/MessageEntity.java | 9 +-
.../standalone/broker/task/SubScribeTask.java | 27 ++--
.../standalone/consumer/StandaloneConsumer.java | 92 ++++-------
.../consumer/StandaloneConsumerAdaptor.java | 172 +++++++++++++++++++++
.../StandaloneMeshMQPushConsumerAdaptor.java | 135 ----------------
.../producer/StandaloneMeshMQProducerAdaptor.java | 125 ---------------
.../standalone/producer/StandaloneProducer.java | 131 ++++++++++------
.../producer/StandaloneProducerAdaptor.java | 112 ++++++++++++++
... => org.apache.eventmesh.api.consumer.Consumer} | 2 +-
... => org.apache.eventmesh.api.producer.Producer} | 2 +-
.../standalone/broker/StandaloneBrokerTest.java | 11 +-
.../eventmeshmessage}/AsyncPublishInstance.java | 9 +-
.../AsyncSyncRequestInstance.java | 2 +-
.../eventmeshmessage}/SyncRequestInstance.java | 47 +++---
.../http/demo/sub/controller/SubController.java | 7 +-
.../http/demo/sub/service/SubService.java | 2 +-
.../eventmesh/tcp/common/EventMeshTestUtils.java | 62 ++++----
.../{ => pub/eventmeshmessage}/AsyncPublish.java | 24 +--
.../eventmeshmessage}/AsyncPublishBroadcast.java | 34 ++--
.../{ => pub/eventmeshmessage}/SyncRequest.java | 37 ++---
.../{ => sub/eventmeshmessage}/AsyncSubscribe.java | 20 ++-
.../eventmeshmessage}/AsyncSubscribeBroadcast.java | 20 +--
.../{ => sub/eventmeshmessage}/SyncResponse.java | 18 ++-
.../client/http/producer/OpenMessageProducer.java | 7 +-
.../eventmesh/client/tcp/common/MessageUtils.java | 39 ++---
.../EventMeshMessageTCPPubClient.java | 3 +-
.../EventMeshMessageTCPSubClient.java | 2 +-
.../client/http/demo/SyncRequestInstance.java | 45 +++---
.../client/http/util/HttpLoadBalanceUtilsTest.java | 17 +-
.../client/tcp/common/EventMeshTestUtils.java | 63 ++++----
.../eventmesh/client/tcp/demo/SyncRequest.java | 25 +--
36 files changed, 696 insertions(+), 717 deletions(-)
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
index a6fafe0..87b3ac7 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
@@ -28,8 +28,6 @@ import java.util.Properties;
import io.cloudevents.CloudEvent;
-
-
/**
* Consumer Interface.
*/
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
index 04662e9..7a9e032 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
@@ -17,7 +17,11 @@
package org.apache.eventmesh.api.producer;
-import org.apache.eventmesh.api.*;
+import org.apache.eventmesh.api.LifeCycle;
+import org.apache.eventmesh.api.RRCallback;
+import org.apache.eventmesh.api.RequestReplyCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
import org.apache.eventmesh.spi.EventMeshExtensionType;
import org.apache.eventmesh.spi.EventMeshSPI;
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/build.gradle
b/eventmesh-connector-plugin/eventmesh-connector-standalone/build.gradle
index 32b1372..e857140 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-standalone/build.gradle
+++ b/eventmesh-connector-plugin/eventmesh-connector-standalone/build.gradle
@@ -16,7 +16,7 @@
*/
dependencies {
- compileOnly project(":eventmesh-common")
- compileOnly project(":eventmesh-connector-plugin:eventmesh-connector-api")
+ implementation project(":eventmesh-common")
+ implementation
project(":eventmesh-connector-plugin:eventmesh-connector-api")
}
\ No newline at end of file
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/MessagingAccessPointImpl.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/MessagingAccessPointImpl.java
deleted file mode 100644
index 6701e81..0000000
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/MessagingAccessPointImpl.java
+++ /dev/null
@@ -1,91 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.eventmesh.connector.standalone;
-
-import io.openmessaging.api.Consumer;
-import io.openmessaging.api.MessagingAccessPoint;
-import io.openmessaging.api.Producer;
-import io.openmessaging.api.PullConsumer;
-import io.openmessaging.api.batch.BatchConsumer;
-import io.openmessaging.api.order.OrderConsumer;
-import io.openmessaging.api.order.OrderProducer;
-import io.openmessaging.api.transaction.LocalTransactionChecker;
-import io.openmessaging.api.transaction.TransactionProducer;
-import org.apache.eventmesh.connector.standalone.consumer.StandaloneConsumer;
-import org.apache.eventmesh.connector.standalone.producer.StandaloneProducer;
-
-import java.util.Properties;
-
-public class MessagingAccessPointImpl implements MessagingAccessPoint {
-
- private Properties accessPointProperties;
-
- public MessagingAccessPointImpl(Properties accessPointProperties) {
- this.accessPointProperties = accessPointProperties;
- }
-
- @Override
- public String version() {
- return null;
- }
-
- @Override
- public Properties attributes() {
- return accessPointProperties;
- }
-
- @Override
- public Producer createProducer(Properties properties) {
- return new StandaloneProducer(properties);
- }
-
- @Override
- public OrderProducer createOrderProducer(Properties properties) {
- return null;
- }
-
- @Override
- public TransactionProducer createTransactionProducer(Properties
properties, LocalTransactionChecker checker) {
- return null;
- }
-
- @Override
- public TransactionProducer createTransactionProducer(Properties
properties) {
- return null;
- }
-
- @Override
- public Consumer createConsumer(Properties properties) {
- return new StandaloneConsumer(properties);
- }
-
- @Override
- public PullConsumer createPullConsumer(Properties properties) {
- return null;
- }
-
- @Override
- public BatchConsumer createBatchConsumer(Properties properties) {
- return null;
- }
-
- @Override
- public OrderConsumer createOrderedConsumer(Properties properties) {
- return null;
- }
-}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBroker.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBroker.java
index 8f35c3a..72a9b58 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBroker.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBroker.java
@@ -17,6 +17,7 @@
package org.apache.eventmesh.connector.standalone.broker;
+import io.cloudevents.CloudEvent;
import io.openmessaging.api.Message;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.eventmesh.connector.standalone.broker.model.MessageEntity;
@@ -53,7 +54,7 @@ public class StandaloneBroker {
* @param message message
* @throws InterruptedException
*/
- public MessageEntity putMessage(String topicName, Message message) throws
InterruptedException {
+ public MessageEntity putMessage(String topicName, CloudEvent message)
throws InterruptedException {
Pair<MessageQueue, AtomicLong> pair = createTopicIfAbsent(topicName);
AtomicLong topicOffset = pair.getRight();
MessageQueue messageQueue = pair.getLeft();
@@ -70,7 +71,7 @@ public class StandaloneBroker {
*
* @param topicName
*/
- public Message takeMessage(String topicName) throws InterruptedException {
+ public CloudEvent takeMessage(String topicName) throws
InterruptedException {
TopicMetadata topicMetadata = new TopicMetadata(topicName);
return messageContainer.computeIfAbsent(topicMetadata, k -> new
MessageQueue()).take().getMessage();
}
@@ -80,7 +81,7 @@ public class StandaloneBroker {
*
* @param topicName
*/
- public Message getMessage(String topicName) {
+ public CloudEvent getMessage(String topicName) {
TopicMetadata topicMetadata = new TopicMetadata(topicName);
MessageEntity head = messageContainer.computeIfAbsent(topicMetadata, k
-> new MessageQueue()).getHead();
if (head == null) {
@@ -96,7 +97,7 @@ public class StandaloneBroker {
* @param offset offset
* @return
*/
- public Message getMessage(String topicName, long offset) {
+ public CloudEvent getMessage(String topicName, long offset) {
TopicMetadata topicMetadata = new TopicMetadata(topicName);
MessageEntity messageEntity =
messageContainer.computeIfAbsent(topicMetadata, k -> new
MessageQueue()).getByOffset(offset);
if (messageEntity == null) {
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/model/MessageEntity.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/model/MessageEntity.java
index d816891..968053a 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/model/MessageEntity.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/model/MessageEntity.java
@@ -17,6 +17,7 @@
package org.apache.eventmesh.connector.standalone.broker.model;
+import io.cloudevents.CloudEvent;
import io.openmessaging.api.Message;
import java.io.Serializable;
@@ -25,13 +26,13 @@ public class MessageEntity implements Serializable {
private TopicMetadata topicMetadata;
- private Message message;
+ private CloudEvent message;
private long offset;
private long createTimeMills;
- public MessageEntity(TopicMetadata topicMetadata, Message message, long
offset, long currentTimeMills) {
+ public MessageEntity(TopicMetadata topicMetadata, CloudEvent message, long
offset, long currentTimeMills) {
this.topicMetadata = topicMetadata;
this.message = message;
this.offset = offset;
@@ -46,11 +47,11 @@ public class MessageEntity implements Serializable {
this.topicMetadata = topicMetadata;
}
- public Message getMessage() {
+ public CloudEvent getMessage() {
return message;
}
- public void setMessage(Message message) {
+ public void setMessage(CloudEvent message) {
this.message = message;
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
index 430e4aa..bfbf167 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
@@ -27,14 +27,14 @@ import java.util.concurrent.atomic.AtomicInteger;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import io.openmessaging.api.Message;
+import io.cloudevents.CloudEvent;
public class SubScribeTask implements Runnable {
- private String topicName;
- private StandaloneBroker standaloneBroker;
- private EventListener listener;
- private volatile boolean isRunning;
+ private String topicName;
+ private StandaloneBroker standaloneBroker;
+ private EventListener listener;
+ private volatile boolean isRunning;
private AtomicInteger offset;
@@ -55,13 +55,13 @@ public class SubScribeTask implements Runnable {
try {
logger.debug("execute subscribe task, topic: {}, offset: {}",
topicName, offset);
if (offset == null) {
- Message message = standaloneBroker.getMessage(topicName);
+ CloudEvent message =
standaloneBroker.getMessage(topicName);
if (message != null) {
- offset = new AtomicInteger((int) message.getOffset());
+ offset = new AtomicInteger((int)
message.getExtension("offset"));
}
}
if (offset != null) {
- Message message = standaloneBroker.getMessage(topicName,
offset.get());
+ CloudEvent message =
standaloneBroker.getMessage(topicName, offset.get());
if (message != null) {
EventMeshAsyncConsumeContext consumeContext = new
EventMeshAsyncConsumeContext() {
@Override
@@ -69,7 +69,8 @@ public class SubScribeTask implements Runnable {
switch (action) {
case CommitMessage:
// update offset
- logger.info("message commit, topic:
{}, current offset:{}", topicName, offset.get());
+ logger.info("message commit, topic:
{}, current offset:{}", topicName,
+ offset.get());
break;
case ReconsumeLater:
// don't update offset
@@ -77,7 +78,8 @@ public class SubScribeTask implements Runnable {
case ManualAck:
// update offset
offset.incrementAndGet();
- logger.info("message ack, topic: {},
current offset:{}", topicName, offset.get());
+ logger
+ .info("message ack, topic: {},
current offset:{}", topicName, offset.get());
break;
default:
@@ -89,13 +91,14 @@ public class SubScribeTask implements Runnable {
}
} catch (Exception ex) {
- logger.error("consumer error, topic: {}, offset: {}",
topicName, offset == null ? null : offset.get(), ex);
+ logger.error("consumer error, topic: {}, offset: {}",
topicName, offset == null ? null : offset.get(),
+ ex);
}
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
logger.error("Thread is interrupted, topic: {}, offset: {}
thread name: {}",
- topicName, offset == null ? null : offset.get(),
Thread.currentThread().getName(), e);
+ topicName, offset == null ? null : offset.get(),
Thread.currentThread().getName(), e);
Thread.currentThread().interrupt();
}
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumer.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumer.java
index d206cea..5a18379 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumer.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumer.java
@@ -17,23 +17,21 @@
package org.apache.eventmesh.connector.standalone.consumer;
-import io.openmessaging.api.AsyncGenericMessageListener;
-import io.openmessaging.api.AsyncMessageListener;
-import io.openmessaging.api.Consumer;
-import io.openmessaging.api.GenericMessageListener;
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageListener;
-import io.openmessaging.api.MessageSelector;
+import org.apache.eventmesh.api.AbstractContext;
+import org.apache.eventmesh.api.EventListener;
+import org.apache.eventmesh.api.consumer.Consumer;
import org.apache.eventmesh.common.ThreadPoolFactory;
import org.apache.eventmesh.connector.standalone.broker.StandaloneBroker;
import org.apache.eventmesh.connector.standalone.broker.model.TopicMetadata;
import org.apache.eventmesh.connector.standalone.broker.task.SubScribeTask;
+import java.util.List;
import java.util.Properties;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.atomic.AtomicBoolean;
-import java.util.function.Predicate;
+
+import io.cloudevents.CloudEvent;
public class StandaloneConsumer implements Consumer {
@@ -50,30 +48,49 @@ public class StandaloneConsumer implements Consumer {
this.subscribeTaskTable = new ConcurrentHashMap<>(16);
this.isStarted = new AtomicBoolean(false);
this.consumeExecutorService =
ThreadPoolFactory.createThreadPoolExecutor(
- Runtime.getRuntime().availableProcessors() * 2,
- Runtime.getRuntime().availableProcessors() * 2,
- "StandaloneConsumerThread"
+ Runtime.getRuntime().availableProcessors() * 2,
+ Runtime.getRuntime().availableProcessors() * 2,
+ "StandaloneConsumerThread"
);
}
@Override
- public void subscribe(String topic, String subExpression, MessageListener
listener) {
+ public boolean isStarted() {
+ return isStarted.get();
}
@Override
- public void subscribe(String topic, MessageSelector selector,
MessageListener listener) {
+ public boolean isClosed() {
+ return !isStarted.get();
}
@Override
- public <T> void subscribe(String topic, String subExpression,
GenericMessageListener<T> listener) {
+ public void start() {
+ isStarted.compareAndSet(false, true);
}
@Override
- public <T> void subscribe(String topic, MessageSelector selector,
GenericMessageListener<T> listener) {
+ public void shutdown() {
+ isStarted.compareAndSet(true, false);
+ subscribeTaskTable.forEach(((topic, subScribeTask) ->
subScribeTask.shutdown()));
+ subscribeTaskTable.clear();
+ }
+
+ @Override
+ public void init(Properties keyValue) throws Exception {
+
}
@Override
- public void subscribe(String topic, String subExpression,
AsyncMessageListener listener) {
+ public void updateOffset(List<CloudEvent> cloudEvents, AbstractContext
context) {
+ cloudEvents.forEach(cloudEvent -> standaloneBroker.updateOffset(
+ new TopicMetadata(cloudEvent.getSubject()), (Long)
cloudEvent.getExtension("offset"))
+ );
+
+ }
+
+ @Override
+ public void subscribe(String topic, EventListener listener) throws
Exception {
if (listener == null) {
throw new IllegalArgumentException("listener cannot be null");
}
@@ -89,18 +106,6 @@ public class StandaloneConsumer implements Consumer {
}
@Override
- public void subscribe(String topic, MessageSelector selector,
AsyncMessageListener listener) {
- }
-
- @Override
- public <T> void subscribe(String topic, String subExpression,
AsyncGenericMessageListener<T> listener) {
- }
-
- @Override
- public <T> void subscribe(String topic, MessageSelector selector,
AsyncGenericMessageListener<T> listener) {
- }
-
- @Override
public void unsubscribe(String topic) {
if (!subscribeTaskTable.containsKey(topic)) {
return;
@@ -111,35 +116,4 @@ public class StandaloneConsumer implements Consumer {
subscribeTaskTable.remove(topic);
}
}
-
- @Override
- public void updateCredential(Properties credentialProperties) {
-
- }
-
- @Override
- public boolean isStarted() {
- return isStarted.get();
- }
-
- @Override
- public boolean isClosed() {
- return !isStarted.get();
- }
-
- @Override
- public void start() {
- isStarted.compareAndSet(false, true);
- }
-
- @Override
- public void shutdown() {
- isStarted.compareAndSet(true, false);
- subscribeTaskTable.forEach(((topic, subScribeTask) ->
subScribeTask.shutdown()));
- subscribeTaskTable.clear();
- }
-
- public void updateOffset(Message message) {
- standaloneBroker.updateOffset(new TopicMetadata(message.getTopic()),
message.getOffset());
- }
}
\ No newline at end of file
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumerAdaptor.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumerAdaptor.java
new file mode 100644
index 0000000..56518a7
--- /dev/null
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumerAdaptor.java
@@ -0,0 +1,172 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.eventmesh.connector.standalone.consumer;
+
+import org.apache.eventmesh.api.AbstractContext;
+import org.apache.eventmesh.api.EventListener;
+import org.apache.eventmesh.api.consumer.Consumer;
+
+import java.util.List;
+import java.util.Properties;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.cloudevents.CloudEvent;
+
+public class StandaloneConsumerAdaptor implements Consumer {
+
+ private final Logger logger =
LoggerFactory.getLogger(StandaloneConsumerAdaptor.class);
+
+ private StandaloneConsumer consumer;
+
+ public StandaloneConsumerAdaptor() {
+ }
+
+ @Override
+ public boolean isStarted() {
+ return false;
+ }
+
+ @Override
+ public boolean isClosed() {
+ return false;
+ }
+
+ @Override
+ public void start() {
+
+ }
+
+ @Override
+ public void shutdown() {
+
+ }
+
+ @Override
+ public void init(Properties keyValue) throws Exception {
+
+ }
+
+ @Override
+ public void updateOffset(List<CloudEvent> cloudEvents, AbstractContext
context) {
+
+ }
+
+ @Override
+ public void subscribe(String topic, EventListener listener) throws
Exception {
+
+ }
+
+ @Override
+ public void unsubscribe(String topic) {
+
+ }
+
+// @Override
+// public void init(Properties keyValue) throws Exception {
+// String producerGroup = keyValue.getProperty("producerGroup");
+//
+// MessagingAccessPointImpl messagingAccessPoint = new
MessagingAccessPointImpl(keyValue);
+// consumer = (StandaloneConsumer)
messagingAccessPoint.createConsumer(keyValue);
+//
+// }
+//
+// @Override
+// public void updateOffset(List<Message> msgs, AbstractContext context) {
+// for(Message message : msgs) {
+// consumer.updateOffset(message);
+// }
+// }
+//
+// @Override
+// public void subscribe(String topic, AsyncMessageListener listener)
throws Exception {
+// // todo: support subExpression
+// consumer.subscribe(topic, "*", listener);
+// }
+//
+// @Override
+// public void unsubscribe(String topic) {
+// consumer.unsubscribe(topic);
+// }
+//
+// @Override
+// public void subscribe(String topic, String subExpression,
MessageListener listener) {
+// throw new UnsupportedOperationException("not supported yet");
+// }
+//
+// @Override
+// public void subscribe(String topic, MessageSelector selector,
MessageListener listener) {
+// throw new UnsupportedOperationException("not supported yet");
+// }
+//
+// @Override
+// public <T> void subscribe(String topic, String subExpression,
GenericMessageListener<T> listener) {
+// throw new UnsupportedOperationException("not supported yet");
+// }
+//
+// @Override
+// public <T> void subscribe(String topic, MessageSelector selector,
GenericMessageListener<T> listener) {
+// throw new UnsupportedOperationException("not supported yet");
+// }
+//
+// @Override
+// public void subscribe(String topic, String subExpression,
AsyncMessageListener listener) {
+// throw new UnsupportedOperationException("not supported yet");
+// }
+//
+// @Override
+// public void subscribe(String topic, MessageSelector selector,
AsyncMessageListener listener) {
+// throw new UnsupportedOperationException("not supported yet");
+// }
+//
+// @Override
+// public <T> void subscribe(String topic, String subExpression,
AsyncGenericMessageListener<T> listener) {
+// throw new UnsupportedOperationException("not supported yet");
+// }
+//
+// @Override
+// public <T> void subscribe(String topic, MessageSelector selector,
AsyncGenericMessageListener<T> listener) {
+// throw new UnsupportedOperationException("not supported yet");
+// }
+//
+// @Override
+// public void updateCredential(Properties credentialProperties) {
+//
+// }
+//
+// @Override
+// public boolean isStarted() {
+// return consumer.isStarted();
+// }
+//
+// @Override
+// public boolean isClosed() {
+// return consumer.isClosed();
+// }
+//
+// @Override
+// public void start() {
+// consumer.start();
+// }
+//
+// @Override
+// public void shutdown() {
+// consumer.shutdown();
+// }
+}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneMeshMQPushConsumerAdaptor.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneMeshMQPushConsumerAdaptor.java
deleted file mode 100644
index 3bd4ec9..0000000
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneMeshMQPushConsumerAdaptor.java
+++ /dev/null
@@ -1,135 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.eventmesh.connector.standalone.consumer;
-
-import io.openmessaging.api.AsyncGenericMessageListener;
-import io.openmessaging.api.AsyncMessageListener;
-import io.openmessaging.api.GenericMessageListener;
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageListener;
-import io.openmessaging.api.MessageSelector;
-import org.apache.eventmesh.api.AbstractContext;
-import org.apache.eventmesh.api.consumer.MeshMQPushConsumer;
-import org.apache.eventmesh.connector.standalone.MessagingAccessPointImpl;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.List;
-import java.util.Properties;
-
-public class StandaloneMeshMQPushConsumerAdaptor implements MeshMQPushConsumer
{
-
- private final Logger logger =
LoggerFactory.getLogger(StandaloneMeshMQPushConsumerAdaptor.class);
-
- private StandaloneConsumer consumer;
-
- public StandaloneMeshMQPushConsumerAdaptor() {
- }
-
- @Override
- public void init(Properties keyValue) throws Exception {
- String producerGroup = keyValue.getProperty("producerGroup");
-
- MessagingAccessPointImpl messagingAccessPoint = new
MessagingAccessPointImpl(keyValue);
- consumer = (StandaloneConsumer)
messagingAccessPoint.createConsumer(keyValue);
-
- }
-
- @Override
- public void updateOffset(List<Message> msgs, AbstractContext context) {
- for(Message message : msgs) {
- consumer.updateOffset(message);
- }
- }
-
- @Override
- public void subscribe(String topic, AsyncMessageListener listener) throws
Exception {
- // todo: support subExpression
- consumer.subscribe(topic, "*", listener);
- }
-
- @Override
- public void unsubscribe(String topic) {
- consumer.unsubscribe(topic);
- }
-
- @Override
- public void subscribe(String topic, String subExpression, MessageListener
listener) {
- throw new UnsupportedOperationException("not supported yet");
- }
-
- @Override
- public void subscribe(String topic, MessageSelector selector,
MessageListener listener) {
- throw new UnsupportedOperationException("not supported yet");
- }
-
- @Override
- public <T> void subscribe(String topic, String subExpression,
GenericMessageListener<T> listener) {
- throw new UnsupportedOperationException("not supported yet");
- }
-
- @Override
- public <T> void subscribe(String topic, MessageSelector selector,
GenericMessageListener<T> listener) {
- throw new UnsupportedOperationException("not supported yet");
- }
-
- @Override
- public void subscribe(String topic, String subExpression,
AsyncMessageListener listener) {
- throw new UnsupportedOperationException("not supported yet");
- }
-
- @Override
- public void subscribe(String topic, MessageSelector selector,
AsyncMessageListener listener) {
- throw new UnsupportedOperationException("not supported yet");
- }
-
- @Override
- public <T> void subscribe(String topic, String subExpression,
AsyncGenericMessageListener<T> listener) {
- throw new UnsupportedOperationException("not supported yet");
- }
-
- @Override
- public <T> void subscribe(String topic, MessageSelector selector,
AsyncGenericMessageListener<T> listener) {
- throw new UnsupportedOperationException("not supported yet");
- }
-
- @Override
- public void updateCredential(Properties credentialProperties) {
-
- }
-
- @Override
- public boolean isStarted() {
- return consumer.isStarted();
- }
-
- @Override
- public boolean isClosed() {
- return consumer.isClosed();
- }
-
- @Override
- public void start() {
- consumer.start();
- }
-
- @Override
- public void shutdown() {
- consumer.shutdown();
- }
-}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneMeshMQProducerAdaptor.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneMeshMQProducerAdaptor.java
deleted file mode 100644
index fb632c7..0000000
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneMeshMQProducerAdaptor.java
+++ /dev/null
@@ -1,125 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.eventmesh.connector.standalone.producer;
-
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageBuilder;
-import io.openmessaging.api.SendCallback;
-import io.openmessaging.api.SendResult;
-import org.apache.eventmesh.api.RRCallback;
-import org.apache.eventmesh.api.producer.MeshMQProducer;
-import org.apache.eventmesh.connector.standalone.MessagingAccessPointImpl;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.Properties;
-import java.util.concurrent.ExecutorService;
-
-public class StandaloneMeshMQProducerAdaptor implements MeshMQProducer {
-
- private final Logger logger =
LoggerFactory.getLogger(StandaloneMeshMQProducerAdaptor.class);
-
- private StandaloneProducer standaloneProducer;
-
- public StandaloneMeshMQProducerAdaptor() {
- }
-
- @Override
- public void init(Properties properties) throws Exception {
- MessagingAccessPointImpl messagingAccessPoint = new
MessagingAccessPointImpl(properties);
- standaloneProducer = (StandaloneProducer)
messagingAccessPoint.createProducer(properties);
- }
-
- @Override
- public void send(Message message, SendCallback sendCallback) throws
Exception {
- standaloneProducer.sendAsync(message, sendCallback);
- }
-
- @Override
- public void request(Message message, RRCallback rrCallback, long timeout)
throws Exception {
- throw new UnsupportedOperationException("not support request-reply
mode when eventstore=standalone");
- }
-
- @Override
- public boolean reply(Message message, SendCallback sendCallback) throws
Exception {
- throw new UnsupportedOperationException("not support request-reply
mode when eventstore=standalone");
- }
-
- @Override
- public void checkTopicExist(String topic) throws Exception {
- boolean exist = standaloneProducer.checkTopicExist(topic);
- if (!exist) {
- throw new RuntimeException(String.format("topic:%s is not exist",
topic));
- }
- }
-
- @Override
- public void setExtFields() {
-
- }
-
- @Override
- public SendResult send(Message message) {
- return standaloneProducer.send(message);
- }
-
- @Override
- public void sendOneway(Message message) {
- standaloneProducer.sendOneway(message);
- }
-
- @Override
- public void sendAsync(Message message, SendCallback sendCallback) {
- standaloneProducer.sendAsync(message, sendCallback);
- }
-
- @Override
- public void setCallbackExecutor(ExecutorService callbackExecutor) {
- standaloneProducer.setCallbackExecutor(callbackExecutor);
- }
-
- @Override
- public void updateCredential(Properties credentialProperties) {
- standaloneProducer.updateCredential(credentialProperties);
- }
-
- @Override
- public boolean isStarted() {
- return standaloneProducer.isStarted();
- }
-
- @Override
- public boolean isClosed() {
- return standaloneProducer.isClosed();
- }
-
- @Override
- public void start() {
- standaloneProducer.start();
- }
-
- @Override
- public void shutdown() {
- standaloneProducer.shutdown();
- }
-
- @Override
- public <T> MessageBuilder<T> messageBuilder() {
- return null;
- }
-}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducer.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducer.java
index 1522f51..57c8a2a 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducer.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducer.java
@@ -17,23 +17,26 @@
package org.apache.eventmesh.connector.standalone.producer;
-import com.google.common.base.Preconditions;
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageBuilder;
-import io.openmessaging.api.OnExceptionContext;
-import io.openmessaging.api.Producer;
-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.api.RequestReplyCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
+import org.apache.eventmesh.api.exception.OnExceptionContext;
+import org.apache.eventmesh.api.producer.Producer;
import org.apache.eventmesh.connector.standalone.broker.StandaloneBroker;
import org.apache.eventmesh.connector.standalone.broker.model.MessageEntity;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
import java.util.Properties;
-import java.util.concurrent.ExecutorService;
import java.util.concurrent.atomic.AtomicBoolean;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.google.common.base.Preconditions;
+
+import io.cloudevents.CloudEvent;
+
public class StandaloneProducer implements Producer {
private Logger logger = LoggerFactory.getLogger(StandaloneProducer.class);
@@ -48,80 +51,110 @@ public class StandaloneProducer implements Producer {
}
@Override
- public SendResult send(Message message) {
- Preconditions.checkNotNull(message);
+ public boolean isStarted() {
+ return isStarted.get();
+ }
+
+ @Override
+ public boolean isClosed() {
+ return !isStarted.get();
+ }
+
+ @Override
+ public void start() {
+ isStarted.compareAndSet(false, true);
+ }
+
+ @Override
+ public void shutdown() {
+ isStarted.compareAndSet(true, false);
+ }
+
+ @Override
+ public void init(Properties properties) throws Exception {
+
+ }
+
+ @Override
+ public SendResult publish(CloudEvent cloudEvent) {
+ Preconditions.checkNotNull(cloudEvent);
try {
- MessageEntity messageEntity =
standaloneBroker.putMessage(message.getTopic(), message);
+ MessageEntity messageEntity =
standaloneBroker.putMessage(cloudEvent.getSubject(), cloudEvent);
SendResult sendResult = new SendResult();
- sendResult.setTopic(message.getTopic());
+ sendResult.setTopic(cloudEvent.getSubject());
sendResult.setMessageId(String.valueOf(messageEntity.getOffset()));
return sendResult;
} catch (Exception e) {
- logger.error("send message error, topic: {}", message.getTopic(),
e);
- throw new OMSRuntimeException(String.format("Send message error,
topic: %s", message.getTopic()));
+ logger.error("send message error, topic: {}",
cloudEvent.getSubject(), e);
+ throw new ConnectorRuntimeException(
+ String.format("Send message error, topic: %s",
cloudEvent.getSubject()));
}
}
@Override
- public void sendOneway(Message message) {
- send(message);
- }
-
- @Override
- public void sendAsync(Message message, SendCallback sendCallback) {
- Preconditions.checkNotNull(message);
+ public void publish(CloudEvent cloudEvent, SendCallback sendCallback)
throws Exception {
+ Preconditions.checkNotNull(cloudEvent);
Preconditions.checkNotNull(sendCallback);
try {
- SendResult sendResult = send(message);
+ SendResult sendResult = publish(cloudEvent);
sendCallback.onSuccess(sendResult);
} catch (Exception ex) {
- OnExceptionContext exceptionContext = new OnExceptionContext();
- exceptionContext.setTopic(message.getTopic());
- exceptionContext.setException(new OMSRuntimeException(ex));
- exceptionContext.setMessageId(message.getMsgID());
- sendCallback.onException(exceptionContext);
+ OnExceptionContext onExceptionContext = new OnExceptionContext();
+ onExceptionContext.setMessageId(cloudEvent.getId());
+ onExceptionContext.setTopic(cloudEvent.getSubject());
+ onExceptionContext.setException(new ConnectorRuntimeException(ex));
+ sendCallback.onException(onExceptionContext);
}
-
}
@Override
- public void setCallbackExecutor(ExecutorService callbackExecutor) {
-
+ public void sendOneway(CloudEvent cloudEvent) {
+ publish(cloudEvent);
}
@Override
- public void updateCredential(Properties credentialProperties) {
-
+ public void sendAsync(CloudEvent cloudEvent, SendCallback sendCallback) {
+ Preconditions.checkNotNull(cloudEvent);
+ Preconditions.checkNotNull(sendCallback);
+ // todo: current is not async
+ try {
+ SendResult sendResult = publish(cloudEvent);
+ sendCallback.onSuccess(sendResult);
+ } catch (Exception ex) {
+ OnExceptionContext onExceptionContext = new OnExceptionContext();
+ onExceptionContext.setMessageId(cloudEvent.getId());
+ onExceptionContext.setTopic(cloudEvent.getSubject());
+ onExceptionContext.setException(new ConnectorRuntimeException(ex));
+ sendCallback.onException(onExceptionContext);
+ }
}
@Override
- public boolean isStarted() {
- return isStarted.get();
+ public void request(CloudEvent cloudEvent, RRCallback rrCallback, long
timeout) throws Exception {
+ throw new ConnectorRuntimeException("Request is not supported");
}
@Override
- public boolean isClosed() {
- return !isStarted.get();
+ public void request(CloudEvent cloudEvent, RequestReplyCallback
rrCallback, long timeout) throws Exception {
+ throw new ConnectorRuntimeException("Request is not supported");
}
@Override
- public void start() {
- isStarted.compareAndSet(false, true);
+ public boolean reply(CloudEvent cloudEvent, SendCallback sendCallback)
throws Exception {
+ throw new ConnectorRuntimeException("Reply is not supported");
}
@Override
- public void shutdown() {
- isStarted.compareAndSet(true, false);
-
+ public void checkTopicExist(String topic) throws Exception {
+ boolean exist = standaloneBroker.checkTopicExist(topic);
+ if (!exist) {
+ throw new ConnectorRuntimeException(String.format("topic:%s is not
exist", topic));
+ }
}
@Override
- public <T> MessageBuilder<T> messageBuilder() {
- return null;
- }
+ public void setExtFields() {
- public boolean checkTopicExist(String topicName) {
- return standaloneBroker.checkTopicExist(topicName);
}
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducerAdaptor.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducerAdaptor.java
new file mode 100644
index 0000000..d60968c
--- /dev/null
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducerAdaptor.java
@@ -0,0 +1,112 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.eventmesh.connector.standalone.producer;
+
+import org.apache.eventmesh.api.RRCallback;
+import org.apache.eventmesh.api.RequestReplyCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.api.producer.Producer;
+
+import java.util.Properties;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.cloudevents.CloudEvent;
+
+public class StandaloneProducerAdaptor implements Producer {
+
+ private final Logger logger =
LoggerFactory.getLogger(StandaloneProducerAdaptor.class);
+
+ private StandaloneProducer standaloneProducer;
+
+ public StandaloneProducerAdaptor() {
+ }
+
+ @Override
+ public boolean isStarted() {
+ return standaloneProducer.isStarted();
+ }
+
+ @Override
+ public boolean isClosed() {
+ return standaloneProducer.isClosed();
+ }
+
+ @Override
+ public void start() {
+ standaloneProducer.start();
+ }
+
+ @Override
+ public void shutdown() {
+ standaloneProducer.shutdown();
+ }
+
+ @Override
+ public void init(Properties properties) throws Exception {
+ standaloneProducer.init(properties);
+ }
+
+ @Override
+ public SendResult publish(CloudEvent cloudEvent) {
+ return standaloneProducer.publish(cloudEvent);
+ }
+
+ @Override
+ public void publish(CloudEvent cloudEvent, SendCallback sendCallback)
throws Exception {
+ standaloneProducer.publish(cloudEvent, sendCallback);
+ }
+
+ @Override
+ public void sendOneway(CloudEvent cloudEvent) {
+ standaloneProducer.sendOneway(cloudEvent);
+ }
+
+ @Override
+ public void sendAsync(CloudEvent cloudEvent, SendCallback sendCallback) {
+ standaloneProducer.sendAsync(cloudEvent, sendCallback);
+ }
+
+ @Override
+ public void request(CloudEvent cloudEvent, RRCallback rrCallback, long
timeout) throws Exception {
+ standaloneProducer.request(cloudEvent, rrCallback, timeout);
+ }
+
+ @Override
+ public void request(CloudEvent cloudEvent, RequestReplyCallback
rrCallback, long timeout) throws Exception {
+ standaloneProducer.request(cloudEvent, rrCallback, timeout);
+ }
+
+ @Override
+ public boolean reply(CloudEvent cloudEvent, SendCallback sendCallback)
throws Exception {
+ return standaloneProducer.reply(cloudEvent, sendCallback);
+ }
+
+ @Override
+ public void checkTopicExist(String topic) throws Exception {
+ standaloneProducer.checkTopicExist(topic);
+ }
+
+ @Override
+ public void setExtFields() {
+ standaloneProducer.setExtFields();
+ }
+
+}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.MeshMQPushConsumer
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.Consumer
similarity index 96%
rename from
eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.MeshMQPushConsumer
rename to
eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.Consumer
index 190a92a..5ceb0ed 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.MeshMQPushConsumer
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.Consumer
@@ -13,4 +13,4 @@
# See the License for the specific language governing permissions and
# limitations under the License.
-standalone=org.apache.eventmesh.connector.standalone.consumer.StandaloneMeshMQPushConsumerAdaptor
\ No newline at end of file
+standalone=org.apache.eventmesh.connector.standalone.consumer.StandaloneConsumerAdaptor
\ No newline at end of file
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.MeshMQProducer
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.Producer
similarity index 96%
rename from
eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.MeshMQProducer
rename to
eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.Producer
index 4bbac31..1b30d98 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.MeshMQProducer
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.Producer
@@ -13,4 +13,4 @@
# See the License for the specific language governing permissions and
# limitations under the License.
-standalone=org.apache.eventmesh.connector.standalone.producer.StandaloneMeshMQProducerAdaptor
\ No newline at end of file
+standalone=org.apache.eventmesh.connector.standalone.producer.StandaloneProducerAdaptor
\ No newline at end of file
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/test/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBrokerTest.java
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/test/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBrokerTest.java
index a63be75..92e4eb7 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/test/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBrokerTest.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/test/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBrokerTest.java
@@ -17,11 +17,14 @@
package org.apache.eventmesh.connector.standalone.broker;
-import io.openmessaging.api.Message;
import org.apache.eventmesh.connector.standalone.broker.model.MessageEntity;
+
import org.junit.Assert;
import org.junit.Test;
+import io.cloudevents.CloudEvent;
+import io.cloudevents.core.builder.CloudEventBuilder;
+
public class StandaloneBrokerTest {
@Test
@@ -32,15 +35,15 @@ public class StandaloneBrokerTest {
@Test
public void putMessage() throws InterruptedException {
StandaloneBroker instance = StandaloneBroker.getInstance();
- MessageEntity messageEntity = instance.putMessage("test-topic", new
Message());
+ MessageEntity messageEntity = instance.putMessage("test-topic",
CloudEventBuilder.v1().build());
Assert.assertNotNull(messageEntity);
}
@Test
public void takeMessage() throws InterruptedException {
StandaloneBroker instance = StandaloneBroker.getInstance();
- instance.putMessage("test-topic", new Message());
- Message message = instance.takeMessage("test-topic");
+ instance.putMessage("test-topic", CloudEventBuilder.v1().build());
+ CloudEvent message = instance.takeMessage("test-topic");
Assert.assertNotNull(message);
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncPublishInstance.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
similarity index 95%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncPublishInstance.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
index 1af175d..8c90d8e 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncPublishInstance.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
@@ -15,7 +15,9 @@
* limitations under the License.
*/
-package org.apache.eventmesh.http.demo;
+package org.apache.eventmesh.http.demo.pub.eventmeshmessage;
+
+import lombok.extern.slf4j.Slf4j;
import org.apache.eventmesh.client.http.conf.EventMeshHttpClientConfig;
import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
@@ -33,10 +35,9 @@ import java.util.Properties;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+@Slf4j
public class AsyncPublishInstance {
- public static Logger logger =
LoggerFactory.getLogger(AsyncPublishInstance.class);
-
// This messageSize is also used in SubService.java (Subscriber)
public static int messageSize = 5;
@@ -45,7 +46,7 @@ public class AsyncPublishInstance {
final String eventMeshIp = properties.getProperty("eventmesh.ip");
final String eventMeshHttpPort =
properties.getProperty("eventmesh.http.port");
- String eventMeshIPPort;
+ final String eventMeshIPPort;
if (StringUtils.isBlank(eventMeshIp) ||
StringUtils.isBlank(eventMeshHttpPort)) {
// if has multi value, can config as:
127.0.0.1:10105;127.0.0.2:10105
eventMeshIPPort = "127.0.0.1:10105";
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncSyncRequestInstance.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
similarity index 98%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncSyncRequestInstance.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
index a007ed1..4284726 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncSyncRequestInstance.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.eventmesh.http.demo;
+package org.apache.eventmesh.http.demo.pub.eventmeshmessage;
import org.apache.eventmesh.client.http.conf.EventMeshHttpClientConfig;
import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/SyncRequestInstance.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
similarity index 70%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/SyncRequestInstance.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
index 58d287d..4d7e1d2 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/SyncRequestInstance.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
@@ -15,16 +15,17 @@
* limitations under the License.
*/
-package org.apache.eventmesh.http.demo;
-
-import org.apache.commons.lang3.StringUtils;
+package org.apache.eventmesh.http.demo.pub.eventmeshmessage;
import org.apache.eventmesh.client.http.conf.EventMeshHttpClientConfig;
import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
-import org.apache.eventmesh.common.utils.IPUtils;
import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.utils.IPUtils;
import org.apache.eventmesh.common.utils.RandomStringUtils;
import org.apache.eventmesh.common.utils.ThreadUtils;
+
+import org.apache.commons.lang3.StringUtils;
+
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -50,38 +51,36 @@ public class SyncRequestInstance {
eventMeshIPPort = "127.0.0.1:10105";
}
- EventMeshHttpClientConfig eventMeshClientConfig = new
EventMeshHttpClientConfig();
- eventMeshClientConfig.setLiteEventMeshAddr(eventMeshIPPort)
- .setProducerGroup("EventMeshTest-producerGroup")
- .setEnv("env")
- .setIdc("idc")
- .setIp(IPUtils.getLocalAddress())
- .setSys("1234")
- .setPid(String.valueOf(ThreadUtils.getPID()));
+ EventMeshHttpClientConfig eventMeshClientConfig =
EventMeshHttpClientConfig.builder()
+ .liteEventMeshAddr(eventMeshIPPort)
+ .producerGroup("EventMeshTest-producerGroup")
+ .env("env")
+ .idc("idc")
+ .ip(IPUtils.getLocalAddress())
+ .sys("1234")
+ .pid(String.valueOf(ThreadUtils.getPID())).build();
eventMeshHttpProducer = new
EventMeshHttpProducer(eventMeshClientConfig);
- eventMeshHttpProducer.start();
long startTime = System.currentTimeMillis();
- EventMeshMessage eventMeshMessage = new EventMeshMessage();
- eventMeshMessage.setBizSeqNo(RandomStringUtils.generateNum(30))
- .setContent("contentStr with special protocal")
- .setTopic(topic)
- .setUniqueId(RandomStringUtils.generateNum(30));
+ EventMeshMessage eventMeshMessage = EventMeshMessage.builder()
+ .bizSeqNo(RandomStringUtils.generateNum(30))
+ .content("contentStr with special protocal")
+ .topic(topic)
+ .uniqueId(RandomStringUtils.generateNum(30)).build();
EventMeshMessage rsp =
eventMeshHttpProducer.request(eventMeshMessage, 10000);
if (logger.isDebugEnabled()) {
- logger.debug("sendmsg : {}, return : {}, cost:{}ms",
eventMeshMessage.getContent(), rsp.getContent(), System.currentTimeMillis() -
startTime);
+ logger.debug("sendmsg : {}, return : {}, cost:{}ms",
eventMeshMessage.getContent(), rsp.getContent(),
+ System.currentTimeMillis() - startTime);
}
} catch (Exception e) {
logger.warn("send msg failed", e);
}
- try {
- Thread.sleep(30000);
- if (eventMeshHttpProducer != null) {
- eventMeshHttpProducer.shutdown();
- }
+ Thread.sleep(30000);
+ try (final EventMeshHttpProducer close = eventMeshHttpProducer) {
+ // close producer
} catch (Exception e1) {
logger.warn("producer shutdown exception", e1);
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/controller/SubController.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/controller/SubController.java
index aa400ce..2f0089c 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/controller/SubController.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/controller/SubController.java
@@ -17,6 +17,8 @@
package org.apache.eventmesh.http.demo.sub.controller;
+import lombok.extern.slf4j.Slf4j;
+
import org.apache.eventmesh.common.utils.JsonUtils;
import org.apache.eventmesh.http.demo.sub.service.SubService;
@@ -31,18 +33,17 @@ import
org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestMethod;
import org.springframework.web.bind.annotation.RestController;
+@Slf4j
@RestController
@RequestMapping("/sub")
public class SubController {
- public static Logger logger = LoggerFactory.getLogger(SubController.class);
-
@Autowired
private SubService subService;
@RequestMapping(value = "/test", method = RequestMethod.POST)
public String subTest(@RequestBody String message) {
- logger.info("=======receive message======= {}",
JsonUtils.serialize(message));
+ log.info("=======receive message======= {}",
JsonUtils.serialize(message));
subService.consumeMessage(message);
Map<String, Object> map = new HashMap<>();
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/service/SubService.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/service/SubService.java
index 6470c77..0c0f2c1 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/service/SubService.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/service/SubService.java
@@ -26,7 +26,7 @@ import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
import org.apache.eventmesh.common.utils.IPUtils;
import org.apache.eventmesh.common.utils.ThreadUtils;
-import org.apache.eventmesh.http.demo.AsyncPublishInstance;
+import
org.apache.eventmesh.http.demo.pub.eventmeshmessage.AsyncPublishInstance;
import org.apache.eventmesh.util.Utils;
import java.util.ArrayList;
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/common/EventMeshTestUtils.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/common/EventMeshTestUtils.java
index 854eb5b..bd7f7cb 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/common/EventMeshTestUtils.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/common/EventMeshTestUtils.java
@@ -22,49 +22,49 @@ import static
org.apache.eventmesh.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PR
import static
org.apache.eventmesh.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_WQ2ClientBroadCast;
import static
org.apache.eventmesh.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_WQ2ClientUniCast;
-import java.util.concurrent.ThreadLocalRandom;
-
import org.apache.eventmesh.common.protocol.tcp.Command;
import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.Header;
import org.apache.eventmesh.common.protocol.tcp.Package;
import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+import java.util.concurrent.ThreadLocalRandom;
+
public class EventMeshTestUtils {
private static final int seqLength = 10;
public static UserAgent generateClient1() {
return UserAgent.builder()
- .env("test")
- .host("127.0.0.1")
- .password(generateRandomString(8))
- .username("PU4283")
- .producerGroup("EventmeshTest-ProducerGroup")
- .consumerGroup("EventmeshTest-ConsumerGroup")
- .path("/data/app/umg_proxy")
- .port(8362)
- .subsystem("5023")
- .pid(32893)
- .version("2.0.11")
- .idc("FT")
- .build();
+ .env("test")
+ .host("127.0.0.1")
+ .password(generateRandomString(8))
+ .username("PU4283")
+ .producerGroup("EventmeshTest-ProducerGroup")
+ .consumerGroup("EventmeshTest-ConsumerGroup")
+ .path("/data/app/umg_proxy")
+ .port(8362)
+ .subsystem("5023")
+ .pid(32893)
+ .version("2.0.11")
+ .idc("FT")
+ .build();
}
public static UserAgent generateClient2() {
return UserAgent.builder()
- .env("test")
- .host("127.0.0.1")
- .password(generateRandomString(8))
- .username("PU4283")
- .producerGroup("EventmeshTest-ProducerGroup")
- .consumerGroup("EventmeshTest-ConsumerGroup")
- .path("/data/app/umg_proxy")
- .port(9362)
- .subsystem("5017")
- .pid(42893)
- .version("2.0.11")
- .idc("FT")
- .build();
+ .env("test")
+ .host("127.0.0.1")
+ .password(generateRandomString(8))
+ .username("PU4283")
+ .producerGroup("EventmeshTest-ProducerGroup")
+ .consumerGroup("EventmeshTest-ConsumerGroup")
+ .path("/data/app/umg_proxy")
+ .port(9362)
+ .subsystem("5017")
+ .pid(42893)
+ .version("2.0.11")
+ .idc("FT")
+ .build();
}
public static Package syncRR() {
@@ -102,7 +102,7 @@ public class EventMeshTestUtils {
return msg;
}
- private static EventMeshMessage generateSyncRRMqMsg() {
+ public static EventMeshMessage generateSyncRRMqMsg() {
EventMeshMessage mqMsg = new EventMeshMessage();
mqMsg.setTopic(TOPIC_PRX_SyncSubscribeTest);
mqMsg.getProperties().put("msgType", "persistent");
@@ -123,7 +123,7 @@ public class EventMeshTestUtils {
return mqMsg;
}
- private static EventMeshMessage generateAsyncEventMqMsg() {
+ public static EventMeshMessage generateAsyncEventMqMsg() {
EventMeshMessage mqMsg = new EventMeshMessage();
mqMsg.setTopic(TOPIC_PRX_WQ2ClientUniCast);
mqMsg.getProperties().put("REPLY_TO",
"10.36.0.109@ProducerGroup-producerPool-9-access#V1_4_0#CI");
@@ -133,7 +133,7 @@ public class EventMeshTestUtils {
return mqMsg;
}
- private static EventMeshMessage generateBroadcastMqMsg() {
+ public static EventMeshMessage generateBroadcastMqMsg() {
EventMeshMessage mqMsg = new EventMeshMessage();
mqMsg.setTopic(TOPIC_PRX_WQ2ClientBroadCast);
mqMsg.getProperties().put("REPLY_TO",
"10.36.0.109@ProducerGroup-producerPool-9-access#V1_4_0#CI");
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublish.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
similarity index 70%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublish.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
index 90bd4c7..52ccc15 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublish.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
@@ -15,14 +15,14 @@
* limitations under the License.
*/
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.pub.eventmeshmessage;
import java.util.Properties;
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
-import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPPubClient;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.UserAgent;
import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
import org.apache.eventmesh.util.Utils;
@@ -33,7 +33,7 @@ public class AsyncPublish {
public static Logger logger = LoggerFactory.getLogger(AsyncPublish.class);
- private static EventMeshTCPClient client;
+ private static EventMeshMessageTCPPubClient client;
public static AsyncPublish handler = new AsyncPublish();
@@ -43,14 +43,20 @@ public class AsyncPublish {
final int eventMeshTcpPort =
Integer.parseInt(properties.getProperty("eventmesh.tcp.port"));
try {
UserAgent userAgent = EventMeshTestUtils.generateClient1();
- client = new DefaultEventMeshTCPClient(eventMeshIp,
eventMeshTcpPort, userAgent);
+ EventMeshTcpClientConfig eventMeshTcpClientConfig =
EventMeshTcpClientConfig.builder()
+ .host(eventMeshIp)
+ .port(eventMeshTcpPort)
+ .userAgent(userAgent)
+ .build();
+ client = new
EventMeshMessageTCPPubClient(eventMeshTcpClientConfig);
client.init();
client.heartbeat();
for (int i = 0; i < 5; i++) {
- Package asyncMsg = EventMeshTestUtils.asyncMessage();
- logger.info("begin send async msg[{}]==================={}",
i, asyncMsg);
- client.publish(asyncMsg,
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
+ EventMeshMessage eventMeshMessage =
EventMeshTestUtils.generateAsyncEventMqMsg();
+
+ logger.info("begin send async msg[{}]==================={}",
i, eventMeshMessage);
+ client.publish(eventMeshMessage,
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
Thread.sleep(1000);
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublishBroadcast.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
similarity index 67%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublishBroadcast.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
index f1dadef..12071de 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublishBroadcast.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
@@ -15,17 +15,18 @@
* limitations under the License.
*/
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.pub.eventmeshmessage;
-import java.util.Properties;
-
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
-import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPPubClient;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.UserAgent;
import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
import org.apache.eventmesh.util.Utils;
+
+import java.util.Properties;
+
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -33,25 +34,26 @@ public class AsyncPublishBroadcast {
public static Logger logger =
LoggerFactory.getLogger(AsyncPublishBroadcast.class);
- private static EventMeshTCPClient client;
-
public static void main(String[] agrs) throws Exception {
Properties properties =
Utils.readPropertiesFile("application.properties");
final String eventMeshIp = properties.getProperty("eventmesh.ip");
final int eventMeshTcpPort =
Integer.parseInt(properties.getProperty("eventmesh.tcp.port"));
- try {
- UserAgent userAgent = EventMeshTestUtils.generateClient1();
- client = new DefaultEventMeshTCPClient(eventMeshIp,
eventMeshTcpPort, userAgent);
+ UserAgent userAgent = EventMeshTestUtils.generateClient1();
+ EventMeshTcpClientConfig eventMeshTcpClientConfig =
EventMeshTcpClientConfig.builder()
+ .host(eventMeshIp)
+ .port(eventMeshTcpPort)
+ .userAgent(userAgent)
+ .build();
+ try (final EventMeshMessageTCPPubClient client = new
EventMeshMessageTCPPubClient(eventMeshTcpClientConfig)) {
client.init();
client.heartbeat();
- Package broadcastMsg = EventMeshTestUtils.broadcastMessage();
- logger.info("begin send broadcast msg============={}",
broadcastMsg);
- client.broadcast(broadcastMsg,
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
+ EventMeshMessage eventMeshMessage =
EventMeshTestUtils.generateBroadcastMqMsg();
+ logger.info("begin send broadcast msg============={}",
eventMeshMessage);
+ client.broadcast(eventMeshMessage,
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
Thread.sleep(2000);
- // release resource and close client
- // client.close();
+
} catch (Exception e) {
logger.warn("AsyncPublishBroadcast failed", e);
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncRequest.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/SyncRequest.java
similarity index 51%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncRequest.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/SyncRequest.java
index 65339de..844b9cf 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncRequest.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/SyncRequest.java
@@ -15,40 +15,41 @@
* limitations under the License.
*/
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.pub.eventmeshmessage;
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPPubClient;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.Package;
import org.apache.eventmesh.common.protocol.tcp.UserAgent;
import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
+import lombok.extern.slf4j.Slf4j;
+@Slf4j
public class SyncRequest {
- public static Logger logger = LoggerFactory.getLogger(SyncRequest.class);
-
- private static EventMeshTCPClient client;
+ private static EventMeshMessageTCPPubClient client;
public static void main(String[] agrs) throws Exception {
- try {
- UserAgent userAgent = EventMeshTestUtils.generateClient1();
- client = new DefaultEventMeshTCPClient("127.0.0.1", 10000,
userAgent);
+ UserAgent userAgent = EventMeshTestUtils.generateClient1();
+ EventMeshTcpClientConfig eventMeshTcpClientConfig =
EventMeshTcpClientConfig.builder()
+ .host("127.0.0.1")
+ .port(10000)
+ .userAgent(userAgent)
+ .build();
+ try (EventMeshMessageTCPPubClient client = new
EventMeshMessageTCPPubClient(eventMeshTcpClientConfig)) {
client.init();
client.heartbeat();
- Package rrMsg = EventMeshTestUtils.syncRR();
- logger.info("begin send rr msg=================={}", rrMsg);
- Package response = client.rr(rrMsg,
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
- logger.info("receive rr reply==================={}", response);
+ EventMeshMessage eventMeshMessage =
EventMeshTestUtils.generateSyncRRMqMsg();
+ log.info("begin send rr msg=================={}",
eventMeshMessage);
+ Package response = client.rr(eventMeshMessage,
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
+ log.info("receive rr reply==================={}", response);
- // release resource and close client
- // client.close();
} catch (Exception e) {
- logger.warn("SyncRequest failed", e);
+ log.warn("SyncRequest failed", e);
}
}
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribe.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribe.java
similarity index 77%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribe.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribe.java
index 18d607a..ec56aa5 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribe.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribe.java
@@ -15,11 +15,11 @@
* limitations under the License.
*/
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.sub.eventmeshmessage;
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPSubClient;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
@@ -36,7 +36,7 @@ import lombok.extern.slf4j.Slf4j;
@Slf4j
public class AsyncSubscribe implements ReceiveMsgHook<EventMeshMessage> {
- private static EventMeshTCPClient client;
+ private static EventMeshMessageTCPSubClient client;
public static AsyncSubscribe handler = new AsyncSubscribe();
@@ -44,14 +44,18 @@ public class AsyncSubscribe implements
ReceiveMsgHook<EventMeshMessage> {
Properties properties =
Utils.readPropertiesFile("application.properties");
final String eventMeshIp = properties.getProperty("eventmesh.ip");
final int eventMeshTcpPort =
Integer.parseInt(properties.getProperty("eventmesh.tcp.port"));
- try {
- UserAgent userAgent = EventMeshTestUtils.generateClient2();
- client = new DefaultEventMeshTCPClient(eventMeshIp,
eventMeshTcpPort, userAgent);
+ UserAgent userAgent = EventMeshTestUtils.generateClient2();
+ EventMeshTcpClientConfig eventMeshTcpClientConfig =
EventMeshTcpClientConfig.builder()
+ .host(eventMeshIp)
+ .port(eventMeshTcpPort)
+ .userAgent(userAgent)
+ .build();
+ try (EventMeshMessageTCPSubClient client = new
EventMeshMessageTCPSubClient(eventMeshTcpClientConfig)) {
client.init();
client.heartbeat();
client.subscribe("TEST-TOPIC-TCP-ASYNC",
SubscriptionMode.CLUSTERING, SubscriptionType.ASYNC);
- client.registerSubBusiHandler(handler);
+ client.registerBusiHandler(handler);
client.listen();
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribeBroadcast.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribeBroadcast.java
similarity index 79%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribeBroadcast.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribeBroadcast.java
index 1d6c6d1..8d642e1 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribeBroadcast.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribeBroadcast.java
@@ -15,11 +15,11 @@
* limitations under the License.
*/
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.sub.eventmeshmessage;
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPSubClient;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
@@ -36,22 +36,24 @@ import lombok.extern.slf4j.Slf4j;
@Slf4j
public class AsyncSubscribeBroadcast implements
ReceiveMsgHook<EventMeshMessage> {
- private static EventMeshTCPClient client;
-
public static AsyncSubscribeBroadcast handler = new
AsyncSubscribeBroadcast();
public static void main(String[] agrs) throws Exception {
Properties properties =
Utils.readPropertiesFile("application.properties");
final String eventMeshIp = properties.getProperty("eventmesh.ip");
final int eventMeshTcpPort =
Integer.parseInt(properties.getProperty("eventmesh.tcp.port"));
- try {
- UserAgent userAgent = EventMeshTestUtils.generateClient2();
- client = new DefaultEventMeshTCPClient(eventMeshIp,
eventMeshTcpPort, userAgent);
+ UserAgent userAgent = EventMeshTestUtils.generateClient2();
+ EventMeshTcpClientConfig eventMeshTcpClientConfig =
EventMeshTcpClientConfig.builder()
+ .host(eventMeshIp)
+ .port(eventMeshTcpPort)
+ .userAgent(userAgent)
+ .build();
+ try (EventMeshMessageTCPSubClient client = new
EventMeshMessageTCPSubClient(eventMeshTcpClientConfig)) {
client.init();
client.heartbeat();
client.subscribe("TEST-TOPIC-TCP-BROADCAST",
SubscriptionMode.BROADCASTING, SubscriptionType.ASYNC);
- client.registerSubBusiHandler(handler);
+ client.registerBusiHandler(handler);
client.listen();
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncResponse.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/SyncResponse.java
similarity index 76%
rename from
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncResponse.java
rename to
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/SyncResponse.java
index 026fa11..b796c4b 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncResponse.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/SyncResponse.java
@@ -15,11 +15,11 @@
* limitations under the License.
*/
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.sub.eventmeshmessage;
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPSubClient;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
@@ -36,15 +36,19 @@ public class SyncResponse implements
ReceiveMsgHook<EventMeshMessage> {
public static SyncResponse handler = new SyncResponse();
public static void main(String[] agrs) throws Exception {
- try {
- UserAgent userAgent = EventMeshTestUtils.generateClient2();
- EventMeshTCPClient client = new
DefaultEventMeshTCPClient("127.0.0.1", 10000, userAgent);
+ UserAgent userAgent = EventMeshTestUtils.generateClient2();
+ EventMeshTcpClientConfig eventMeshTcpClientConfig =
EventMeshTcpClientConfig.builder()
+ .host("127.0.0.1")
+ .port(10000)
+ .userAgent(userAgent)
+ .build();
+ try (EventMeshMessageTCPSubClient client = new
EventMeshMessageTCPSubClient(eventMeshTcpClientConfig)) {
client.init();
client.heartbeat();
client.subscribe("TEST-TOPIC-TCP-SYNC",
SubscriptionMode.CLUSTERING, SubscriptionType.SYNC);
// Synchronize RR messages
- client.registerSubBusiHandler(handler);
+ client.registerBusiHandler(handler);
client.listen();
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/http/producer/OpenMessageProducer.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/http/producer/OpenMessageProducer.java
index b3cadcc..4dcffff 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/http/producer/OpenMessageProducer.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/http/producer/OpenMessageProducer.java
@@ -26,6 +26,10 @@ import lombok.extern.slf4j.Slf4j;
@Slf4j
class OpenMessageProducer extends AbstractHttpClient implements
EventMeshProtocolProducer<Message> {
+ private static final String PROTOCOL_TYPE = "openmessage";
+
+ private static final String PROTOCOL_DESC = "http";
+
public OpenMessageProducer(EventMeshHttpClientConfig
eventMeshHttpClientConfig)
throws EventMeshException {
super(eventMeshHttpClientConfig);
@@ -94,8 +98,9 @@ class OpenMessageProducer extends AbstractHttpClient
implements EventMeshProtoco
requestParam
.addHeader(ProtocolKey.ClientInstanceKey.USERNAME,
eventMeshHttpClientConfig.getUserName())
.addHeader(ProtocolKey.ClientInstanceKey.PASSWD,
eventMeshHttpClientConfig.getPassword())
- .addHeader(ProtocolKey.VERSION, ProtocolVersion.V1.getVersion())
.addHeader(ProtocolKey.LANGUAGE, Constants.LANGUAGE_JAVA)
+ .addHeader(ProtocolKey.PROTOCOL_TYPE, PROTOCOL_TYPE)
+ .addHeader(ProtocolKey.PROTOCOL_DESC, PROTOCOL_DESC)
// todo: add producerGroup to header, set protocol type, protocol
version
.addBody(SendMessageRequestBody.PRODUCERGROUP,
eventMeshHttpClientConfig.getProducerGroup())
.addBody(SendMessageRequestBody.CONTENT,
JsonUtils.serialize(openMessage));
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/MessageUtils.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/MessageUtils.java
index e799fb7..0b96cab 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/MessageUtils.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/MessageUtils.java
@@ -17,19 +17,23 @@
package org.apache.eventmesh.client.tcp.common;
+import org.apache.eventmesh.common.Constants;
+import org.apache.eventmesh.common.protocol.SubscriptionItem;
+import org.apache.eventmesh.common.protocol.SubscriptionMode;
+import org.apache.eventmesh.common.protocol.SubscriptionType;
+import org.apache.eventmesh.common.protocol.tcp.Command;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
+import org.apache.eventmesh.common.protocol.tcp.Header;
+import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.common.protocol.tcp.Subscription;
+import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ThreadLocalRandom;
import io.cloudevents.CloudEvent;
import io.cloudevents.SpecVersion;
-import org.apache.eventmesh.common.Constants;
-import org.apache.eventmesh.common.protocol.SubscriptionType;
-import org.apache.eventmesh.common.protocol.tcp.Subscription;
-import org.apache.eventmesh.common.protocol.SubscriptionItem;
-import org.apache.eventmesh.common.protocol.SubscriptionMode;
-import org.apache.eventmesh.common.protocol.tcp.*;
-import org.apache.eventmesh.common.protocol.tcp.Package;
public class MessageUtils {
private static final int seqLength = 10;
@@ -59,7 +63,8 @@ public class MessageUtils {
return msg;
}
- public static Package subscribe(String topic, SubscriptionMode
subscriptionMode, SubscriptionType subscriptionType) {
+ public static Package subscribe(String topic, SubscriptionMode
subscriptionMode,
+ SubscriptionType subscriptionType) {
Package msg = new Package();
msg.setHeader(new Header(Command.SUBSCRIBE_REQUEST, 0, null,
generateRandomString(seqLength)));
msg.setBody(generateSubscription(topic, subscriptionMode,
subscriptionType));
@@ -81,18 +86,13 @@ public class MessageUtils {
public static Package buildPackage(Object message, Command command) {
Package msg = new Package();
- msg.setHeader(new Header(command, 0,
- null, generateRandomString(seqLength)));
+ msg.setHeader(new Header(command, 0, null,
generateRandomString(seqLength)));
if (message instanceof CloudEvent) {
- msg.getHeader().putProperty(Constants.PROTOCOL_TYPE,
- EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME);
- msg.getHeader().putProperty(Constants.PROTOCOL_VERSION,
- ((CloudEvent) message).getSpecVersion().toString());
+ msg.getHeader().putProperty(Constants.PROTOCOL_TYPE,
EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME);
+ msg.getHeader().putProperty(Constants.PROTOCOL_VERSION,
((CloudEvent) message).getSpecVersion().toString());
} else if (message instanceof EventMeshMessage) {
- msg.getHeader().putProperty(Constants.PROTOCOL_TYPE,
- EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME);
- msg.getHeader().putProperty(Constants.PROTOCOL_VERSION,
- SpecVersion.V1.toString());
+ msg.getHeader().putProperty(Constants.PROTOCOL_TYPE,
EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME);
+ msg.getHeader().putProperty(Constants.PROTOCOL_VERSION,
SpecVersion.V1.toString());
} else {
// unsupported protocol for server
return msg;
@@ -159,7 +159,8 @@ public class MessageUtils {
.build();
}
- private static Subscription generateSubscription(String topic,
SubscriptionMode subscriptionMode, SubscriptionType subscriptionType) {
+ private static Subscription generateSubscription(String topic,
SubscriptionMode subscriptionMode,
+ SubscriptionType
subscriptionType) {
Subscription subscription = new Subscription();
List<SubscriptionItem> subscriptionItems = new ArrayList<>();
subscriptionItems.add(new SubscriptionItem(topic, subscriptionMode,
subscriptionType));
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
index 1f4e774..baa5013 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
@@ -9,9 +9,9 @@ import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
import org.apache.eventmesh.client.tcp.common.RequestContext;
import org.apache.eventmesh.client.tcp.common.TcpClient;
import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
-import org.apache.eventmesh.common.EventMeshMessage;
import org.apache.eventmesh.common.exception.EventMeshException;
import org.apache.eventmesh.common.protocol.tcp.Command;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.Package;
import org.apache.eventmesh.common.protocol.tcp.UserAgent;
@@ -82,6 +82,7 @@ public class EventMeshMessageTCPPubClient extends TcpClient
implements EventMesh
}
}
+ // todo: Maybe use org.apache.eventmesh.common.EvetMesh here is better
@Override
public Package rr(EventMeshMessage eventMeshMessage, long timeout) throws
EventMeshException {
try {
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
index cdce845..5891e4b 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
@@ -7,12 +7,12 @@ import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
import org.apache.eventmesh.client.tcp.common.RequestContext;
import org.apache.eventmesh.client.tcp.common.TcpClient;
import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
-import org.apache.eventmesh.common.EventMeshMessage;
import org.apache.eventmesh.common.exception.EventMeshException;
import org.apache.eventmesh.common.protocol.SubscriptionItem;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
import org.apache.eventmesh.common.protocol.tcp.Command;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.Package;
import org.apache.eventmesh.common.protocol.tcp.UserAgent;
diff --git
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
index 0b26161..30e7142 100644
---
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
+++
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
@@ -17,14 +17,15 @@
package org.apache.eventmesh.client.http.demo;
-import org.apache.commons.lang3.StringUtils;
-
import org.apache.eventmesh.client.http.conf.EventMeshHttpClientConfig;
import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
-import org.apache.eventmesh.common.utils.IPUtils;
import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.utils.IPUtils;
import org.apache.eventmesh.common.utils.RandomStringUtils;
import org.apache.eventmesh.common.utils.ThreadUtils;
+
+import org.apache.commons.lang3.StringUtils;
+
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -45,38 +46,36 @@ public class SyncRequestInstance {
eventMeshIPPort = "127.0.0.1:10105";
}
- EventMeshHttpClientConfig eventMeshClientConfig = new
EventMeshHttpClientConfig();
- eventMeshClientConfig.setLiteEventMeshAddr(eventMeshIPPort)
- .setProducerGroup("EventMeshTest-producerGroup")
- .setEnv("env")
- .setIdc("idc")
- .setIp(IPUtils.getLocalAddress())
- .setSys("1234")
- .setPid(String.valueOf(ThreadUtils.getPID()));
+ EventMeshHttpClientConfig eventMeshClientConfig =
EventMeshHttpClientConfig.builder()
+ .liteEventMeshAddr(eventMeshIPPort)
+ .producerGroup("EventMeshTest-producerGroup")
+ .env("env")
+ .idc("idc")
+ .ip(IPUtils.getLocalAddress())
+ .sys("1234")
+ .pid(String.valueOf(ThreadUtils.getPID())).build();
eventMeshHttpProducer = new
EventMeshHttpProducer(eventMeshClientConfig);
- eventMeshHttpProducer.start();
long startTime = System.currentTimeMillis();
- EventMeshMessage eventMeshMessage = new EventMeshMessage();
- eventMeshMessage.setBizSeqNo(RandomStringUtils.generateNum(30))
- .setContent("contentStr with special protocal")
- .setTopic(topic)
- .setUniqueId(RandomStringUtils.generateNum(30));
+ EventMeshMessage eventMeshMessage = EventMeshMessage.builder()
+ .bizSeqNo(RandomStringUtils.generateNum(30))
+ .content("contentStr with special protocal")
+ .topic(topic)
+ .uniqueId(RandomStringUtils.generateNum(30)).build();
EventMeshMessage rsp =
eventMeshHttpProducer.request(eventMeshMessage, 10000);
if (logger.isDebugEnabled()) {
- logger.debug("sendmsg : {}, return : {}, cost:{}ms",
eventMeshMessage.getContent(), rsp.getContent(), System.currentTimeMillis() -
startTime);
+ logger.debug("sendmsg : {}, return : {}, cost:{}ms",
eventMeshMessage.getContent(), rsp.getContent(),
+ System.currentTimeMillis() - startTime);
}
} catch (Exception e) {
logger.warn("send msg failed", e);
}
- try {
- Thread.sleep(30000);
- if (eventMeshHttpProducer != null) {
- eventMeshHttpProducer.shutdown();
- }
+ Thread.sleep(30000);
+ try (final EventMeshHttpProducer closed = eventMeshHttpProducer){
+ // close producer
} catch (Exception e1) {
logger.warn("producer shutdown exception", e1);
}
diff --git
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/util/HttpLoadBalanceUtilsTest.java
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/util/HttpLoadBalanceUtilsTest.java
index 4c03ee1..aabafb1 100644
---
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/util/HttpLoadBalanceUtilsTest.java
+++
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/util/HttpLoadBalanceUtilsTest.java
@@ -28,8 +28,9 @@ public class HttpLoadBalanceUtilsTest {
@Test
public void testCreateRandomSelector() throws EventMeshException {
- EventMeshHttpClientConfig eventMeshHttpClientConfig = new
EventMeshHttpClientConfig()
- .setLiteEventMeshAddr("127.0.0.1:1001;127.0.0.2:1002");
+ EventMeshHttpClientConfig eventMeshHttpClientConfig =
EventMeshHttpClientConfig.builder()
+ .liteEventMeshAddr("127.0.0.1:1001;127.0.0.2:1002")
+ .build();
LoadBalanceSelector<String> randomSelector = HttpLoadBalanceUtils
.createEventMeshServerLoadBalanceSelector(eventMeshHttpClientConfig);
Assert.assertEquals(LoadBalanceType.RANDOM, randomSelector.getType());
@@ -37,9 +38,9 @@ public class HttpLoadBalanceUtilsTest {
@Test
public void testCreateWeightRoundRobinSelector() throws EventMeshException
{
- EventMeshHttpClientConfig eventMeshHttpClientConfig = new
EventMeshHttpClientConfig()
- .setLiteEventMeshAddr("127.0.0.1:1001:1;127.0.0.2:1001:2")
- .setLoadBalanceType(LoadBalanceType.WEIGHT_ROUND_ROBIN);
+ EventMeshHttpClientConfig eventMeshHttpClientConfig =
EventMeshHttpClientConfig.builder()
+ .liteEventMeshAddr("127.0.0.1:1001:1;127.0.0.2:1001:2")
+ .loadBalanceType(LoadBalanceType.WEIGHT_ROUND_ROBIN).build();
LoadBalanceSelector<String> weightRoundRobinSelector =
HttpLoadBalanceUtils
.createEventMeshServerLoadBalanceSelector(eventMeshHttpClientConfig);
Assert.assertEquals(LoadBalanceType.WEIGHT_ROUND_ROBIN,
weightRoundRobinSelector.getType());
@@ -47,9 +48,9 @@ public class HttpLoadBalanceUtilsTest {
@Test
public void testCreateWeightRandomSelector() throws EventMeshException {
- EventMeshHttpClientConfig eventMeshHttpClientConfig = new
EventMeshHttpClientConfig()
- .setLiteEventMeshAddr("127.0.0.1:1001:1;127.0.0.2:1001:2")
- .setLoadBalanceType(LoadBalanceType.WEIGHT_RANDOM);
+ EventMeshHttpClientConfig eventMeshHttpClientConfig =
EventMeshHttpClientConfig.builder()
+ .liteEventMeshAddr("127.0.0.1:1001:1;127.0.0.2:1001:2")
+ .loadBalanceType(LoadBalanceType.WEIGHT_RANDOM).build();
LoadBalanceSelector<String> weightRoundRobinSelector =
HttpLoadBalanceUtils
.createEventMeshServerLoadBalanceSelector(eventMeshHttpClientConfig);
Assert.assertEquals(LoadBalanceType.WEIGHT_RANDOM,
weightRoundRobinSelector.getType());
diff --git
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/common/EventMeshTestUtils.java
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/common/EventMeshTestUtils.java
index 78e1016..f6397b8 100644
---
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/common/EventMeshTestUtils.java
+++
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/common/EventMeshTestUtils.java
@@ -20,53 +20,48 @@ package org.apache.eventmesh.client.tcp.common;
import static
org.apache.eventmesh.client.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_SyncSubscribeTest;
import static
org.apache.eventmesh.client.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_WQ2ClientBroadCast;
import static
org.apache.eventmesh.client.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_WQ2ClientUniCast;
-
-import java.util.concurrent.ThreadLocalRandom;
+import static
org.apache.eventmesh.common.protocol.tcp.Command.RESPONSE_TO_SERVER;
import org.apache.eventmesh.common.protocol.tcp.Command;
import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.Header;
-import org.apache.eventmesh.common.protocol.tcp.UserAgent;
import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.common.protocol.tcp.UserAgent;
-import static
org.apache.eventmesh.common.protocol.tcp.Command.RESPONSE_TO_SERVER;
-
-
-
+import java.util.concurrent.ThreadLocalRandom;
public class EventMeshTestUtils {
private static final int seqLength = 10;
public static UserAgent generateClient1() {
- UserAgent user = new UserAgent();
- user.setHost("127.0.0.1");
- user.setPassword(generateRandomString(8));
- user.setUsername("PU4283");
- user.setConsumerGroup("EventmeshTest-ConsumerGroup");
- user.setProducerGroup("EventmeshTest-ProducerGroup");
- user.setPath("/data/app/umg_proxy");
- user.setPort(8362);
- user.setSubsystem("5023");
- user.setPid(32893);
- user.setVersion("2.0.11");
- user.setIdc("FT");
- return user;
+ return UserAgent.builder()
+ .host("127.0.0.1")
+ .password(generateRandomString(8))
+ .username("PU4283")
+ .consumerGroup("EventmeshTest-ConsumerGroup")
+ .producerGroup("EventmeshTest-ProducerGroup")
+ .path("/data/app/umg_proxy")
+ .port(8362)
+ .subsystem("5023")
+ .pid(32893)
+ .version("2.0.11")
+ .idc("FT")
+ .build();
}
public static UserAgent generateClient2() {
- UserAgent user = new UserAgent();
- user.setHost("127.0.0.1");
- user.setPassword(generateRandomString(8));
- user.setUsername("PU4283");
- user.setConsumerGroup("EventmeshTest-ConsumerGroup");
- user.setProducerGroup("EventmeshTest-ProducerGroup");
- user.setPath("/data/app/umg_proxy");
- user.setPort(9362);
- user.setSubsystem("5017");
- user.setPid(42893);
- user.setVersion("2.0.11");
- user.setIdc("FT");
- return user;
+ return UserAgent.builder()
+ .host("127.0.0.1")
+ .password(generateRandomString(8))
+ .username("PU4283")
+ .consumerGroup("EventmeshTest-ConsumerGroup")
+ .producerGroup("EventmeshTest-ProducerGroup")
+ .path("/data/app/umg_proxy")
+ .port(9362)
+ .subsystem("5017")
+ .pid(42893)
+ .version("2.0.11")
+ .idc("FT").build();
}
public static Package syncRR() {
@@ -104,7 +99,7 @@ public class EventMeshTestUtils {
return msg;
}
- private static EventMeshMessage generateSyncRRMqMsg() {
+ public static EventMeshMessage generateSyncRRMqMsg() {
EventMeshMessage mqMsg = new EventMeshMessage();
mqMsg.setTopic(TOPIC_PRX_SyncSubscribeTest);
mqMsg.getProperties().put("msgType", "persistent");
diff --git
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/demo/SyncRequest.java
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/demo/SyncRequest.java
index 272ece0..7b6b4ca 100644
---
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/demo/SyncRequest.java
+++
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/demo/SyncRequest.java
@@ -17,12 +17,14 @@
package org.apache.eventmesh.client.tcp.demo;
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.client.tcp.common.EventMeshTestUtils;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
-import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPPubClient;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -30,18 +32,23 @@ public class SyncRequest {
public static Logger logger = LoggerFactory.getLogger(SyncRequest.class);
- private static EventMeshTCPClient client;
+ private static EventMeshMessageTCPPubClient client;
- public static void main(String[] agrs) throws Exception {
+ public static void main(String[] agrs) {
try {
UserAgent userAgent = EventMeshTestUtils.generateClient1();
- client = new DefaultEventMeshTCPClient("127.0.0.1", 10000,
userAgent);
+ EventMeshTcpClientConfig eventMeshTcpClientConfig =
EventMeshTcpClientConfig.builder()
+ .host("127.0.0.1")
+ .port(10000)
+ .userAgent(userAgent)
+ .build();
+ client = new
EventMeshMessageTCPPubClient(eventMeshTcpClientConfig);
client.init();
client.heartbeat();
- Package rrMsg = EventMeshTestUtils.syncRR();
- logger.info("begin send rr msg=================={}", rrMsg);
- Package response = client.rr(rrMsg,
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
+ EventMeshMessage eventMeshMessage =
EventMeshTestUtils.generateSyncRRMqMsg();
+ logger.info("begin send rr msg=================={}",
eventMeshMessage);
+ Package response = client.rr(eventMeshMessage,
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
logger.info("receive rr reply==================={}", response);
// release resource and close client
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]