TrainExecution
ITrainExecutionService provides programmatic train execution by name. Instead of resolving a specific train interface, you pass the train's service type name and a JSON string. The service handles discovery, deserialization, and dispatch. Train lookup matches by fully-qualified canonical name first (ServiceType.FullName), then by friendly name (ServiceTypeName). There is no short-name fallback.
It supports two execution paths:
- Queue: creates a WorkQueue entry for asynchronous dispatch by the scheduler.
- Run: executes the train synchronously through the registered
IRunExecutor: on the current machine viaITrainBusby default, or on a remote endpoint underUseRemoteRunorUseLambdaRun.
A third method, PrepareAsync, does only the steps both paths start with (lookup, authorization, input reading), for a surface that submits the work some other way.
Registered automatically by AddMediator() as a scoped service.
ITrainExecutionService
public interface ITrainExecutionService
{
Task<QueueTrainResult> QueueAsync(
string trainName,
string? inputJson,
int priority = 0,
DateTime? scheduledAt = null,
CancellationToken ct = default
);
Task<RunTrainResult> RunAsync(
string trainName,
string inputJson,
CancellationToken ct = default
);
Task<PreparedTrain> PrepareAsync(
string trainName,
string? inputJson,
CancellationToken ct = default
);
}QueueAsync
Creates a WorkQueue entry for asynchronous execution. The scheduler picks up the entry on its next polling cycle and dispatches the train.
scheduledAtwas added beforect, and a JSONnullinput now throwsJsonException. A call passing the token positionally afterpriorityno longer compiles; name it (ct: ct). See Enqueue and Outcome Changes.
Task<QueueTrainResult> QueueAsync(
string trainName,
string? inputJson,
int priority = 0,
DateTime? scheduledAt = null,
CancellationToken ct = default
)| Parameter | Type | Required | Default | Description |
|---|---|---|---|---|
trainName | string | Yes | N/A | Train name, matched by canonical name (ServiceType.FullName), then friendly name (ServiceTypeName). Prefer the fully-qualified interface name (e.g. "MyApp.Trains.IProcessOrderTrain"). |
inputJson | string? | Yes | N/A | JSON-serialized input matching the train's InputType. Null or blank is read as an empty object, {}, which is refused when the input type needs values (see Throws). |
priority | int | No | 0 | Dispatch priority (0-31, higher runs first) |
scheduledAt | DateTime? | No | null | Earliest time the entry may be dispatched, stored as UTC. A Local value is converted; an Unspecified one is taken to already be UTC, which is how a timestamp without an offset arrives from JSON. Null dispatches as soon as a worker is free. |
ct | CancellationToken | No | default | Cancellation token |
Returns: QueueTrainResult
| Property | Type | Description |
|---|---|---|
WorkQueueId | long | Database ID of the created WorkQueue entry |
ExternalId | string | External ID assigned to the entry |
Throws:
TrainNotFoundException(anInvalidOperationException) if no train is registered with the given name. UseITrainDiscoveryService.DiscoverTrains()to list available trains.AmbiguousTrainNameExceptionif the name matches more than one train's friendly name.TrainInputValidationExceptionifinputJsonexceeds the configured size cap (WithMaxInputJsonBytes, 256 KiB by default), or if the input as it would be stored (step 4) is larger thanTrainInputReader.StoredInputGrowthFactor(4) times that cap. For the second,MaxBytesis the stored cap andObservedByteshow much of the stored form had been written when writing stopped, which is more than the cap but not the stored form's full size.JsonExceptionifinputJsondoes not deserialize to the train's input type, or names a property twice ({"amount":1,"Amount":999}). That includes a null or blankinputJsonfor an input type that needs values: a constructor parameter with no default (a positional record such asrecord RenamePlayer(string Id, string NewName)) or arequiredmember. It fails here, at enqueue, rather than queueing a run whose input is full of nulls.Unit, an input with only settable properties, and parameters with defaults are built from{}as before, and an explicit"{}"is read like any other input.JsonExceptionifinputJsonis the JSON literalnull, which is well-formed but is not an input, or is blank and the input type needs values, as forQueueAsync.TrainAuthorizationException(from Trax.Api,Trax.Api.Exceptions, thrown by itsITrainAuthorizationService) if the train has[TraxAuthorize]requirements the caller does not meet. Authorization runs before the input is read, and applies to every caller-built enqueue, including the operations surface (queueTrain,requeueExecution). The dashboard's queue dialog and re-queue button also route through this method, but inside a trusted execution scope, so per-train requirements are skipped there; the dashboard is gated by its host.Trax.Docs/adr/0017records why.TrainAuthorizationNotConfiguredException(anInvalidOperationException, carryingTrainName) if the train declares[TraxAuthorize]and noITrainAuthorizationServiceis registered, unless the call runs inside a trusted execution scope. It is a host misconfiguration rather than a refusal of the caller's input, and the type lets a caller report it that way without reading the message. In a hosted app this rarely fires, because the mediator's startup authorization check (an internal hosted service) already refuses to start such a host; the runtime check covers hosts where hosted services do not run, such as the Lambda runner. The check fails closed; a host that serves no API submissions opts out withAddMediator(m => m.AllowMissingAuthorizationService()), after which the missing service is a no-op.- Any exception thrown by the train's
QueueSubjectKeyoverride. A key that cannot be computed aborts the enqueue rather than becoming null. InvalidOperationExceptionifQueueSubjectKeyreturns an empty string, a key that is only whitespace, a key containing a NUL character or an unpaired surrogate, or a key longer than 512 Unicode characters (an emoji counts once, although it is two UTF-16 units; from Trax.Mediator 1.23.0, which leaves these rules toWorkQueue.Createand wraps itsArgumentExceptionas the inner exception). Return null for an entry that should not be serialized.- Any exception thrown by the train's
OnQueuehook, if the train overrides it. A throw aborts the enqueue and leaves no entry behind: on the default path the hook runs before the entry is committed, and for a train that defers promotion the already-staged entry is removed (only if it is still staged, never once promoted or dispatched), whether or not the caller has cancelled. A failure to remove it does not replace the hook's exception; the stale staged entry sweep resolves an entry left behind. QueueHookTimeoutException(anInvalidOperationException, carryingTrainNameandLimit) from Trax.Mediator 1.23.0, when a train that does not defer promotion runs itsOnQueuehook longer thanMaxQueueHookDuration(30 seconds by default;AddMediator(m => m.WithMaxQueueHookDuration(TimeSpan))changes it, andTimeout.InfiniteTimeSpanremoves it). The hook's token is cancelled at the limit and the enqueue stops waiting whether or not the hook stops: it rolls back, so no entry is written and nothing the hook wrote onIEnqueueContextAccessor.Currentis kept, even a write the hook had already saved, and the connection goes back to the pool. A hook that ignores its token keeps running, but an enqueue it starts after that is refused rather than committed on its own.Trax.Mediator/docs/adr/0004records the reasoning.QueuedWorkCancelledException(anInvalidOperationException, carryingWorkQueueIdandTrainName) for a train that defers promotion, when its staged entry was cancelled while the hook ran (by an operator, or by the stale staged entry sweep because the hook outlivedStaleStagedEntryTimeout). If the sweep promoted it instead (PromoteStaleStagedEntries()), the entry will run andQueueAsyncsucceeds. When it throws, the work will not run, but the hook's side-effect may already have been applied, and the message says so.
QueueAsync with QueueTrainOptions
Task<QueueTrainResult> QueueAsync(
string trainName,
string? inputJson,
QueueTrainOptions options,
CancellationToken ct = default
)Queues exactly as the overload above does, with the options it has no parameter for. It reads the input, authorizes, and throws the same exceptions for the same reasons.
QueueTrainOptions property | Type | Default | Description |
|---|---|---|---|
Priority | int | 0 | Dispatch priority (0-31, higher runs first) |
ScheduledAt | DateTime? | null | Earliest dispatch time, read as scheduledAt above |
ReplayDecisionsOf | long? | null | The metadata id of an earlier run whose recorded decisions the new run replays, so it takes the tracks that run took instead of asking its deciders again. IOperationsService.RequeueExecutionAsync sets it, to the run being re-queued, when that run has decisions to replay. |
An ITrainExecutionService written before this overload existed, a custom one or a decorator
around the mediator's, gets a default implementation that queues through the overload above. When
ReplayDecisionsOf is set it throws DecisionReplayNotSupportedException (in
Trax.Mediator.Exceptions, deriving from NotSupportedException, carrying the implementation's
type as ImplementationType) rather than queue a run that would ask afresh. It is a host
misconfiguration, not a refusal of the caller: implement this overload and pass
ReplayDecisionsOf through.
What it does
- Looks up the train by
trainNameviaITrainDiscoveryService. - Authorizes the caller against the train's requirements, failing closed as described under Throws.
- Deserializes
inputJsonto the train'sInputTypethroughTrainInputReader.Read, reading null or blank as{}, soOnQueueandQueueSubjectKeyalways receive a real input. Property names are matched whatever their case, so{"Amount":5}and{"amount":5}are the same input, and a property given twice, in the same or another casing, is refused withJsonException(Trax.Docs/adr/0023). JSON reference metadata is not honoured: a list written as{"$id":"1","$values":[...]}is refused withJsonException, and$refis read as an unknown property, not as a reference to another part of the input. A run's saved input carries that metadata, so a re-queue resolves it first withTrainInputReader.ResolveSavedInputand hands this method the plain tree. - Re-serializes the input using manifest serialization options (normalizes the JSON: indented, every member written), and refuses the enqueue with
TrainInputValidationExceptionwhen that stored form is larger than 4 timesMaxInputJsonBytes. The bytes are counted as they are written and writing stops the moment the cap is crossed, so a stored form far over the cap is never built in full. The check runs before the entry exists, so nothing is written. - Creates a
WorkQueueentry with the train name, serialized input, input type name, priority, andscheduledAtconverted to UTC. - Stamps the entry's subject key from the train's
QueueSubjectKeyoverride, if it has one. An exception fromQueueSubjectKey, or an empty, whitespace-only or over-long key, aborts the enqueue, so no entry is written. - Tracks the entry, then (if the train overrides
OnQueue) enters the enqueue context and invokes the hook with the entry'sExternalIdand the input, then saves and commits, all in one transaction. A throw rolls the whole thing back, so nothing the hook tracked onIEnqueueContextAccessor.Currentsurvives either. Trains that do not override the hook are never resolved here, enter no context, and open no transaction: their enqueue is a single write. - Returns the entry's ID and external ID.
The train that steps 6 and 7 call is resolved once, on first use, in a DI scope the enqueue creates and disposes before it returns, never from the caller's scope. A caller that holds its scope for a long time, such as a Blazor circuit, therefore keeps no train alive between enqueues, and the scoped services a hook takes are fresh for each enqueue.
Tracking the entry before the hook runs does not insert it (Track is change tracking only), so the hook still runs before the row exists, as its contract states.
When the train sets DeferQueuePromotion, the shape changes to three steps instead: the entry is committed unconfirmed and undispatchable, the hook runs outside that transaction, and a second commit promotes it. A throwing hook removes the staged entry if it is still staged, so the observable contract is the same. Once the hook has returned the mutation counts as accepted, so the promotion runs even if the caller cancels. If the entry was cancelled while the hook ran, the promotion finds nothing to confirm and QueueAsync throws QueuedWorkCancelledException instead of reporting success; if something else already confirmed it (the sweep, with promotion opted in), it will run and QueueAsync succeeds. IEnqueueContextAccessor.Current is null inside such a hook, because the entry is already committed and there is no transaction to join. A crash between the two commits leaves the entry unconfirmed, and the scheduler's stale staged entry sweep cancels it (or promotes it, if the host opted in) once it is older than StaleStagedEntryTimeout. The sweep runs in the ManifestManager, so it does not run while the ManifestManager is disabled (SchedulerConfiguration.ManifestManagerEnabled = false, also the dashboard's Server Settings switch); some host sharing the database must run it.
An enqueue started from inside another train's OnQueue hook, while that enqueue's transaction is open, takes a different path: it tracks its entry on the outer enqueue's context, runs its own hook, and flushes the entry inside the outer transaction, so it commits or rolls back with the outer entry and uses no connection of its own. A deferring train on this path is written confirmed rather than staged. If the nested enqueue fails, the outer one fails too, even when the hook catches the exception; if the hook returns while a nested enqueue it started is still running, the outer enqueue throws InvalidOperationException. See OnQueue.
On the in-memory provider, beginning the transaction succeeds but returns one whose commit and rollback do nothing: the provider ignores EF's TransactionIgnoredWarning. The queue row and anything the hook tracked on IEnqueueContextAccessor.Current still land together in one SaveChanges, and because a transaction object exists, an enqueue nested in a hook finds one to join as it would on a relational provider. A provider whose BeginTransaction throws InvalidOperationException or NotSupportedException gets no transaction at all, with the same single SaveChanges.
RunAsync
Executes a train synchronously through the registered IRunExecutor (locally via ITrainBus unless a remote executor is configured). This is a blocking call that returns when the train completes. It creates no work queue entry, so it does not fire OnQueue and does not consult QueueSubjectKey: a synchronous run can overlap queued work for the same subject.
Task<RunTrainResult> RunAsync(
string trainName,
string inputJson,
CancellationToken ct = default
)| Parameter | Type | Required | Default | Description |
|---|---|---|---|---|
trainName | string | Yes | N/A | Train name (matched by canonical name, then friendly name) |
inputJson | string | Yes | N/A | JSON-serialized input matching the train's InputType. Blank is read the way QueueAsync reads a missing input |
ct | CancellationToken | No | default | Cancellation token forwarded to the train's Run |
Returns: RunTrainResult
| Property | Type | Description |
|---|---|---|
MetadataId | long | Database ID of the metadata record for this execution |
Output | object? | The train's typed output. null for Unit trains; the actual output object for trains with a typed TOut parameter. |
Throws:
TrainNotFoundException(anInvalidOperationException) if no train is registered with the given name, orAmbiguousTrainNameExceptionif the name matches more than one train's friendly name.TrainInputValidationExceptionifinputJsonexceeds the configured size cap.JsonExceptionifinputJsonis the JSON literalnull, which is well-formed but is not an input.TrainExceptionif the train itself fails during execution (propagated fromITrainBus).InvalidOperationExceptionfrom the DI container if the train cannot be built, for example because a constructor dependency is not registered. The train is resolved before its metadata is written, so this leaves noPendingrecord behind.NotSupportedExceptionif the host registered its ownITrainBusand it does not implementRunByNameAsync, keeping the interface's default body. It is thrown before any record is written.TrainAuthorizationExceptionif the train has[TraxAuthorize]requirements the caller does not meet.TrainAuthorizationNotConfiguredException(anInvalidOperationException) if the train declares[TraxAuthorize]and noITrainAuthorizationServiceis registered, unless the call runs inside a trusted execution scope or the host calledAllowMissingAuthorizationService(). The same fail-closed rule, and the same trusted-scope exemption, asQueueAsync.
What it does
- Looks up the train by
trainNameviaITrainDiscoveryService. - Authorizes the caller against the train's requirements, failing closed as described under Throws.
- Deserializes
inputJsonto the train'sInputType, reading blank, casing, repeated properties and reference metadata asQueueAsyncdoes. - Resolves the train found in step 1 by its canonical name, in a child DI scope. The train that runs is the one that was authorized, even when another train takes the same input type. A train that cannot be built throws here, before anything is written.
- Creates a
Metadatarecord with a generated external ID and persists it in thePendingstate. - Runs the train as that record, for the train's
OutputType, invoked by reflection. - Returns the metadata ID and the train's output (or
nullforUnittrains). The output is read through a generic method closed overOutputType, so an output type that is not public (aninternalrecord in the consumer's assembly, say) is returned like any other.
Steps 4 to 6 are the default LocalRunExecutor with the default ITrainBus. A host that registers its own ITrainBus gets the record written first, then ITrainBus.RunByNameAsync<TOut>(trainName, input, ct, metadata) called with it; a bus that does not implement RunByNameAsync is refused before the record is written. Anything that throws after the record is written and before the train takes it over (a replaced bus that cannot build the train, a cancellation in that gap) leaves the record Failed with the exception recorded, or Cancelled for an OperationCanceledException, rather than Pending with nothing to move it on. A record the train already took over is left alone, because the train records its own outcome. A remote executor (UseRemoteRun, UseLambdaRun) takes over steps 4 to 7 and does them its own way.
PrepareAsync
Resolves a train by name, authorizes the current caller for it, and reads the caller's input into the train's input type. These are steps 1 to 3 of QueueAsync and RunAsync, which call the same code, so a surface that submits work itself (the scheduler's run operation, say) accepts and refuses exactly what a queue or a run does. Nothing is written.
Task<PreparedTrain> PrepareAsync(
string trainName,
string? inputJson,
CancellationToken ct = default
)
public sealed class PreparedTrain
{
public TrainRegistration Registration { get; }
public object Input { get; } // an instance of Registration.InputType, never null
}Authorization runs before the input is read, so a caller who may not use the train learns nothing about its input from a parse error. The input is read as QueueAsync reads it: null or blank as {}, property names in any case, a repeated property refused, and the MaxInputJsonBytes cap applied.
PreparedTrain has no public constructor: the only way to get one is from PrepareAsync, so code holding one knows the authorization check ran. An implementation of ITrainExecutionService written before this method existed inherits a default that throws NotSupportedException rather than skipping the check.
Throws
The same as RunAsync: TrainNotFoundException, AmbiguousTrainNameException, UnauthorizedAccessException, TrainInputValidationException, JsonException, and InvalidOperationException for a [TraxAuthorize] train on a host with no ITrainAuthorizationService outside a trusted scope.
TrainInputReader
The reading rules above, as a public static class in Trax.Mediator.Services.TrainExecution. QueueAsync, RunAsync and PrepareAsync all read input through it; a host or package that takes train input JSON on a path of its own should call it rather than copy the rules, so it accepts and refuses the same JSON.
public static class TrainInputReader
{
public const int StoredInputGrowthFactor = 4;
public static object Read(
string? inputJson,
TrainRegistration registration,
int maxInputJsonBytes
);
public static string ResolveSavedInput(
string savedInputJson,
TrainRegistration registration,
int maxInputJsonBytes
);
}Read
| Parameter | Type | Required | Default | Description |
|---|---|---|---|---|
inputJson | string? | Yes | N/A | The caller's JSON. Null or blank is read as {} |
registration | TrainRegistration | Yes | N/A | The train whose InputType is read |
maxInputJsonBytes | int | Yes | N/A | The size cap in UTF-8 bytes, normally MediatorConfiguration.MaxInputJsonBytes |
Returns: an instance of registration.InputType, never null.
Throws: TrainInputValidationException if the JSON is larger than maxInputJsonBytes (checked before parsing); JsonException if it cannot be read as the input type, names a property twice, uses $id/$values reference metadata where the type expects a list, is the literal null, or is blank for an input type that needs values.
It does not authorize. Call it after the caller has been authorized for the train, as PrepareAsync does, so a caller who may not use the train learns nothing about its input from a parse error.
StoredInputGrowthFactor is how many times MaxInputJsonBytes a queued input's stored form may be (see step 4 of QueueAsync).
ResolveSavedInput
Turns a run's saved input (the Input column SaveTrainParameters() writes) into JSON that Read reads as the input the run was given. Anything that reads a saved input back as a train's input, as a re-queue does, passes it through here first.
SaveTrainParameters() writes with TraxJsonSerializationOptions.Default, which preserves references: every object carries an $id, a list is written as {"$id":"2","$values":[...]}, and a second occurrence of the same object as {"$ref":"3"}. Read directly, the list is refused and the $ref reads back as an object with every member at its default. ResolveSavedInput rewrites a saved input whose root object starts with $id (the form such a writer always gives it) as the plain tree it stands for: each $id dropped, each $values object written as its array, and each $ref written as a full copy of the value it names, so no two members of the input share an object. Any other input is returned unchanged.
| Parameter | Type | Required | Default | Description |
|---|---|---|---|---|
savedInputJson | string | Yes | N/A | The saved input |
registration | TrainRegistration | Yes | N/A | The train the input was saved for |
maxInputJsonBytes | int | Yes | N/A | The cap Read will hold the result to, normally MediatorConfiguration.MaxInputJsonBytes |
Returns: the input as a plain JSON tree, or savedInputJson unchanged when it carries no reference metadata at its root.
Throws: JsonException if the saved input is not JSON, its reference metadata is malformed or sits where a reference-preserving writer would not have put it, one $id is given to two values, a $ref names an $id that is not there, a value refers to a value that contains it (a cycle has no plain tree), or the resolved input nests more than 64 levels deep; TrainInputValidationException if the plain form is larger than maxInputJsonBytes. Writing stops the moment the result passes the cap, so references that copy one value many times over never make a small saved input into a large one.
Nothing it does depends on the input type, so it reveals nothing about the input and may run before the caller is authorized. Its result still goes through Read, and through authorization, like any other input.
Examples
Queue a train for async dispatch
public class OrderController(ITrainExecutionService execution) : ControllerBase
{
[HttpPost("orders/queue")]
public async Task<IActionResult> QueueOrder(
[FromBody] JsonElement input,
CancellationToken ct)
{
var result = await execution.QueueAsync(
"MyApp.Trains.IProcessOrderTrain",
input.GetRawText(),
priority: 5,
ct: ct);
return Accepted(new { result.WorkQueueId, result.ExternalId });
}
}Run a train synchronously
public class OrderController(ITrainExecutionService execution) : ControllerBase
{
[HttpPost("orders/run")]
public async Task<IActionResult> RunOrder(
[FromBody] JsonElement input,
CancellationToken ct)
{
var result = await execution.RunAsync(
"MyApp.Trains.IProcessOrderTrain",
input.GetRawText(),
ct);
return Ok(new { result.MetadataId });
}
}Discover available trains first
public class TrainController(
ITrainDiscoveryService discovery,
ITrainExecutionService execution
) : ControllerBase
{
[HttpPost("trains/{trainName}/run")]
public async Task<IActionResult> RunByName(
string trainName,
[FromBody] JsonElement input,
CancellationToken ct)
{
// Validate the train exists before attempting execution
var trains = discovery.DiscoverTrains();
var match = trains.FirstOrDefault(t => t.ServiceType.FullName == trainName);
if (match is null)
return NotFound($"No train registered with name '{trainName}'");
var result = await execution.RunAsync(trainName, input.GetRawText(), ct);
return Ok(new { result.MetadataId });
}
}Concurrency Limiting
RunAsync supports per-train, per-principal and global concurrency limits to prevent overloading remote backends. When a limit is reached, additional requests wait in-process until a slot opens. No requests are rejected.
See Concurrency Limiting for configuration details.
Package
dotnet add package Trax.Mediator