using System; using System.Collections.Generic; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Qualifier; using GB5Shared.QueryExecutor; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; namespace GB5Shared.GOP.Qualifier { // ============================================================ // BulkQualifierRunBLL — "test this Qualifier against a batch of records" runner. // // Reuses the real Qualifier engine (IQualifierFacadeImpl.ExecuteStageAsync) directly // per record rather than going through QualifierFacade's PreviewAsync (which discards // everything except Error-severity findings and has no way to select a single // QualifierId out of the full entity+stage+scope DAG). Findings from the whole bound // DAG are filtered down to the one QualifierId this run targets -- correct data-wise // (findings are stamped per-qualifier by the engine already) even though sibling // qualifiers in the same DAG still execute, since the target qualifier may genuinely // depend on FactBag output one of them produces. // ============================================================ public class BulkQualifierRunBLL : IBulkQualifierRunBLL { private readonly IQualifierBulkRunDAL _Dal; private readonly IQualifierDAL _QualifierDal; private readonly IQualifierFacadeImpl _Engine; private readonly ILogger _Logger; public BulkQualifierRunBLL( IQualifierBulkRunDAL dal, IQualifierDAL qualifierDal, IQualifierFacadeImpl engine, ILogger logger) { _Dal = dal; _QualifierDal = qualifierDal; _Engine = engine; _Logger = logger; } public async Task EnqueueRunAsync( int qualifierId, string entityCode, byte stage, byte scope, IReadOnlyList records, LoginDTO login, CancellationToken ct = default) { if (records.Count == 0) throw new ArgumentException("At least one record is required."); GB5Trace.Step("enqueue-qualifier-bulkrun", new { qualifierId, entityCode, stage, scope, recordCount = records.Count }); // QUALIFIERBULKRUNID reads as a plain INT in GOP_DDL.sql's own CREATE TABLE text, // but is a live IDENTITY column on GB5DEMO (verified via sys.columns) — DB-generated, // no AutoNumber involved. var runId = await _Dal.InsertRunAsync(new QualifierBulkRunDTO { QualifierId = qualifierId, EntityCode = entityCode, Stage = stage, Scope = scope, TotalRecords = records.Count, CreatedById = login.UserId }, login, ct); return runId; } public async Task ExecuteBulkRunAsync( int qualifierBulkRunId, int qualifierId, string entityCode, byte stage, byte scope, IReadOnlyList records, LoginDTO login, CancellationToken ct) { await _Dal.MarkRunningAsync(qualifierBulkRunId, login, ct).ConfigureAwait(false); var qualifierDef = await _QualifierDal.GetQualifierById(qualifierId, login, ct).ConfigureAwait(false); var qualifierCode = qualifierDef?.QualifierCode ?? qualifierId.ToString(); var processed = 0; var errorCount = 0; var warningCount = 0; var findingsBuffer = new List(); try { foreach (var record in records) { ct.ThrowIfCancellationRequested(); JsonDocument document; try { document = JsonDocument.Parse(string.IsNullOrWhiteSpace(record.DataJson) ? "{}" : record.DataJson); } catch (JsonException ex) { findingsBuffer.Add(new QualifierBulkFindingDTO { QualifierBulkRunId = qualifierBulkRunId, QualifierId = qualifierId, QualifierCode = qualifierCode, RecordKey = record.RecordKey, Severity = (byte)FindingSeverity.Error, FindingCode = "INVALID_JSON", FindingMessage = $"Record's DataJson is not valid JSON: {ex.Message}", }); errorCount++; processed++; continue; } using (document) { var qCtx = new QualifierExecutionContext { ClientId = login.ClientId, EntityId = 0, EntityCode = entityCode, Stage = stage, Scope = scope, Document = document, Facts = new QualifierFactBag(), // fresh per record — records are independent CorrelationId = record.RecordKey, CancellationToken = ct }; QualifierResult result; try { result = await _Engine.ExecuteStageAsync(qCtx, login, ct).ConfigureAwait(false); } catch (Exception ex) { findingsBuffer.Add(new QualifierBulkFindingDTO { QualifierBulkRunId = qualifierBulkRunId, QualifierId = qualifierId, QualifierCode = qualifierCode, RecordKey = record.RecordKey, Severity = (byte)FindingSeverity.Error, FindingCode = "ENGINE_ERROR", FindingMessage = ex.Message, }); errorCount++; processed++; continue; } foreach (var finding in result.Findings.Where(f => f.QualifierId == qualifierId)) { findingsBuffer.Add(new QualifierBulkFindingDTO { QualifierBulkRunId = qualifierBulkRunId, QualifierId = finding.QualifierId, QualifierCode = finding.QualifierCode, RecordKey = record.RecordKey, Severity = (byte)finding.Severity, FindingCode = finding.Code, FindingMessage = finding.Message, DocumentPath = finding.DocumentPath, }); if (finding.Severity == FindingSeverity.Error) errorCount++; else if (finding.Severity == FindingSeverity.Warning) warningCount++; } } processed++; // Flush periodically so a long run's progress is visible before it finishes, // and so the findings buffer never grows unbounded for a very large batch. if (findingsBuffer.Count >= 500) { await _Dal.BulkInsertFindingsAsync(findingsBuffer, login, ct).ConfigureAwait(false); findingsBuffer.Clear(); await _Dal.UpdateProgressAsync(qualifierBulkRunId, processed, errorCount, warningCount, login, ct).ConfigureAwait(false); } } if (findingsBuffer.Count > 0) await _Dal.BulkInsertFindingsAsync(findingsBuffer, login, ct).ConfigureAwait(false); await _Dal.CompleteRunAsync(qualifierBulkRunId, runStatus: 2 /* Completed */, processed, errorCount, warningCount, errorMessage: null, login, ct).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("qualifier-bulkrun-failed", ex); _Logger.LogError(ex, "ExecuteBulkRunAsync failed for QualifierBulkRunId {QualifierBulkRunId}", qualifierBulkRunId); if (findingsBuffer.Count > 0) await _Dal.BulkInsertFindingsAsync(findingsBuffer, login, ct).ConfigureAwait(false); await _Dal.CompleteRunAsync(qualifierBulkRunId, runStatus: 3 /* Failed */, processed, errorCount, warningCount, errorMessage: ex.Message, login, CancellationToken.None).ConfigureAwait(false); } } public async Task> GetHistoryAsync( string entityCode, int page, int pageSize, LoginDTO login, CancellationToken ct = default) { var (offset, size) = NormalizePaging(page, pageSize); return await _Dal.GetHistoryAsync(entityCode, offset, size, login, ct).ConfigureAwait(false); } public async Task GetDetailAsync(int qualifierBulkRunId, LoginDTO login, CancellationToken ct = default) => await _Dal.GetByIdAsync(qualifierBulkRunId, login, ct).ConfigureAwait(false); public async Task> GetFindingsAsync( int qualifierBulkRunId, byte? severity, int page, int pageSize, LoginDTO login, CancellationToken ct = default) { var (offset, size) = NormalizePaging(page, pageSize); return await _Dal.GetFindingsAsync(qualifierBulkRunId, severity, offset, size, login, ct).ConfigureAwait(false); } private static (int Offset, int PageSize) NormalizePaging(int page, int pageSize) { var size = pageSize <= 0 ? 20 : Math.Min(pageSize, 200); var safePage = page <= 0 ? 1 : page; return ((safePage - 1) * size, size); } } }