caishen

Message Bus — Distributed Event Backbone

RISE Framework Specification

Spec ID: 77 Version: 1.0 Document ID: caishen-rise-messagebus-v1.0 Last Updated: 2026-08-01 Depends on: 73 (events & hooks) Related: 60–63 (state machineries), 70 (PDSP architecture)

Spec 73 specifies which events exist and how they propagated in the shipped system (in-process delegates + shell hooks). This spec specifies the message bus — the durable, typed, cross-process transport that was built, proven in one subsystem, and left unwired everywhere else. It is the realization of requirements R1–R6 in spec 73.


Creative Intent

What The Message Bus Enables:

Desired Outcomes:

  1. A price refresh completing is announced, not dispatched to a known list
  2. A subscriber that was offline catches up rather than missing the event
  3. The set of subscribers is discoverable and changeable at runtime
  4. Publishing costs the publisher a bounded, non-blocking amount of work
  5. A message carries enough identity to be traced across the whole chain

What Was Actually Built

This is not a sketch. A working bus exists, with a proven end-to-end publisher/subscriber pair.

Layer Implementation
Messaging library Rebus, consumed through a private repackaging (PS.Common.Services.RebusKit 1.18.116.808)
Transport MSMQ — Microsoft Message Queuing, one queue per endpoint
Subscription store SQL Server, centralized, database messagebus on host zeus, table <endpoint>_subscriptions
Routing By .NET type. Publish-subscribe, topic-per-message-type
Envelope SMEventArgsBase — the same base type every state machine event derives from
Serialization Rebus default (JSON), Newtonsoft
Logging Serilog, with coloured-console / rolling-file / event-log sinks

The critical design decision: the bus message type and the state machine event type are the same type. An event raised by a component’s state machine can be published to the bus unchanged, and a subscriber in another process receives it as the same strongly-typed object. There is no separate DTO, no mapping layer, no schema drift between the in-process and cross-process representations.

That is the single most valuable property of this design and it must survive any re-implementation.


The Message Envelope

Every bus message derives from a common base carrying four identity fields.

MessageEnvelope:
    header:
        name       : string     -- the event's business name, e.g. "SECDBCreated"
        typeName   : string     -- the concrete type's name
        timeStamp  : instant    -- when the event was constructed

    eventId               : uuid    -- STABLE identity of the event KIND
    messageId             : uuid    -- unique per emission
    correlationId         : uuid    -- caller-assigned, links related messages
    correlationAssemblyId : uuid    -- identity of the emitting process instance

    busPublish            : bool    -- per-instance publication opt-in

    <payload fields, defined per event type>

eventId vs messageId — the distinction that matters

Requirements

Payload discipline

Event payloads in the shipped code carry domain objects directly — ChaosDataBuilderDTO, List<string>, PriceSeries. Two consequences:

  1. The payload’s serialized shape is whatever the domain object happens to serialize to, and it changes silently when the domain object changes.
  2. Large payloads travel through the queue. A chart-data event carrying a full builder object is orders of magnitude larger than a notification.

Requirement M6: distinguish notification events from document events.

Default to notification. The bar-identity design (spec 70) makes this nearly free — a povTlid is a complete, unambiguous reference to exactly the data the subscriber needs.


Publisher Contract

Interface BusPublisher:
    Publish<TMessage>(message: TMessage) -> void
        where TMessage : MessageEnvelope

One method. Deliberately minimal — the publisher does not name a destination, a topic, or a subscriber. Routing is a property of the message type.

The per-instance opt-in

Every event carries a busPublish flag set at construction. The generated dispatch site checks it:

if (busPublisher != null)
    if (event.busPublish)
        fireAndForget(() => busPublisher.Publish(event))

// then the in-process path proceeds regardless
dispatchToStateMachine(event)
raiseInProcessEvent(event)

Three properties follow, and they are worth understanding before changing anything:

  1. The bus is optional. A component with no publisher attached works normally; the check is a null test. A component can be used in a unit test, a desktop app, or a headless service with no infrastructure at all.
  2. Publication is per-event-instance, not per-event-type. The same event type can be published in one situation and kept local in another, decided by whoever constructs it. The generator’s own comment records the author wrestling with this: “How do I tell it not to publish it or do it after having generated the code / A dictionary of type of message with a flag true or false”.
  3. The in-process path never depends on the bus. Local subscribers get the event whether or not it is published.

Assessment. Properties 1 and 3 are correct and must be preserved. Property 2 is a mistake in this form. Whether an event crosses a process boundary is an architectural property of the event kind, not a per-call decision — and scattering it across every construction site guarantees inconsistency. In the shipped code every generated default is false, so most events are never published at all.

Requirement P1: declare publication in the model (spec 60), as an attribute of the event kind, alongside its eventId. Allow a runtime override for genuine special cases, but make the declared default authoritative and visible in one place.

Fire-and-forget is not free

The dispatch site wraps publication in an unobserved background task with no continuation, no error handling, and no awaiting:

Task t = Task.Run(() => { this.BusPublisher.PublishBusMessage(e); });

The intent is right — publication must not block or fail the producer. The implementation loses every failure silently, and the task is never observed.

Requirement P2: publication is asynchronous and never fails the producer, and every failure is recorded. Use an explicit outbound queue with a bounded buffer, a retry policy, and a failure counter, not an unobserved task.

Requirement P3 — the outbox. Where an event announces a committed state change (spec 72’s SeriesUpdateCompleted is the canonical case), publication must not be able to succeed while the commit fails, or vice versa. Write the outbound message to an outbox table in the same transaction as the data change, and publish from the outbox. This is the difference between “we usually notify” and “we always notify”, and it is what spec 73 R1 requires.


Subscriber Contract

Three subscription styles exist. They are not redundant; each fits a different consumer shape.

1. Receiver-based (the primary style)

Interface BusSubscriber:
    SubscribeWithReceiver<TMessage>(subscriptionId?) -> MessageReceiver<TMessage>

MessageReceiver<TMessage>:
    event MessageReceived(sender, message)
    Unsubscribe()

Subscribing returns a receiver object that raises an event when a message arrives, and knows how to unsubscribe itself. Usage:

var receiver = hub.SubscribeWithReceiver<ContentChangedEvent>()
receiver.MessageReceived += OnContentChanged
...
receiver.Unsubscribe()

This is the good design. The subscription is a value — it can be held, passed, and disposed. Compare with subscribing by registering a global handler, where unsubscribing means remembering what you registered.

Defect: subscriptionId is accepted, defaulted to a fresh GUID when empty — and then never used. The subscription cannot be named, so it cannot be resumed or inspected. Requirement S1: a subscription has a stable, caller-supplied name. That name is the durable-subscription key: a consumer restarting under the same name resumes its position; a new name starts fresh. This is the mechanism that makes spec 73 R1 (catch-up on restart) actually work.

2. Callback-based

Interface BusSubscriberCallback:
    SubscribeByCallback<TMessage>(callback: MessageReceivedCallback<TMessage>)

For consumers that prefer implementing an interface over wiring an event. Thin sugar over style 1.

3. “Custom” — messages outside the envelope

A parallel set of operations (PublishCustom, SubscribeWithReceiverCustom, MessageReceiverCustom, UnSubscribeCustom) exists for message types that do not derive from the envelope base. It was added so plain domain messages could travel on the bus without being retrofitted as state machine events.

Requirement S2: do not carry this forward as a parallel API. It doubles the surface and loses every identity field — no messageId, no correlationId, no eventId, so such messages are untraceable and undedupable. Instead, wrap: a plain payload is carried inside an envelope. One API, one envelope, arbitrary payloads.

Delivery semantics

The shipped configuration sets no retry policy, no dead-letter queue, and no error handling in the receive path. A handler that throws loses the message to whatever the transport’s default happens to be.

Requirements:


Topology

Endpoints, topics, and channels

Concept Meaning Original naming
Endpoint A queue owned by one process. Its address. <appTag>.<channel>, or <FQDN>.<channel>
Topic A logical grouping; determines the subscription table GDSTopic, WDSWritingGuide, + "_Subscriptions"
Message type The actual routing key the .NET type

Endpoint names are sanitized before use — dots and dashes stripped — because the transport constrains queue names. The WDSServiceHub composes the endpoint from the host’s fully-qualified name plus a channel name, so two instances on different hosts get distinct queues automatically. That is correct and worth keeping.

Requirement T1: endpoint identity = (service, instance). Two instances of one service must have distinct endpoints for competing-consumer work and a shared subscription name for load balancing. The original conflates these.

The hub pattern

WDSServiceHub wraps the bus with an explicit role:

HubRole = Publisher | Subscriber | Both | None

A single binary can run as either side, chosen at startup. The WDS store service demonstrates it: -tst runs it as a publisher that emits one message, -db runs it as a subscriber that waits. Same code, same message types, opposite roles.

Requirement T2: keep the explicit role. It makes the topology declarable and testable — a smoke test is “start the same binary twice, once each way”. None is genuinely useful: it disables the bus entirely for local development.

Transport variants built

Implementation Transport Status
MyBus171208 Rebus over MSMQ, SQL Server subscriptions Working. The real one
BusZure Azure Service Bus topics Publish only; subscribe throws
BDSBus, BDSBusSimple, BDSBusMsgIndependent Rebus/MSMQ, POV-specialized Earlier generations
DTSBus — Stub; Publish has a commented-out body
SharedServiceBus — Three fields, no behaviour

BusZure matters more than its completeness suggests: someone already established that the publisher contract can be satisfied by a cloud broker without changing a single caller. The abstraction is sound; only the implementation is missing.

Requirement T3: the transport is a swappable adapter behind the publisher/subscriber contracts. Prove it with at least two adapters — one in-memory for tests, one real broker.


Where It Was Wired

Subsystem Bus usage
WDS (writing guide) Fully wired. File watcher publishes content changes; a store service subscribes and persists. The one proven end-to-end path.
GDS Bus class + callback wrapper built; wired in the prototype app
STC Eventing prototype
Experimental apps DTSBusServiceExperimentalApp — the original proving ground
SE, IDS, TDS, SDS State machines expose BusPublisher and contain the publish dispatch — nothing assigns a publisher
PDSP (spec 70) Same: property present, never assigned

So: the bus works, and the trading path has the wiring point built into every generated component — an unassigned BusPublisher property and a publish dispatch already in the code path. Connecting the trading subsystems is assignment, not construction.

The generator’s own note records the intended final step:

“Might be interesting to think about using the default App bus like: DTSApp.PublishBusMessage(message) so we would be setting DTSApp as the default bus if none chosen”

Requirement W1: provide a resolvable default publisher so components get one without every construction site wiring it, while still allowing explicit injection and a null (bus-disabled) mode.


Target Architecture

Keep the contracts. Replace the transport and close the gaps.

   Producer
      │  commits state change + outbox row   (one transaction)          [P3]
      ▼
   ┌──────────────┐
   │   Outbox     │──── relay ────▶ ┌──────────────────────────┐
   └──────────────┘                 │   Broker                  │
                                    │   durable, partitioned,   │  [R1]
                                    │   replayable log          │
                                    └───────────┬──────────────┘
                                                │
                    ┌───────────────────────────┼───────────────────────┐
                    ▼                           ▼                       ▼
            named subscription          named subscription       hook adapter
            (indicator engine)          (chart data service)     (operator scripts,
                    │                           │                 spec 73 R5)
              idempotent [S4]             idempotent [S4]

Requirements summary

# Requirement
M1–M6 Envelope: keep four identity fields; populate correlationId; route on eventId; add schemaVersion; unambiguous instants; notification-by-default payloads
P1–P3 Publisher: declare publication in the model; async-but-observed; transactional outbox
S1–S6 Subscriber: named durable subscriptions; one envelope API; at-least-once; idempotent consumers; retry + dead-letter; poison isolation
T1–T3 Topology: distinct endpoints vs shared subscriptions; explicit roles; swappable transport adapters
W1 A resolvable default publisher

Transport selection

MSMQ is the blocking constraint: Windows-only, deprecated, and requiring a queue service on every participating host. Any broker satisfying the requirements works — the contracts do not name one, and BusZure already demonstrated a second binding. Selection criteria, in priority order:

  1. Durable, replayable storage (R1, S1)
  2. Named durable subscriptions with independent per-consumer position
  3. Dead-letter support (S5)
  4. Cross-language clients — the eventId design (M3) exists precisely so a non-.NET subscriber can participate
  5. Operable at the deployment’s scale, which is small: dozens of events per minute, not thousands per second

Requirement 5 deserves emphasis. This is a single-operator trading system. A broker that needs a cluster to stay healthy is the wrong trade. The original’s choice — a queue service plus a table in the database already being run — was proportionate, and its replacement should be too.

Subscription store

The original stores subscriptions centrally in SQL Server, on the same host as the trading data (zeus). Two problems: the price database and the bus share a failure domain, and the credentials are inline in configuration (see Security below).

Requirement: the bus’s durable state is separate from the trading data — a different database at minimum, a different host preferably. A bus outage must not be able to take down price persistence, and vice versa.


Relationship to the State Machines

The bus and the state machine framework (specs 60–63) are one design:

Requirement X1: preserve the generation relationship. Events, their stable identities, their payload shapes, and their publication policy should all come from one model, so the wire contract cannot drift from the code. This is what made the “bus message == state machine event” property possible, and it is the thing most easily lost by hand-writing a new bus.

Requirement X2: the model is the schema registry. Emit it in a language-neutral form (JSON Schema, Protobuf, Avro — whatever the target prefers) so subscribers in other languages bind against the same definitions.


Defects To Fix, Not Port

# Defect Location Fix
1 Guaranteed null dereference on subscribe. var m = default(TEventMessage) on a reference type is null, then m.EventId.ToString() is called. Any use of SubscribeByCallBack throws. MyBus171208.fireSubscriptionADdedNotification Use typeof(TEventMessage) metadata; never default(T) for a reference type
2 Unobserved fire-and-forget publish. No await, no continuation, no error handling. Failures are invisible. generated FSM dispatch, all components P2 + P3
3 Sync-over-async. .Wait() on subscribe/unsubscribe/publish. Deadlock risk under any synchronization context. MyBus171208, BDSBus Async throughout
4 correlationId never populated. Slot exists, always empty. everywhere M2
5 subscriptionId accepted and discarded. Subscriptions cannot be named or resumed. SubscribeWithReceiver S1
6 Credentials inline in configuration (bus SQL host/user/password, alongside the trading DB credentials). DTS.Common.Configuration.Console/App.config Environment or secret store; rotate the existing values
7 Parallel “custom” API with no identity fields. MyBus171208, WDSBus S2 — wrap, don’t parallel
8 Three copies of the bus interfaces (usr/StateMachine, PS.Common.StateTransitionMachineries, ...Std) with divergent content — the usr/ copy lacks eventId entirely. — One definition
9 Two copies of the envelope base (PS.Common.Messages, ...Std). — One definition
10 Dead abstractions shipped. DTSBus.Publish has a commented-out body; SharedServiceBus has three fields and no methods; DTSApp.PublishEvent constructs an unused object then throws NotImplementedException.   Delete
11 BusZure subscribe throws. Publish-only adapter presented as an implementation. AzureBus/BusZure.cs Complete or remove
12 Thread.Sleep used for bus readiness (300 ms, 700 ms, 2,455 ms at various sites). MyBus171208, hubs, apps Await a real readiness signal
13 No retry, dead-letter, or poison handling configured. configureBus S5, S6
14 Console writes on every publish and receive. MyBus171208 Structured logging at an appropriate level
15 Header assumed non-null on publish (e.Header.TypeName), but the parameterless envelope constructor leaves it null. Generated events always use the parameterized form, so this is latent — until a hand-written event is published. MyBus171208.Publish Make header non-optional

Migration Path

The trading subsystems already contain the wiring point. Suggested order:

  1. Envelope first. One definition, four identity fields, schemaVersion, correlationId populated. Everything else depends on this.
  2. In-memory adapter. Prove the contracts with no infrastructure; convert the existing in-process events (spec 73) to flow through it. No behaviour change, full test coverage.
  3. Broker adapter + outbox. Now events survive process death.
  4. Wire PDSP. SeriesUpdateCompleted (spec 72, step 8) becomes the first real published event. Its consumer is the indicator engine.
  5. Convert the hook chain. The shell hooks of spec 73 become a subscriber adapter — operator scripts keep working, but now with delivery guarantees and correlation. Do not remove the file-drop mechanism; it is genuinely good.
  6. Retire the remaining direct process launches (spec 73’s chain, spec 74’s RequestRefresh).

Each step is independently valuable and independently revertible.


Verification

# Scenario Expected
T1 Publish with no publisher attached No error; in-process path unaffected
T2 Publish with the bus down Producer completes; message queued in the outbox; delivered on recovery (P3)
T3 Producer’s transaction rolls back No message is published (P3)
T4 Subscriber offline during a publish, restarts under the same subscription name Message delivered on restart (S1, R1)
T5 Subscriber offline, restarts under a new name Starts fresh; no backlog
T6 Same message delivered twice Consumer produces an identical result (S4)
T7 Handler throws Retried per policy, then dead-lettered with envelope and reason (S5)
T8 Poison message in a queue Following messages still delivered (S6)
T9 Two instances of one service, shared subscription name Each message handled once across the pair
T10 Two instances, distinct subscription names Each message handled by both
T11 Event emitted while handling another correlationId inherited (M2)
T12 Non-.NET subscriber Binds by eventId and the emitted schema; receives correctly (M3, X2)
T13 Payload field added, old subscriber still deployed Old subscriber continues working (M4)
T14 Transport swapped for the in-memory adapter No caller changes (T3)
T15 Price refresh completes → indicators → chart data → strategy Whole chain shares one correlationId and is reconstructible from the log

Traceability

Concept Original artifact
Publisher/subscriber contracts Common/PS.Common.StateTransitionMachineries.Std/StateForge.StateMachine/CreasTypes/IBusPublisher.cs (canonical); duplicates in PS.Common.StateTransitionMachineries/ and usr/StateMachine/
Receiver types .../CreasTypes/IMessageReceiver.cs
Message envelope Common/PS.Common.Messages/SMEventArgsBase.cs (canonical, has EventId); Common/PS.Common.Messages.Std/ duplicate; usr/StateMachine/SMEventArgsBase.cs older stripped copy
The working bus Common/PS.Common.Services.ServiceBus/MyBus171208.cs (+ MyBus171208Fsm.cs, MyBus171208StateEnum.cs, IMyBus171208Context.cs)
Transport + subscription config MyBus171208.cs → configureBus, initDbConfig, ConnString
Earlier generations BDSBus.cs, BDSBusSimple.cs, BDSBusMsgIndependent.cs, BDSBusFsm.cs
Azure adapter AzureBus/BusZure.cs, BusZure.creas.cs, BusZureFsm.cs
Sync primitive BusSyncBases/BusSyncBase.cs
Bus logging PSLogger.cs
Dependencies (Rebus, Serilog) Common/PS.Common.Services.ServiceBus/packages.config
Domain wrappers GDS/GDS.PrototypeApp1712/GDSBus.cs, WDS/WDS.WritingGuide.WDSLib/Services/WDSBus.cs
Hub + role pattern WDS/WDS.WritingGuide.WDSLib/Services/WDSServiceHub.cs, WDSServiceHubFsm.cs
Working end-to-end example WDS/WDS.WriterGuide.StoreServiceApp/Program.cs (subscriber + publisher roles); WDS/WDS.WritingGuide.Prototype171231.FSWatcherApp/ (the file-watch publisher)
Generated publish dispatch any *Fsm.cs, e.g. IDS/IDS.Services/IDSServiceFsm.cs, SE/SE.Business/SEComponentFsm.cs, PDSP/PDSP.Business/PDSPPriceComponentFsm.cs
Model-assigned event identity *Fsm.cs, EventIdug: comments + SetEventId(eventIdug)
Bus configuration keys Common/DTS.Common.Configuration.Console/App.config (BusServiceSqlHost, BusServiceSqlDatabaseName, BusServiceSqlDatabaseUser, BusServiceSqlDatabaseUserPassword)
Dead abstractions Common/PS.Common.Services.ServiceBus/DTSBus.cs, Common/PS.Common.Framework/SharedServiceBus.cs, Common/PS.Common.Framework/App/DTS.cs → PublishEvent
Experimental origin Experiments/DTSBusServiceExperimentalApp/, Experiments/DTSBusServiceExperimentalApp.Data/