shoc-backend/SeaHaven.DataServices/Implementation/WorkOrderReconciliationDataService.cs

201 lines
7.1 KiB
C#
Raw Normal View History

using Data.SeaHavenIndustries;
using Microsoft.EntityFrameworkCore;
using SeaHaven.DataServices.Interfaces;
namespace SeaHaven.DataServices.Implementation
{
public sealed class WorkOrderReconciliationDataService : IWorkOrderReconciliationDataService
{
private const int SingletonId = 1;
private readonly ApplicationDbContext _context;
public WorkOrderReconciliationDataService(ApplicationDbContext context)
{
_context = context;
}
public async Task<ReconciliationJobSnapshot> EnqueueAsync(
string reason,
DateTimeOffset now,
CancellationToken cancellationToken)
{
var job = await GetOrCreateAsync(cancellationToken);
if (job.State is not ("Pending" or "Running"))
{
job.RunId = Guid.NewGuid();
job.FenceToken = null;
job.State = "Pending";
job.Reason = reason[..Math.Min(reason.Length, 128)];
job.RequestedAt = now;
job.StartedAt = null;
job.LeaseExpiresAt = null;
job.CompletedAt = null;
job.WorkOrdersProcessed = 0;
job.CommentsProcessed = 0;
job.ErrorCode = null;
try
{
await _context.SaveChangesAsync(cancellationToken);
}
catch (DbUpdateConcurrencyException)
{
_context.ChangeTracker.Clear();
job = await _context.WorkOrderReconciliationJobs
.SingleAsync(j => j.Id == SingletonId, cancellationToken);
}
}
return Snapshot(job);
}
public async Task<ReconciliationLease?> TryAcquirePendingAsync(
DateTimeOffset now,
TimeSpan leaseDuration,
CancellationToken cancellationToken)
{
var job = await GetOrCreateAsync(cancellationToken);
var leaseExpired = job.State == "Running"
&& job.LeaseExpiresAt.HasValue
&& job.LeaseExpiresAt <= now;
if (job.State != "Pending" && !leaseExpired)
return null;
job.State = "Running";
job.StartedAt ??= now;
job.LeaseExpiresAt = now.Add(leaseDuration);
job.FenceToken = Guid.NewGuid();
try
{
await _context.SaveChangesAsync(cancellationToken);
return new ReconciliationLease(job.RunId, job.FenceToken.Value);
}
catch (DbUpdateConcurrencyException)
{
_context.ChangeTracker.Clear();
return null;
}
}
public async Task<ReconciliationJobSnapshot> GetStatusAsync(
CancellationToken cancellationToken)
{
var job = await _context.WorkOrderReconciliationJobs
.AsNoTracking()
.SingleOrDefaultAsync(j => j.Id == SingletonId, cancellationToken);
return job == null
? new ReconciliationJobSnapshot(null, "Idle", null, null, null, null, 0, 0, null)
: Snapshot(job);
}
public async Task<bool> RenewLeaseAsync(
Guid runId,
Guid fenceToken,
DateTimeOffset now,
TimeSpan leaseDuration,
CancellationToken cancellationToken)
{
var job = await _context.WorkOrderReconciliationJobs.SingleOrDefaultAsync(
j => j.Id == SingletonId
&& j.RunId == runId
&& j.FenceToken == fenceToken
&& j.State == "Running",
cancellationToken);
if (job == null)
return false;
job.LeaseExpiresAt = now.Add(leaseDuration);
try
{
await _context.SaveChangesAsync(cancellationToken);
return true;
}
catch (DbUpdateConcurrencyException)
{
_context.ChangeTracker.Clear();
return false;
}
}
public Task CompleteAsync(
Guid runId,
Guid fenceToken,
DateTimeOffset completedAt,
int workOrders,
int comments,
CancellationToken cancellationToken) =>
FinishAsync(runId, fenceToken, completedAt, "Succeeded", workOrders, comments, null, cancellationToken);
public Task FailAsync(
Guid runId,
Guid fenceToken,
DateTimeOffset completedAt,
string errorCode,
CancellationToken cancellationToken) =>
FinishAsync(runId, fenceToken, completedAt, "Failed", 0, 0, errorCode, cancellationToken);
private async Task FinishAsync(
Guid runId,
Guid fenceToken,
DateTimeOffset completedAt,
string state,
int workOrders,
int comments,
string? errorCode,
CancellationToken cancellationToken)
{
var job = await _context.WorkOrderReconciliationJobs
.SingleOrDefaultAsync(
j => j.Id == SingletonId
&& j.RunId == runId
&& j.FenceToken == fenceToken
&& j.State == "Running",
cancellationToken);
if (job == null)
return;
job.State = state;
job.CompletedAt = completedAt;
job.LeaseExpiresAt = null;
job.WorkOrdersProcessed = workOrders;
job.CommentsProcessed = comments;
job.ErrorCode = errorCode;
await _context.SaveChangesAsync(cancellationToken);
}
private async Task<WorkOrderReconciliationJob> GetOrCreateAsync(
CancellationToken cancellationToken)
{
var job = await _context.WorkOrderReconciliationJobs
.SingleOrDefaultAsync(j => j.Id == SingletonId, cancellationToken);
if (job != null)
return job;
job = new WorkOrderReconciliationJob { Id = SingletonId, RunId = Guid.NewGuid() };
_context.WorkOrderReconciliationJobs.Add(job);
try
{
await _context.SaveChangesAsync(cancellationToken);
return job;
}
catch (DbUpdateException)
{
_context.ChangeTracker.Clear();
return await _context.WorkOrderReconciliationJobs
.SingleAsync(j => j.Id == SingletonId, cancellationToken);
}
}
private static ReconciliationJobSnapshot Snapshot(WorkOrderReconciliationJob job) =>
new(
job.RunId,
job.State,
job.Reason,
job.RequestedAt,
job.StartedAt,
job.CompletedAt,
job.WorkOrdersProcessed,
job.CommentsProcessed,
job.ErrorCode);
}
}