TransportMessageBus

abstract class TransportMessageBus<M : Message<*, *>, AggregateIdCapable, NamedAggregate, E : MessageExchange<*, M>>(val transport: Transport, topicNaming: TopicNaming, decodeFailureHandler: TransportDecodeFailureHandler = TransportDecodeFailureHandler.FAIL) : DistributedMessageBus<M, E>

A distributed message bus over a Transport: the one place that names topics, encodes and decodes messages, validates received records and builds exchanges, whatever the broker.

The wire format is the v9 one on every backend: the topic is topicNaming's (memoised per aggregate), the key is the aggregate ID, the timestamp is the message's creation time and the payload is the message's JSON. A received record must decode to a messageType whose aggregate ID equals the record key (when the backend has keys) and whose topic is the record's; anything else goes to decodeFailureHandler. Decoded messages are read-only.

Inheritors

Constructors

Link copied to clipboard
constructor(transport: Transport, topicNaming: TopicNaming, decodeFailureHandler: TransportDecodeFailureHandler = TransportDecodeFailureHandler.FAIL)

Types

Link copied to clipboard
object Companion

Properties

Link copied to clipboard
abstract val messageType: Class<M>
Link copied to clipboard

The transport's resources for the runtime to close (Transport.runtimeResource).

Link copied to clipboard

Functions

Link copied to clipboard
open override fun close()

Closes the message bus and releases any resources. Default implementation does nothing.

Link copied to clipboard
fun decode(record: TransportRecord): M

The message record carries, read-only.

Link copied to clipboard
fun encode(message: M): TransportMessage

message as it goes on the wire. Makes message read-only.

Link copied to clipboard
open override fun receiver(subscription: MessageSubscription): MessageReceiver<E>

The one receive entry: a single message source for subscription, with an explicit transport readiness boundary and processing admission (see MessageReceiver).

Link copied to clipboard
open override fun send(message: M): Mono<Void>

Sends a message through the message bus.

Link copied to clipboard
fun topicOf(namedAggregate: NamedAggregate): String