/
afanasevn
/
MaxSystems
Обзор
Документация
Войти
/
afanasevn
/
MaxSystems
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
src/MaxSystems.Infrastructure/Services/IntegrationDeliveryJournalService.cs
334 строки
11 KB
IBS\NAfanasev
CRON Maintenance events
22 июн 2026, 18:24
22 июн 2026, 18:24
3913607
Код
Авторство
О чём код?
using MaxSystems.Application.Abstractions; using MaxSystems.Domain.Enums; using MaxSystems.Infrastructure.Data; using Microsoft.EntityFrameworkCore; namespace MaxSystems.Infrastructure.Services; /// <summary> /// EF-реализация обновления журнала доставки по паре (<see cref="Domain.Entities.IntegrationDeliveryLog.MessageId"/>, очередь). /// </summary> public sealed class IntegrationDeliveryJournalService(AppDbContext db) : IIntegrationDeliveryJournalService { /// <summary>Максимум попыток доставки до перевода в очередь недоставленных.</summary> public const int MaxDeliveryAttempts = 3; /// <inheritdoc /> public async Task MarkProcessingAsync(Guid messageId, string queueName, CancellationToken ct = default) { var row = await FindRowAsync(messageId, queueName, ct); if (row is null) return; row.Status = IntegrationDeliveryStatus.Processing; row.UpdatedAt = DateTime.UtcNow; await db.SaveChangesAsync(ct); } /// <inheritdoc /> public async Task MarkDeliveredAsync( Guid messageId, string queueName, string deliverySummary, CancellationToken ct = default) { var row = await FindRowAsync(messageId, queueName, ct); if (row is null) return; var now = DateTime.UtcNow; row.Status = IntegrationDeliveryStatus.Delivered; row.PayloadSummary = deliverySummary; row.ErrorMessage = null; row.ProcessedAt = now; row.UpdatedAt = now; await db.SaveChangesAsync(ct); } /// <inheritdoc /> public async Task MarkFailedAsync( Guid messageId, string queueName, string errorMessage, CancellationToken ct = default) { var row = await FindRowAsync(messageId, queueName, ct); if (row is null) return; row.AttemptCount += 1; row.ErrorMessage = TruncateError(errorMessage); if (row.AttemptCount >= MaxDeliveryAttempts) { row.Status = IntegrationDeliveryStatus.DeadLettered; row.ProcessedAt = DateTime.UtcNow; } else { row.Status = IntegrationDeliveryStatus.Failed; } row.UpdatedAt = DateTime.UtcNow; await db.SaveChangesAsync(ct); } /// <inheritdoc /> public async Task MarkDeadLetteredAsync( Guid messageId, string queueName, string errorMessage, CancellationToken ct = default) { var row = await FindRowAsync(messageId, queueName, ct); if (row is null) return; var now = DateTime.UtcNow; row.Status = IntegrationDeliveryStatus.DeadLettered; row.ErrorMessage = TruncateError(errorMessage); row.ProcessedAt = now; row.UpdatedAt = now; await db.SaveChangesAsync(ct); } /// <inheritdoc /> public Task MarkProcessingForWorkQueueAsync(Guid workId, string queueName, CancellationToken ct = default) => UpdateWorkQueueRowAsync( workId, queueName, row => { row.Status = IntegrationDeliveryStatus.Processing; row.UpdatedAt = DateTime.UtcNow; }, ct); /// <inheritdoc /> public Task MarkDeliveredForWorkQueueAsync( Guid workId, string queueName, string deliverySummary, CancellationToken ct = default) => UpdateWorkQueueRowAsync( workId, queueName, row => { var now = DateTime.UtcNow; row.Status = IntegrationDeliveryStatus.Delivered; row.PayloadSummary = deliverySummary; row.ErrorMessage = null; row.ProcessedAt = now; row.UpdatedAt = now; }, ct); /// <inheritdoc /> public Task MarkFailedForWorkQueueAsync( Guid workId, string queueName, string errorMessage, CancellationToken ct = default) => UpdateWorkQueueRowAsync( workId, queueName, row => { row.AttemptCount += 1; row.ErrorMessage = TruncateError(errorMessage); if (row.AttemptCount >= MaxDeliveryAttempts) { row.Status = IntegrationDeliveryStatus.DeadLettered; row.ProcessedAt = DateTime.UtcNow; } else { row.Status = IntegrationDeliveryStatus.Failed; } row.UpdatedAt = DateTime.UtcNow; }, ct); /// <inheritdoc /> public Task MarkDeadLetteredForWorkQueueAsync( Guid workId, string queueName, string errorMessage, CancellationToken ct = default) => UpdateWorkQueueRowAsync( workId, queueName, row => { var now = DateTime.UtcNow; row.Status = IntegrationDeliveryStatus.DeadLettered; row.ErrorMessage = TruncateError(errorMessage); row.ProcessedAt = now; row.UpdatedAt = now; }, ct); /// <inheritdoc /> public Task MarkProcessingForMaintenanceQueueAsync( Guid equipmentId, Guid maintenanceRegulationId, string eventType, string queueName, CancellationToken ct = default) => UpdateMaintenanceQueueRowAsync( equipmentId, maintenanceRegulationId, eventType, queueName, row => { row.Status = IntegrationDeliveryStatus.Processing; row.UpdatedAt = DateTime.UtcNow; }, ct); /// <inheritdoc /> public Task MarkDeliveredForMaintenanceQueueAsync( Guid equipmentId, Guid maintenanceRegulationId, string eventType, string queueName, string deliverySummary, CancellationToken ct = default) => UpdateMaintenanceQueueRowAsync( equipmentId, maintenanceRegulationId, eventType, queueName, row => { var now = DateTime.UtcNow; row.Status = IntegrationDeliveryStatus.Delivered; row.PayloadSummary = deliverySummary; row.ErrorMessage = null; row.ProcessedAt = now; row.UpdatedAt = now; }, ct); /// <inheritdoc /> public Task MarkFailedForMaintenanceQueueAsync( Guid equipmentId, Guid maintenanceRegulationId, string eventType, string queueName, string errorMessage, CancellationToken ct = default) => UpdateMaintenanceQueueRowAsync( equipmentId, maintenanceRegulationId, eventType, queueName, row => { row.AttemptCount += 1; row.ErrorMessage = TruncateError(errorMessage); if (row.AttemptCount >= MaxDeliveryAttempts) { row.Status = IntegrationDeliveryStatus.DeadLettered; row.ProcessedAt = DateTime.UtcNow; } else { row.Status = IntegrationDeliveryStatus.Failed; } row.UpdatedAt = DateTime.UtcNow; }, ct); /// <summary>Обновляет активную строку журнала по workId и очереди.</summary> private async Task UpdateWorkQueueRowAsync( Guid workId, string queueName, Action<Domain.Entities.IntegrationDeliveryLog> update, CancellationToken ct) { var row = await FindActiveWorkQueueRowAsync(workId, queueName, ct); if (row is null) return; update(row); await db.SaveChangesAsync(ct); } /// <summary>Ищет последнюю незавершённую строку журнала для work + очередь.</summary> private async Task<Domain.Entities.IntegrationDeliveryLog?> FindActiveWorkQueueRowAsync( Guid workId, string queueName, CancellationToken ct) => await db.IntegrationDeliveryLogs .Where(x => x.WorkId == workId && x.QueueName == queueName && x.Status != IntegrationDeliveryStatus.Delivered && x.Status != IntegrationDeliveryStatus.DeadLettered) .OrderByDescending(x => x.UpdatedAt) .FirstOrDefaultAsync(ct); /// <summary>Обновляет активную строку журнала по оборудованию, регламенту и типу события.</summary> private async Task UpdateMaintenanceQueueRowAsync( Guid equipmentId, Guid maintenanceRegulationId, string eventType, string queueName, Action<Domain.Entities.IntegrationDeliveryLog> update, CancellationToken ct) { var row = await FindActiveMaintenanceQueueRowAsync( equipmentId, maintenanceRegulationId, eventType, queueName, ct); if (row is null) return; update(row); await db.SaveChangesAsync(ct); } /// <summary>Ищет последнюю незавершённую строку журнала для maintenance-события.</summary> private async Task<Domain.Entities.IntegrationDeliveryLog?> FindActiveMaintenanceQueueRowAsync( Guid equipmentId, Guid maintenanceRegulationId, string eventType, string queueName, CancellationToken ct) => await db.IntegrationDeliveryLogs .Where(x => x.EquipmentId == equipmentId && x.MaintenanceRegulationId == maintenanceRegulationId && x.EventType == eventType && x.QueueName == queueName && x.Status != IntegrationDeliveryStatus.Delivered && x.Status != IntegrationDeliveryStatus.DeadLettered) .OrderByDescending(x => x.UpdatedAt) .FirstOrDefaultAsync(ct); /// <summary>Ищет строку журнала по messageId и имени очереди.</summary> private async Task<Domain.Entities.IntegrationDeliveryLog?> FindRowAsync( Guid messageId, string queueName, CancellationToken ct) { return await db.IntegrationDeliveryLogs .FirstOrDefaultAsync( x => x.MessageId == messageId && x.QueueName == queueName, ct); } /// <summary>Обрезает текст ошибки до 2000 символов для хранения в БД.</summary> private static string TruncateError(string errorMessage) => errorMessage.Length <= 2000 ? errorMessage : errorMessage[..2000]; }