-
Notifications
You must be signed in to change notification settings - Fork 0
Server Bus
Package com.evento.server.bus.* in evento-server. This is the broker: it accepts bundle
connections, authenticates them, tracks who can handle what, routes messages, and guarantees
exactly-once delivery across reconnects.
It replaced v1's single 1 099-line MessageBus with composable pieces.
| Class | Role |
|---|---|
BusLifecycle |
The orchestrator. start(port) / stop(Duration). Owns handshake, routing, correlation, reconnect buffer |
BusFacade |
Spring SPI used by every server-side consumer (Dashboard, ClusterStatus, AutoDiscovery, Consumer, Bundle) |
BusLifecycleFacade |
BusFacade implementation wrapping BusLifecycle
|
BusFacadeConfiguration |
Unconditional Spring @Configuration; wires BusLifecycleFacade as the only path |
BusConfiguration + BusProperties
|
Spring Boot auto-config and the evento.server.bus.* properties |
ConnectionRegistry |
ConcurrentHashMap<NodeAddress, (Connection, token)>. register supersedes and closes the old transport; unregister is token-guarded |
ClusterRegistry |
payloadType → Set<NodeAddress>, with RANDOM / FIRST pick strategy |
CorrelationStore |
ConcurrentHashMap<UUID, PendingCorrelation>. Bounded shutdown, scheduler-driven expiry |
ForwardingTable |
correlationId → (originatorAddress, destinationAddress) |
ForwardingDedupCache |
LRU, 5 min TTL, 50 k entries. Deduplicates retried Request at the broker |
HandshakeHandler |
Validates Hello → TokenValidator → BusLifecycle.onHello
|
TokenValidator |
SPI: acceptAll() (default) or sharedSecret(token) with constant-time compare |
MessageRouter |
OCP dispatcher — Map<Class<? extends Message>, Handler> registered in start()
|
BundleSession |
Per-connection mutable state: NodeAddress, registered handlers, enable/disable flag |
BusEvent |
Sealed hierarchy: BundleRegistered, BundleDiscovered, BundleLeft, AdminNotification, … |
BusEventBus |
subscribe(Consumer<BusEvent>) — replaces v1's four parallel listener lists |
This distinction is load-bearing and easy to conflate:
CorrelationStore |
ForwardingTable |
|
|---|---|---|
| Used for | Server-initiated requests | Bundle-A → server → bundle-B relays |
| Holds | A CompletableFuture to await |
Just the originator and destination addresses |
| Why | The server itself is waiting for the answer | The server is a postman, not a caller |
On disconnect the forwarding table drains via drainByDestination(addr), which removes only
destination-side entries. Originator-side entries are deliberately left alive so that a response
arriving for a bundle that has since disconnected can still be delivered when it reconnects.
Draining "everything involving this address" would silently break that guarantee.
Three mechanisms compose so that a caller bundle can retry with the same correlationId and receive
exactly one side-effect plus the same response back:
| Layer | Mechanism |
|---|---|
| Broker, incoming dedup |
ForwardingDedupCache — LRU, 5 min TTL, 50 k entries. A duplicate Request with the same correlationId replays the cached Response; an in-flight duplicate is silently dropped |
| Bundle handler side |
ProcessedRequestCache — resolveOrClaim returns Claimed / InFlight / Replay. The handler runs at most once per correlationId
|
| Reconnect delivery |
reconnectBuffer on BusLifecycle buffers responses for disconnected originators (keyed by instanceId, 2 min TTL); deliverPendingResponses replays them on re-handshake |
Two fixes here have subtle failure modes worth knowing about:
Superseded-session race (Fix A). When a bundle reconnects, the new session supersedes the old one
in ConnectionRegistry. If the old transport's disconnect callback then fired and blindly cleaned
up, it would wipe the handlers of the freshly reconnected bundle — the bundle looks connected but
receives nothing. So onTransportDisconnected returns early when ConnectionRegistry.unregister
returns empty (an empty return means token mismatch, i.e. this session was superseded). The skip is
logged as event=disconnect_superseded_skip.
In-flight responses across a disconnect (Fix D). When a handler replies while the originator is
disconnected, BusLifecycle.onResponse buffers the Response in the reconnect buffer rather than
dropping it. onHello calls deliverPendingResponses on re-handshake.
DEGRADED is not disconnected (Fix C). A channel over its high-water mark transitions to
DEGRADED, but ConnectionState.canSend() still returns true — the TCP socket is alive and
backpressure is advisory. Treating DEGRADED as unsendable dropped in-flight responses.
CommandBrokerHandler (com.evento.server.es) is a Spring @Component subscribing to
BusEvent.BundleDiscovered. For each AggregateCommandHandler and service CommandHandler payload
type it registers a LocalRequestHandler in BusLifecycle — so a command is intercepted by the
server rather than blindly relayed. The full flow is in
Architecture Overview § 5.
PgDistributedLock guards aggregate command execution. When the DataSource is null it falls back
to a JVM-only semaphore, which keeps single-JVM and embedded scenarios (integration tests) working.
Request handling runs on BusBusinessExecutor, which grows to max before it queues — a plain
ThreadPoolExecutor does the opposite and would leave the configured maximum unreachable under
exactly the load it exists for. This is the single most important thing to understand about server
capacity, and it has a dedicated page: Throughput and Capacity.
- Wire Protocol — the frames the bus routes
- Bundle Client — the other end of every connection
- Observability — the meters and log events this layer emits
Evento Framework — Copyright 2020–2026 © Gabor Galazzo. Dual-licensed under AGPL-3.0 and a commercial licence.
This wiki documents the implementation; the repository is authoritative where the two disagree. Found something out of date? Open an issue.
Getting oriented
Internals
Operations
- Server Configuration
- Throughput and Capacity
- Observability
- Security Model
- Server REST API
- Troubleshooting
Project