using FrameworkBLL.SchedulerTaskGenerator;
using FrameworkDAL.DTO.SchedulerTaskGenerator;
using GB5Shared.ActionProcessor;
using GB5Shared.DTO.Framework.Login;
using GB5Shared.DTO.JobEngine;
using GB5Shared.DTO.PubSub;
using GB5Shared.QueueReader;
using GB5Shared.Telemetry;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using System;
using System.Diagnostics;
using System.Net.Http;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
namespace FrameworkSL.Controllers.SchedulerTaskGenerator
{
public interface ISchedulerExecutionDeliveryService
{
Task DeliverAsync(SchedulerTaskDTO item, LoginDTO login, CancellationToken ct);
}
///
/// Makes the actual HTTP call for one TJOBEXECUTION row claimed by
/// SchedulerExecutionPoller, then writes the result back via
/// ISchedulerTaskServiceBLL.UpdateExecutionResultAsync, which itself decides
/// retry-vs-DLQ (SchedulerTaskServiceBLL.ApplyRetryOrDeadLetterAsync). On success,
/// also enqueues a SCHEDULER_ACTION_EVENT to TJOBQUEUE if the job has a configured
/// MACTION (see EnqueueActionEventIfConfiguredAsync) — TJOBQUEUE's job as of
/// 2026-08-20 is purely this action-dispatch queue, not webhook delivery.
///
/// Replaces SchedulerWebServiceCallHandler (the old TJOBQUEUE IJobQueueHandler for
/// webhook delivery). TJOBEXECUTION.RETRYNUMBER is the only retry counter for the
/// webhook delivery itself — there is no second, TJOBQUEUE-level retry/DLQ decision
/// to keep in sync with it any more (see SchedulerExecutionPoller's header comment).
/// The webhook URL is resolved fresh on every attempt (including retries) rather
/// than baked in once and carried through a queue payload — MJOBDEFINE's
/// WebServiceId/UriParameterValue are already available on the claimed row.
///
public sealed class SchedulerExecutionDeliveryService : ISchedulerExecutionDeliveryService
{
private const int HttpTimeoutSeconds = 120;
private static readonly JsonSerializerOptions RequestPayloadJsonOptions = new()
{
DefaultIgnoreCondition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingNull
};
///
/// Snapshot of the exact outbound request for one delivery attempt — saved to
/// TJOBEXECUTION.REQUESTPAYLOAD right before the HTTP call, so the request itself
/// (not just SUCCESSMESSAGE/ERRORMESSAGE, which only capture the response) is
/// inspectable per execution. PostData is null (and so omitted from the JSON) for
/// GET, since no body is ever sent for that method.
///
private sealed class SchedulerRequestPayload
{
public int JobId { get; set; }
public long JobExecutionId { get; set; }
public int TenantId { get; set; }
public int WebServiceId { get; set; }
public string EndpointUrl { get; set; } = string.Empty;
public string HttpMethod { get; set; } = string.Empty;
public string HttpMethodName { get; set; } = string.Empty;
public string? PostData { get; set; }
public string? CorrelationId { get; set; }
}
private readonly IHttpClientFactory _httpFactory;
private readonly IWebhookEndpointResolver _urlResolver;
private readonly IJobQueueEnqueuer _jobQueueEnqueuer;
private readonly IServiceScopeFactory _scopeFactory;
private readonly ILogger _logger;
public SchedulerExecutionDeliveryService(
IHttpClientFactory httpFactory,
IWebhookEndpointResolver urlResolver,
IJobQueueEnqueuer jobQueueEnqueuer,
IServiceScopeFactory scopeFactory,
ILogger logger)
{
_httpFactory = httpFactory;
_urlResolver = urlResolver;
_jobQueueEnqueuer = jobQueueEnqueuer;
_scopeFactory = scopeFactory;
_logger = logger;
}
public async Task DeliverAsync(SchedulerTaskDTO item, LoginDTO login, CancellationToken ct)
{
// Reconstruct the "JobScheduler" root span's context (persisted to
// TJOBEXECUTION.TRACEPARENT/TRACESTATE when Quartz fired) so this delivery —
// running in a completely separate async continuation (a later poll tick) —
// still lands as a child span of the SAME trace instead of starting a new one.
var deliverParentContext =
item.TraceParent is not null
&& ActivityContext.TryParse(item.TraceParent, item.TraceState, isRemote: true, out var deliverCtx)
? deliverCtx
: default;
using var deliverActivity = deliverParentContext != default
? GB5ActivitySources.Scheduler.StartActivity("deliver", ActivityKind.Internal, deliverParentContext)
: GB5ActivitySources.Scheduler.StartActivity("deliver", ActivityKind.Internal);
deliverActivity?.SetTag("gb5.job.id", item.JobId);
deliverActivity?.SetTag("gb5.job_execution.id", item.JobExecutionId);
deliverActivity?.SetTag("gb5.tenant.id", item.TenantId);
using var scope = _scopeFactory.CreateScope();
var service = scope.ServiceProvider.GetRequiredService();
var hub = scope.ServiceProvider.GetRequiredService();
await hub.PushJobDispatched(item.TenantId, new JobDispatchedDTO
{
JobId = item.JobId,
JobExecutionId = item.JobExecutionId,
CorrelationId = item.CorrelationId,
ActionCount = 1,
DispatchedAtUtc = DateTime.UtcNow
}, ct).ConfigureAwait(false);
// Build the execution login from the RunAs context stored in MJOBDEFINE at scheduling time.
// Clone via JSON round-trip to copy all tenant-routing fields (ClientId, DatabaseName,
// ServerConfigId, DatabaseType, etc.) from the system login, then overwrite only the
// user-identity fields with the scheduling user's saved context so the target endpoint
// sees the correct OU/branch/role — not a generic system user (UserId=-1).
// IsSchedulerRun=1 marks this as an automated execution (not an interactive user request).
var executionLogin = JsonSerializer.Deserialize(JsonSerializer.Serialize(login))!;
executionLogin.IsSchedulerRun = 1;
if (item.RunAsUserId > 0)
{
executionLogin.UserId = item.RunAsUserId;
executionLogin.UserName = item.RunAsUserName ?? "SchedulerService";
executionLogin.RoleId = item.RunAsRoleId;
executionLogin.WorkOUId = item.RunAsOUId;
executionLogin.WorkPartyBranchId = item.RunAsBranchId;
executionLogin.WorkPartyId = item.RunAsWorkPartyId;
executionLogin.WorkStoreId = item.RunAsStoreId;
executionLogin.WorkPeriodId = item.RunAsPeriodId;
}
var started = DateTime.UtcNow;
bool success;
string? errorMessage = null;
string? resultMessage = null;
string? responseBody = null;
var resolvedUrl = string.Empty;
try
{
var baseResolvedUrl = await _urlResolver.ResolveAsync(
item.WebServiceId, uriParameterValue: null, item.TenantId, login.DatabaseName, ct)
.ConfigureAwait(false);
if (string.IsNullOrWhiteSpace(baseResolvedUrl))
throw new InvalidOperationException(
$"SchedulerExecutionDeliveryService: WebhookEndpointResolver could not resolve a URL for " +
$"JobId={item.JobId} WebServiceId={item.WebServiceId} — check SysJobSettings:Gb4ServiceBaseUrl");
resolvedUrl = ApplyNamedParameters(baseResolvedUrl, item.UriParameterValue);
// Validated explicitly (rather than letting the request fail deep inside
// HttpClient) so a malformed URL reports WHAT failed and WHY — .NET's own
// "Invalid URI: The hostname could not be parsed." gives no clue which job/URL
// it came from. Uri.TryCreate alone doesn't reliably catch either failure mode
// below — it's lenient enough to accept a literal unresolved "{Name}" token or
// an accidental doubled "http://http://" scheme.
string? invalidReason = null;
if (resolvedUrl.Contains("{BaseURI}"))
invalidReason = "the '{BaseURI}' token was never resolved — check SysJobSettings:Gb4ServiceBaseUrl";
else if (resolvedUrl.Contains('{') && resolvedUrl.Contains('}'))
invalidReason = "one or more '{Name}' placeholders were left unresolved — check MJOBDEFINE.UriParameterValue";
else if (System.Text.RegularExpressions.Regex.IsMatch(resolvedUrl, @"^[a-zA-Z][a-zA-Z0-9+.\-]*://[a-zA-Z][a-zA-Z0-9+.\-]*://"))
invalidReason = "the URL has a doubled scheme (e.g. 'http://http://...') — check SysJobSettings:Gb4ServiceBaseUrl for a leftover 'http://' prefix";
else if (!Uri.TryCreate(resolvedUrl, UriKind.Absolute, out _))
invalidReason = "the resolved string is not a well-formed absolute URI";
if (invalidReason is not null)
throw new InvalidOperationException(
$"Could not build a valid request URL for JobId={item.JobId} " +
$"WebServiceId={item.WebServiceId} — {invalidReason}. Resolved='{resolvedUrl}'");
var validatedUri = new Uri(resolvedUrl, UriKind.Absolute);
// MWEBSERVICE.METHODTYPE: 0=GET, 1=POST, 2=DELETE (per GB5URITemplate
// convention — see WebServiceQB.cs's GET_WEBSERVICEID_BY_ENTITY_AND_METHOD
// comment). NOT the same mapping SysJobExecutorQuartzJob.CallWebServiceAsync
// uses ("1"=GET) — that convention is specific to that pipeline's own historical
// usage and does not match MWEBSERVICE's actual documented semantics.
var method = item.HttpMethod switch
{
"1" => HttpMethod.Post,
"2" => HttpMethod.Delete,
_ => HttpMethod.Get
};
// Snapshot of exactly what's about to be requested — saved to
// TJOBEXECUTION.REQUESTPAYLOAD so the actual request (not just the
// response) is inspectable per execution. PostData is only present
// (even if empty) for methods that could carry a body — omitted
// entirely for GET, matching the fact that no body is ever sent today.
var requestPayload = JsonSerializer.Serialize(new SchedulerRequestPayload
{
JobId = item.JobId,
JobExecutionId = item.JobExecutionId,
TenantId = item.TenantId,
WebServiceId = item.WebServiceId,
EndpointUrl = resolvedUrl,
HttpMethod = item.HttpMethod,
HttpMethodName = method.Method,
PostData = method == System.Net.Http.HttpMethod.Post ? string.Empty : null,
CorrelationId = item.CorrelationId
}, RequestPayloadJsonOptions);
await service.SaveRequestPayloadAsync((int)item.JobExecutionId, requestPayload, login, ct)
.ConfigureAwait(false);
var http = _httpFactory.CreateClient("sysjob");
var request = new HttpRequestMessage(method, validatedUri);
// GB5-native FastEndpoints targets resolve LoginDTO server-side off this
// header — without it every such call fails with 400 "This header is
// missing from the request!" (see SysJobExecutorQuartzJob.CallWebServiceAsync,
// the same convention for the pre-existing SysJob pipeline).
// Use executionLogin (the scheduling user's identity) so the endpoint sees the
// correct OU/branch/role context — not the generic system login (UserId=-1).
request.Headers.Add("Login", JsonSerializer.Serialize(executionLogin));
// For POST endpoints with a configured TCRITERIACONFIG: build and attach the
// criteria JSON as the request body so the endpoint receives the correct filter
// parameters. Variable field tokens ({WorkOUID}, {TodayFrom}, etc.) are resolved
// against executionLogin's org context at delivery time — not at scheduling time —
// so date-relative filters always use the actual execution date.
if (method == HttpMethod.Post && item.CriteriaConfigId > 0)
{
var criteriaJson = await service.LoadCriteriaForExecutionAsync(
item.CriteriaConfigId, executionLogin, ct).ConfigureAwait(false);
if (!string.IsNullOrWhiteSpace(criteriaJson))
request.Content = new StringContent(criteriaJson, Encoding.UTF8, "application/json");
}
using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
cts.CancelAfter(TimeSpan.FromSeconds(HttpTimeoutSeconds));
using var response = await http.SendAsync(
request, HttpCompletionOption.ResponseHeadersRead, cts.Token).ConfigureAwait(false);
// Full body, no preview truncation — TJOBEXECUTION.SUCCESSMESSAGE/ERRORMESSAGE
// are NVARCHAR(MAX), so there is no column-size reason to cut this short, and a
// truncated response is not verifiable as "the actual API response".
var body = await response.Content.ReadAsStringAsync(cts.Token).ConfigureAwait(false);
if (response.IsSuccessStatusCode)
{
success = true;
resultMessage = $"HTTP {(int)response.StatusCode} {response.ReasonPhrase}: {body}";
// Raw response body, with no "HTTP 200 OK:" prefix — this is what gets
// carried into the action-event Payload (and so into TJOBQUEUE.PAYLOAD),
// not the JobId/JobExecutionId bookkeeping object that was here before.
responseBody = body;
}
else
{
errorMessage = $"HTTP {(int)response.StatusCode}: {body}";
resultMessage = errorMessage;
success = false;
}
}
catch (Exception ex)
{
success = false;
errorMessage = ex.Message;
resultMessage = errorMessage;
}
var completed = DateTime.UtcNow;
await service.UpdateExecutionResultAsync(
(int)item.JobExecutionId, success, resultMessage, login, ct).ConfigureAwait(false);
if (success)
{
_logger.LogInformation(
"SchedulerExecutionDeliveryService: SUCCESS | JobId={JobId} JobExecutionId={Id} Url={Url}",
item.JobId, item.JobExecutionId, resolvedUrl);
await hub.PushJobCompleted(item.TenantId, new JobCompletedDTO
{
JobId = item.JobId,
JobExecutionId = item.JobExecutionId,
CorrelationId = item.CorrelationId,
SuccessCount = 1,
TotalCount = 1,
DurationMs = (long)(completed - started).TotalMilliseconds,
CompletedAtUtc = completed
}, ct).ConfigureAwait(false);
await EnqueueActionEventIfConfiguredAsync(item, executionLogin, service, responseBody, ct).ConfigureAwait(false);
return;
}
_logger.LogWarning(
"SchedulerExecutionDeliveryService: FAILED | JobId={JobId} JobExecutionId={Id} Url={Url} Error={Error}",
item.JobId, item.JobExecutionId, resolvedUrl, errorMessage);
await hub.PushJobFailed(item.TenantId, new JobFailedDTO
{
JobId = item.JobId,
JobExecutionId = item.JobExecutionId,
CorrelationId = item.CorrelationId,
ErrorMessage = errorMessage ?? "Unknown error",
FailureCount = 1,
TotalCount = 1,
FailedAtUtc = completed
}, ct).ConfigureAwait(false);
// Deliberately not rethrown — there is no queue-level catch left to retry this.
// UpdateExecutionResultAsync above already routed the failure through
// SchedulerTaskServiceBLL.ApplyRetryOrDeadLetterAsync (STATUS=0 + backoff, or
// STATUS=5 DLQ), which is now the only retry/DLQ decision for this job.
}
///
/// On a successful job completion, checks whether this job has an active MACTION row
/// linked via the scheduler's fixed EventTypeId (SchedulerActionEventConstants.
/// SchedulerJobEventTypeId) and, if so, enqueues a SCHEDULER_ACTION_EVENT to TJOBQUEUE
/// carrying an OutboxEventMessage — the same message shape OutBox.PublishPendingEventsAsync
/// produces for every other domain event. SchedulerActionEventQueueProcessor/
/// SchedulerActionEventPublishHandler pick it up and Dapr-publish it to topic
/// "EVENTTYPEID:{id}", which EventActionSubscribeController/EventSubBLL already know
/// how to resolve back to this same MACTION row (via MEVENTTYPEACTION/
/// MEVENTTYPEACTIONDETAIL, wired by AutoSchedulerDAL when the action was saved) and
/// dispatch (mail/SMS/webhook/etc). Only fires on success — a failed/DLQ'd job does
/// not enqueue an action event.
///
private async Task EnqueueActionEventIfConfiguredAsync(
SchedulerTaskDTO item, LoginDTO executionLogin, ISchedulerTaskServiceBLL service, string? responseBody, CancellationToken ct)
{
bool hasAction;
try
{
hasAction = await service.HasActiveSchedulerActionAsync(item.JobId, executionLogin, ct)
.ConfigureAwait(false);
}
catch (Exception ex)
{
Activity.Current?.SetTag("gb5.scheduler.has_action", false);
Activity.Current?.SetTag("gb5.scheduler.has_action_check_error", ex.Message);
_logger.LogWarning(ex,
"SchedulerExecutionDeliveryService: HasActiveSchedulerActionAsync failed — skipping action event | JobId={JobId}",
item.JobId);
return;
}
// Visible directly on the "deliver" span — GB5.Scheduler.HasActiveSchedulerAction
// in Zipkin — so "did this job have an action?" is answered by the trace itself,
// no separate MACTION query needed.
using (var actionCheckActivity = GB5ActivitySources.Scheduler.StartActivity(
"HasActiveSchedulerAction", ActivityKind.Internal))
{
actionCheckActivity?.SetTag("gb5.job.id", item.JobId);
actionCheckActivity?.SetTag("gb5.scheduler.event_type_id", SchedulerActionEventConstants.SchedulerJobEventTypeId);
actionCheckActivity?.SetTag("gb5.scheduler.has_action", hasAction);
}
Activity.Current?.SetTag("gb5.scheduler.has_action", hasAction);
if (!hasAction)
return;
var message = new OutboxEventMessage
{
EventTypeId = SchedulerActionEventConstants.SchedulerJobEventTypeId,
ObjectId = item.JobId,
TenantId = item.TenantId,
UserId = executionLogin.UserId,
// The actual webservice response for this job's HTTP call — not a
// JobId/JobExecutionId bookkeeping object. Subscribers (mail/webhook/etc.)
// receive exactly what the API returned, matching TOUTBOX's own convention
// of carrying the real domain payload rather than a synthetic wrapper.
Payload = responseBody ?? string.Empty,
CorrelationKey = item.CorrelationId ?? string.Empty,
ConnectionName = executionLogin.DatabaseName ?? string.Empty,
// Carries the "deliver" span (itself a child of the original "JobScheduler"
// root) forward — SchedulerActionEventPublishHandler reconstructs this as the
// parent for its own Dapr-publish span, keeping the whole chain in one trace.
TraceParent = Activity.Current?.Id,
TraceState = Activity.Current?.TraceStateString
};
using var enqueueActivity = GB5ActivitySources.Scheduler.StartActivity(
"enqueue TJOBQUEUE", ActivityKind.Producer);
enqueueActivity?.SetTag("gb5.job.id", item.JobId);
enqueueActivity?.SetTag("gb5.job_execution.id", item.JobExecutionId);
enqueueActivity?.SetTag("messaging.destination", "TJOBQUEUE");
enqueueActivity?.SetTag("messaging.message_type", SchedulerActionEventConstants.MessageType);
try
{
var queueId = await _jobQueueEnqueuer.EnqueueAsync(
SchedulerActionEventConstants.MessageType,
JsonSerializer.Serialize(message),
executionLogin,
messageId: $"job-action-{item.JobExecutionId}",
correlationId: item.CorrelationId,
sourceService: "SchedulerExecutionDelivery",
jobExecutionId: item.JobExecutionId,
ct: ct).ConfigureAwait(false);
enqueueActivity?.SetTag("gb5.tjobqueue.queue_id", queueId);
_logger.LogInformation(
"SchedulerExecutionDeliveryService: action event enqueued | JobId={JobId} JobExecutionId={Id} QueueId={QueueId}",
item.JobId, item.JobExecutionId, queueId);
}
catch (Exception ex)
{
// Non-fatal: the job itself already completed successfully and that result is
// already durably recorded. Losing the action-dispatch trigger here should not
// fail the job, only be visible in logs.
_logger.LogError(ex,
"SchedulerExecutionDeliveryService: failed to enqueue action event | JobId={JobId} JobExecutionId={Id}",
item.JobId, item.JobExecutionId);
}
}
///
/// Parses "{Name}=Value,{Name2}=Value2" (GB5's established URI-parameter convention —
/// see AutoSchedulerDTO's own default "{FirstNumber}=1,{MaxResult}=10") and substitutes
/// each named token into the template individually. Unlike IWebhookEndpointResolver's
/// own placeholder handling (a blanket single-value replace of every {...} token), this
/// maps each {Name} to its own value — required whenever a template has more than one
/// distinct placeholder remaining after {BaseURI} is resolved.
///
internal static string ApplyNamedParameters(string template, string? uriParameterValue)
{
if (string.IsNullOrWhiteSpace(uriParameterValue))
return template;
foreach (var pair in uriParameterValue.Split(',', StringSplitOptions.RemoveEmptyEntries))
{
var eq = pair.IndexOf('=');
if (eq <= 0) continue; // malformed segment — no "{Name}=" prefix, skip rather than throw
var token = pair[..eq].Trim(); // e.g. "{TLeaveId}"
var value = pair[(eq + 1)..].Trim();
if (token.Length > 0)
template = template.Replace(token, value);
}
return template;
}
}
}