Transport Consumer¶
Verified by tests
TransportConsumerBuilderExtensionsTests, TransportConsumerBuilderExtensionsServiceNameTests, TransportConsumerWorkerResilienceTests, TransportConsumerWorkerConnectionRecoveryTests, SubscriptionRetryHelperTests, SubscriptionHealthCheckTests, ServiceInstanceProviderTests — library CI run #31657041675 (2026-08-13)
The transport consumer automatically subscribes to message broker destinations and processes incoming messages. When combined with WithRouting(), subscriptions are auto-generated from your routing configuration.
Overview¶
The AddTransportConsumer() extension method:
- Auto-generates subscriptions from
RoutingOptionsconfigured viaWithRouting() - Registers
TransportConsumerOptionswith populated destinations - Starts
TransportConsumerWorkeras a hosted service
Auto-Configuration¶
The recommended approach chains WithRouting() and AddTransportConsumer():
Auto-Configuration
services.AddWhizbang()
.WithRouting(routing => {
routing
.OwnDomains("myapp.orders.commands")
.SubscribeTo("myapp.payments.events")
.Inbox.UseSharedTopic("inbox");
})
.WithEFCore<OrderDbContext>()
.WithDriver.Postgres
.AddTransportConsumer();
This auto-generates subscriptions:
- Inbox subscription from OwnDomains() - Filters commands by namespace pattern
- Event subscriptions from SubscribeTo() - Subscribes to each namespace topic
- Auto-discovered events from perspectives and receptors
What Gets Generated¶
For the configuration above, AddTransportConsumer() generates:
| Destination | Address | Routing Key |
|---|---|---|
| Inbox | inbox |
myapp.orders.commands.# |
| Payment Events | myapp.payments.events |
# |
If your service has perspectives or receptors that handle events from other namespaces, those are automatically discovered and added.
Additional Destinations¶
Add custom destinations beyond auto-generated ones:
Additional Destinations
services.AddWhizbang()
.WithRouting(routing => {
routing.OwnDomains("myapp.orders.commands");
})
.AddTransportConsumer(config => {
// Add custom destination (address + routing key)
config.AdditionalDestinations.Add(
new TransportDestination("custom-topic", "custom.events.#"));
// Add multiple custom destinations
config.AdditionalDestinations.Add(
new TransportDestination("audit-events", "#"));
});
Additional destinations are appended after auto-generated ones.
Complete Worker Setup¶
A typical worker service includes transport registration, routing, and consumer:
Complete Worker Setup
var builder = Host.CreateApplicationBuilder(args);
// 1. Register transport (Azure Service Bus or RabbitMQ)
var serviceBusConnection = builder.Configuration.GetConnectionString("servicebus")
?? throw new InvalidOperationException("Connection string not found");
builder.Services.AddAzureServiceBusTransport(serviceBusConnection);
// 2. Configure Whizbang with routing and consumer
builder.Services.AddWhizbang()
.WithRouting(routing => {
routing
.OwnDomains("myapp.orders.commands")
.SubscribeTo("myapp.payments.events", "myapp.users.events")
.Inbox.UseSharedTopic("inbox");
})
.WithEFCore<OrderDbContext>()
.WithDriver.Postgres
.AddTransportConsumer();
// 3. Register generated services
builder.Services.AddReceptors();
builder.Services.AddWhizbangDispatcher();
var host = builder.Build();
host.Run();
Transport Independence¶
The consumer configuration is transport-agnostic. The same WithRouting() and AddTransportConsumer() calls work with:
- Azure Service Bus - Creates topics and subscriptions
- RabbitMQ - Creates exchanges and queues
- In-Memory (testing) - Direct message routing
Transport-specific behavior is handled by the transport implementation registered separately.
Error Handling¶
When WithRouting() is not called before AddTransportConsumer():
Error Handling
// This throws InvalidOperationException at runtime
services.AddWhizbang()
.AddTransportConsumer(); // Error: WithRouting() must be called first
The error occurs when resolving TransportConsumerOptions from the service provider, not at registration time.
Subscription Resilience¶
New
Added in v1.0.0
By default, the transport consumer includes built-in resilience for subscription failures. Subscriptions retry forever until success or cancellation - critical for production systems where transient broker issues should not cause permanent failures.
Core Types¶
SubscriptionResilienceOptions: Configuration options for subscription retry behavior. Controls exponential backoff, retry limits, and partial subscription handling.
SubscriptionStatus:
Enumeration tracking subscription states:
- Pending - Waiting to subscribe
- Recovering - Retrying after failure
- Healthy - Successfully subscribed
- Failed - Failed (when RetryIndefinitely = false)
SubscriptionState: Tracks the current state of each subscription, including destination, status, error messages, and retry count.
SubscriptionRetryHelper: Internal helper class that implements exponential backoff logic and retry coordination.
Retry Behavior¶
The retry system uses exponential backoff:
| Property | Default | Description |
|---|---|---|
InitialRetryDelay |
1 second | Starting delay between retries |
MaxRetryDelay |
120 seconds | Cap on exponential backoff |
BackoffMultiplier |
2.0 | Delay multiplier per attempt |
InitialRetryAttempts |
5 | Attempts before reducing log verbosity |
RetryIndefinitely |
true | Never give up (recommended) |
HealthCheckInterval |
1 minute | Interval for health check sweeps |
AllowPartialSubscriptions |
true | Start with partial subscriptions |
How Exponential Backoff Works:
Attempt 1: 1 second delay
Attempt 2: 2 seconds delay (1 * 2.0)
Attempt 3: 4 seconds delay (2 * 2.0)
Attempt 4: 8 seconds delay (4 * 2.0)
Attempt 5: 16 seconds delay (8 * 2.0)
Attempt 6: 32 seconds delay (16 * 2.0)
Attempt 7: 64 seconds delay (32 * 2.0)
Attempt 8: 120 seconds delay (64 * 2.0, capped at MaxRetryDelay)
Attempt 9+: 120 seconds delay (continues at max)
Configuration¶
Customize resilience behavior through TransportConsumerConfiguration:
Configuration
services.AddWhizbang()
.WithRouting(routing => {
routing.OwnDomains("myapp.orders.commands");
})
.AddTransportConsumer(config => {
// Customize retry behavior
config.ResilienceOptions.InitialRetryDelay = TimeSpan.FromSeconds(2);
config.ResilienceOptions.MaxRetryDelay = TimeSpan.FromMinutes(5);
config.ResilienceOptions.BackoffMultiplier = 1.5;
// Allow partial failures (some subscriptions can fail)
config.ResilienceOptions.AllowPartialSubscriptions = true;
// Custom health check interval
config.ResilienceOptions.HealthCheckInterval = TimeSpan.FromSeconds(30);
});
Health Monitoring¶
When resilience is enabled, a health check is automatically registered:
Health Monitoring
Health check results: - Healthy: All subscriptions healthy - Degraded: Mixed status (some subscriptions recovering, pending, or failed) - Unhealthy: All subscriptions failed
The health check includes diagnostic data:
- failed_destinations: List of failed subscription addresses
- recovering_destinations: List of subscriptions currently retrying
Connection Recovery¶
For transports that support connection recovery (RabbitMQ, Azure Service Bus), subscriptions are automatically re-established after connection loss:
- Transport detects connection recovery
- Worker receives recovery notification via
ITransportWithRecovery - All subscriptions are reset to pending
- Retry loop re-establishes each subscription
This ensures subscriptions survive both initial failures and runtime connection issues.
Observability¶
Monitor subscription state through:
Logging:
- Initial retries (attempts 1 through InitialRetryAttempts): Warning level with full details
- Indefinite retries (beyond InitialRetryAttempts): logged only every 10th attempt to reduce noise
Metrics (if using health checks): - Total subscriptions count - Active subscriptions count - Failed subscriptions count - Recovering subscriptions count
Diagnostic Data: Access subscription state programmatically through the health check data dictionary.
Service Name Resolution¶
The consumer uses service name for subscription naming:
IServiceInstanceProvider- If registered, usesServiceNameproperty- Entry Assembly - Falls back to assembly name
- Default - Uses "UnknownService" as ultimate fallback
Register a custom provider for explicit control:
Service Name Resolution
builder.Services.AddSingleton<IServiceInstanceProvider>(
new ServiceInstanceProvider(
instanceId: TrackedGuid.NewMedo(),
serviceName: "MyOrderService",
hostName: Environment.MachineName,
processId: Environment.ProcessId));
Alternatively, set the Whizbang:ServiceName (or ServiceName) configuration key and use the default ServiceInstanceProvider(IConfiguration?) constructor, which resolves the name from configuration before falling back to the entry assembly name.
Service Name via Configuration
builder.Services.AddSingleton<IServiceInstanceProvider>(sp =>
new ServiceInstanceProvider(sp.GetRequiredService<IConfiguration>()));
Prerequisites¶
Before calling AddTransportConsumer():
- Transport Registration - Call
AddAzureServiceBusTransport()orAddRabbitMQTransport() - Routing Configuration - Call
WithRouting()to configure routing options - Receptors (optional) - Call
AddReceptors()for message handlers
Related Documentation¶
- Routing - Namespace-based routing configuration
- Inbox/Outbox - Message persistence and delivery guarantees
- Workers - Background processing workers
- RabbitMQ Transport - RabbitMQ transport configuration
- Azure Service Bus Transport - Azure Service Bus transport configuration