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 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 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 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 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 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); } }