NexJob.Trigger.Salesforce
5.9.0
dotnet add package NexJob.Trigger.Salesforce --version 5.9.0
NuGet\Install-Package NexJob.Trigger.Salesforce -Version 5.9.0
<PackageReference Include="NexJob.Trigger.Salesforce" Version="5.9.0" />
<PackageVersion Include="NexJob.Trigger.Salesforce" Version="5.9.0" />
<PackageReference Include="NexJob.Trigger.Salesforce" />
paket add NexJob.Trigger.Salesforce --version 5.9.0
#r "nuget: NexJob.Trigger.Salesforce, 5.9.0"
#:package NexJob.Trigger.Salesforce@5.9.0
#addin nuget:?package=NexJob.Trigger.Salesforce&version=5.9.0
#tool nuget:?package=NexJob.Trigger.Salesforce&version=5.9.0
NexJob.Trigger.Salesforce
Salesforce Pub/Sub API trigger for NexJob. 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.
Features
- Salesforce Pub/Sub API (gRPC): Native streaming support using official Protocol Buffers and bidirectional flow control.
- Apache Avro Deserialization: In-memory caching of schema fingerprints (
SchemaId) and binary decoding to JSON. - Resilient Replay ID Checkpointing:
IReplayIdStorewith atomic file-based persistence (FileReplayIdStore) and memory-only storage (InMemoryReplayIdStore). - Configurable Fallback Policies:
ReplayFallbackPolicy.FailFast,ResetToLatest, andResetToEarliesthandle expired offsets gracefully. - OAuth2 Token Caching: Automatic Client Credentials flow with thread-safe token caching and refresh ahead of expiration.
- W3C Distributed Tracing: Automatic extraction of
traceparentheaders mapped directly toJobRecord.TraceParent. - Operational Visibility: Integrates with
IListenerRegistryto report connection status (Starting,Listening,Reconnecting,Faulted,Stopped) directly to Dashboard/listenersand Cluster Topology Map. - All 5 Trigger Guarantees: Guaranteed at-least-once processing, offset commit only after enqueue, idempotency via native event IDs, and dead-letter routing.
Installation
dotnet add package NexJob.Trigger.Salesforce
Quick Start
1. Basic Registration
Register the trigger to consume from a CDC topic using the default payload job or a custom job:
using NexJob;
using NexJob.Trigger.Salesforce;
var builder = WebApplication.CreateBuilder(args);
builder.Services.AddNexJob();
builder.Services.AddSalesforceTrigger(options =>
{
options.Topic = "/data/ChangeEvents"; // Standard CDC topic
options.ClientId = builder.Configuration["Salesforce:ClientId"]!;
options.ClientSecret = builder.Configuration["Salesforce:ClientSecret"]!;
options.TargetQueue = "salesforce-events";
});
2. Custom Strongly-Typed Job Handler
using NexJob;
using NexJob.Trigger.Salesforce;
builder.Services.AddSalesforceTrigger<ProcessAccountChangeJob>(options =>
{
options.Topic = "/data/AccountChangeEvent";
options.ClientId = builder.Configuration["Salesforce:ClientId"]!;
options.ClientSecret = builder.Configuration["Salesforce:ClientSecret"]!;
options.FallbackPolicy = ReplayFallbackPolicy.ResetToLatest;
});
public sealed class ProcessAccountChangeJob : IJob<string>
{
public async Task ExecuteAsync(string payloadJson, CancellationToken cancellationToken)
{
// payloadJson contains decoded Avro fields as standard JSON
using var doc = JsonDocument.Parse(payloadJson);
var root = doc.RootElement;
// Process account modification...
}
}
3. Fluent NexJobBuilder API
builder.Services.AddNexJob(options =>
{
options.Queues = ["default", "salesforce-events"];
})
.AddSalesforceTrigger(options =>
{
options.Topic = "/event/OrderEvent__e";
options.ClientId = "3MVG9...";
options.ClientSecret = "secret...";
});
Configuration Options
| Property | Type | Description | Default |
|---|---|---|---|
Topic |
string |
Salesforce topic name starting with / (e.g. /data/ChangeEvents, /event/OrderEvent__e) |
Required |
ClientId |
string |
Connected App OAuth2 Client ID (Consumer Key) | Required |
ClientSecret |
string |
Connected App OAuth2 Client Secret (Consumer Secret) | Required |
AuthEndpoint |
string |
Salesforce OAuth2 token endpoint URL | https://login.salesforce.com/services/oauth2/token |
PubSubEndpoint |
string |
Salesforce Pub/Sub gRPC endpoint | api.pubsub.salesforce.com:7443 |
TargetQueue |
string |
Target NexJob queue for enqueued jobs | salesforce-events |
JobPriority |
JobPriority |
Execution priority for enqueued jobs | Normal |
TenantId |
string? |
18-character Salesforce organization ID (derived if omitted) | null |
ReplayPreset |
SalesforceReplayPreset |
Starting offset when no checkpoint exists (Latest, Earliest, Custom) |
Latest |
CustomReplayId |
byte[]? |
Specific binary Replay ID when preset is Custom |
null |
FallbackPolicy |
ReplayFallbackPolicy |
Action when saved Replay ID is expired (FailFast, ResetToLatest, ResetToEarliest) |
FailFast |
BatchSize |
int |
Flow control batch size requested per stream roundtrip | 100 |
ReplayStoreDirectory |
string |
Disk directory for FileReplayIdStore checkpoints |
./.nexjob/salesforce |
DeadLetterQueue |
string? |
Queue to route events when initial enqueue fails | null |
Replay ID Checkpointing
The trigger maintains your position in the Salesforce event stream by checkpointing the binary ReplayId received with each event:
FileReplayIdStore(Default): Writes Replay IDs to disk using atomic rename operations (File.Move(..., overwrite: true)), ensuring zero corruption across abrupt restarts.InMemoryReplayIdStore: Thread-safe in-memory store suitable for containerized workers using external offsets or testing.- Custom
IReplayIdStore: Implement theIReplayIdStoreinterface to persist checkpoints to PostgreSQL, Redis, DynamoDB, or any external store.
Replay Fallback Policies
Salesforce event retention is typically 72 hours. If a worker is offline longer than the retention window, the saved Replay ID expires:
FailFast(Default): Logs a critical alert and terminates the service to prevent silent event loss.ResetToLatest: Clears the invalid checkpoint and resumes streaming from current real-time events.ResetToEarliest: Clears the invalid checkpoint and replays all events still available in the retention window.
Core Guarantees
- Never Silently Drop: All stream errors, deserialization failures, and enqueue faults are logged and routed to the configured
DeadLetterQueue. - Idempotency: Broker event ID (
ConsumerEvent.Event.Id) is assigned toJobRecord.IdempotencyKey. - Trace Propagation: Extracts W3C
traceparentfromConsumerEvent.Event.Headersdirectly intoJobRecord.TraceParent. - Signal After Enqueue: Dispatcher wake-up notification is handled automatically by
IScheduler.EnqueueAsync. - Ack Only After Enqueue: The binary Replay ID is committed to
IReplayIdStoreonly after successful storage persistence.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 was computed. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 was computed. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net8.0
- Apache.Avro (>= 1.12.0)
- Google.Protobuf (>= 3.26.1)
- Grpc.Net.Client (>= 2.62.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 8.0.1)
- Microsoft.Extensions.Http (>= 8.0.1)
- Microsoft.Extensions.Logging.Abstractions (>= 8.0.3)
- Microsoft.Extensions.Options (>= 8.0.2)
- Microsoft.Extensions.Options.DataAnnotations (>= 8.0.0)
- NexJob (>= 5.9.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 5.9.0 | 0 | 10/5/2026 |
| 5.8.0 | 41 | 10/4/2026 |
| 5.7.0 | 59 | 10/1/2026 |
| 5.6.2 | 59 | 10/1/2026 |
| 5.6.1 | 75 | 9/30/2026 |
| 5.6.0 | 68 | 9/30/2026 |
| 5.5.0 | 84 | 9/25/2026 |
| 5.4.1 | 81 | 9/24/2026 |
| 5.4.0 | 83 | 9/24/2026 |
| 5.3.0 | 102 | 9/21/2026 |
| 5.2.0 | 101 | 9/20/2026 |
| 5.1.0 | 97 | 9/17/2026 |
| 5.0.0 | 102 | 9/17/2026 |
| 0.0.0-alpha.0 | 53 | 9/17/2026 |