Job Engine — Developer Guide

GB5 Enterprise Scheduler & Queue Processor · .NET 9 · Quartz 3.13 · Dapper · FastEndpoints

Architecture Overview

The Job Engine is a standalone GB5 microservice (port 5400) responsible for two concerns:

SL
JobEngineSL 17 FastEndpoints JobEngineHub (SignalR) QuartzJobExecutor QuartzSyncService OutBoxPollerService
↕
BLL
JobDefineBLL SchedulerBLL JobExecutionBLL JobQueueBLL
↕
DAL
JobDefineDAL SchedulerDAL JobExecutionDAL JobQueueDAL JobDefineQB SchedulerQB JobExecutionQB JobQueueQB
Key rules

SL → BLL → DAL only. BLL never references ASP.NET types. DAL contains only Dapper SQL — no business logic. All queries filter by TENANTID.

Module Structure

GB5Solution/JobEngine/
├── JobEngineDAL/
│   ├── DTOs/          JobDefineDTO, SchedulerDTO, JobExecutionDTO, JobQueueItemDTO, JobDashboardDTO
│   ├── Query/         JobDefineQB, SchedulerQB, JobExecutionQB, JobQueueQB      ← SQL constants only
│   ├── Interfaces/    IJobDefineDAL, ISchedulerDAL, IJobExecutionDAL, IJobQueueDAL
│   └── Implementations/
├── JobEngineBLL/
│   ├── Interfaces/    IJobDefineBLL, ISchedulerBLL, IJobExecutionBLL, IJobQueueBLL, IQuartzSyncService
│   └── Implementations/
└── JobEngineSL/
    ├── EndPoints/     17 FastEndpoints (GET/POST/PUT/DELETE)
    ├── Hubs/          JobEngineHub.cs, IJobEngineClient.cs
    ├── Services/      QuartzJobExecutor, QuartzSyncService, OutBoxPollerService
    ├── Migrations/    JobEngine_Migration.sql
    └── Program.cs

Database Schema

Existing tables (extended via ALTER TABLE)

TableNew columns added
MJOBDEFINEMAXRETRIES, TIMEOUTSECONDS, NEXTRUNTIME, ISCONCURRENT, JOBCATEGORY, SOURCEMODULE, SOURCEOBJECTID, DAPRPUBSUBTOPIC, audit columns
TSCHEDULERSCHEDULERTYPE (CRON/INTERVAL/ONE_TIME), INTERVALSECONDS, SCHEDULEDDATE, DESCRIPTION, SCHEDULERSTATUS, audit columns

New tables (created by migration)

TJOBEXECUTION
ColumnTypeNotes
EXECUTIONIDBIGINT IDENTITY PK
JOBIDINT NULLFK to MJOBDEFINE; NULL for queue-only runs
TENANTIDINT NOT NULLTenant isolation
JOBNAMENVARCHAR(200)
TRIGGERTYPENVARCHAR(20)SCHEDULED | MANUAL | QUEUE | API
STATUSNVARCHAR(20)QUEUED | RUNNING | SUCCESS | FAILED | TIMEOUT | CANCELLED
STARTEDAT / COMPLETEDATDATETIME2
DURATIONMSBIGINT NULLPopulated on completion
CORRELATIONIDNVARCHAR(100)GUID set by QuartzJobExecutor; carried to SignalR events
ERRORDETAILSNVARCHAR(MAX)Full exception string on failure
HOSTINSTANCENVARCHAR(100)Pod hostname — for cluster attribution
TJOBQUEUE
ColumnTypeNotes
QUEUEIDBIGINT IDENTITY PK
TENANTIDINT
MESSAGETYPENVARCHAR(100)CASHBACK | LOYALTY | PAYSLIP | IMPORT | WORKFLOW | …
PAYLOADNVARCHAR(MAX)JSON body; handler-specific schema
PRIORITYTINYINT default 5Higher = processed first
STATUSNVARCHAR(20)PENDING | PROCESSING | COMPLETED | FAILED | DLQ | CANCELLED
MAXRETRIESINT default 3
RETRYATTEMPTINTIncremented on each claim
NEXTATTEMPTAFTERDATETIME2 NULLBackoff ceiling: 30s → 2m → 8m → 30m → 2h → DLQ
MESSAGEIDNVARCHAR(200)Idempotency key — unique per (MESSAGEID, TENANTID)

Migration script: JobEngineSL/Migrations/JobEngine_Migration.sql. Run once; idempotent (ALTER TABLE only adds columns that don't exist).

Quartz tables

QRTZ_* tables (QRTZ_JOB_DETAILS, QRTZ_TRIGGERS, QRTZ_SCHEDULER_STATE, etc.) are not in the migration. Run the bundled Quartz SQL Server DDL once against the system database before first start.

Quartz.NET Integration

QuartzSyncService — job lifecycle management

Singleton registered in DI. BLL calls it via IQuartzSyncService so BLL never depends on Quartz types directly.

MethodWhen calledWhat it does
LoadAllJobsAsync()StartupReads all MJOBDEFINE STATUS=1 and schedules each in Quartz. Called from Program.cs after app.Build().
SyncJobAsync(jobId)After SaveJob, PauseJob, ResumeJobIf Status≠1: deletes Quartz job. Otherwise: reschedules from fresh DB read.
RemoveJobAsync(jobId)After DeleteJobRemoves trigger and job detail from Quartz cluster store.
TriggerJobNowAsync(jobId)After TriggerJob endpointFires the job immediately regardless of its next scheduled time.

QuartzJobExecutor — per-execution IJob

A new instance is created for every execution. Execution flow:

1
Read job definition
Reads MJOBDEFINE + MWEBSERVICE join via IJobDefineDAL.GetJobDetailAsync(). JobId and TenantId come from Quartz JobDataMap.
2
Insert TJOBEXECUTION STATUS=RUNNING
Generates correlationId (GUID), inserts execution row, notes STARTEDAT.
3
Broadcast ReceiveJobStarted
Pushes to SignalR group jobengine:client:{ClientId}.
4
Execute
If DaprPubSubTopic set → Dapr publish. Otherwise → HTTP call to UriTemplate (GET or POST based on METHODTYPE). {Param} in template substituted with UriParameterValue.
5
Update execution + NEXTRUNTIME
Sets STATUS=SUCCESS, DURATIONMS. Updates MJOBDEFINE.NEXTRUNTIME and LASTRUNTIME.
6
Broadcast ReceiveJobCompleted
Pushes completion event with duration to SignalR tenant group.
On failure

Updates execution STATUS=FAILED, broadcasts ReceiveJobFailed, then throws JobExecutionException(refireImmediately: false). Quartz handles retries based on trigger configuration. GB5Trace.MarkFailed() is called in every catch block.

API Endpoints

All endpoints follow the GB5 FastEndpoint pattern. Registered base path: /JobEngine. Gateway prefix: /je/.

Job Definition

GET/JobEngine/GetJobsAll jobs for tenant (MJOBDEFINE + TSCHEDULER join)
GET/JobEngine/GetJobDetail?JobId=Single job with scheduler + webservice details
POST/JobEngine/SaveJobCreate or update — hot-reloads Quartz immediately
DELETE/JobEngine/DeleteJob?JobId=Soft-delete + remove from Quartz
POST/JobEngine/TriggerJob?JobId=Manual immediate fire (not waiting for schedule)
PUT/JobEngine/PauseJob?JobId=STATUS=0 + remove from Quartz
PUT/JobEngine/ResumeJob?JobId=STATUS=1 + re-add to Quartz

Scheduler

GET/JobEngine/GetSchedulersAll TSCHEDULER rows for tenant
POST/JobEngine/SaveSchedulerCreate / update TSCHEDULER row

Execution History

GET/JobEngine/GetExecutionHistoryPaged — filter by JobId, Status, DateFrom, DateTo
GET/JobEngine/GetExecutionDetail?ExecutionId=Full single execution row

Queue + DLQ

GET/JobEngine/GetQueueStatsAggregate queue depth by STATUS + MESSAGETYPE
GET/JobEngine/GetQueuePaged — filter by Status, MessageType
POST/JobEngine/EnqueueMessageInsert TJOBQUEUE row (idempotency key required)
PUT/JobEngine/RequeueItem?QueueId=DLQ → PENDING, reset RETRYATTEMPT
DELETE/JobEngine/CancelQueueItem?QueueId=STATUS=CANCELLED
DELETE/JobEngine/PurgeDLQ?DaysOld=Hard-delete DLQ rows older than N days
GET/JobEngine/GetDashboardSummary: running, pending, 7-day stats, avg durations, DLQ

SignalR Events

Hub: JobEngineHub mapped to /hubs/jobengine. All events broadcast to group jobengine:client:{ClientId}. Transport: LongPolling (WebSocket blocked by custom Login header).

ReceiveJobStarted
(int jobId, string jobName, string correlationId)
Fired by QuartzJobExecutor immediately after inserting TJOBEXECUTION STATUS=RUNNING.
ReceiveJobCompleted
(int jobId, string jobName, long durationMs, string? resultSummary)
Fired after execution SUCCESS. resultSummary is the HTTP response body (truncated).
ReceiveJobFailed
(int jobId, string jobName, string error, int retryAttempt, bool isDlq)
Fired when execution fails. isDlq=true means RETRYATTEMPT exceeded MAXRETRIES.
ReceiveQueueDepth
(int pending, int processing, int dlq, string messageType)
Fired by BaseJobQueueProcessor after each batch claim to update live KPI cards.
LoginDTO in Hub

Always reconstructed from Context.GetHttpContext() claims — never accepted from the client. Group name pattern: jobengine:client:{ClientId}. Group membership added in OnConnectedAsync and removed in OnDisconnectedAsync.

Queue Processor (BaseJobQueueProcessor)

Module-specific processors inherit BaseJobQueueProcessor (in GB5Shared/QueueReader/). They specify which MessageType values they own.

Implemented processors

Processor classModuleMessageTypes
MMJobQueueProcessorMMSLIMPORT, WORKFLOW, TEMPLATECREATION
PAYJobQueueProcessorPayRollSLCASHBACK, LOYALTY, PAYSLIP, PAYRUN

Atomic claim SQL

The claim query uses UPDATE TOP(@BatchSize) … OUTPUT INSERTED.* WITH (ROWLOCK, READPAST) to guarantee no two processor instances ever claim the same row, even when running in a Kubernetes cluster.

Retry backoff schedule

RETRYATTEMPTNEXTATTEMPTAFTER delay
0 (first fail)30 seconds
12 minutes
28 minutes
330 minutes
42 hours
≥ MAXRETRIESSTATUS → DLQ (no retry)

Adding a New Scheduled Job

Two paths depending on who owns the job:

Path A — Admin UI (via Job Define screen)

  1. Create a TSCHEDULER row via Schedule Manager screen (type: CRON / INTERVAL / ONE_TIME).
  2. In Job Define, pick the scheduler, the target MWEBSERVICE URL, and optional URI parameter.
  3. Save → BLL calls IQuartzSyncService.SyncJobAsync() which immediately registers the trigger in Quartz.

Path B — Module-driven registration

BLL calls IJobDefineBLL.RegisterJobAsync() which does an upsert by (SourceModule, SourceObjectId, TenantId):

await _jobDefineBLL.RegisterJobAsync(new JobDefineDTO
{
    JobName       = $"Report:{dto.ReportName}",
    JobCategory   = "REPORT",
    SourceModule  = "REPORTS",
    SourceObjectId = dto.ReportId,
    WebServiceId  = dto.ReportWebServiceId,
    // Scheduler is created inline by RegisterJobAsync if not provided
    TenantId      = login.ClientId,
    Status        = 1
}, login, ct);

Adding a New Queue Message Handler

Implement IJobQueueHandler in the relevant module's BLL or a shared library:

public class MyNewHandler : IJobQueueHandler
{
    public string MessageType => "MY_MESSAGE_TYPE";

    public async Task HandleAsync(string payload, LoginDTO login, CancellationToken ct)
    {
        var dto = JsonSerializer.Deserialize<MyPayloadDTO>(payload)!;
        // ... business logic ...
    }
}

Register in the module's Program.cs:

builder.Services.AddScoped<IJobQueueHandler, MyNewHandler>();

Then add "MY_MESSAGE_TYPE" to the relevant processor's MessageTypes list. The base processor resolves handlers by type at runtime.

Processor Registration in Program.cs

// MMSL/Program.cs
builder.Services.AddHostedService<MMJobQueueProcessor>();

// PayRollSL/Program.cs
builder.Services.AddHostedService<PAYJobQueueProcessor>();
AddHostedService, not AddSingleton

Queue processors must use AddHostedService<T>(). Using AddSingleton registers the class in DI but does NOT start its ExecuteAsync loop.

ActionProcessor Integration

The ActionProcessor system (in FrameworkBLL/ActionProcessor/) handles event-driven action delivery — email, SMS, webhook, and SignalR notifications. It is complementary to (not part of) the Job Engine.

Integration point: when QuartzJobExecutor completes a job, it can publish a SchedulerTaskDTO to the Scheduler.Ready RabbitMQ queue. SchedulerActionGeneratorBLL picks this up, evaluates MEVENTACTION rules (JSONPath conditions), and queues matching delivery actions in TEVENTACTIONRUN + TACTIONOUTBOX.

ActionTypeHandler classStatus
0 — EmailEmailActionHandler (SendGrid / SMTP)✅ Complete
1 — SMSSmsActionHandler (Twilio)✅ Complete
2 — WebhookWebhookActionHandler✅ Complete
3 — In-app notificationNotificationHandler⬜ Skeleton only
4 — ReportReportHandler⬜ Not implemented

OpenTelemetry Instrumentation

Every BLL save method must follow this pattern:

public async Task<string> SaveJobAsync(JobDefineDTO dto, LoginDTO login, CancellationToken ct)
{
    GB5Trace.Step("validate-job-define", new { dto.JobId });
    _validation.NotEmpty(dto.JobName, nameof(dto.JobName));

    var isNew = dto.JobId == 0;
    GB5Trace.Step("save-job-define", new { dto.JobId, isNew });
    var savedId = await _dal.SaveJobAsync(dto, login, ct).ConfigureAwait(false);

    GB5Trace.Step("sync-quartz", new { savedId });
    await _quartzSync.SyncJobAsync(savedId, ct).ConfigureAwait(false);

    return isNew
        ? $"{SuccessResponse.SaveSuccessMessage} {savedId}"
        : SuccessResponse.UpdateSuccess;
}
catch (Exception ex)
{
    GB5Trace.MarkFailed("save-job-define-failed", ex);
    _logger.LogError(ex, "SaveJob failed for {JobName}", dto.JobName);
    throw;
}

New Job / Handler Checklist

ItemWhere
✅ Create DTO in JobEngineDAL/DTOs/JobEngineDAL
✅ Add SQL constants in relevant *QB.csJobEngineDAL/Query/
✅ Add interface + implementation in DALJobEngineDAL/Interfaces/ + Implementations/
✅ Add interface + implementation in BLLJobEngineBLL/Interfaces/ + Implementations/
✅ Add FastEndpoint + Parameters classJobEngineSL/EndPoints/
✅ Return null from GetCacheKey() on mutationsFastEndpoint Configure()
✅ GB5Trace.Step() before validate, before DAL, before event publishBLL method
✅ GB5Trace.MarkFailed() in every catchBLL catch block
✅ ConfigureAwait(false) on all awaits in BLL/DALBLL + DAL
✅ All queries filter by TENANTID or DatabaseNameQB SQL strings
✅ IJobQueueHandler registered as Scoped in module Program.csModule SL Program.cs
✅ AddHostedService<T>() for processor registrationModule SL Program.cs