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 handoffKey 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 executionA 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 workThis 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.
| Stage | Work | Concurrency boundary | Primary concern |
|---|---|---|---|
| Ingestion | Async I/O + state application | Bounded stream stages | Ordering and input bursts |
| Graph / topology | In-memory mutation | Owner-local state → typed jobs | Avoid unnecessary rebuilds |
| Route lookup | Indexed memory lookup | Bounded candidate scope | Prevent scope explosion |
| Route preparation | Metadata + request construction | live_route_prep | Freshness and admission |
| Evaluation | CPU-heavy simulation | live_route_evaluation → bounded workers | CPU saturation and stale work |
| Execution | External I/O + lifecycle state | Bounded execution_stage | Dedupe, retry, external latency |
| Metrics | Serialization + file I/O | Independent bounded writer | Never 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 IDsHydration 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 referencesThe 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.
| Policy | Meaning | Typical use |
|---|---|---|
BlockProducer | Wait for capacity | Ordered or durable work that must be preserved |
DropNewest | Reject new work immediately | Replaceable diagnostics or work safe to shed |
LatestWins | Application-owned replacement/coalescing | State 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 mergeThe 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:
| Primitive | Salus responsibility | Rule |
|---|---|---|
Tokio mpsc | bounded stage handoff | move ownership |
AtomicU64 | latest block / generation | one independent monotonic fact |
AtomicBool | cooperative cancellation | one independent flag |
std::sync::Mutex | queue/in-flight metadata | short compound invariant; never across .await |
Tokio Mutex / RwLock | async mutable services or read-heavy caches | use only when state genuinely spans async work |
Arc | shared lifetime-managed state | ownership, not synchronization |
| scoped OS threads | numerical evaluation | bounded 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:
- start from a feasible probe;
- expand within a configured ceiling;
- refine the best observed region;
- 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 waitFigure 4. End-to-end decision latency is decomposed across queue wait, preparation, CPU evaluation, finalization, execution handoff, and external waiting.
| Evidence | Likely issue | First response |
|---|---|---|
| Queue wait rises; service time stable | admission or worker-capacity pressure | inspect queue depth, caps, and producer rate |
| Service time rises; CPU utilization high | evaluator bottleneck | inspect simulation, search breadth, allocations |
| Evaluation-to-submit low; submission RPC high | external provider/network | improve the boundary, not the evaluator |
| Throughput high; stale completion high | wrong scheduling policy | reduce 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 measurement | Workload / profile | What it supports | What it does not support |
|---|---|---|---|
| 2,401,108 persisted four-hop Ethereum routes; 3,415 graph tokens; 4,552 components | retained Ethereum topology snapshot | why full-universe scanning is the wrong steady-state operation | current live route count, heap usage, or universal graph size |
| affected component-scope lookup: 12 ms p50, 42 ms p99 | retained Ethereum route review | indexed in-memory scope lookup stayed small beside numerical evaluation | an 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 evaluations | 2026-07-08 Ethereum no-throttle trace | workload-specific active evaluator throughput and health under that profile | universal capacity, block-cadence readiness, or all-chain behavior |
| 4 ms evaluation-to-submit; 252 ms evaluation-to-broadcast; 246 ms submission RPC | retained forced-route trace | most observed broadcast delay in that trace was external submission RPC | general 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 p95 | separate retained HL-CARB runs | the qualifying collection met its stated overhead gate | a 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
| Primitive | Current Salus use | Why it fits |
|---|---|---|
Tokio mpsc | typed bounded stage work, including prep/evaluation/execution | transfers ownership, expresses finite capacity, and supports async coordination |
AtomicU64 | latest block publication and monotonic generations | a single independent fact is readable in hot stale checks without queue-state locking |
AtomicBool | cooperative cancellation and one-time lifecycle flags | cheaply shares a cancellation signal across chunks/tasks |
AtomicUsize | chunk allocation and processed/profitable counters | independent monotonically changing counters |
std::sync::Mutex | queue metrics, queued/in-flight metadata, writer state | maintains compound invariants with short non-await guards |
Tokio Mutex | asynchronous execution runtime and services crossing .await | serializes mutable async state safely where a guard spans async work |
Tokio RwLock | read-heavy price/gas/flash-balance cache surfaces | supports a read-heavy cache when its measured trade-off is justified |
Arc | shared metrics, queue state, liveness registry, immutable protocol handles | clear shared ownership/lifetime without deep-copying state |
| Scoped OS worker thread | bounded CPU evaluation pool | prevents CPU fan-out from becoming unlimited Tokio task fan-out |
spawn_blocking | blocking route discovery and other blocking boundaries | moves bounded blocking work off Tokio executor workers |
JoinHandle | queue workers and writer lifecycle | gives 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(®istry);
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), ®istry, 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.