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–R6in spec 73.
What The Message Bus Enables:
Desired Outcomes:
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.
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 matterseventId is assigned by the model, not at runtime. The state machine
definition (spec 60) assigns each event kind a GUID, and the code generator
emits it as a constant that the event’s constructor stamps in:
// PriceLoadingCompletedEventArgs, EventIdug:78b8fd47-3ff3-4881-9419-7bfe84ce89cb
private void initEventArgs() { SetEventId(eventIdug); ... }
So the identity of “price loading completed” is stable across processes,
across languages, and across renames of the class. This is a schema
registry expressed as model metadata, and it is the right answer to the
question “how does a Python subscriber know which event this is” — it does
not need the .NET type name, it needs the eventId.
messageId is a fresh GUID per emission — the deduplication key.
correlationAssemblyId is derived from the emitting assembly’s execution
identity, so every message from one process run shares it. That answers “which
run of which service produced this”.
correlationId is caller-assigned and defaults to empty. It is empty
everywhere in the shipped code. This is the field spec 73 R3 asks to be
populated; the slot already exists and nothing fills it.
M1 Preserve all four identity fields. They are not redundant — each answers
a different question (what kind, which emission, which chain, which process
run).M2 Populate correlationId. A message emitted while handling another
message inherits that message’s correlationId; a message that starts a chain
generates one. This single change makes the whole processing chain traceable.M3 Make eventId the wire-level discriminator, not the type name. Publish it
in the envelope and route on it. This is what allows a non-.NET subscriber to
participate.M4 Add schemaVersion to the envelope. The original has none, so an event’s
payload cannot evolve without breaking every deployed subscriber.M5 timeStamp must be an unambiguous instant (spec 70, Time semantics).
The original uses local wall-clock with no zone.Event payloads in the shipped code carry domain objects directly —
ChaosDataBuilderDTO, List<string>, PriceSeries. Two consequences:
Requirement M6: distinguish notification events from document
events.
EUR/USD_H4 updated through
bar EUR-USD_H4__22031512”. Subscribers fetch what they need. Small, cheap,
always safe to publish.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.
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.
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:
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.
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.
Three subscription styles exist. They are not redundant; each fits a different consumer shape.
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.
Interface BusSubscriberCallback:
SubscribeByCallback<TMessage>(callback: MessageReceivedCallback<TMessage>)
For consumers that prefer implementing an interface over wiring an event. Thin sugar over style 1.
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.
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:
S3 At-least-once delivery. Combined with P3’s outbox, this is the
achievable and correct guarantee.S4 Idempotent consumers. Every consumer must tolerate redelivery.
messageId is the dedup key. For the trading chain this is nearly free —
every downstream stage is keyed on (instrument, timeframe, barKey) and
recomputing is harmless (spec 72’s convergence property).S5 Explicit retry and dead-letter policy. Bounded retries with backoff,
then dead-letter with the original envelope and the failure reason preserved.S6 Poison-message isolation. One undeliverable message must not block
its queue.| 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.
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.
| 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.
| 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.
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]
| # | 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 |
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:
R1, S1)S5)eventId design (M3) exists precisely so a
non-.NET subscriber can participateRequirement 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.
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.
The bus and the state machine framework (specs 60–63) are one design:
eventId comes from the model’s event GUIDMyBus171208Fsm, 1,210 lines)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.
| # | 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 |
The trading subsystems already contain the wiring point. Suggested order:
schemaVersion,
correlationId populated. Everything else depends on this.SeriesUpdateCompleted (spec 72, step 8) becomes the first
real published event. Its consumer is the indicator engine.RequestRefresh).Each step is independently valuable and independently revertible.
| # | 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 |
| 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/ |