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; } } }