Skip to content

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.

Point-to-point request with a typed response (queue semantics — single consumer).

[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.

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.

Pragmatic.Messaging.Jobs enables future message delivery backed by Pragmatic.Jobs:

msg.EnableScheduledMessages(); // requires Pragmatic.Messaging.Jobs

This 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 hours
await scheduler.ScheduleAsync(new CheckoutReminder(reservationId), delay: TimeSpan.FromHours(24));
// Deliver at a specific time
await scheduler.ScheduleAsync(new CheckoutReminder(reservationId), scheduledAt: checkout.AddHours(-2));

Internally it creates a PublishMessageJob executed by the Jobs scheduler.

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)
});
StorePackageUse case
InMemoryAuditStorePragmatic.Messaging.AuditingDevelopment (default)
EfCoreAuditStorePragmatic.Messaging.EFCoreProduction