using System.Diagnostics;
using System.Text;
using System.Text.Json;
using GB5Shared.DTO.Framework.Login;
using GB5Shared.DTO.ReportOrchestration;
using GB5Shared.Enums.ReportOrchestration;
using GB5Shared.QueryExecutor;
using FrameworkBLL.SystemJob;
using FrameworkDAL.DTO.SystemJob;
using FrameworkDAL.Query.SystemJob;
using Microsoft.Extensions.Caching.Distributed;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.Logging;
namespace FrameworkBLL.ReportOrchestration
{
///
/// Central 6-stage orchestration pipeline for all report and analysis execution.
/// Stage 1: Load MREPORTCONFIG
/// Stage 2: Validate ExportFormat
/// Stage 3: Resolve DataSourceType (Auto → concrete)
/// Stage 4: Resolve IDbConnection
/// Stage 5: Evaluate async/sync decision
/// Stage 6a (sync): execute → optional pivot → export → return data/stream
/// Stage 6b (async): serialize context → submit SysJob → return job ID
/// Never throws — all exceptions are caught and wrapped as Failed results.
///
public sealed class QueryOrchestrator : IQueryOrchestrator
{
private readonly IQueryExecutor _queryExecutor;
private readonly IReportConnectionResolver _connectionResolver;
private readonly IDataSourceRuleEvaluator _ruleEvaluator;
private readonly IAsyncDecisionEngine _asyncDecision;
private readonly IReportEndpointRegistry _registry;
private readonly ISystemJobBLL _systemJobBLL;
private readonly ILogger _logger;
private readonly IConfiguration _configuration;
private readonly IDistributedCache _resultCache;
public QueryOrchestrator(
IQueryExecutor queryExecutor,
IReportConnectionResolver connectionResolver,
IDataSourceRuleEvaluator ruleEvaluator,
IAsyncDecisionEngine asyncDecision,
IReportEndpointRegistry registry,
ISystemJobBLL systemJobBLL,
ILogger logger,
IConfiguration configuration,
IDistributedCache resultCache)
{
_queryExecutor = queryExecutor;
_connectionResolver = connectionResolver;
_ruleEvaluator = ruleEvaluator;
_asyncDecision = asyncDecision;
_registry = registry;
_systemJobBLL = systemJobBLL;
_logger = logger;
_configuration = configuration;
_resultCache = resultCache;
}
// =====================================================================
// PRIMARY ENTRY POINT
// =====================================================================
public async Task ExecuteReportAsync(ReportExecutionContext context)
{
var sw = Stopwatch.StartNew();
try
{
// ── Stage 1: Load MREPORTCONFIG ──────────────────────────────
await LoadReportConfigAsync(context).ConfigureAwait(false);
// ── Stage 2: Validate requested format ───────────────────────
ValidateExportFormat(context);
// ── Stage 3: Resolve DataSourceType ──────────────────────────
var resolvedSource = await _ruleEvaluator
.EvaluateAsync(context, context.CancellationToken).ConfigureAwait(false);
context.ResolvedDataSource = resolvedSource;
// ── Stage 4: Resolve connection config ───────────────────────
var serverConfig = await _connectionResolver
.ResolveConfigAsync(context.LoginDTO, resolvedSource, context.CancellationToken)
.ConfigureAwait(false);
context.ResolvedServerConfig = serverConfig;
// Apply NOLOCK and timeout from config
if (context.ReportConfig is not null)
{
context.ApplyNoLock = context.ReportConfig.IsNoLockAllowed
&& serverConfig.ProviderType == DatabaseProviderType.SqlServer
&& serverConfig.IsDedicatedAnalyticsDb;
context.QueryTimeoutSeconds = context.ReportConfig.QueryTimeoutSeconds;
}
// ── Stage 5: Async decision ───────────────────────────────────
bool runAsync = await _asyncDecision
.ShouldRunAsyncAsync(context, context.CancellationToken).ConfigureAwait(false);
context.IsAsync = runAsync;
// ── Stage 6b: Async path ──────────────────────────────────────
if (runAsync)
return await SubmitAsyncJobAsync(context, sw).ConfigureAwait(false);
// ── Stage 6a: Sync path ───────────────────────────────────────
return await ExecuteSyncAsync(context, sw).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
_logger.LogInformation(
"Report execution cancelled. CorrelationId={Id}", context.CorrelationId);
return ReportExecutionResult.Failure(
"Report execution was cancelled.",
"CANCELLED",
context.CorrelationId,
sw.Elapsed);
}
catch (Exception ex)
{
_logger.LogError(ex,
"Report execution failed. CorrelationId={Id} ReportId={ReportId}",
context.CorrelationId, context.ReportId);
return ReportExecutionResult.Failure(
"An error occurred while generating the report. Please try again or contact support.",
"EXECUTION_ERROR",
context.CorrelationId,
sw.Elapsed,
ex);
}
}
// =====================================================================
// ASYNC WORKER ENTRY POINT
// =====================================================================
public async Task ExecuteAsyncJobAsync(
int sysJobId,
LoginDTO loginDTO,
CancellationToken cancellationToken = default)
{
_logger.LogInformation(
"ExecuteAsyncJob: Starting SysJobId={Id}", sysJobId);
// ── Step 1: Load the TSYSJOB row ──────────────────────────────────
var jobs = await _queryExecutor
.QueryAsync(
loginDTO,
SystemJobQB.GET_SYSTEMJOB,
new { SysJobId = sysJobId, SysJobTenantId = loginDTO.ClientId },
cancellationToken: cancellationToken)
.ConfigureAwait(false);
var job = jobs.FirstOrDefault()
?? throw new InvalidOperationException(
$"SysJob {sysJobId} not found for tenant {loginDTO.ClientId}.");
// ── Step 2: Claim the job atomically (Pending=0 → InProgress=1) ──
int claimed = await _queryExecutor
.ExecuteAsync(
loginDTO,
SystemJobQB.CLAIM_JOB,
new { SysJobId = sysJobId, SysJobTenantId = loginDTO.ClientId,
SysJobModifiedById = loginDTO.UserId },
cancellationToken: cancellationToken)
.ConfigureAwait(false);
if (claimed == 0)
throw new InvalidOperationException(
$"SysJob {sysJobId} could not be claimed — it may already be running or cancelled.");
string resultLocation;
try
{
// ── Step 3: Reconstruct context from stored Parameters JSON ────
if (string.IsNullOrWhiteSpace(job.Parameters))
throw new InvalidOperationException(
$"SysJob {sysJobId} has no Parameters JSON to reconstruct context from.");
var stored = Newtonsoft.Json.JsonConvert
.DeserializeObject(job.Parameters)
?? throw new InvalidOperationException(
$"Failed to deserialize Parameters JSON for SysJob {sysJobId}.");
var context = new ReportExecutionContext
{
LoginDTO = loginDTO,
ReportId = stored.ReportId,
AnalysisId = stored.AnalysisId,
ReportViewId = stored.ReportViewId,
CallType = (ReportCallType)stored.CallType,
Parameters = stored.Parameters ?? new Dictionary(),
RequestedFormat = (ExportFormat)stored.RequestedFormat,
CorrelationId = stored.CorrelationId,
CancellationToken = cancellationToken
};
// ── Step 4: Run pipeline stages 1–4 ─────────────────────────
await LoadReportConfigAsync(context).ConfigureAwait(false);
ValidateExportFormat(context);
var resolvedSource = await _ruleEvaluator
.EvaluateAsync(context, cancellationToken).ConfigureAwait(false);
context.ResolvedDataSource = resolvedSource;
var serverConfig = await _connectionResolver
.ResolveConfigAsync(loginDTO, resolvedSource, cancellationToken)
.ConfigureAwait(false);
context.ResolvedServerConfig = serverConfig;
if (context.ReportConfig is not null)
{
context.ApplyNoLock = context.ReportConfig.IsNoLockAllowed
&& serverConfig.ProviderType == DatabaseProviderType.SqlServer
&& serverConfig.IsDedicatedAnalyticsDb;
// Use the longer async timeout for worker execution
context.QueryTimeoutSeconds = context.ReportConfig.AsyncQueryTimeoutSeconds;
}
// Stage 5 skipped — always run sync inside a worker
context.IsAsync = false;
// ── Step 5: Execute sync pipeline (stage 6a) ─────────────────
var sw = Stopwatch.StartNew();
var result = await ExecuteSyncAsync(context, sw).ConfigureAwait(false);
if (result.ResultType == ReportResultType.Failed)
throw new InvalidOperationException(
result.ErrorMessage ?? "Report execution returned a failure result.");
// ── Step 6: Store result in distributed cache ─────────────────
string ext = context.RequestedFormat switch
{
ExportFormat.Excel => ".xlsx",
ExportFormat.Pdf => ".pdf",
ExportFormat.Csv => ".csv",
_ => ".json"
};
string redisKey = $"report:{loginDTO.ClientId}:{sysJobId}{ext}";
resultLocation = $"redis:{redisKey}";
var cacheOptions = new DistributedCacheEntryOptions
{
AbsoluteExpirationRelativeToNow = TimeSpan.FromHours(24)
};
if (result.ExportStream is not null)
{
result.ExportStream.Seek(0, SeekOrigin.Begin);
using var ms = new System.IO.MemoryStream();
await result.ExportStream.CopyToAsync(ms, cancellationToken).ConfigureAwait(false);
await _resultCache.SetAsync(redisKey, ms.ToArray(), cacheOptions, cancellationToken)
.ConfigureAwait(false);
}
else
{
// Grid result — serialise as JSON so GetSystemJobResult can stream it
string json = JsonSerializer.Serialize(new
{
Data = result.GridData,
TotalRows = result.TotalRows,
DataSource = result.DataSourceUsed.ToString()
});
await _resultCache.SetStringAsync(redisKey, json, cacheOptions, cancellationToken)
.ConfigureAwait(false);
}
// ── Step 7: Mark job Completed (RUNSTATUS=2) ─────────────────
await _queryExecutor
.ExecuteAsync(
loginDTO,
SystemJobQB.COMPLETE_JOB,
new { SysJobId = sysJobId, SysJobTenantId = loginDTO.ClientId,
ResultLocation = resultLocation,
SysJobModifiedById = loginDTO.UserId },
cancellationToken: CancellationToken.None)
.ConfigureAwait(false);
_logger.LogInformation(
"ExecuteAsyncJob completed. SysJobId={Id} ResultLocation={Path}",
sysJobId, resultLocation);
}
catch (Exception ex)
{
_logger.LogError(ex,
"ExecuteAsyncJob failed. SysJobId={Id}", sysJobId);
// Mark Failed (RUNSTATUS=3) or reset to Pending (RUNSTATUS=0) for retry
int runStatus = (job.RetryCount < job.MaxRetries) ? 0 : 3;
string errorMsg = ex.Message.Length > 2000
? ex.Message[..2000]
: ex.Message;
await _queryExecutor
.ExecuteAsync(
loginDTO,
SystemJobQB.FAIL_OR_RETRY_JOB,
new { SysJobId = sysJobId, SysJobTenantId = loginDTO.ClientId,
RunStatus = runStatus, ErrorMessage = errorMsg,
SysJobModifiedById = loginDTO.UserId },
cancellationToken: CancellationToken.None)
.ConfigureAwait(false);
throw;
}
return resultLocation;
}
// =====================================================================
// CONTEXT BUILDER FACTORY
// =====================================================================
public IReportContextBuilder ForStandardReport(
LoginDTO loginDTO,
int reportId,
int viewId = -1)
=> new ReportContextBuilder(
loginDTO, reportId, analysisId: -1, viewId, ReportCallType.StandardReport);
public IReportContextBuilder ForAnalysisReport(
LoginDTO loginDTO,
int analysisId,
int queryId = -1)
=> new ReportContextBuilder(
loginDTO, reportId: -1, analysisId, viewId: queryId, ReportCallType.AnalysisReport);
// =====================================================================
// PRIVATE PIPELINE HELPERS
// =====================================================================
private async Task LoadReportConfigAsync(ReportExecutionContext context)
{
if (context.ReportId == -1)
return; // Analysis reports — no MREPORTCONFIG
try
{
var config = await _queryExecutor
.QuerySingleAsync(
context.LoginDTO,
ReportOrchestrationQB.GET_REPORT_CONFIG,
new { ReportId = context.ReportId },
cancellationToken: context.CancellationToken)
.ConfigureAwait(false);
context.ReportConfig = config;
}
catch
{
// MREPORTCONFIG row not found — orchestrator uses safe defaults
_logger.LogDebug(
"No MREPORTCONFIG found for ReportId={Id}; using defaults.",
context.ReportId);
context.ReportConfig = new ReportConfigDTO
{
ReportId = context.ReportId,
AsyncMode = AsyncMode.AlwaysSync
};
}
}
private static void ValidateExportFormat(ReportExecutionContext context)
{
if (context.ReportConfig is null)
return;
var allowed = context.ReportConfig.AllowedExportFormats
.Split(',', StringSplitOptions.RemoveEmptyEntries)
.Select(s => byte.TryParse(s.Trim(), out var b) ? b : (byte)0)
.ToHashSet();
if (!allowed.Contains((byte)context.RequestedFormat))
throw new InvalidOperationException(
$"Export format '{context.RequestedFormat}' is not allowed for this report. " +
$"Allowed formats: {context.ReportConfig.AllowedExportFormats}");
}
// ── Sync execution ────────────────────────────────────────────────────
private async Task ExecuteSyncAsync(
ReportExecutionContext context,
Stopwatch sw)
{
// Look up report code for registry
string? reportCode = null;
if (context.ReportId != -1)
{
try
{
reportCode = await _queryExecutor
.QuerySingleAsync(
context.LoginDTO,
ReportOrchestrationQB.GET_REPORT_CODE,
new { ReportId = context.ReportId },
cancellationToken: context.CancellationToken)
.ConfigureAwait(false);
}
catch
{
throw new InvalidOperationException(
$"MREPORT row not found for ReportId={context.ReportId}.");
}
}
IReportEndpoint? endpoint = reportCode is not null
? _registry.Resolve(reportCode)
: null;
if (endpoint is null)
{
// No registered endpoint — could be an unregistered legacy report URI
// or an analysis report. Return a structured error.
throw new InvalidOperationException(
$"No registered IReportEndpoint found for ReportCode='{reportCode}'. " +
"Ensure the module assembly is scanned during startup.");
}
// Open the resolved connection
using var connection = await _connectionResolver
.ResolveConnectionAsync(
context.LoginDTO,
context.ResolvedDataSource,
context.CancellationToken)
.ConfigureAwait(false);
await ((System.Data.Common.DbConnection)connection)
.OpenAsync(context.CancellationToken).ConfigureAwait(false);
// Collect rows from the endpoint
var rows = new List>();
await foreach (var row in endpoint
.GetReportDataAsync(context, connection, context.CancellationToken)
.ConfigureAwait(false))
{
rows.Add(row);
}
sw.Stop();
return context.RequestedFormat switch
{
ExportFormat.Grid or ExportFormat.Json =>
ReportExecutionResult.SyncGrid(
rows, rows.Count,
context.ResolvedDataSource,
context.CorrelationId,
sw.Elapsed),
ExportFormat.Excel =>
await ExportExcelAsync(context, rows, sw.Elapsed).ConfigureAwait(false),
ExportFormat.Pdf =>
await ExportPdfAsync(context, rows, sw.Elapsed).ConfigureAwait(false),
ExportFormat.Csv =>
ExportCsv(context, rows, sw.Elapsed),
_ => throw new NotSupportedException(
$"Export format '{context.RequestedFormat}' is not handled by the orchestrator.")
};
}
// ── Export helpers ────────────────────────────────────────────────────
private static async Task ExportExcelAsync(
ReportExecutionContext context,
List> rows,
TimeSpan elapsed)
{
// Pivot path delegated to PivotExcelExport; flat path to IExcelExport
// Both are handled by CommonReportBLL for the existing report endpoints.
// The orchestrator exposes the raw grid; the calling endpoint applies
// the appropriate exporter using context.ReportConfig.PivotEngine.
// Full implementation is a follow-up task when module endpoints migrate.
await Task.CompletedTask.ConfigureAwait(false);
var csvBytes = BuildCsvBytes(rows);
var stream = new MemoryStream(csvBytes);
return ReportExecutionResult.SyncExport(
stream,
"application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
$"Report_{DateTime.UtcNow:yyyyMMdd_HHmmss}.xlsx",
ExportFormat.Excel,
context.ResolvedDataSource,
context.CorrelationId,
elapsed);
}
private static async Task ExportPdfAsync(
ReportExecutionContext context,
List> rows,
TimeSpan elapsed)
{
// PDF rendering delegates to existing PivotPdfExport / PdfExport
// when module endpoints are migrated. Placeholder for now.
await Task.CompletedTask.ConfigureAwait(false);
var csvBytes = BuildCsvBytes(rows);
var stream = new MemoryStream(csvBytes);
return ReportExecutionResult.SyncExport(
stream,
"application/pdf",
$"Report_{DateTime.UtcNow:yyyyMMdd_HHmmss}.pdf",
ExportFormat.Pdf,
context.ResolvedDataSource,
context.CorrelationId,
elapsed);
}
private static ReportExecutionResult ExportCsv(
ReportExecutionContext context,
List> rows,
TimeSpan elapsed)
{
var bytes = BuildCsvBytes(rows);
var stream = new MemoryStream(bytes);
return ReportExecutionResult.SyncExport(
stream,
"text/csv",
$"Report_{DateTime.UtcNow:yyyyMMdd_HHmmss}.csv",
ExportFormat.Csv,
context.ResolvedDataSource,
context.CorrelationId,
elapsed);
}
private static byte[] BuildCsvBytes(List> rows)
{
if (rows.Count == 0)
return Array.Empty();
var sb = new StringBuilder();
var headers = rows[0].Keys.ToList();
sb.AppendLine(string.Join(",", headers.Select(EscapeCsv)));
foreach (var row in rows)
{
sb.AppendLine(string.Join(",",
headers.Select(h => EscapeCsv(row.TryGetValue(h, out var v)
? v?.ToString() ?? string.Empty
: string.Empty))));
}
return Encoding.UTF8.GetBytes(sb.ToString());
}
private static string EscapeCsv(string? value)
{
if (value is null) return string.Empty;
if (value.Contains(',') || value.Contains('"') || value.Contains('\n'))
return $"\"{value.Replace("\"", "\"\"")}\"";
return value;
}
// Deserialization target for TSYSJOB.PARAMETERS stored by SubmitAsyncJobAsync.
private sealed class AsyncJobStoredParameters
{
public int ReportId { get; set; }
public int AnalysisId { get; set; }
public int ReportViewId { get; set; }
public byte CallType { get; set; }
public Dictionary Parameters { get; set; } = new();
public byte RequestedFormat { get; set; }
public Guid CorrelationId { get; set; }
}
// ── Async job submission ──────────────────────────────────────────────
private async Task SubmitAsyncJobAsync(
ReportExecutionContext context,
Stopwatch sw)
{
// Serialize the execution context for the SysJob worker
string parametersJson = JsonSerializer.Serialize(new
{
ReportId = context.ReportId,
AnalysisId = context.AnalysisId,
ReportViewId = context.ReportViewId,
CallType = (byte)context.CallType,
Parameters = context.Parameters,
RequestedFormat = (byte)context.RequestedFormat,
CorrelationId = context.CorrelationId
});
// Submit via existing SysJob framework
var jobDto = new SystemJobDTO
{
SysJobType = "REPORT",
Parameters = parametersJson,
SubmittedById = context.LoginDTO.UserId,
Priority = context.ReportConfig?.DefaultPriority ?? 5,
MaxRetries = 1,
ExpiryAt = DateTime.UtcNow.AddHours(
context.ReportConfig?.ResultExpiryHours ?? 24),
SysJobTenantId = context.LoginDTO.ClientId
};
string jobResult = await _systemJobBLL
.SaveSystemJob(jobDto, context.LoginDTO).ConfigureAwait(false);
sw.Stop();
// Parse the returned SysJobId from the BLL response string
// (BLL returns a JSON string; the SysJobId is included in the message)
_logger.LogInformation(
"Async report job submitted. CorrelationId={Id} JobResult={Result}",
context.CorrelationId, jobResult);
return ReportExecutionResult.AsyncQueued(
sysJobId: -1, // SysJobId extraction from jobResult is service-specific
context.ResolvedDataSource,
context.CorrelationId,
sw.Elapsed);
}
}
}