mirror of
https://github.com/Sea-Haven-Industries/shoc-backend.git
synced 2026-09-30 06:03:12 +00:00
201 lines
7.1 KiB
C#
201 lines
7.1 KiB
C#
|
|
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);
|
||
|
|
}
|
||
|
|
}
|