Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -212,7 +212,7 @@ protected override ISender CreateSender(IWolverineRuntime runtime)
return tenantedSender;
}

if (Mode == EndpointMode.Inline)
if (SendsInline)
{
return new InlineSnsSender(runtime, this);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -704,7 +704,7 @@ protected override ISender CreateSender(IWolverineRuntime runtime)
return tenantedSender;
}

if (Mode == EndpointMode.Inline)
if (SendsInline)
{
return new InlineSqsSender(runtime, this);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ private ISender buildSenderForTopic(IWolverineRuntime runtime, AzureServiceBusTo
{
var sender = BusClient.CreateSender(topic.TopicName);

if (topic.Mode == EndpointMode.Inline)
if (topic.SendsInline)
{
var inlineSender = new InlineAzureServiceBusSender(topic, mapper, sender,
runtime.LoggerFactory.CreateLogger<InlineAzureServiceBusSender>(), runtime.Cancellation);
Expand Down Expand Up @@ -120,7 +120,7 @@ private ISender buildSenderForQueue(IWolverineRuntime runtime, AzureServiceBusQu
{
var sender = BusClient.CreateSender(queue.QueueName);

if (queue.Mode == EndpointMode.Inline)
if (queue.SendsInline)
{
var inlineSender = new InlineAzureServiceBusSender(queue, mapper, sender,
runtime.LoggerFactory.CreateLogger<InlineAzureServiceBusSender>(), runtime.Cancellation);
Expand Down
2 changes: 1 addition & 1 deletion src/Transports/Kafka/Wolverine.Kafka/KafkaTopic.cs
Original file line number Diff line number Diff line change
Expand Up @@ -331,7 +331,7 @@ protected override ISender CreateSender(IWolverineRuntime runtime)
return tenantedSender;
}

return Mode == EndpointMode.Inline
return SendsInline
? new InlineKafkaSender(this)
: new BatchedSender(this, new KafkaSenderProtocol(this), runtime.Cancellation,
runtime.LoggerFactory.CreateLogger<KafkaSenderProtocol>());
Expand Down
2 changes: 1 addition & 1 deletion src/Transports/MQTT/Wolverine.MQTT/MqttTopic.cs
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,7 @@ protected override ISender CreateSender(IWolverineRuntime runtime)

// Inline keeps the historical immediate/fire-and-forget path (this itself, via SendAsync
// below) so explicit inline usage is unaffected.
if (Mode == EndpointMode.Inline)
if (SendsInline)
{
return this;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -341,7 +341,7 @@ protected override ISender CreateSender(IWolverineRuntime runtime)
return tenantedSender;
}

return Mode == EndpointMode.Inline
return SendsInline
? new InlineRedisStreamSender(_transport, this, runtime)
: new BatchedSender(this, new RedisSenderProtocol(_transport, this), runtime.Cancellation,
runtime.LoggerFactory.CreateLogger<RedisSenderProtocol>());
Expand Down
20 changes: 20 additions & 0 deletions src/Wolverine/Configuration/Endpoint.cs
Original file line number Diff line number Diff line change
Expand Up @@ -282,6 +282,26 @@ public int MaxDegreeOfParallelism
/// </summary>
internal bool ModeIgnoresParallelism => Mode == EndpointMode.Inline;

/// <summary>
/// GH-4061. Does this endpoint send synchronously, awaiting the transport, rather than through a batching
/// sender?
///
/// <para>
/// Transports use this to choose between their inline sender and a <c>BatchedSender</c>. It is one property
/// because getting it wrong is silent and total: a <c>BatchedSender</c> only functions once the sending agent
/// has registered a callback on it, and <c>InlineSendingAgent</c> is NOT an <c>ISenderCallback</c>, so the
/// registration in <c>EndpointCollection.CreateSendingAgent</c> is skipped for it. A transport that asks
/// <c>Mode == EndpointMode.Inline</c> directly therefore hands a NativeAck endpoint an unregistered
/// <c>BatchedSender</c> whose every send throws "This sender has not been registered."
/// </para>
///
/// <para>
/// Ask this rather than comparing against <see cref="EndpointMode.Inline"/>, so that a future mode which also
/// sends inline does not require finding all nine call sites again.
/// </para>
/// </summary>
public bool SendsInline => Mode is EndpointMode.Inline or EndpointMode.NativeAck;

/// <summary>
/// GH-3712. Render <see cref="MaxDegreeOfParallelism"/> for diagnostics, saying "n/a" rather than
/// printing a dead number for a mode that never reads it.
Expand Down
Loading