Advanced Patterns
Beyond basic publish/subscribe (covered in Getting Started) and Sagas, Pragmatic.Messaging supports point-to-point request/reply, isolated named buses, scheduled (future) delivery, and message auditing.
Request/Reply
Section titled “Request/Reply”Point-to-point request with a typed response (queue semantics — single consumer).
Define a handler
Section titled “Define a handler”[RequestHandler]public sealed partial class GetPriceHandler : IRequestHandler<GetPriceRequest, PriceResponse>{ public Task<PriceResponse> HandleAsync( GetPriceRequest request, MessageContext context, CancellationToken ct) { return Task.FromResult(new PriceResponse(request.ProductId, 29.99m)); }}var response = await bus.RequestAsync<GetPriceRequest, PriceResponse>( new GetPriceRequest(productId));For cross-service calls, consider RemoteBoundary instead (typed HTTP invokers) — see
Composition.
Multi-Bus
Section titled “Multi-Bus”Isolate message flows by routing handlers to named buses (e.g. keep analytics traffic off the business-critical bus).
[MessageHandler][OnBus("analytics")]public sealed partial class PageViewedHandler : IMessageHandler<PageViewed>{ public async Task HandleAsync(PageViewed message, MessageContext context, CancellationToken ct) { // Runs on the "analytics" bus, isolated from business-critical traffic }}msg.AddBus("analytics", bus =>{ // Each named bus can have its own transport // bus.UseRabbitMq(rmq => rmq.ConnectionString = "amqp://analytics-rmq:5672");});Handlers without [OnBus] consume from the default bus. Named buses are fully independent, with their
own transport and routing.
Scheduled Messages (Jobs Bridge)
Section titled “Scheduled Messages (Jobs Bridge)”Pragmatic.Messaging.Jobs enables future message delivery backed by Pragmatic.Jobs:
msg.EnableScheduledMessages(); // requires Pragmatic.Messaging.JobsThis registers IMessageScheduler:
public interface IMessageScheduler{ Task<Guid> ScheduleAsync<T>(T message, TimeSpan delay, CancellationToken ct = default) where T : notnull; Task<Guid> ScheduleAsync<T>(T message, DateTimeOffset scheduledAt, CancellationToken ct = default) where T : notnull; Task CancelAsync(Guid scheduleId, CancellationToken ct = default);}// Deliver in 24 hoursawait scheduler.ScheduleAsync(new CheckoutReminder(reservationId), delay: TimeSpan.FromHours(24));
// Deliver at a specific timeawait scheduler.ScheduleAsync(new CheckoutReminder(reservationId), scheduledAt: checkout.AddHours(-2));Internally it creates a PublishMessageJob executed by the Jobs scheduler.
Auditing
Section titled “Auditing”Record lifecycle events (Published / Handled / Failed / DeadLettered) for every message.
msg.EnableAuditing(a =>{ a.IncludePayload = false; // exclude the message body from audit entries (default: true)});AuditMiddleware (Order -100) wraps every handler and records duration, handler name, correlation ID,
and tenant. Query the trail:
var entries = await auditStore.QueryAsync(new AuditQuery{ CorrelationId = "order-123", FromDate = DateTimeOffset.UtcNow.AddDays(-7)});| Store | Package | Use case |
|---|---|---|
InMemoryAuditStore | Pragmatic.Messaging.Auditing | Development (default) |
EfCoreAuditStore | Pragmatic.Messaging.EFCore | Production |