Skip to content

Commit f5dcdcd

Browse files
authored
Merge pull request #246 from Vulthil/feature/testharness-fidelity
Test harness fidelity: consumer retry, Fault<T>, and request timeout
2 parents d85c339 + b48bc23 commit f5dcdcd

14 files changed

Lines changed: 175 additions & 43 deletions

File tree

docs/articles/testing.md

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -292,8 +292,9 @@ var result = await requester.RequestAsync<GetWeatherRequest, WeatherForecast>(ne
292292
result.Value.TemperatureC.ShouldBe(20);
293293
```
294294

295-
A request with neither a responder nor a registered request consumer completes with a
296-
`Messaging.Request.NoConsumer` failure; a request consumer that throws surfaces as a `Messaging.Request.Failure`.
295+
A request with neither a responder nor a registered request consumer **times out** with a
296+
`Messaging.Request.Timeout` failure — just as it would against a real broker, where no consumer means no reply;
297+
a request consumer that throws surfaces as a `Messaging.Request.Failure`.
297298

298299
### Swapping the transport in integration tests
299300

src/Vulthil.Messaging.RabbitMq/Consumers/MessageContextFactory.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,7 @@ private static MessageContext BuildMetadata(BasicDeliverEventArgs ea)
8686
return new MessageContext
8787
{
8888
MessageId = props.MessageId,
89-
CorrelationId = props.CorrelationId ?? string.Empty,
89+
CorrelationId = props.CorrelationId,
9090
RequestId = props.CorrelationId,
9191
RoutingKey = ea.RoutingKey,
9292
Headers = headers.ToDictionary(),
@@ -120,7 +120,7 @@ private static MessageContext<TMessage> BuildTypedMetadata<TMessage>(
120120
SendEndpointProvider = sendEndpointProvider,
121121
CancellationToken = cancellationToken,
122122
MessageId = props.MessageId,
123-
CorrelationId = props.CorrelationId ?? string.Empty,
123+
CorrelationId = props.CorrelationId,
124124
RequestId = props.CorrelationId,
125125
RoutingKey = ea.RoutingKey,
126126
Headers = headers.ToDictionary(),

src/Vulthil.Messaging.TestHarness/ITestHarness.cs

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,10 @@ namespace Vulthil.Messaging.TestHarness;
1010
/// <remarks>
1111
/// The harness dispatches synchronously: by the time a publish, send, or request call completes, every
1212
/// consumer (and registered <see cref="Handle{TMessage}"/>/<see cref="Respond{TRequest, TResponse}"/> stub)
13-
/// it triggered has run, so assertions need no polling. An exception thrown by a one-way consumer propagates
14-
/// to the caller of publish/send; a request consumer's exception is surfaced as a failed request result.
13+
/// it triggered has run, so assertions need no polling. A one-way consumer that throws is retried per its
14+
/// configured policy and, once the attempts are exhausted, a <c>Fault&lt;T&gt;</c> is published — the publish or
15+
/// send itself still completes — mirroring the broker transport; a request consumer's exception is surfaced as a
16+
/// failed request result.
1517
/// </remarks>
1618
public interface ITestHarness
1719
{

src/Vulthil.Messaging.TestHarness/InMemoryContext.cs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,14 +20,14 @@ public static TMessage Deserialize<TMessage>(IServiceProvider scope, MessageEnve
2020
?? throw new InvalidOperationException($"The in-memory transport could not deserialize a '{envelope.MessageType}' payload.");
2121
}
2222

23-
public static MessageContext<TMessage> Create<TMessage>(IServiceProvider scope, TMessage message, MessageEnvelope envelope, CancellationToken cancellationToken)
23+
public static MessageContext<TMessage> Create<TMessage>(IServiceProvider scope, TMessage message, MessageEnvelope envelope, CancellationToken cancellationToken, int retryCount = 0)
2424
where TMessage : notnull
2525
=> MessageContext.CreateFromEnvelope(
2626
message,
2727
envelope,
2828
routingKey: string.Empty,
29-
redelivered: false,
30-
retryCount: 0,
29+
redelivered: retryCount > 0,
30+
retryCount: retryCount,
3131
replyToFallback: null,
3232
scope.GetRequiredService<IPublisher>(),
3333
scope.GetRequiredService<ISendEndpointProvider>(),

src/Vulthil.Messaging.TestHarness/InMemoryHandlerFactory.cs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,14 +19,14 @@ internal sealed class InMemoryHandlerFactory : IMessageHandlerFactory<InMemoryHa
1919
.GetMethod(nameof(InMemoryMessageHandlers.ForRequestConsumer), BindingFlags.Public | BindingFlags.Static)
2020
?? throw new InvalidOperationException($"{nameof(InMemoryMessageHandlers)}.{nameof(InMemoryMessageHandlers.ForRequestConsumer)} not found.");
2121

22-
private readonly ConcurrentDictionary<(Type Consumer, Type Message), Func<InMemoryHandler>> _consumerCache = new();
22+
private readonly ConcurrentDictionary<(Type Consumer, Type Message), Func<RetryPolicyDefinition?, InMemoryHandler>> _consumerCache = new();
2323
private readonly ConcurrentDictionary<(Type Consumer, Type Request, Type Response), Func<InMemoryHandler>> _requestCache = new();
2424

2525
public HandlerEntry<InMemoryHandler> ForConsumer(Type consumerType, Type messageType, RetryPolicyDefinition? retryPolicy)
2626
{
2727
var factory = _consumerCache.GetOrAdd((consumerType, messageType), static key =>
28-
_consumerMethod.MakeGenericMethod(key.Consumer, key.Message).CreateDelegate<Func<InMemoryHandler>>());
29-
return new HandlerEntry<InMemoryHandler>(factory(), HandlerKind.Consumer);
28+
_consumerMethod.MakeGenericMethod(key.Consumer, key.Message).CreateDelegate<Func<RetryPolicyDefinition?, InMemoryHandler>>());
29+
return new HandlerEntry<InMemoryHandler>(factory(retryPolicy), HandlerKind.Consumer);
3030
}
3131

3232
public HandlerEntry<InMemoryHandler> ForRequestConsumer(Type consumerType, Type requestType, Type responseType, RetryPolicyDefinition? retryPolicy)

src/Vulthil.Messaging.TestHarness/InMemoryMessageHandlers.cs

Lines changed: 92 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
using System.Text.Json;
22
using Microsoft.Extensions.DependencyInjection;
33
using Vulthil.Messaging.Abstractions.Consumers;
4+
using Vulthil.Messaging.Queues;
45
using Vulthil.Messaging.Transport;
56

67
namespace Vulthil.Messaging.TestHarness;
@@ -12,23 +13,53 @@ namespace Vulthil.Messaging.TestHarness;
1213
/// </summary>
1314
internal static class InMemoryMessageHandlers
1415
{
15-
/// <summary>Builds a handler for a one-way <see cref="IConsumer{TMessage}"/>.</summary>
16-
public static InMemoryHandler ForConsumer<TConsumer, TMessage>()
16+
/// <summary>
17+
/// Builds a handler for a one-way <see cref="IConsumer{TMessage}"/>. A throwing consumer is retried in-process
18+
/// per <paramref name="retryPolicy"/> — a fresh scope per attempt, mirroring the broker transport but without
19+
/// the real back-off delays — and once the attempts are exhausted a <see cref="Fault{TMessage}"/> is published
20+
/// and the delivery completes normally, so the originating publish/send succeeds just as it would against a
21+
/// real broker.
22+
/// </summary>
23+
public static InMemoryHandler ForConsumer<TConsumer, TMessage>(RetryPolicyDefinition? retryPolicy)
1724
where TConsumer : class, IConsumer<TMessage>
1825
where TMessage : notnull
19-
=> new(HandlerKind.Consumer, async (scope, message, envelope, ct) =>
26+
=> new(HandlerKind.Consumer, async (scope, message, envelope, cancellationToken) =>
2027
{
21-
var consumer = scope.GetRequiredService<TConsumer>();
28+
var scopeFactory = scope.GetRequiredService<IServiceScopeFactory>();
2229
var harness = scope.GetRequiredService<TestHarness>();
23-
var context = InMemoryContext.Create(scope, (TMessage)message, envelope, ct);
30+
var maxRetries = Math.Max(0, retryPolicy?.MaxRetryCount ?? 0);
31+
var ignoredExceptions = retryPolicy?.GetIgnoredExceptionTypes();
2432

25-
var pipeline = ConsumePipelineFactory.Build<TMessage>(scope, terminal: c =>
33+
Exception? lastError = null;
34+
for (var attempt = 0; attempt <= maxRetries; attempt++)
2635
{
27-
harness.RecordConsumed((TMessage)message, envelope);
28-
return consumer.ConsumeAsync(c, c.CancellationToken);
29-
});
36+
await using var attemptScope = scopeFactory.CreateAsyncScope();
37+
var serviceProvider = attemptScope.ServiceProvider;
38+
var consumer = serviceProvider.GetRequiredService<TConsumer>();
39+
var context = InMemoryContext.Create(serviceProvider, (TMessage)message, envelope, cancellationToken, attempt);
40+
41+
try
42+
{
43+
var pipeline = ConsumePipelineFactory.Build<TMessage>(serviceProvider, terminal: async c =>
44+
{
45+
await consumer.ConsumeAsync(c, c.CancellationToken);
46+
harness.RecordConsumed((TMessage)message, envelope);
47+
});
3048

31-
await pipeline(context);
49+
await pipeline(context);
50+
return null;
51+
}
52+
catch (Exception ex) when (ex is not OperationCanceledException)
53+
{
54+
lastError = ex;
55+
if (ignoredExceptions is not null && ignoredExceptions.Contains(ex.GetType()))
56+
{
57+
break;
58+
}
59+
}
60+
}
61+
62+
await PublishFaultAsync<TMessage>(scope, (TMessage)message, envelope, lastError!, maxRetries, cancellationToken);
3263
return null;
3364
});
3465

@@ -73,4 +104,55 @@ public static InMemoryHandler ForRequestConsumer<TConsumer, TRequest, TResponse>
73104
return InMemoryReply.BuildFault(ex, options, envelope);
74105
}
75106
});
107+
108+
/// <summary>
109+
/// Publishes a <see cref="Fault{TMessage}"/> for a terminally-failed one-way delivery, mirroring the broker
110+
/// transport: the fault is captured (so tests can assert it) and delivered in-process to any consumer bound to
111+
/// it. Best-effort — it never disrupts completing the original delivery.
112+
/// </summary>
113+
private static async Task PublishFaultAsync<TMessage>(
114+
IServiceProvider scope,
115+
TMessage message,
116+
MessageEnvelope envelope,
117+
Exception error,
118+
int retryCount,
119+
CancellationToken cancellationToken)
120+
where TMessage : notnull
121+
{
122+
var provider = scope.GetRequiredService<IMessageConfigurationProvider>();
123+
var transport = scope.GetRequiredService<InMemoryTransport>();
124+
var harness = scope.GetRequiredService<TestHarness>();
125+
var context = InMemoryContext.Create(scope, message, envelope, cancellationToken, retryCount);
126+
127+
var fault = new Fault<TMessage>
128+
{
129+
Message = message,
130+
ExceptionMessage = error.Message,
131+
StackTrace = error.StackTrace,
132+
ExceptionType = error.GetType().FullName ?? "Unknown",
133+
FaultedAt = DateTimeOffset.UtcNow,
134+
OriginalContext = CreateSnapshot(context),
135+
};
136+
137+
var faultEnvelope = OutgoingEnvelope.Build(provider, fault, new PublishContext());
138+
harness.RecordPublished(fault, faultEnvelope);
139+
await transport.DeliverAsync(faultEnvelope, cancellationToken);
140+
}
141+
142+
private static MessageContextSnapshot CreateSnapshot<TMessage>(MessageContext<TMessage> context)
143+
where TMessage : notnull
144+
=> new()
145+
{
146+
MessageId = context.MessageId,
147+
RequestId = context.RequestId,
148+
CorrelationId = context.CorrelationId,
149+
ConversationId = context.ConversationId,
150+
InitiatorId = context.InitiatorId,
151+
SourceAddress = context.SourceAddress,
152+
DestinationAddress = context.DestinationAddress,
153+
ResponseAddress = context.ResponseAddress,
154+
FaultAddress = context.FaultAddress,
155+
RoutingKey = context.RoutingKey,
156+
RetryCount = context.RetryCount,
157+
};
76158
}

src/Vulthil.Messaging.TestHarness/InMemoryRequester.cs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -56,9 +56,9 @@ private Result<TResponse> MapReply<TResponse>(MessageEnvelope? reply)
5656
{
5757
if (reply is null)
5858
{
59-
return Result.Failure<TResponse>(Error.NotFound(
60-
"Messaging.Request.NoConsumer",
61-
"No request consumer or responder is registered for the request type."));
59+
return Result.Failure<TResponse>(Error.Failure(
60+
"Messaging.Request.Timeout",
61+
"Request timed out — no consumer or responder is registered for the request type."));
6262
}
6363

6464
var options = _provider.JsonSerializerOptions;

src/Vulthil.Messaging/IMessageConfigurationProvider.cs

Lines changed: 0 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -16,13 +16,6 @@ public interface IMessageConfigurationProvider
1616
/// <returns>The resolved <see cref="MessageConfiguration"/> instance.</returns>
1717
MessageConfiguration GetMessageConfiguration(Type messageType);
1818

19-
/// <summary>
20-
/// Gets the message configuration for the specified generic message type.
21-
/// </summary>
22-
/// <typeparam name="TMessage">The message CLR type.</typeparam>
23-
/// <returns>The resolved <see cref="MessageConfiguration"/> instance.</returns>
24-
MessageConfiguration GetMessageConfiguration<TMessage>() where TMessage : class;
25-
2619
/// <summary>
2720
/// Gets the stable wire URN for the supplied message type. Equivalent to
2821
/// <c>GetMessageConfiguration(messageType).Urn</c> — provided for clarity at call sites.

src/Vulthil.Messaging/MessagingOptions.cs

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -50,10 +50,6 @@ public MessageConfiguration GetMessageConfiguration(Type messageType)
5050
return fresh;
5151
}
5252

53-
/// <inheritdoc />
54-
public MessageConfiguration GetMessageConfiguration<TMessage>() where TMessage : class
55-
=> GetMessageConfiguration(typeof(TMessage));
56-
5753
/// <inheritdoc />
5854
public Uri GetUrn(Type messageType) => GetMessageConfiguration(messageType).Urn;
5955

src/Vulthil.Messaging/PublicAPI.Unshipped.txt

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -125,7 +125,6 @@ Vulthil.Messaging.IMessageConfigurationProvider.ConsumeFilters.get -> Vulthil.Me
125125
Vulthil.Messaging.IMessageConfigurationProvider.DefaultTimeout.get -> System.TimeSpan
126126
Vulthil.Messaging.IMessageConfigurationProvider.FaultExchangeName.get -> string!
127127
Vulthil.Messaging.IMessageConfigurationProvider.GetMessageConfiguration(System.Type! messageType) -> Vulthil.Messaging.MessageConfiguration!
128-
Vulthil.Messaging.IMessageConfigurationProvider.GetMessageConfiguration<TMessage>() -> Vulthil.Messaging.MessageConfiguration!
129128
Vulthil.Messaging.IMessageConfigurationProvider.GetMessageType(System.Uri! urn) -> System.Type?
130129
Vulthil.Messaging.IMessageConfigurationProvider.GetPartition(System.Type! messageType) -> Vulthil.Messaging.PartitionSpec?
131130
Vulthil.Messaging.IMessageConfigurationProvider.GetUrn(System.Type! messageType) -> System.Uri!
@@ -352,7 +351,7 @@ Vulthil.Messaging.Transport.MessageContext.CancellationToken.get -> System.Threa
352351
Vulthil.Messaging.Transport.MessageContext.CancellationToken.init -> void
353352
Vulthil.Messaging.Transport.MessageContext.ConversationId.get -> string?
354353
Vulthil.Messaging.Transport.MessageContext.ConversationId.init -> void
355-
Vulthil.Messaging.Transport.MessageContext.CorrelationId.get -> string!
354+
Vulthil.Messaging.Transport.MessageContext.CorrelationId.get -> string?
356355
Vulthil.Messaging.Transport.MessageContext.CorrelationId.init -> void
357356
Vulthil.Messaging.Transport.MessageContext.DestinationAddress.get -> System.Uri?
358357
Vulthil.Messaging.Transport.MessageContext.DestinationAddress.init -> void

0 commit comments

Comments
 (0)