Multi-Server Concurrency

Trax.Core's scheduler is safe to run across multiple server instances sharing the same PostgreSQL database. Each polling service uses a different concurrency strategy matched to its semantics: an advisory lock for leader election, row-level locking for parallel dispatch, a per-subject advisory lock for work that names the same subject, and idempotent operations where none of these is needed.

This page documents the concurrency model, the guarantees it provides, and the implications for multi-server deployments.

Overview

ServiceStrategyParallelismGuarantee
ManifestManagerPollingServiceAdvisory lock (single-leader)One server per cycleNo duplicate WorkQueue entries
JobDispatcherPollingServiceFOR UPDATE SKIP LOCKED (per-entry), plus an advisory lock per subject keyAll servers dispatch concurrentlyNo duplicate Metadata or double-dispatch; at most one run in flight per subject
LocalWorkerServiceFOR UPDATE SKIP LOCKED (per-job)All servers execute concurrentlyNo duplicate job execution
MetadataCleanupPollingServiceNone (idempotent)All servers run concurrentlyDeleting already-deleted rows is a no-op

ManifestManager: Advisory Lock

The Problem

The ManifestManager evaluates which manifests are "due" for execution and writes WorkQueue entries. If two servers run the ManifestManager simultaneously, they both load the same manifests, both evaluate them as due, and both insert duplicate WorkQueue entries, causing the same train to be dispatched twice.

The race window exists between LoadManifestsJunction (which reads HasQueuedWork = false) and CreateWorkQueueEntriesJunction (which inserts the entry). This is a classic time-of-check-time-of-use (TOCTOU) bug.

The Solution

The ManifestManagerPollingService acquires a PostgreSQL transaction-scoped advisory lock before running the train:

SELECT pg_try_advisory_xact_lock(hashtext('trax_manifest_manager'))

This is a non-blocking try-lock: if another server already holds the lock, the current server skips the cycle and waits for the next polling tick. No server ever blocks waiting for the lock.

Server A                              Server B
────────                              ────────
BEGIN TRANSACTION
pg_try_advisory_xact_lock → true ✓
  LoadManifests                       BEGIN TRANSACTION
  ReapFailedJobs                      pg_try_advisory_xact_lock → false ✗
  DetermineJobsToQueue                "Another server is running ManifestManager,
  CreateWorkQueueEntries               skipping cycle"
COMMIT (releases lock)                ROLLBACK
                                      (waits for next polling tick)

How Advisory Locks Work

PostgreSQL advisory locks are application-level locks managed by the database but not tied to any table or row. They come in two flavors:

  • Session-level (pg_advisory_lock): held until explicitly released or the connection closes. Risky with connection pooling, if the connection returns to the pool with the lock held, it stays held until the connection is eventually closed.
  • Transaction-scoped (pg_try_advisory_xact_lock): automatically released when the transaction commits or rolls back. This is what Trax.Core uses. No risk of leaked locks.

The lock key is hashtext('trax_manifest_manager'), which produces a stable 32-bit integer from the string. It uses the single-key form of the advisory lock functions. Another application's single-key advisory lock on the same database would conflict only if its key hashed to the same integer, which is unlikely with a descriptive string. The dispatcher's subject lock uses the two-key form, which Postgres keeps in a separate space, so it can never contend with the leader lock.

Each lock name is hashed on its own. Trax.Effect 1.57.2 and earlier put the name inside SQL quotes, so Postgres hashed the text of EF's parameter placeholder and every leader lock, whatever its name, was the same key. The ManifestManager is the only caller, so nothing contended, but the key changed in 1.57.3: an old instance and a new one take different keys, and both run the ManifestManager at once for as long as both are up.

Upgrading from Trax.Effect 1.57.2 or earlier: stop every scheduler host before you start one on the new version. A rolling deploy is not safe across this change. The one-queued-entry-per-manifest index still refuses a duplicate queue row, but two leaders each run a full cycle (reaping, dead letters, queue entries) against the same manifests. Later upgrades can roll as usual.

Transaction Scope

The advisory lock wraps the entire ManifestManager train in a single transaction. This has two implications:

  1. Atomicity: All SaveChanges() calls the junctions make on the train's own data context (ReapFailedJobsJunction, CreateWorkQueueEntriesJunction) are buffered within the transaction. If the train fails partway through, those roll back together. No partial state (e.g., dead letters created but WorkQueue entries missing). One write is outside it: ResolveStaleStagedEntriesJunction cancels or promotes stranded staged entries through IWorkQueuePromotion, which opens its own context and commits immediately, so it is not undone if the train fails later in the cycle. That is safe because the update only touches entries that are still unconfirmed and queued, so repeating it on the next cycle changes nothing already resolved.

  2. Visibility delay: WorkQueue entries created by CreateWorkQueueEntriesJunction are not visible to the JobDispatcher until the ManifestManager transaction commits. This is typically a few milliseconds of additional latency. The JobDispatcher picks them up on its next polling tick. No work is lost.

Non-Postgres Providers

The advisory lock is only acquired when the IDataContext is backed by Entity Framework Core (i.e., it inherits from DbContext). When using the InMemory provider for tests, the lock is skipped entirely and the train runs directly. This is safe because InMemory implies a single-server, single-process setup.

Defense-in-Depth: Unique Partial Index

As an additional safety net, a unique partial index on the work_queue table prevents duplicate Queued entries for the same manifest at the database level:

CREATE UNIQUE INDEX ix_work_queue_unique_queued_manifest
    ON trax.work_queue (manifest_id)
    WHERE status = 'queued' AND manifest_id IS NOT NULL;

The index also meets a manual trigger: TriggerAsync, TriggerGroupAsync and the dashboard's trigger buttons queue an entry for the manifest, with its manifest_id set, and can land between the cycle loading a manifest as due and writing its entry. CreateWorkQueueEntriesJunction saves each entry on its own; when one insert fails, it logs the error and detaches that entry, so the rest of the cycle, the reapers' and dead-letter writes included, still saves. EF Core takes a savepoint before each save inside the leader transaction and rolls back to it on failure, so the transaction itself stays usable. No crash, no corruption.

A trigger checks for a queued entry first and, when the manifest already has one, marks that entry as the triggered run (bringing it forward if it was due later) instead of inserting another; an insert that loses a race is recognised the same way. So a trigger never fails on this index. Entries with no manifest (manifest_id IS NULL, such as a queueTrain enqueue) are excluded from it, and any number may be queued.

JobDispatcher: Row-Level Locking

The Problem

The JobDispatcher loads Queued WorkQueue entries and dispatches them, creating Metadata records, updating entry status to Dispatched, and enqueuing to the job submitter. If two servers load the same entries simultaneously, both would create Metadata records for the same entry and dispatch the train twice.

The Solution

The DispatchJobsJunction uses PostgreSQL's FOR UPDATE SKIP LOCKED to atomically claim each WorkQueue entry before dispatching it. Each entry is processed within its own DI scope and database transaction:

SELECT * FROM trax.work_queue
WHERE id = :entry_id AND status = 'queued' AND confirmed_at IS NOT NULL
  -- and, for an entry with a subject key, no run in flight for that subject
FOR UPDATE SKIP LOCKED

If the entry has already been claimed by another server (either locked in another transaction or already updated to Dispatched), the query returns no rows and the dispatcher skips it.

Server A                              Server B
────────                              ────────
Load queued entries [1, 2, 3]         Load queued entries [1, 2, 3]
 
BEGIN TRANSACTION                     BEGIN TRANSACTION
SELECT ... WHERE id=1 FOR UPDATE      SELECT ... WHERE id=1 FOR UPDATE
  SKIP LOCKED → row returned ✓          SKIP LOCKED → skipped (locked) ✗
  Create Metadata                     SELECT ... WHERE id=2 FOR UPDATE
  Update status → Dispatched            SKIP LOCKED → row returned ✓
COMMIT                                  Create Metadata
Enqueue to job submitter                  Update status → Dispatched
                                      COMMIT
BEGIN TRANSACTION                     Enqueue to job submitter
SELECT ... WHERE id=2 FOR UPDATE
  SKIP LOCKED → skipped (already      BEGIN TRANSACTION
  dispatched, status ≠ 'queued') ✗    SELECT ... WHERE id=3 FOR UPDATE
SELECT ... WHERE id=3 FOR UPDATE        SKIP LOCKED → row returned ✓
  SKIP LOCKED → skipped (locked) ✗      ...
                                      COMMIT
                                      Enqueue to job submitter

Why Not a Leader Lock?

Unlike the ManifestManager, the JobDispatcher benefits from parallel dispatch across servers. Each server can claim and dispatch different entries simultaneously, increasing throughput. A single dispatch-wide advisory lock, like the ManifestManager's leader lock, would serialize all dispatch activity to a single server, wasteful when the work queue has many entries.

The FOR UPDATE SKIP LOCKED pattern allows fine-grained, per-entry parallelism: multiple servers work through the queue concurrently, each atomically claiming the next available entry. This is the same pattern used by the LocalWorkerService for job execution.

The dispatcher does take an advisory lock, but a much narrower one: per subject key, and only for entries that have one.

Subject Serialization: Advisory Lock per Subject

A train can name the subject its queued work touches, and entries for one subject must not run at the same time. The claim query refuses an entry whose subject already has a dispatched run that is Pending or InProgress. That check alone has a race that FOR UPDATE SKIP LOCKED does not close:

Server A                              Server B
────────                              ────────
BEGIN; claim entry 1 (subject S)      BEGIN; claim entry 2 (subject S)
  no dispatched sibling for S ✓         no dispatched sibling for S ✓
  lock row 1                            lock row 2 (a different row, no contention)
COMMIT                                COMMIT
→ two runs for S in flight

The two entries are different rows, so row locks never make them contend, and while both are still queued neither transaction can see the other's dispatched sibling. So the claim first takes a transaction-scoped advisory lock on the subject:

SELECT pg_advisory_xact_lock(hashtext('trax_subject'), hashtext(:subject_key))

This is the blocking two-key form. The second claimant waits until the first commits, then runs its claim query, sees the dispatched sibling, and gets no row. The fixed class key hashtext('trax_subject') keeps subject locks in the two-key space, apart from the single-key space the ManifestManager's leader lock and any consumer's own single-key locks use, whatever a subject hashes to. The lock is held only for the claim transaction, which commits before the job is submitted, so nothing remote happens while it is held. Entries without a subject key take no lock.

SQLite needs no lock: it has a single writer, so two claims cannot interleave. Its dialect's lock is a no-op.

Candidate loading also drops entries whose subject is busy and keeps only the first queued entry per subject in a cycle, so entries the claim will refuse do not consume MaxActiveJobs slots. See JobDispatcher.

Intra-Cycle Parallelism

In addition to multi-server parallelism, a single server can dispatch multiple entries concurrently within a single polling cycle via MaxConcurrentDispatch. This is particularly important when using UseRemoteWorkers(), where each dispatch blocks on an HTTP POST. Without intra-cycle parallelism, dispatching 50 entries at 2 seconds each takes ~100 seconds, blocking the polling service from starting the next cycle.

With MaxConcurrentDispatch(10), the same 50 entries are dispatched in ~10 seconds (5 batches of 10). The FOR UPDATE SKIP LOCKED pattern prevents conflicts: concurrent dispatches within the same cycle cannot claim the same entry, just as concurrent dispatches across servers cannot.

See JobDispatcher. Parallel Dispatch for configuration details.

Per-Entry DI Scope

Each entry is dispatched within its own DI scope, following the same pattern as the LocalWorkerService. This provides:

  1. Clean change tracker: each entry gets a fresh IDataContext with no stale tracked entities from previous iterations.
  2. Transaction isolation: if one entry fails, its transaction is rolled back without affecting others.
  3. Commit-then-enqueue: the claim transaction (Metadata creation + WorkQueue status update) is committed before calling EnqueueAsync on the job submitter. This makes the Metadata record visible to the job submitter when it begins execution, necessary because the InMemoryJobSubmitter executes trains synchronously within EnqueueAsync. If the enqueue fails after commit, the dispatcher fails that Metadata with a write that matches it only while it is still Pending. If a runner already started it (a remote runner that ran the job and answered with an error, or is still running it when the HTTP call times out), the write matches nothing, the entry stays Dispatched and the train is not run again. Otherwise, while the entry has dispatch attempts left (MaxDispatchAttempts, default 5), the failed run is recorded with FailureException = "DispatchRequeued" (not counted toward the manifest's retries) and the WorkQueue entry is reset to Queued, held back by a backoff (5 seconds, doubling, up to 5 minutes), so a later cycle dispatches it again with a new Metadata. Once the attempts are used up, the entry stays Dispatched and that last failure counts once. See Dispatch failures.

Capacity Limit Approximation

With multiple servers, MaxActiveJobs enforcement is approximate. Each server independently counts active Metadata records in LoadDispatchCapacityJunction. Between the count and the actual dispatch, other servers may have dispatched entries of their own, so the total can exceed the configured limit.

This is a deliberate tradeoff. MaxActiveJobs is a soft limit to prevent overwhelming the system. Not a strict concurrency semaphore. The alternative (a global advisory lock for the entire dispatch cycle) would serialize all dispatch activity, defeating the purpose of multi-server deployment.

Each server can dispatch up to the whole limit from the same count, so with N dispatching servers the number of active jobs can reach N times MaxActiveJobs (and N times a group's MaxActiveJobs). Size the limit for that, or run the dispatcher on one host if the limit must hold exactly.

LocalWorkerService: Already Safe

The LocalWorkerService has used FOR UPDATE SKIP LOCKED since its introduction. Multiple worker threads (across one or many servers) atomically claim jobs from the background_job table. Each claim is a separate transaction: lock the row, set fetched_at, commit. Other workers skip locked rows and move to the next available job.

See Job Submission. Worker Lifecycle for the full dequeue SQL and crash recovery details.

MetadataCleanupPollingService: Idempotent

The cleanup service deletes old Metadata records based on a retention period. Multiple servers can run cleanup concurrently without conflict, deleting an already-deleted row is a no-op. No locking is needed.

Deployment Considerations

Minimum Configuration

No configuration changes are needed for multi-server deployments. The concurrency controls are always active, advisory locks and row-level locking work correctly even with a single server (the leader lock is always acquired, the FOR UPDATE always succeeds).

Polling Interval Tuning

With multiple servers, consider the polling interval for the ManifestManager. Since only one server runs the ManifestManager per cycle (advisory lock), having many servers poll frequently means many lock acquisition attempts that return false. This is cheap (a single SQL call that returns immediately), but if you want to reduce noise in logs, you can increase the interval:

.AddScheduler(scheduler => scheduler
    .ManifestManagerPollingInterval(TimeSpan.FromSeconds(10))
)

The JobDispatcher polling interval doesn't need adjustment, multiple servers processing the work queue concurrently is the desired behavior.

Database Connection Pooling

Advisory locks are transaction-scoped, so they don't interact with connection pooling. When a transaction commits or rolls back, the lock is released regardless of what happens to the underlying connection. No special pooling configuration is needed.

Monitoring

In a multi-server deployment, you'll see these log messages:

# Server that acquires the lock:
ManifestManager polling cycle starting
ManifestManager polling cycle completed
 
# Servers that skip:
Another server is running ManifestManager, skipping cycle
 
# JobDispatcher (on any server):
Work queue entry {id} already claimed by another server, skipping

These are Debug-level messages. In production, set the log level to Information or higher to suppress them.

Summary of Guarantees

ScenarioGuaranteeMechanism
Two servers evaluate the same manifest as "due"Only one creates a WorkQueue entryAdvisory lock + unique partial index
Two servers try to dispatch the same WorkQueue entryOnly one creates the Metadata and enqueuesFOR UPDATE SKIP LOCKED
Two servers try to dispatch different entries for the same subjectAt most one run for the subject is in flightPer-subject advisory lock + in-flight check in the claim
Two workers try to execute the same BackgroundJobOnly one claims and runs itFOR UPDATE SKIP LOCKED
Two servers run metadata cleanup concurrentlyBoth succeed, no side effectsIdempotent deletes
A server crashes mid-ManifestManager cycleTransaction rolls back, lock released, no partial stateTransaction-scoped advisory lock
A server crashes mid-dispatch of a WorkQueue entryTransaction rolls back, entry remains Queued for next cyclePer-entry transaction
A worker crashes mid-execution of a BackgroundJobVisibility timeout expires, job reclaimed by another workerfetched_at timestamp, refreshed only while the job runs
One job is delivered twice (SQS redelivery, retried dispatch, re-claimed job)Only one delivery runs the train; the other completes without running it or recording anythingConditional Pending → InProgress claim on the run's row

SDK Reference

AddScheduler | UseRemoteWorkers