Skip to content
Draft
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 @@ -25,19 +25,18 @@ public async ValueTask InvokeAsync(IReceiveContext context, ReceiveDelegate next
feature.ProcessMessageEventArgs?.EntityPath
?? feature.ProcessSessionMessageEventArgs?.EntityPath
?? string.Empty;
var cancellationToken = context.CancellationToken;

try
{
await next(context);

await CompleteAsync(feature.Actions, context.Services, feature.Message, entityPath, cancellationToken);
await CompleteAsync(feature.Actions, context.Services, feature.Message, entityPath, CancellationToken.None);
}
catch
{
try
{
await AbandonAsync(feature.Actions, context.Services, feature.Message, entityPath, cancellationToken);
await AbandonAsync(feature.Actions, context.Services, feature.Message, entityPath, CancellationToken.None);
}
catch
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,10 @@ await ExecuteAsync(
item.Envelope,
cancellationToken);
}
catch (Exception) when (cancellationToken.IsCancellationRequested)
{
// Message processing was cancelled during shutdown.
}
catch (Exception ex)
{
logger.LogCritical(ex, "Error processing message");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -415,6 +415,40 @@ FROM updated_message
public Task ReleaseMessageAsync(Guid transportMessageId, CancellationToken cancellationToken)
=> ReleaseMessageAsync(transportMessageId, errorInfo: null, cancellationToken);

/// <summary>
/// Releases the messages that are still leased by the specified consumer back to the queue.
/// Messages that were deleted or already released are not affected.
/// </summary>
public async Task ReleaseLeasedMessagesAsync(
Guid[] transportMessageIds,
Guid consumerId,
CancellationToken cancellationToken)
{
await using var connection = await _connectionManager.OpenConnectionAsync(cancellationToken);
await using var command = connection.CreateCommand();

command.CommandText = $"""
WITH released AS (
UPDATE {_schemaOptions.MessageTable}
SET consumer_id = NULL
WHERE transport_message_id = ANY(@ids)
AND consumer_id = @consumer_id
RETURNING queue_id
)
SELECT pg_notify(
'{_schemaOptions.NotificationChannel}',
q.name::text
)
FROM (SELECT DISTINCT queue_id FROM released) r
JOIN {_schemaOptions.QueueTable} q ON r.queue_id = q.id;
""";

command.Parameters.Add(
new NpgsqlParameter("ids", NpgsqlDbType.Array | NpgsqlDbType.Uuid) { Value = transportMessageIds });
command.Parameters.Add(new NpgsqlParameter("consumer_id", NpgsqlDbType.Uuid) { Value = consumerId });
await command.ExecuteNonQueryAsync(cancellationToken);
}

/// <summary>
/// Appends error information to the <c>error_reason</c> JSONB array on a message.
/// Called when a message processing attempt fails but has not yet exceeded the retry limit.
Expand Down
57 changes: 48 additions & 9 deletions src/Mocha/src/Mocha.Transport.Postgres/PostgresReceiveEndpoint.cs
Original file line number Diff line number Diff line change
Expand Up @@ -117,14 +117,24 @@ private async Task PollMessagesAsync(ILogger logger, CancellationToken cancellat
continue;
}

await Parallel.ForEachAsync(
batch.Messages,
new ParallelOptions
try
{
await Parallel.ForEachAsync(
batch.Messages,
new ParallelOptions
{
MaxDegreeOfParallelism = _maxConcurrency,
CancellationToken = cancellationToken
},
(message, ct) => new ValueTask(ProcessMessageAsync(message, logger, ct)));
}
finally
{
if (cancellationToken.IsCancellationRequested)
{
MaxDegreeOfParallelism = _maxConcurrency,
CancellationToken = cancellationToken
},
(message, ct) => new ValueTask(ProcessMessageAsync(message, logger, ct)));
await ReleaseLeasedMessagesAsync(batch, logger);
}
}

// If we got a full batch, there may be more messages
hasMore = batch.Count >= _maxBatchSize;
Expand Down Expand Up @@ -190,9 +200,14 @@ await ExecuteAsync(
message,
cancellationToken);

// The message was handled, so it is deleted even when the endpoint is stopping.
await transport.MessageStore.DeleteMessageAsync(
message.TransportMessageId,
cancellationToken);
CancellationToken.None);
}
catch (Exception) when (cancellationToken.IsCancellationRequested)
{
// The interrupted message is released with the rest of its batch.
}
catch (Exception ex)
{
Expand All @@ -204,7 +219,7 @@ await transport.MessageStore.DeleteMessageAsync(
await transport.MessageStore.ReleaseMessageAsync(
message.TransportMessageId,
errorInfo,
cancellationToken);
CancellationToken.None);
}
catch (Exception releaseEx)
{
Expand All @@ -213,6 +228,27 @@ await transport.MessageStore.ReleaseMessageAsync(
}
}

private async Task ReleaseLeasedMessagesAsync(PostgresMessageBatch batch, ILogger logger)
{
var transportMessageIds = new Guid[batch.Count];
for (var i = 0; i < transportMessageIds.Length; i++)
{
transportMessageIds[i] = batch.Messages[i].TransportMessageId;
}

try
{
await transport.MessageStore.ReleaseLeasedMessagesAsync(
transportMessageIds,
_consumerId,
CancellationToken.None);
}
catch (Exception ex)
{
logger.LeasedMessagesReleaseFailed(ex, Queue.Name);
}
}

private async Task UpdateScheduledTriggerAsync(CancellationToken cancellationToken)
{
var scheduledAt = await transport.MessageStore.GetNextScheduledTimeAsync(
Expand Down Expand Up @@ -301,6 +337,9 @@ internal static partial class Logs
[LoggerMessage(LogLevel.Error, "Error releasing message {TransportMessageId}.")]
public static partial void MessageReleaseFailed(this ILogger logger, Exception exception, Guid transportMessageId);

[LoggerMessage(LogLevel.Error, "Error releasing the leased messages of queue {QueueName} while stopping.")]
public static partial void LeasedMessagesReleaseFailed(this ILogger logger, Exception exception, string queueName);

[LoggerMessage(LogLevel.Warning, "Database is unreachable for queue {QueueName}, waiting for connectivity to resume.")]
public static partial void WaitingForDatabase(this ILogger logger, string queueName);
}
Original file line number Diff line number Diff line change
Expand Up @@ -374,9 +374,9 @@ private async Task DisconnectAsync(CancellationToken cancellationToken)
/// </summary>
public async ValueTask DisposeAsync()
{
await DisconnectAsync(CancellationToken.None);
Manager.RemoveConsumer(this);
await _consumerCts.CancelAsync();
await DisconnectAsync(CancellationToken.None);
_consumerCts.Dispose();
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,22 +19,21 @@ public async ValueTask InvokeAsync(IReceiveContext context, ReceiveDelegate next
var feature = context.Features.GetOrSet<RabbitMQReceiveFeature>();
var channel = feature.Channel;
var eventArgs = feature.EventArgs;
var cancellationToken = context.CancellationToken;

try
{
await next(context);

if (channel.IsOpen)
{
await channel.BasicAckAsync(eventArgs.DeliveryTag, false, cancellationToken);
await channel.BasicAckAsync(eventArgs.DeliveryTag, false, CancellationToken.None);
}
}
catch
{
if (channel.IsOpen)
{
await channel.BasicNackAsync(eventArgs.DeliveryTag, false, true, cancellationToken);
await channel.BasicNackAsync(eventArgs.DeliveryTag, false, true, CancellationToken.None);
}

throw;
Expand Down
27 changes: 27 additions & 0 deletions src/Mocha/src/Mocha/Consumers/Batching/BatchCollector.cs
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,33 @@ public async ValueTask<BufferedEntry<TEvent>> Add(IConsumeContext<TEvent> contex
return entry;
}

/// <summary>
/// Removes an entry that has not been emitted in a batch yet.
/// </summary>
/// <returns><c>true</c> if the entry was removed; <c>false</c> if it was already emitted.</returns>
public bool TryRemove(BufferedEntry<TEvent> entry)
{
lock (_sync)
{
for (var i = 0; i < _buffer.Count; i++)
{
if (_buffer[i].Task == entry.Task)
{
_buffer.RemoveAt(i);

if (_buffer.Count == 0)
{
_delay.Cancel();
}

return true;
}
}
}

return false;
}

private async ValueTask OnDelayElapsed()
{
MessageBatch<TEvent>? batch;
Expand Down
2 changes: 1 addition & 1 deletion src/Mocha/src/Mocha/Consumers/Batching/BatchOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ public sealed class BatchOptions
public TimeSpan BatchTimeout { get; set; } = TimeSpan.FromSeconds(1);

/// <summary>
/// Gets or sets the maximum number of batches that can be processed concurrently.
/// Gets or sets the maximum number of batches that each receive endpoint processes concurrently.
/// Higher values improve throughput when batch processing is slow relative to message
/// arrival rate, at the cost of losing ordering guarantees between batches. Default: 1.
/// </summary>
Expand Down
Loading
Loading