Refactor MediatR commands and update solution structure

- Consolidated commands and handlers into single files for better organization.
- Updated file naming conventions for commands and queries.
- Added explicit Git operation rules to prevent automatic commits/pushes.
- Introduced new projects and restructured solution file (`legacy` folder).
- Refactored `CreateEmailAccountCommand`, `ProcessEmailCommand`, and others to use `IUnitOfWork`.
- Enhanced `ProcessEmailCommandHandler` with attachment validation and error handling.
- Removed redundant handler files after consolidation.
- Improved code consistency and added `TODO` comments for future enhancements.
This commit is contained in:
2026-07-09 14:37:12 +02:00
parent 45654796b7
commit 50c21ee628
12 changed files with 267 additions and 242 deletions

View File

@@ -1,3 +1,5 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Domain.Entities;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailAccounts.Commands;
@@ -22,3 +24,35 @@ public record CreateEmailAccountCommand : IRequest<int>
public string? EncryptedClientSecret { get; init; } // Already encrypted by client
public bool IsActive { get; init; } = true;
}
public class CreateEmailAccountCommandHandler(IUnitOfWork unitOfWork)
: IRequestHandler<CreateEmailAccountCommand, int>
{
public async Task<int> Handle(CreateEmailAccountCommand request, CancellationToken cancellationToken)
{
var account = new EmailAccount
{
AccountName = request.AccountName,
Username = request.Username,
ImapServer = request.ImapServer,
ImapPort = request.ImapPort,
ImapUseSsl = request.ImapUseSsl,
SmtpServer = request.SmtpServer,
SmtpPort = request.SmtpPort,
SmtpUseSsl = request.SmtpUseSsl,
UseOAuth2 = request.UseOAuth2,
EncryptedPassword = request.EncryptedPassword,
TenantId = request.TenantId,
ClientId = request.ClientId,
EncryptedClientSecret = request.EncryptedClientSecret,
IsActive = request.IsActive,
AddedWhen = DateTime.Now,
AddedWho = "System" // TODO: Get from current user context
};
var createdAccount = await unitOfWork.EmailAccounts.AddAsync(account, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
return createdAccount.Id;
}
}

View File

@@ -1,37 +0,0 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Domain.Entities;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailAccounts.Commands;
public class CreateEmailAccountCommandHandler(IUnitOfWork unitOfWork)
: IRequestHandler<CreateEmailAccountCommand, int>
{
public async Task<int> Handle(CreateEmailAccountCommand request, CancellationToken cancellationToken)
{
var account = new EmailAccount
{
AccountName = request.AccountName,
Username = request.Username,
ImapServer = request.ImapServer,
ImapPort = request.ImapPort,
ImapUseSsl = request.ImapUseSsl,
SmtpServer = request.SmtpServer,
SmtpPort = request.SmtpPort,
SmtpUseSsl = request.SmtpUseSsl,
UseOAuth2 = request.UseOAuth2,
EncryptedPassword = request.EncryptedPassword,
TenantId = request.TenantId,
ClientId = request.ClientId,
EncryptedClientSecret = request.EncryptedClientSecret,
IsActive = request.IsActive,
AddedWhen = DateTime.Now,
AddedWho = "System" // TODO: Get from current user context
};
var createdAccount = await unitOfWork.EmailAccounts.AddAsync(account, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
return createdAccount.Id;
}
}

View File

@@ -1,3 +1,10 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Application.Interfaces.Services;
using DigitalData.EmailProfiler.Domain.Entities;
using DigitalData.EmailProfiler.Domain.Enums;
using DigitalData.EmailProfiler.Domain.Events;
using DigitalData.EmailProfiler.Domain.Exceptions;
using DigitalData.EmailProfiler.Domain.Services;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailProcessing.Commands;
@@ -26,3 +33,120 @@ public record AttachmentData(
byte[] Content,
string ContentType,
long SizeBytes);
public class ProcessEmailCommandHandler(
IUnitOfWork unitOfWork,
IPublisher publisher,
MessageIdGenerator messageIdGenerator,
IPdfProcessingService pdfService,
IDmsService dmsService)
: IRequestHandler<ProcessEmailCommand, int>
{
public async Task<int> Handle(ProcessEmailCommand request, CancellationToken cancellationToken)
{
// 1. Get profile with related entities
var profile = await unitOfWork.EmailProfiles.GetWithRelatedEntitiesAsync(request.ProfileId, cancellationToken)
?? throw new DomainException($"Profile with ID {request.ProfileId} not found");
// 2. Generate message ID hash for duplicate detection
var messageId = messageIdGenerator.Generate(
request.MessageId,
request.Sender,
request.ReceivedDate,
request.Subject);
// 3. Check for duplicates
var isDuplicate = await unitOfWork.EmailHistories.IsDuplicateAsync(messageId.Hash, cancellationToken);
if (isDuplicate)
{
throw new ValidationException("Email already processed (duplicate detected)", ErrorCode.DuplicateMessageId);
}
// 4. Create email history record
var emailHistory = new EmailHistory
{
ProfileId = profile.Id,
MessageIdHash = messageId.Hash,
OriginalMessageId = request.MessageId,
SenderAddress = request.Sender,
EmailDate = request.ReceivedDate,
Subject = request.Subject,
EmailBodyText = request.BodyText,
EmailBodyHtml = request.BodyHtml,
Status = EmailStatus.Processing.ToString(),
AddedWhen = DateTime.Now
};
var createdHistory = await unitOfWork.EmailHistories.AddAsync(emailHistory, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
try
{
// 5. Process attachments
foreach (var attachmentData in request.Attachments)
{
var attachment = new EmailAttachment
{
EmailHistoryId = createdHistory.Id,
OriginalFileName = attachmentData.FileName,
SavedFileName = attachmentData.FileName, // TODO: Generate unique name
FilePath = string.Empty, // TODO: Save to disk and get path
FileSize = attachmentData.SizeBytes,
Extension = Path.GetExtension(attachmentData.FileName),
Status = AttachmentStatus.Pending.ToString(),
AddedWhen = DateTime.Now
};
// Validate PDF attachments
if (attachmentData.ContentType.Contains("pdf", StringComparison.OrdinalIgnoreCase))
{
using var stream = new MemoryStream(attachmentData.Content);
var isValidPdf = await pdfService.IsValidPdfAsync(stream, cancellationToken);
if (isValidPdf)
{
attachment.MarkAsValid();
}
else
{
attachment.MarkAsCorrupt(ErrorCode.PdfStructureInvalid, "Invalid PDF structure");
}
}
createdHistory.Attachments.Add(attachment);
}
// 6. Archive to DMS if configured
if (profile.EmailProcess != null && profile.EmailProcess.EnableWindreamImport)
{
// TODO: Implement DMS archiving with indexing steps
// This will be implemented based on ProcessSteps and IndexingSteps
}
// 7. Mark as processed
createdHistory.MarkAsProcessed();
await unitOfWork.EmailHistories.UpdateAsync(createdHistory, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
// 8. Publish domain event
await publisher.Publish(
new EmailProcessedEvent(
createdHistory.Id,
profile.Id,
messageId.Hash,
EmailStatus.Processed),
cancellationToken);
return createdHistory.Id;
}
catch (Exception ex)
{
// Mark as failed
createdHistory.MarkAsFailed(ErrorCode.AttachmentExtractionFailed, ex.Message);
await unitOfWork.EmailHistories.UpdateAsync(createdHistory, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
throw;
}
}
}

View File

@@ -1,127 +0,0 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Application.Interfaces.Services;
using DigitalData.EmailProfiler.Domain.Entities;
using DigitalData.EmailProfiler.Domain.Enums;
using DigitalData.EmailProfiler.Domain.Events;
using DigitalData.EmailProfiler.Domain.Exceptions;
using DigitalData.EmailProfiler.Domain.Services;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailProcessing.Commands;
public class ProcessEmailCommandHandler(
IUnitOfWork unitOfWork,
IPublisher publisher,
MessageIdGenerator messageIdGenerator,
IPdfProcessingService pdfService,
IDmsService dmsService)
: IRequestHandler<ProcessEmailCommand, int>
{
public async Task<int> Handle(ProcessEmailCommand request, CancellationToken cancellationToken)
{
// 1. Get profile with related entities
var profile = await unitOfWork.EmailProfiles.GetWithRelatedEntitiesAsync(request.ProfileId, cancellationToken)
?? throw new DomainException($"Profile with ID {request.ProfileId} not found");
// 2. Generate message ID hash for duplicate detection
var messageId = messageIdGenerator.Generate(
request.MessageId,
request.Sender,
request.ReceivedDate,
request.Subject);
// 3. Check for duplicates
var isDuplicate = await unitOfWork.EmailHistories.IsDuplicateAsync(messageId.Hash, cancellationToken);
if (isDuplicate)
{
throw new ValidationException("Email already processed (duplicate detected)", ErrorCode.DuplicateMessageId);
}
// 4. Create email history record
var emailHistory = new EmailHistory
{
ProfileId = profile.Id,
MessageIdHash = messageId.Hash,
OriginalMessageId = request.MessageId,
SenderAddress = request.Sender,
EmailDate = request.ReceivedDate,
Subject = request.Subject,
EmailBodyText = request.BodyText,
EmailBodyHtml = request.BodyHtml,
Status = EmailStatus.Processing.ToString(),
AddedWhen = DateTime.Now
};
var createdHistory = await unitOfWork.EmailHistories.AddAsync(emailHistory, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
try
{
// 5. Process attachments
foreach (var attachmentData in request.Attachments)
{
var attachment = new EmailAttachment
{
EmailHistoryId = createdHistory.Id,
OriginalFileName = attachmentData.FileName,
SavedFileName = attachmentData.FileName, // TODO: Generate unique name
FilePath = string.Empty, // TODO: Save to disk and get path
FileSize = attachmentData.SizeBytes,
Extension = Path.GetExtension(attachmentData.FileName),
Status = AttachmentStatus.Pending.ToString(),
AddedWhen = DateTime.Now
};
// Validate PDF attachments
if (attachmentData.ContentType.Contains("pdf", StringComparison.OrdinalIgnoreCase))
{
using var stream = new MemoryStream(attachmentData.Content);
var isValidPdf = await pdfService.IsValidPdfAsync(stream, cancellationToken);
if (isValidPdf)
{
attachment.MarkAsValid();
}
else
{
attachment.MarkAsCorrupt(ErrorCode.PdfStructureInvalid, "Invalid PDF structure");
}
}
createdHistory.Attachments.Add(attachment);
}
// 6. Archive to DMS if configured
if (profile.EmailProcess != null && profile.EmailProcess.EnableWindreamImport)
{
// TODO: Implement DMS archiving with indexing steps
// This will be implemented based on ProcessSteps and IndexingSteps
}
// 7. Mark as processed
createdHistory.MarkAsProcessed();
await unitOfWork.EmailHistories.UpdateAsync(createdHistory, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
// 8. Publish domain event
await publisher.Publish(
new EmailProcessedEvent(
createdHistory.Id,
profile.Id,
messageId.Hash,
EmailStatus.Processed),
cancellationToken);
return createdHistory.Id;
}
catch (Exception ex)
{
// Mark as failed
createdHistory.MarkAsFailed(ErrorCode.AttachmentExtractionFailed, ex.Message);
await unitOfWork.EmailHistories.UpdateAsync(createdHistory, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
throw;
}
}
}

View File

@@ -1,3 +1,5 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Domain.Entities;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailProfiles.Commands;
@@ -14,3 +16,27 @@ public record CreateEmailProfileCommand : IRequest<int>
public int PollIntervalMinutes { get; init; } = 15;
public bool IsActive { get; init; } = true;
}
public class CreateEmailProfileCommandHandler(IUnitOfWork unitOfWork)
: IRequestHandler<CreateEmailProfileCommand, int>
{
public async Task<int> Handle(CreateEmailProfileCommand request, CancellationToken cancellationToken)
{
var profile = new EmailProfile
{
ProfileName = request.ProfileName,
EmailAccountId = request.EmailAccountId,
ProcessId = request.ProcessId,
ValidationSql = request.ValidationSql,
PollIntervalMinutes = request.PollIntervalMinutes,
IsActive = request.IsActive,
AddedWhen = DateTime.Now,
AddedWho = "System" // TODO: Get from current user context
};
var createdProfile = await unitOfWork.EmailProfiles.AddAsync(profile, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
return createdProfile.Id;
}
}

View File

@@ -1,29 +0,0 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Domain.Entities;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailProfiles.Commands;
public class CreateEmailProfileCommandHandler(IUnitOfWork unitOfWork)
: IRequestHandler<CreateEmailProfileCommand, int>
{
public async Task<int> Handle(CreateEmailProfileCommand request, CancellationToken cancellationToken)
{
var profile = new EmailProfile
{
ProfileName = request.ProfileName,
EmailAccountId = request.EmailAccountId,
ProcessId = request.ProcessId,
ValidationSql = request.ValidationSql,
PollIntervalMinutes = request.PollIntervalMinutes,
IsActive = request.IsActive,
AddedWhen = DateTime.Now,
AddedWho = "System" // TODO: Get from current user context
};
var createdProfile = await unitOfWork.EmailProfiles.AddAsync(profile, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
return createdProfile.Id;
}
}

View File

@@ -1,3 +1,5 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Domain.Exceptions;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailProfiles.Commands;
@@ -6,3 +8,18 @@ namespace DigitalData.EmailProfiler.Application.Features.EmailProfiles.Commands;
/// Command to delete an email profile.
/// </summary>
public record DeleteEmailProfileCommand(int Id) : IRequest<Unit>;
public class DeleteEmailProfileCommandHandler(IUnitOfWork unitOfWork)
: IRequestHandler<DeleteEmailProfileCommand, Unit>
{
public async Task<Unit> Handle(DeleteEmailProfileCommand request, CancellationToken cancellationToken)
{
var profile = await unitOfWork.EmailProfiles.GetByIdAsync(request.Id, cancellationToken)
?? throw new DomainException($"Email profile with ID {request.Id} not found");
await unitOfWork.EmailProfiles.DeleteAsync(profile, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
return Unit.Value;
}
}

View File

@@ -1,20 +0,0 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Domain.Exceptions;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailProfiles.Commands;
public class DeleteEmailProfileCommandHandler(IUnitOfWork unitOfWork)
: IRequestHandler<DeleteEmailProfileCommand, Unit>
{
public async Task<Unit> Handle(DeleteEmailProfileCommand request, CancellationToken cancellationToken)
{
var profile = await unitOfWork.EmailProfiles.GetByIdAsync(request.Id, cancellationToken)
?? throw new DomainException($"Email profile with ID {request.Id} not found");
await unitOfWork.EmailProfiles.DeleteAsync(profile, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
return Unit.Value;
}
}

View File

@@ -1,3 +1,5 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Domain.Exceptions;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailProfiles.Commands;
@@ -13,3 +15,25 @@ public record UpdateEmailProfileCommand : IRequest<Unit>
public int PollIntervalMinutes { get; init; }
public bool IsActive { get; init; }
}
public class UpdateEmailProfileCommandHandler(IUnitOfWork unitOfWork)
: IRequestHandler<UpdateEmailProfileCommand, Unit>
{
public async Task<Unit> Handle(UpdateEmailProfileCommand request, CancellationToken cancellationToken)
{
var profile = await unitOfWork.EmailProfiles.GetByIdAsync(request.Id, cancellationToken)
?? throw new DomainException($"Email profile with ID {request.Id} not found");
profile.ProfileName = request.ProfileName;
profile.ValidationSql = request.ValidationSql;
profile.PollIntervalMinutes = request.PollIntervalMinutes;
profile.IsActive = request.IsActive;
profile.ChangedWhen = DateTime.Now;
profile.ChangedWho = "System"; // TODO: Get from current user context
await unitOfWork.EmailProfiles.UpdateAsync(profile, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
return Unit.Value;
}
}

View File

@@ -1,27 +0,0 @@
using DigitalData.EmailProfiler.Application.Interfaces.Repositories;
using DigitalData.EmailProfiler.Domain.Exceptions;
using MediatR;
namespace DigitalData.EmailProfiler.Application.Features.EmailProfiles.Commands;
public class UpdateEmailProfileCommandHandler(IUnitOfWork unitOfWork)
: IRequestHandler<UpdateEmailProfileCommand, Unit>
{
public async Task<Unit> Handle(UpdateEmailProfileCommand request, CancellationToken cancellationToken)
{
var profile = await unitOfWork.EmailProfiles.GetByIdAsync(request.Id, cancellationToken)
?? throw new DomainException($"Email profile with ID {request.Id} not found");
profile.ProfileName = request.ProfileName;
profile.ValidationSql = request.ValidationSql;
profile.PollIntervalMinutes = request.PollIntervalMinutes;
profile.IsActive = request.IsActive;
profile.ChangedWhen = DateTime.Now;
profile.ChangedWho = "System"; // TODO: Get from current user context
await unitOfWork.EmailProfiles.UpdateAsync(profile, cancellationToken);
await unitOfWork.SaveChangesAsync(cancellationToken);
return Unit.Value;
}
}