shoc-backend/SeaHaven.DataServices/Implementation/UpliftDataService.cs
Alexandre Brandizzi f40ad7c1ef fix(uplifts): pending exposure header sums the row deltas
The approvals header summed the whole RequestedNTE for work-order-path
requests, while each Pending row shows RequestedNTE - CurrentNTE. When a
work order already had an NTE the header overstated exposure by that NTE.

The header now sums the same Delta the rows display, over the same rows
the Pending tab lists (non-deleted request on a non-deleted dispatch).
The unused duplicate aggregate is removed so one definition remains.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-25 02:10:41 -03:00

639 lines
34 KiB
C#

using System.Collections.Concurrent;
using Data.SeaHavenIndustries;
using Data.SeaHavenIndustries.Enums;
using Microsoft.EntityFrameworkCore;
using SeaHaven.DataServices.Exceptions;
using SeaHaven.DataServices.Interfaces;
namespace SeaHaven.DataServices.Implementation
{
public class UpliftDataService : IUpliftDataService
{
private static readonly ConcurrentDictionary<int, SemaphoreSlim> WorkOrderGates = new();
private readonly ApplicationDbContext _context;
public UpliftDataService(ApplicationDbContext context)
{
_context = context;
}
public async Task<(int TotalCount, IReadOnlyList<UpliftListItemData> Items)> GetPagedAsync(
string? status, int? tier, int page, int pageSize, CancellationToken cancellationToken)
{
var query = from u in _context.DispatchUpliftRequests
join d in _context.Dispatches on u.DispatchId equals d.Id
join v in _context.Vendors on d.VendorId equals v.Id into vendors
from v in vendors.DefaultIfEmpty()
join ev in _context.VendorCompletionDocuments on u.EvidenceDocumentId equals ev.Id into evidences
from ev in evidences.DefaultIfEmpty()
// SH-207: resolve the work order through the same primary-plus-linked
// convention the sibling reads use (GetWorkOrderIdForUpliftAsync,
// GetApprovedExposureForWorkOrdersAsync). A dispatch whose work order is
// linked only through DispatchWorkOrders still surfaces its work-order id,
// flattened WO context, the closed flag, and approved exposure. The
// effective id feeds a single work-order lookup so every field resolves
// from one row.
let effectiveWorkOrderId = d.WorkOrderId
?? d.DispatchWorkOrders!.Select(link => (int?)link.WorkOrderId).FirstOrDefault()
let workOrder = _context.workOrders
.FirstOrDefault(candidate => candidate.Id == effectiveWorkOrderId)
join reqUser in _context.Users on u.createdby equals reqUser.Id into requestUsers
from reqUser in requestUsers.DefaultIfEmpty()
join decUser in _context.Users on u.DecidedByUserId equals decUser.Id into decUsers
from decUser in decUsers.DefaultIfEmpty()
where (u.IsDeleted == null || u.IsDeleted == false)
&& (d.IsDeleted == null || d.IsDeleted == false)
select new { u, d, v, ev, effectiveWorkOrderId, workOrder, reqUser, decUser };
if (!string.IsNullOrWhiteSpace(status))
query = query.Where(x => x.u.Status == status);
if (tier.HasValue)
query = query.Where(x => x.u.RequiredTier == tier.Value);
var total = await query.CountAsync(cancellationToken);
// Approval queue read contract: the actionable queue (Pending) surfaces the
// oldest request first; the decision log (Approved) surfaces the most
// recently decided first. Every other read keeps the historical
// newest-request-first order. Id is the deterministic tiebreaker.
if (string.Equals(status, "Pending", StringComparison.Ordinal))
query = query.OrderBy(x => x.u.CreatedDate).ThenBy(x => x.u.Id);
else if (string.Equals(status, "Approved", StringComparison.Ordinal))
query = query.OrderByDescending(x => x.u.DecidedAt).ThenByDescending(x => x.u.Id);
else
query = query.OrderByDescending(x => x.u.CreatedDate).ThenByDescending(x => x.u.Id);
var items = await query
.Skip((Math.Max(page, 1) - 1) * pageSize)
.Take(pageSize)
.Select(x => new UpliftListItemData
{
Id = x.u.Id,
DispatchId = x.d.Id,
DispatchNumber = x.d.DispatchNumber,
PONumber = x.d.PONumber,
WorkOrderId = x.effectiveWorkOrderId,
VendorCompanyName = x.v != null ? x.v.CompanyName : x.u.RequestedByVendorName,
// The queue renders the internal display number (InternalWONumber, what the
// board shows); WorkerOrderNumber holds the CRM external id on synced
// work orders and must not leak into the queue display.
WorkOrderNumber = x.workOrder != null ? x.workOrder.InternalWONumber : null,
WorkOrderSiteCode = x.workOrder != null ? x.workOrder.SiteCode : null,
WorkOrderService = x.workOrder != null ? x.workOrder.Service : null,
// SH-209: the detail modal reads its work-order context from this row,
// so it never depends on a separate work-order fetch. Technician
// follows the board convention (vendor contact) for the requesting
// dispatch's vendor, the same vendor as VendorCompanyName.
WorkOrderScheduledDate = x.workOrder != null ? x.workOrder.ScheduledDate : null,
WorkOrderDispatcherFirstName = x.workOrder != null && x.workOrder.AssignToUser != null
? x.workOrder.AssignToUser.FirstName
: null,
WorkOrderDispatcherLastName = x.workOrder != null && x.workOrder.AssignToUser != null
? x.workOrder.AssignToUser.LastName
: null,
TechnicianName = x.v != null ? x.v.ContactName : null,
// SH-208: a work order is closed for uplift decisions once its
// lifecycle reaches a terminal state; mirrors the SH-196 revoke guard.
WorkOrderClosed = x.workOrder != null
&& (x.workOrder.LifecycleStatus == LifecycleStatus.Completed
|| x.workOrder.LifecycleStatus == LifecycleStatus.Canceled),
// SH-207: every non-deleted UpliftEvidence document on the dispatch is
// an attachment reviewers can open; the linked evidence is one of them.
AttachmentCount = x.d.CompletionDocuments.Count(doc =>
doc.Purpose == "UpliftEvidence"
&& (doc.IsDeleted == null || doc.IsDeleted == false)),
RequestedByVendorName = x.u.RequestedByVendorName,
RequestedByFirstName = x.reqUser != null ? x.reqUser.FirstName : null,
RequestedByLastName = x.reqUser != null ? x.reqUser.LastName : null,
DecidedByFirstName = x.decUser != null ? x.decUser.FirstName : null,
DecidedByLastName = x.decUser != null ? x.decUser.LastName : null,
CurrentNTE = x.u.CurrentNTE,
RequestedNTE = x.u.RequestedNTE,
VendorReason = x.u.VendorReason,
RequiredTier = x.u.RequiredTier,
Status = x.u.Status,
CreatedDate = x.u.CreatedDate,
DecidedAt = x.u.DecidedAt,
DecisionNote = x.u.DecisionNote,
EvidenceDocumentId = x.u.EvidenceDocumentId,
EvidenceFileName = x.ev != null ? x.ev.OriginalFileName : null,
EvidenceContentType = x.ev != null ? x.ev.ContentType : null,
EvidenceSizeBytes = x.ev != null ? x.ev.SizeBytes : null,
ExpiresAt = x.u.ExpiresAt,
NotificationStatus = x.u.NotificationStatus,
NotificationError = x.u.NotificationError,
})
.ToListAsync(cancellationToken);
return (total, items);
}
public async Task<IReadOnlyList<UpliftForDispatchData>> GetForDispatchAsync(int dispatchId, CancellationToken cancellationToken)
{
return await (from u in _context.DispatchUpliftRequests
where u.DispatchId == dispatchId && (u.IsDeleted == null || u.IsDeleted == false)
join dec in _context.Users on u.DecidedByUserId equals dec.Id into decs
from dec in decs.DefaultIfEmpty()
join ev in _context.VendorCompletionDocuments on u.EvidenceDocumentId equals ev.Id into evidences
from ev in evidences.DefaultIfEmpty()
orderby u.CreatedDate descending
select new UpliftForDispatchData
{
Id = u.Id,
DispatchId = u.DispatchId,
CurrentNTE = u.CurrentNTE,
RequestedNTE = u.RequestedNTE,
VendorReason = u.VendorReason,
Status = u.Status,
RequiredTier = u.RequiredTier,
RequestedByVendorName = u.RequestedByVendorName,
CreatedDate = u.CreatedDate,
DecidedAt = u.DecidedAt,
DecisionNote = u.DecisionNote,
DecidedByFirstName = dec != null ? dec.FirstName : null,
DecidedByLastName = dec != null ? dec.LastName : null,
EvidenceDocumentId = u.EvidenceDocumentId,
ExpiresAt = u.ExpiresAt,
NotificationStatus = u.NotificationStatus,
NotificationError = u.NotificationError,
EvidenceFileName = ev != null ? ev.OriginalFileName : null,
EvidenceContentType = ev != null ? ev.ContentType : null,
EvidenceSizeBytes = ev != null ? ev.SizeBytes : null,
EvidenceScanPassed = ev != null && ev.ScanStatus == "Passed"
}).ToListAsync(cancellationToken);
}
public async Task<IReadOnlyList<UpliftForWorkOrderData>> GetForWorkOrderAsync(int workOrderId, CancellationToken cancellationToken)
{
return await (from u in _context.DispatchUpliftRequests
where (u.IsDeleted == null || u.IsDeleted == false)
&& u.Dispatch != null
&& (u.Dispatch.IsDeleted == null || u.Dispatch.IsDeleted == false)
&& (
u.Dispatch.WorkOrderId == workOrderId
|| u.Dispatch.DispatchWorkOrders!.Any(link => link.WorkOrderId == workOrderId))
join dec in _context.Users on u.DecidedByUserId equals dec.Id into decs
from dec in decs.DefaultIfEmpty()
join ev in _context.VendorCompletionDocuments on u.EvidenceDocumentId equals ev.Id into evidences
from ev in evidences.DefaultIfEmpty()
orderby u.CreatedDate descending
select new UpliftForWorkOrderData
{
Id = u.Id,
DispatchId = u.DispatchId,
CurrentNTE = u.CurrentNTE,
RequestedNTE = u.RequestedNTE,
VendorReason = u.VendorReason,
Status = u.Status,
RequiredTier = u.RequiredTier,
RequestedByVendorName = u.RequestedByVendorName,
CreatedDate = u.CreatedDate,
DecidedAt = u.DecidedAt,
DecisionNote = u.DecisionNote,
DecidedByFirstName = dec != null ? dec.FirstName : null,
DecidedByLastName = dec != null ? dec.LastName : null,
EvidenceDocumentId = u.EvidenceDocumentId,
ExpiresAt = u.ExpiresAt,
NotificationStatus = u.NotificationStatus,
NotificationError = u.NotificationError,
EvidenceFileName = ev != null ? ev.OriginalFileName : null,
EvidenceContentType = ev != null ? ev.ContentType : null,
EvidenceSizeBytes = ev != null ? ev.SizeBytes : null,
EvidenceScanPassed = ev != null && ev.ScanStatus == "Passed",
CreatedByUserId = u.createdby
}).ToListAsync(cancellationToken);
}
public async Task<DispatchUpliftRequest?> GetByIdAndWorkOrderAsync(
int requestId,
int workOrderId,
CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.FirstOrDefaultAsync(u =>
u.Id == requestId
&& (u.IsDeleted == null || u.IsDeleted == false)
&& u.Dispatch != null
&& (u.Dispatch.IsDeleted == null || u.Dispatch.IsDeleted == false)
&& (
u.Dispatch.WorkOrderId == workOrderId
|| u.Dispatch.DispatchWorkOrders!.Any(link => link.WorkOrderId == workOrderId)),
cancellationToken);
}
public async Task<IReadOnlyList<PortalUpliftData>> GetForVendorDispatchAsync(int dispatchId, CancellationToken cancellationToken)
{
return await (from u in _context.DispatchUpliftRequests
where u.DispatchId == dispatchId && (u.IsDeleted == null || u.IsDeleted == false)
join dec in _context.Users on u.DecidedByUserId equals dec.Id into decs
from dec in decs.DefaultIfEmpty()
join ev in _context.VendorCompletionDocuments on u.EvidenceDocumentId equals ev.Id into evidences
from ev in evidences.DefaultIfEmpty()
orderby u.CreatedDate descending
select new PortalUpliftData
{
Id = u.Id,
CurrentNTE = u.CurrentNTE,
RequestedNTE = u.RequestedNTE,
VendorReason = u.VendorReason,
Status = u.Status,
RequiredTier = u.RequiredTier,
RequestedByVendorName = u.RequestedByVendorName,
CreatedDate = u.CreatedDate,
DecidedAt = u.DecidedAt,
DecisionNote = u.DecisionNote,
DecidedByFirstName = dec != null ? dec.FirstName : null,
DecidedByLastName = dec != null ? dec.LastName : null,
EvidenceDocumentId = u.EvidenceDocumentId,
ExpiresAt = u.ExpiresAt,
NotificationStatus = u.NotificationStatus,
NotificationError = u.NotificationError,
EvidenceFileName = ev != null ? ev.OriginalFileName : null,
EvidenceContentType = ev != null ? ev.ContentType : null,
EvidenceSizeBytes = ev != null ? ev.SizeBytes : null,
EvidenceScanPassed = ev != null && ev.ScanStatus == "Passed"
}).ToListAsync(cancellationToken);
}
public async Task<DispatchUpliftRequest?> GetByIdAsync(int id, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.FirstOrDefaultAsync(u => u.Id == id, cancellationToken);
}
public async Task<int?> GetWorkOrderIdForUpliftAsync(
int upliftRequestId,
CancellationToken cancellationToken)
{
var link = await _context.DispatchUpliftRequests
.Where(u => u.Id == upliftRequestId
&& (u.IsDeleted == null || u.IsDeleted == false)
&& u.Dispatch != null
&& (u.Dispatch.IsDeleted == null || u.Dispatch.IsDeleted == false))
.Select(u => new
{
PrimaryWorkOrderId = u.Dispatch!.WorkOrderId,
LinkedWorkOrderId = u.Dispatch.DispatchWorkOrders!
.Select(dispatchWorkOrder => (int?)dispatchWorkOrder.WorkOrderId)
.FirstOrDefault()
})
.FirstOrDefaultAsync(cancellationToken);
return link?.PrimaryWorkOrderId ?? link?.LinkedWorkOrderId;
}
public async Task<DispatchUpliftRequest?> GetByIdAndDispatchAsync(int requestId, int dispatchId, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.FirstOrDefaultAsync(u => u.Id == requestId && u.DispatchId == dispatchId, cancellationToken);
}
// SH-101: internal evidence download. The inner join on EvidenceDocumentId together
// with the DispatchId equality filter enforces server-side request/document linkage:
// a row is returned only when the document is the one linked to this exact request
// and dispatch. Soft-deleted requests/documents never resolve.
public async Task<UpliftEvidenceDownloadData?> GetEvidenceForInternalDownloadAsync(int upliftRequestId, CancellationToken cancellationToken)
{
return await (from u in _context.DispatchUpliftRequests
where u.Id == upliftRequestId && (u.IsDeleted == null || u.IsDeleted == false)
join ev in _context.VendorCompletionDocuments on u.EvidenceDocumentId equals ev.Id
where ev.DispatchId == u.DispatchId && (ev.IsDeleted == null || ev.IsDeleted == false)
select new UpliftEvidenceDownloadData
{
Id = u.Id,
DispatchId = u.DispatchId,
VendorId = ev.VendorId,
RequiredTier = u.RequiredTier,
EvidenceDocumentId = u.EvidenceDocumentId,
StoredFileName = ev.StoredFileName,
OriginalFileName = ev.OriginalFileName,
ContentType = ev.ContentType,
Purpose = ev.Purpose,
ScanStatus = ev.ScanStatus
}).FirstOrDefaultAsync(cancellationToken);
}
public async Task<bool> HasPendingAsync(int dispatchId, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.AnyAsync(u => u.DispatchId == dispatchId
&& u.Status == "Pending"
&& (u.IsDeleted == null || u.IsDeleted == false), cancellationToken);
}
public Task<bool> HasPendingForWorkOrderAsync(int workOrderId, CancellationToken cancellationToken)
{
return ForWorkOrder(workOrderId)
.AnyAsync(u => u.Status == "Pending" || u.Status == "ChangesRequested", cancellationToken);
}
public async Task<Dispatch?> GetUpliftDispatchForWorkOrderAsync(
int workOrderId,
int? primaryDispatchId,
CancellationToken cancellationToken)
{
var candidates = await _context.Dispatches
.Where(d =>
(d.IsDeleted == null || d.IsDeleted == false)
&& (
d.WorkOrderId == workOrderId
|| d.DispatchWorkOrders!.Any(link => link.WorkOrderId == workOrderId)))
.OrderByDescending(d => d.DispatchedAt ?? d.CreatedDate)
.ThenByDescending(d => d.Id)
.ToListAsync(cancellationToken);
return candidates.FirstOrDefault(d => d.Id == primaryDispatchId)
?? candidates.FirstOrDefault();
}
public Task<decimal> SumAutoApprovedAmountForWorkOrderAsync(int workOrderId, CancellationToken cancellationToken)
{
return ForWorkOrder(workOrderId)
.Where(u => u.Status == "NoApprovalRequired")
.SumAsync(u => u.RequestedNTE, cancellationToken);
}
// Approval queue read contract: set-based exposure aggregation. Dispatches map to
// work orders through the server-derived linkage (primary work order plus linked
// work orders), matching the ForWorkOrder scope; approved amounts are summed per
// dispatch in SQL and folded into per-work-order totals in memory (three set-based
// round trips, no query-per-work-order).
public async Task<IReadOnlyList<WorkOrderUpliftExposureData>> GetApprovedExposureForWorkOrdersAsync(
IReadOnlyCollection<int> workOrderIds, CancellationToken cancellationToken)
{
var workOrderIdsScope = workOrderIds.Distinct().ToList();
if (workOrderIdsScope.Count == 0)
return Array.Empty<WorkOrderUpliftExposureData>();
var primaryPairs = await _context.Dispatches
.Where(d => (d.IsDeleted == null || d.IsDeleted == false)
&& d.WorkOrderId != null
&& workOrderIdsScope.Contains(d.WorkOrderId.Value))
.Select(d => new { DispatchId = d.Id, WorkOrderId = d.WorkOrderId!.Value })
.ToListAsync(cancellationToken);
var linkedPairs = await _context.DispatchWorkOrders
.Where(l => workOrderIdsScope.Contains(l.WorkOrderId)
&& l.Dispatch != null
&& (l.Dispatch.IsDeleted == null || l.Dispatch.IsDeleted == false))
.Select(l => new { l.DispatchId, l.WorkOrderId })
.ToListAsync(cancellationToken);
var workOrdersByDispatch = primaryPairs
.Concat(linkedPairs)
.GroupBy(p => p.DispatchId)
.ToDictionary(g => g.Key, g => g.Select(p => p.WorkOrderId).ToHashSet());
if (workOrdersByDispatch.Count == 0)
return Array.Empty<WorkOrderUpliftExposureData>();
var dispatchIds = workOrdersByDispatch.Keys.ToList();
var sums = await _context.DispatchUpliftRequests
.Where(u => (u.IsDeleted == null || u.IsDeleted == false)
&& dispatchIds.Contains(u.DispatchId)
&& (u.Status == "Approved" || u.Status == "NoApprovalRequired"))
.GroupBy(u => u.DispatchId)
.Select(g => new
{
DispatchId = g.Key,
// The two creation paths store different meanings in RequestedNTE.
// Vendor-portal rows store the requested new NTE total, so the
// granted amount is RequestedNTE - CurrentNTE; work-order-path rows
// store the granted increment directly. Vendor sessions have no
// identity user, so createdby is null only on vendor-portal rows. Summing
// granted amounts keeps sequential approvals from double-counting whole
// NTE totals. This applies to both buckets.
AutoApproved = g.Where(x => x.Status == "NoApprovalRequired")
.Sum(x => (decimal?)(x.createdby == null
? x.RequestedNTE - (x.CurrentNTE ?? 0m)
: x.RequestedNTE)),
AdminApproved = g.Where(x => x.Status == "Approved")
.Sum(x => (decimal?)(x.createdby == null
? x.RequestedNTE - (x.CurrentNTE ?? 0m)
: x.RequestedNTE))
})
.ToListAsync(cancellationToken);
var totals = new Dictionary<int, WorkOrderUpliftExposureData>();
foreach (var sum in sums)
{
if (!workOrdersByDispatch.TryGetValue(sum.DispatchId, out var linkedWorkOrders))
continue;
foreach (var workOrderId in linkedWorkOrders)
{
if (!totals.TryGetValue(workOrderId, out var total))
{
total = new WorkOrderUpliftExposureData { WorkOrderId = workOrderId };
totals[workOrderId] = total;
}
total.AutoApprovedTotal += sum.AutoApproved ?? 0m;
total.AdminApprovedTotal += sum.AdminApproved ?? 0m;
}
}
return totals.Values.ToList();
}
// Queue-wide pending exposure for the approvals header: the sum of the Delta each
// Pending queue row displays (RequestedNTE - CurrentNTE), over the same rows the
// Pending tab lists (non-deleted request on a non-deleted dispatch), independent
// of page or filters.
public async Task<decimal> GetPendingExposureTotalAsync(CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.Where(u => (u.IsDeleted == null || u.IsDeleted == false)
&& u.Status == "Pending"
&& u.Dispatch != null
&& (u.Dispatch.IsDeleted == null || u.Dispatch.IsDeleted == false))
.SumAsync(u => (decimal?)(u.RequestedNTE - (u.CurrentNTE ?? 0m)), cancellationToken) ?? 0m;
}
public Task<List<DispatchUpliftRequest>> GetPendingForWorkOrderAsync(
int workOrderId,
CancellationToken cancellationToken)
{
return ForWorkOrder(workOrderId)
.Include(u => u.Dispatch)
.Where(u => u.Status == "Pending" || u.Status == "ChangesRequested")
.ToListAsync(cancellationToken);
}
public async Task<bool> HasActiveAsync(int dispatchId, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.AnyAsync(u => u.DispatchId == dispatchId
&& (u.Status == "Pending" || u.Status == "ChangesRequested")
&& (u.IsDeleted == null || u.IsDeleted == false), cancellationToken);
}
public async Task<DispatchUpliftRequest?> GetActiveRequestAsync(int dispatchId, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.FirstOrDefaultAsync(u => u.DispatchId == dispatchId
&& (u.Status == "Pending" || u.Status == "ChangesRequested")
&& (u.IsDeleted == null || u.IsDeleted == false), cancellationToken);
}
public async Task<DispatchUpliftRequest?> GetByRequestKeyAsync(int dispatchId, string requestKey, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.FirstOrDefaultAsync(u => u.DispatchId == dispatchId
&& u.RequestKey == requestKey
&& (u.IsDeleted == null || u.IsDeleted == false), cancellationToken);
}
public async Task<List<DispatchUpliftRequest>> GetDueForExpiryAsync(DateTime utcNow, int count, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.Where(u => (u.Status == "Pending" || u.Status == "ChangesRequested")
&& u.ExpiresAt != null && u.ExpiresAt <= utcNow
&& (u.IsDeleted == null || u.IsDeleted == false))
.OrderBy(u => u.ExpiresAt)
.Take(count)
.ToListAsync(cancellationToken);
}
public async Task<List<DispatchUpliftRequest>> GetDueForInitialNotificationAsync(int count, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.Where(u => (u.Status == "Pending" || u.Status == "ChangesRequested")
&& u.InitialNotificationSentAt == null
&& (u.IsDeleted == null || u.IsDeleted == false))
.OrderBy(u => u.CreatedDate)
.Take(count)
.ToListAsync(cancellationToken);
}
public async Task<List<DispatchUpliftRequest>> GetDueForEscalationAsync(int count, CancellationToken cancellationToken)
{
return await _context.DispatchUpliftRequests
.Where(u => (u.Status == "Pending" || u.Status == "ChangesRequested")
&& u.InitialNotificationSentAt != null
&& u.EscalatedAt == null
&& (u.IsDeleted == null || u.IsDeleted == false))
.OrderBy(u => u.InitialNotificationSentAt)
.Take(count)
.ToListAsync(cancellationToken);
}
private IQueryable<DispatchUpliftRequest> ForWorkOrder(int workOrderId)
{
return _context.DispatchUpliftRequests.Where(u =>
(u.IsDeleted == null || u.IsDeleted == false)
&& u.Dispatch != null
&& (u.Dispatch.IsDeleted == null || u.Dispatch.IsDeleted == false)
&& (
u.Dispatch.WorkOrderId == workOrderId
|| u.Dispatch.DispatchWorkOrders!.Any(link => link.WorkOrderId == workOrderId)));
}
public async Task StageAsync(DispatchUpliftRequest request, CancellationToken cancellationToken)
{
await _context.DispatchUpliftRequests.AddAsync(request, cancellationToken);
}
public async Task SaveChangesAsync(CancellationToken cancellationToken)
{
try
{
await _context.SaveChangesAsync(cancellationToken);
}
catch (DbUpdateException exception) when (IsActiveDispatchUniqueIndexViolation(exception))
{
throw new UpliftDispatchConflictException();
}
}
private bool IsActiveDispatchUniqueIndexViolation(DbUpdateException exception)
{
var provider = _context.Database.ProviderName;
var databaseException = exception.GetBaseException();
var exceptionType = databaseException.GetType();
if (provider?.Contains("SqlServer", StringComparison.OrdinalIgnoreCase) == true
&& exceptionType.FullName == "Microsoft.Data.SqlClient.SqlException")
{
var number = exceptionType.GetProperty("Number")?.GetValue(databaseException) as int?;
return number is 2601 or 2627
&& IsActiveDispatchIndexViolationMessage(databaseException.Message);
}
if (provider?.Contains("Sqlite", StringComparison.OrdinalIgnoreCase) == true
&& exceptionType.FullName == "Microsoft.Data.Sqlite.SqliteException")
{
var extendedCode = exceptionType.GetProperty("SqliteExtendedErrorCode")?.GetValue(databaseException) as int?;
const string uniqueConstraintPrefix = "UNIQUE constraint failed: ";
var message = databaseException.Message;
var prefixIndex = message.IndexOf(uniqueConstraintPrefix, StringComparison.OrdinalIgnoreCase);
if (extendedCode != 2067 || prefixIndex < 0)
return false;
var columnsStart = prefixIndex + uniqueConstraintPrefix.Length;
var columnsEnd = message.IndexOf('\'', columnsStart);
if (columnsEnd < 0)
columnsEnd = message.Length;
var columns = message[columnsStart..columnsEnd].Trim();
return string.Equals(
columns,
"DispatchUpliftRequests.DispatchId",
StringComparison.OrdinalIgnoreCase);
}
return false;
}
internal static bool IsActiveDispatchIndexViolationMessage(string message)
{
const string activeDispatchIndex = "'IX_DispatchUpliftRequests_DispatchId'";
return message.Contains(activeDispatchIndex, StringComparison.OrdinalIgnoreCase);
}
public async Task<T> ExecuteWorkOrderMutationAsync<T>(
int workOrderId,
Func<CancellationToken, Task<T>> work,
CancellationToken cancellationToken)
{
var gate = WorkOrderGates.GetOrAdd(workOrderId, _ => new SemaphoreSlim(1, 1));
await gate.WaitAsync(cancellationToken);
try
{
await using var transaction = _context.Database.IsRelational()
? await _context.Database.BeginTransactionAsync(cancellationToken)
: null;
try
{
await LockWorkOrderRowAsync(workOrderId, cancellationToken);
var result = await work(cancellationToken);
if (transaction is not null)
await transaction.CommitAsync(cancellationToken);
return result;
}
catch
{
if (transaction is not null)
await transaction.RollbackAsync(cancellationToken);
throw;
}
}
finally
{
gate.Release();
}
}
private async Task LockWorkOrderRowAsync(int workOrderId, CancellationToken cancellationToken)
{
if (_context.Database.ProviderName?.Contains("SqlServer", StringComparison.OrdinalIgnoreCase) != true)
return;
await _context.workOrders
.FromSqlRaw(
"SELECT * FROM [workOrders] WITH (UPDLOCK, ROWLOCK, HOLDLOCK) WHERE [Id] = {0}",
workOrderId)
.Select(workOrder => workOrder.Id)
.FirstOrDefaultAsync(cancellationToken);
}
}
}