From 7fc10bd8bb7948af6bc66332bcacbb23b60b0da6 Mon Sep 17 00:00:00 2001 From: radams Date: Fri, 31 Jul 2026 14:37:56 -0400 Subject: [PATCH] - remove sqlite internal message counter; rely on db as source of truth - deprecated/vulnerable package update --- LegacyTestProject/packages.config | 126 +++++++++--------- .../SqliteMessageQueue.cs | 39 +++--- MessageQueue.Http/MessageQueue.Http.csproj | 2 +- 3 files changed, 81 insertions(+), 86 deletions(-) diff --git a/LegacyTestProject/packages.config b/LegacyTestProject/packages.config index 4c65997..aeb12b0 100644 --- a/LegacyTestProject/packages.config +++ b/LegacyTestProject/packages.config @@ -1,66 +1,66 @@  - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/MessageQueue.Database.Sqlite/SqliteMessageQueue.cs b/MessageQueue.Database.Sqlite/SqliteMessageQueue.cs index fbae2b5..588b5b0 100644 --- a/MessageQueue.Database.Sqlite/SqliteMessageQueue.cs +++ b/MessageQueue.Database.Sqlite/SqliteMessageQueue.cs @@ -42,11 +42,12 @@ public SqliteMessageQueue(ILogger> logger, IOptions } using var dbContext = GetDatabaseContext(); - (_currentQueueSize, _sequenceNumber) = GetCurrentStats(dbContext); + var currentQueueSize = GetCurrentQueueSize(dbContext); + _sequenceNumber = GetCurrentSequenceNumber(dbContext); Name = opts.Name ?? nameof(SqliteMessageQueue); - _logger.LogTrace($"{Name} initialized with {_currentQueueSize} stored messages"); + _logger.LogTrace($"{Name} initialized with {currentQueueSize} stored messages"); } private readonly SqliteMessageQueueOptions _options; @@ -59,7 +60,6 @@ public SqliteMessageQueue(ILogger> logger, IOptions private readonly IMessageFormatter _messageFormatter; private readonly int? _maxQueueSize; - private int _currentQueueSize; private long _sequenceNumber; private readonly CancellationTokenSource _cancellationSource = new(); @@ -73,16 +73,20 @@ public SqliteMessageQueue(ILogger> logger, IOptions public int MaxReadCount { get; } - private static (int QueueSize, long SequenceNumber) GetCurrentStats(SqliteDatabaseContext dbContext) + private static int GetCurrentQueueSize(SqliteDatabaseContext dbContext) { var queueSize = dbContext.SqliteQueueMessages.Count(); + return queueSize; + } + private static long GetCurrentSequenceNumber(SqliteDatabaseContext dbContext) + { var sequenceNumber = dbContext.SqliteQueueMessages .Select(x => x.SequenceNumber) .Max() ?? 0L; - return (queueSize, sequenceNumber); + return sequenceNumber; } @@ -164,9 +168,11 @@ public async Task PostManyMessagesAsync(IEnumerable<(TMessage message, MessageAt throw new InvalidOperationException($"Message count exceeds max write count of {MaxWriteCount}"); } + using var dbContext = GetDatabaseContext(); if (_maxQueueSize is { } maxQueueSize) { - if (_currentQueueSize >= maxQueueSize) + + if (GetCurrentQueueSize(dbContext) >= maxQueueSize) { _logger.LogError($"{Name} {nameof(PostManyMessagesAsync)} exceeded maximum queue size of {{MaxQueueSize}}", maxQueueSize); throw new InvalidOperationException($"{Name} {nameof(PostManyMessagesAsync)} exceeded maximum queue size of {maxQueueSize}"); @@ -199,13 +205,8 @@ public async Task PostManyMessagesAsync(IEnumerable<(TMessage message, MessageAt var messageString = messageCount == 1 ? sqlMessages[0].Body : $"{messageCount} messages"; _logger.LogTrace($"{Name} {nameof(PostManyMessagesAsync)} posting to store, Message: {{Message}}", messageString); - using (var dbContext = GetDatabaseContext()) - { - dbContext.SqliteQueueMessages.AddRange(sqlMessages); - _ = await dbContext.SaveChangesAsync(linkedCancellation.Token).ConfigureAwait(false); - } - - _ = Interlocked.Add(ref _currentQueueSize, sqlMessages.Count); + dbContext.SqliteQueueMessages.AddRange(sqlMessages); + _ = await dbContext.SaveChangesAsync(linkedCancellation.Token).ConfigureAwait(false); } public Task> GetReaderAsync(MessageQueueReaderOptions options, CancellationToken cancellationToken) @@ -249,12 +250,6 @@ public Task> GetReaderAsync(MessageQueueReaderOpti { linkedCancellation.Token.ThrowIfCancellationRequested(); - if (_currentQueueSize == 0) - { - shouldWait = true; - continue; - } - using var dbContext = GetDatabaseContext(); var dbMessages = await dbContext.SqliteQueueMessages @@ -310,7 +305,9 @@ public Task> GetReaderAsync(MessageQueueReaderOpti { try { - (_currentQueueSize, _sequenceNumber) = GetCurrentStats(dbContext); + using var recoveryDbContext = GetDatabaseContext(); + var recoverySequenceNumber = GetCurrentSequenceNumber(recoveryDbContext); + _ = Interlocked.Exchange(ref _sequenceNumber, recoverySequenceNumber); } catch (Exception inner) { @@ -319,8 +316,6 @@ public Task> GetReaderAsync(MessageQueueReaderOpti throw; } - - _ = Interlocked.Add(ref _currentQueueSize, -itemsToRemove.Count); } return (completionResult, result); diff --git a/MessageQueue.Http/MessageQueue.Http.csproj b/MessageQueue.Http/MessageQueue.Http.csproj index d87ff3b..ece5ee9 100644 --- a/MessageQueue.Http/MessageQueue.Http.csproj +++ b/MessageQueue.Http/MessageQueue.Http.csproj @@ -43,7 +43,7 @@ - +