Requester and Notifier¶
Dispatch requests and notifications through handler pipelines with reusable cross-cutting behaviors.
- Requester and Notifier
- Overview
- Challenges
- Solution
- Key Features
- Architecture
- Use Cases
- Basic Usage
- Part 1: Requester
- Basic usage
- Configuring behavior options
- Available options
- Configuration methods
- RetryOptions configuration
- TimeoutOptions configuration
- CircuitBreakerOptions configuration
- ChaosOptions configuration
- DatabaseTransactionOptions configuration (Entity Framework)
- Generic BehaviorOptions configuration
- Fluent API chaining
- Configuration priority
- Notifier options configuration
- Examples
- Part 2: Notifier
- Part 3: pipeline behaviors
- Appendix A: CQS pattern
- Appendix B: comparison with MediatR
- Appendix C: Generic handlers in Requester and Notifier
- Appendix D: Source-generated commands, queries, and events
Overview¶
The Requester and Notifier systems dispatch requests and notifications through configurable
handler pipelines. Requester selects one handler for a command or query. Notifier publishes one
notification to every matching handler. Both use the same behavior contract for concerns such as
validation, retries, and timeouts.
Two kinds of messages are supported:
- Request/Response Messages (via
Requester): Dispatched to a single handler, typically for commands (actions that modify state) or queries (requests that retrieve data). - Notification Messages (via
Notifier): Dispatched to multiple handlers in a publish/subscribe (pub/sub) model, typically for domain events where multiple components need to react to an event (e.g., user registration, system updates).
Using these components (requests for commands/queries and notifications for domain events) helps decrease coupling in the application by isolating message handling logic from the business logic of the application. For instance, a service can publish a domain event via the Notifier without knowing which components will handle it, and a command can be dispatched via the Requester without the caller needing to know the implementation details of the handler.
Some of the cross-cutting behaviors described on this page rely on shared common infrastructure, especially Common Caching for cache-related behaviors and Common Observability / Tracing for broader instrumentation conventions around message processing flows.
Challenges¶
Managing the core mechanics of message dispatching and handler execution presents several challenges:
- Inconsistent Dispatching: Without a standardized mechanism, dispatching requests or notifications to handlers can vary across the application, leading to unpredictable behavior.
- Error Propagation Complexity: Propagating errors from handlers through multiple layers while preserving context is difficult, often resulting in lost information.
- Coupled concerns: Handlers often mix business logic with error handling and logging. This coupling makes handlers harder to maintain, test, and change.
- Extensibility Limitations: Adding new functionality (e.g., validation, retries) typically requires modifying existing handlers, increasing complexity and risk.
- Type Safety Issues: Ensuring type-safe handling of requests, notifications, and their results, especially for operations with no meaningful return value, can be error-prone.
- Message Tracking: Tracking message metadata (e.g., IDs, timestamps) for debugging and auditing is often ad hoc, leading to inconsistent monitoring.
- Multiple Handler Coordination: For notifications, coordinating multiple handlers (e.g., sequential, concurrent, or fire-and-forget execution) adds complexity, especially when ensuring consistent error handling.
Requester and Notifier address these concerns through common dispatch, result, metadata, discovery, and behavior contracts.
Solution¶
The systems provide:
- Standardizing message dispatching through central interfaces (
IRequesterfor requests,INotifierfor notifications), ensuring consistent behavior across the application. - Enabling structured error propagation with
Result<TValue>for request handlers andResultfor notification handlers.INotifier.PublishAsync(...)exposes the aggregate asIResult. - Decoupling concerns by using a pipeline of behaviors to handle technical aspects (e.g., validation, retries) separately from business logic, and reducing coupling between components through the use of requests (for commands/queries) and notifications (for domain events).
- Supporting extensibility through shared behaviors that can be added without modifying handlers.
- Providing type-safe handling with
Result<Unit>for commands with no meaningful return value and non-genericResultfor notifications. - Including built-in message metadata (
RequestId,RequestTimestampfor requests;NotificationId,NotificationTimestampfor notifications) for tracking and auditing. - Supporting multiple handler coordination for notifications with configurable execution modes (sequential, concurrent, fire-and-forget).
Key Features¶
- Type-Safe Handling:
Requester: UsesResult<TResponse>for requests.Notifier: Uses non-genericResultbecause notifications do not return values.- DI Integration: Scoped handler lifetimes for both systems.
- Message Metadata:
RequestBase<TResponse>: ProvidesRequestId,RequestTimestamp.NotificationBase: ProvidesNotificationId,NotificationTimestamp.- Asynchronous dispatch contracts:
RequestHandlerBase<TRequest, TResponse>returnsTask<Result<TResponse>>.NotificationHandlerBase<TNotification>returnsTask<Result>.- Source-generated
[Handle]methods may be synchronous or asynchronous; generated bridges use the asynchronous contracts. - Shared Pipeline Behaviors: For validation, retry, timeout, and custom logic, applicable to both systems.
- Execution Modes (Notifier): Sequential (default), concurrent, or fire-and-forget, configurable per notification.
- Per-Handler Policies: Timeout, retry, and chaos injection policies via attributes.
- Progress Reporting: Via
IProgress<ProgressReport>inSendOptions(requests) andPublishOptions(notifications). - Automatic Discovery: Handler and validator discovery via assembly scanning for both systems.
Architecture¶
IRequester dispatches a request to one handler, while INotifier dispatches a notification to
multiple handlers. Manual messages normally derive from RequestBase<TValue> or NotificationBase;
source-generated messages receive the corresponding base type. Handlers return Result<TValue> or
Result, and shared IPipelineBehavior<TRequest, TResponse> implementations wrap dispatch. The
builders register handlers, behaviors, and discovered validators.
classDiagram
class IRequester {
<<interface>>
+SendAsync(TRequest, SendOptions, CancellationToken)
+GetRegistrationInformation()
}
class INotifier {
<<interface>>
+PublishAsync(TNotification, PublishOptions, CancellationToken)
+GetRegistrationInformation()
}
class IResult {
<<interface>>
+IReadOnlyList Messages
+IReadOnlyList Errors
+bool IsSuccess
+bool IsFailure
}
class IResultGeneric {
<<IResult of T>>
+TValue Value
}
class IResultError {
<<interface>>
+string Message
}
class IRequest {
<<interface>>
+Guid RequestId
+DateTimeOffset RequestTimestamp
}
class INotification {
<<interface>>
+Guid NotificationId
+DateTimeOffset NotificationTimestamp
}
class IPipelineBehavior {
<<interface>>
+HandleAsync(TRequest, object, Type, Func, CancellationToken)
}
class IHandlerCache {
<<interface>>
+TryAdd(Type, Type)
}
class IRequestHandler {
<<interface>>
+HandleAsync(TRequest, SendOptions, CancellationToken)
}
class INotificationHandler {
<<interface>>
+HandleAsync(TNotification, PublishOptions, CancellationToken)
}
class Requester {
+SendAsync(TRequest, SendOptions, CancellationToken)
+GetRegistrationInformation()
}
class Notifier {
+PublishAsync(TNotification, PublishOptions, CancellationToken)
+GetRegistrationInformation()
}
class RequestBase {
+Guid RequestId
+DateTimeOffset RequestTimestamp
}
class NotificationBase {
+Guid NotificationId
+DateTimeOffset NotificationTimestamp
}
class PipelineBehaviorBase {
+HandleAsync(TRequest, object, Type, Func, CancellationToken)
}
class HandlerCache {
+TryAdd(Type, Type)
}
class RequestHandlerBase {
+HandleAsync(TRequest, SendOptions, CancellationToken)
}
class NotificationHandlerBase {
+HandleAsync(TNotification, PublishOptions, CancellationToken)
}
IRequester <|.. Requester
INotifier <|.. Notifier
IResult <|-- IResultGeneric
IResult ..> IResultError : contains
class ResultGeneric {
<<Result of T>>
}
IResultGeneric <|.. ResultGeneric
IRequest <|.. RequestBase
INotification <|.. NotificationBase
IPipelineBehavior <|.. PipelineBehaviorBase
IHandlerCache <|.. HandlerCache
IRequestHandler <|.. RequestHandlerBase
INotificationHandler <|.. NotificationHandlerBase
Requester --> IPipelineBehavior : uses
Requester --> IHandlerCache : uses
Requester --> IRequestHandler : uses
Notifier --> IPipelineBehavior : uses
Notifier --> IHandlerCache : uses
Notifier --> INotificationHandler : uses
RequestHandlerBase --> IResultGeneric : returns
NotificationHandlerBase --> IResult : returns
Use Cases¶
- Requester:
- Creating a customer (
Result<Unit>). - Updating a customer's email (
Result<string>). - Fetching customer details (
Result<CustomerDto>). - Applying retry/timeout policies to critical operations.
- Processing generic entities (e.g.,
SaveEntityRequest<TEntity>). - Notifier:
- Sending email notifications after user registration (handled by multiple components like email sender, logger).
- Logging domain events to multiple destinations (database, file, external service).
- Updating caches in a fire-and-forget manner after data changes.
- Broadcasting system events to multiple subscribers (e.g., audit logging, metrics collection).
Basic Usage¶
Register handler discovery and the behaviors needed by each dispatcher:
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(RetryPipelineBehavior<,>));
services.AddNotifier()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(RetryPipelineBehavior<,>));
Dispatch a request and publish a notification, checking each result before using it:
var requestResult = await requester.SendAsync(
new DoSomethingCommand { Message = "Index customer records" },
cancellationToken: cancellationToken);
if (requestResult.IsFailure)
{
Console.Error.WriteLine(string.Join(
Environment.NewLine,
requestResult.Errors.Select(error => error.Message)));
return;
}
var notificationResult = await notifier.PublishAsync(
new UserRegisteredNotification
{
UserId = Guid.Parse("5f6b5ba2-85d5-44bb-87a7-f876a65cdb09"),
Email = "ada@example.test"
},
cancellationToken: cancellationToken);
if (notificationResult.IsFailure)
{
Console.Error.WriteLine(string.Join(
Environment.NewLine,
notificationResult.Errors.Select(error => error.Message)));
return;
}
Console.WriteLine("Request and notification completed.");
The success path ends with:
Request and notification completed.
Part 1: Requester¶
Basic usage¶
Request dispatching¶
// Creating a command
var command = new DoSomethingCommand { Message = "Performing task..." };
// Dispatching the command
var requester = provider.GetRequiredService<IRequester>();
var result = await requester.SendAsync(command);
// Checking result status
if (result.IsSuccess)
{
Console.WriteLine("Task performed successfully.");
}
else
{
Console.WriteLine($"Failed: {result.Errors.FirstOrDefault()?.Message}");
}
// Dispatching a query with a value
var query = new GetUserQuery { UserId = Guid.NewGuid() };
var queryResult = await requester.SendAsync(query);
if (queryResult.IsSuccess)
{
Console.WriteLine($"User found: {queryResult.Value.Username}");
}
Command with unit result¶
public class DoSomethingCommand : RequestBase<Unit>
{
public string Message { get; set; }
// Nested FluentValidation validator for automatic discovery
public class Validator : AbstractValidator<DoSomethingCommand>
{
public Validator()
{
RuleFor(x => x.Message).NotEmpty().WithMessage("Message cannot be empty.");
}
}
}
[HandlerRetry(3, 200)] // Retry 3 times with 200ms delay
public class DoSomethingCommandHandler : RequestHandlerBase<DoSomethingCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(DoSomethingCommand request, SendOptions options, CancellationToken cancellationToken)
{
await Task.Delay(100, cancellationToken); // Simulate async operation
return Result<Unit>.Success(Unit.Value);
}
}
Query with value result¶
public class User : IEntity
{
public Guid Id { get; set; }
public string Username { get; set; }
public bool HasIdentity() => this.Id != Guid.Empty;
}
public class GetUserQuery : RequestBase<User>
{
public Guid UserId { get; set; }
// Nested FluentValidation validator for automatic discovery
public class Validator : AbstractValidator<GetUserQuery>
{
public Validator()
{
RuleFor(x => x.UserId).NotEmpty().WithMessage("UserId cannot be empty.");
}
}
}
[HandlerTimeout(5000)] // Timeout after 5 seconds
public class GetUserQueryHandler : RequestHandlerBase<GetUserQuery, User>
{
private readonly IGenericReadOnlyRepository<User> userRepository;
public GetUserQueryHandler(IGenericReadOnlyRepository<User> userRepository)
{
this.userRepository = userRepository;
}
protected override async Task<Result<User>> HandleAsync(GetUserQuery request, SendOptions options, CancellationToken cancellationToken)
{
var user = await this.userRepository.FindOneAsync(request.UserId, cancellationToken: cancellationToken);
return user != null
? Result<User>.Success(user)
: Result<User>.Failure().WithMessage($"User with ID {request.UserId} not found.");
}
}
FluentValidation setup in requests¶
Requests can include a nested Validator class that extends FluentValidation.AbstractValidator<TRequest>. The RequesterBuilder automatically discovers these validators during assembly scanning and registers them for use with the ValidationPipelineBehavior.
public class CreateCustomerCommand : RequestBase<string>
{
public string Email { get; set; }
public class Validator : AbstractValidator<CreateCustomerCommand>
{
public Validator()
{
RuleFor(x => x.Email)
.NotEmpty().WithMessage("Email cannot be empty.")
.EmailAddress().WithMessage("Invalid email format.");
}
}
}
Validation behavior¶
The ValidationPipelineBehavior is a pipeline behavior that automatically validates requests using FluentValidation if a validator is registered. It runs before the handler, ensuring validation errors are caught early.
// Ensure the ValidationPipelineBehavior is registered
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>));
// Example request with validation
public class UpdateEmailCommand : RequestBase<string>
{
public string Email { get; set; }
public class Validator : AbstractValidator<UpdateEmailCommand>
{
public Validator()
{
RuleFor(x => x.Email)
.NotEmpty().WithMessage("Email cannot be empty.")
.EmailAddress().WithMessage("Invalid email format.");
}
}
}
public class UpdateEmailCommandHandler : RequestHandlerBase<UpdateEmailCommand, string>
{
protected override async Task<Result<string>> HandleAsync(UpdateEmailCommand request, SendOptions options, CancellationToken cancellationToken)
{
// ValidationPipelineBehavior ensures Email is valid before this executes
return Result<string>.Success(request.Email);
}
}
Using SendOptions¶
var requester = provider.GetRequiredService<IRequester>();
var query = new GetUserQuery { UserId = Guid.NewGuid() };
// Configure SendOptions to throw exceptions
var options = new SendOptions
{
HandleExceptionsAsResultError = false
};
try
{
var result = await requester.SendAsync(query, options);
if (result.IsSuccess)
{
Console.WriteLine($"User found: {result.Value.Username}");
}
}
catch (Exception ex)
{
Console.WriteLine($"Error: {ex.Message}");
}
Request metadata and cancellation¶
public class CancelableCommand : RequestBase<Unit>
{
public Guid TaskId { get; set; }
}
public class CancelableCommandHandler : RequestHandlerBase<CancelableCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(CancelableCommand request, SendOptions options, CancellationToken cancellationToken)
{
Console.WriteLine($"Request ID: {request.RequestId}, Timestamp: {request.RequestTimestamp}");
await Task.Delay(5000, cancellationToken); // Long-running operation
return Result<Unit>.Success(Unit.Value);
}
}
var cts = new CancellationTokenSource(TimeSpan.FromSeconds(2));
var result = await requester.SendAsync(new CancelableCommand { TaskId = Guid.NewGuid() }, cancellationToken: cts.Token);
Practices¶
- Early Returns: Check
result.IsSuccessearly to avoid unnecessary processing.
var result = await requester.SendAsync(request);
if (result.IsFailure)
{
return result;
}
- Meaningful Messages: Include context in error messages for better debugging.
return Result<Unit>.Failure($"Failed to process task {request.TaskId}: Invalid input");
- Order Behaviors: Place critical behaviors (e.g., validation, transactions) early in the pipeline.
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(DatabaseTransactionPipelineBehavior<,>))
.WithBehavior(typeof(RetryPipelineBehavior<,>));
Configuring behavior options¶
The Requester and Notifier systems support configuring default options for various pipeline behaviors through a fluent API. These options can be configured globally to provide default values that are used when handler-specific attributes don't specify them. This allows for centralized configuration of cross-cutting concerns like retries, timeouts, and circuit breakers.
Available options¶
The following options can be configured:
- RetryOptions: Configure default retry behavior (count and delay).
- TimeoutOptions: Configure default timeout duration.
- CircuitBreakerOptions: Configure default circuit breaker behavior (attempts, break duration, backoff).
- ChaosOptions: Configure default chaos injection behavior (injection rate, enabled state).
- DatabaseTransactionOptions (EntityFramework): Configure default database transaction behavior (context name).
Configuration methods¶
Each option type has two configuration methods:
- Parameter-based: Pass specific values directly.
- Action-based: Pass a configuration action for more flexibility.
The generic WithBehaviorOptions<TOptions> method configures an option type through an action delegate.
RetryOptions configuration¶
Configure default retry behavior for handlers with the HandlerRetryAttribute:
RetryPipelineBehavior retries thrown exceptions. A handler that returns a failed Result is not
retried.
// Using parameters
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(RetryPipelineBehavior<,>))
.WithRetryOptions(defaultCount: 3, defaultDelay: 100);
// Using action
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(RetryPipelineBehavior<,>))
.WithRetryOptions(options =>
{
options.DefaultCount = 3;
options.DefaultDelay = 100;
});
// Handler usage: attribute values take precedence, otherwise defaults are used
[HandlerRetry] // Uses defaults from RetryOptions
public class MyCommandHandler : RequestHandlerBase<MyCommand, Unit>
{
protected override Task<Result<Unit>> HandleAsync(MyCommand request, SendOptions options, CancellationToken cancellationToken)
{
// Exceptions from this handler are retried; failed Result values are returned as-is
return Task.FromResult(Result<Unit>.Success(Unit.Value));
}
}
TimeoutOptions configuration¶
Configure default timeout duration for handlers with the HandlerTimeoutAttribute:
The pessimistic timeout bounds how long the pipeline awaits the handler, but it does not replace or cancel the handler's original token. The underlying operation can continue after a timeout unless it has another cancellation mechanism.
// Using parameter
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(TimeoutPipelineBehavior<,>))
.WithTimeoutOptions(defaultDuration: 5000); // 5 seconds
// Using action
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(TimeoutPipelineBehavior<,>))
.WithTimeoutOptions(options => options.DefaultDuration = 5000);
// Handler usage
[HandlerTimeout] // Uses default from TimeoutOptions
public class MyQueryHandler : RequestHandlerBase<MyQuery, User>
{
protected override async Task<Result<User>> HandleAsync(MyQuery request, SendOptions options, CancellationToken cancellationToken)
{
// Handler logic with automatic timeout
return Result<User>.Success(new User());
}
}
CircuitBreakerOptions configuration¶
Configure default circuit breaker behavior for handlers with the HandlerCircuitBreakerAttribute:
// Using parameters
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(CircuitBreakerPipelineBehavior<,>))
.WithCircuitBreakerOptions(
defaultAttempts: 5,
defaultBreakDurationSeconds: 60,
defaultBackoffMilliseconds: 1000,
defaultBackoffExponential: true);
// Using action
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(CircuitBreakerPipelineBehavior<,>))
.WithCircuitBreakerOptions(options =>
{
options.DefaultAttempts = 5;
options.DefaultBreakDurationSeconds = 60;
options.DefaultBackoffMilliseconds = 1000;
options.DefaultBackoffExponential = true;
});
// Handler usage
[HandlerCircuitBreaker] // Uses defaults from CircuitBreakerOptions
public class ExternalApiCommandHandler : RequestHandlerBase<ExternalApiCommand, Unit>
{
protected override Task<Result<Unit>> HandleAsync(ExternalApiCommand request, SendOptions options, CancellationToken cancellationToken)
{
// Handler logic with circuit breaker protection
return Task.FromResult(Result<Unit>.Success(Unit.Value));
}
}
ChaosOptions configuration¶
Configure default chaos injection behavior for handlers with the HandlerChaosAttribute:
// Using parameters
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ChaosPipelineBehavior<,>))
.WithChaosOptions(defaultInjectionRate: 0.1, defaultEnabled: false); // 10% injection rate, disabled by default
// Using action
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ChaosPipelineBehavior<,>))
.WithChaosOptions(options =>
{
options.DefaultInjectionRate = 0.1;
options.DefaultEnabled = false; // Disable chaos injection by default
});
// Handler usage
[HandlerChaos] // Uses defaults from ChaosOptions
public class TestCommandHandler : RequestHandlerBase<TestCommand, Unit>
{
protected override Task<Result<Unit>> HandleAsync(TestCommand request, SendOptions options, CancellationToken cancellationToken)
{
// Handler logic with chaos injection for testing resilience
return Task.FromResult(Result<Unit>.Success(Unit.Value));
}
}
DatabaseTransactionOptions configuration (Entity Framework)¶
Configure default database transaction behavior for handlers with the HandlerDatabaseTransactionAttribute:
// Using parameter
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(DatabaseTransactionPipelineBehavior<,>))
.WithDatabaseTransactionOptions(defaultContextName: "Core"); // Default to "CoreDbContext"
// Using action
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(DatabaseTransactionPipelineBehavior<,>))
.WithDatabaseTransactionOptions(options => options.DefaultContextName = "Core");
// Handler usage
[HandlerDatabaseTransaction] // Uses default "Core" from DatabaseTransactionOptions
public class UpdateUserCommandHandler : RequestHandlerBase<UpdateUserCommand, Unit>
{
private readonly IGenericRepository<User> userRepository;
public UpdateUserCommandHandler(IGenericRepository<User> userRepository)
{
this.userRepository = userRepository;
}
protected override async Task<Result<Unit>> HandleAsync(UpdateUserCommand request, SendOptions options, CancellationToken cancellationToken)
{
// Handler logic wrapped in database transaction
var user = await this.userRepository.FindOneAsync(request.UserId, cancellationToken: cancellationToken);
user.Username = request.Username;
await this.userRepository.UpdateAsync(user, cancellationToken);
return Result<Unit>.Success(Unit.Value);
}
}
Generic BehaviorOptions configuration¶
Use the generic WithBehaviorOptions<TOptions> method to configure any option type:
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(RetryPipelineBehavior<,>))
.WithBehaviorOptions<RetryOptions>(options =>
{
options.DefaultCount = 3;
options.DefaultDelay = 100;
})
.WithBehavior(typeof(TimeoutPipelineBehavior<,>))
.WithBehaviorOptions<TimeoutOptions>(options => options.DefaultDuration = 5000);
Fluent API chaining¶
All option configuration methods return the builder instance, enabling fluent API chaining:
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(RetryPipelineBehavior<,>))
.WithRetryOptions(3, 100)
.WithBehavior(typeof(TimeoutPipelineBehavior<,>))
.WithTimeoutOptions(5000)
.WithBehavior(typeof(CircuitBreakerPipelineBehavior<,>))
.WithCircuitBreakerOptions(5, 60, 1000, true)
.WithBehavior(typeof(ChaosPipelineBehavior<,>))
.WithChaosOptions(0.1, false)
.WithBehavior(typeof(DatabaseTransactionPipelineBehavior<,>))
.WithDatabaseTransactionOptions("Core");
Configuration priority¶
When both default options and handler attributes are present, the handler attribute values take precedence:
// Configure defaults
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(RetryPipelineBehavior<,>))
.WithRetryOptions(3, 100); // Default: 3 retries with 100ms delay
// Handler with specific values
[HandlerRetry(5, 200)] // Overrides defaults: 5 retries with 200ms delay
public class ImportantCommandHandler : RequestHandlerBase<ImportantCommand, Unit>
{
protected override Task<Result<Unit>> HandleAsync(ImportantCommand request, SendOptions options, CancellationToken cancellationToken)
{
return Task.FromResult(Result<Unit>.Success(Unit.Value));
}
}
// Handler using defaults
[HandlerRetry] // Uses defaults: 3 retries with 100ms delay
public class StandardCommandHandler : RequestHandlerBase<StandardCommand, Unit>
{
protected override Task<Result<Unit>> HandleAsync(StandardCommand request, SendOptions options, CancellationToken cancellationToken)
{
return Task.FromResult(Result<Unit>.Success(Unit.Value));
}
}
Notifier options configuration¶
The same options configuration methods are available for the Notifier system:
services.AddNotifier()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(RetryPipelineBehavior<,>))
.WithRetryOptions(3, 100)
.WithBehavior(typeof(TimeoutPipelineBehavior<,>))
.WithTimeoutOptions(5000);
Examples¶
Command handling example¶
public class DoSomethingCommand : RequestBase<Unit>
{
public string Message { get; set; }
public class Validator : AbstractValidator<DoSomethingCommand>
{
public Validator()
{
RuleFor(x => x.Message).NotEmpty().WithMessage("Message cannot be empty.");
}
}
}
[HandlerRetry(3, 200)] // Retry 3 times with 200ms delay
public class DoSomethingCommandHandler : RequestHandlerBase<DoSomethingCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(DoSomethingCommand request, SendOptions options, CancellationToken cancellationToken)
{
await Task.Delay(100, cancellationToken); // Simulate async operation
return Result<Unit>.Success(Unit.Value);
}
}
var command = new DoSomethingCommand { Message = "Performing task..." };
var result = await requester.SendAsync(command);
if (result.IsSuccess)
{
Console.WriteLine("Task performed successfully.");
}
else
{
Console.WriteLine($"Failed: {result.Errors.FirstOrDefault()?.Message}");
}
Query handling example¶
public class User : IEntity
{
public Guid Id { get; set; }
public string Username { get; set; }
public bool HasIdentity() => this.Id != Guid.Empty;
}
public class GetUserQuery : RequestBase<User>
{
public Guid UserId { get; set; }
public class Validator : AbstractValidator<GetUserQuery>
{
public Validator()
{
RuleFor(x => x.UserId).NotEmpty().WithMessage("UserId cannot be empty.");
}
}
}
[HandlerTimeout(5000)] // Timeout after 5 seconds
public class GetUserQueryHandler : RequestHandlerBase<GetUserQuery, User>
{
private readonly IGenericReadOnlyRepository<User> userRepository;
public GetUserQueryHandler(IGenericReadOnlyRepository<User> userRepository)
{
this.userRepository = userRepository;
}
protected override async Task<Result<User>> HandleAsync(GetUserQuery request, SendOptions options, CancellationToken cancellationToken)
{
var user = await this.userRepository.FindOneAsync(request.UserId, cancellationToken: cancellationToken);
return user != null
? Result<User>.Success(user)
: Result<User>.Failure().WithMessage($"User with ID {request.UserId} not found.");
}
}
var query = new GetUserQuery { UserId = Guid.NewGuid() };
var result = await requester.SendAsync(query);
if (result.IsSuccess)
{
Console.WriteLine($"User found: {result.Value.Username}");
}
else
{
Console.WriteLine($"Failed: {result.Errors.FirstOrDefault()?.Message}");
}
Part 2: Notifier¶
Notifier overview¶
The Notifier system complements the Requester with a publish/subscribe model. Unlike requests,
notifications are dispatched to multiple handlers. The notifier supports sequential, concurrent, and
fire-and-forget execution and returns IResult. Sequential and concurrent modes aggregate completed
handler outcomes; fire-and-forget returns success after scheduling the handlers.
Challenges addressed by Notifier¶
- Multiple Handler Coordination: Coordinating multiple handlers for a single notification, with options for sequential, concurrent, or fire-and-forget execution.
- Consistent Error Handling: Aggregating errors from multiple handlers into a single
Resultwhile preserving context. - Execution Flexibility: Allowing different execution modes to suit various use cases (e.g., ordered processing, parallel execution, non-blocking dispatch).
- Shared Behaviors: Reusing the same pipeline behaviors as the
Requesterfor consistency (e.g., validation, retries).
Solution¶
The Notifier system addresses these challenges by:
- Providing a central
INotifierinterface for publishing notifications to multiple handlers. - Aggregating handler outcomes into an
IResultfor sequential and concurrent publication. - Supporting three execution modes (sequential, concurrent, fire-and-forget) configurable via
PublishOptions. - Reusing the
Requesterpipeline behaviors for the same cross-cutting behavior.
Key features specific to Notifier¶
- Pub/Sub Model: Dispatches notifications to multiple handlers.
- Execution Modes:
- Sequential: Handlers run one after another, stopping on the first failure (default).
- Concurrent: Handlers run in parallel, with results aggregated.
- Fire-and-Forget: Handlers are invoked without awaiting completion, returning immediately.
- Non-generic result: Returns
IResult; completed handler outcomes are represented byResult. - Shared Behaviors: Uses the same pipeline behaviors as
Requester(e.g., validation, retries).
Basic usage¶
Notification dispatching¶
// Creating a notification
var notification = new UserRegisteredNotification { UserId = Guid.NewGuid(), Email = "user@example.com" };
// Dispatching the notification
var notifier = provider.GetRequiredService<INotifier>();
var result = await notifier.PublishAsync(notification);
// Checking result status
if (result.IsSuccess)
{
Console.WriteLine("Notification processed successfully by all handlers.");
}
else
{
Console.WriteLine($"Failed: {result.Errors.FirstOrDefault()?.Message}");
}
Notification with multiple handlers¶
public class UserRegisteredNotification : NotificationBase
{
public Guid UserId { get; set; }
public string Email { get; set; }
// Nested FluentValidation validator for automatic discovery
public class Validator : AbstractValidator<UserRegisteredNotification>
{
public Validator()
{
RuleFor(x => x.UserId).NotEmpty().WithMessage("UserId cannot be empty.");
RuleFor(x => x.Email)
.NotEmpty().WithMessage("Email cannot be empty.")
.EmailAddress().WithMessage("Invalid email format.");
}
}
}
// Handler to send an email
public class SendEmailNotificationHandler : NotificationHandlerBase<UserRegisteredNotification>
{
protected override async Task<Result> HandleAsync(UserRegisteredNotification notification, PublishOptions options, CancellationToken cancellationToken)
{
// Simulate sending an email
await Task.Delay(100, cancellationToken);
Console.WriteLine($"Email sent to {notification.Email} for user {notification.UserId}");
return Result.Success();
}
}
// Handler to log the event
public class LogUserRegistrationHandler : NotificationHandlerBase<UserRegisteredNotification>
{
protected override async Task<Result> HandleAsync(UserRegisteredNotification notification, PublishOptions options, CancellationToken cancellationToken)
{
// Simulate logging
await Task.Delay(50, cancellationToken);
Console.WriteLine($"Logged registration for user {notification.UserId}");
return Result.Success();
}
}
// Dispatching the notification
var result = await notifier.PublishAsync(new UserRegisteredNotification
{
UserId = Guid.NewGuid(),
Email = "user@example.com"
});
if (result.IsSuccess)
{
Console.WriteLine("User registration notification processed successfully.");
}
else
{
Console.WriteLine($"Failed: {result.Errors.FirstOrDefault()?.Message}");
}
FluentValidation setup in notifications¶
Notifications, like requests, can include a nested Validator class that extends FluentValidation.AbstractValidator<TNotification>. The NotifierBuilder automatically discovers these validators during assembly scanning and registers them for use with the ValidationPipelineBehavior.
public class EmailSentNotification : NotificationBase
{
public string EmailAddress { get; set; }
public class Validator : AbstractValidator<EmailSentNotification>
{
public Validator()
{
RuleFor(x => x.EmailAddress)
.NotEmpty().WithMessage("Email cannot be empty.")
.EmailAddress().WithMessage("Invalid email format.");
}
}
}
Using PublishOptions¶
var notifier = provider.GetRequiredService<INotifier>();
var notification = new UserRegisteredNotification { UserId = Guid.NewGuid(), Email = "user@example.com" };
// Configure PublishOptions for concurrent execution
var options = new PublishOptions
{
ExecutionMode = ExecutionMode.Concurrent,
HandleExceptionsAsResultError = true
};
var result = await notifier.PublishAsync(notification, options);
if (result.IsSuccess)
{
Console.WriteLine("Notification processed successfully by all handlers.");
}
else
{
Console.WriteLine($"Failed: {result.Errors.FirstOrDefault()?.Message}");
}
Notification metadata and cancellation¶
public class SystemEventNotification : NotificationBase
{
public string EventType { get; set; }
}
public class SystemEventLoggerHandler : NotificationHandlerBase<SystemEventNotification>
{
protected override async Task<Result> HandleAsync(SystemEventNotification notification, PublishOptions options, CancellationToken cancellationToken)
{
Console.WriteLine($"Notification ID: {notification.NotificationId}, Timestamp: {notification.NotificationTimestamp}");
await Task.Delay(5000, cancellationToken); // Long-running operation
return Result.Success();
}
}
var cts = new CancellationTokenSource(TimeSpan.FromSeconds(2));
var result = await notifier.PublishAsync(new SystemEventNotification { EventType = "SystemStarted" }, cancellationToken: cts.Token);
Practices¶
- Execution Mode Selection:
- Use
Sequentialfor ordered processing when one handler's failure must stop the remaining handlers. - Use
Concurrentfor performance when handlers can run independently. - Use
FireAndForgetonly when the caller does not need completion or failure information from handlers.
var options = new PublishOptions { ExecutionMode = ExecutionMode.Concurrent };
var result = await notifier.PublishAsync(notification, options);
- Error Aggregation:
- Check
result.Errorsto handle failures from multiple handlers.
if (result.IsFailure)
{
foreach (var error in result.Errors)
{
Console.WriteLine($"Handler error: {error.Message}");
}
}
- Order Behaviors: Ensure critical behaviors (e.g., validation) run early, as with the
Requester.
services.AddNotifier()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(RetryPipelineBehavior<,>));
Examples¶
Notification with sequential execution¶
public class UserRegisteredNotification : NotificationBase
{
public Guid UserId { get; set; }
public string Email { get; set; }
public class Validator : AbstractValidator<UserRegisteredNotification>
{
public Validator()
{
RuleFor(x => x.UserId).NotEmpty().WithMessage("UserId cannot be empty.");
RuleFor(x => x.Email)
.NotEmpty().WithMessage("Email cannot be empty.")
.EmailAddress().WithMessage("Invalid email format.");
}
}
}
public class SendEmailNotificationHandler : NotificationHandlerBase<UserRegisteredNotification>
{
protected override async Task<Result> HandleAsync(UserRegisteredNotification notification, PublishOptions options, CancellationToken cancellationToken)
{
await Task.Delay(100, cancellationToken);
Console.WriteLine($"Email sent to {notification.Email} for user {notification.UserId}");
return Result.Success();
}
}
public class LogUserRegistrationHandler : NotificationHandlerBase<UserRegisteredNotification>
{
protected override async Task<Result> HandleAsync(UserRegisteredNotification notification, PublishOptions options, CancellationToken cancellationToken)
{
await Task.Delay(50, cancellationToken);
Console.WriteLine($"Logged registration for user {notification.UserId}");
return Result.Success();
}
}
// Register the Notifier with ValidationPipelineBehavior
services.AddNotifier()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>));
var notification = new UserRegisteredNotification
{
UserId = Guid.NewGuid(),
Email = "user@example.com"
};
var result = await notifier.PublishAsync(notification);
if (result.IsSuccess)
{
Console.WriteLine("User registration notification processed successfully.");
}
else
{
Console.WriteLine($"Failed: {result.Errors.FirstOrDefault()?.Message}");
}
Notification with concurrent execution¶
// Using the same UserRegisteredNotification as above
var options = new PublishOptions { ExecutionMode = ExecutionMode.Concurrent };
var result = await notifier.PublishAsync(new UserRegisteredNotification
{
UserId = Guid.NewGuid(),
Email = "user@example.com"
}, options);
if (result.IsSuccess)
{
Console.WriteLine("User registration notification processed concurrently by all handlers.");
}
else
{
Console.WriteLine($"Failed: {result.Errors.FirstOrDefault()?.Message}");
}
Notification with fire-and-forget execution¶
public class CacheUpdateNotification : NotificationBase
{
public Guid EntityId { get; set; }
}
public class CacheUpdaterHandler : NotificationHandlerBase<CacheUpdateNotification>
{
protected override async Task<Result> HandleAsync(CacheUpdateNotification notification, PublishOptions options, CancellationToken cancellationToken)
{
await Task.Delay(500, cancellationToken); // Simulate cache update
Console.WriteLine($"Cache updated for entity {notification.EntityId}");
return Result.Success();
}
}
var options = new PublishOptions { ExecutionMode = ExecutionMode.FireAndForget };
var result = await notifier.PublishAsync(new CacheUpdateNotification
{
EntityId = Guid.NewGuid()
}, options);
Console.WriteLine("Notification dispatched in fire-and-forget mode.");
The returned result confirms dispatch only. It does not contain handler failures, and the supplied cancellation token must remain valid long enough for scheduled handlers that observe it.
Part 3: pipeline behaviors¶
Basic usage¶
Configuring behaviors¶
Pipeline behaviors can be added to both the Requester and Notifier systems to handle cross-cutting concerns. They are registered using the fluent API during system configuration and are shared between the two systems for consistency.
// Adding validation, retry, and timeout behaviors for both Requester and Notifier
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(RetryPipelineBehavior<,>))
.WithBehavior(typeof(TimeoutPipelineBehavior<,>));
services.AddNotifier()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(RetryPipelineBehavior<,>))
.WithBehavior(typeof(TimeoutPipelineBehavior<,>));
Creating a custom behavior¶
You can create custom pipeline behaviors to add your own cross-cutting concerns. To do so, inherit from PipelineBehaviorBase<TRequest, TResponse> and implement the required methods. The behavior can be used by both Requester and Notifier.
public class LoggingPipelineBehavior<TRequest, TResponse> : PipelineBehaviorBase<TRequest, TResponse>
where TRequest : class
where TResponse : IResult
{
public LoggingPipelineBehavior(ILoggerFactory loggerFactory) : base(loggerFactory) { }
protected override bool CanProcess(TRequest request, Type handlerType)
{
return true;
}
protected override async Task<TResponse> Process(TRequest request, Type handlerType, Func<Task<TResponse>> next, CancellationToken cancellationToken)
{
this.Logger.LogInformation("Processing message {MessageType}", request.GetType().Name);
var result = await next();
this.Logger.LogInformation("Message processed (success={IsSuccess})", result.IsSuccess);
return result;
}
public override bool IsHandlerSpecific() => false;
}
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(LoggingPipelineBehavior<,>));
services.AddNotifier()
.AddHandlers()
.WithBehavior(typeof(LoggingPipelineBehavior<,>));
Examples¶
Progress reporting¶
- Allows handlers and behaviors to report progress during message processing using the
SendOptions.Progressproperty. - Useful for long-running operations where you want to provide feedback to the caller.
public class LongRunningCommand : RequestBase<Unit> { }
public class LongRunningCommandHandler : RequestHandlerBase<LongRunningCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(LongRunningCommand request, SendOptions options, CancellationToken cancellationToken)
{
for (int i = 0; i <= 100; i += 10)
{
options.Progress?.Report(new ProgressReport("LongRunningCommand", new[] { $"Processing {i}%" }, i));
await Task.Delay(100, cancellationToken);
}
return Result<Unit>.Success(Unit.Value);
}
}
// Dispatch with progress reporting
var options = new SendOptions
{
Progress = new Progress<ProgressReport>(report => Console.WriteLine($"Progress: {report.Messages.First()} ({report.PercentageComplete}%)"))
};
var result = await requester.SendAsync(new LongRunningCommand(), options);
DatabaseTransactionPipelineBehavior¶
- Wraps message handling in a database transaction to ensure data consistency.
- Configured with
HandlerDatabaseTransactionAttributeto specify isolation level, rollback behavior, and theDbContexttype.
// Register the behavior
builder.Services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(DatabaseTransactionPipelineBehavior<,>));
// Command to update a user
public class UpdateUserCommand : RequestBase<Unit>
{
public Guid UserId { get; set; }
public string Username { get; set; }
public class Validator : AbstractValidator<UpdateUserCommand>
{
public Validator()
{
RuleFor(x => x.UserId).NotEmpty().WithMessage("UserId cannot be empty.");
RuleFor(x => x.Username).NotEmpty().WithMessage("Username cannot be empty.");
}
}
}
[HandlerDatabaseTransaction(contextName: "Core")] // Resolves CoreDbContext by its logical name
public class UpdateUserCommandHandler : RequestHandlerBase<UpdateUserCommand, Unit>
{
private readonly IGenericRepository<User> userRepository;
public UpdateUserCommandHandler(IGenericRepository<User> userRepository)
{
this.userRepository = userRepository;
}
protected override async Task<Result<Unit>> HandleAsync(UpdateUserCommand request, SendOptions options, CancellationToken cancellationToken)
{
var user = await this.userRepository.FindOneAsync(request.UserId, cancellationToken: cancellationToken);
if (user == null)
{
return Result<Unit>.Failure().WithMessage($"User with ID {request.UserId} not found.");
}
user.Username = request.Username;
await this.userRepository.UpdateAsync(user, cancellationToken);
return Result<Unit>.Success(Unit.Value);
}
}
// Register the behavior
services.AddSqlServerDbContext<CoreDbContext>("connection_string");
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>))
.WithBehavior(typeof(DatabaseTransactionPipelineBehavior<,>));
ChaosPipelineBehavior¶
- Injects random failures to test system resilience.
- Configure
HandlerChaosAttributewith the injection rate and enabled state.
// Command to test chaos
public class TestChaosCommand : RequestBase<Unit> { }
[HandlerChaos(0.1, true)] // 10% chance of failure
public class TestChaosCommandHandler : RequestHandlerBase<TestChaosCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(TestChaosCommand request, SendOptions options, CancellationToken cancellationToken)
{
await Task.Delay(50, cancellationToken); // Simulate some work
return Result<Unit>.Success(Unit.Value);
}
}
// Register the behavior
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(ChaosPipelineBehavior<,>));
CircuitBreakerPipelineBehavior¶
- Implements a circuit breaker pattern to prevent repeated failures.
- Configured with
HandlerCircuitBreakerAttributeto specify attempts, break duration, and backoff settings.
The current implementation wraps the breaker in WaitAndRetryForeverAsync. Exception paths continue
retrying with the configured backoff until an attempt succeeds or the caller cancels. Failed Result
values do not count as breaker failures.
// Command for a critical operation
public class CriticalOperationCommand : RequestBase<Unit> { }
[HandlerCircuitBreaker(3, 30, 200, true)] // 3 attempts, 30s break, 200ms backoff, exponential
public class CriticalOperationCommandHandler : RequestHandlerBase<CriticalOperationCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(CriticalOperationCommand request, SendOptions options, CancellationToken cancellationToken)
{
await Task.Delay(50, cancellationToken); // Simulate critical operation
return Result<Unit>.Success(Unit.Value);
}
}
// Register the behavior
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(CircuitBreakerPipelineBehavior<,>));
CacheInvalidatePipelineBehavior¶
- Invalidates cache entries after message processing.
- Configured with
HandlerCacheInvalidateAttributeto specify the cache key to invalidate.
// Command to clear user cache
public class ClearUserCacheCommand : RequestBase<Unit> { }
[HandlerCacheInvalidate("user-cache")] // Invalidate cache entries starting with "user-cache"
public class ClearUserCacheCommandHandler : RequestHandlerBase<ClearUserCacheCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(ClearUserCacheCommand request, SendOptions options, CancellationToken cancellationToken)
{
return Result<Unit>.Success(Unit.Value);
}
}
// Register the behavior
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(CacheInvalidatePipelineBehavior<,>));
AuthorizationPolicyPipelineBehavior¶
- Enforces named authorization policies before handler execution.
- Configured with
HandlerAuthorizePolicyAttributeon the handler class; all listed policies must succeed.
// Require both policies to pass before executing the handler
[HandlerAuthorizePolicy("Customers.Write", "Audited")]
public class CreateCustomerCommandHandler
: RequestHandlerBase<CreateCustomerCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(
CreateCustomerCommand request,
SendOptions options,
CancellationToken cancellationToken)
{
// Handler logic only runs if policies succeed
return Result<Unit>.Success(Unit.Value);
}
}
// Register the behavior early (before retry/timeout)
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(AuthorizationPolicyPipelineBehavior<,>));
Both authorization behaviors are provided by Presentation.Web and require an
ICurrentUserAccessor that exposes the current principal. Policy authorization also requires
ASP.NET Core authorization services.
Current limitation: the denial path creates a non-generic
Result. That is compatible with a NotifierIResultpipeline but not with Requester's typedIResult<T>pipeline. Do not rely on the Requester examples below to return a typed unauthorized or forbidden result until that conversion is corrected in the implementation.
AuthorizationRolesPipelineBehavior¶
- Enforces role-based authorization (OR semantics) before handler execution.
- Configured with
HandlerAuthorizeRolesAttributeon the handler class; user must be in at least one of the specified roles.
// Allow either Admin or Manager to execute the handler
[HandlerAuthorizeRoles("Admin", "Manager")]
public class DeleteCustomerCommandHandler
: RequestHandlerBase<DeleteCustomerCommand, Unit>
{
protected override async Task<Result<Unit>> HandleAsync(
DeleteCustomerCommand request,
SendOptions options,
CancellationToken cancellationToken)
{
// Handler logic only runs if the user has one of the roles
return Result<Unit>.Success(Unit.Value);
}
}
// Register the behavior early (before retry/timeout)
services.AddRequester()
.AddHandlers()
.WithBehavior(typeof(AuthorizationRolesPipelineBehavior<,>));
Appendix A: CQS pattern¶
The Command-Query Separation (CQS) pattern is a design principle that separates operations into two distinct categories: commands and queries. Introduced by Bertrand Meyer as part of his work on the Eiffel programming language, CQS aims to improve code clarity, maintainability, and predictability by enforcing a clear distinction between operations that modify state and those that retrieve data.
General pattern¶
In the CQS pattern:
- Commands: Operations that change system state without returning a value. For example,
UpdateUser(userId, newName)changes the user's name and does not return data. - Queries: Operations that retrieve data without changing system state. For example,
GetUser(userId)returns user data without updating it.
The key principle of CQS is that a method should either be a command or a query, but not both. This separation ensures that:
- Commands do not return values: They focus solely on state modification, making their intent clear.
- Queries do not have side effects: They are safe to call without worrying about unintended state changes.
This distinction leads to several benefits:
- Predictability: The request type shows whether an operation changes state or returns data.
- Testability: Query tests compare output with the supplied input and system state. Command tests verify the resulting state change.
- Maintainability: Separating concerns makes the codebase easier to navigate and modify, as commands and queries have distinct roles.
- Scalability: Queries can be optimized (e.g., through caching) without affecting commands, and commands can be audited or logged independently of queries.
CQS in practice¶
In practice, adhering strictly to CQS can sometimes be challenging, especially in scenarios where a command might need to return a value (e.g., the ID of a newly created resource). To address this, variations like Command-Query Responsibility Segregation (CQRS) extend CQS by allowing commands to return minimal data (e.g., a success status or identifier) while maintaining separation between state-changing and data-retrieving operations.
The Requester feature aligns with the CQS pattern by providing a structured way to handle commands and queries:
- Commands: Represented as requests that return
Result<Unit>, focusing on state modification without returning meaningful data. For example, aCreateUserCommandmight update a database and returnResult<Unit>.Success(Unit.Value)to indicate success. - Queries: Requests that return
Result<TValue>without modifying state. For example,GetUserQuerycan returnResult<User>with the user details.
This alignment offers several benefits:
- Clarity: Commands and queries have distinct roles, making the code easier to understand and maintain.
- Predictability: Commands modify state without returning data, while queries return data without modifying state, reducing side effects.
- Scalability: Separating concerns allows for better optimization, such as caching query results or applying different behaviors to commands and queries.
Requester supports a CQS organization but does not enforce whether a request changes state. Application
code remains responsible for keeping command and query semantics distinct. Behaviors such as
ValidationPipelineBehavior can apply to both categories.
Appendix B: comparison with MediatR¶
Requester and Notifier share several concepts with MediatR, but their public contracts and configuration differ:
- Similarities:
- Both support request/response (
IRequest<TResponse>) and pub/sub (INotification) patterns for dispatching messages to handlers. - Both use a pipeline behavior model to handle cross-cutting concerns (e.g., validation, logging).
-
Both integrate with dependency injection for handler resolution.
-
Differences:
- Execution modes (Notifier): Notifier selects sequential, concurrent, or fire-and-forget execution per
PublishOptions. MediatR configures its notification publisher, defaulting toForeachAwaitPublisherand allowing alternatives such asTaskWhenAllPublisher; that selection is service configuration rather than a per-publication option. - Built-in Behaviors:
RequesterandNotifierprovide pre-built behaviors (e.g., retry, timeout, chaos injection) via attributes, whereas MediatR requires custom implementation of such features. - Result Type:
RequesterusesResult<TValue>andNotifierusesResultfor consistent error handling, while MediatR allows handlers to return any type, leaving error handling to the user. - Metadata:
RequestBase<TValue>andNotificationBaseinclude built-in metadata (RequestId,NotificationId, timestamps), which MediatR lacks by default.
Appendix C: Generic handlers in Requester and Notifier¶
Overview¶
Generic handlers in the Requester and Notifier systems enable handling of generic request/notification types (e.g., GenericRequest<TData>, GenericNotification<TData>) with a single handler, reducing code duplication. They are registered using AddGenericHandlers, which discovers open generic handlers, validates constraints, and registers closed handlers (e.g., GenericDataProcessor<UserData>).
Key features¶
- Automatic Discovery:
AddGenericHandlersscans assemblies for open generic handlers and discovers type arguments based on constraints. - Constraint Validation: Ensures type arguments meet constraints (e.g.,
where TData : class, IDataItem).
Setup with AddGenericHandlers¶
Requester¶
services.AddRequester()
.AddHandlers()
.AddGenericHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>));
Notifier¶
services.AddNotifier()
.AddHandlers()
.AddGenericHandlers()
.WithBehavior(typeof(ValidationPipelineBehavior<,>));
AddGenericHandlers discovers open generic handlers (e.g., GenericDataProcessor<TData>), finds type arguments (e.g., UserData, OrderData) that satisfy constraints, and registers closed handlers.
Examples¶
Generic request handler (Requester)¶
public class ProcessDataRequest<TData> : RequestBase<string>
where TData : class, IDataItem
{
public TData Data { get; set; }
public class Validator : AbstractValidator<ProcessDataRequest<TData>>
{
public Validator()
{
RuleFor(x => x.Data).NotNull();
}
}
}
public class GenericDataProcessor<TData> : RequestHandlerBase<ProcessDataRequest<TData>, string>
where TData : class, IDataItem
{
protected override async Task<Result<string>> HandleAsync(
ProcessDataRequest<TData> request,
SendOptions options,
CancellationToken cancellationToken)
{
return Result<string>.Success($"Processed: {request.Data.Id}");
}
}
public interface IDataItem { string Id { get; set; } }
public class UserData : IDataItem { public string Id { get; set; } public string Name { get; set; } }
// Dispatching
var userRequest = new ProcessDataRequest<UserData> { Data = new UserData { Id = "user123" } };
var result = await requester.SendAsync(userRequest); // "Processed: user123"
Generic notification handler (Notifier)¶
public class GenericNotification<TData> : NotificationBase
where TData : class, IDataItem
{
public TData Data { get; set; }
public class Validator : AbstractValidator<GenericNotification<TData>>
{
public Validator()
{
RuleFor(x => x.Data).NotNull();
}
}
}
public class GenericNotificationHandler<TData> : NotificationHandlerBase<GenericNotification<TData>>
where TData : class, IDataItem
{
protected override async Task<Result> HandleAsync(
GenericNotification<TData> notification,
PublishOptions options,
CancellationToken cancellationToken)
{
Console.WriteLine($"Handled: {notification.Data.Id}");
return Result.Success();
}
}
// Dispatching
var userNotification = new GenericNotification<UserData> { Data = new UserData { Id = "user123" } };
var result = await notifier.PublishAsync(userNotification); // Logs: "Handled: user123"
Appendix D: Source-generated commands, queries, and events¶
The source-generated authoring model lets you write a command, query, or event as a single partial type. This is a convenient way to define simple request/notification flows without needing separate message and handler classes. The source generators create the necessary boilerplate code, allowing you to focus on the business logic.
Quick start¶
Add the
BridgingIT.DevKit.Common.Utilities.CodeGenpackage as an analyzer reference in the application project that contains the commands and queries.
- Use
[Command]for commands. - Use
[Query]for queries. - Use
[Event]for notifications published throughINotifier. - Add one or more
[Handle]methods with your business logic. - Use validation attributes for simple property rules.
- Use
[Validate]when the validation needs full FluentValidation code.
The [Handle] method is an instance method, so message properties can be used directly inside the method body.
Services can be added as parameters and are resolved from DI.
For a package reference, use:
<PackageReference Include="BridgingIT.DevKit.Common.Utilities.CodeGen"
Version="x.y.z"
PrivateAssets="all" />
Command without response¶
[Command]
[HandlerRetry(3, 200)]
public partial class DoSomethingCommand
{
[ValidateNotEmpty("Message cannot be empty.")]
public string Message { get; set; }
[Handle]
private async Task<Result<Unit>> HandleAsync(CancellationToken cancellationToken)
{
await Task.Delay(100, cancellationToken);
return Success();
}
}
Command with response¶
public sealed class CreateUserCommandResult
{
public Guid UserId { get; set; }
public string Username { get; set; }
}
[Command]
[HandlerDatabaseTransaction(contextName: "Core")]
public partial class CreateUserCommand
{
[ValidateNotEmpty("Username cannot be empty.")]
public string Username { get; set; }
[Handle]
private async Task<Result<CreateUserCommandResult>> HandleAsync(
IGenericRepository<User> userRepository, // DI services can be injected as parameters
CancellationToken cancellationToken)
{
var user = new User
{
Id = Guid.NewGuid(),
Username = Username
};
await userRepository.InsertAsync(user, cancellationToken);
return Success(new CreateUserCommandResult
{
UserId = user.Id,
Username = user.Username
});
}
}
Query with response¶
[Query]
[HandlerTimeout(5000)]
public partial class GetUserQuery
{
[ValidateNotEmpty("UserId cannot be empty.")]
public Guid UserId { get; set; }
[Handle]
private async Task<Result<User>> HandleAsync(
IGenericReadOnlyRepository<User> userRepository, // DI services can be injected as parameters
CancellationToken cancellationToken)
{
var user = await userRepository.FindOneAsync(
UserId,
cancellationToken: cancellationToken);
return user != null
? Success(user)
: Failure($"User with ID {UserId} not found.");
}
}
Complex validation with [Validate]¶
When a rule is more complex than a simple property validator, use [Validate]:
[Command]
public partial class CreateOrderCommand
{
[ValidateNotEmpty]
public List<OrderItem> Items { get; init; }
[Validate]
private static void Validate(InlineValidator<CreateOrderCommand> validator)
{
validator.RuleFor(x => x.Items) // Supports validation of the entire request object graph
.Must(items => items.Count <= 100)
.WithMessage("A maximum of 100 items is allowed.");
}
[Handle]
private Result<Unit> Handle()
{
return Success();
}
}
Notes¶
- The response type is inferred from the
Result<T>returned by[Handle]. Success(...)andFailure(...)can be used directly inside[Handle].- Existing handler policy attributes such as retry, timeout, authorization, and transactions still apply.
- Generated validators continue to run through the normal Requester validation pipeline.
Event quick start¶
Use [Event] to author a notification as a single partial type. Each [Handle] method becomes a generated NotificationHandlerBase<TEvent> implementation.
[Event]
[HandlerRetry(2, 300)]
public partial class UserRegisteredEvent
{
[ValidateNotEmpty("Email is required.")]
[ValidateEmail("Email must be valid.")]
public string Email { get; init; }
[Handle]
private Result Audit()
{
Console.WriteLine($"Audit user registration for {Email}");
return Success();
}
[Handle]
private async Task<Result> SendWelcomeEmailAsync(
IEmailService emailService,
CancellationToken cancellationToken)
{
await emailService.SendAsync(Email, cancellationToken);
return Success();
}
}
Event validation with [Validate]¶
For more complex event validation, use [Validate] with an InlineValidator<TEvent>:
[Event]
public partial class OrderImportedEvent
{
[ValidateNotEmpty]
public List<string> OrderIds { get; init; }
[Validate]
private static void Validate(InlineValidator<OrderImportedEvent> validator)
{
validator.RuleFor(x => x.OrderIds)
.Must(ids => ids.Count <= 100)
.WithMessage("A maximum of 100 order ids is allowed.");
}
[Handle]
private Result Handle()
{
return Success();
}
}
Combining generated and manual event handlers¶
Generated handlers do not replace the normal notifier pub/sub model. You can use multiple generated [Handle] methods and still add more manual INotificationHandler<TEvent> subscribers.
[Event]
public partial class CustomerSignedUpEvent
{
public string Email { get; init; }
[Handle]
private Result Audit()
{
Console.WriteLine($"Audit: {Email}");
return Success();
}
[Handle]
private Result UpdateReadModel()
{
Console.WriteLine($"Projection updated for {Email}");
return Success();
}
}
public class CustomerSignedUpMetricsHandler : NotificationHandlerBase<CustomerSignedUpEvent>
{
protected override Task<Result> HandleAsync(CustomerSignedUpEvent notification, PublishOptions options, CancellationToken cancellationToken)
{
Console.WriteLine($"Metrics tracked for {notification.Email}");
return Task.FromResult(Result.Success());
}
}
Event notes¶
[Event]requires a top-level, non-generic,partialclass.- Events can declare one or more
[Handle]methods. [Handle]methods must returnResultorTask<Result>.PublishOptionsandCancellationTokencan be declared as[Handle]parameters when needed.- Any other
[Handle]parameters are resolved from DI. - If the event does not explicitly inherit
NotificationBase, the generator adds it automatically. - Class-level notifier policy attributes such as retry, timeout, authorization, chaos, circuit breaker, and cache invalidation are copied to each generated handler.
- Generated validators continue to run through the normal Notifier validation pipeline.
- Manual
INotificationHandler<TEvent>implementations continue to work alongside generated handlers.