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 deee538  [Feature #562] support cloudevents api in 
eventmesh-connector-api module (#578)
deee538 is described below

commit deee538f5ed0f65f16f387de42e68421b2dd7f68
Author: mike_xwm <[email protected]>
AuthorDate: Wed Nov 10 10:23:32 2021 +0800

    [Feature #562] support cloudevents api in eventmesh-connector-api module 
(#578)
    
    * [Feature #564] Support CloudEvents protocols for pub/sub in 
EventMesh-feature design
    
    * support cloudevents api in eventmesh-connector-api module
    
    * fix checkStyle
    
    * fix checkStyle
    
    * fix checkStyle
    
    * 1.support LifeCycle.java
    2.update Consumer and Producer
    
    * fix remove the extra blank line
---
 .../eventmesh-connector-api/build.gradle           |  1 +
 ...entMeshAction.java => AsyncConsumeContext.java} |  8 ++--
 .../{EventMeshAction.java => EventListener.java}   | 20 ++++++--
 .../org/apache/eventmesh/api/EventMeshAction.java  |  1 +
 .../api/EventMeshAsyncConsumeContext.java          |  1 +
 .../api/{EventMeshAction.java => LifeCycle.java}   | 19 ++++++--
 .../{EventMeshAction.java => SendCallback.java}    | 16 +++++--
 .../api/{EventMeshAction.java => SendResult.java}  | 28 +++++++++--
 .../Consumer.java}                                 | 40 +++++++++-------
 .../ConnectorRuntimeException.java}                | 23 +++++++--
 .../OnExceptionContext.java}                       | 44 ++++++++++++------
 .../apache/eventmesh/api/producer/Producer.java    | 54 ++++++++++++++++++++++
 12 files changed, 199 insertions(+), 56 deletions(-)

diff --git a/eventmesh-connector-plugin/eventmesh-connector-api/build.gradle 
b/eventmesh-connector-plugin/eventmesh-connector-api/build.gradle
index 19cb54c..a721189 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-api/build.gradle
+++ b/eventmesh-connector-plugin/eventmesh-connector-api/build.gradle
@@ -18,6 +18,7 @@
 dependencies {
     implementation project(":eventmesh-spi")
     implementation project(":eventmesh-common")
+    api 'io.cloudevents:cloudevents-core'
     api 'io.openmessaging:openmessaging-api'
     api 'io.dropwizard.metrics:metrics-core'
     api "io.dropwizard.metrics:metrics-healthchecks"
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/AsyncConsumeContext.java
similarity index 89%
copy from 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/AsyncConsumeContext.java
index 4fda6d0..7c1c739 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/AsyncConsumeContext.java
@@ -14,12 +14,12 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.eventmesh.api;
 
-public enum EventMeshAction {
-    CommitMessage,
 
-    ReconsumeLater,
+public abstract class AsyncConsumeContext {
+
+    public abstract void commit(EventMeshAction action);
 
-    ManualAck
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventListener.java
similarity index 67%
copy from 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventListener.java
index 4fda6d0..eede41c 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventListener.java
@@ -14,12 +14,24 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.eventmesh.api;
 
-public enum EventMeshAction {
-    CommitMessage,
 
-    ReconsumeLater,
+import io.cloudevents.CloudEvent;
+
+/**
+ * Event listener, registered for consume messages by consumer.
+ *
+ * <p>
+ * <strong>
+ * Thread safe requirements: this interface will be invoked by multi threads,
+ * so users should keep thread safe during the consume process.
+ * </strong>
+ * </p>
+ */
+public interface EventListener {
+
+    void consume(final CloudEvent cloudEvent, final AsyncConsumeContext 
context);
 
-    ManualAck
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
index 4fda6d0..f783201 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
@@ -14,6 +14,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.eventmesh.api;
 
 public enum EventMeshAction {
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
index c7e4e7f..c123ca2 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
@@ -14,6 +14,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.eventmesh.api;
 
 import io.openmessaging.api.Action;
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/LifeCycle.java
similarity index 69%
copy from 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/LifeCycle.java
index 4fda6d0..f5d61de 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/LifeCycle.java
@@ -14,12 +14,23 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.eventmesh.api;
 
-public enum EventMeshAction {
-    CommitMessage,
+import org.apache.eventmesh.api.consumer.Consumer;
+import org.apache.eventmesh.api.producer.Producer;
+
+/**
+ * The {@code LifeCycle} defines a lifecycle interface for a OMS related 
service endpoint,
+ * like {@link Producer}, {@link Consumer}, and so on.
+ */
+public interface LifeCycle {
+
+    boolean isStarted();
+
+    boolean isClosed();
 
-    ReconsumeLater,
+    void start();
 
-    ManualAck
+    void shutdown();
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendCallback.java
similarity index 68%
copy from 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendCallback.java
index 4fda6d0..c955f5d 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendCallback.java
@@ -14,12 +14,20 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.eventmesh.api;
 
-public enum EventMeshAction {
-    CommitMessage,
+import org.apache.eventmesh.api.exception.OnExceptionContext;
+import org.apache.eventmesh.api.producer.Producer;
+
+import io.cloudevents.CloudEvent;
+
+/**
+ * Call back interface used in {@link Producer#sendAsync(CloudEvent, 
SendCallback)}.
+ */
+public interface SendCallback {
 
-    ReconsumeLater,
+    void onSuccess(final SendResult sendResult);
 
-    ManualAck
+    void onException(final OnExceptionContext context);
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendResult.java
similarity index 62%
copy from 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendResult.java
index 4fda6d0..1a68132 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendResult.java
@@ -14,12 +14,32 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.eventmesh.api;
 
-public enum EventMeshAction {
-    CommitMessage,
+public class SendResult {
+    private String messageId;
+
+    private String topic;
+
+    public String getMessageId() {
+        return messageId;
+    }
+
+    public void setMessageId(String messageId) {
+        this.messageId = messageId;
+    }
+
+    public String getTopic() {
+        return topic;
+    }
 
-    ReconsumeLater,
+    public void setTopic(String topic) {
+        this.topic = topic;
+    }
 
-    ManualAck
+    @Override
+    public String toString() {
+        return "SendResult[topic=" + topic + ", messageId=" + messageId + ']';
+    }
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
similarity index 50%
copy from 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
copy to 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
index c7e4e7f..a6fafe0 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
@@ -14,27 +14,33 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
-package org.apache.eventmesh.api;
 
-import io.openmessaging.api.Action;
-import io.openmessaging.api.AsyncConsumeContext;
+package org.apache.eventmesh.api.consumer;
 
-public abstract class EventMeshAsyncConsumeContext extends AsyncConsumeContext 
{
+import org.apache.eventmesh.api.AbstractContext;
+import org.apache.eventmesh.api.EventListener;
+import org.apache.eventmesh.api.LifeCycle;
+import org.apache.eventmesh.spi.EventMeshExtensionType;
+import org.apache.eventmesh.spi.EventMeshSPI;
 
-    private AbstractContext abstractContext;
+import java.util.List;
+import java.util.Properties;
 
-    public AbstractContext getAbstractContext() {
-        return abstractContext;
-    }
+import io.cloudevents.CloudEvent;
 
-    public void setAbstractContext(AbstractContext abstractContext) {
-        this.abstractContext = abstractContext;
-    }
 
-    public abstract void commit(EventMeshAction action);
 
-    @Override
-    public void commit(Action action) {
-        throw new UnsupportedOperationException("not support yet");
-    }
-}
\ No newline at end of file
+/**
+ * Consumer Interface.
+ */
+@EventMeshSPI(isSingleton = false, eventMeshExtensionType = 
EventMeshExtensionType.CONNECTOR)
+public interface Consumer extends LifeCycle {
+
+    void init(Properties keyValue) throws Exception;
+
+    void updateOffset(List<CloudEvent> cloudEvents, AbstractContext context);
+
+    void subscribe(String topic, final EventListener listener) throws 
Exception;
+
+    void unsubscribe(String topic);
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/ConnectorRuntimeException.java
similarity index 63%
copy from 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/ConnectorRuntimeException.java
index 4fda6d0..0efe513 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/ConnectorRuntimeException.java
@@ -14,12 +14,25 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
-package org.apache.eventmesh.api;
 
-public enum EventMeshAction {
-    CommitMessage,
+package org.apache.eventmesh.api.exception;
 
-    ReconsumeLater,
+public class ConnectorRuntimeException extends RuntimeException {
+
+    public ConnectorRuntimeException() {
+
+    }
+
+    public ConnectorRuntimeException(String message) {
+        super(message);
+    }
+
+    public ConnectorRuntimeException(Throwable throwable) {
+        super(throwable);
+    }
+
+    public ConnectorRuntimeException(String message, Throwable throwable) {
+        super(message, throwable);
+    }
 
-    ManualAck
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/OnExceptionContext.java
similarity index 53%
copy from 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
copy to 
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/OnExceptionContext.java
index c7e4e7f..1237e5f 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/OnExceptionContext.java
@@ -14,27 +14,43 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
-package org.apache.eventmesh.api;
 
-import io.openmessaging.api.Action;
-import io.openmessaging.api.AsyncConsumeContext;
+package org.apache.eventmesh.api.exception;
 
-public abstract class EventMeshAsyncConsumeContext extends AsyncConsumeContext 
{
+public class OnExceptionContext {
 
-    private AbstractContext abstractContext;
+    private String messageId;
 
-    public AbstractContext getAbstractContext() {
-        return abstractContext;
+    private String topic;
+
+    /**
+     * Detailed exception stack information.
+     */
+    private ConnectorRuntimeException exception;
+
+
+    public String getMessageId() {
+        return messageId;
+    }
+
+    public void setMessageId(String messageId) {
+        this.messageId = messageId;
     }
 
-    public void setAbstractContext(AbstractContext abstractContext) {
-        this.abstractContext = abstractContext;
+    public String getTopic() {
+        return topic;
     }
 
-    public abstract void commit(EventMeshAction action);
+    public void setTopic(String topic) {
+        this.topic = topic;
+    }
+
+
+    public ConnectorRuntimeException getException() {
+        return exception;
+    }
 
-    @Override
-    public void commit(Action action) {
-        throw new UnsupportedOperationException("not support yet");
+    public void setException(ConnectorRuntimeException exception) {
+        this.exception = exception;
     }
-}
\ No newline at end of file
+}
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
new file mode 100644
index 0000000..08d8998
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
@@ -0,0 +1,54 @@
+/*
+ * 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.api.producer;
+
+import org.apache.eventmesh.api.LifeCycle;
+import org.apache.eventmesh.api.RRCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.spi.EventMeshExtensionType;
+import org.apache.eventmesh.spi.EventMeshSPI;
+
+import java.util.Properties;
+
+import io.cloudevents.CloudEvent;
+
+/**
+ * Producer Interface.
+ */
+@EventMeshSPI(isSingleton = false, eventMeshExtensionType = 
EventMeshExtensionType.CONNECTOR)
+public interface Producer extends LifeCycle {
+
+    void init(Properties properties) throws Exception;
+
+    SendResult publish(final CloudEvent cloudEvent);
+
+    void publish(CloudEvent cloudEvent, SendCallback sendCallback) throws 
Exception;
+
+    void sendOneway(final CloudEvent cloudEvent);
+
+    void sendAsync(final CloudEvent cloudEvent, final SendCallback 
sendCallback);
+
+    void request(CloudEvent cloudEvent, RRCallback rrCallback, long timeout) 
throws Exception;
+
+    boolean reply(final CloudEvent cloudEvent, final SendCallback 
sendCallback) throws Exception;
+
+    void checkTopicExist(String topic) throws Exception;
+
+    void setExtFields();
+}

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to