+ * A {@code Client} interface contains all the common behaviors of producer and consumer. which can be used to achieve + * some basic interaction with the server. + *
+ * + * @version OMS 1.0.0 + * @since OMS 1.0.0 + */ +public interface Client { + /** + * Get the extension method, and this interface is optional, Therefore, users need to check whether this interface + * has been implemented by vendors. + *
+ *
+ * @return the implementation of {@link Extension}
+ */
+ @io.openmessaging.annotation.Optional
+ Optional After this method returns, subsequent calls to {@link #isDone} will
+ * always return {@code true}. Subsequent calls to {@link #isCancelled}
+ * will always return {@code true} if this method returned {@code true}.
+ *
+ * @param mayInterruptIfRunning {@code true} if the thread executing this
+ * task should be interrupted; otherwise, in-progress tasks are allowed
+ * to complete
+ * @return {@code false} if the task could not be cancelled,
+ * typically because it has already completed normally;
+ * {@code true} otherwise
+ */
+ boolean cancel(boolean mayInterruptIfRunning);
+
/**
* Returns {@code true} if this task was cancelled before it completed normally.
*
diff --git a/openmessaging-api/src/main/java/io/openmessaging/KeyValue.java b/openmessaging-api/src/main/java/io/openmessaging/KeyValue.java
index d521bb6e..9c617f2f 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/KeyValue.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/KeyValue.java
@@ -34,6 +34,15 @@
* @since OMS 1.0.0
*/
public interface KeyValue {
+
+ /**
+ * Inserts or replaces {@code boolean} value for the specified key.
+ *
+ * @param key the key to be placed into this {@code KeyValue} object
+ * @param value the value corresponding to key
+ */
+ KeyValue put(String key, boolean value);
+
/**
* Inserts or replaces {@code short} value for the specified key.
*
@@ -74,6 +83,27 @@ public interface KeyValue {
*/
KeyValue put(String key, String value);
+ /**
+ * Searches for the {@code boolean} property with the specified key in this {@code KeyValue} object. If the key is
+ * not found in this property list, false is returned.
+ *
+ * @param key the property key
+ * @return the value in this {@code KeyValue} object with the specified key value
+ * @see #put(String, boolean)
+ */
+ boolean getBoolean(String key);
+
+ /**
+ * Searches for the {@code boolean} property with the specified key in this {@code KeyValue} object. If the key is
+ * not found in this property list, false is returned.
+ *
+ * @param key the property key
+ * @param defaultValue a default value
+ * @return the value in this {@code KeyValue} object with the specified key value
+ * @see #put(String, boolean)
+ */
+ boolean getBoolean(String key, boolean defaultValue);
+
/**
* Searches for the {@code short} property with the specified key in this {@code KeyValue} object. If the key is not
* found in this property list, zero is returned.
@@ -82,7 +112,7 @@ public interface KeyValue {
* @return the value in this {@code KeyValue} object with the specified key value
* @see #put(String, short)
*/
- int getShort(String key);
+ short getShort(String key);
/**
* Searches for the {@code short} property with the specified key in this {@code KeyValue} object. If the key is not
@@ -93,7 +123,7 @@ public interface KeyValue {
* @return the value in this {@code KeyValue} object with the specified key value
* @see #put(String, short)
*/
- int getShort(String key, short defaultValue);
+ short getShort(String key, short defaultValue);
/**
* Searches for the {@code int} property with the specified key in this {@code KeyValue} object. If the key is not
diff --git a/openmessaging-api/src/main/java/io/openmessaging/Message.java b/openmessaging-api/src/main/java/io/openmessaging/Message.java
deleted file mode 100644
index 9d06f889..00000000
--- a/openmessaging-api/src/main/java/io/openmessaging/Message.java
+++ /dev/null
@@ -1,363 +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 io.openmessaging;
-
-import io.openmessaging.exception.OMSMessageFormatException;
-
-/**
- * The {@code Message} interface is the root interface of all OMS messages, and the most commonly used OMS message is
- * {@link Message}.
- *
- * Most message-oriented middleware (MOM) products treat messages as lightweight entities that consist of header and
- * body and is used by separate applications to exchange a piece of information, like Apache RocketMQ.
- *
- * The header contains fields used by the messaging system that describes the message's meta information, while the body
- * contains the application data being transmitted.
- *
- * As for the message header, OMS defines two kinds types: headers {@link Headers} and properties {@link KeyValue}, with
- * respect to flexibility in vendor implementation and user usage.
- *
- * In BytesMessage, the body is just a byte array, may be compressed and uncompressed in the transmitting process by the
- * messaging system. The application is responsible for explaining the concrete content and format of the message body,
- * OMS is never aware of that.
- *
- * The body part is placed in the implementation classes of {@code Message}.
- *
- * @version OMS 1.0.0
- * @since OMS 1.0.0
- */
-public interface Message {
-
- interface Headers {
-
- /**
- * The {@code DESTINATION} header field contains the destination to which the message is being sent.
- *
- * When a message is sent this value is set to the right {@code Queue}, then the message will be sent to the
- * specified destination.
- *
- * When a message is received, its destination is equivalent to the {@code Queue} where the message resides in.
- */
- Headers setDestination(String destination);
-
- /**
- * The {@code MESSAGE_ID} header field contains a value that uniquely identifies each message sent by a {@code
- * Producer}.
- */
- Headers setMessageId(String messageId);
-
- /**
- * The {@code BORN_TIMESTAMP} header field contains the time a message was handed off to a {@code Producer} to
- * be sent.
- *
- * When a message is sent, BORN_TIMESTAMP will be set with current timestamp as the born timestamp of a message
- * in client side, on return from the send method, the message's BORN_TIMESTAMP header field contains this
- * value.
- *
- * When a message is received its, BORN_TIMESTAMP header field contains this same value.
- *
- * This filed is a {@code long} value, measured in milliseconds.
- */
- Headers setBornTimestamp(long bornTimestamp);
-
- /**
- * The {@code BORN_HOST} header field contains the born host info of a message in client side.
- *
- * When a message is sent, BORN_HOST will be set with the local host info, on return from the send method, the
- * message's BORN_HOST header field contains this value.
- *
- * When a message is received, its BORN_HOST header field contains this same value.
- */
- Headers setBornHost(String bornHost);
-
- /**
- * The {@code STORE_TIMESTAMP} header field contains the store timestamp of a message in server side.
- *
- * When a message is sent, STORE_TIMESTAMP is ignored.
- *
- * When the send method returns it contains a server-assigned value.
- *
- * This filed is a {@code long} value, measured in milliseconds.
- */
- Headers setStoreTimestamp(long storeTimestamp);
-
- /**
- * The {@code STORE_HOST} header field contains the store host info of a message in server side.
- *
- * When a message is sent, STORE_HOST is ignored.
- *
- * When the send method returns it contains a server-assigned value.
- */
- Headers setStoreHost(String storeHost);
-
- /**
- * The {@code DELAY_TIME} header field contains a number that represents the delayed times in milliseconds.
- *
- * The {@code EXPIRE_TIME} represents a relative valid interval that a message can be delivered in it. If the
- * EXPIRE_TIME field is specified as zero, that indicates the message does not expire.
- *
- * When an undelivered message's expiration time is reached, the message should be destroyed. OMS does not
- * define a notification of message expiration.
- *
- * OMS defines a ten level priority value with 1 as the lowest priority and 10 as the highest, and the default
- * priority is 5. The priority beyond this region will be ignored.
- *
- * OMS does not require or provide any guarantee that the message should be delivered in priority order
- * strictly, but the vendor should provide a best effort to deliver expedited messages ahead of normal
- * messages.
- *
- * If PRIORITY field isn't set explicitly, use {@code 5} as the default priority.
- */
- Headers setPriority(short priority);
-
- /**
- * The {@code RELIABILITY} header field contains the reliability level of a message, the vendor should guarantee
- * the reliability level for a message.
- *
- * OMS defines two modes of message delivery:
- *
- * This key is a customer identifier for a class of messages, and this key may be used for server to hash or
- * dispatch messages, or even can use this key to implement order message.
- *
- */
- Headers setMessageKey(String messageKey);
-
- /**
- * The {@code TRACE_ID} header field contains the trace ID of a message, which represents a global and unique
- * identification, to associate key events in the whole lifecycle of a message, like sent by who, stored at
- * where, and received by who.
- *
+ * A {@code Optional} is an annotation to mark some certain methods ,interfaces and etc. this annotation represented
+ * these methods or interfaces are not mandatory in OpenMessaging.
+ *
+ * If these methods or interfaces adopted by more and more vendors and end users, they may be become the mandatory
+ * interface in the future. Of course, if they are used very little, they may be removed.
+ *
+ * A message listener should handle different types of {@code BatchMessage}.
+ *
+ * @param batchMessage the received batchMessage.
+ */
+ void onReceived(List
+ * Messages that have been received but not acknowledged may be redelivered.
+ *
+ * @throws OMSRuntimeException if the consumer fails to acknowledge the messages due to some internal error.
+ */
+ void success(MessageReceipt... messages);
+
+ /**
+ * Acknowledges all messages in this batch, which is related to this {@code MessageContext}.
+ *
+ *
+ * @throws OMSRuntimeException if the consumer fails to acknowledge the messages due to some internal error.
+ */
+ void ack();
+ }
+}
diff --git a/openmessaging-api/src/main/java/io/openmessaging/consumer/Consumer.java b/openmessaging-api/src/main/java/io/openmessaging/consumer/Consumer.java
index 948ec068..4ea4ea61 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/consumer/Consumer.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/consumer/Consumer.java
@@ -17,106 +17,33 @@
package io.openmessaging.consumer;
-import io.openmessaging.Message;
+import io.openmessaging.Client;
import io.openmessaging.MessagingAccessPoint;
import io.openmessaging.ServiceLifecycle;
import io.openmessaging.exception.OMSDestinationException;
import io.openmessaging.exception.OMSRuntimeException;
import io.openmessaging.exception.OMSSecurityException;
-import io.openmessaging.exception.OMSTimeOutException;
import io.openmessaging.interceptor.ConsumerInterceptor;
+import io.openmessaging.message.Message;
+import java.util.Collection;
+import java.util.List;
+import java.util.Set;
/**
* A {@code PushConsumer} receives messages from multiple queues, these messages are pushed from MOM server to {@code
- * PushConsumer} client.
+ * Consumer} client.
*
* @version OMS 1.0.0
- * @see MessagingAccessPoint#createConsumer().
* @since OMS 1.0.0
*/
-public interface Consumer extends ServiceLifecycle {
+public interface Consumer extends ServiceLifecycle, Client {
/**
- * Resumes the {@code Consumer} in push model after a suspend.
- *
- * This method resumes the {@code Consumer} instance after it was suspended. The instance will not receive new
- * messages between the suspend and resume calls.
- *
- * @throws OMSRuntimeException if the instance has not been suspended.
- * @see Consumer#suspend()
- */
- void resume();
-
- /**
- * Suspends the {@code Consumer} in push model for later resumption.
- *
- * This method suspends the consumer until it is resumed. The consumer will not receive new messages between the
- * suspend and resume calls.
- *
- * This method behaves exactly as if it simply performs the call {@code suspend(0)}.
- *
- * @throws OMSRuntimeException if the instance is not currently running.
- * @see Consumer#resume()
- */
- void suspend();
-
- /**
- * Suspends the {@code Consumer} in push model for later resumption.
- *
- * This method suspends the consumer until it is resumed or a specified amount of time has elapsed. The consumer
- * will not receive new messages during the suspended state.
- *
- * This method is similar to the {@link #suspend()} method, but it allows finer control over the amount of time to
- * suspend, and the consumer will be suspended until it is resumed if the timeout is zero.
- *
- * @param timeout the maximum time to suspend in milliseconds.
- * @throws OMSRuntimeException if the instance is not currently running.
- */
- void suspend(long timeout);
-
- /**
- * This method is used to find out whether the {@code Consumer} in push model is suspended.
- *
- * @return true if this {@code Consumer} is suspended, false otherwise.
- */
- boolean isSuspended();
-
- /**
- * Bind the {@code Consumer} to a specified queue in pull model, user can use {@link Consumer#receive(long)} to get
- * message from bind queue.
- *
- * {@link MessageListener#onReceived(Message, MessageListener.Context)} will be called when new delivered message is
- * coming.
- *
- * @param queueName a specified queue.
- * @throws OMSSecurityException when have no authority to bind to this queue.
- * @throws OMSDestinationException when have no given destination in the server.
- * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
- */
- void bindQueue(String queueName);
-
- /**
- * Bind the {@code Consumer} to a specified queue, with a {@code MessageListener}.
- *
- * {@link MessageListener#onReceived(Message, MessageListener.Context)} will be called when new delivered message is
- * coming.
+ * This method is used to find out the collection of queues bind to {@code Consumer}.
*
- * @param queueName a specified queue.
- * @param listener a specified listener to receive new message.
- * @throws OMSSecurityException when have no authority to bind to this queue.
- * @throws OMSDestinationException when have no given destination in the server.
- * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ * @return the queues this consumer is bind, or null if the consumer is not bind queue.
*/
- void bindQueue(String queueName, MessageListener listener);
-
- /**
- * Unbind the {@code Consumer} from a specified queue.
- *
- * After the success call, this consumer won't receive new message from the specified queue any more.
- *
- * @param queueName a specified queue.
- */
- void unbindQueue(String queueName);
+ Set
- * This call blocks indefinitely until a message is arrives, the timeout expires, or until this {@code PullConsumer}
- * is shut down.
- *
- * @param timeout receive message will blocked at most
* Messages that have been received but not acknowledged may be redelivered.
*
- * @param receiptHandle the receipt handle associated with the consumed message.
+ * @param receipt the receipt handle associated with the consumed message.
*/
- void ack(String receiptHandle);
+ void ack(MessageReceipt receipt);
+
}
\ No newline at end of file
diff --git a/openmessaging-api/src/main/java/io/openmessaging/consumer/MessageListener.java b/openmessaging-api/src/main/java/io/openmessaging/consumer/MessageListener.java
index 97a6f414..81a3fa8e 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/consumer/MessageListener.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/consumer/MessageListener.java
@@ -17,7 +17,7 @@
package io.openmessaging.consumer;
-import io.openmessaging.Message;
+import io.openmessaging.message.Message;
import io.openmessaging.exception.OMSRuntimeException;
/**
diff --git a/openmessaging-api/src/main/java/io/openmessaging/consumer/MessageReceipt.java b/openmessaging-api/src/main/java/io/openmessaging/consumer/MessageReceipt.java
new file mode 100644
index 00000000..26a75395
--- /dev/null
+++ b/openmessaging-api/src/main/java/io/openmessaging/consumer/MessageReceipt.java
@@ -0,0 +1,27 @@
+/*
+ * 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 io.openmessaging.consumer;
+
+/**
+ * A {@code MessageReceipt} is a {@code Message} with a {@code Receipt}.
+ *
+ * @version OMS 1.0.0
+ * @since OMS 1.0.0
+ */
+public interface MessageReceipt {
+}
diff --git a/openmessaging-api/src/main/java/io/openmessaging/consumer/PullConsumer.java b/openmessaging-api/src/main/java/io/openmessaging/consumer/PullConsumer.java
new file mode 100755
index 00000000..82f7814a
--- /dev/null
+++ b/openmessaging-api/src/main/java/io/openmessaging/consumer/PullConsumer.java
@@ -0,0 +1,147 @@
+/*
+ * 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 io.openmessaging.consumer;
+
+import io.openmessaging.MessagingAccessPoint;
+import io.openmessaging.annotation.Optional;
+import io.openmessaging.exception.OMSDestinationException;
+import io.openmessaging.exception.OMSRuntimeException;
+import io.openmessaging.exception.OMSSecurityException;
+import io.openmessaging.exception.OMSTimeOutException;
+import io.openmessaging.extension.QueueMetaData;
+import io.openmessaging.message.Message;
+import java.util.Collection;
+import java.util.List;
+
+/**
+ * A {@code PullConsumer} pulls messages from the specified queue, and supports submit the consume result by
+ * acknowledgement.
+ *
+ * @version OMS 1.0.0
+ * @see MessagingAccessPoint#createPullConsumer()
+ * @since OMS 1.0.0
+ */
+public interface PullConsumer extends Consumer {
+
+ /**
+ * Bind the {@code Consumer} to a collection of queue in pull model, user can use {@link PullConsumer#receive(long)}
+ * to get messages from a collection of queue.
+ *
+ *
+ * @param queueNames a collection of queues.
+ * @throws OMSSecurityException when have no authority to bind to this queue.
+ * @throws OMSDestinationException when have no given destination in the server.
+ * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ */
+ void bindQueue(Collection
+ * After the success call, this consumer won't receive new message from the specified queue any more.
+ *
+ * @param queueNames a collection of queues.
+ */
+ void unbindQueue(Collection
+ * This call blocks indefinitely until a message is arrives, the timeout expires, or until this {@code PullConsumer}
+ * is shut down.
+ *
+ * @return the next message received from the attached queues, or null if the consumer is concurrently shut down or
+ * the timeout expires
+ * @throws OMSRuntimeException if the consumer fails to pull the next message due to some internal error.
+ */
+ Message receive();
+
+ /**
+ * Receives the next message from the bind queues of this consumer in pull model.
+ *
+ * This call blocks indefinitely until a message is arrives, the timeout expires, or until this {@code PullConsumer}
+ * is shut down.
+ *
+ * @param timeout receive message will blocked at most
+ * This call blocks indefinitely until a message is arrives, the timeout expires, or until this {@code PullConsumer}
+ * is shut down.
+ *
+ * @param queueName receive message from which queueName in Message Queue.
+ * @param queueMetaData receive message from which partition in Message Queue.
+ * @param messageReceipt receive message from which receipt position in Message Queue.
+ * @param timeout receive message will blocked at most
+ *
+ * @param timeout receive messages will blocked at most
+ *
+ * @param queueName receive message from which queueName in Message Queue.
+ * @param queueMetaData receive message from which partition in Message Queue.
+ * @param messageReceipt receive message from which receipt position in Message Queue.
+ * @param timeout receive messages will blocked at most
+ * Messages that have been received but not acknowledged may be redelivered.
+ *
+ * @param receiptHandle the receipt handle associated with the consumed message
+ * @throws OMSRuntimeException if the consumer fails to acknowledge the messages due to some internal error.
+ */
+ void ack(MessageReceipt receiptHandle);
+
+}
\ No newline at end of file
diff --git a/openmessaging-api/src/main/java/io/openmessaging/consumer/PushConsumer.java b/openmessaging-api/src/main/java/io/openmessaging/consumer/PushConsumer.java
new file mode 100644
index 00000000..3eaf5ae6
--- /dev/null
+++ b/openmessaging-api/src/main/java/io/openmessaging/consumer/PushConsumer.java
@@ -0,0 +1,120 @@
+/*
+ * 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 io.openmessaging.consumer;
+
+import io.openmessaging.KeyValue;
+import io.openmessaging.MessagingAccessPoint;
+import io.openmessaging.exception.OMSDestinationException;
+import io.openmessaging.exception.OMSRuntimeException;
+import io.openmessaging.exception.OMSSecurityException;
+import io.openmessaging.message.Message;
+import java.util.Collection;
+import java.util.List;
+
+/**
+ * A {@code PushConsumer} receives messages from multiple queues, these messages are pushed from
+ * MOM server to {@code PushConsumer} client.
+ *
+ * @version OMS 1.0.0
+ * @see MessagingAccessPoint#createPushConsumer()
+ * @since OMS 1.0.0
+ */
+public interface PushConsumer extends Consumer{
+ /**
+ * Resumes the {@code Consumer} in push model after a suspend.
+ *
+ * This method resumes the {@code Consumer} instance after it was suspended. The instance will not receive new
+ * messages between the suspend and resume calls.
+ *
+ * @throws OMSRuntimeException if the instance has not been suspended.
+ * @see PushConsumer#suspend()
+ */
+ void resume();
+
+ /**
+ * Suspends the {@code Consumer} in push model for later resumption.
+ *
+ * This method suspends the consumer until it is resumed. The consumer will not receive new messages between the
+ * suspend and resume calls.
+ *
+ * This method behaves exactly as if it simply performs the call {@code suspend(0)}.
+ *
+ * @throws OMSRuntimeException if the instance is not currently running.
+ * @see PushConsumer#resume()
+ */
+ void suspend();
+
+ /**
+ * Suspends the {@code Consumer} in push model for later resumption.
+ *
+ * This method suspends the consumer until it is resumed or a specified amount of time has elapsed. The consumer
+ * will not receive new messages during the suspended state.
+ *
+ * This method is similar to the {@link #suspend()} method, but it allows finer control over the amount of time to
+ * suspend, and the consumer will be suspended until it is resumed if the timeout is zero.
+ *
+ * @param timeout the maximum time to suspend in milliseconds.
+ * @throws OMSRuntimeException if the instance is not currently running.
+ */
+ void suspend(long timeout);
+
+ /**
+ * This method is used to find out whether the {@code Consumer} in push model is suspended.
+ *
+ * @return true if this {@code Consumer} is suspended, false otherwise.
+ */
+ boolean isSuspended();
+
+ /**
+ * Bind the {@code Consumer} to a collection of queue, with a {@code MessageListener}.
+ *
+ * {@link MessageListener#onReceived(Message, MessageListener.Context)} will be called when new delivered message is
+ * coming.
+ *
+ * @param queueNames a collection of queues.
+ * @param listener a specified listener to receive new message.
+ * @throws OMSSecurityException when have no authority to bind to this queue.
+ * @throws OMSDestinationException when have no given destination in the server.
+ * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ */
+ void bindQueue(Collection
+ * {@link BatchMessageListener#onReceived(List, BatchMessageListener.Context)} will be called when new delivered
+ * messages is coming.
+ *
+ * @param queueNames a collection of queues.
+ * @param listener a specified listener to receive new messages.
+ * @throws OMSSecurityException when have no authority to bind to this queue.
+ * @throws OMSDestinationException when have no given destination in the server.
+ * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ */
+ void bindQueue(Collection
+ * After the success call, this consumer won't receive new message from the specified queue any more.
+ *
+ * @param queueNames a collection of queues.
+ */
+ void unbindQueue(Collection
+ * This interface contains some methods are used for getting configurations related implementation. but this interface
+ * are not mandatory.
+ *
+ *
+ * @param queueName Queue name, message destination.
+ * @return {@link QueueMetaData} Queue config in the server
+ * @throws OMSSecurityException when have no authority to send messages to a given destination.
+ * @throws OMSTimeOutException when the given timeout elapses before the send operation completes.
+ * @throws OMSDestinationException when have no given destination in the server.
+ * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ */
+ Set
+ * The {@code ExtensionHeader} interface contains extended properties for common implementations in current messaging
+ * and streaming field, such as the queue-based partitioning implementation, but the related properties in this
+ * interface are not mandatory.
+ *
+ *
+ * When a {@link Message} is set with this value, this message will be delivered to specified partition, but the
+ * premise is that the implementation of the server side is dependent on the partition or a queue-like storage
+ * mechanism.
+ *
+ *
+ * @param partition The specified partition will be sent to.
+ */
+ ExtensionHeader setPartition(int partition);
+
+ /**
+ * This method is only called by the server. and {@code OFFSET} represents this message offset in partition.
+ *
+ *
+ * @param offset The offset in the current partition, used to quickly get this message in the queue
+ */
+ ExtensionHeader setOffset(long offset);
+
+ /**
+ * A client can use the {@code CORRELATION_ID} field to link one message with another. A typical use is to link a
+ * response message with its request message.
+ */
+ ExtensionHeader setCorrelationId(String correlationId);
+
+ /**
+ * This field {@code TRANSACTION_ID} is used in transactional message, and it can be used to trace a transaction.
+ *
+ * When a message is sent, STORE_TIMESTAMP is ignored.
+ *
+ * When the send method returns it contains a server-assigned value.
+ *
+ * This filed is a {@code long} value, measured in milliseconds.
+ */
+ ExtensionHeader setStoreTimestamp(long storeTimestamp);
+
+ /**
+ * The {@code STORE_HOST} header field contains the store host info of a message in server side.
+ *
+ * When a message is sent, STORE_HOST is ignored.
+ *
+ * When the send method returns it contains a server-assigned value.
+ */
+ ExtensionHeader setStoreHost(String storeHost);
+
+ /**
+ * The {@code messagekey} header field contains the custom key of a message.
+ *
+ * This key is a customer identifier for a class of messages, and this key may be used for server to hash or
+ * dispatch messages, or even can use this key to implement order message.
+ *
+ */
+ ExtensionHeader setMessageKey(String messageKey);
+
+ /**
+ * The {@code TRACE_ID} header field contains the trace ID of a message, which represents a global and unique
+ * identification, to associate key events in the whole lifecycle of a message, like sent by who, stored at where,
+ * and received by who.
+ *
+ * The {@code EXPIRE_TIME} represents a relative valid interval that a message can be delivered in it. If the
+ * EXPIRE_TIME field is specified as zero, that indicates the message does not expire.
+ *
+ * When an undelivered message's expiration time is reached, the message should be destroyed. OMS does not define a
+ * notification of message expiration.
+ *
+ *
+ * @return The {@code PARTITION} to which the message belongs
+ */
+ int getPartiton();
+
+ /**
+ * This method will return the {@code OFFSET} in the partition to which the message belongs to, but the premise is
+ * that the implementation of the server side is dependent on the partition or a queue-like storage mechanism.
+ *
+ * @return The offset of the partition to which the message belongs.
+ */
+ long getOffset();
+
+ /**
+ * See {@link ExtensionHeader#setCorrelationId(String)}
+ *
+ * @return correlationId
+ */
+ String getCorrelationId();
+
+ /**
+ * See {@link ExtensionHeader#setTransactionId(String)}
+ *
+ * @return transactionId
+ */
+ String getTransactionId();
+
+ /**
+ * See {@link ExtensionHeader#setStoreTimestamp(long)}
+ *
+ * @return storeTimestamp
+ */
+ long getStoreTimestamp();
+
+ /**
+ * See {@link ExtensionHeader#setStoreHost(String)}
+ *
+ * @return storeHost
+ */
+ String getStoreHost();
+
+ /**
+ * See {@link ExtensionHeader#setDelayTime(long)}
+ *
+ * @return delayTime
+ */
+ long getDelayTime();
+
+ /**
+ * See {@link ExtensionHeader#setExpireTime(long)}
+ *
+ * @return expireTime
+ */
+ long getExpireTime();
+
+ /**
+ * See {@link ExtensionHeader#setMessageKey(String)}
+ *
+ * @return messageKey
+ */
+ String getMessageKey();
+
+ /**
+ * See {@link ExtensionHeader#setTraceId(String)}
+ *
+ * @return traceId
+ */
+ String getTraceId();
+}
diff --git a/openmessaging-api/src/main/java/io/openmessaging/extension/QueueMetaData.java b/openmessaging-api/src/main/java/io/openmessaging/extension/QueueMetaData.java
new file mode 100644
index 00000000..0a0409b7
--- /dev/null
+++ b/openmessaging-api/src/main/java/io/openmessaging/extension/QueueMetaData.java
@@ -0,0 +1,66 @@
+/*
+ * 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 io.openmessaging.extension;
+
+import io.openmessaging.annotation.Optional;
+import java.util.List;
+import java.util.Set;
+
+/**
+ * This interface {@code QueueMetaData} contains methods are used for getting configurations related some certain
+ * implementation. but this interface are not mandatory.
+ *
+ *
+ * In order to improve performance, in some scenarios where message persistence is required, some message middleware
+ * will store messages on multiple partitions in multi servers.
+ *
+ *
+ * In some scenarios, it is very useful to get the relevant partitions meta data for a queue.
+ *
+ * @version OMS 1.0.0
+ * @since OMS 1.0.0
+ */
+@Optional
+public interface QueueMetaData {
+
+ /**
+ * Set queueName to this Message Queue.
+ * @param queueName
+ */
+ void setQueueName(String queueName);
+
+ /**
+ * Set the specified partition.
+ * @param partitionId
+ */
+ void setPartitionId(int partitionId);
+
+ /**
+ * Get partition identifier of target queue.
+ *
+ * @return Partition identifier
+ */
+ int partitionId();
+
+ /**
+ * Queue name
+ *
+ *
+ * @return Queue name.
+ */
+ String queueName();
+}
diff --git a/openmessaging-api/src/main/java/io/openmessaging/interceptor/ConsumerInterceptor.java b/openmessaging-api/src/main/java/io/openmessaging/interceptor/ConsumerInterceptor.java
index 440974d2..98aba203 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/interceptor/ConsumerInterceptor.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/interceptor/ConsumerInterceptor.java
@@ -17,7 +17,7 @@
package io.openmessaging.interceptor;
-import io.openmessaging.Message;
+import io.openmessaging.message.Message;
import io.openmessaging.consumer.MessageListener;
/**
diff --git a/openmessaging-api/src/main/java/io/openmessaging/interceptor/ProducerInterceptor.java b/openmessaging-api/src/main/java/io/openmessaging/interceptor/ProducerInterceptor.java
index eaefde35..e41f60a0 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/interceptor/ProducerInterceptor.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/interceptor/ProducerInterceptor.java
@@ -17,13 +17,12 @@
package io.openmessaging.interceptor;
-import io.openmessaging.Message;
+import io.openmessaging.message.Message;
/**
* A {@code ProducerInterceptor} is used to intercept send operations of producer.
*
- * The interceptor is able to view or modify the message being transmitted and collect
- * the send record.
+ * The interceptor is able to view or modify the message being transmitted and collect the send record.
*
* @version OMS 1.0.0
* @since OMS 1.0.0
@@ -46,5 +45,5 @@ public interface ProducerInterceptor {
* @param attributes the extensible attributes delivered to the intercept thread.
*/
void postSend(Message message, Context attributes);
-
+
}
diff --git a/openmessaging-api/src/main/java/io/openmessaging/internal/DefaultKeyValue.java b/openmessaging-api/src/main/java/io/openmessaging/internal/DefaultKeyValue.java
index e2f6477b..952e17c6 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/internal/DefaultKeyValue.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/internal/DefaultKeyValue.java
@@ -32,13 +32,35 @@ public class DefaultKeyValue implements KeyValue {
private Map
+ * The header contains fields used by the messaging system that describes the message's meta information, while the body
+ * contains the application data being transmitted.
+ *
+ * As for the message header, OMS defines three kinds types: headers {@link Header} {@link ExtensionHeader} and
+ * properties {@link KeyValue}, with respect to flexibility in vendor implementation and user usage.
+ *
+ *
+ * The header part is placed in the implementation classes of {@code Message}.
+ *
+ * @version OMS 1.0.0
+ * @since OMS 1.0.0
+ */
+public interface Header {
+ /**
+ * The {@code DESTINATION} header field contains the destination to which the message is being sent.
+ *
+ * When a message is set to the {@code Queue}, then the message will be sent to the specified destination.
+ *
+ * When a message is received, its destination is equivalent to the {@code Queue} where the message resides in.
+ */
+ Header setDestination(String destination);
+
+ /**
+ * The {@code MESSAGE_ID} header field contains a value that uniquely identify each message sent by a {@code
+ * Producer}. this identifier is generated by producer.
+ */
+ Header setMessageId(String messageId);
+
+ /**
+ * The {@code BORN_TIMESTAMP} header field contains the time a message was handed off to a {@code Producer} to be
+ * sent.
+ *
+ * When a message is sent, BORN_TIMESTAMP will be set with current timestamp as the born timestamp of a message in
+ * client side, on return from the send method, the message's BORN_TIMESTAMP header field contains this value.
+ *
+ * When a message is received its, BORN_TIMESTAMP header field contains this same value.
+ *
+ * This filed is a {@code long} value, measured in milliseconds.
+ */
+ Header setBornTimestamp(long bornTimestamp);
+
+ /**
+ * The {@code BORN_HOST} header field contains the born host info of a message in client side.
+ *
+ * When a message is sent, BORN_HOST will be set with the local host info, on return from the send method, the
+ * message's BORN_HOST header field contains this value.
+ *
+ * When a message is received, its BORN_HOST header field contains this same value.
+ */
+ Header setBornHost(String bornHost);
+
+ /**
+ * The {@code PRIORITY} header field contains the priority level of a message, a message with a higher priority
+ * value should be delivered preferentially.
+ *
+ * OMS defines a ten level priority value with 1 as the lowest priority and 10 as the highest, and the default
+ * priority is 5. The priority beyond this region will be ignored.
+ *
+ * OMS does not require or provide any guarantee that the message should be delivered in priority order strictly,
+ * but the vendor should provide a best effort to deliver expedited messages ahead of normal messages.
+ *
+ * If PRIORITY field isn't set explicitly, use {@code 5} as the default priority.
+ */
+ Header setPriority(short priority);
+
+ /**
+ * The {@code DURABILITY} header field contains the persistent level of a message, the vendor should guarantee the
+ * reliability level for a message.
+ *
+ * OMS defines two modes of message delivery:
+ *
+ * Most message-oriented middleware (MOM) products treat messages as lightweight entities that consist of header and
+ * body and is used by separate applications to exchange a piece of information, like Apache RocketMQ.
+ *
+ * The header contains fields used by the messaging system that describes the message's meta information, while the body
+ * contains the application data being transmitted.
+ *
+ * As for the message header, OMS defines three kinds types: headers {@link Header} {@link ExtensionHeader} and
+ * properties {@link KeyValue}, with respect to flexibility in vendor implementation and user usage.
+ *
+ * In BytesMessage, the body is just a byte array, may be compressed and uncompressed in the transmitting process by the
+ * messaging system. The application is responsible for explaining the concrete content and format of the message body,
+ * OMS is never aware of that.
+ *
+ * The body part is placed in the implementation classes of {@code Message}.
+ *
+ * @version OMS 1.0.0
+ * @since OMS 1.0.0
+ */
+public interface Message {
+ /**
+ * Returns all the system header fields of the {@code Message} object as a {@code KeyValue}.
+ *
+ * @return the system headers of a {@code Message}
+ */
+ Header header();
+
+ /**
+ * This interface is optional, Therefore, users need to check whether the interface is implemented and the
+ * correctness of its implementation.
+ *
+ *
+ * @return The implementation of {@link ExtensionHeader}
+ */
+ @Optional
+ ExtensionHeader extensionHeader();
+
+ /**
+ * Returns all the customized user header fields of the {@code Message} object as a {@code KeyValue}.
+ *
+ * @return the user properties of a {@code Message}
+ */
+ KeyValue properties();
+
+ /**
+ * Get data from message body
+ *
+ * @return message body
+ * @throws OMSMessageFormatException if the message body cannot be assigned to the specified type
+ */
+ byte[] getData();
+
+ /**
+ * Set data to message body
+ *
+ * @param data set message body in binary stream
+ */
+ void setData(byte[] data);
+
+ /**
+ * Get the {@code MessageReceipt} of this Message, which will be used to acknowledge this message.
+ *
+ * @see Consumer#ack(io.openmessaging.consumer.MessageReceipt)
+ * @see MessageListener.Context#ack()
+ * @see BatchMessageListener.Context#success(io.openmessaging.consumer.MessageReceipt...)
+ */
+ MessageReceipt getMessageReceipt();
+
+}
\ No newline at end of file
diff --git a/openmessaging-api/src/main/java/io/openmessaging/MessageFactory.java b/openmessaging-api/src/main/java/io/openmessaging/message/MessageFactory.java
similarity index 95%
rename from openmessaging-api/src/main/java/io/openmessaging/MessageFactory.java
rename to openmessaging-api/src/main/java/io/openmessaging/message/MessageFactory.java
index 536f8ce1..4902d4b6 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/MessageFactory.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/message/MessageFactory.java
@@ -15,9 +15,10 @@
* limitations under the License.
*/
-package io.openmessaging;
+package io.openmessaging.message;
import io.openmessaging.exception.OMSMessageFormatException;
+import io.openmessaging.message.Message;
/**
* A factory interface for creating {@code Message} objects.
diff --git a/openmessaging-api/src/main/java/io/openmessaging/producer/Producer.java b/openmessaging-api/src/main/java/io/openmessaging/producer/Producer.java
index 8ebd35d5..3dbc10d4 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/producer/Producer.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/producer/Producer.java
@@ -17,10 +17,9 @@
package io.openmessaging.producer;
+import io.openmessaging.Client;
import io.openmessaging.Future;
import io.openmessaging.FutureListener;
-import io.openmessaging.Message;
-import io.openmessaging.MessageFactory;
import io.openmessaging.MessagingAccessPoint;
import io.openmessaging.ServiceLifecycle;
import io.openmessaging.exception.OMSDestinationException;
@@ -30,6 +29,8 @@
import io.openmessaging.exception.OMSTimeOutException;
import io.openmessaging.exception.OMSTransactionException;
import io.openmessaging.interceptor.ProducerInterceptor;
+import io.openmessaging.message.Message;
+import io.openmessaging.message.MessageFactory;
import java.util.List;
/**
@@ -49,11 +50,11 @@
* @version OMS 1.0.0
* @since OMS 1.0.0
*/
-public interface Producer extends MessageFactory, ServiceLifecycle {
+public interface Producer extends MessageFactory, ServiceLifecycle, Client {
/**
* Sends a message to the specified destination synchronously, the destination should be preset to {@link
- * Message#headers()}, other header fields as well.
+ * Message#header()}, other header fields as well.
*
* @param message a message will be sent.
* @return the successful {@code SendResult}.
@@ -67,7 +68,7 @@ public interface Producer extends MessageFactory, ServiceLifecycle {
/**
* Sends a message to the specified destination asynchronously, the destination should be preset to {@link
- * Message#headers()}, other header fields as well.
+ * Message#header()}, other header fields as well.
*
* The returned {@code Promise} will have the result once the operation completes, and the registered {@code
* FutureListener} will be notified, either because the operation was successful or because of an error.
@@ -96,6 +97,29 @@ public interface Producer extends MessageFactory, ServiceLifecycle {
*/
void send(List
+ * The returned {@code Promise} will have the result once the operation completes, and the registered {@code
+ * FutureListener} will be notified, either because the operation was successful or because of an error.
+ *
+ * @param messages a batch messages will be sent.
+ * @return the {@code Promise} of an asynchronous messages send operation.
+ * @see Future
+ * @see FutureListener
+ */
+ Future
+ * There is no {@code Promise} related or {@code RuntimeException} thrown. The calling thread doesn't care about the
+ * send result and also have no context to get the result.
+ *
+ * @param messages a batch message will be sent.
+ */
+ void sendOneway(List
* A transactional send result will be exposed to consumer if this prepare message send success, and then, you can
* execute your local transaction, when local transaction execute success, users can use {@link
@@ -127,7 +151,9 @@ public interface Producer extends MessageFactory, ServiceLifecycle {
* @throws OMSTimeOutException when the given timeout elapses before the send operation completes.
* @throws OMSDestinationException when have no given destination in the server.
* @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
- * @throws OMSTransactionException when used normal producer which haven't register {@link TransactionStateCheckListener}.
+ * @throws OMSTransactionException when used normal producer which haven't register {@link
+ * TransactionStateCheckListener}.
*/
TransactionalResult prepare(Message message);
+
}
\ No newline at end of file
diff --git a/openmessaging-api/src/main/java/io/openmessaging/producer/TransactionStateCheckListener.java b/openmessaging-api/src/main/java/io/openmessaging/producer/TransactionStateCheckListener.java
index 99665544..efa0d5f1 100644
--- a/openmessaging-api/src/main/java/io/openmessaging/producer/TransactionStateCheckListener.java
+++ b/openmessaging-api/src/main/java/io/openmessaging/producer/TransactionStateCheckListener.java
@@ -17,7 +17,7 @@
package io.openmessaging.producer;
-import io.openmessaging.Message;
+import io.openmessaging.message.Message;
/**
* Each executor will be associated with a transactional message, can be used to execute local transaction branch and
diff --git a/openmessaging-api/src/test/java/io/openmessaging/internal/MessagingAccessPointAdapterTest.java b/openmessaging-api/src/test/java/io/openmessaging/internal/MessagingAccessPointAdapterTest.java
index 14cb8d0b..b348996e 100644
--- a/openmessaging-api/src/test/java/io/openmessaging/internal/MessagingAccessPointAdapterTest.java
+++ b/openmessaging-api/src/test/java/io/openmessaging/internal/MessagingAccessPointAdapterTest.java
@@ -22,7 +22,10 @@
import io.openmessaging.OMS;
import io.openmessaging.OMSBuiltinKeys;
import io.openmessaging.consumer.Consumer;
+import io.openmessaging.consumer.PullConsumer;
+import io.openmessaging.consumer.PushConsumer;
import io.openmessaging.manager.ResourceManager;
+import io.openmessaging.message.MessageFactory;
import io.openmessaging.producer.Producer;
import io.openmessaging.producer.TransactionStateCheckListener;
import org.junit.Test;
@@ -46,27 +49,48 @@ class TestVendor implements MessagingAccessPoint {
public TestVendor(KeyValue keyValue) {
}
- @Override public Producer createProducer(TransactionStateCheckListener transactionStateCheckListener) {
+ @Override
+ public Producer createProducer(TransactionStateCheckListener transactionStateCheckListener) {
return null;
}
@Override
- public String version() {
- return OMS.specVersion;
+ public PushConsumer createPushConsumer() {
+ return null;
}
@Override
- public KeyValue attributes() {
+ public PullConsumer createPullConsumer() {
return null;
}
@Override
- public Producer createProducer() {
+ public PushConsumer createPushConsumer(KeyValue attributes) {
return null;
}
@Override
- public Consumer createConsumer() {
+ public PullConsumer createPullConsumer(KeyValue attributes) {
+ return null;
+ }
+
+ @Override
+ public MessageFactory messageFactory() {
+ return null;
+ }
+
+ @Override
+ public String version() {
+ return OMS.specVersion;
+ }
+
+ @Override
+ public KeyValue attributes() {
+ return null;
+ }
+
+ @Override
+ public Producer createProducer() {
return null;
}
diff --git a/pom.xml b/pom.xml
index 895934f5..0ef1d346 100644
--- a/pom.xml
+++ b/pom.xml
@@ -10,7 +10,7 @@
- *
- * The body contains the application data being transmitted, which is generally ignored by the messaging system and
- * simply transmitted to its destination.
- *
- *
- */
- Headers setDurability(short durability);
-
- /**
- * The {@code messagekey} header field contains the custom key of a message.
- * timeout milliseconds.
- * @return the next message received from the bind queues, or null if the consumer is concurrently shut down.
- * @throws OMSSecurityException when have no authority to receive messages from this queue.
- * @throws OMSTimeOutException when the given timeout elapses before the send operation completes.
- * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
- */
- Message receive(long timeout);
-
/**
* Acknowledges the specified and consumed message with the unique message receipt handle, in the scenario of using
* manual commit.
* timeout milliseconds.
+ * @return the next message received from the bind queues, or null if the consumer is concurrently shut down.
+ * @throws OMSSecurityException when have no authority to receive messages from this queue.
+ * @throws OMSTimeOutException when the given timeout elapses before the send operation completes.
+ * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ */
+ Message receive(long timeout);
+
+ /**
+ * Receives the next message from the which bind queue,partition and receiptId of this consumer in pull model.
+ * timeout milliseconds.
+ * @return the next message received from the bind queues, or null if the consumer is concurrently shut down.
+ * @throws OMSSecurityException when have no authority to receive messages from this queue.
+ * @throws OMSTimeOutException when the given timeout elapses before the send operation completes.
+ * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ */
+ @Optional
+ Message receive(String queueName, QueueMetaData queueMetaData, MessageReceipt messageReceipt, long timeout);
+
+ /**
+ * Receive message in asynchronous way. This call doesn't block user's thread, and user's message resolve logic
+ * should implement in the {@link MessageListener}.
+ * timeout milliseconds.
+ * @return the next batch messages received from the bind queues, or null if the consumer is concurrently shut down.
+ * @throws OMSSecurityException when have no authority to receive messages from this queue.
+ * @throws OMSTimeOutException when the given timeout elapses before the send operation completes.
+ * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ */
+ Listtimeout milliseconds.
+ * @return the next batch messages received from the bind queues, or null if the consumer is concurrently shut down.
+ * @throws OMSSecurityException when have no authority to receive messages from this queue.
+ * @throws OMSTimeOutException when the given timeout elapses before the send operation completes.
+ * @throws OMSRuntimeException when the {@code Producer} fails to send the message due to some internal error.
+ */
+ @Optional
+ List
+ *
+ * The body contains the application data being transmitted, which is generally ignored by the messaging system and
+ * simply transmitted to its destination.
+ *
+ *
+ */
+ Header setDurability(short durability);
+
+ /**
+ * The {@code DELIVERY_COUNT} header field contains a number, which represents the count of the message delivery.
+ */
+ Header setDeliveryCount(int deliveryCount);
+
+ /**
+ * The field {@code COMPRESSION} in headers represents the message body compress algorithm. vendors are free to
+ * choose the compression algorithm, but must ensure that the decompressed message is delivered to the user.
+ */
+ Header setCompression(short compression);
+
+ /**
+ * See {@link Header#setDestination(String)}
+ *
+ * @return destination
+ */
+ String getDestination();
+
+ /**
+ * See {@link Header#setMessageId(String)}
+ *
+ * @return messageId
+ */
+ String getMessageId();
+
+ /**
+ * See {@link Header#setBornTimestamp(long)}
+ *
+ * @return bornTimestamp
+ */
+ long getBornTimestamp();
+
+ /**
+ * See {@link Header#setBornHost(String)}
+ *
+ * @return bornHost
+ */
+ String getBornHost();
+
+ /**
+ * See {@link Header#setPriority(short)}
+ *
+ * @return priority
+ */
+ short getPriority();
+
+ /**
+ * See {@link Header#setDurability(short)}
+ *
+ * @return durability
+ */
+ short getDurability();
+
+ /**
+ * See {@link Header#setDeliveryCount(int)}
+ *
+ * @return deliveryCount
+ */
+ int getDeliveryCount();
+
+ /**
+ * See {@link Header#setCompression(short)}
+ *
+ * @return compression
+ */
+ short getCompression();
+}
diff --git a/openmessaging-api/src/main/java/io/openmessaging/message/Message.java b/openmessaging-api/src/main/java/io/openmessaging/message/Message.java
new file mode 100644
index 00000000..a43d23ed
--- /dev/null
+++ b/openmessaging-api/src/main/java/io/openmessaging/message/Message.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 io.openmessaging.message;
+
+import io.openmessaging.KeyValue;
+import io.openmessaging.annotation.Optional;
+import io.openmessaging.consumer.BatchMessageListener;
+import io.openmessaging.consumer.Consumer;
+import io.openmessaging.consumer.MessageListener;
+import io.openmessaging.consumer.MessageReceipt;
+import io.openmessaging.exception.OMSMessageFormatException;
+import io.openmessaging.extension.ExtensionHeader;
+
+/**
+ * The {@code Message} interface is the root interface of all OMS messages, and the most commonly used OMS message is
+ * {@link Message}.
+ *
+ *
+ * The body contains the application data being transmitted, which is generally ignored by the messaging system and
+ * simply transmitted to its destination.
+ *