diff --git a/.travis.yml b/.travis.yml index 7047779b..633b3b73 100644 --- a/.travis.yml +++ b/.travis.yml @@ -1,24 +1,18 @@ -language: java +dist: trusty -addons: - apt: - packages: - - openjdk-6-jdk +notifications: + email: + recipients: + - openmessaging-members@googlegroups.com + on_success: change + on_failure: always + + +language: java -install: - - echo "Downloading Maven 3.0"; - - wget https://archive.apache.org/dist/maven/binaries/apache-maven-3.0-bin.zip || travis_terminate 1 - - unzip -qq apache-maven-3.0-bin.zip || travis_terminate 1 - - export M2_HOME=$PWD/apache-maven-3.0 - - export PATH=$M2_HOME/bin:$PATH - - mvn -version - - mvn clean package install -DskipTests -Dgpg.skip - jdk: -- oraclejdk8 -- oraclejdk9 -- openjdk7 -- openjdk6 + - oraclejdk8 + - oraclejdk9 script: mvn install env: @@ -26,4 +20,4 @@ env: secure: hmPdcALAi6qE3TqJDRqdVCqZftd/i2hWLCyZbIAcRzu38nO94JZYKSZjfif1FvXTJYotFW25JClXNyvOMwMjjK3OPQINfYFZIp6LLeOmXGbUcktwQ8TIoKZ7IOvvWiZK054H7zNKapz+ke3OPN/5WTmMBezV0Ct4+bSf9udKVnQSMG2sJ8YJ/SeZkh7RTlqO+zTkh+yq8Hk0BdaWEOK8RtEoWgcUFGVfkycvjgvna+TbDp3K7vjmhYBBqACsNKxXPgIumStbCGW4vwjoVkCOGIJKWnuQEVHxiqBUH3pp81bxnt+RIcMuZMR2HnDSpHyAIulTJNHVo3VFAAiy9HMdP8Wfy/OVdjBSZ8xIOoQvFijo+yGNNn8v4hILcX4IpumQeyjpG134BOWVbMLhKH7qWR3Z8TGgijSd4lYYjabCJ564E93KvqK1u2CuS9u89N8J7AKFYMbknH1DP8E5tCD+VI3Gwut9YNofywj3Jln8uCOP4I//8p61j9A9QF7ORpY59Ru4RNzxYrFn2QSTltMfaBfVZchh5AqURUamcJd+1orZfz/v+6yH9FOW+MAG8EJdzHDsqzP1NXrt+4VtF6yqOnhBxnKVNEwFwjsinW9PFi9dXyzdEd33jKGL7UO8Old5XlBoA7idWIDH4GKKSlBRZhEKWMe4ZfxpQVg3VPz2Qqo= after_success: -- bash .utility/push-javadoc-to-gh-pages.sh + - bash .utility/push-javadoc-to-gh-pages.sh diff --git a/README.md b/README.md index 76594031..d7a9539f 100644 --- a/README.md +++ b/README.md @@ -1,15 +1,13 @@ ## ![logo](assets/images/logo-color.png) -[![Build Status](https://travis-ci.org/openmessaging/openmessaging-java.svg?branch=master)](https://travis-ci.org/openmessaging/openmessaging-java) [![Maven Central](https://maven-badges.herokuapp.com/maven-central/io.openmessaging/openmessaging-api/badge.svg)](http://search.maven.org/#search%7Cga%7C1%7Copenmessaging) [![Gitter chat](https://badges.gitter.im/gitterHQ/gitter.png)](https://gitter.im/openmessaging/public) [![License](https://img.shields.io/badge/license-Apache%202-4EB1BA.svg)](https://www.apache.org/licenses/LICENSE-2.0.html) +[![Build Status](https://travis-ci.org/openmessaging/openmessaging-java.svg?branch=master)](https://travis-ci.org/openmessaging/openmessaging-java) [![Maven Central](https://maven-badges.herokuapp.com/maven-central/io.openmessaging/openmessaging-api/badge.svg)](http://search.maven.org/#search%7Cga%7C1%7Copenmessaging) [![Slack chat](https://img.shields.io/badge/chat-on%20slack-green.svg)](https://openmessaging.herokuapp.com/) [![License](https://img.shields.io/badge/license-Apache%202-4EB1BA.svg)](https://www.apache.org/licenses/LICENSE-2.0.html) ### A vendor-neutral open standard for distributed messaging and streaming OpenMessaging, which includes the establishment of industry guidelines and messaging, streaming specifications to provide a common framework for finance, e-commerce, IoT and big-data area. The design principles are the cloud-oriented, simplicity, flexibility, and language independent in distributed heterogeneous environments. Conformance to these specifications will make it possible to develop a heterogeneous messaging applications across all major platforms and operating systems. -## The domain architecture -![domain-design](assets/images/OpenMessaging-V0.3.0-alpha.png) ## Doc [API Doc](https://openmessaging.github.io/openmessaging-java/). -## ![Powered by Linux Foundation](http://openmessaging.cloud/images/linux-foundation-logo.png) \ No newline at end of file +## ![Powered by Linux Foundation](http://openmessaging.cloud/images/linux-foundation-logo.png) diff --git a/openmessaging-admin/pom.xml b/openmessaging-admin/pom.xml index b0d46f67..debe8f13 100644 --- a/openmessaging-admin/pom.xml +++ b/openmessaging-admin/pom.xml @@ -2,7 +2,7 @@ io.openmessaging parent - 1.0.0-Preview + 1.0.0-beta-SNAPSHOT 4.0.0 diff --git a/openmessaging-api-samples/pom.xml b/openmessaging-api-samples/pom.xml index 2672ae1f..d6bc37dc 100644 --- a/openmessaging-api-samples/pom.xml +++ b/openmessaging-api-samples/pom.xml @@ -2,25 +2,26 @@ io.openmessaging parent - 1.0.0-Preview + 1.0.0-beta-SNAPSHOT 4.0.0 jar openmessaging-api-samples + 1.0.0-beta-SNAPSHOT openmessaging-api-samples ${project.version} junit junit - 4.11 + 4.13.1 test ${project.groupId} openmessaging-api - ${project.version} + 1.0.0-beta-SNAPSHOT org.slf4j diff --git a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/consumer/PullConsumerApp.java b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/consumer/PullConsumerApp.java index b60f89e0..7f9a2e85 100644 --- a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/consumer/PullConsumerApp.java +++ b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/consumer/PullConsumerApp.java @@ -17,11 +17,11 @@ package io.openmessaging.samples.consumer; -import io.openmessaging.Message; import io.openmessaging.MessagingAccessPoint; import io.openmessaging.OMS; -import io.openmessaging.consumer.Consumer; -import io.openmessaging.manager.ResourceManager; +import io.openmessaging.consumer.PullConsumer; +import io.openmessaging.message.Message; +import java.util.Arrays; public class PullConsumerApp { public static void main(String[] args) { @@ -29,13 +29,8 @@ public static void main(String[] args) { final MessagingAccessPoint messagingAccessPoint = OMS.getMessagingAccessPoint("oms:rocketmq://alice@rocketmq.apache.org/us-east"); - //Fetch a ResourceManager to create Queue resource. - ResourceManager resourceManager = messagingAccessPoint.resourceManager(); - resourceManager.createQueue("NS://HELLO_QUEUE"); - //Start a PullConsumer to receive messages from the specific queue. - final Consumer consumer = messagingAccessPoint.createConsumer(); - consumer.start(); + final PullConsumer consumer = messagingAccessPoint.createPullConsumer(); //Register a shutdown hook to close the opened endpoints. Runtime.getRuntime().addShutdownHook(new Thread(new Runnable() { @@ -44,10 +39,15 @@ public void run() { consumer.stop(); } })); - consumer.bindQueue("NS://HELLO_QUEUE"); + + consumer.bindQueue(Arrays.asList("NS://HELLO_QUEUE")); + consumer.start(); + Message message = consumer.receive(1000); System.out.println("Received message: " + message); //Acknowledge the consumed message - consumer.ack(message.headers().getMessageId()); + consumer.ack(message.getMessageReceipt()); + consumer.stop(); + } } diff --git a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/consumer/PushConsumerApp.java b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/consumer/PushConsumerApp.java index 0fed2032..47ea14eb 100644 --- a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/consumer/PushConsumerApp.java +++ b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/consumer/PushConsumerApp.java @@ -17,12 +17,14 @@ package io.openmessaging.samples.consumer; -import io.openmessaging.Message; import io.openmessaging.MessagingAccessPoint; import io.openmessaging.OMS; import io.openmessaging.consumer.Consumer; import io.openmessaging.consumer.MessageListener; +import io.openmessaging.consumer.PushConsumer; import io.openmessaging.manager.ResourceManager; +import io.openmessaging.message.Message; +import java.util.Arrays; public class PushConsumerApp { public static void main(String[] args) { @@ -33,7 +35,7 @@ public static void main(String[] args) { //Fetch a ResourceManager to create Queue resource. ResourceManager resourceManager = messagingAccessPoint.resourceManager(); resourceManager.createNamespace("NS://XXXX"); - final Consumer consumer = messagingAccessPoint.createConsumer(); + final PushConsumer consumer = messagingAccessPoint.createPushConsumer(); consumer.start(); //Register a shutdown hook to close the opened endpoints. @@ -49,7 +51,7 @@ public void run() { resourceManager.createQueue(simpleQueue); //This queue doesn't has a source queue, so only the message delivered to the queue directly can //be consumed by this consumer. - consumer.bindQueue(simpleQueue, new MessageListener() { + consumer.bindQueue(Arrays.asList(simpleQueue), new MessageListener() { @Override public void onReceived(Message message, Context context) { System.out.println("Received one message: " + message); @@ -58,7 +60,7 @@ public void onReceived(Message message, Context context) { }); - consumer.unbindQueue(simpleQueue); + consumer.unbindQueue(Arrays.asList(simpleQueue)); consumer.stop(); } diff --git a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/producer/ProducerApp.java b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/producer/ProducerApp.java index 2da775fc..9084b8a8 100644 --- a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/producer/ProducerApp.java +++ b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/producer/ProducerApp.java @@ -18,11 +18,11 @@ package io.openmessaging.samples.producer; import io.openmessaging.Future; -import io.openmessaging.Message; import io.openmessaging.MessagingAccessPoint; import io.openmessaging.OMS; import io.openmessaging.interceptor.Context; import io.openmessaging.interceptor.ProducerInterceptor; +import io.openmessaging.message.Message; import io.openmessaging.producer.Producer; import io.openmessaging.producer.SendResult; import java.nio.charset.Charset; @@ -35,17 +35,19 @@ public static void main(String[] args) { OMS.getMessagingAccessPoint("oms:rocketmq://alice@rocketmq.apache.org/us-east"); final Producer producer = messagingAccessPoint.createProducer(); - producer.start(); ProducerInterceptor interceptor = new ProducerInterceptor() { @Override public void preSend(Message message, Context attributes) { + System.out.println("PreSend message: " + message); } @Override public void postSend(Message message, Context attributes) { + System.out.println("PostSend message: " + message); } }; producer.addInterceptor(interceptor); + producer.start(); //Register a shutdown hook to close the opened endpoints. Runtime.getRuntime().addShutdownHook(new Thread(new Runnable() { @@ -55,9 +57,11 @@ public void run() { } })); - //Sends a message to the specified destination synchronously. + //Send a message to the specified destination synchronously. Message message = producer.createMessage( - "NS://HELLO_QUEUE", "HELLO_BODY".getBytes(Charset.forName("UTF-8"))); + "NS://HELLO_QUEUE1", "HELLO_BODY".getBytes(Charset.forName("UTF-8"))); + message.header().setBornHost("127.0.0.1").setDurability((short) 0); + message.extensionHeader().setPartition(1); SendResult sendResult = producer.send(message); System.out.println("SendResult: " + sendResult); @@ -75,6 +79,7 @@ public void run() { Message msg = producer.createMessage("NS://HELLO_QUEUE", ("Hello" + i).getBytes()); messages.add(msg); } + producer.send(messages); producer.removeInterceptor(interceptor); producer.stop(); diff --git a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/producer/TransactionProducerApp.java b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/producer/TransactionProducerApp.java index f6db2730..bed57642 100644 --- a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/producer/TransactionProducerApp.java +++ b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/producer/TransactionProducerApp.java @@ -17,7 +17,7 @@ package io.openmessaging.samples.producer; -import io.openmessaging.Message; +import io.openmessaging.message.Message; import io.openmessaging.MessagingAccessPoint; import io.openmessaging.OMS; import io.openmessaging.producer.Producer; diff --git a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/routing/RoutingApp.java b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/routing/RoutingApp.java index 3f3a9a66..298655cc 100644 --- a/openmessaging-api-samples/src/main/java/io/openmessaging/samples/routing/RoutingApp.java +++ b/openmessaging-api-samples/src/main/java/io/openmessaging/samples/routing/RoutingApp.java @@ -17,13 +17,15 @@ package io.openmessaging.samples.routing; -import io.openmessaging.Message; +import io.openmessaging.consumer.PushConsumer; +import io.openmessaging.message.Message; import io.openmessaging.MessagingAccessPoint; import io.openmessaging.OMS; import io.openmessaging.consumer.Consumer; import io.openmessaging.consumer.MessageListener; import io.openmessaging.manager.ResourceManager; import io.openmessaging.producer.Producer; +import java.util.Arrays; public class RoutingApp { public static void main(String[] args) { @@ -54,10 +56,10 @@ public static void main(String[] args) { producer.send(message); //Consume messages from the queue behind the routing. - final Consumer consumer = messagingAccessPoint.createConsumer(); + final PushConsumer consumer = messagingAccessPoint.createPushConsumer(); consumer.start(); - consumer.bindQueue(destinationQueue, new MessageListener() { + consumer.bindQueue(Arrays.asList(destinationQueue), new MessageListener() { @Override public void onReceived(Message message, Context context) { //The message sent to the sourceQueue will be delivered to anotherConsumer by the routing rule diff --git a/openmessaging-api/pom.xml b/openmessaging-api/pom.xml index 196eba5a..6f969895 100644 --- a/openmessaging-api/pom.xml +++ b/openmessaging-api/pom.xml @@ -2,7 +2,7 @@ io.openmessaging parent - 1.0.0-Preview + 1.0.0-beta-SNAPSHOT 4.0.0 @@ -21,7 +21,7 @@ junit junit - 4.11 + 4.13.1 test diff --git a/openmessaging-api/src/main/java/io/openmessaging/Client.java b/openmessaging-api/src/main/java/io/openmessaging/Client.java new file mode 100644 index 00000000..ed20ebad --- /dev/null +++ b/openmessaging-api/src/main/java/io/openmessaging/Client.java @@ -0,0 +1,41 @@ +/* + * 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.extension.Extension; +import java.util.Optional; + +/** + *

+ * 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 getExtension(); +} diff --git a/openmessaging-api/src/main/java/io/openmessaging/Future.java b/openmessaging-api/src/main/java/io/openmessaging/Future.java index b22ad541..e0973549 100644 --- a/openmessaging-api/src/main/java/io/openmessaging/Future.java +++ b/openmessaging-api/src/main/java/io/openmessaging/Future.java @@ -32,6 +32,29 @@ * @since OMS 1.0.0 */ public interface Future { + /** + * Attempts to cancel execution of this task. This attempt will + * fail if the task has already completed, has already been cancelled, + * or could not be cancelled for some other reason. If successful, + * and this task has not started when {@code cancel} is called, + * this task should never run. If the task has already started, + * then the {@code mayInterruptIfRunning} parameter determines + * whether the thread executing this task should be interrupted in + * an attempt to stop the task. + * + *

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. - *

    - *
  • - * System Headers, OMS defines some standard attributes that represent the characteristics of the message. - *
  • - *
  • - * User properties, some OMS vendors may require enhanced extra attributes of the message or some users may want to - * clarify some customized attributes to draw the body. OMS provides the improved scalability for these scenarios. - *
  • - *
- * The body contains the application data being transmitted, which is generally ignored by the messaging system and - * simply transmitted to its destination. - *

- * 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 message will be delivered after delayTime milliseconds starting from {@CODE BORN_TIMESTAMP} . When this - * filed isn't set explicitly, this means this message should be delivered immediately. - */ - Headers setDelayTime(long delayTime); - - /** - * The {@code EXPIRE_TIME} header field contains the expiration time, it represents a time-to-live value. - *

- * 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. - *

- */ - Headers setExpireTime(long expireTime); - - /** - * 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. - */ - 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: - *

    - *
  • - * PERSISTENT, the persistent mode instructs the vendor should provide stable storage to ensure the message - * won't be lost. - *
  • - *
  • - * NON_PERSISTENT, this mode does not require the message be logged to stable storage, in most cases, the memory - * storage is enough for better performance and lower cost. - *
  • - *
- */ - Headers setDurability(short durability); - - /** - * 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. - *

- */ - 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. - *

- * And, the messaging system only plays exchange role in a distributed system in most cases, so the TraceID can - * be used to trace the whole call link with other parts in the whole system. - */ - Headers setTraceId(String traceId); - - /** - * The {@code DELIVERY_COUNT} header field contains a number, which represents the count of the message - * delivery. - */ - Headers setDeliveryCount(int deliveryCount); - - /** - * This field {@code TRANSACTION_ID} is used in transactional message, and it can be used to trace a - * transaction. - *

- * So the same {@code TRANSACTION_ID} will be appeared not only in prepare message, but also in commit message, - * and consumer received message also contains this field. - */ - Headers setTransactionId(String transactionId); - - /** - * 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. - */ - Headers setCorrelationId(String correlationId); - - /** - * 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. - */ - Headers setCompression(short compression); - - /** - * See {@link Headers#setDestination(String)} - * - * @return destination - */ - String getDestination(); - - /** - * See {@link Headers#setMessageId(String)} - * - * @return messageId - */ - String getMessageId(); - - /** - * See {@link Headers#setBornTimestamp(long)} - * - * @return bornTimestamp - */ - long getBornTimestamp(); - - /** - * See {@link Headers#setBornHost(String)} - * - * @return bornHost - */ - String getBornHost(); - - /** - * See {@link Headers#setStoreTimestamp(long)} - * - * @return storeTimestamp - */ - long getStoreTimestamp(); - - /** - * See {@link Headers#setStoreHost(String)} - * - * @return storeHost - */ - String getStoreHost(); - - /** - * See {@link Headers#setDelayTime(long)} - * - * @return delayTime - */ - long getDelayTime(); - - /** - * See {@link Headers#setExpireTime(long)} - * - * @return expireTime - */ - long getExpireTime(); - - /** - * See {@link Headers#setPriority(short)} - * - * @return priority - */ - short getPriority(); - - /** - * See {@link Headers#setDurability(short)} - * - * @return durability - */ - short getDurability(); - - /** - * See {@link Headers#setMessageKey(String)} - * - * @return messageKey - */ - String getMessageKey(); - - /** - * See {@link Headers#setTraceId(String)} - * - * @return traceId - */ - String getTraceId(); - - /** - * See {@link Headers#setDeliveryCount(int)} - * - * @return deliveryCount - */ - int getDeliveryCount(); - - /** - * See {@link Headers#setTransactionId(String)} - * - * @return transactionId - */ - String getTransactionId(); - - /** - * See {@link Headers#setCorrelationId(String)} - * - * @return correlationId - */ - String getCorrelationId(); - - /** - * See {@link Headers#setCompression(short)} - * - * @return compression - */ - short getCompression(); - - } - - /** - * Returns all the system header fields of the {@code Message} object as a {@code KeyValue}. - * - * @return the system headers of a {@code Message} - */ - Headers headers(); - - /** - * 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); - -} \ No newline at end of file diff --git a/openmessaging-api/src/main/java/io/openmessaging/MessagingAccessPoint.java b/openmessaging-api/src/main/java/io/openmessaging/MessagingAccessPoint.java index 96cf2a20..80bee1dd 100644 --- a/openmessaging-api/src/main/java/io/openmessaging/MessagingAccessPoint.java +++ b/openmessaging-api/src/main/java/io/openmessaging/MessagingAccessPoint.java @@ -19,11 +19,15 @@ import io.openmessaging.consumer.Consumer; import io.openmessaging.consumer.MessageListener; +import io.openmessaging.consumer.PullConsumer; +import io.openmessaging.consumer.PushConsumer; import io.openmessaging.exception.OMSRuntimeException; import io.openmessaging.exception.OMSSecurityException; import io.openmessaging.manager.ResourceManager; +import io.openmessaging.message.MessageFactory; import io.openmessaging.producer.Producer; import io.openmessaging.producer.TransactionStateCheckListener; +import java.util.Collection; /** * An instance of {@code MessagingAccessPoint} may be obtained from {@link OMS}, which is capable of creating {@code @@ -90,15 +94,43 @@ public interface MessagingAccessPoint { Producer createProducer(TransactionStateCheckListener transactionStateCheckListener); /** - * Creates a new {@code PushConsumer} for the specified {@code MessagingAccessPoint}. The returned {@code Consumer} - * isn't bind to any queue, uses {@link Consumer#bindQueue(String, MessageListener)} to bind queues. + * Creates a new {@code PushConsumer} for the specified {@code MessagingAccessPoint}. + * The returned {@code PushConsumer} isn't attached to any queue, + * uses {@link PushConsumer#bindQueue(Collection, MessageListener)} to attach queues. * * @return the created {@code PushConsumer} - * @throws OMSRuntimeException if the {@code MessagingAccessPoint} fails to handle this request due to some internal - * error - * @throws OMSSecurityException if have no authority to create a consumer. + * @throws OMSRuntimeException if the {@code MessagingAccessPoint} fails to handle this request + * due to some internal error + */ + PushConsumer createPushConsumer(); + + /** + * Creates a new {@code PullConsumer} for the specified {@code MessagingAccessPoint}. + * + * @return the created {@code PullConsumer} + * @throws OMSRuntimeException if the {@code MessagingAccessPoint} fails to handle this request + * due to some internal error + */ + PullConsumer createPullConsumer(); + + /** + * Creates a new {@code PushConsumer} for the specified {@code MessagingAccessPoint} with some preset attributes. + * + * @param attributes the preset attributes + * @return the created {@code PushConsumer} + * @throws OMSRuntimeException if the {@code MessagingAccessPoint} fails to handle this request + * due to some internal error */ - Consumer createConsumer(); + PushConsumer createPushConsumer(KeyValue attributes); + + /** + * Creates a new {@code PullConsumer} for the specified {@code MessagingAccessPoint}. + * + * @return the created {@code PullConsumer} + * @throws OMSRuntimeException if the {@code MessagingAccessPoint} fails to handle this request + * due to some internal error + */ + PullConsumer createPullConsumer(KeyValue attributes); /** * Gets a lightweight {@code ResourceManager} instance from the specified {@code MessagingAccessPoint}. @@ -109,4 +141,13 @@ public interface MessagingAccessPoint { * @throws OMSSecurityException if have no authority to obtain a resource manager. */ ResourceManager resourceManager(); + + /** + * Gets a {@link MessageFactory} instance from the specified {@code MessagingAccessPoint}. + * + * @return the resource manger + * @throws OMSRuntimeException if the {@code MessagingAccessPoint} fails to handle this request due to some internal + * error + */ + MessageFactory messageFactory(); } diff --git a/openmessaging-api/src/main/java/io/openmessaging/OMSBuiltinKeys.java b/openmessaging-api/src/main/java/io/openmessaging/OMSBuiltinKeys.java index 5702724b..7d47d9a6 100644 --- a/openmessaging-api/src/main/java/io/openmessaging/OMSBuiltinKeys.java +++ b/openmessaging-api/src/main/java/io/openmessaging/OMSBuiltinKeys.java @@ -41,6 +41,11 @@ public interface OMSBuiltinKeys { */ String ACCOUNT_ID = "ACCOUNT_ID"; + /** + * The {@code ACCOUNT_KEY} key shows the specified account key in OMS attribute. + */ + String ACCOUNT_KEY = "ACCOUNT_KEY"; + /** * The {@code REGION} key shows the specified region in OMS driver schema. */ diff --git a/openmessaging-api/src/main/java/io/openmessaging/ServiceLifeState.java b/openmessaging-api/src/main/java/io/openmessaging/ServiceLifeState.java index 66e90965..eafdb18f 100644 --- a/openmessaging-api/src/main/java/io/openmessaging/ServiceLifeState.java +++ b/openmessaging-api/src/main/java/io/openmessaging/ServiceLifeState.java @@ -43,10 +43,10 @@ public enum ServiceLifeState { /** * Service is stopping. */ - STOPING, + STOPPING, /** * Service has been stopped. */ - STOPED, + STOPPED, } diff --git a/openmessaging-api/src/main/java/io/openmessaging/ServiceLifecycle.java b/openmessaging-api/src/main/java/io/openmessaging/ServiceLifecycle.java index 78a17418..3066cd56 100644 --- a/openmessaging-api/src/main/java/io/openmessaging/ServiceLifecycle.java +++ b/openmessaging-api/src/main/java/io/openmessaging/ServiceLifecycle.java @@ -18,6 +18,7 @@ package io.openmessaging; import io.openmessaging.consumer.Consumer; +import io.openmessaging.extension.Extension; import io.openmessaging.producer.Producer; /** @@ -32,7 +33,7 @@ * @version OMS 1.0.0 * @since OMS 1.0.0 */ -public interface ServiceLifecycle { +public interface ServiceLifecycle extends Extension { /** * Used for startup or initialization of a service endpoint. A service endpoint instance will be in a ready state * after this method has been completed. diff --git a/openmessaging-api/src/main/java/io/openmessaging/annotation/Optional.java b/openmessaging-api/src/main/java/io/openmessaging/annotation/Optional.java new file mode 100644 index 00000000..b5b9d490 --- /dev/null +++ b/openmessaging-api/src/main/java/io/openmessaging/annotation/Optional.java @@ -0,0 +1,43 @@ +/* + * 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.annotation; + +import java.lang.annotation.Documented; +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +/** + *

+ * 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. + *

+ * + * @version OMS 1.0.0 + * @since OMS 1.0.0 + */ +@Documented +@Retention(RetentionPolicy.RUNTIME) +@Target({ElementType.PACKAGE, ElementType.TYPE, ElementType.FIELD, ElementType.METHOD, ElementType.PARAMETER, ElementType.LOCAL_VARIABLE}) +public @interface Optional { +} diff --git a/openmessaging-api/src/main/java/io/openmessaging/consumer/BatchMessageListener.java b/openmessaging-api/src/main/java/io/openmessaging/consumer/BatchMessageListener.java new file mode 100644 index 00000000..09044d65 --- /dev/null +++ b/openmessaging-api/src/main/java/io/openmessaging/consumer/BatchMessageListener.java @@ -0,0 +1,59 @@ +/* + * 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.exception.OMSRuntimeException; +import io.openmessaging.message.Message; +import java.util.List; + +/** + * A message listener can implement this {@code BatchMessageListener} interface and register itself to a consumer + * instance to asynchronously receive messages. + * + * @version OMS 1.0.0 + * @since OMS 1.0.0 + */ +public interface BatchMessageListener { + /** + * Callback method to receive incoming messages. + *

+ * A message listener should handle different types of {@code BatchMessage}. + * + * @param batchMessage the received batchMessage. + */ + void onReceived(List batchMessage, Context context); + + interface Context { + /** + * Acknowledges the specified and consumed message, which is related to this {@code MessageContext}. + *

+ * 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 getBindQueues(); /** * Adds a {@code ConsumerInterceptor} instance to this consumer. @@ -132,27 +59,14 @@ public interface Consumer extends ServiceLifecycle { */ void removeInterceptor(ConsumerInterceptor interceptor); - /** - * 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 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. *

* 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 queueNames); + + /** + * Unbind the {@code Consumer} from a collection of queues. + *

+ * 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 queueNames); + + /** + * Receives the next message from the attached queues of this consumer. + *

+ * 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 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. + *

+ * 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 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}. + *

+ * + * @param timeout receive messages will blocked at most 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. + */ + List batchReceive(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}. + *

+ * + * @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 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. + */ + @Optional + List batchReceive(String queueName, QueueMetaData queueMetaData, MessageReceipt messageReceipt, + long timeout); + + /** + * Acknowledges the specified and consumed message with the unique message receipt handle. + *

+ * 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 queueNames, MessageListener listener); + + + /** + * Bind the {@code Consumer} to a collection of queue, with a {@code BatchMessageListener}. + *

+ * {@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 queueNames, BatchMessageListener listener); + + /** + * Unbind the {@code Consumer} from a collection of queues. + *

+ * 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 queueNames); + +} diff --git a/openmessaging-api/src/main/java/io/openmessaging/exception/OMSUnsupportException.java b/openmessaging-api/src/main/java/io/openmessaging/exception/OMSUnsupportException.java new file mode 100644 index 00000000..5720d452 --- /dev/null +++ b/openmessaging-api/src/main/java/io/openmessaging/exception/OMSUnsupportException.java @@ -0,0 +1,49 @@ +/* + * 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.exception; + +import io.openmessaging.annotation.Optional; + +/** + * The {@code OMSUnsupportException} must be thrown when the specified methods, headers or properties have not been + * provided by vendors, these methods or headers are usually marked by {@link Optional}. + * + * @version OMS 1.0.0 + * @since OMS 1.0.0 + */ +public class OMSUnsupportException extends OMSRuntimeException { + /** + * @see OMSUnsupportException#OMSUnsupportException(int, String) + */ + public OMSUnsupportException(int errorCode, String message) { + super(errorCode, message); + } + + /** + * @see OMSUnsupportException#OMSUnsupportException(int, Throwable) + */ + public OMSUnsupportException(int errorCode, Throwable cause) { + super(errorCode, cause); + } + + /** + * @see OMSUnsupportException#OMSUnsupportException(int, String, Throwable) + */ + public OMSUnsupportException(int errorCode, String message, Throwable cause) { + super(errorCode, message, cause); + } +} diff --git a/openmessaging-api/src/main/java/io/openmessaging/extension/Extension.java b/openmessaging-api/src/main/java/io/openmessaging/extension/Extension.java new file mode 100644 index 00000000..56edceb9 --- /dev/null +++ b/openmessaging-api/src/main/java/io/openmessaging/extension/Extension.java @@ -0,0 +1,51 @@ +/* + * 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 io.openmessaging.exception.OMSDestinationException; +import io.openmessaging.exception.OMSRuntimeException; +import io.openmessaging.exception.OMSSecurityException; +import io.openmessaging.exception.OMSTimeOutException; +import java.util.Set; + +/** + *

+ * This interface contains some methods are used for getting configurations related implementation. but this interface + * are not mandatory. + *

+ * + * @version OMS 1.0.0 + * @since OMS 1.0.0 + */ +@Optional +public interface Extension { + + /** + * This method used for getting the related queue's meta data, and this method is optional, vendors may not provide + * this method based on their implementation. + *

+ * + * @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 getQueueMetaData(String queueName); +} diff --git a/openmessaging-api/src/main/java/io/openmessaging/extension/ExtensionHeader.java b/openmessaging-api/src/main/java/io/openmessaging/extension/ExtensionHeader.java new file mode 100644 index 00000000..20e2b32e --- /dev/null +++ b/openmessaging-api/src/main/java/io/openmessaging/extension/ExtensionHeader.java @@ -0,0 +1,201 @@ +/* + * 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 io.openmessaging.message.Message; + +/** + *

+ * 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. + *

+ * + * @version OMS 1.0.0 + * @since OMS 1.0.0 + */ +@Optional +public interface ExtensionHeader { + /** + * The {@code PARTITION} in extension header field contains the partition of target destination which the message + * is being sent. + *

+ * + * 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. + *

+ * So the same {@code TRANSACTION_ID} will be appeared not only in prepare message, but also in commit message, and + * consumer received message also contains this field. + */ + ExtensionHeader setTransactionId(String transactionId); + + /** + * 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. + */ + 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. + *

+ * And, the messaging system only plays exchange role in a distributed system in most cases, so the TraceID can be + * used to trace the whole call link with other parts in the whole system. + */ + ExtensionHeader setTraceId(String traceId); + + /** + * The {@code DELAY_TIME} header field contains a number that represents the delayed times in milliseconds. + *

+ * The message will be delivered after delayTime milliseconds starting from {@code BORN_TIMESTAMP} . When this filed + * isn't set explicitly, this means this message should be delivered immediately. + */ + ExtensionHeader setDelayTime(long delayTime); + + /** + * The {@code EXPIRE_TIME} header field contains the expiration time, it represents a time-to-live value. + *

+ * 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. + *

+ */ + ExtensionHeader setExpireTime(long expireTime); + + /** + * This method will return the partition of this message belongs. + *

+ * + * @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 properties; @Override - public int getShort(String key) { - return 0; + public KeyValue put(String key, boolean value) { + properties.put(key, String.valueOf(value)); + return this; + } + + @Override + public boolean getBoolean(String key) { + if (!properties.containsKey(key)) { + return false; + } + return Boolean.valueOf(properties.get(key)); + } + + @Override + public boolean getBoolean(String key, boolean defaultValue) { + return properties.containsKey(key) ? getBoolean(key) : defaultValue; + } + + @Override + public short getShort(String key) { + if (!properties.containsKey(key)) { + return 0; + } + return Short.valueOf(properties.get(key)); } @Override - public int getShort(String key, short defaultValue) { - return 0; + public short getShort(String key, short defaultValue) { + return properties.containsKey(key) ? getShort(key) : defaultValue; } @Override diff --git a/openmessaging-api/src/main/java/io/openmessaging/message/Header.java b/openmessaging-api/src/main/java/io/openmessaging/message/Header.java new file mode 100644 index 00000000..1017c687 --- /dev/null +++ b/openmessaging-api/src/main/java/io/openmessaging/message/Header.java @@ -0,0 +1,183 @@ +/* + * 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.extension.ExtensionHeader; + +/** + * The {@code Header} interface is the root interface of all OMS messages, and the most commonly used by OMS message + * {@link Message}. + *

+ * 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. + *

    + *
  • + * System Headers, OMS defines some standard attributes that represent the characteristics of the message. + *
  • + * + *
+ * The body contains the application data being transmitted, which is generally ignored by the messaging system and + * simply transmitted to its destination. + *

+ * + * 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: + *

    + *
  • + * PERSISTENT, the persistent mode instructs the vendor should provide stable storage to ensure the message won't be + * lost. + *
  • + *
  • + * NON_PERSISTENT, this mode does not require the message be logged to stable storage, in most cases, the memory + * storage is enough for better performance and lower cost. + *
  • + *
+ */ + 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}. + *

+ * 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. + *

    + *
  • + * System Headers, OMS defines some standard attributes that represent the characteristics of the message. + *
  • + *
  • + * User properties, some OMS vendors may require enhanced extra attributes of the message or some users may want to + * clarify some customized attributes to draw the body. OMS provides the improved scalability for these scenarios. + *
  • + *
+ * The body contains the application data being transmitted, which is generally ignored by the messaging system and + * simply transmitted to its destination. + *

+ * 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 messages); + /** + * Send messages to the specified destination asynchronously, the destination should be preset to {@link + * 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. + * + * @param messages a batch messages will be sent. + * @return the {@code Promise} of an asynchronous messages send operation. + * @see Future + * @see FutureListener + */ + Future sendAsync(List messages); + + /** + *

+ * 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 messages); + /** * Adds a {@code ProducerInterceptor} to intercept send operations of producer. * @@ -104,7 +128,7 @@ public interface Producer extends MessageFactory, ServiceLifecycle { void addInterceptor(ProducerInterceptor interceptor); /** - * Removes a {@code ProducerInterceptor}. + * Remove a {@code ProducerInterceptor}. * * @param interceptor a producer interceptor will be removed. */ @@ -112,7 +136,7 @@ public interface Producer extends MessageFactory, ServiceLifecycle { /** * Sends a transactional message to the specified destination synchronously, the destination should be preset to - * {@link Message#headers()}, other header fields as well. + * {@link Message#header()}, other header fields as well. *

* 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 @@ io.openmessaging parent - 1.0.0-Preview + 1.0.0-beta-SNAPSHOT pom openmessaging @@ -51,8 +51,8 @@ UTF-8 - 1.6 - 1.6 + 1.8 + 1.8