shoc-backend/SeaHaven.DataServices/Implementation/WorkOrderDomainJobDataService.cs
Arthur Bassi f977483c17 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-10 11:01:52 -03:00

166 lines
6.1 KiB
C#

using Data.SeaHavenIndustries;
using Data.SeaHavenIndustries.Enums;
using Microsoft.EntityFrameworkCore;
using SeaHaven.DataServices.Interfaces;
namespace SeaHaven.DataServices.Implementation
{
public class WorkOrderDomainJobDataService : IWorkOrderDomainJobDataService
{
private static readonly LifecycleStatus?[] TerminalStatuses =
{
LifecycleStatus.Completed,
LifecycleStatus.Canceled
};
private readonly ApplicationDbContext _context;
public WorkOrderDomainJobDataService(ApplicationDbContext context)
{
_context = context;
}
public async Task<IReadOnlyList<int>> GetWeekRolledCandidatesAsync(
DateOnly sourceWeekStart,
DateOnly sourceWeekEnd,
int batchSize = 500,
CancellationToken cancellationToken = default)
{
var weekStartDate = sourceWeekStart.ToDateTime(TimeOnly.MinValue);
var weekEndDate = sourceWeekEnd.ToDateTime(TimeOnly.MinValue);
// Always return the first page of remaining unprocessed candidates.
// Callers must not offset-page: ledger writes shrink this set after each batch.
return await _context.workOrders
.AsNoTracking()
.Where(w => w.IsDeleted != true)
.Where(w => w.istemplate != true)
.Where(w => w.ScheduledDate != null)
.Where(w =>
w.ScheduledDate!.Value.Date >= weekStartDate.Date
&& w.ScheduledDate.Value.Date <= weekEndDate.Date)
.Where(w => !TerminalStatuses.Contains(w.LifecycleStatus))
.Where(w => !_context.WorkOrderWeekRolledLedgers.Any(l =>
l.WorkOrderId == w.Id && l.SourceWeekStart == sourceWeekStart))
.OrderBy(w => w.Id)
.Take(batchSize)
.Select(w => w.Id)
.ToListAsync(cancellationToken);
}
public async Task<WeekRolledProcessOutcome> TryProcessWeekRolledAsync(
int workOrderId,
DateOnly sourceWeekStart,
string correlationId,
CancellationToken cancellationToken = default)
{
await using var transaction = await _context.Database.BeginTransactionAsync(cancellationToken);
var alreadyProcessed = await _context.WorkOrderWeekRolledLedgers
.AnyAsync(l => l.WorkOrderId == workOrderId && l.SourceWeekStart == sourceWeekStart, cancellationToken);
if (alreadyProcessed)
{
await transaction.RollbackAsync(cancellationToken);
return new WeekRolledProcessOutcome { Skipped = true };
}
var workOrder = await _context.workOrders
.FirstOrDefaultAsync(w => w.Id == workOrderId, cancellationToken);
if (workOrder == null)
{
await transaction.RollbackAsync(cancellationToken);
return new WeekRolledProcessOutcome { Skipped = true };
}
_context.WorkOrderWeekRolledLedgers.Add(new WorkOrderWeekRolledLedger
{
WorkOrderId = workOrderId,
SourceWeekStart = sourceWeekStart,
ProcessedAt = DateTime.UtcNow,
CorrelationId = correlationId
});
var oldCarriedOver = workOrder.CarriedOver;
workOrder.CarriedOver = oldCarriedOver + 1;
try
{
await _context.SaveChangesAsync(cancellationToken);
await transaction.CommitAsync(cancellationToken);
}
catch (DbUpdateException)
{
await transaction.RollbackAsync(cancellationToken);
return new WeekRolledProcessOutcome { Skipped = true };
}
return new WeekRolledProcessOutcome
{
Processed = true,
OldCarriedOver = oldCarriedOver,
NewCarriedOver = workOrder.CarriedOver
};
}
public async Task<PastDueRefreshResult> RefreshPastDueFlagsAsync(
DateTime utcNow,
int batchSize = 500,
CancellationToken cancellationToken = default)
{
var today = utcNow.Date;
var setCount = 0;
var clearedCount = 0;
var examinedCount = 0;
var skip = 0;
while (true)
{
var batch = await _context.workOrders
.Where(w => w.IsDeleted != true)
.Where(w => w.istemplate != true)
.Where(w => w.ScheduledDate != null)
.OrderBy(w => w.Id)
.Skip(skip)
.Take(batchSize)
.ToListAsync(cancellationToken);
if (batch.Count == 0)
break;
foreach (var workOrder in batch)
{
examinedCount++;
var shouldBePastDue = workOrder.ScheduledDate!.Value.Date < today
&& !TerminalStatuses.Contains(workOrder.LifecycleStatus);
var hasFlag = workOrder.OperationalFlags.HasFlag(OperationalFlags.PastDue);
if (shouldBePastDue && !hasFlag)
{
workOrder.OperationalFlags |= OperationalFlags.PastDue;
setCount++;
}
else if (!shouldBePastDue && hasFlag)
{
workOrder.OperationalFlags &= ~OperationalFlags.PastDue;
clearedCount++;
}
}
await _context.SaveChangesAsync(cancellationToken);
skip += batchSize;
if (batch.Count < batchSize)
break;
}
return new PastDueRefreshResult
{
SetCount = setCount,
ClearedCount = clearedCount,
ExaminedCount = examinedCount
};
}
}
}