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
2 changes: 1 addition & 1 deletion src/MongoDb/CleanUp/Cleaner.cs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ internal sealed class Cleaner<TMessage>(
MongoDbCleanUpSettings settings,
MongoDbCleanUpCollectionSettings<TMessage> collectionSettings,
IMongoDatabase db,
TimeProvider timeProvider) : IOutboxCleaner where TMessage : IMessage
TimeProvider timeProvider) : IOutboxCleaner where TMessage : class, IMessage
{
private readonly IMongoCollection<TMessage> _collection = db.GetCollection<TMessage>(collectionSettings.Name);

Expand Down
2 changes: 1 addition & 1 deletion src/MongoDb/CleanUp/MongoDbCleanUpSettings.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ internal sealed record MongoDbCleanUpSettings
public TimeSpan MaxAge { get; init; } = TimeSpan.FromDays(1);
}

internal sealed record MongoDbCleanUpCollectionSettings<TMessage> where TMessage : IMessage
internal sealed record MongoDbCleanUpCollectionSettings<TMessage> where TMessage : class, IMessage
{
public required string Name { get; init; }
public required Expression<Func<TMessage, DateTime?>> ProcessedAtSelector { get; init; }
Expand Down
2 changes: 1 addition & 1 deletion src/MongoDb/Polling/BatchCompleter.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@

namespace YakShaveFx.OutboxKit.MongoDb.Polling;

internal sealed class BatchCompleter<TMessage, TId> : IBatchCompleteRetrier where TMessage : IMessage
internal sealed class BatchCompleter<TMessage, TId> : IBatchCompleteRetrier where TMessage : class, IMessage
{
private readonly IMongoCollection<TMessage> _collection;
private readonly TimeProvider _timeProvider;
Expand Down
13 changes: 6 additions & 7 deletions src/MongoDb/Polling/BatchFetcher.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
namespace YakShaveFx.OutboxKit.MongoDb.Polling;

// ReSharper disable once ClassNeverInstantiated.Global - automagically instantiated by DI
internal sealed class BatchFetcher<TMessage, TId> : IBatchFetcher where TMessage : IMessage
internal sealed class BatchFetcher<TMessage, TId> : IBatchFetcher where TMessage : class, IMessage
{
private readonly IMongoCollection<TMessage> _collection;
private readonly DistributedLockThingy _lockThingy;
Expand Down Expand Up @@ -77,14 +77,13 @@ public async Task<IBatchContext> FetchAndHoldAsync(CancellationToken ct)

private async Task<IReadOnlyCollection<IMessage>> FetchMessagesAsync(CancellationToken ct)
{
var messages = await _collection.Find(_findFilter).Sort(_sort).Limit(_batchSize).ToListAsync(ct);
var cast = new IMessage[messages.Count];
for (var i = 0; i < messages.Count; i++)
var messages = new List<TMessage>(_batchSize);
using var cursor = await _collection.Find(_findFilter).Sort(_sort).Limit(_batchSize).ToCursorAsync(ct);
while (await cursor.MoveNextAsync(ct))
{
cast[i] = messages[i];
messages.AddRange(cursor.Current);
}
Comment thread
joaofbantunes marked this conversation as resolved.

return cast;
return messages;
}

private Task<bool> HasNextAsync(CancellationToken ct) => _collection.Find(_findFilter).Limit(1).AnyAsync(ct);
Expand Down
4 changes: 2 additions & 2 deletions src/MongoDb/Polling/ConfigurationExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,7 @@ public interface IMongoDbPollingOutboxKitConfigurator
/// <returns>The <see cref="IMongoDbPollingOutboxKitConfigurator"/> instance for chaining calls.</returns>
IMongoDbPollingOutboxKitConfigurator WithCollection<TMessage, TId>(
Action<IMongoDbOutboxCollectionConfigurator<TMessage, TId>> configure)
where TMessage : IMessage;
where TMessage : class, IMessage;

/// <summary>
/// Configures the outbox polling interval.
Expand Down Expand Up @@ -99,7 +99,7 @@ IMongoDbPollingOutboxKitConfigurator WithCollection<TMessage, TId>(
/// </summary>
/// <typeparam name="TMessage">The type of the message in the outbox.</typeparam>
/// <typeparam name="TId">The type of the id of the message in the outbox.</typeparam>
public interface IMongoDbOutboxCollectionConfigurator<TMessage, TId> where TMessage : IMessage
public interface IMongoDbOutboxCollectionConfigurator<TMessage, TId> where TMessage : class, IMessage
{
/// <summary>
/// Configures the name of the outbox collection.
Expand Down
10 changes: 5 additions & 5 deletions src/MongoDb/Polling/ConfigurationImplementation.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ namespace YakShaveFx.OutboxKit.MongoDb.Polling;
internal interface IGetMongoDbOutboxCollectionConfigured
{
void ConfigureCollection<TMessage, TId>(MongoDbPollingCollectionSettings<TMessage, TId> collectionSettings)
where TMessage : IMessage;
where TMessage : class, IMessage;
}

internal interface IMongoDbOutboxCollectionConfigurator
Expand All @@ -24,7 +24,7 @@ internal interface IMongoDbOutboxCollectionConfigurator
internal sealed class MongoDbOutboxCollectionConfigurator<TMessage, TId>(
MongoDbPollingCollectionSettings<TMessage, TId>? defaults)
: IMongoDbOutboxCollectionConfigurator, IMongoDbOutboxCollectionConfigurator<TMessage, TId>
where TMessage : IMessage
where TMessage : class, IMessage
{
private MongoDbPollingCollectionSettings<TMessage, TId> _settings = defaults ?? new();

Expand Down Expand Up @@ -89,7 +89,7 @@ public IMongoDbPollingOutboxKitConfigurator WithDatabaseFactory(
}

public IMongoDbPollingOutboxKitConfigurator WithCollection<TMessage, TId>(
Action<IMongoDbOutboxCollectionConfigurator<TMessage, TId>> configure) where TMessage : IMessage
Action<IMongoDbOutboxCollectionConfigurator<TMessage, TId>> configure) where TMessage : class, IMessage
{
ArgumentNullException.ThrowIfNull(configure);
var configurator = new MongoDbOutboxCollectionConfigurator<TMessage, TId>(null);
Expand Down Expand Up @@ -189,7 +189,7 @@ private sealed class GetMongoDbOutboxCollectionConfigured(
{
public void ConfigureCollection<TMessage, TId>(
MongoDbPollingCollectionSettings<TMessage, TId> collectionSettings)
where TMessage : IMessage
where TMessage : class, IMessage
{
if (string.IsNullOrWhiteSpace(collectionSettings.Name))
throw new InvalidOperationException("Collection name must be set");
Expand Down Expand Up @@ -321,7 +321,7 @@ internal sealed record MongoDbPollingSettings
public CompletionMode CompletionMode { get; init; } = CompletionMode.Delete;
}

internal sealed record MongoDbPollingCollectionSettings<TMessage, TId> where TMessage : IMessage
internal sealed record MongoDbPollingCollectionSettings<TMessage, TId> where TMessage : class, IMessage
{
public string Name { get; init; } = "outbox_messages";
public Expression<Func<TMessage, TId>> IdSelector { get; init; } = null!;
Expand Down
Loading