using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using GB5Shared.Telemetry; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using System; using System.Collections.Generic; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; namespace FrameworkSL.Jobs { // Consumer/worker for TJOBQUEUE — the generic durable job queue defined for the // separate JobEngine microservice (C:\gb5\GB5Solution\JobEngine). That service's // write side (EnqueueMessage) is fully built but has no runnable worker anywhere // (ClaimBatchAsync has zero callers) and cannot even start (missing QuartzStore // connection string). This worker re-implements the same claim/complete/retry/DLQ // SQL semantics as JobEngineDAL.Query.JobQueueQB directly against the same table, // so TJOBQUEUE is genuinely processed end-to-end inside the app you already run. public sealed class JobQueueWorkerService : BackgroundService { private readonly IServiceScopeFactory _scopeFactory; private readonly IConfiguration _config; private readonly ILogger _logger; private const int BatchSize = 5; public JobQueueWorkerService( IServiceScopeFactory scopeFactory, IConfiguration config, ILogger logger) { _scopeFactory = scopeFactory; _config = config; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { int intervalSecs = _config.GetValue("JobQueueWorker:PollIntervalSeconds", 5); _logger.LogInformation("JobQueueWorkerService starting — polling TJOBQUEUE every {Seconds}s", intervalSecs); using var timer = new PeriodicTimer(TimeSpan.FromSeconds(intervalSecs)); try { while (await timer.WaitForNextTickAsync(stoppingToken).ConfigureAwait(false)) { try { await ProcessBatchAsync(stoppingToken).ConfigureAwait(false); } catch (Exception ex) when (ex is not OperationCanceledException) { _logger.LogError(ex, "JobQueueWorkerService tick failed"); } } } catch (OperationCanceledException) { _logger.LogInformation("JobQueueWorkerService stopping (shutdown requested)."); } } private async Task ProcessBatchAsync(CancellationToken ct) { using var scope = _scopeFactory.CreateScope(); var qe = scope.ServiceProvider.GetRequiredService(); var login = BuildSystemLogin(); var claimed = (await qe.QueryAsync(login, ClaimBatchSql, new { TenantId = login.ClientId, BatchSize }, cancellationToken: ct).ConfigureAwait(false)).ToList(); if (claimed.Count == 0) return; _logger.LogInformation("JobQueueWorker: claimed {Count} item(s) from TJOBQUEUE", claimed.Count); foreach (var item in claimed) await ProcessItemAsync(qe, login, item, ct).ConfigureAwait(false); } private async Task ProcessItemAsync(IQueryExecutor qe, LoginDTO login, ClaimedItem item, CancellationToken ct) { using var section = GB5Trace.BeginSection("jobqueue-process"); GB5Trace.Step("claimed", new { item.QueueId, item.MessageType, item.RetryAttempt, item.MaxRetries }); _logger.LogInformation( "JobQueueWorker [{QueueId}] Processing | Type={Type} Attempt={Attempt}/{Max} Payload={Payload}", item.QueueId, item.MessageType, item.RetryAttempt + 1, item.MaxRetries, item.Payload); bool shouldFail = false; try { using var doc = JsonDocument.Parse(item.Payload); if (doc.RootElement.TryGetProperty("forceFail", out var ff) && ff.ValueKind == JsonValueKind.True) { var failCount = doc.RootElement.TryGetProperty("failCount", out var fc) ? fc.GetInt32() : int.MaxValue; shouldFail = item.RetryAttempt < failCount; } } catch (JsonException) { // Non-JSON payload — treat as an always-succeed demo message. } if (!shouldFail) { await qe.ExecuteAsync(login, CompleteSql, new { item.QueueId, item.TenantId }, cancellationToken: ct).ConfigureAwait(false); _logger.LogInformation("JobQueueWorker [{QueueId}] COMPLETED", item.QueueId); GB5Trace.Step("completed", new { item.QueueId }); return; } var nextAttempt = item.RetryAttempt + 1; if (nextAttempt >= item.MaxRetries) { await qe.ExecuteAsync(login, DlqSql, new { item.QueueId, item.TenantId, ErrorDetails = "Simulated permanent failure (demo)" }, cancellationToken: ct).ConfigureAwait(false); _logger.LogWarning("JobQueueWorker [{QueueId}] MOVED TO DLQ after {Attempts} attempt(s)", item.QueueId, nextAttempt); GB5Trace.MarkFailed("moved-to-dlq"); return; } var backoffSeconds = 8 * nextAttempt; await qe.ExecuteAsync(login, RetrySql, new { item.QueueId, item.TenantId, NextAttemptAfter = DateTime.UtcNow.AddSeconds(backoffSeconds), ErrorDetails = "Simulated transient failure (demo)" }, cancellationToken: ct).ConfigureAwait(false); _logger.LogWarning( "JobQueueWorker [{QueueId}] TRANSIENT FAILURE — retry {Attempt}/{Max} in {Backoff}s", item.QueueId, nextAttempt, item.MaxRetries, backoffSeconds); GB5Trace.Step("retry-scheduled", new { item.QueueId, nextAttempt, backoffSeconds }); } private LoginDTO BuildSystemLogin() => new() { ClientId = _config.GetValue("JobQueueWorker:TenantId", -1499999887), DatabaseName = _config["JobQueueWorker:ConnectionName"] ?? "UNISOFTGB5", UserId = -1 }; private sealed class ClaimedItem { public long QueueId { get; set; } public int TenantId { get; set; } public string MessageType { get; set; } = ""; public string Payload { get; set; } = ""; public byte Priority { get; set; } public int MaxRetries { get; set; } public int RetryAttempt { get; set; } public string MessageId { get; set; } = ""; public string? CorrelationId { get; set; } public string? SourceService { get; set; } } // Same claim/complete/retry/DLQ semantics as JobEngineDAL.Query.JobQueueQB — // re-implemented here because FrameworkSL doesn't reference the separate // GB5Solution/JobEngine projects. private const string ClaimBatchSql = @" UPDATE TOP(@BatchSize) Q SET STATUS = 'PROCESSING', PROCESSINGSTARTED = GETUTCDATE(), UPDATEDDATE = GETUTCDATE() OUTPUT INSERTED.QUEUEID AS QueueId, INSERTED.TENANTID AS TenantId, INSERTED.MESSAGETYPE AS MessageType, INSERTED.PAYLOAD AS Payload, INSERTED.PRIORITY AS Priority, INSERTED.MAXRETRIES AS MaxRetries, INSERTED.RETRYATTEMPT AS RetryAttempt, INSERTED.MESSAGEID AS MessageId, INSERTED.CORRELATIONID AS CorrelationId, INSERTED.SOURCESERVICE AS SourceService FROM TJOBQUEUE Q WITH (ROWLOCK, READPAST) WHERE Q.TENANTID = @TenantId AND Q.STATUS IN ('PENDING','FAILED') AND Q.SCHEDULEDFOR <= GETUTCDATE() AND (Q.NEXTATTEMPTAFTER IS NULL OR Q.NEXTATTEMPTAFTER <= GETUTCDATE())"; private const string CompleteSql = @" UPDATE TJOBQUEUE SET STATUS = 'COMPLETED', UPDATEDDATE = GETUTCDATE() WHERE QUEUEID = @QueueId AND TENANTID = @TenantId"; private const string RetrySql = @" UPDATE TJOBQUEUE SET STATUS = 'FAILED', RETRYATTEMPT = RETRYATTEMPT + 1, NEXTATTEMPTAFTER = @NextAttemptAfter, ERRORDETAILS = @ErrorDetails, UPDATEDDATE = GETUTCDATE() WHERE QUEUEID = @QueueId AND TENANTID = @TenantId"; private const string DlqSql = @" UPDATE TJOBQUEUE SET STATUS = 'DLQ', ERRORDETAILS = @ErrorDetails, UPDATEDDATE = GETUTCDATE() WHERE QUEUEID = @QueueId AND TENANTID = @TenantId"; } }