diff --git a/src/MongoDb/CleanUp/Cleaner.cs b/src/MongoDb/CleanUp/Cleaner.cs index c7be57b..b4c218b 100644 --- a/src/MongoDb/CleanUp/Cleaner.cs +++ b/src/MongoDb/CleanUp/Cleaner.cs @@ -9,7 +9,7 @@ internal sealed class Cleaner( MongoDbCleanUpSettings settings, MongoDbCleanUpCollectionSettings collectionSettings, IMongoDatabase db, - TimeProvider timeProvider) : IOutboxCleaner where TMessage : IMessage + TimeProvider timeProvider) : IOutboxCleaner where TMessage : class, IMessage { private readonly IMongoCollection _collection = db.GetCollection(collectionSettings.Name); diff --git a/src/MongoDb/CleanUp/MongoDbCleanUpSettings.cs b/src/MongoDb/CleanUp/MongoDbCleanUpSettings.cs index cc38f6d..34e604d 100644 --- a/src/MongoDb/CleanUp/MongoDbCleanUpSettings.cs +++ b/src/MongoDb/CleanUp/MongoDbCleanUpSettings.cs @@ -8,7 +8,7 @@ internal sealed record MongoDbCleanUpSettings public TimeSpan MaxAge { get; init; } = TimeSpan.FromDays(1); } -internal sealed record MongoDbCleanUpCollectionSettings where TMessage : IMessage +internal sealed record MongoDbCleanUpCollectionSettings where TMessage : class, IMessage { public required string Name { get; init; } public required Expression> ProcessedAtSelector { get; init; } diff --git a/src/MongoDb/Polling/BatchCompleter.cs b/src/MongoDb/Polling/BatchCompleter.cs index 6a7be84..8dc19b7 100644 --- a/src/MongoDb/Polling/BatchCompleter.cs +++ b/src/MongoDb/Polling/BatchCompleter.cs @@ -5,7 +5,7 @@ namespace YakShaveFx.OutboxKit.MongoDb.Polling; -internal sealed class BatchCompleter : IBatchCompleteRetrier where TMessage : IMessage +internal sealed class BatchCompleter : IBatchCompleteRetrier where TMessage : class, IMessage { private readonly IMongoCollection _collection; private readonly TimeProvider _timeProvider; diff --git a/src/MongoDb/Polling/BatchFetcher.cs b/src/MongoDb/Polling/BatchFetcher.cs index e3ca38e..db1cc4d 100644 --- a/src/MongoDb/Polling/BatchFetcher.cs +++ b/src/MongoDb/Polling/BatchFetcher.cs @@ -6,7 +6,7 @@ namespace YakShaveFx.OutboxKit.MongoDb.Polling; // ReSharper disable once ClassNeverInstantiated.Global - automagically instantiated by DI -internal sealed class BatchFetcher : IBatchFetcher where TMessage : IMessage +internal sealed class BatchFetcher : IBatchFetcher where TMessage : class, IMessage { private readonly IMongoCollection _collection; private readonly DistributedLockThingy _lockThingy; @@ -77,14 +77,13 @@ public async Task FetchAndHoldAsync(CancellationToken ct) private async Task> 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(_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); } - - return cast; + return messages; } private Task HasNextAsync(CancellationToken ct) => _collection.Find(_findFilter).Limit(1).AnyAsync(ct); diff --git a/src/MongoDb/Polling/ConfigurationExtensions.cs b/src/MongoDb/Polling/ConfigurationExtensions.cs index e45fb8b..766993e 100644 --- a/src/MongoDb/Polling/ConfigurationExtensions.cs +++ b/src/MongoDb/Polling/ConfigurationExtensions.cs @@ -61,7 +61,7 @@ public interface IMongoDbPollingOutboxKitConfigurator /// The instance for chaining calls. IMongoDbPollingOutboxKitConfigurator WithCollection( Action> configure) - where TMessage : IMessage; + where TMessage : class, IMessage; /// /// Configures the outbox polling interval. @@ -99,7 +99,7 @@ IMongoDbPollingOutboxKitConfigurator WithCollection( /// /// The type of the message in the outbox. /// The type of the id of the message in the outbox. -public interface IMongoDbOutboxCollectionConfigurator where TMessage : IMessage +public interface IMongoDbOutboxCollectionConfigurator where TMessage : class, IMessage { /// /// Configures the name of the outbox collection. diff --git a/src/MongoDb/Polling/ConfigurationImplementation.cs b/src/MongoDb/Polling/ConfigurationImplementation.cs index a272851..c3cfc86 100644 --- a/src/MongoDb/Polling/ConfigurationImplementation.cs +++ b/src/MongoDb/Polling/ConfigurationImplementation.cs @@ -13,7 +13,7 @@ namespace YakShaveFx.OutboxKit.MongoDb.Polling; internal interface IGetMongoDbOutboxCollectionConfigured { void ConfigureCollection(MongoDbPollingCollectionSettings collectionSettings) - where TMessage : IMessage; + where TMessage : class, IMessage; } internal interface IMongoDbOutboxCollectionConfigurator @@ -24,7 +24,7 @@ internal interface IMongoDbOutboxCollectionConfigurator internal sealed class MongoDbOutboxCollectionConfigurator( MongoDbPollingCollectionSettings? defaults) : IMongoDbOutboxCollectionConfigurator, IMongoDbOutboxCollectionConfigurator - where TMessage : IMessage + where TMessage : class, IMessage { private MongoDbPollingCollectionSettings _settings = defaults ?? new(); @@ -89,7 +89,7 @@ public IMongoDbPollingOutboxKitConfigurator WithDatabaseFactory( } public IMongoDbPollingOutboxKitConfigurator WithCollection( - Action> configure) where TMessage : IMessage + Action> configure) where TMessage : class, IMessage { ArgumentNullException.ThrowIfNull(configure); var configurator = new MongoDbOutboxCollectionConfigurator(null); @@ -189,7 +189,7 @@ private sealed class GetMongoDbOutboxCollectionConfigured( { public void ConfigureCollection( MongoDbPollingCollectionSettings collectionSettings) - where TMessage : IMessage + where TMessage : class, IMessage { if (string.IsNullOrWhiteSpace(collectionSettings.Name)) throw new InvalidOperationException("Collection name must be set"); @@ -321,7 +321,7 @@ internal sealed record MongoDbPollingSettings public CompletionMode CompletionMode { get; init; } = CompletionMode.Delete; } -internal sealed record MongoDbPollingCollectionSettings where TMessage : IMessage +internal sealed record MongoDbPollingCollectionSettings where TMessage : class, IMessage { public string Name { get; init; } = "outbox_messages"; public Expression> IdSelector { get; init; } = null!;