diff options
author | Jakob Odersky <jakob@odersky.com> | 2018-07-31 12:13:47 -0700 |
---|---|---|
committer | GitHub <noreply@github.com> | 2018-07-31 12:13:47 -0700 |
commit | e6552f3b31b55396c652c196c5c3a9c3a6cfed71 (patch) | |
tree | f6826eac8bff8470683547006d1e64d2bc425d55 /src/main/scala/xyz/driver/core/messaging/Bus.scala | |
parent | db0c9bebee4cbc587d4da0a624f671ffcf7a649f (diff) | |
download | driver-core-e6552f3b31b55396c652c196c5c3a9c3a6cfed71.tar.gz driver-core-e6552f3b31b55396c652c196c5c3a9c3a6cfed71.tar.bz2 driver-core-e6552f3b31b55396c652c196c5c3a9c3a6cfed71.zip |
Add message bus and topic abstractions (#181)v1.12.0
Diffstat (limited to 'src/main/scala/xyz/driver/core/messaging/Bus.scala')
-rw-r--r-- | src/main/scala/xyz/driver/core/messaging/Bus.scala | 70 |
1 files changed, 70 insertions, 0 deletions
diff --git a/src/main/scala/xyz/driver/core/messaging/Bus.scala b/src/main/scala/xyz/driver/core/messaging/Bus.scala new file mode 100644 index 0000000..e2ee76a --- /dev/null +++ b/src/main/scala/xyz/driver/core/messaging/Bus.scala @@ -0,0 +1,70 @@ +package xyz.driver.core +package messaging + +import scala.concurrent._ + +/** Base trait for representing message buses. + * + * Message buses are expected to provide "at least once" delivery semantics and are + * expected to retry delivery when a message remains unacknowledged for a reasonable + * amount of time. */ +trait Bus { + + /** Type of unique message identifiers. Usually a string or UUID. */ + type MessageId + + /** Most general kind of message. Any implementation of a message bus must provide + * the fields and methods specified in this trait. */ + trait BasicMessage[A] { + + /** All messages must have unique IDs so that they can be acknowledged unambiguously. */ + def id: MessageId + + /** Actual message content. */ + def data: A + + } + + /** Actual type of messages provided by this bus. This must be a subtype of BasicMessage + * (as that defines the minimal required fields of a messages), but may be refined to + * provide bus-specific additional data. */ + type Message[A] <: BasicMessage[A] + + /** Type of a bus-specific configuration object can be used to tweak settings of subscriptions. */ + type SubscriptionConfig + + /** Default value for a subscription configuration. It is such that any service will have a unique subscription + * for every topic, shared among all its instances. */ + val defaultSubscriptionConfig: SubscriptionConfig + + /** Maximum amount of messages handled in a single retrieval call. */ + val defaultMaxMessages = 64 + + /** Retrieve any new messages in the mailbox of a subscription. + * + * Any retrieved messages become "outstanding" and should not be returned by this function + * again until a reasonable (bus-specific) amount of time has passed and they remain unacknowledged. + * In that case, they will again be considered new and will be returned by this function. + * + * Note that although outstanding and acknowledged messages will eventually be removed from + * mailboxes, no guarantee can be made that a message will be delivered only once. */ + def fetchMessages[A]( + topic: Topic[A], + config: SubscriptionConfig = defaultSubscriptionConfig, + maxMessages: Int = defaultMaxMessages): Future[Seq[Message[A]]] + + /** Acknowledge that a given series of messages has been handled and should + * not be delivered again. + * + * Note that messages become eventually acknowledged and hence may be delivered more than once. + * @see fetchMessages() + */ + def acknowledgeMessages(messages: Seq[MessageId]): Future[Unit] + + /** Send a series of messages to a topic. + * + * The returned future will complete once messages have been accepted to the underlying bus. + * No guarantee can be made of their delivery to subscribers. */ + def publishMessages[A](topic: Topic[A], messages: Seq[A]): Future[Unit] + +} |