From 3be6e284777326d9786605f92bc724a0a81b44d9 Mon Sep 17 00:00:00 2001 From: TekH Date: Thu, 23 Jul 2026 11:19:59 +0200 Subject: [PATCH] refactor(infrastructure): Refactor RabbitMqEmailQueue and remove InMemoryEmailQueue, update DbContext --- .../Persistence/EmailProfilerDbContext.cs | 10 ---- .../Queue/InMemoryEmailQueue.cs | 59 ------------------- .../Queue/RabbitMqEmailQueue.cs | 20 +++---- 3 files changed, 10 insertions(+), 79 deletions(-) delete mode 100644 src/DigitalData.EmailProfiler.Infrastructure/Queue/InMemoryEmailQueue.cs diff --git a/src/DigitalData.EmailProfiler.Infrastructure/Persistence/EmailProfilerDbContext.cs b/src/DigitalData.EmailProfiler.Infrastructure/Persistence/EmailProfilerDbContext.cs index bedf352..de19223 100644 --- a/src/DigitalData.EmailProfiler.Infrastructure/Persistence/EmailProfilerDbContext.cs +++ b/src/DigitalData.EmailProfiler.Infrastructure/Persistence/EmailProfilerDbContext.cs @@ -1,4 +1,3 @@ -using DigitalData.EmailProfiler.Domain.Entities; using Microsoft.EntityFrameworkCore; namespace DigitalData.EmailProfiler.Infrastructure.Persistence; @@ -9,13 +8,4 @@ namespace DigitalData.EmailProfiler.Infrastructure.Persistence; /// public class EmailProfilerDbContext(DbContextOptions options) : DbContext(options) { - // DbSets for all entities - public DbSet EmailAccounts { get; set; } - public DbSet EmailProfiles { get; set; } - public DbSet EmailHistories { get; set; } - public DbSet EmailAttachments { get; set; } - public DbSet EmailProcesses { get; set; } - public DbSet ProcessSteps { get; set; } - public DbSet IndexingSteps { get; set; } - public DbSet EmailOutbox { get; set; } } diff --git a/src/DigitalData.EmailProfiler.Infrastructure/Queue/InMemoryEmailQueue.cs b/src/DigitalData.EmailProfiler.Infrastructure/Queue/InMemoryEmailQueue.cs deleted file mode 100644 index 88ff6ca..0000000 --- a/src/DigitalData.EmailProfiler.Infrastructure/Queue/InMemoryEmailQueue.cs +++ /dev/null @@ -1,59 +0,0 @@ -using System.Threading.Channels; -using DigitalData.EmailProfiler.Application.Common.Interfaces; -using DigitalData.EmailProfiler.Domain.Entities; - -namespace DigitalData.EmailProfiler.Infrastructure.Queue; - -/// -/// In-memory email queue implementation using System.Threading.Channels. -/// Thread-safe, high-performance queue for outgoing emails. -/// -/// NOTE: This class is OBSOLETE. Use RabbitMqEmailQueue for production. -/// InMemoryEmailQueue does not persist messages and will lose data on application restart. -/// -[Obsolete("InMemoryEmailQueue is obsolete. Use RabbitMqEmailQueue for production deployment.")] -public class InMemoryEmailQueue : IEmailQueue -{ - private readonly Channel _channel; - - public InMemoryEmailQueue() - { - var options = new BoundedChannelOptions(1000) - { - FullMode = BoundedChannelFullMode.Wait - }; - - _channel = Channel.CreateBounded(options); - } - - public async Task EnqueueAsync(EmailOutbox email, CancellationToken cancellationToken = default) - { - await _channel.Writer.WriteAsync(email, cancellationToken); - } - - public async Task DequeueAsync(CancellationToken cancellationToken = default) - { - if (await _channel.Reader.WaitToReadAsync(cancellationToken)) - { - if (_channel.Reader.TryRead(out var email)) - { - return email; - } - } - - return null; - } - - public Task GetQueueDepthAsync(CancellationToken cancellationToken = default) - { - return Task.FromResult(_channel.Reader.Count); - } - - /// - /// NOT IMPLEMENTED - InMemoryEmailQueue does not support event-driven consumers - /// - public Task StartConsumerAsync(Func onMessageReceived, CancellationToken cancellationToken = default) - { - throw new NotSupportedException("InMemoryEmailQueue does not support StartConsumerAsync. Use RabbitMqEmailQueue for event-driven consumers."); - } -} diff --git a/src/DigitalData.EmailProfiler.Infrastructure/Queue/RabbitMqEmailQueue.cs b/src/DigitalData.EmailProfiler.Infrastructure/Queue/RabbitMqEmailQueue.cs index 8752656..3d326ec 100644 --- a/src/DigitalData.EmailProfiler.Infrastructure/Queue/RabbitMqEmailQueue.cs +++ b/src/DigitalData.EmailProfiler.Infrastructure/Queue/RabbitMqEmailQueue.cs @@ -1,13 +1,13 @@ using System.Text; using System.Text.Json; using DigitalData.EmailProfiler.Application.Common.Interfaces; -using DigitalData.EmailProfiler.Domain.Common; -using DigitalData.EmailProfiler.Domain.Entities; +using DigitalData.EmailProfiler.Application.EmailSending.Commands; using DigitalData.EmailProfiler.Infrastructure.Messaging; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using RabbitMQ.Client; using RabbitMQ.Client.Events; +using DigitalData.EmailProfiler.Application.Common.Events; namespace DigitalData.EmailProfiler.Infrastructure.Queue; @@ -126,11 +126,11 @@ public class RabbitMqEmailQueue : IEmailQueue, IDisposable await _initializationTask.Value; } - public async Task EnqueueAsync(EmailOutbox email, CancellationToken cancellationToken = default) + public async Task EnqueueAsync(OutgoingEmailEvent outgoingEmailEvent, CancellationToken cancellationToken = default) { await EnsureInitializedAsync(); // Initialize on first call - var json = JsonSerializer.Serialize(email); + var json = JsonSerializer.Serialize(outgoingEmailEvent); var body = Encoding.UTF8.GetBytes(json); var properties = new BasicProperties @@ -149,7 +149,7 @@ public class RabbitMqEmailQueue : IEmailQueue, IDisposable cancellationToken: cancellationToken); } - public async Task DequeueAsync(CancellationToken cancellationToken = default) + public async Task DequeueAsync(CancellationToken cancellationToken = default) { await EnsureInitializedAsync(); // Initialize on first call @@ -161,7 +161,7 @@ public class RabbitMqEmailQueue : IEmailQueue, IDisposable try { var json = Encoding.UTF8.GetString(result.Body.ToArray()); - var email = JsonSerializer.Deserialize(json); + var email = JsonSerializer.Deserialize(json); // Acknowledge message after successful deserialization await _channel.BasicAckAsync(result.DeliveryTag, false, cancellationToken); @@ -188,7 +188,7 @@ public class RabbitMqEmailQueue : IEmailQueue, IDisposable /// Start event-driven consumer that processes messages as they arrive /// public async Task StartConsumerAsync( - Func onMessageReceived, + Func onMailReceived, CancellationToken cancellationToken = default) { await EnsureInitializedAsync(); // Initialize on first call @@ -200,14 +200,14 @@ public class RabbitMqEmailQueue : IEmailQueue, IDisposable try { var json = Encoding.UTF8.GetString(args.Body.ToArray()); - var email = JsonSerializer.Deserialize(json); + var email = JsonSerializer.Deserialize(json); if (email != null) { _logger.LogDebug("Received email message: To={To}, Subject={Subject}", email.Recipient, email.Subject); // Process message via callback - await onMessageReceived(email); + await onMailReceived(email); // Acknowledge message after successful processing await _channel.BasicAckAsync(args.DeliveryTag, false, cancellationToken); @@ -225,7 +225,7 @@ public class RabbitMqEmailQueue : IEmailQueue, IDisposable // TODO: Error Reporting Strategy // Option 1: Separate RabbitMQ Queue (emailprofiler.errors) - // - Create EmailErrorReport entity { EmailOutboxId, Exception, StackTrace, Timestamp, RetryAttempt } + // - Create EmailErrorReport entity { OutgoingEmailEventId, Exception, StackTrace, Timestamp, RetryAttempt } // - Publish to error queue: await _errorQueue.EnqueueAsync(errorReport) // - Separate worker processes error queue → Log to DB/File/External monitoring //