Package-level declarations

Types

Link copied to clipboard

A Transport inside one JVM with broker semantics: every consumer group on a topic receives each record, and the members of a group share them round-robin. Records go through their encoded form, so a bus over it exercises the same encode, decode and validation as one over Kafka or Redis.

Link copied to clipboard
fun interface TopicNaming

The topic that carries one kind of message of an aggregate on one backend. Topic names are frozen wire format: Kafka wow.<context>.<aggregate>.command, Redis <context>.<aggregate>:command, and likewise per kind.

Link copied to clipboard

A message broker as Wow uses it: strings in, strings out.

Link copied to clipboard
open class TransportCommandBus(transport: Transport, topicNaming: TopicNaming, decodeFailureHandler: TransportDecodeFailureHandler = TransportDecodeFailureHandler.FAIL) : TransportMessageBus<CommandMessage<*>, ServerCommandExchange<*>> , DistributedCommandBus

The distributed command bus over a Transport.

Link copied to clipboard

The receive stream's error for TransportDecodeFailureAction.FAIL. It names the record and the cause's type but carries neither the cause nor the payload, since a decoder's message can quote the payload.

Link copied to clipboard
class TransportDecodeFailure(val record: TransportRecord, val group: String, val messageType: Class<*>, val cause: Exception)

A record that TransportMessageBus could not decode: its payload is missing or not a messageType, or its key or topic does not match the decoded message.

Link copied to clipboard

What the bus does with an undecodable record once TransportDecodeFailureHandler has seen it.

Link copied to clipboard

Decides what happens to a record the bus cannot decode. An error from handle fails the receive stream like TransportDecodeFailureAction.FAIL.

Link copied to clipboard
open class TransportDomainEventBus(transport: Transport, topicNaming: TopicNaming, decodeFailureHandler: TransportDecodeFailureHandler = TransportDecodeFailureHandler.FAIL) : TransportMessageBus<DomainEventStream, EventStreamExchange> , DistributedDomainEventBus

The distributed domain event bus over a Transport.

Link copied to clipboard
class TransportEventStreamExchange(val message: DomainEventStream, val record: TransportRecord, val attributes: MutableMap<String, Any> = ConcurrentHashMap()) : EventStreamExchange

An event stream received from a Transport; acknowledging it acknowledges its record.

Link copied to clipboard
class TransportFailurePolicy(val receiveRetry: Retry = receiveRetry())

How a Transport reacts when its receive stream fails, the same for every backend.

Link copied to clipboard
data class TransportMessage(val topic: String, val key: String, val payload: String, val timestamp: Long, val headers: Map<String, String> = emptyMap())

One outbound record. Records with one key stay in order; timestamp is the message's creation time in epoch milliseconds. A backend writes what it supports: Kafka writes all of it, Redis Streams only topic and payload.

Link copied to clipboard
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.

Link copied to clipboard

The consumer Transport.open returns.

Link copied to clipboard

One received record. ack confirms it to the broker (Kafka offset commit, Redis XACK); nack gives it up without confirming, so the broker's own redelivery applies (a Kafka offset that is never committed, a Redis entry that stays pending).

Link copied to clipboard

A record whose payload decodes but whose key or topic does not match the decoded message: it was published under another aggregate's key or topic.

Link copied to clipboard
class TransportServerCommandExchange<C : Any>(val message: CommandMessage<C>, val record: TransportRecord, val attributes: MutableMap<String, Any> = ConcurrentHashMap()) : ServerCommandExchange<C>

A command received from a Transport; acknowledging it acknowledges its record.

Link copied to clipboard
open class TransportStateEventBus(transport: Transport, topicNaming: TopicNaming, decodeFailureHandler: TransportDecodeFailureHandler = TransportDecodeFailureHandler.FAIL) : TransportMessageBus<StateEvent<*>, StateEventExchange<*>> , DistributedStateEventBus

The distributed state event bus over a Transport.

Link copied to clipboard
class TransportStateEventExchange<S : Any>(val message: StateEvent<S>, val record: TransportRecord, val attributes: MutableMap<String, Any> = ConcurrentHashMap()) : StateEventExchange<S>

A state event received from a Transport; acknowledging it acknowledges its record.

Properties

Link copied to clipboard

Functions

Link copied to clipboard
fun TopicNaming.memoized(maxAggregates: Int = DEFAULT_MEMOIZED_AGGREGATES): TopicNaming

This naming computed once per aggregate, which holds because a topic depends only on the context and aggregate names. At most maxAggregates aggregates are kept, so records naming arbitrary aggregates cannot grow it without bound; beyond that a topic is computed on every call.