External Triggers¶
NexJob supports external triggers — broker messages that automatically enqueue NexJob jobs. This enables event-driven job scheduling from various message brokers.
How Triggers Work¶
Triggers follow a standard pipeline:
[broker] → trigger package → JobRecordFactory → IScheduler.EnqueueAsync → dispatcher
- The Trigger Package consumes a message from the broker.
- It resolves the Job Type (see Which job runs) and extracts the Trace Context from headers/attributes.
- It uses
JobRecordFactoryto build aJobRecordusing the message body as input. The body is stored as a JSON string, so a trigger-boundIJob<string>receives it verbatim as text (JSON, XML, CSV or plain text); the job decides how to parse it. - It calls
IScheduler.EnqueueAsyncto persist the job. - It Acknowledge (Ack) the message only after a successful enqueue.
Message Contract¶
The message Body is used as the job input. Since broker triggers are generic, the input type is always string (usually JSON). Your job handler should deserialize the body as needed.
The optional traceparent header/attribute carries the W3C trace context for distributed tracing.
Which job runs for a message¶
Most triggers resolve the job with the same three-step precedence, so publishers do not have to know NexJob types:
- The
nexjob.job_typeheader/attribute of the message (assembly-qualified name, for exampleMyApp.Jobs.ProcessOrderJob, MyApp), if present. - The subscriber's configured job:
options.JobType, or the generic overloadAdd{Broker}Trigger<TJob>(), which sets it for you. - Neither is present: the message can never become a job. It is treated as a permanent failure (see Error handling).
| Trigger | Resolves the job from | Idempotency key |
|---|---|---|
| Kafka | header nexjob.job_type, then options.JobType |
kafka:{topic}:{partition}:{offset} |
| RabbitMQ | header nexjob.job_type, then options.JobType |
MessageId (none when blank) |
| Azure Service Bus | application property nexjob.job_type, then options.JobType |
MessageId |
| Google Pub/Sub | attribute nexjob.job_type, then options.JobType |
message id |
| AWS SQS | configured only: options.JobName, or AddNexJobAwsSqsTrigger<TJob>() (the message attributes are not consulted) |
message id |
| Salesforce Pub/Sub API | configured options.JobType, otherwise the built-in SalesforceEventJob |
event id (replay id, hex, when absent) |
| Salesforce Streaming | configured options.JobType, otherwise the built-in SalesforceStreamingEventJob |
{channel}:{eventId} ({channel}:{replayId} when absent) |
Broker Guarantees¶
All NexJob triggers satisfy 5 core guarantees:
1. At-least-once delivery: Messages are never silently dropped before enqueue.
2. Idempotency: Uses the broker's native message identity as idempotencyKey to prevent duplicate jobs (see the table above for what each broker uses).
3. Trace propagation: Extracts traceparent from headers to maintain the trace across systems.
4. Signal after enqueue: Enqueueing a job automatically signals the dispatcher (no manual wake-up needed).
5. Ack only after success: Messages are acknowledged only after IScheduler.EnqueueAsync completes successfully.
Azure Service Bus¶
Installation:
dotnet add package NexJob.Trigger.AzureServiceBus
Usage:
using NexJob.Trigger.AzureServiceBus;
builder.Services.AddNexJobAzureServiceBusTrigger(options =>
{
options.ConnectionString = "Endpoint=sb://...";
options.QueueOrTopicName = "my-queue"; // or my-topic
options.SubscriptionName = "my-sub"; // required for topics
});
To run one job for every message regardless of the nexjob.job_type property, use the generic overload (or set options.JobType):
builder.Services.AddNexJobAzureServiceBusTrigger<ProcessOrderJob>(options =>
{
options.ConnectionString = "Endpoint=sb://...";
options.QueueOrTopicName = "orders";
});
AWS SQS¶
Installation:
dotnet add package NexJob.Trigger.AwsSqs
Usage:
using NexJob.Trigger.AwsSqs;
builder.Services.AddNexJobAwsSqsTrigger(options =>
{
options.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789/my-queue";
options.JobName = typeof(ProcessOrderJob).AssemblyQualifiedName!; // SQS does not read a job type from the message
});
// Equivalent, and clearer: bind the trigger to one job type
builder.Services.AddNexJobAwsSqsTrigger<ProcessOrderJob>(options =>
{
options.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789/my-queue";
});
RabbitMQ¶
[!NOTE] RabbitMQ capabilities have evolved from a simple trigger into a dedicated, unified package:
NexJob.RabbitMQ. It provides both RabbitMQ Triggers (Consumers) and a Resilient Outbox Producer with Publisher Confirms. For full configuration, outbox producer examples, and environment variables, see the dedicated guide: RabbitMQ Integration (21-RabbitMQ.md).
Installation:
dotnet add package NexJob.RabbitMQ
Usage (Consumer / Trigger):
using NexJob;
builder.Services.AddNexJob()
.AddRabbitMqTrigger(options =>
{
options.HostName = "localhost";
options.QueueName = "nexjob-trigger";
options.UserName = "guest";
options.Password = "guest";
});
AddNexJobRabbitMqTrigger remains supported for backward compatibility).
Kafka¶
[!NOTE] Kafka capabilities have evolved from a simple trigger into a dedicated, unified package:
NexJob.Kafka. It provides both Kafka Triggers (Consumers) and a Resilient Outbox Producer. For full configuration, environment variables, and producer patterns, see the dedicated guide: Kafka Integration (20-Kafka.md).
Installation:
dotnet add package NexJob.Kafka
Usage (Consumer / Trigger):
using NexJob.Kafka;
builder.Services.AddNexJob()
.AddKafkaTrigger(options =>
{
options.BootstrapServers = "localhost:9092";
options.Topic = "nexjob-jobs";
options.GroupId = "nexjob-consumer-group";
});
AddNexJobKafkaTrigger remains supported for backward compatibility).
Google Pub/Sub¶
Installation:
dotnet add package NexJob.Trigger.GooglePubSub
Usage:
using NexJob.Trigger.GooglePubSub;
builder.Services.AddNexJobGooglePubSubTrigger(options =>
{
options.ProjectId = "my-project";
options.SubscriptionId = "my-subscription";
// options.JobType = typeof(ProcessOrderJob).AssemblyQualifiedName; // used when a message has no nexjob.job_type attribute
});
The generic overload AddNexJobGooglePubSubTrigger<TJob>() sets that fallback job type for you.
Salesforce Pub/Sub API¶
Installation:
dotnet add package NexJob.Trigger.Salesforce
Consumes Salesforce Change Data Capture (CDC) events and custom Platform Events over high-throughput bidirectional gRPC streams, automatically decodes Apache Avro binary payloads to JSON, manages Replay ID checkpointing, and enqueues background jobs with zero message loss.
Usage:
using NexJob;
using NexJob.Trigger.Salesforce;
// Register trigger with default SalesforceEventJob
builder.Services.AddNexJob()
.AddSalesforceTrigger(options =>
{
options.Topic = "/data/ChangeEvents"; // Standard CDC or Platform Event topic
options.ClientId = "3MVG9...";
options.ClientSecret = "secret...";
options.TargetQueue = "salesforce-events";
options.ReplayPreset = SalesforceReplayPreset.Latest;
options.FallbackPolicy = ReplayFallbackPolicy.ResetToLatest;
});
// Or register with a strongly typed custom job
builder.Services.AddNexJob()
.AddSalesforceTrigger<ProcessAccountChangeJob>(options =>
{
options.Topic = "/data/AccountChangeEvent";
options.ClientId = "3MVG9...";
options.ClientSecret = "secret...";
});
Key features:
- Bi-directional gRPC Streaming: Uses official Salesforce Pub/Sub API protobufs and flow control.
- Apache Avro binary decoding: In-memory schema caching (ISalesforceSchemaService) and decoding into JSON payloads.
- Replay ID Checkpointing: IReplayIdStore with atomic file-based persistence (FileReplayIdStore) and in-memory store (InMemoryReplayIdStore).
- Resilient Fallback Policies: ReplayFallbackPolicy.FailFast, ResetToLatest, and ResetToEarliest handle expired offsets gracefully.
- OAuth2 Token Caching: Automatic token acquisition, caching, and refresh ahead of expiration.
- Distributed Tracing: Extracts W3C traceparent from event headers into JobRecord.TraceParent.
Salesforce Streaming API (Legacy / CometD)¶
Installation:
dotnet add package NexJob.Trigger.SalesforceStreaming
Consumes Salesforce PushTopic events, Change Data Capture (CDC), and Platform Events over HTTP long-polling using the CometD/Bayeux protocol. Designed for legacy environments and orgs connecting via standard Bayeux endpoints without gRPC/HTTP2 requirements.
Usage:
using NexJob;
using NexJob.Trigger.SalesforceStreaming;
// Register trigger with default SalesforceStreamingEventJob
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;
});
// Or register with a custom strongly-typed job
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 features:
- CometD/Bayeux Protocol: Standard Bayeux handshake, subscription with replay extension, long-polling connect loop, and graceful disconnect.
- Multi-Authentication Support: OAuth 2.0 Username-Password, OAuth 2.0 Client Credentials, and direct Session ID / Bearer token.
- Replay ID Checkpointing: IStreamingReplayIdStore with atomic file-based persistence (FileStreamingReplayIdStore) and memory store (InMemoryStreamingReplayIdStore).
- Session Expiry Resilience: Automatic token invalidation and re-handshake upon 403::Unknown client session expiry.
- Exponential Backoff: Resilient reconnection loop with configurable delays and backoff multipliers.
- All 5 Trigger Guarantees: Guaranteed at-least-once enqueue, broker-native idempotency keys, W3C traceparent propagation, dispatcher signal, and commit replay ID only after enqueue.
Error handling¶
Enqueue errors are classified by the Kafka, RabbitMQ and Azure Service Bus triggers:
Permanent failure (missing nexjob.job_type, malformed payload):
The message can never become a job, so redelivery is pointless. No job is created and the message is not
redelivered: Kafka produces it to the dead-letter topic and commits (or, with no topic configured, logs an
Error and commits so the partition is not blocked), RabbitMQ nacks it with requeue: false (dead-letter
exchange if configured) and Azure Service Bus dead-letters it.
Job type not found in DI: The trigger enqueues the job record. The dispatcher will fail the job on execution with a clear error. Retries apply normally.
Transient failure (storage unavailable, network error, timeout):
The message is never lost and never dead-lettered because of it. Kafka retries the same record in place with a
1 s, 2 s, 5 s, ... 30 s backoff and does not consume the next record meanwhile; RabbitMQ nacks with
requeue: true after a one-second pause; Azure Service Bus abandons the message so it is delivered again
(MaxDeliveryCount decides when it is dead-lettered if the failure never clears); SQS and Pub/Sub leave the
message unacknowledged. Combined with idempotency keys, this prevents duplicate jobs even under partial failures.
A job that was created from a message and then exhausted its retries is not sent to the broker's dead-letter: the message was acknowledged when it became a job, so the failure stays inside NexJob (Failed in the dashboard). To forward a copy to a Kafka topic or a RabbitMQ exchange, see Forwarding a dead-lettered job.
Active Triggers & Listener Registry (IListenerRegistry)¶
NexJob includes a centralized, thread-safe registry (IListenerRegistry) that tracks the real-time operational status of all registered broker consumers and triggers.
Each trigger registers its endpoint, broker type, consumer group, and dynamically updates its status throughout its lifecycle:
- Starting: Initializing broker connection and subscribing.
- Listening: Successfully connected and actively polling/consuming messages.
- Reconnecting: Transient connection loss or broker error; undergoing automatic reconnect/retry.
- Faulted: Fatal unrecoverable broker error (includes error message).
- Stopped: Host shutdown or graceful deregistration.
Dashboard Visibility¶
All registered triggers and their live connection states, endpoints, target queues, and uptimes are rendered in the dashboard at the /listeners route and summarized in the Cluster Pipeline Topology Map on the Overview page.
Next Steps¶
- Idempotency — Learn about
DuplicatePolicy - OpenTelemetry — Trace jobs across brokers
- Storage Providers — Where jobs are persisted