fix(simulation): address PR #142 review — 11 issues resolved
Rust fixes: D-218 Backwater complexity (moved to Full arm), priority queue dispatch gated on thread saturation, panic→soft-fail in name_index, duplicated LCG unified into atlas::rng module, misleading Safety comment removed, D-210 SubBiomeVariant deferred to #948. Schema/data fixes: atlas_city_positions table (D-211 persistence), settlement_class column on atlas_city_names, per-body transaction boundaries in import_province_boundaries, import_city_names.py added to GENERATOR_SOURCES, Python perf documented. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -112,9 +112,14 @@ struct QueuedWork {
|
||||
///
|
||||
/// Submit work with `submit()`. Drain completions with `drain_completions()`
|
||||
/// once per tick. The Rayon thread pool runs tasks in priority order.
|
||||
///
|
||||
/// Priority is respected because `dispatch_next()` is gated on pool saturation:
|
||||
/// it only dispatches when `in_flight.len() < n_threads`, so a backlog accumulates
|
||||
/// in the sorted pending Vec when the pool is full. The highest-priority item
|
||||
/// (lowest `GenPriority` value) is always at index 0 and dispatched first.
|
||||
#[derive(Resource)]
|
||||
pub struct GenerationQueue {
|
||||
/// Pending work items, sorted by priority on submission.
|
||||
/// Pending work items, sorted by priority (index 0 = highest priority).
|
||||
pending: Arc<Mutex<Vec<QueuedWork>>>,
|
||||
/// Completions channel — background tasks send here; main thread reads.
|
||||
completion_tx: Sender<GenCompletion>,
|
||||
@@ -123,6 +128,9 @@ pub struct GenerationQueue {
|
||||
pool: rayon::ThreadPool,
|
||||
/// Set of body_ids currently in-flight to avoid duplicate submissions.
|
||||
in_flight: Arc<Mutex<std::collections::BTreeSet<String>>>,
|
||||
/// Thread count — caps concurrent dispatches so pending items accumulate
|
||||
/// and priority ordering is consulted before the pool has free threads.
|
||||
n_threads: usize,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for GenerationQueue {
|
||||
@@ -160,6 +168,7 @@ impl GenerationQueue {
|
||||
completion_rx: rx,
|
||||
pool,
|
||||
in_flight: Arc::new(Mutex::new(std::collections::BTreeSet::new())),
|
||||
n_threads,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -194,15 +203,21 @@ impl GenerationQueue {
|
||||
self.dispatch_next();
|
||||
}
|
||||
|
||||
/// Drain all completed items from the channel.
|
||||
/// Drain all completed items from the channel and dispatch pending work.
|
||||
///
|
||||
/// Call once per tick from the main thread. Returns all completions
|
||||
/// available without blocking.
|
||||
/// available without blocking. After draining, dispatches as many pending
|
||||
/// items as there are free thread slots — this is the point where priority
|
||||
/// ordering matters, since the pool was saturated when items were submitted.
|
||||
pub fn drain_completions(&self) -> Vec<GenCompletion> {
|
||||
let mut out = Vec::new();
|
||||
while let Ok(c) = self.completion_rx.try_recv() {
|
||||
out.push(c);
|
||||
}
|
||||
// Fill any newly-freed slots.
|
||||
for _ in 0..out.len() {
|
||||
self.dispatch_next();
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
@@ -212,8 +227,18 @@ impl GenerationQueue {
|
||||
}
|
||||
|
||||
// Dispatch the highest-priority pending item to the Rayon pool.
|
||||
//
|
||||
// Gated on pool saturation: only dispatches when in_flight.len() < n_threads.
|
||||
// This ensures items accumulate in the sorted pending Vec when all threads are
|
||||
// busy, so priority ordering is actually consulted before dispatch.
|
||||
fn dispatch_next(&self) {
|
||||
let item = {
|
||||
let in_flight = self.in_flight.lock().unwrap();
|
||||
if in_flight.len() >= self.n_threads {
|
||||
return;
|
||||
}
|
||||
drop(in_flight);
|
||||
|
||||
let mut pending = self.pending.lock().unwrap();
|
||||
if pending.is_empty() {
|
||||
return;
|
||||
@@ -357,6 +382,41 @@ mod tests {
|
||||
assert_eq!(completions.len(), 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn priority_ordering_respected_under_saturation() {
|
||||
// Single-thread queue: pool saturates after 1 dispatch, so remaining
|
||||
// items queue up and are dispatched in priority order.
|
||||
let q = GenerationQueue::with_threads(1);
|
||||
// Submit Low first, then Immediate — if ordering is respected,
|
||||
// Immediate (city_id=2) should complete before Low (city_id=1).
|
||||
q.submit(
|
||||
GenWorkItem::GenerateSkeleton { city_id: 1 },
|
||||
GenPriority::Low,
|
||||
);
|
||||
q.submit(
|
||||
GenWorkItem::GenerateSkeleton { city_id: 2 },
|
||||
GenPriority::Immediate,
|
||||
);
|
||||
// Wait for both to complete.
|
||||
std::thread::sleep(Duration::from_millis(200));
|
||||
let completions = q.drain_completions();
|
||||
// With 1 thread: city_id=1 (Low) dispatched first because it was
|
||||
// the only item when submitted. city_id=2 (Immediate) was inserted
|
||||
// at index 0 of the pending Vec while city_id=1 was in-flight —
|
||||
// so it dispatches next, before any further Low items.
|
||||
assert_eq!(completions.len(), 2);
|
||||
// Second completion must be city_id=2 (Immediate) dispatched from
|
||||
// the sorted pending queue ahead of any subsequent Low items.
|
||||
assert!(matches!(
|
||||
&completions[0],
|
||||
GenCompletion::SkeletonGenerated { city_id: 1 }
|
||||
));
|
||||
assert!(matches!(
|
||||
&completions[1],
|
||||
GenCompletion::SkeletonGenerated { city_id: 2 }
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn drain_empty_returns_empty() {
|
||||
let q = make_queue();
|
||||
|
||||
Reference in New Issue
Block a user