External Triggers
NexJob triggers bridge external message brokers and your background job pipeline. Each trigger package subscribes to a broker source — a queue, topic, or channel — and converts every incoming message into a persisted NexJob job, delivering it to the dispatcher for execution. The pipeline is: broker → trigger → JobRecordFactory → IScheduler.EnqueueAsync → dispatcher. Your job code receives the raw message body as a string; it decides how to parse it.
Broker Guarantees¶
All NexJob triggers satisfy five core guarantees regardless of the broker:
- At-least-once delivery — Messages are never silently dropped before enqueue.
- Idempotency — The broker's native message identity is used as the
idempotencyKey, preventing duplicate jobs on redelivery. - Trace propagation — W3C
traceparentis extracted from message headers or attributes and attached to the job record for end-to-end distributed tracing. - Signal after enqueue — Enqueueing a job automatically signals the dispatcher; no manual wake-up is needed.
- Ack only after success — Messages are acknowledged only after
IScheduler.EnqueueAsynccompletes successfully. This means the job is visible before the acknowledgement (or the Salesforce replay id) is written. A stop or crash in that short window delivers the message again, and the idempotency key keeps it from becoming a second job. A Salesforce commit that follows a successful enqueue is not cancelled by a stop.
Message Contract¶
The message body is the job input. Because broker triggers are generic, the input type is always string. Your job handler implements IJob<string> and deserializes the body as needed:
public sealed class ProcessOrderJob : IJob<string>
{
public async Task ExecuteAsync(string messageBody, CancellationToken ct)
{
// messageBody contains the raw broker payload (JSON, XML, CSV, or plain text)
var order = JsonSerializer.Deserialize<OrderDto>(messageBody);
// Process...
}
}
Note
Non-JSON message bodies (XML, CSV, plain text) are passed verbatim to your job. The trigger does not attempt to parse or transform the payload.
An optional traceparent header or attribute carries the W3C trace context for distributed tracing.
Job Type Resolution¶
Most triggers resolve which job to run using the same three-step precedence:
nexjob.job_typeheader/attribute on the message — set by the publisher as the assembly-qualified name (e.g.MyApp.Jobs.ProcessOrderJob, MyApp).- Subscriber's configured job type —
options.JobType, or the generic overloadAdd{Broker}Trigger<TJob>()which sets it for you. - Neither present — permanent failure: the message can never become a job and is dead-lettered or committed without redelivery (see Error Handling below).
Idempotency Keys by Broker¶
| Trigger | Idempotency Key |
|---|---|
| Kafka | kafka:{topic}:{partition}:{offset} |
| RabbitMQ | MessageId (no deduplication when blank) |
| Azure Service Bus | MessageId |
| Google Pub/Sub | Message ID |
| AWS SQS | Message ID |
| Salesforce Pub/Sub API | Event ID (Replay ID as hex; timestamp when absent) |
| Salesforce Streaming API | {channel}:{eventId} (or {channel}:{replayId} when absent) |
Supported Triggers¶
Install the package¶
dotnet add package NexJob.Trigger.AzureServiceBus
Register the trigger¶
Use the generic overload to bind all messages on a queue or topic to a single job type:
using NexJob.Trigger.AzureServiceBus;
builder.Services.AddNexJobAzureServiceBusTrigger<ProcessOrderJob>(options =>
{
options.ConnectionString = "Endpoint=sb://...";
options.QueueOrTopicName = "orders";
});
To route messages to different job types based on the nexjob.job_type application property, omit the generic argument:
builder.Services.AddNexJobAzureServiceBusTrigger(options =>
{
options.ConnectionString = "Endpoint=sb://...";
options.QueueOrTopicName = "my-topic";
options.SubscriptionName = "my-sub"; // required for topics
});
The idempotency key is the Service Bus MessageId. Transient failures abandon the message so it is redelivered; MaxDeliveryCount on the entity controls when it is dead-lettered. Permanent failures (unresolvable job type) dead-letter the message immediately.
Install the package¶
dotnet add package NexJob.Trigger.AwsSqs
Register the trigger¶
SQS message attributes are not read for the job type — you must configure it explicitly:
using NexJob.Trigger.AwsSqs;
// Preferred: generic overload binds queue to one job type
builder.Services.AddNexJobAwsSqsTrigger<ProcessOrderJob>(options =>
{
options.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789/my-queue";
});
// Equivalent explicit form:
// builder.Services.AddNexJobAwsSqsTrigger(options =>
// {
// options.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789/my-queue";
// options.JobName = typeof(ProcessOrderJob).AssemblyQualifiedName!;
// });
The idempotency key is the SQS message ID. Unacknowledged messages become visible again after the queue's visibility timeout, so transient failures are retried automatically by SQS.
Install the package¶
dotnet add package NexJob.Trigger.GooglePubSub
Register the trigger¶
using NexJob.Trigger.GooglePubSub;
// Preferred: generic overload
builder.Services.AddNexJobGooglePubSubTrigger<ProcessOrderJob>(options =>
{
options.ProjectId = "my-project";
options.SubscriptionId = "my-subscription";
});
// Dynamic routing via nexjob.job_type attribute:
// builder.Services.AddNexJobGooglePubSubTrigger(options =>
// {
// options.ProjectId = "my-project";
// options.SubscriptionId = "my-subscription";
// options.JobType = typeof(DefaultEventJob).AssemblyQualifiedName;
// });
The idempotency key is the Pub/Sub message ID assigned by the service. Messages remain unacknowledged on transient failures and are redelivered by Pub/Sub.
Install the package¶
dotnet add package NexJob.Trigger.Salesforce
Register the trigger¶
The Salesforce Pub/Sub trigger consumes Change Data Capture (CDC) events and custom Platform Events over bidirectional gRPC streams. It decodes Apache Avro binary payloads to JSON automatically:
using NexJob.Trigger.Salesforce;
// Use the built-in SalesforceEventJob handler
builder.Services.AddNexJob()
.AddSalesforceTrigger(options =>
{
options.Topic = "/data/ChangeEvents";
options.ClientId = "3MVG9...";
options.ClientSecret = "secret...";
options.TargetQueue = "salesforce-events";
options.ReplayPreset = SalesforceReplayPreset.Latest;
options.FallbackPolicy = ReplayFallbackPolicy.ResetToLatest;
});
// Or bind to a strongly typed custom job
builder.Services.AddNexJob()
.AddSalesforceTrigger<ProcessAccountChangeJob>(options =>
{
options.Topic = "/data/AccountChangeEvent";
options.ClientId = "3MVG9...";
options.ClientSecret = "secret...";
});
Key capabilities:
- Bi-directional gRPC streaming using the official Salesforce Pub/Sub API protobufs with flow control.
- Apache Avro decoding with in-memory schema caching via ISalesforceSchemaService.
- Replay ID checkpointing via IReplayIdStore (atomic file-based or in-memory).
- Resilient fallback policies: FailFast, ResetToLatest, and ResetToEarliest for expired offsets.
- OAuth2 token caching with automatic refresh ahead of expiration.
Install the package¶
dotnet add package NexJob.Trigger.SalesforceStreaming
Register the trigger¶
The Salesforce Streaming trigger consumes PushTopic events, CDC, and Platform Events over HTTP long-polling using the CometD/Bayeux protocol. Use this for legacy environments without gRPC/HTTP2:
using NexJob.Trigger.SalesforceStreaming;
// Use the built-in SalesforceStreamingEventJob handler
builder.Services.AddNexJob()
.AddSalesforceStreamingTrigger(options =>
{
options.Channel = "/data/Order__ChangeEvent";
options.Authentication.AuthType = SalesforceStreamingAuthType.OAuth2ClientCredentials;
options.Authentication.AuthEndpoint = "https://login.salesforce.com/services/oauth2/token";
options.Authentication.ClientId = "3MVG9...";
options.Authentication.ClientSecret = "secret...";
options.TargetQueue = "salesforce-events";
options.ReplayPreset = SalesforceStreamingReplayPreset.Latest;
});
// Bind to a custom job with username/password auth
builder.Services.AddNexJob()
.AddSalesforceStreamingTrigger<ProcessSalesforceOrderJob>(options =>
{
options.Channel = "/topic/InvoiceUpdates";
options.Authentication.AuthType = SalesforceStreamingAuthType.OAuth2UsernamePassword;
options.Authentication.ClientId = "3MVG9...";
options.Authentication.ClientSecret = "secret...";
options.Authentication.Username = "integration@company.com";
options.Authentication.Password = "Password123";
options.Authentication.SecurityToken = "TokenXYZ";
options.DeadLetterQueue = "salesforce-dlq";
});
Key capabilities:
- CometD/Bayeux protocol with replay extension and graceful disconnect.
- Multi-auth: OAuth 2.0 Username-Password, OAuth 2.0 Client Credentials, and direct Session ID / Bearer token.
- Replay ID checkpointing via IStreamingReplayIdStore (file-based or in-memory).
- Session expiry resilience: automatic token invalidation and re-handshake on 403::Unknown client.
- Exponential backoff reconnection loop with configurable delays and multipliers.
Error Handling¶
Trigger enqueue failures fall into three categories:
Permanent failure (missing nexjob.job_type, malformed payload): The message can never become a job and redelivery is pointless. No job is created. The trigger dead-letters the message (or commits without redelivery if no dead-letter destination is configured) so the consumer does not stall.
Job type not found in DI: The trigger enqueues the job record successfully. The dispatcher fails the job on execution with a clear error. NexJob's normal retry policy applies.
Transient failure (storage unavailable, network error, timeout): The message is never lost and never dead-lettered. Kafka retries in place with 1 s → 2 s → 5 s → 10 s → 20 s → 30 s backoff without consuming the next record; RabbitMQ nacks with requeue: true after a one-second pause; Azure Service Bus abandons the message for redelivery; SQS and Pub/Sub leave the message unacknowledged. Combined with idempotency keys, this prevents duplicate jobs even under partial failures.
Live Listener Registry¶
NexJob maintains a centralized, thread-safe IListenerRegistry that tracks the real-time operational status of every registered trigger. Each trigger updates its status throughout its lifecycle:
| Status | Meaning |
|---|---|
Starting |
Initializing the broker connection and subscribing. |
Listening |
Successfully connected and actively consuming messages. |
Reconnecting |
Transient connection loss; undergoing automatic reconnect. |
Faulted |
Fatal unrecoverable broker error (includes error message). |
Stopped |
Host shutdown or graceful deregistration. |
All registered triggers and their live states are visible in the dashboard at the /listeners route and summarized in the Cluster Pipeline Topology Map on the overview page.
Try it¶
The NexJob.Sample.CloudTriggers sample wires AWS SQS, Azure Service Bus, Google Pub/Sub and Salesforce triggers, with /simulate/* endpoints to produce messages without the real services.