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 : TraxEffectBuilder

The 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).

ParameterTypeRequiredDescription
builderTBuilderYesThe effect configuration builder
configureAction<BroadcasterBuilder>YesCallback to select a transport (e.g., UseRabbitMq())

What It Registers

ComponentDescription
Broadcast lifecycle hookInternal lifecycle hook that publishes events to ITrainEventBroadcaster
Broadcast change sinkInternal sink that forwards coalesced data-change signals (onDataChanged) to other processes over the same transport
TrainEventReceiverServiceBackgroundService 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:

FieldTypeDescription
MetadataIdlongDatabase metadata row ID (0 on a DataChanged signal)
ExternalIdstringExternal identifier for the execution (empty on a DataChanged signal)
TrainNamestringCanonical train name, the train interface's full name (empty on a DataChanged signal)
TrainStatestringCurrent state (serialized as string for transport)
TimestampDateTimeWhen the event occurred
FailureJunctionstring?Junction that failed (if applicable)
FailureReasonstring?Failure message (if applicable)
EventTypestringSee the event types below
Executorstring?Assembly name of the process that broadcast the event, for display
Outputstring?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
HostNamestring?Machine name of the host that ran the train
HostEnvironmentstring?Environment name of the host that ran the train
ChangeDomainstring?The changed domain on a DataChanged signal (WorkQueue, DeadLetter, Manifest, ManifestGroup, SchedulerConfig, Execution); null otherwise
InstanceIdstring?Id of the host that published the message, used for de-duplication (see above); null from a publisher that predates it
FailureExceptionstring?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 typePublished when
StartedA run's row is saved as in progress
CompletedA run completes
FailedA run fails
CancelledA run is cancelled
StateChangedAfter 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
DataChangedA 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"))
ParameterTypeRequiredDefaultDescription
connectionStringstringYesN/AAMQP connection URI
configureAction<RabbitMqBroadcasterOptions>?NonullOptional callback to customize options

Options:

PropertyTypeDefaultDescription
ConnectionStringstringN/AAMQP connection URI
ExchangeNamestring"trax.lifecycle"Fanout exchange name
PrefetchCountushort64How 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 delivery

The 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:

  1. Implement ITrainEventBroadcaster and ITrainEventReceiver
  2. Create an extension method on BroadcasterBuilder that 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