UseBroadcaster
Enables cross-process lifecycle event broadcasting. When trains execute on a remote worker process, their lifecycle events (started, completed, failed, cancelled) are published to a message bus and delivered to hub processes, where GraphQL subscriptions forward them to connected clients.
Without UseBroadcaster(), subscriptions only fire for trains that execute in the same process as the GraphQL API. With it, subscriptions work regardless of which process executes the train.
Signature
public static TBuilder UseBroadcaster<TBuilder>(
this TBuilder builder,
Action<BroadcasterBuilder> configure
)
where TBuilder : TraxEffectBuilderThe generic type parameter TBuilder is inferred by the compiler, so callers just write .UseBroadcaster(...). This preserves the concrete builder type through chaining (e.g., TraxEffectBuilderWithData stays as TraxEffectBuilderWithData).
| Parameter | Type | Required | Description |
|---|---|---|---|
builder | TBuilder | Yes | The effect configuration builder |
configure | Action<BroadcasterBuilder> | Yes | Callback to select a transport (e.g., UseRabbitMq()) |
What It Registers
| Component | Description |
|---|---|
| Broadcast lifecycle hook | Internal lifecycle hook that publishes events to ITrainEventBroadcaster |
| Broadcast change sink | Internal sink that forwards coalesced data-change signals (onDataChanged) to other processes over the same transport |
TrainEventReceiverService | BackgroundService that consumes events from ITrainEventReceiver and dispatches to ITrainEventHandler instances |
The transport-specific ITrainEventBroadcaster and ITrainEventReceiver are registered by the callback (e.g., UseRabbitMq()). The hook and the sink are internal types: UseBroadcaster() is the only way to add them.
Connection Resilience
The TrainEventReceiverService automatically retries if the transport connection fails (e.g., RabbitMQ is unavailable at startup). It uses exponential backoff starting at 5 seconds, capping at 2 minutes. The service will not crash the host. It logs a warning and keeps retrying until the transport becomes available or the host shuts down.
De-duplication
When a train runs locally on the hub (via a run mutation), the GraphQLSubscriptionHook fires directly and notifies subscribers. The same event is also published to the message bus by the broadcast lifecycle hook. UseBroadcaster() gives each host an instance id (a GUID, one per service provider), the broadcast hook and change sink stamp it on every message as InstanceId, and the TrainEventReceiverService skips only messages carrying its own host's id. This prevents double-notification without dropping anything another host published.
Because the id is per host rather than per application, replicas of one app behind a load balancer (which share an entry assembly, and so an Executor) each receive the others' events. A message from a publisher on a version without InstanceId has none, and is always delivered.
The Executor field is kept for display. It is stamped by the broadcasting process (via Assembly.GetEntryAssembly()), not copied from metadata.Executor, because metadata may be pre-created by a different process (e.g., the API pre-creates metadata for queued jobs that execute on a worker).
Abstractions
The broadcaster system is built on four interfaces that allow alternative transport implementations:
// Publishes lifecycle events to a message bus
public interface ITrainEventBroadcaster
{
Task PublishAsync(TrainLifecycleEventMessage message, CancellationToken ct);
}
// Receives lifecycle events from a message bus
public interface ITrainEventReceiver : IAsyncDisposable
{
Task StartAsync(
Func<TrainLifecycleEventMessage, CancellationToken, Task> handler,
CancellationToken ct
);
Task StopAsync(CancellationToken ct);
}
// Handles received events (e.g., forwarding to GraphQL subscriptions)
public interface ITrainEventHandler
{
Task HandleAsync(TrainLifecycleEventMessage message, CancellationToken ct);
}The TrainLifecycleEventMessage is a serializable record containing:
| Field | Type | Description |
|---|---|---|
MetadataId | long | Database metadata row ID (0 on a DataChanged signal) |
ExternalId | string | External identifier for the execution (empty on a DataChanged signal) |
TrainName | string | Canonical train name, the train interface's full name (empty on a DataChanged signal) |
TrainState | string | Current state (serialized as string for transport) |
Timestamp | DateTime | When the event occurred |
FailureJunction | string? | Junction that failed (if applicable) |
FailureReason | string? | Failure message (if applicable) |
EventType | string | See the event types below |
Executor | string? | Assembly name of the process that broadcast the event, for display |
Output | string? | The completed train's output as JSON, with [TraxSensitive] members masked. With SaveTrainParameters it follows the stored copy: null for an output excluded by SaveOutputs = false, ExcludeOutput or ShouldSaveOutputs, and bounded by MaxParameterBytes. Without it, every completed train's output is serialized for the hooks, up to 1 MiB. An output over its ceiling is replaced by a {"_truncated": true, ...} placeholder. null on any event but a completion's Completed and StateChanged |
HostName | string? | Machine name of the host that ran the train |
HostEnvironment | string? | Environment name of the host that ran the train |
ChangeDomain | string? | The changed domain on a DataChanged signal (WorkQueue, DeadLetter, Manifest, ManifestGroup, SchedulerConfig, Execution); null otherwise |
InstanceId | string? | Id of the host that published the message, used for de-duplication (see above); null from a publisher that predates it |
FailureException | string? | Type name of the exception a failed run recorded (Metadata.FailureException), such as TrainException; null otherwise, and from a publisher that predates it |
Output, HostName, HostEnvironment, ChangeDomain, InstanceId and FailureException are optional on the wire, so a message from an older publisher still deserializes.
EventType is one of:
| Event type | Published when |
|---|---|
Started | A run's row is saved as in progress |
Completed | A run completes |
Failed | A run fails |
Cancelled | A run is cancelled |
StateChanged | After each of the four above. Every transition is published twice, once under its own type and once as StateChanged, so a subscriber to the aggregate stream (the GraphQL onTrainStateChanged subscription) is fed on every host |
DataChanged | A coalesced data-change signal, not a train event. Only ChangeDomain, Timestamp, Executor and InstanceId are set. A handler that only cares about trains ignores it (TrainLifecycleEventMessage.DataChangedEventType) |
Transports
RabbitMQ
effects.UseBroadcaster(b => b.UseRabbitMq("amqp://guest:guest@localhost:5672"))| Parameter | Type | Required | Default | Description |
|---|---|---|---|---|
connectionString | string | Yes | N/A | AMQP connection URI |
configure | Action<RabbitMqBroadcasterOptions>? | No | null | Optional callback to customize options |
Options:
| Property | Type | Default | Description |
|---|---|---|---|
ConnectionString | string | N/A | AMQP connection URI |
ExchangeName | string | "trax.lifecycle" | Fanout exchange name |
PrefetchCount | ushort | 64 | How many received events the receiver may hold unacknowledged at once. The broker holds the rest until the handlers acknowledge one, so a slow handler leaves events queued on the broker rather than in the receiving process. 0 (no limit) is refused: UseRabbitMq throws ArgumentException, and a directly constructed receiver throws InvalidOperationException from StartAsync. |
The RabbitMQ transport uses a fanout exchange so all connected hub instances receive every event. Each receiver creates its own exclusive, auto-delete queue.
Publishing never waits on the broker. A lifecycle hook is awaited inside the train, so the broadcaster writes each event to a bounded queue (1024 events) and returns; one background sender publishes the queue in order.
Each publish waits for the broker's publisher confirm, for at most 5 seconds. An event the broker has not confirmed is sent again, so a receiver can occasionally see the same event twice, but an event lost on a connection that died while still reporting open is not counted as sent. An attempt that times out discards the connection as well as the channel. The exchange is declared on every channel the sender opens, so one deleted, or lost with a broker restart, is recreated.
While the broker is unreachable the sender retries the event it holds with a growing delay (1 second, doubling to 30), each connection attempt bounded at 5 seconds, and further events wait in the queue. An event the broker refuses, by closing the channel on it with a channel-level error (403, 404, 405, 406) or by rejecting the publish, is tried 3 times, then dropped and logged as an Error, so a refusal that does not clear, such as an exchange of the same name declared with another type, does not hold up the events behind it.
When the queue is full, non-terminal events give way first. A new Started, StateChanged or DataChanged event is dropped. A new Completed, Failed or Cancelled event takes the place of the oldest queued non-terminal event, and is dropped itself only when every queued event is terminal. So a run whose Started reached subscribers keeps its outcome, although a subscriber can see an outcome without the Started before it. The first drop logs a Warning, followed by a second one giving the count when the queue drains.
A connection the broker closed is disposed before it is replaced. On shutdown the broadcaster waits up to 5 seconds for queued events to be sent, and does not wait at all while it is failing to reach the broker. Delivery is best effort: a consumer that must not miss an event reads it from the store.
Stopping the receiver tolerates a connection the broker has already closed, for example when the broker restarts or another host sharing it shuts down first. When TrainEventReceiverService retries a receiver that failed to start, each new start closes and disposes the connection the previous attempt opened, and a start that fails partway releases what it opened.
effects.UseBroadcaster(b =>
b.UseRabbitMq("amqp://localhost", opts =>
opts.ExchangeName = "my-app.lifecycle"
)
)SignalR
effects.UseBroadcaster(b => b.UseSignalRHub())SignalR is a sink, not a transport. It pushes events to connected browser or JS clients in real time. Compose it alongside a transport like RabbitMQ for cross-process delivery, or use it on its own when producer and consumer run in the same process.
See UseSignalRHub for filtering and projection options, and MapTraxTrainEventHub for the endpoint mapping that browsers connect to.
Example: Distributed Workers
Both the hub (API + scheduler) and worker processes call UseBroadcaster() with the same RabbitMQ connection:
Hub (Program.cs):
builder.Services.AddTrax(trax =>
trax.AddEffects(effects =>
effects
.UsePostgres(connectionString)
.AddJson()
.UseBroadcaster(b => b.UseRabbitMq(rabbitMqConnectionString))
)
.AddMediator(typeof(Program).Assembly)
.AddScheduler(scheduler => scheduler /* ... */)
);
// AddTraxGraphQL() auto-detects the broadcaster and registers
// GraphQLTrainEventHandler to forward remote events to subscriptions
builder.Services.AddTraxGraphQL();Worker (Program.cs):
builder.Services.AddTrax(trax =>
trax.AddEffects(effects =>
effects
.UsePostgres(connectionString)
.AddJson()
.UseBroadcaster(b => b.UseRabbitMq(rabbitMqConnectionString))
)
.AddMediator(typeof(Program).Assembly)
);
builder.Services.AddTraxWorker(opts => { opts.WorkerCount = 4; });GraphQL Integration
When AddTraxGraphQL() detects that ITrainEventReceiver is registered (via UseBroadcaster()), it automatically registers GraphQLTrainEventHandler as an ITrainEventHandler. This handler maps received TrainLifecycleEventMessage records to TrainLifecycleEvent DTOs and sends them to HotChocolate's ITopicEventSender, making them available to all connected WebSocket subscribers.
No additional configuration is needed. Just call UseBroadcaster() in your effects and AddTraxGraphQL() as usual.
Architecture
Worker Process Hub Process
───────────── ────────────
Train.Run() GraphQL Subscription Clients
→ LifecycleHookRunner ↑
→ BroadcastLifecycleHook TrainEventReceiverService
→ ITrainEventBroadcaster → GraphQLTrainEventHandler
→ RabbitMQ Exchange ─────────→ → ITopicEventSender
→ WebSocket deliveryThe database remains the single source of truth for all train data. The broadcaster carries lifecycle events, and a completion (Completed, and the StateChanged after it) carries the train's output (see Output above), so it reaches every consumer of the exchange and every SignalR and GraphQL subscriber the events are forwarded to. Keep an output you do not want broadcast out with ExcludeOutput or [TraxSensitive]. All metadata, logs, manifests, and train state are always persisted to and queried from PostgreSQL.
Implementing a Custom Transport
To implement a transport other than RabbitMQ:
- Implement
ITrainEventBroadcasterandITrainEventReceiver - Create an extension method on
BroadcasterBuilderthat registers both:
public static BroadcasterBuilder UseMyTransport(
this BroadcasterBuilder builder,
string connectionString
)
{
builder.ServiceCollection.AddSingleton<ITrainEventBroadcaster>(
new MyTransportBroadcaster(connectionString)
);
builder.ServiceCollection.AddSingleton<ITrainEventReceiver>(
new MyTransportReceiver(connectionString)
);
return builder;
}Packages
dotnet add package Trax.Effect # Abstractions
dotnet add package Trax.Effect.Broadcaster.RabbitMQ # RabbitMQ transport
dotnet add package Trax.Effect.Broadcaster.SignalR # SignalR sink for browsers