LocalFirstDistributedCopies

class LocalFirstDistributedCopies(name: String = "LocalFirstDistributedCopies", metrics: WowMetrics = WowMetrics.NONE, backlogHighWaterMark: Int = DEFAULT_BACKLOG_HIGH_WATER_MARK, handOffTimeout: Duration = DEFAULT_HAND_OFF_TIMEOUT) : RuntimeComponent

The distributed copies of a LocalFirstMessageBus, sent in send order per aggregate without the local-first sender waiting for a local receiver's demand (decision D3).

Every local-first send enqueues its copy in the queue of its aggregate (by aggregate ID). A copy is sent only after every earlier copy of that aggregate was sent, so the distributed bus sees an aggregate's messages in the order they were sent, as before 9.3.0. Each copy waits, in its turn, for the local decision (Copy.decide): a handed-off message waits for its admission (local_first=true when admitted, false when rejected), a refused one goes out at once with local_first=false. A failed copy is logged and counted by the distributed bus's send metrics; only a refused sender, which waits for its own copy, sees the failure.

The backlog — copies not yet sent — is a gauge per aggregate type (wow.local_first.backlog); reaching backlogHighWaterMark logs a warning. A slow local consumer makes it, and the in-process local sink, grow.

name must be unique per meter registry (the starter names one per bus): the gauge of a second instance with the same name would read the first one's backlog.

As a RuntimeComponent stopped after the dispatchers and before the transports' runtime resources, it lets every queued copy be sent within the runtime's shutdown deadline (stopGracefully); a force stop cancels the rest, which other services then never receive. Only an instance registered with the runtime (the starter registers the ones it creates) is awaited at shutdown; one created by hand is not, unless it is added to the runtime's components.

Constructors

Link copied to clipboard
constructor(name: String = "LocalFirstDistributedCopies", metrics: WowMetrics = WowMetrics.NONE, backlogHighWaterMark: Int = DEFAULT_BACKLOG_HIGH_WATER_MARK, handOffTimeout: Duration = DEFAULT_HAND_OFF_TIMEOUT)

Types

Link copied to clipboard
object Companion
Link copied to clipboard
inner class Copy

One queued copy; decide tells it the local hand-off result.

Properties

Link copied to clipboard

The copies not yet sent or failed.

Functions

Link copied to clipboard
fun backlog(namedAggregate: NamedAggregate): Int

The copies of namedAggregate's aggregates not yet sent.

Link copied to clipboard
fun enqueue(key: Any, namedAggregate: NamedAggregate, context: ContextView = Context.empty(), description: () -> String, send: (admitted: Boolean) -> Mono<Void>): LocalFirstDistributedCopies.Copy

Queues the copy of a message of the aggregate key (its aggregate ID), sent by send with the local decision once every earlier copy of that aggregate was sent, in the sender's context.

Link copied to clipboard
open override fun forceStop()

Releases resources promptly without blocking, and remains safe before prepare and across repeated or overlapping calls.

Link copied to clipboard
open override fun prepare(runtimeContext: RuntimeContext): Mono<Void>

Prepares this component without opening message processing.

Link copied to clipboard
open override fun start()
Link copied to clipboard
open override fun stopGracefully(): Mono<Void>
Link copied to clipboard
open override fun toString(): String