using Microsoft.EntityFrameworkCore; using System.Text.Json.Nodes; namespace ZymonicServices; public class ZymonicGatewayQueueStore : IZymonicGatewayQueueStore { private readonly IZymonicDbContextFactory _contextFactory; public ZymonicGatewayQueueStore(IZymonicDbContextFactory contextFactory) { _contextFactory = contextFactory; } public void InitialiseDatabase() { using var context = _contextFactory.CreateDbContext(); context.StartupTasks(); } public async Task RequeueInProgressRequestsAsync(CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); var requests = await context.Set() .Where(request => request.Status == ZymonicGatewayQueueStatus.InProgress) .ToListAsync(ct); foreach (var request in requests) { request.Status = ZymonicGatewayQueueStatus.Queued; request.StartedAtUtc = null; request.CompletedAtUtc = null; request.LastAttemptedAtUtc = null; request.Error = "Requeued after gateway startup; previous processing did not complete."; } if (requests.Count > 0) { await context.SaveChangesAsync(ct); } return requests.Count; } public async Task QueueRequestAsync(ZymonicGatewayQueuedRequest request, CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); context.Set().Add(request); await context.SaveChangesAsync(ct); return request.RequestId; } public async Task QueueFailedRequestAsync(ZymonicGatewayQueuedRequest request, string error, CancellationToken ct = default) { request.Status = ZymonicGatewayQueueStatus.Failed; request.Error = error; request.CompletedAtUtc = DateTime.UtcNow; request.LastAttemptedAtUtc = DateTime.UtcNow; using var context = _contextFactory.CreateDbContext(); context.Set().Add(request); await context.SaveChangesAsync(ct); return request.RequestId; } public async Task GetQueuedRequestAsync(Guid requestId, CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); return await context.Set() .AsNoTracking() .SingleOrDefaultAsync(request => request.RequestId == requestId, ct); } public async Task> GetQueuedRequestsAsync(string? status = null, int limit = 100, CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); var query = context.Set() .AsNoTracking() .AsQueryable(); if (!string.IsNullOrWhiteSpace(status)) { query = query.Where(request => request.Status == status); } return await query .OrderByDescending(request => request.LastAttemptedAtUtc ?? request.QueuedAtUtc) .Take(Math.Clamp(limit, 1, 1000)) .ToListAsync(ct); } public async Task ClaimNextQueuedRequestAsync(int maxAttempts, TimeSpan retryDelay, CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); var retryBefore = DateTime.UtcNow - retryDelay; var request = await context.Set() .Where(request => request.Status == ZymonicGatewayQueueStatus.Queued && request.AttemptCount < maxAttempts && (request.LastAttemptedAtUtc == null || request.LastAttemptedAtUtc <= retryBefore)) .OrderBy(request => request.QueuedAtUtc) .FirstOrDefaultAsync(ct); if (request is null) { return null; } request.Status = ZymonicGatewayQueueStatus.InProgress; request.StartedAtUtc = DateTime.UtcNow; request.LastAttemptedAtUtc = DateTime.UtcNow; request.AttemptCount++; await context.SaveChangesAsync(ct); return request; } public async Task CompleteRequestAsync(Guid requestId, CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); var request = await context.Set().FindAsync([requestId], ct); if (request is null) { return; } request.Status = ZymonicGatewayQueueStatus.Completed; request.CompletedAtUtc = DateTime.UtcNow; request.Error = null; await context.SaveChangesAsync(ct); } public async Task UpdateQueuedRequestAsync(Guid requestId, ZymonicGatewayQueuedRequest update, CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); var request = await context.Set().FindAsync([requestId], ct); if (request is null) { return false; } request.Source = update.Source; request.Route = update.Route; request.EffectiveUser = update.EffectiveUser; request.AuthMethod = update.AuthMethod; request.RequestPayloadJson = update.RequestPayloadJson; request.RequestMetadataJson = update.RequestMetadataJson; await context.SaveChangesAsync(ct); return true; } public async Task RetryRequestAsync(Guid requestId, int maxAttempts, bool debugMode, CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); var request = await context.Set().FindAsync([requestId], ct); if (request is null) { return false; } request.Status = ZymonicGatewayQueueStatus.Queued; request.QueuedAtUtc = DateTime.UtcNow; request.StartedAtUtc = null; request.CompletedAtUtc = null; request.LastAttemptedAtUtc = null; request.AttemptCount = Math.Max(0, maxAttempts - 1); request.RequestMetadataJson = WithManualRetryMetadata(request.RequestMetadataJson, debugMode); request.Error = null; await context.SaveChangesAsync(ct); return true; } private static string WithManualRetryMetadata(string? metadataJson, bool debugMode) { JsonObject metadata; try { metadata = string.IsNullOrWhiteSpace(metadataJson) ? new JsonObject() : JsonNode.Parse(metadataJson) as JsonObject ?? new JsonObject(); } catch { metadata = new JsonObject(); } metadata["manualRetry"] = true; metadata["manualRetryDebugMode"] = debugMode; metadata["manualRetryQueuedAtUtc"] = DateTime.UtcNow.ToString("O"); return metadata.ToJsonString(); } public async Task FailRequestAsync(Guid requestId, string error, int maxAttempts, CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); var request = await context.Set().FindAsync([requestId], ct); if (request is null) { return; } request.Status = request.AttemptCount >= maxAttempts ? ZymonicGatewayQueueStatus.Failed : ZymonicGatewayQueueStatus.Queued; request.Error = error; await context.SaveChangesAsync(ct); } public async Task GetStatusSnapshotAsync(CancellationToken ct = default) { using var context = _contextFactory.CreateDbContext(); var now = DateTime.UtcNow; var queuedRequests = await context.Set() .Where(request => request.Status == ZymonicGatewayQueueStatus.Queued) .Select(request => request.QueuedAtUtc) .ToListAsync(ct); return new ZymonicGatewayQueueStatusSnapshot { QueuedRequestCount = queuedRequests.Count, QueuedAgeBuckets = new Dictionary { ["under_1_minute"] = queuedRequests.Count(queuedAt => now - queuedAt < TimeSpan.FromMinutes(1)), ["1_to_5_minutes"] = queuedRequests.Count(queuedAt => now - queuedAt >= TimeSpan.FromMinutes(1) && now - queuedAt < TimeSpan.FromMinutes(5)), ["5_to_30_minutes"] = queuedRequests.Count(queuedAt => now - queuedAt >= TimeSpan.FromMinutes(5) && now - queuedAt < TimeSpan.FromMinutes(30)), ["30_to_120_minutes"] = queuedRequests.Count(queuedAt => now - queuedAt >= TimeSpan.FromMinutes(30) && now - queuedAt < TimeSpan.FromHours(2)), ["over_2_hours"] = queuedRequests.Count(queuedAt => now - queuedAt >= TimeSpan.FromHours(2)) } }; } }