Skip to content

Async Concurrency in Rust for Low Latency Trading Systems

12 min read

Overview

Low-latency trading systems operate on market state that can change while work is still in flight. A calculation can therefore be correct for the state that produced it and useless by the time it completes.

The business challenge is not simply doing more work faster---it is prioritising the right work, avoiding effort on opportunities that have already disappeared, and understanding where time is being spent when decisions are delayed.

A trading runtime has to balance four things:

  • latency — complete a useful decision quickly;
  • throughput — process enough useful work to keep pace with the market;
  • freshness — stop obsolete work consuming the next decision budget; and
  • correctness — preserve explicit state, ownership, and execution boundaries.

This article uses Salus, a Rust-based trading and solving infrastructure project I built for decentralized markets, to examine those trade-offs. The examples use Rust and decentralized finance (DeFi), but the same principles apply to latency-sensitive trading and distributed systems in other languages.

What work is still useful, who owns it, what happens when the system is saturated, and can we explain where the latency went?

Salus

Salus turns changing decentralized-market state into bounded evaluation work and guarded execution decisions. It separates durable topology from live state, keeps known routes in memory, maps changed liquidity components directly to affected routes, evaluates those routes against current state, and keeps execution as a separate authority boundary.

Tokio coordinates asynchronous input/output (I/O), streams, timers, and channels. CPU-heavy simulation runs on bounded workers. Atomics carry narrow freshness facts, locks protect compound state, and channels transfer ownership between stages.

Figure 1. Salus runtime architecture: external market state enters a bounded runtime hot path, while persistence, execution, configuration, and observability retain separate ownership.

market / blockchain state
        ↓
update current in-memory state
        ↓
changed component → affected route lookup
        ↓
bounded admission
        ↓
CPU-heavy evaluation
        ↓
freshness check and selection
        ↓
guarded execution handoff

Key Design Decisions

The decisions below are ordered roughly by their impact on useful work and the complexity they remove from the hot path.

1. Evaluate only routes affected by changed market state

  • Decision: Hydrate known routes into memory and maintain a reverse component → route index, so a pool update resolves directly to affected routes rather than scanning the full catalogue.
  • Impact: This changes the steady-state problem from "inspect every route" to "inspect routes referenced by what changed." A retained 2.4-million-route Ethereum workload measured affected-scope lookup at 12 ms p50 and 42 ms p99. Deep dive →

2. Treat freshness as part of correctness

  • Decision: Publish the latest state independently, coalesce replaceable work, cancel stale evaluation cooperatively, and reject stale results before selection.
  • Impact: CPU capacity is spent on the newest useful decision rather than finishing work that is correct only for an obsolete market state. Deep dive →

3. Make overload behavior explicit

  • Decision: Use bounded queues with stage-specific policies: block where work must be preserved, drop where work is replaceable, and implement latest-state replacement where freshness matters.
  • Impact: Producer/consumer mismatch becomes visible backpressure or an intentional business decision instead of hidden memory growth and queue age. Deep dive →

4. Separate asynchronous coordination from CPU-heavy simulation

  • Decision: Use Tokio for I/O and orchestration, but run numerical route evaluation as finite chunks on a bounded CPU worker model.
  • Impact: More async tasks do not create more CPU. Bounding compute prevents simulation from starving ingestion, timers, channels, and other asynchronous work. Deep dive →

5. Keep durable storage out of the decision path

  • Decision: PostgreSQL remains the durable recovery authority, while startup hydration builds the in-memory route catalogue used by the live runtime.
  • Impact: Normal block processing does not pay database latency to rediscover stable route topology. Deep dive →

6. Measure latency by stage

  • Decision: Record queue wait, preparation, compute, worker utilization, stale work, search effort, execution handoff, and external submission timing separately.
  • Impact: A slow decision can be attributed to admission pressure, CPU work, stale scheduling, or an external provider instead of being reduced to one misleading average. Deep dive →

7. Bound amount search as well as route count

  • Decision: Search for a useful trade size within explicit probe, ceiling, tolerance, and iteration limits rather than claiming an unconstrained global optimum.
  • Impact: Optimization quality cannot consume an unlimited share of the decision budget. Deep dive →

8. Keep evaluation separate from execution authority

  • Decision: Evaluation produces a typed candidate with state, amount, represented costs, and freshness context. Selection can accept or abstain; execution remains a separately guarded handoff.
  • Impact: A model result cannot silently become a transaction or a commercial outcome claim. Deep dive →

9. Choose synchronization from the invariant

  • Decision: Use atomics for independent facts, short locks for compound invariants, channels for ownership transfer, and shared ownership only where lifetime sharing is actually required.
  • Impact: Synchronization follows the data contract rather than a preference for "lock-free" or "async" code. Deep dive →

System Architecture

The runtime narrows authority as work moves forward: ingestion owns current state, discovery owns known route membership, evaluation owns calculation, policy owns selection, and execution owns external handoff.

Figure 2. Salus separates asynchronous state coordination, bounded admission, CPU-heavy evaluation, freshness, and guarded execution into explicit runtime stages.

state ingestion
      ↓
in-memory topology update
      ↓
changed components
      ↓
affected-route index
      ↓
coverage + bounded admission
      ↓
bounded evaluation queue
      ↓
CPU worker chunks
      ↓
freshness + selection
      ↓
guarded execution

A route can stop at any boundary: no affected membership, missing state coverage, admission limits, stale generation, invalid simulation, non-positive modeled result, or failed execution preflight. That is intentional. The pipeline is designed to abstain early rather than make later stages authoritative by accident.

End-to-End Runtime Pipeline

Figure 3. Changed components resolve through the in-memory route index before bounded preparation, coverage validation, and evaluation.

Tycho / blockchain stream
        ↓
canonical block and state ingestion
        ↓
in-memory graph and topology update
        ↓
changed-component → known-route resolution
        ↓
route preparation, coverage, admission, and request construction
        ↓
live route-evaluation queue
        ↓
evaluation coordinator
        ↓
deterministic route chunks → bounded scoped CPU workers
        ↓
deterministic result merge and finalization
        ↓
selection → execution-stage handoff → external submission / receipt work

This is a guide to the runtime ownership boundaries, not a claim that every block takes every path. Empty blocks, unavailable state, admission caps, stale work, and configuration can stop work before evaluation.

StageWorkConcurrency boundaryPrimary concern
IngestionAsync I/O + state applicationBounded stream stagesOrdering and input bursts
Graph / topologyIn-memory mutationOwner-local state → typed jobsAvoid unnecessary rebuilds
Route lookupIndexed memory lookupBounded candidate scopePrevent scope explosion
Route preparationMetadata + request constructionlive_route_prepFreshness and admission
EvaluationCPU-heavy simulationlive_route_evaluation → bounded workersCPU saturation and stale work
ExecutionExternal I/O + lifecycle stateBounded execution_stageDedupe, retry, external latency
MetricsSerialization + file I/OIndependent bounded writerNever stall the hot path

Reduce the Work Before Optimizing It

The largest performance win is often not faster code; it is doing less work.

PostgreSQL durable topology
        ↓ startup hydration
in-memory route catalogue
        ↓ reverse index
changed component → affected route IDs

Hydration removes normal database reads from block-critical route selection. Indexing removes repeated full-catalogue scans. Either one without the other leaves unnecessary work in the hot path.

The runtime's RouteIndex maintains deterministic reverse membership:

struct RouteIndex {
    routes_by_id: BTreeMap<String, RouteCandidate>,
    route_ids_by_component: BTreeMap<ComponentId, BTreeSet<String>>,
}

When a component changes, the runtime looks up its route IDs and unions them with the memberships of other changed components. Candidate caps and preparation budgets bound that union before expensive simulation begins.

R1: WETH → USDC → WETH           components: P1, P2
R2: WETH → USDC → cbBTC → WETH  components: P1, P3, P4
R3: WETH → DAI → WETH            components: P5, P6

P1 changes → {R1, R2}
P4 changes → {R2}

affected routes → {R1, R2}

R3 is never inspected for that block.

Naive:
changed components × total routes × route length

Indexed:
component lookups + union of affected route references

The exact tree operations still carry logarithmic factors, but the important architectural change is the scope of work: Salus selects from the affected portion of the topology instead of repeatedly scanning the full route universe.

On a retained Ethereum route review with 2,401,108 persisted routes, affected component-scope lookup measured 12 ms p50 and 42 ms p99. That is workload-specific evidence for the indexed runtime path, not a database benchmark or universal latency guarantee.

Topology and live state remain separate. PostgreSQL can recover what routes are known; ingestion must still provide current state for every route leg before evaluation. An optimization cache may accelerate lookup, but it does not become an independent source of truth.

Backpressure and Freshness

A full queue means producers are creating work faster than the next stage can consume it, and the system must decide what remains valuable.

PolicyMeaningTypical use
BlockProducerWait for capacityOrdered or durable work that must be preserved
DropNewestReject new work immediatelyReplaceable diagnostics or work safe to shed
LatestWinsApplication-owned replacement/coalescingState where newer work supersedes older pending work

The generic sender makes the policy visible:

pub async fn send(&self, item: T) -> Result<(), QueueSendError> {
    match self.config.overflow_policy {
        QueueOverflowPolicy::BlockProducer => self.send_with_backpressure(item).await,
        QueueOverflowPolicy::DropNewest => self.try_send_now(item),
        QueueOverflowPolicy::LatestWins => self.try_send_now(item),
    }
}

LatestWins is a label, not a generic drop-oldest implementation. Replacement semantics belong to the stage that understands whether old work is still valid.

Freshness extends the same idea to work already in flight.

Figure 3. New market state drives admission control, pending-work replacement, cooperative cancellation, and a final stale-result guard.

The newest block is published as an independent atomic fact:

self.latest_block_number
    .fetch_max(block_number, Ordering::AcqRel);

Workers can read that value without waiting for compound queue bookkeeping. If a newer block supersedes their source state, they set a cooperative cancellation flag. A final stale check runs again before selection.

A result can be correct for the block that produced it and still be useless after the market moves.

Discarding that work preserves capacity for the next decision.

Async I/O vs CPU-Bound Work

Tokio is valuable for streams, network/storage I/O, timers, channels, and lifecycle coordination. It does not create additional CPU capacity.

Tokio: ingestion / channels / timers / async dependencies
                         ↓
                 bounded admission
                         ↓
CPU workers: deterministic route-evaluation chunks
                         ↓
                 deterministic merge

The worker count is finite and never exceeds the available chunks. Workers claim chunk indexes atomically, check cancellation between chunks, and return owned results for reconciliation.

let worker_count = worker_count.min(chunk_count).max(1);
let next_chunk_index = AtomicUsize::new(0);
 
let chunk_index =
    next_chunk_index.fetch_add(1, Ordering::Relaxed);

The point is not that "thread pools should be bounded." The design decision is that CPU-heavy evaluation has its own finite capacity contract and must not consume the executor responsible for keeping market state moving.

The same principle guides synchronization:

PrimitiveSalus responsibilityRule
Tokio mpscbounded stage handoffmove ownership
AtomicU64latest block / generationone independent monotonic fact
AtomicBoolcooperative cancellationone independent flag
std::sync::Mutexqueue/in-flight metadatashort compound invariant; never across .await
Tokio Mutex / RwLockasync mutable services or read-heavy cachesuse only when state genuinely spans async work
Arcshared lifetime-managed stateownership, not synchronization
scoped OS threadsnumerical evaluationbounded CPU ownership
Use an atomic for one independent fact, a lock for a compound invariant, and a channel when ownership of work should move.

Simulation and Bounded Search

A route is not a static weighted graph edge. Price impact, fees, rounding, concentrated liquidity, and protocol state depend on both the input amount and the current state.

Salus evaluates an exact input through each ordered route leg. The output of one transition becomes the input to the next; missing or invalid state produces a typed rejection rather than a partial quote.

The best input amount is also state-dependent. Salus uses bounded search:

  1. start from a feasible probe;
  2. expand within a configured ceiling;
  3. refine the best observed region;
  4. stop at explicit tolerance and iteration limits.

The result is a bounded best-observed candidate, not a mathematical claim of a global optimum.

This matters for concurrency because search breadth changes work per route. One route can become more expensive if the amount search requires more probes or protocol calls. Salus therefore records search effort alongside route counts and service time.

Observability and Latency Attribution

Performance work becomes guesswork if every delay is collapsed into one end-to-end number.

decision latency
= queue wait
+ preparation
+ evaluation compute
+ finalization / selection
+ execution handoff
+ external submission / network / receipt wait

Figure 4. End-to-end decision latency is decomposed across queue wait, preparation, CPU evaluation, finalization, execution handoff, and external waiting.

EvidenceLikely issueFirst response
Queue wait rises; service time stableadmission or worker-capacity pressureinspect queue depth, caps, and producer rate
Service time rises; CPU utilization highevaluator bottleneckinspect simulation, search breadth, allocations
Evaluation-to-submit low; submission RPC highexternal provider/networkimprove the boundary, not the evaluator
Throughput high; stale completion highwrong scheduling policyreduce obsolete work

Telemetry itself is isolated behind a bounded writer. Diagnostic data may be dropped under extreme pressure rather than blocking route evaluation. Durable correctness or execution evidence has a different contract.

Instrumentation showed that the problem was not simply "the evaluator is slow." Route scope, request preparation, queue admission, stale work, CPU service, and external submission could each dominate different traces. The resulting optimizations were targeted: affected-route indexing, bounded admission, freshness-first cancellation, bounded CPU workers, block-scoped reuse, search gates, and external-latency attribution.

Retained Performance Evidence

Every number below is attached to its original retained workload. These are workload-specific engineering measurements, not production service-level objectives or universal Salus capacity claims.

Retained measurementWorkload / profileWhat it supportsWhat it does not support
2,401,108 persisted four-hop Ethereum routes; 3,415 graph tokens; 4,552 componentsretained Ethereum topology snapshotwhy full-universe scanning is the wrong steady-state operationcurrent live route count, heap usage, or universal graph size
affected component-scope lookup: 12 ms p50, 42 ms p99retained Ethereum route reviewindexed in-memory scope lookup stayed small beside numerical evaluationan end-to-end latency guarantee or a database benchmark
1,690,260 routes, 100/100 completed evaluation blocks, 33,578.21 active routes/sec, 503.38 ms average evaluation-service time, zero stale/discarded evaluations2026-07-08 Ethereum no-throttle traceworkload-specific active evaluator throughput and health under that profileuniversal capacity, block-cadence readiness, or all-chain behavior
4 ms evaluation-to-submit; 252 ms evaluation-to-broadcast; 246 ms submission RPCretained forced-route tracemost observed broadcast delay in that trace was external submission RPCgeneral live trading latency or execution quality
earlier one-batch pilot: combined p95 1,046.319 ms to 119.080 ms across different windows; later qualifying 42-batch collection: combined p95 22,270 us, 0.850974% of 2.617 s evaluator-service p95separate retained HL-CARB runsthe qualifying collection met its stated overhead gatea directly comparable end-to-end reduction between those runs, production selection authority, or a universal ranking cost

Sources: Ethereum route review, performance hardening, the earlier HL-CARB pilot, and HL-CARB capture closeout.

Implementation Examples

Three short excerpts capture most of the concurrency story.

Bounded ownership transfer

pub fn new(config: QueueConfig) -> Self {
    let (sender, receiver) = mpsc::channel(config.capacity);
    let metrics = Arc::new(Mutex::new(QueueMetrics::default()));
    let monitored_sender = MonitoredSender {
        inner: sender, config, metrics: Arc::clone(&metrics),
    };
    Self { sender: monitored_sender, receiver, metrics }
}

The queue has finite capacity from construction. If producers outrun consumers, the runtime must expose pressure or apply an explicit policy rather than hide the mismatch in an unbounded backlog.

Freshness without a bookkeeping lock

pub async fn begin_block(&self, block_number: u64) -> RouteEvaluationBlockRollover {
    self.latest_block_number
        .fetch_max(block_number, Ordering::AcqRel);
 
    match self.shared_state.try_lock() {
        Ok(mut state) => state.begin_block(block_number, self.keep_previous_blocks),
        Err(TryLockError::WouldBlock) => RouteEvaluationBlockRollover::default(),
        Err(TryLockError::Poisoned(_)) => panic!("route evaluation shared state mutex poisoned"),
    }
}

Workers need to know immediately when their work is stale. That independent freshness fact should not wait for the mutex protecting multi-field queue state.

Bounded CPU ownership

let worker_count = worker_count.min(chunk_count).max(1);
let next_chunk_index = AtomicUsize::new(0);
 
thread::scope(|scope| {
    for _ in 0..worker_count {
        scope.spawn(|| {
            loop {
                let index = next_chunk_index.fetch_add(1, Ordering::Relaxed);
                if index >= chunk_count { break; }
                if cancel_flag.load(Ordering::Acquire) { continue; }
                run_chunk(index);
            }
        });
    }
});

Route simulation is CPU-bound. Salus caps parallelism, lets workers claim finite chunks, checks cancellation between chunks, and reconciles the results after the workers join.

Conclusion

Low-latency concurrency is an ownership and prioritization problem before it is a task-count problem.

Salus reduces the route universe before simulation, bounds the work that can enter expensive stages, stops obsolete work from consuming the next decision budget, separates asynchronous coordination from finite CPU capacity, and measures each stage independently.

The individual mechanisms---indexes, bounded queues, atomics, mutexes, worker pools, and telemetry---are familiar. The important design work is deciding where each belongs, what business invariant it protects, and what happens when capacity is exhausted.

That leads back to the central question:

What work is still useful, who owns it, what happens when the system is saturated, and can we explain where the latency went?

Further reading

Salus builds on Tycho for foundational market-data ingestion, simulation, and execution infrastructure. Thanks to the broader DeFi community and the Tycho Build Community for the systems and ideas that made this work possible.

Sources and Evidence

This canonical article preserves the supplied V4 author source and reconciles it only with already approved public Salus material. This publication pass did not inspect the private implementation repository. Implementation names, capacities, and historical measurements retained from V4 remain workload-scoped source material; they do not establish current implementation, production, commercial, or execution outcomes.

Appendix A — Synchronization Types in the Salus Runtime

The synchronization model is intentionally semantic: choose the primitive based on the invariant being protected and whether ownership should be shared or transferred.

The practical rule is:

Use an atomic for one independent fact, a lock for a compound invariant, and a channel when ownership of work should move.

The following compact examples show the principal synchronization and lifecycle primitives used across Salus.

use std::{sync::{Arc, Mutex as StdMutex, atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering}}, thread};
use tokio::{sync::{mpsc, Mutex as TokioMutex, RwLock as TokioRwLock, Notify, oneshot}, task::JoinHandle};
 
// Tokio mpsc — bounded ownership transfer. Salus: prep/evaluation/execution.
let (tx, mut rx) = mpsc::channel::<EvaluationJob>(4);
let worker: JoinHandle<()> = tokio::spawn(async move {
    while let Some(job) = rx.recv().await { process(job).await; }
});
 
// AtomicU64 — one monotonic freshness fact. Salus: latest observed block.
let latest_block = Arc::new(AtomicU64::new(0));
latest_block.fetch_max(block_number, Ordering::AcqRel);
let current = latest_block.load(Ordering::Acquire);
 
// AtomicBool — cooperative cancellation. Salus: stale evaluation chunks.
let cancelled = Arc::new(AtomicBool::new(false));
if source_block < current { cancelled.store(true, Ordering::Release); }
if cancelled.load(Ordering::Acquire) { return; }
 
// AtomicUsize — independent progress/work claiming. Salus: chunk indexes.
let next_chunk = AtomicUsize::new(0);
let chunk_index = next_chunk.fetch_add(1, Ordering::Relaxed);
 
// std::sync::Mutex — short compound invariant. Salus: queue/in-flight state.
let state = Arc::new(StdMutex::new(QueueState::default()));
{ let mut s = state.lock().expect("queue state"); s.queued -= 1; s.in_flight += 1; }
// Guard is released before any .await.
 
// Tokio Mutex — protected mutation that genuinely spans async work.
// Salus: asynchronous execution/service boundaries.
let execution = Arc::new(TokioMutex::new(ExecutionState::default()));
{ let mut s = execution.lock().await; submit_using_nonce(s.nonce).await?; s.nonce += 1; }
 
// Tokio RwLock — read-heavy async cache. Salus: price/gas/balance caches.
let gas_price = Arc::new(TokioRwLock::new(None::<u128>));
if let Some(v) = *gas_price.read().await { use_cached_price(v); }
*gas_price.write().await = Some(new_price);
 
// Arc — shared ownership without deep-copying. Salus: handles/metrics/liveness.
let cancel = Arc::new(AtomicBool::new(false));
let worker_cancel = Arc::clone(&cancel);
 
// Scoped OS thread — bounded CPU ownership. Salus: numerical evaluation.
thread::scope(|scope| { scope.spawn(|| evaluate_cpu_chunk()); });
 
// spawn_blocking — explicit blocking/CPU boundary. Salus: route discovery.
let routes = tokio::task::spawn_blocking(|| discover_routes()).await??;
 
// JoinHandle — explicit lifecycle. Salus: queue workers/writer tasks.
let handle: JoinHandle<()> = tokio::spawn(async { run_queue_worker().await });
handle.await?;
 
// Notify — signal app-owned pending work. Salus: route-refresh wakeup.
let ready = Arc::new(Notify::new()); ready.notify_one(); ready.notified().await;
 
// oneshot — one request/one result. Salus: blocking-discovery result handoff.
let (result_tx, result_rx) = oneshot::channel();
result_tx.send(discovery_result)?; let result = result_rx.await?;

The important distinction is semantic rather than syntactic. AtomicU64 fits the latest block because it is one independent monotonic fact. Queued requests, in-flight work, counters, and cancellation metadata form a compound invariant and belong behind a clear owner or short mutex. Channels are preferable when ownership should move instead of mutation being shared.

Primitive map

PrimitiveCurrent Salus useWhy it fits
Tokio mpsctyped bounded stage work, including prep/evaluation/executiontransfers ownership, expresses finite capacity, and supports async coordination
AtomicU64latest block publication and monotonic generationsa single independent fact is readable in hot stale checks without queue-state locking
AtomicBoolcooperative cancellation and one-time lifecycle flagscheaply shares a cancellation signal across chunks/tasks
AtomicUsizechunk allocation and processed/profitable countersindependent monotonically changing counters
std::sync::Mutexqueue metrics, queued/in-flight metadata, writer statemaintains compound invariants with short non-await guards
Tokio Mutexasynchronous execution runtime and services crossing .awaitserializes mutable async state safely where a guard spans async work
Tokio RwLockread-heavy price/gas/flash-balance cache surfacessupports a read-heavy cache when its measured trade-off is justified
Arcshared metrics, queue state, liveness registry, immutable protocol handlesclear shared ownership/lifetime without deep-copying state
Scoped OS worker threadbounded CPU evaluation poolprevents CPU fan-out from becoming unlimited Tokio task fan-out
spawn_blockingblocking route discovery and other blocking boundariesmoves bounded blocking work off Tokio executor workers
JoinHandlequeue workers and writer lifecyclegives owners an explicit shutdown/join result rather than detached work

Source-draft implementation note (pending separately authorized implementation verification). The normal queue pipeline does not use a Semaphore or broadcast as its primary work-queue mechanism; route-refresh uses Notify with its app-owned pending slot, and oneshot returns individual blocking discovery results. Read the current queue architecture table before generalising this inventory.

Appendix B — Queue and Worker Code

The main article explains why Salus bounds work and treats freshness as part of correctness. These excerpts retain the implementation shape for readers who want the concrete Rust mechanics.

B.1 Bounded generic queue construction

pub fn new(config: QueueConfig) -> Self {
    let (sender, receiver) = mpsc::channel(config.capacity);
    let metrics = Arc::new(Mutex::new(QueueMetrics::default()));
    let monitored_sender = MonitoredSender {
        inner: sender, config, metrics: Arc::clone(&metrics),
    };
    Self { sender: monitored_sender, receiver, metrics }
}

Why: Queue capacity is explicit at construction. Producer/consumer mismatch therefore becomes measurable backpressure or a deliberate drop policy instead of an unbounded backlog.

B.2 Explicit overflow-policy dispatch

pub async fn send(&self, item: T) -> Result<(), QueueSendError> {
    match self.config.overflow_policy {
        QueueOverflowPolicy::BlockProducer => self.send_with_backpressure(item).await,
        QueueOverflowPolicy::DropNewest => self.try_send_now(item),
        QueueOverflowPolicy::LatestWins => self.try_send_now(item),
    }
}

Why: Different stages value queued work differently. Durable work may need backpressure; replaceable diagnostics or stale market work may be dropped. LatestWins still requires application-owned replacement or coalescing.

B.3 Latest block as an independent freshness fact

pub async fn begin_block(&self, block_number: u64) -> RouteEvaluationBlockRollover {
    self.latest_block_number
        .fetch_max(block_number, Ordering::AcqRel);
    match self.shared_state.try_lock() {
        Ok(mut shared_state) => shared_state.begin_block(block_number, self.keep_previous_blocks),
        Err(TryLockError::WouldBlock) => RouteEvaluationBlockRollover::default(),
        Err(TryLockError::Poisoned(_)) => panic!("route evaluation shared state mutex poisoned"),
    }
}

Why: Workers need to see that their work is stale immediately. That fact should not wait for the mutex protecting compound queue metadata.

B.4 Bounded CPU worker ownership

let worker_count = worker_count.min(chunk_count).max(1);
let next_chunk_index = AtomicUsize::new(0);
thread::scope(|scope| {
    let mut worker_handles = Vec::with_capacity(worker_count);
    for _ in 0..worker_count {
        let next_chunk_index = &next_chunk_index;
        let registry = Arc::clone(&registry);
        let cancel_flag = Arc::clone(&cancel_flag);
        worker_handles.push(scope.spawn(move || {
            let runtime = match Builder::new_current_thread().enable_all().build() {
                Ok(runtime) => runtime,
                Err(error) => return vec![Err(format!("worker runtime: {error}"))],
            };
            let mut worker_results = Vec::new();
            loop {
                let chunk_index = next_chunk_index.fetch_add(1, Ordering::Relaxed);
                if chunk_index >= chunk_count {
                    break;
                }
                if cancel_flag.load(Ordering::Acquire) {
                    registry.record_cancelled(chunk_index);
                    continue;
                }
                registry.record_worker_started(chunk_index);
                let result = block_on_chunk(&runtime, (run_chunk)(chunk_index), &registry, chunk_index);
                // Record completion/failure, then retain this chunk result.
                worker_results.push(result);
            }
            worker_results
        }));
    }
});

Why: Route simulation is CPU-bound. Salus bounds parallelism, lets workers claim finite chunks, checks cancellation between chunks, and reconciles results after the workers join rather than creating a task per route.

Appendix C — Telemetry Code

Telemetry has its own bounded ownership path so measurement does not become the bottleneck it is trying to diagnose.

C.1 Queue-pressure snapshot

fields.insert("capacity".to_owned(), json!(snapshot.capacity));
fields.insert("current_depth".to_owned(), json!(snapshot.current_depth));
fields.insert("peak_depth".to_owned(), json!(snapshot.peak_depth));
fields.insert("sent".to_owned(), json!(snapshot.sent_count));
fields.insert("dropped".to_owned(), json!(snapshot.dropped_count));
fields.insert("blocked".to_owned(), json!(snapshot.blocked_send_count));

Why: Capacity, depth, sends, drops, and blocked sends need to be observed together to distinguish a slow evaluator from producer/consumer pressure.

C.2 Nonblocking recorder handoff

match self.inner.writer.try_send(line) {
    Ok(()) => true,
    Err(QueueSendError::DroppedByPolicy { .. }) => {
        self.inner.writer.log_queue_drop_once();
        false
    }
    Err(QueueSendError::Closed { queue_name }) => {
        self.disable_once(format!("writer_queue_closed queue={queue_name}"));
        false
    }
}

Why: Diagnostic file I/O should not park route evaluation. A full metrics queue can drop a diagnostic event; a closed writer disables metrics explicitly.

C.3 Single owner for blocking file I/O

while let Some(line) = receiver.blocking_recv() {
    write_runtime_metrics_line(&mut file, &state, &line);
}
if let Err(error) = file.flush() {
    disable_runtime_metrics_state(&state, format!("flush_error error={error}"));
} else {
    state.lock().expect("runtime metrics mutex not poisoned")
        .final_flush_succeeded = Some(true);
}

Why: One writer owns blocking file operations and final flush state, isolating storage latency from evaluation while keeping telemetry loss observable.