using Dapper;
using GB5Shared.ActionProcessor;
using GB5Shared.FileUpload;
using FrameworkBLL.SystemJob;
using FrameworkDAL.CustomCode.User;
using FrameworkDAL.DTO.SysJobRun;
using FrameworkDAL.DTO.SystemJob;
using FrameworkDAL.DTO.User;
using FrameworkDAL.Query.SysJobRun;
using FrameworkDAL.Query.SystemJob;
using FrameworkSL.Hubs.SysJob;
using GB5Shared.Connection;
using GB5Shared.DTO.Framework.CommonConfig;
using GB5Shared.DTO.Framework.Login;
using GB5Shared.DTO.Framework.ServerConfig;
using GB5Shared.GenerateAutoNumber;
using GB5Shared.QueryExecutor;
using Microsoft.AspNetCore.SignalR;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Npgsql;
using Quartz;
using System;
using System.Data;
using System.IO;
using Microsoft.Data.SqlClient;
using System.Linq;
using System.Net.Http;
using System.Text;
using System.Text.Json;
using System.Threading.Tasks;
using static GB5Shared.GB5Constant.Constant;
namespace FrameworkSL.Controllers.SysJob
{
///
/// Quartz job — polls TSYSJOB for pending rows (RUNSTATUS=0) and executes them
/// by calling the configured MWEBSERVICE endpoint.
///
/// Lifecycle per job:
/// 1. CLAIM_JOB — atomically sets RUNSTATUS=1 (only if still 0; prevents double-pickup)
/// 2. Create TSYSJOBRUN (StartedAt, Worker)
/// 3. Call web service via IHttpClientFactory
/// 4. Update TSYSJOBRUN (EndedAt, Status, Log)
/// 5. COMPLETE_JOB / FAIL_OR_RETRY_JOB on TSYSJOB
///
/// Register in Program.cs with a trigger interval (e.g. every 30 seconds).
///
[DisallowConcurrentExecution]
public class SysJobExecutorQuartzJob : IJob
{
private const int BatchSize = 10;
private const int HttpTimeoutS = 120;
// Mirrors ActionProcessorWorker.MaxParallel — same hardcoded-constant convention this
// repo already uses for concurrency caps (no IConfiguration wiring elsewhere either).
private const int MaxParallelJobs = 4;
private static readonly string WorkerName =
$"{Environment.MachineName}:{Environment.ProcessId}";
private readonly IHttpClientFactory _httpFactory;
private readonly IApplicationConnection _appConnection;
private readonly IOptionsSnapshot _databaseDTO;
private readonly IQueryExecutor _queryExecutor;
private readonly AutoNumber _autoNumber;
private readonly IConfiguration _config;
private readonly IHubContext _hubContext;
private readonly IFileUploadBLL _fileUploadBLL;
private readonly ISysJobNotificationDispatcher _notificationDispatcher;
private readonly IUserDAL _userDAL;
private readonly IWebhookEndpointResolver _webhookEndpointResolver;
private readonly ILogger _logger;
public SysJobExecutorQuartzJob(
IHttpClientFactory httpFactory,
IApplicationConnection appConnection,
IOptionsSnapshot databaseDTO,
IQueryExecutor queryExecutor,
AutoNumber autoNumber,
IConfiguration config,
IHubContext hubContext,
IFileUploadBLL fileUploadBLL,
ISysJobNotificationDispatcher notificationDispatcher,
IUserDAL userDAL,
IWebhookEndpointResolver webhookEndpointResolver,
ILogger logger)
{
_httpFactory = httpFactory;
_appConnection = appConnection;
_databaseDTO = databaseDTO;
_queryExecutor = queryExecutor;
_autoNumber = autoNumber;
_config = config;
_hubContext = hubContext;
_fileUploadBLL = fileUploadBLL;
_notificationDispatcher = notificationDispatcher;
_userDAL = userDAL;
_webhookEndpointResolver = webhookEndpointResolver;
_logger = logger;
}
public async Task Execute(IJobExecutionContext context)
{
var ct = context.CancellationToken;
try
{
var systemConn = await _appConnection.Gb5SystemConnectionString().ConfigureAwait(false);
int dbType = _databaseDTO.Value.DataBaseType;
const string tenantSql = @"
SELECT
SERVERCONFIG1.CLIENTID AS ClientId,
SERVERCONFIG1.DATABASENAME AS DatabaseName,
SERVERCONFIG1.DATABASETYPE AS DbType,
SERVERCONFIG1.CONNECTIONNAME AS ConnectionName
FROM MSERVERCONFIG SERVERCONFIG1
JOIN MSERVER SERVER1 ON SERVERCONFIG1.SERVERID = SERVER1.SERVERID
WHERE SERVERCONFIG1.STATUS = 1
AND SERVERCONFIG1.CONNECTIONNAME <> 'ACTIVITI'";
using IDbConnection conn = dbType switch
{
DBTYPE.SQL => new SqlConnection(systemConn),
DBTYPE.POSTGRESQL => new NpgsqlConnection(systemConn),
_ => throw new NotSupportedException($"Unsupported DB type: {dbType}")
};
var tenants = (await SqlMapper.QueryAsync(conn, tenantSql)
.ConfigureAwait(false)).ToList();
foreach (var tenant in tenants)
{
try
{
await _appConnection.DBConnectionStringCached(tenant.ConnectionName)
.ConfigureAwait(false);
}
catch (Exception connEx)
{
// Was previously a bare `catch { LogWarning(...) }` with no exception detail
// at all — every real failure reason (bad password decrypt, missing server
// config row, transient cache/Vault hiccup) was indistinguishable from any
// other. Logging the real exception is required to diagnose a specific
// tenant's connection failure instead of guessing.
_logger.LogWarning(connEx,
"SysJobExecutor: skipping tenant {ClientId} ({ConnectionName}) — connection not configured",
tenant.ClientId, tenant.ConnectionName);
continue;
}
var login = BuildLogin(tenant);
try
{
await ProcessTenantAsync(login, ct).ConfigureAwait(false);
}
catch (Microsoft.Data.SqlClient.SqlException sqlEx) when (sqlEx.Number == 208)
{
// 208 = Invalid object name — TSYSJOB table not installed for this tenant, skip silently
_logger.LogWarning(
"SysJobExecutorQuartzJob: skipping tenant {Id} — TSYSJOB table not found (schema not installed)",
tenant.ClientId);
}
catch (Exception ex) when (ex.Message.Contains("No server configuration found") ||
ex.InnerException?.Message.Contains("No server configuration found") == true)
{
_logger.LogWarning(
"SysJobExecutorQuartzJob: skipping tenant {Id} ({Conn}) — connection not configured",
tenant.ClientId, tenant.ConnectionName);
}
catch (Exception ex)
{
_logger.LogError(ex, "SysJobExecutorQuartzJob failed for tenant {Id}", tenant.ClientId);
}
}
}
catch (Exception ex)
{
_logger.LogCritical(ex, "SysJobExecutorQuartzJob outer loop failed");
}
}
// ── Per-tenant processing ─────────────────────────────────────────────
private async Task ProcessTenantAsync(LoginDTO login, System.Threading.CancellationToken ct)
{
var pending = (await _queryExecutor.QueryAsync(
login,
SystemJobQB.GET_PENDING_JOBS,
new { BatchSize, SysJobTenantId = login.ClientId },
cancellationToken: ct).ConfigureAwait(false)).ToList();
if (pending.Count == 0)
return;
_logger.LogInformation("SysJobExecutor: {Count} pending jobs for tenant {Id}",
pending.Count, login.ClientId);
// Bounded concurrency, same shape as ActionProcessorWorker's own throttle — a plain
// sequential foreach here meant one slow report (up to HttpTimeoutS=120s) blocked
// every other pending job behind it, tenant-wide. GET_PENDING_JOBS already orders by
// Priority DESC then SubmittedOn ASC, and SemaphoreSlim grants permits in that same
// FIFO order, so higher-priority jobs still get a worker slot first — this only adds
// throughput, it doesn't change which job goes first. Each job's own exception is
// caught individually so one bad job can't fault the others via Task.WhenAll.
using var throttle = new System.Threading.SemaphoreSlim(MaxParallelJobs, MaxParallelJobs);
var jobTasks = pending.Select(async job =>
{
await throttle.WaitAsync(ct).ConfigureAwait(false);
try
{
ct.ThrowIfCancellationRequested();
await ExecuteJobAsync(job, login, ct).ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogError(ex, "SysJobExecutor: unhandled error | SysJobId={Id}", job.SysJobId);
}
finally
{
throttle.Release();
}
});
await Task.WhenAll(jobTasks).ConfigureAwait(false);
}
private async Task ExecuteJobAsync(PendingSysJobDTO job, LoginDTO login,
System.Threading.CancellationToken ct)
{
// ── 1. Atomically claim the job ──────────────────────────────────
int claimed = await _queryExecutor.ExecuteAsync(
login,
SystemJobQB.CLAIM_JOB,
new { job.SysJobId, SysJobTenantId = login.ClientId, SysJobModifiedById = login.UserId },
cancellationToken: ct).ConfigureAwait(false);
if (claimed == 0)
{
// Could be claimed by another worker or cancelled — check status to log appropriately
_logger.LogInformation("SysJobExecutor: SysJobId={Id} already claimed or cancelled — skipping",
job.SysJobId);
return;
}
// ── 2. Create TSYSJOBRUN entry ───────────────────────────────────
var runAutoNumber = await _autoNumber.GetAutoNumber(1, "SYSJOBRUN", login).ConfigureAwait(false);
var sysJobRunId = runAutoNumber.StartNumber;
var startedAt = DateTime.UtcNow;
var runDto = new SysJobRunDTO
{
SysJobRunId = sysJobRunId,
SysJobId = job.SysJobId,
StartedAt = startedAt,
EndedAt = null,
Status = 1, // InProgress
Worker = WorkerName,
RunIdentifier = sysJobRunId,
CreatedById = login.UserId,
CreatedOn = startedAt,
ModifiedById = login.UserId,
ModifiedOn = startedAt,
TenantId = login.ClientId
};
await _queryExecutor.ExecuteAsync(login, SysJobRunQB.SAVE_SYSJOBRUN, runDto,
cancellationToken: ct).ConfigureAwait(false);
_logger.LogInformation(
"SysJobExecutor: started | SysJobId={Id} SysJobRunId={RunId} Worker={Worker}",
job.SysJobId, sysJobRunId, WorkerName);
await _hubContext.Clients
.Group(SysJobHub.TenantGroup(login.ClientId))
.JobStarted(job.SysJobId, sysJobRunId, WorkerName, startedAt)
.ConfigureAwait(false);
// ── 3. Execute the web service call ──────────────────────────────
string? resultLocation = null;
string? errorMessage = null;
bool success = false;
try
{
(success, resultLocation, errorMessage) =
await CallWebServiceAsync(job, login, ct).ConfigureAwait(false);
}
catch (Exception ex)
{
errorMessage = ex.Message;
_logger.LogError(ex, "SysJobExecutor: web service call threw | SysJobId={Id}", job.SysJobId);
}
var endedAt = DateTime.UtcNow;
// ── 4. Update TSYSJOBRUN ─────────────────────────────────────────
runDto.EndedAt = endedAt;
runDto.Status = success ? 2 : 3;
runDto.Log = errorMessage ?? (success ? "Completed" : "Failed");
runDto.ModifiedOn = endedAt;
await _queryExecutor.ExecuteAsync(login, SysJobRunQB.UPDATE_SYSJOBRUN, runDto,
cancellationToken: ct).ConfigureAwait(false);
// ── 5. Update TSYSJOB ────────────────────────────────────────────
// COMPLETE_JOB and FAIL_OR_RETRY_JOB both include AND RUNSTATUS = 1 in their
// WHERE clause. If a cancel arrived while the HTTP call was in flight, RUNSTATUS
// is already 4 and 0 rows are updated — we must NOT overwrite the cancel.
if (success)
{
int updated = await _queryExecutor.ExecuteAsync(
login,
SystemJobQB.COMPLETE_JOB,
new { job.SysJobId, SysJobTenantId = login.ClientId, ResultLocation = resultLocation,
SysJobModifiedById = login.UserId },
cancellationToken: ct).ConfigureAwait(false);
if (updated == 0)
{
_logger.LogInformation(
"SysJobExecutor: job cancelled during execution | SysJobId={Id} SysJobRunId={RunId}",
job.SysJobId, sysJobRunId);
await _hubContext.Clients
.Group(SysJobHub.TenantGroup(login.ClientId))
.JobCancelled(job.SysJobId, endedAt)
.ConfigureAwait(false);
return;
}
_logger.LogInformation("SysJobExecutor: completed | SysJobId={Id} SysJobRunId={RunId}",
job.SysJobId, sysJobRunId);
await _hubContext.Clients
.Group(SysJobHub.TenantGroup(login.ClientId))
.JobCompleted(job.SysJobId, sysJobRunId, resultLocation, endedAt)
.ConfigureAwait(false);
if (job.DeliveryChannel.HasValue)
await _notificationDispatcher.DispatchAsync(
job.SysJobId, sysJobRunId, job.DeliveryChannel.Value, job.DeliveryDestination,
job.SubmittedById, runStatus: 2, login, ct).ConfigureAwait(false);
}
else
{
bool shouldRetry = job.RetryCount < job.MaxRetries;
int nextStatus = shouldRetry ? 0 : 3; // 0=retry-pending, 3=final-fail
int updated = await _queryExecutor.ExecuteAsync(
login,
SystemJobQB.FAIL_OR_RETRY_JOB,
new { job.SysJobId, SysJobTenantId = login.ClientId,
RunStatus = nextStatus, ErrorMessage = errorMessage,
SysJobModifiedById = login.UserId },
cancellationToken: ct).ConfigureAwait(false);
if (updated == 0)
{
_logger.LogInformation(
"SysJobExecutor: job cancelled during execution (on failure path) | SysJobId={Id} SysJobRunId={RunId}",
job.SysJobId, sysJobRunId);
await _hubContext.Clients
.Group(SysJobHub.TenantGroup(login.ClientId))
.JobCancelled(job.SysJobId, endedAt)
.ConfigureAwait(false);
return;
}
_logger.LogWarning(
"SysJobExecutor: {Outcome} | SysJobId={Id} SysJobRunId={RunId} RetryCount={R}/{Max} Error={E}",
shouldRetry ? "retry-queued" : "permanently-failed",
job.SysJobId, sysJobRunId, job.RetryCount + 1, job.MaxRetries, errorMessage);
if (shouldRetry)
{
await _hubContext.Clients
.Group(SysJobHub.TenantGroup(login.ClientId))
.JobRetrying(job.SysJobId, sysJobRunId, job.RetryCount + 1, job.MaxRetries,
errorMessage ?? "Unknown error", endedAt)
.ConfigureAwait(false);
}
else
{
await _hubContext.Clients
.Group(SysJobHub.TenantGroup(login.ClientId))
.JobFailed(job.SysJobId, sysJobRunId,
errorMessage ?? "Unknown error", endedAt)
.ConfigureAwait(false);
if (job.DeliveryChannel.HasValue)
await _notificationDispatcher.DispatchAsync(
job.SysJobId, sysJobRunId, job.DeliveryChannel.Value, job.DeliveryDestination,
job.SubmittedById, runStatus: 3, login, ct).ConfigureAwait(false);
}
}
}
// ── HTTP call ─────────────────────────────────────────────────────────
private async Task<(bool Success, string? ResultLocation, string? Error)> CallWebServiceAsync(
PendingSysJobDTO job, LoginDTO login, System.Threading.CancellationToken ct)
{
// Delegates to the same resolver ActionProcessor/webhooks already use — it prefers
// MWEBSERVICE.SECONDURITEMPLATE over URITEMPLATE when set (several real reports are
// wired this way: URITEMPLATE holds the legacy GB4 route, SECONDURITEMPLATE the
// GB5-native replacement) and substitutes the literal {BaseURI} token every real
// MWEBSERVICE row carries — this executor's own prior "StartsWith(\"http\")" check
// never did that substitution, so any row using the standard {BaseURI} convention
// would have produced a literally-invalid host.
string? url = await _webhookEndpointResolver.ResolveAsync(
job.WebServiceId, uriParameterValue: null, login.ClientId, login.ConnectionDatabaseName, ct)
.ConfigureAwait(false);
if (string.IsNullOrWhiteSpace(url))
return (false, null, $"No MWEBSERVICE configured for WebServiceId={job.WebServiceId}");
// MWEBSERVICE.METHODTYPE only means something for the legacy GB4 URITEMPLATE side —
// it was never updated when SECONDURITEMPLATE (the GB5-native replacement) was added,
// and every BaseReportEndpoint-derived report route in this codebase is POST (confirmed
// across all ~28 of them), never GET. A real report's METHODTYPE=1 (GET), live-tested
// against AccountProfitLossDashboardReport, produced a hard HTTP 405 for exactly this
// reason. Once a usable SECONDURITEMPLATE exists, always POST — ignore METHODTYPE.
bool isUsableTemplate(string? template) =>
!string.IsNullOrWhiteSpace(template)
&& !string.Equals(template, "NONE", StringComparison.OrdinalIgnoreCase);
var method = isUsableTemplate(job.WebServiceSecondUriTemplate)
? HttpMethod.Post
: (job.WebServiceMethodType == 1 ? HttpMethod.Get : HttpMethod.Post);
var http = _httpFactory.CreateClient("sysjob");
var request = new HttpRequestMessage(method, url);
// GB5-native FastEndpoints targets (e.g. BaseReportEndpoint-derived report routes)
// require a Login header to resolve LoginDTO server-side — without it every such
// call fails with 400 "This header is missing from the request!". Legacy GB4 .svc
// targets ignore unknown headers, so this is safe to send unconditionally.
// Use the submitting user's own WorkOUId/DateFormat/etc — the bare tenant-level
// `login` has none of these, so any report needing company/branch letterhead data
// (e.g. PDF export "standard fields") fails; the job should render as if the
// submitting user ran it themselves, not as a contextless system account.
var targetLogin = await BuildTargetLoginAsync(job, login, ct).ConfigureAwait(false);
request.Headers.Add("Login", JsonSerializer.Serialize(targetLogin));
// BaseReportEndpoint (every GB5-native report route) picks its export format off a
// ReportFormat header — without it a Job submitted for a PDF/Excel export just gets
// the default JSON/paged response, not a file. job.ResultFormat uses the frontend's
// own scheme (gbjobworker.component.ts RESULT_FORMAT_BY_OUTPUT), which does not match
// GB5Shared's ReportFormat codes, so it must be translated, not forwarded as-is.
int? reportFormatCode = MapResultFormatToReportFormatHeader(job.ResultFormat);
if (reportFormatCode.HasValue)
request.Headers.Add("ReportFormat", reportFormatCode.Value.ToString());
else if (job.ResultFormat != 0)
_logger.LogWarning(
"SysJobExecutor: SysJobId={Id} requested unsupported ResultFormat={Format} — " +
"no ReportFormat header sent, target will return its default JSON response",
job.SysJobId, job.ResultFormat);
if (method == HttpMethod.Post && !string.IsNullOrWhiteSpace(job.Parameters))
request.Content = new StringContent(job.Parameters, Encoding.UTF8, "application/json");
using var cts = System.Threading.CancellationTokenSource.CreateLinkedTokenSource(ct);
cts.CancelAfter(TimeSpan.FromSeconds(HttpTimeoutS));
// ResponseHeadersRead: check Content-Type before reading the body so we
// can pick the right read path (binary vs text) without double-buffering.
using var response = await http.SendAsync(
request, HttpCompletionOption.ResponseHeadersRead, cts.Token).ConfigureAwait(false);
if (!response.IsSuccessStatusCode)
{
string errorBody = await response.Content.ReadAsStringAsync(cts.Token).ConfigureAwait(false);
return (false, null, $"HTTP {(int)response.StatusCode}: {errorBody[..Math.Min(errorBody.Length, 2000)]}");
}
_logger.LogInformation(
"SysJobExecutor: HTTP {Status} | SysJobId={Id} Url={Url}",
(int)response.StatusCode, job.SysJobId, url);
// ── File response — upload to DMS centrally ──────────────────────
// The report service just returns the file (Content-Type: application/pdf,
// Content-Disposition: attachment; filename="..."). No DMS knowledge needed
// on the service side. The executor owns all file storage concerns.
if (IsFileResponse(response))
{
byte[] fileBytes = await response.Content.ReadAsByteArrayAsync(cts.Token).ConfigureAwait(false);
string mimeType = response.Content.Headers.ContentType?.MediaType ?? "application/octet-stream";
string fileName = ExtractFileName(response, job);
string token = await SysJobAttachmentHelper.UploadAsync(
_fileUploadBLL, fileBytes, fileName, mimeType, job.SysJobId, login, ct)
.ConfigureAwait(false);
_logger.LogInformation(
"SysJobExecutor: file stored in DMS | SysJobId={Id} FileName={Name} Size={Bytes}B Token={Token}",
job.SysJobId, fileName, fileBytes.Length, token);
return (true, token, null);
}
// ── Text / JSON response ─────────────────────────────────────────
// Fire-and-forget jobs return nothing useful. Services that still want to
// return a pre-built token (att:/redis:/filesystem path) can do so as a
// short plain-text body — kept for backward compatibility.
string body = await response.Content.ReadAsStringAsync(cts.Token).ConfigureAwait(false);
// BaseReportEndpoint's export path never sends a raw file response — it always
// wraps the export bytes as a base64 string inside a ResponseStandardDTO