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:
- Scheduled execution — Quartz.NET clustered persistent store drives cron / interval / one-time jobs, calling HTTP endpoints (MWEBSERVICE) or Dapr topics.
- Async message queue — Persistent
TJOBQUEUEtable replaces MSMQ. Module-specificBackgroundServicesubclasses claim, process, and retry batches.
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)
| Table | New columns added |
|---|---|
MJOBDEFINE | MAXRETRIES, TIMEOUTSECONDS, NEXTRUNTIME, ISCONCURRENT, JOBCATEGORY, SOURCEMODULE, SOURCEOBJECTID, DAPRPUBSUBTOPIC, audit columns |
TSCHEDULER | SCHEDULERTYPE (CRON/INTERVAL/ONE_TIME), INTERVALSECONDS, SCHEDULEDDATE, DESCRIPTION, SCHEDULERSTATUS, audit columns |
New tables (created by migration)
| Column | Type | Notes |
|---|---|---|
| EXECUTIONID | BIGINT IDENTITY PK | |
| JOBID | INT NULL | FK to MJOBDEFINE; NULL for queue-only runs |
| TENANTID | INT NOT NULL | Tenant isolation |
| JOBNAME | NVARCHAR(200) | |
| TRIGGERTYPE | NVARCHAR(20) | SCHEDULED | MANUAL | QUEUE | API |
| STATUS | NVARCHAR(20) | QUEUED | RUNNING | SUCCESS | FAILED | TIMEOUT | CANCELLED |
| STARTEDAT / COMPLETEDAT | DATETIME2 | |
| DURATIONMS | BIGINT NULL | Populated on completion |
| CORRELATIONID | NVARCHAR(100) | GUID set by QuartzJobExecutor; carried to SignalR events |
| ERRORDETAILS | NVARCHAR(MAX) | Full exception string on failure |
| HOSTINSTANCE | NVARCHAR(100) | Pod hostname — for cluster attribution |
| Column | Type | Notes |
|---|---|---|
| QUEUEID | BIGINT IDENTITY PK | |
| TENANTID | INT | |
| MESSAGETYPE | NVARCHAR(100) | CASHBACK | LOYALTY | PAYSLIP | IMPORT | WORKFLOW | … |
| PAYLOAD | NVARCHAR(MAX) | JSON body; handler-specific schema |
| PRIORITY | TINYINT default 5 | Higher = processed first |
| STATUS | NVARCHAR(20) | PENDING | PROCESSING | COMPLETED | FAILED | DLQ | CANCELLED |
| MAXRETRIES | INT default 3 | |
| RETRYATTEMPT | INT | Incremented on each claim |
| NEXTATTEMPTAFTER | DATETIME2 NULL | Backoff ceiling: 30s → 2m → 8m → 30m → 2h → DLQ |
| MESSAGEID | NVARCHAR(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).
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.
| Method | When called | What it does |
|---|---|---|
LoadAllJobsAsync() | Startup | Reads all MJOBDEFINE STATUS=1 and schedules each in Quartz. Called from Program.cs after app.Build(). |
SyncJobAsync(jobId) | After SaveJob, PauseJob, ResumeJob | If Status≠1: deletes Quartz job. Otherwise: reschedules from fresh DB read. |
RemoveJobAsync(jobId) | After DeleteJob | Removes trigger and job detail from Quartz cluster store. |
TriggerJobNowAsync(jobId) | After TriggerJob endpoint | Fires the job immediately regardless of its next scheduled time. |
QuartzJobExecutor — per-execution IJob
A new instance is created for every execution. Execution flow:
IJobDefineDAL.GetJobDetailAsync(). JobId and TenantId come from Quartz JobDataMap.jobengine:client:{ClientId}.DaprPubSubTopic set → Dapr publish. Otherwise → HTTP call to UriTemplate (GET or POST based on METHODTYPE). {Param} in template substituted with UriParameterValue.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
Scheduler
Execution History
Queue + 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).
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 class | Module | MessageTypes |
|---|---|---|
MMJobQueueProcessor | MMSL | IMPORT, WORKFLOW, TEMPLATECREATION |
PAYJobQueueProcessor | PayRollSL | CASHBACK, 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
| RETRYATTEMPT | NEXTATTEMPTAFTER delay |
|---|---|
| 0 (first fail) | 30 seconds |
| 1 | 2 minutes |
| 2 | 8 minutes |
| 3 | 30 minutes |
| 4 | 2 hours |
| ≥ MAXRETRIES | STATUS → DLQ (no retry) |
Adding a New Scheduled Job
Two paths depending on who owns the job:
Path A — Admin UI (via Job Define screen)
- Create a
TSCHEDULERrow via Schedule Manager screen (type: CRON / INTERVAL / ONE_TIME). - In Job Define, pick the scheduler, the target
MWEBSERVICEURL, and optional URI parameter. - 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>();
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.
| ActionType | Handler class | Status |
|---|---|---|
| 0 — Email | EmailActionHandler (SendGrid / SMTP) | ✅ Complete |
| 1 — SMS | SmsActionHandler (Twilio) | ✅ Complete |
| 2 — Webhook | WebhookActionHandler | ✅ Complete |
| 3 — In-app notification | NotificationHandler | ⬜ Skeleton only |
| 4 — Report | ReportHandler | ⬜ 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
| Item | Where |
|---|---|
| ✅ Create DTO in JobEngineDAL/DTOs/ | JobEngineDAL |
| ✅ Add SQL constants in relevant *QB.cs | JobEngineDAL/Query/ |
| ✅ Add interface + implementation in DAL | JobEngineDAL/Interfaces/ + Implementations/ |
| ✅ Add interface + implementation in BLL | JobEngineBLL/Interfaces/ + Implementations/ |
| ✅ Add FastEndpoint + Parameters class | JobEngineSL/EndPoints/ |
| ✅ Return null from GetCacheKey() on mutations | FastEndpoint Configure() |
| ✅ GB5Trace.Step() before validate, before DAL, before event publish | BLL method |
| ✅ GB5Trace.MarkFailed() in every catch | BLL catch block |
| ✅ ConfigureAwait(false) on all awaits in BLL/DAL | BLL + DAL |
| ✅ All queries filter by TENANTID or DatabaseName | QB SQL strings |
| ✅ IJobQueueHandler registered as Scoped in module Program.cs | Module SL Program.cs |
| ✅ AddHostedService<T>() for processor registration | Module SL Program.cs |