shoc-backend/SeaHaven.Services/Implementation/WorkOrderWeekRolledService.cs
Arthur Bassi 93ef4f578b fix(work-orders): process all WeekRolled candidates across batches
Offset paging skipped remaining WOs after ledger writes; always refetch the first page of unprocessed candidates and cover with a >BatchSize test.
2026-07-24 10:13:17 -03:00

110 lines
4.2 KiB
C#

using System.Diagnostics;
using Microsoft.Extensions.Logging;
using SeaHaven.DataServices.Interfaces;
using SeaHaven.Services.DTOs;
using SeaHaven.Services.Helpers;
using SeaHaven.Services.Interfaces;
namespace SeaHaven.Services.Implementation
{
public class WorkOrderWeekRolledService : IWorkOrderWeekRolledService
{
private const int BatchSize = 500;
private readonly IWorkOrderDomainJobDataService _dataService;
private readonly IWorkOrderAuditService _auditService;
private readonly ILogger<WorkOrderWeekRolledService> _logger;
public WorkOrderWeekRolledService(
IWorkOrderDomainJobDataService dataService,
IWorkOrderAuditService auditService,
ILogger<WorkOrderWeekRolledService> logger)
{
_dataService = dataService;
_auditService = auditService;
_logger = logger;
}
public async Task<WeekRolledJobResult> ProcessWeekRolledAsync(
DateOnly sourceWeekStart,
CancellationToken cancellationToken = default)
{
var stopwatch = Stopwatch.StartNew();
var sourceWeekEnd = WorkOrderOperationalWeek.GetOperationalWeekEnd(sourceWeekStart);
var correlationId = WorkOrderOperationalWeek.BuildWeekCorrelationId(sourceWeekStart);
var processed = 0;
var skipped = 0;
var failed = 0;
_logger.LogInformation(
"WeekRolled job started. CorrelationId={CorrelationId}, SourceWeekStart={SourceWeekStart}, SourceWeekEnd={SourceWeekEnd}",
correlationId, sourceWeekStart, sourceWeekEnd);
// Always fetch the first page of remaining candidates. Processed rows write ledger
// entries that exclude them from the next query; offset paging would skip unprocessed rows.
while (!cancellationToken.IsCancellationRequested)
{
var candidates = await _dataService.GetWeekRolledCandidatesAsync(
sourceWeekStart, sourceWeekEnd, BatchSize, cancellationToken);
if (candidates.Count == 0)
break;
foreach (var workOrderId in candidates)
{
if (cancellationToken.IsCancellationRequested)
break;
try
{
var outcome = await _dataService.TryProcessWeekRolledAsync(
workOrderId, sourceWeekStart, correlationId, cancellationToken);
if (outcome.Skipped)
{
skipped++;
continue;
}
if (outcome.Processed)
{
await _auditService.LogWeekRolledAsync(
workOrderId, outcome.OldCarriedOver, outcome.NewCarriedOver, correlationId);
processed++;
}
}
catch (Exception ex)
{
failed++;
_logger.LogError(ex,
"WeekRolled failed for WorkOrderId={WorkOrderId}, CorrelationId={CorrelationId}",
workOrderId, correlationId);
}
}
if (candidates.Count < BatchSize)
break;
}
stopwatch.Stop();
var result = new WeekRolledJobResult
{
Processed = processed,
Skipped = skipped,
Failed = failed,
CorrelationId = correlationId,
SourceWeekStart = sourceWeekStart,
SourceWeekEnd = sourceWeekEnd,
DurationMs = stopwatch.ElapsedMilliseconds
};
_logger.LogInformation(
"WeekRolled job completed. CorrelationId={CorrelationId}, Processed={Processed}, Skipped={Skipped}, Failed={Failed}, DurationMs={DurationMs}",
correlationId, processed, skipped, failed, result.DurationMs);
return result;
}
}
}