A simple IPC proxy for MediatR requests.
MediatR.IPC provides two public interfaces: a client and a server. The client interface implements ISender. The server takes a dependency on ISender, usually resolved via a DI container. Messages from the client to the server are serialized with protobuf-net. The IPC transport can use Named Pipes or Unix Domain Sockets.
All IPC requests need to be registered. This can be done via assembly scanning or explicit registration. Since IPCs have different app domains, the registration will need to be done on the client and server. I recommend using a shared assembly which does this for you.
[ProtoContract(ImplicitFields = ImplicitFields.AllPublic)]
public record MyFancyCommand : IRequest<bool>
{
public string Message { get; init; } = string.Empty;
}
public static async Task Main(string[] args)
{
// Register the IPC requests on startup
IPCMediator.UseTransport(IPCTransport.NamedPipe);
IPCMediator.RegisterAssemblyTypes(Assembly.GetExecutingAssembly())
.WithAttribute<IPCRequestAttribute>()
.Where(...);
IPCMediator.RegisterType<MyFancyCommand>();
// Use the default runtime type model from protobuf-net
IPCMediator.TypeModel = RuntimeTypeModel.Default;
// Resolve the ISender, so the server can proxy incoming requests.
ISender sender = MyContainer.Resolve<ISender>();
// Run the server until the application is closed.
var server = new MediatorServerPool(sender, poolName: "MyRequestPool", poolSize: 8);
await server.Run();
}public static async Task Main(string[] args)
{
// Register the requests, just like in Process 1
IPCMediator.UseTransport(IPCTransport.NamedPipe);
...
IPCMediator.RegisterType<MyFancyCommand>();
// Create an ISender, which can be consumed by the application.
ISender ipcSender = new MediatorClientPool(poolName: "MyRequestPool", poolSize: 8);
// All requests are sent via ipcSender are sent to and handled by Process 1,
// and the response is sent back to Process 2.
bool result = await ipcSender.Send<MyDto>(new MyFancyCommand { Message = "Hello!" });
Console.WriteLine(result);
}MediatR.IPC always propagates the current W3C System.Diagnostics.Activity with requests and notifications. The receiver creates a new Activity around deserialization and handler dispatch, even when no ActivityListener or tracing exporter is installed. This makes ambient correlation available to handlers through Activity.Current.
For a propagated call:
- The receiver
TraceIdis the same as the senderTraceId. - The receiver
ParentSpanIdis the senderSpanId. - Activities started by the handler become children of the receiver Activity.
- W3C baggage is available through
Activity.Current.GetBaggageItem(...).
When the sender has no current W3C Activity, the receiver starts a new W3C root. Legacy hierarchical Activity IDs are not converted. Propagation metadata travels with the request only; responses do not return the receiver SpanId.
Microsoft.Extensions.Logging can add the ambient IDs and baggage to structured logging scopes:
builder.Logging.Configure(options =>
{
options.ActivityTrackingOptions =
ActivityTrackingOptions.TraceId |
ActivityTrackingOptions.SpanId |
ActivityTrackingOptions.ParentId |
ActivityTrackingOptions.Baggage;
});The selected logging provider must support scopes. Do not add trace or span IDs as metric dimensions because their unbounded cardinality creates a separate time series per operation. Metrics systems that support exemplars can use the ambient Activity for trace-to-metric correlation instead.
An unparseable traceparent is rejected before handler dispatch, because honouring it would attach the operation to a parent that never existed: the request fails with an IPCException, and an invalid one-way notification is dropped because notifications have no acknowledgement channel. tracestate and baggage are carried opaquely, so an oversized value is discarded rather than failing the message. tracestate is limited to 32 members and 512 bytes; baggage is limited to 64 members and 8192 bytes. Baggage may contain sensitive application data and should not cross trust boundaries without review.
All requests sent to the client need to be registered on the client and the server. This allows for quick resolution of types, and little overhead during runtime.
There are two supported transport types: Named Pipes and Unix Domain Sockets. Transport is specified with
IPCMediator.UseTransport(IPCTransport.UnixDomainSocket)
.WithOptions(new UnixDomainSocketOptions
{
SocketPrefix = "/tmp/",
SocketSuffix = ".sock",
});You could also implement your own transport; a TCP transport for instance.
However, this can easily be expanded upon! Simply create an implementation of IStreamStratergy.
Notifications are a work in progress. NotificationHandlers need to be registered in the DI container. IPC notifications would also need to be routed to the MediatorServer. Pull requests are most welcome!
Currently, exceptions thrown by request handlers are not serialized, and no type information is preserved. Exceptions can be huge, and there is no garantue that the client process has a reference to the Exception-type thrown. Ideas and suggestions are welcome!
All forms of contribution are welcome! Here is a list of some much needed features.
- Request cancellation with
CancellationToken - Dynamic buffers for requests in
MediatorServerBase -
IPublisherimplementation - Routing of
INotificationtoINotificationHandlerdesignated for IPC via DI container.