diff --git a/docs/ISSUES.md b/docs/ISSUES.md
index 0fb401bf..274a5457 100644
--- a/docs/ISSUES.md
+++ b/docs/ISSUES.md
@@ -24,6 +24,37 @@ What does NOT go here:
- Every session: scan OPEN issues at start; promote/close anything we touched during the session before ending.
- Promoting to a Phase: mark as `DONE (promoted to Phase X)` + commit SHA where the Phase entry landed.
+## #418 — Login world load takes ~27 s: publication advances at a flat 32 blocks/s
+
+**Status:** IN-PROGRESS 2026-08-17 — producer half landed (this commit's
+striped `LandblockStreamer` worker pool); the pacer measurably remains on the
+consumer side. **Symptom:** login holds the portal tunnel ~27 s while the
+25×25 window (625 landblocks) drips in at exactly 32 blocks/s
+(`ACDREAM_PROBE_REVEAL_TIMING=1`, probe `695a27b4`; baseline
+totalMs=26728), then render/composites/collision/gate/materialization all
+flip ready in the same millisecond. **Evidence chain:** an A/B with every
+`StreamingWorkBudgetOptions` env ceiling cranked 25–64x changed nothing →
+read as producer-limited (ONE `acdream.streaming.worker` thread,
+~31 ms/block serial). This commit parallelized the producer
+(min(cores−2, 8) striped workers, per-landblock ordering preserved) and
+**disproved that reading**: with 8 workers ALL 625 builds complete in
+~203 ms of wall clock (`ACDREAM_PROBE_TELEPORT=1` BUILD lines
+t=3475390→3475593, dat-lock waited ≤ 12 ms, held 0–13 ms), yet `loaded=`
+still advances at exactly +32/1000 ms and totalMs measured 27395 / 27503
+across two runs of the new binary. **Hypothesis:** the 32/s cadence lives in
+the update-thread admission/publication path (`StreamingController` meter →
+`LandblockPresentationPipeline`), is frame-quantized (an exactly-integer
+per-second rate held for 14+ consecutive seconds — N update ticks per
+landblock at a stable tick rate, e.g. 2 ticks × 64 Hz), and is NOT governed
+by the budget env ceilings (the original A/B and this change now agree on
+that). **Next step:** instrument per-frame meter operations/yields + the
+tunnel frame rate, find which operation stage eats the ensured-progress
+floor, then lift the actual limiter. The producer pool stays: it takes the
+builds off the critical path (23 s → 0.2 s) and is quality-neutral
+(per-landblock ordering, ClearLoads, priority, and disposal semantics
+preserved; a pool of 1 reproduces the old serial behavior,
+regression-tested in `LandblockStreamerPoolTests`).
+
## #417 — World ambience keeps playing (and re-firing) on the character-select screen after the in-world logoff
**Status:** ✅ FIXED 2026-08-17 (logout-audio round; fix + tests in the same
diff --git a/src/AcDream.App/Streaming/LandblockStreamer.cs b/src/AcDream.App/Streaming/LandblockStreamer.cs
index 234b90e1..600f7b07 100644
--- a/src/AcDream.App/Streaming/LandblockStreamer.cs
+++ b/src/AcDream.App/Streaming/LandblockStreamer.cs
@@ -16,19 +16,24 @@ namespace AcDream.App.Streaming;
/// per OnUpdate.
///
///
-/// Thread model (Phase A.5 T11+): spawns a
-/// dedicated background worker thread. and
-/// write non-blocking to the inbox
-/// ; the worker drains it and posts
-/// records to the outbox.
+/// Thread model (#418): spawns a small pool of
+/// dedicated background worker threads (default
+/// ; the single-worker degenerate case is
+/// the pre-#418 Phase A.5 T11+ shape). Jobs are striped across per-worker
+/// lanes by landblock id, so every job for one landblock id executes on
+/// one lane in enqueue order — a Load and Unload for the same id can never
+/// race on two workers. and
+/// write non-blocking to the owning lane's
+/// inbox ; each worker drains its lane and posts
+/// records to the shared outbox.
///
///
///
/// DatCollection thread safety is provided by the caller:
/// GameWindow's _datLock (Phase A.5 T10) serialises all
/// DatCollection.Get<T> calls. Both factory closures passed at
-/// construction acquire that lock before reading dats. The worker never
-/// touches DatCollection directly — it only calls the factories.
+/// construction acquire that lock before reading dats. The workers never
+/// touch DatCollection directly — they only call the factories.
///
///
///
@@ -40,7 +45,10 @@ namespace AcDream.App.Streaming;
///
/// Threading: must be called from a single
/// consumer thread (the render thread in production). All other public
-/// methods are thread-safe.
+/// methods are thread-safe. Completion arrival order is preserved per
+/// landblock id (per lane); across different landblocks it is arbitrary,
+/// which the single consumer already tolerates — the StreamingController's
+/// admission classifies each result independently into per-priority FIFOs.
///
///
public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
@@ -52,14 +60,23 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
///
public const int DefaultDrainBatchSize = 4;
+ ///
+ /// Default build worker pool size: leave two cores for the render and
+ /// update threads, cap at 8 (the login window's ~625 builds saturate
+ /// well before that), floor 1 (the pre-#418 single-worker shape).
+ ///
+ public static int DefaultWorkerCount =>
+ Math.Max(1, Math.Min(Environment.ProcessorCount - 2, 8));
+
private readonly Func _loadLandblock;
private readonly bool _supportsRequestOrigin;
private readonly Func _buildMeshOrNull;
- private readonly Channel _inbox;
+ private readonly Channel[] _lanes;
private readonly Channel _outbox;
private readonly CancellationTokenSource _cancel = new();
private readonly object _inboxGate = new();
- private Thread? _worker;
+ private Thread[]? _workers;
+ private int _activeWorkers;
private Exception? _workerFailure;
private int _completionBacklog;
private int _disposed;
@@ -74,17 +91,31 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
private LandblockStreamer(
Func loadLandblock,
Func? buildMeshOrNull,
- bool supportsRequestOrigin)
+ bool supportsRequestOrigin,
+ int? workerCount)
{
+ if (workerCount is < 1)
+ {
+ throw new ArgumentOutOfRangeException(
+ nameof(workerCount),
+ workerCount,
+ "The landblock build pool needs at least one worker.");
+ }
_loadLandblock = loadLandblock;
_supportsRequestOrigin = supportsRequestOrigin;
// Default: no mesh build (returns null → Failed result). Production
// wires in LandblockMesh.Build via the T12 construction site.
_buildMeshOrNull = buildMeshOrNull ?? ((_, _) => null);
- _inbox = Channel.CreateUnbounded(
- new UnboundedChannelOptions { SingleReader = true, SingleWriter = false });
+ int lanes = workerCount ?? DefaultWorkerCount;
+ _lanes = new Channel[lanes];
+ for (int i = 0; i < lanes; i++)
+ {
+ _lanes[i] = Channel.CreateUnbounded(
+ new UnboundedChannelOptions { SingleReader = true, SingleWriter = false });
+ }
+ // SingleWriter = false: every pool worker posts completions.
_outbox = Channel.CreateUnbounded(
- new UnboundedChannelOptions { SingleReader = true, SingleWriter = true });
+ new UnboundedChannelOptions { SingleReader = true, SingleWriter = false });
}
///
@@ -94,8 +125,9 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
///
public static LandblockStreamer CreateForRequests(
Func loadLandblock,
- Func? buildMeshOrNull = null) =>
- new(loadLandblock, buildMeshOrNull, supportsRequestOrigin: true);
+ Func? buildMeshOrNull = null,
+ int? workerCount = null) =>
+ new(loadLandblock, buildMeshOrNull, supportsRequestOrigin: true, workerCount);
///
/// Compatibility constructor for build factories that predate the
@@ -103,13 +135,15 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
///
public LandblockStreamer(
Func loadLandblock,
- Func? buildMeshOrNull = null)
+ Func? buildMeshOrNull = null,
+ int? workerCount = null)
: this(
request => loadLandblock(request.LandblockId, request.Kind) is { } build
? build
: null,
buildMeshOrNull,
- supportsRequestOrigin: false)
+ supportsRequestOrigin: false,
+ workerCount)
{
}
@@ -120,13 +154,15 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
///
public LandblockStreamer(
Func loadLandblock,
- Func? buildMeshOrNull = null)
+ Func? buildMeshOrNull = null,
+ int? workerCount = null)
: this(
request => loadLandblock(request.LandblockId, request.Kind) is { } landblock
? new LandblockBuild(landblock, Origin: request.Origin)
: null,
buildMeshOrNull,
- supportsRequestOrigin: false)
+ supportsRequestOrigin: false,
+ workerCount)
{
}
@@ -138,19 +174,40 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
///
public LandblockStreamer(
Func loadLandblock,
- Func? buildMeshOrNull = null)
+ Func? buildMeshOrNull = null,
+ int? workerCount = null)
: this(
request => loadLandblock(request.LandblockId) is { } landblock
? new LandblockBuild(landblock, Origin: request.Origin)
: null,
buildMeshOrNull,
- supportsRequestOrigin: false)
+ supportsRequestOrigin: false,
+ workerCount)
{
}
+ /// Configured pool size (also the lane count).
+ internal int WorkerCount => _lanes.Length;
+
///
- /// Activate the dedicated background worker thread. Idempotent and
- /// thread-safe: concurrent callers will only spawn one worker; subsequent
+ /// Lane (worker) that owns every job for .
+ /// Landblock ids carry their identity in the high 16 bits (0xXXYYFFFF —
+ /// the low word is constant across ids), so the id is mixed with a
+ /// Knuth multiplicative hash before reduction; a bare modulo would map
+ /// every id to one lane. Internal so tests can construct deterministic
+ /// per-lane contention.
+ ///
+ internal int LaneFor(uint landblockId)
+ {
+ if (_lanes.Length == 1)
+ return 0;
+ uint mixed = (landblockId >> 16) * 2654435761u;
+ return (int)(mixed % (uint)_lanes.Length);
+ }
+
+ ///
+ /// Activate the dedicated background worker threads. Idempotent and
+ /// thread-safe: concurrent callers will only spawn one pool; subsequent
/// calls are no-ops. Serialized with disposal so a worker can never start
/// after the owning DAT lifetime has been released.
///
@@ -160,16 +217,27 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
{
if (System.Threading.Volatile.Read(ref _disposed) != 0)
throw new ObjectDisposedException(nameof(LandblockStreamer));
- if (_worker is not null)
+ if (_workers is not null)
return;
- var worker = new Thread(WorkerLoop)
+ var workers = new Thread[_lanes.Length];
+ // Armed before any worker can run so the last worker to exit —
+ // however fast — is the one that completes the outbox.
+ System.Threading.Volatile.Write(ref _activeWorkers, workers.Length);
+ for (int i = 0; i < workers.Length; i++)
{
- IsBackground = true,
- Name = "acdream.streaming.worker",
- };
- worker.Start();
- _worker = worker;
+ int lane = i;
+ workers[i] = new Thread(() => WorkerLoop(lane))
+ {
+ IsBackground = true,
+ Name = workers.Length == 1
+ ? "acdream.streaming.worker"
+ : $"acdream.streaming.worker.{lane}",
+ };
+ }
+ foreach (Thread worker in workers)
+ worker.Start();
+ _workers = workers;
}
}
@@ -236,12 +304,15 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
///
/// Cancel every queued-but-not-started Load. Posts a
- /// control job which the worker
- /// honours at read time, dropping all pending Loads from both priority
- /// queues (Unloads survive). Used on the dungeon-entry edge to abort the
- /// in-flight 25×25 neighbor window so the ~129 ocean-grid dungeons never
- /// finish loading (#133 FPS). Loads the worker has ALREADY dequeued still
- /// complete; the StreamingController's collapsed-sweep unloads those few.
+ /// control job to EVERY lane
+ /// (one atomic broadcast under the enqueue gate, so a load enqueued
+ /// before this call is ordered ahead of its lane's ClearLoads copy and
+ /// dropped) which each worker honours at read time, dropping all pending
+ /// Loads from both of its priority queues (Unloads survive). Used on the
+ /// dungeon-entry edge to abort the in-flight 25×25 neighbor window so the
+ /// ~129 ocean-grid dungeons never finish loading (#133 FPS). Loads a
+ /// worker has ALREADY dequeued still complete (up to one per worker);
+ /// the StreamingController's collapsed-sweep unloads those few.
///
public void ClearPendingLoads()
{
@@ -250,18 +321,28 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
private void WriteJob(LandblockStreamJob job)
{
- // Serialize the writer's terminal transition with enqueue. Without
- // this small lifecycle gate a worker crash could complete the channel
+ // Serialize the writers' terminal transition with enqueue. Without
+ // this small lifecycle gate a worker crash could complete a lane
// between the caller's state check and TryWrite, silently dropping a
- // landblock request. Enqueues are destination-boundary events, not a
- // per-frame hot path.
+ // landblock request. The same gate makes the ClearLoads broadcast
+ // atomic with respect to concurrent enqueues. Enqueues are
+ // destination-boundary events, not a per-frame hot path.
lock (_inboxGate)
{
if (System.Threading.Volatile.Read(ref _disposed) != 0)
throw new ObjectDisposedException(nameof(LandblockStreamer));
if (_workerFailure is { } failure)
- throw new InvalidOperationException("The landblock streaming worker has terminated.", failure);
- if (!_inbox.Writer.TryWrite(job))
+ throw new InvalidOperationException("A landblock streaming worker has terminated.", failure);
+ if (job is LandblockStreamJob.ClearLoads)
+ {
+ foreach (Channel lane in _lanes)
+ {
+ if (!lane.Writer.TryWrite(job))
+ throw new InvalidOperationException("The landblock streaming inbox is no longer accepting work.");
+ }
+ return;
+ }
+ if (!_lanes[LaneFor(job.LandblockId)].Writer.TryWrite(job))
throw new InvalidOperationException("The landblock streaming inbox is no longer accepting work.");
}
}
@@ -308,8 +389,9 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
System.Threading.Interlocked.Increment(ref _completionBacklog);
}
- private void WorkerLoop()
+ private void WorkerLoop(int laneIndex)
{
+ ChannelReader inbox = _lanes[laneIndex].Reader;
var highPriority = new Queue();
var lowPriority = new Queue();
@@ -325,12 +407,12 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
{
if (highPriority.Count == 0 &&
lowPriority.Count == 0 &&
- !_inbox.Reader.WaitToReadAsync(_cancel.Token).AsTask().GetAwaiter().GetResult())
+ !inbox.WaitToReadAsync(_cancel.Token).AsTask().GetAwaiter().GetResult())
{
break;
}
- while (_inbox.Reader.TryRead(out var job))
+ while (inbox.TryRead(out var job))
{
if (job is LandblockStreamJob.ClearLoads)
{
@@ -358,17 +440,39 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
catch (Exception ex)
{
// Last-ditch: surface via outbox so the caller at least sees
- // something. We never retry a crashed worker.
+ // something. We never retry a crashed worker; any worker crash
+ // terminates the pool (matching the single-worker contract that
+ // a crash ends all job processing). A sibling that merely
+ // observed the crashed worker's lane completion (the wrapped
+ // ChannelClosedException) exits quietly instead of reporting a
+ // second, spurious crash.
+ bool cascade;
lock (_inboxGate)
{
- _workerFailure = ex;
- _inbox.Writer.TryComplete(ex);
+ cascade = _workerFailure is not null && ex is ChannelClosedException;
+ _workerFailure ??= ex;
+ foreach (Channel lane in _lanes)
+ lane.Writer.TryComplete(ex);
+ }
+ if (!cascade)
+ {
+ PublishResult(new LandblockStreamResult.WorkerCrashed(
+ _lanes.Length == 1
+ ? ex.ToString()
+ : $"worker {laneIndex}: {ex}"));
+ // Stop the sibling workers. Safe against Dispose: its
+ // CTS disposal only happens after every worker (this one
+ // included) has been joined.
+ _cancel.Cancel();
}
- PublishResult(new LandblockStreamResult.WorkerCrashed(ex.ToString()));
}
finally
{
- _outbox.Writer.TryComplete();
+ // The outbox has N writers; only the last worker out may
+ // complete it, or a crashed worker would silently drop the
+ // still-running workers' completions.
+ if (Interlocked.Decrement(ref _activeWorkers) == 0)
+ _outbox.Writer.TryComplete();
}
}
@@ -386,8 +490,9 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
// older queued LoadFar for the same landblock: LoadNear obviously
// loads everything, and PromoteToNear now carries mesh data so the
// render thread can run the full near-tier apply side effects. If a
- // LoadFar is already being processed, the single worker naturally
- // finishes it before the promotion is dequeued.
+ // LoadFar is already being processed, the owning lane's worker
+ // naturally finishes it before the promotion is dequeued (every
+ // job for one landblock id lives on one lane).
RemoveLowPriorityJobsForLandblock(
lowPriority,
high.LandblockId,
@@ -540,11 +645,18 @@ public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource
System.Threading.Interlocked.Exchange(ref _disposed, 1);
_cancel.Cancel();
lock (_inboxGate)
- _inbox.Writer.TryComplete();
+ {
+ foreach (Channel lane in _lanes)
+ lane.Writer.TryComplete();
+ }
// The owner releases the memory-mapped DAT immediately after this
- // object. Join the actual worker without a grace-period timeout so
- // no native read can survive into that teardown.
- _worker?.Join();
+ // object. Join every actual worker without a grace-period timeout
+ // so no native read can survive into that teardown.
+ if (_workers is { } workers)
+ {
+ foreach (Thread worker in workers)
+ worker.Join();
+ }
_cancel.Dispose();
_disposeCompleted = true;
}
diff --git a/tests/AcDream.App.Tests/Streaming/LandblockBuildOriginTests.cs b/tests/AcDream.App.Tests/Streaming/LandblockBuildOriginTests.cs
index 678a5452..4684e93a 100644
--- a/tests/AcDream.App.Tests/Streaming/LandblockBuildOriginTests.cs
+++ b/tests/AcDream.App.Tests/Streaming/LandblockBuildOriginTests.cs
@@ -118,13 +118,17 @@ public sealed class LandblockBuildOriginTests
var observed = new List();
var firstOrigin = new LandblockBuildOrigin(0xA9, 0xB4);
var secondOrigin = new LandblockBuildOrigin(0x71, 0xEC);
+ // workerCount: 1 — the ordered `observed`/completion asserts span two
+ // DIFFERENT landblock ids; only the serial degenerate pool (#418)
+ // guarantees that global order (per-id order is the pool contract).
using var streamer = LandblockStreamer.CreateForRequests(
loadLandblock: request =>
{
observed.Add(request);
return EmptyBuild(request.LandblockId, request.Origin);
},
- buildMeshOrNull: (_, _) => EmptyMesh());
+ buildMeshOrNull: (_, _) => EmptyMesh(),
+ workerCount: 1);
streamer.EnqueueLoad(new LandblockBuildRequest(
0xA9B4FFFFu,
diff --git a/tests/AcDream.App.Tests/Streaming/LandblockStreamerPoolTests.cs b/tests/AcDream.App.Tests/Streaming/LandblockStreamerPoolTests.cs
new file mode 100644
index 00000000..75caac17
--- /dev/null
+++ b/tests/AcDream.App.Tests/Streaming/LandblockStreamerPoolTests.cs
@@ -0,0 +1,431 @@
+using System.Collections.Concurrent;
+using AcDream.App.Streaming;
+using AcDream.Core.World;
+using DatReaderWriter.DBObjs;
+
+namespace AcDream.App.Tests.Streaming;
+
+///
+/// #418 build-worker-pool contract tests. The pool stripes jobs across
+/// per-worker lanes by landblock id, so the contract is: per-landblock
+/// enqueue order is execution (and completion-arrival) order, ClearLoads
+/// supersedes every load queued before it on every lane, Near-tier jobs run
+/// before Far-tier ones within a lane, a pool of one reproduces the serial
+/// pre-#418 behavior, and disposal joins every worker.
+///
+public sealed class LandblockStreamerPoolTests
+{
+ private const int SpinTimeoutMs = 10_000;
+ private const int SpinStepMs = 10;
+
+ private static AcDream.Core.Terrain.LandblockMeshData StubMesh() =>
+ new(
+ Array.Empty(),
+ Array.Empty());
+
+ private static LoadedLandblock StubLandblock(uint id) =>
+ new(id, new LandBlock(), Array.Empty());
+
+ ///
+ /// Produce a landblock id (retail 0xXXYYFFFF shape) owned by
+ /// , distinct from every id in
+ /// .
+ ///
+ private static uint IdInLane(LandblockStreamer streamer, int lane, HashSet taken)
+ {
+ for (uint x = 0; x < 256; x++)
+ {
+ for (uint y = 0; y < 256; y++)
+ {
+ uint id = (x << 24) | (y << 16) | 0xFFFFu;
+ if (streamer.LaneFor(id) == lane && taken.Add(id))
+ return id;
+ }
+ }
+ throw new InvalidOperationException($"No landblock id maps to lane {lane}.");
+ }
+
+ private static async Task> DrainCountAsync(
+ LandblockStreamer streamer,
+ int count,
+ int timeoutMs = SpinTimeoutMs)
+ {
+ var results = new List(count);
+ for (int i = 0; i < timeoutMs / SpinStepMs && results.Count < count; i++)
+ {
+ results.AddRange(streamer.DrainCompletions(count - results.Count));
+ if (results.Count < count)
+ await Task.Delay(SpinStepMs);
+ }
+
+ Assert.Equal(count, results.Count);
+ return results;
+ }
+
+ [Fact]
+ public void WorkerCount_DefaultsBoundedFloorsAtOneAndRejectsZero()
+ {
+ Assert.InRange(LandblockStreamer.DefaultWorkerCount, 1, 8);
+
+ using var defaulted = new LandblockStreamer(loadLandblock: _ => null);
+ Assert.Equal(LandblockStreamer.DefaultWorkerCount, defaulted.WorkerCount);
+
+ using var explicitThree = new LandblockStreamer(
+ loadLandblock: _ => null,
+ buildMeshOrNull: null,
+ workerCount: 3);
+ Assert.Equal(3, explicitThree.WorkerCount);
+
+ Assert.Throws(() => new LandblockStreamer(
+ loadLandblock: _ => null,
+ buildMeshOrNull: null,
+ workerCount: 0));
+ }
+
+ [Fact]
+ public void LaneAssignment_SpreadsRealWindowIdsAcrossLanes()
+ {
+ // Landblock ids carry their identity in the high word (0xXXYYFFFF);
+ // a bare modulo of the raw id would put EVERY id in one lane. Guard
+ // the mixed hash: a realistic 25x25 login window must engage more
+ // than one lane of an 8-lane pool.
+ using var streamer = new LandblockStreamer(
+ loadLandblock: _ => null,
+ buildMeshOrNull: null,
+ workerCount: 8);
+
+ var lanes = new HashSet();
+ for (uint x = 0xA0; x < 0xA0 + 25; x++)
+ for (uint y = 0xB0; y < 0xB0 + 25; y++)
+ lanes.Add(streamer.LaneFor((x << 24) | (y << 16) | 0xFFFFu));
+
+ Assert.True(
+ lanes.Count > 1,
+ $"25x25 window ids collapsed into {lanes.Count} lane(s).");
+ }
+
+ [Fact]
+ public async Task PerLandblockJobs_ExecuteInEnqueueOrder_UnderPoolContention()
+ {
+ const int idCount = 12;
+ const int jobsPerId = 8; // alternating LoadFar / Unload
+ int loaderCalls = 0;
+
+ using var streamer = new LandblockStreamer(
+ loadLandblock: (uint id, LandblockStreamJobKind _) =>
+ {
+ // Shake worker scheduling so a per-landblock ordering bug
+ // would actually interleave.
+ if (Interlocked.Increment(ref loaderCalls) % 3 == 0)
+ Thread.Sleep(1);
+ return StubLandblock(id);
+ },
+ buildMeshOrNull: (_, _) => StubMesh(),
+ workerCount: 4);
+ streamer.Start();
+
+ var ids = new uint[idCount];
+ for (uint i = 0; i < idCount; i++)
+ ids[i] = ((0x30u + i) << 24) | ((0x40u + i) << 16) | 0xFFFFu;
+
+ // Round-robin across ids while the pool is already running, so jobs
+ // for the same id repeatedly queue behind and race with other lanes.
+ ulong generation = 0;
+ for (int job = 0; job < jobsPerId; job++)
+ {
+ foreach (uint id in ids)
+ {
+ generation++;
+ if (job % 2 == 0)
+ streamer.EnqueueLoad(id, LandblockStreamJobKind.LoadFar, generation);
+ else
+ streamer.EnqueueUnload(id, generation);
+ }
+ }
+
+ List results =
+ await DrainCountAsync(streamer, idCount * jobsPerId);
+
+ foreach (uint id in ids)
+ {
+ var perId = results.Where(result => result.LandblockId == id).ToList();
+ Assert.Equal(jobsPerId, perId.Count);
+ // Completion arrival order per landblock id must be the enqueue
+ // order: alternating Loaded/Unloaded with strictly increasing
+ // generations.
+ for (int i = 0; i < perId.Count; i++)
+ {
+ if (i % 2 == 0)
+ Assert.IsType(perId[i]);
+ else
+ Assert.IsType(perId[i]);
+ if (i > 0)
+ {
+ Assert.True(
+ perId[i].Generation > perId[i - 1].Generation,
+ $"LB 0x{id:X8}: result {i} (gen {perId[i].Generation}) arrived " +
+ $"after gen {perId[i - 1].Generation} out of enqueue order.");
+ }
+ }
+ }
+ }
+
+ [Fact]
+ public async Task ClearPendingLoads_DropsQueuedLoadsAcrossAllLanes()
+ {
+ const int workerCount = 4;
+ using var release = new ManualResetEventSlim();
+ using var entered = new CountdownEvent(workerCount);
+ var loadedIds = new ConcurrentBag();
+ var blockerIds = new HashSet();
+
+ using var streamer = new LandblockStreamer(
+ loadLandblock: (uint id, LandblockStreamJobKind _) =>
+ {
+ loadedIds.Add(id);
+ bool isBlocker;
+ lock (blockerIds)
+ isBlocker = blockerIds.Contains(id);
+ if (isBlocker)
+ {
+ entered.Signal();
+ release.Wait();
+ }
+ return StubLandblock(id);
+ },
+ buildMeshOrNull: (_, _) => StubMesh(),
+ workerCount: workerCount);
+
+ var taken = new HashSet();
+ var victimIds = new List();
+ var unloadIds = new List();
+ var survivorIds = new List();
+ lock (blockerIds)
+ {
+ for (int lane = 0; lane < workerCount; lane++)
+ {
+ blockerIds.Add(IdInLane(streamer, lane, taken));
+ victimIds.Add(IdInLane(streamer, lane, taken));
+ victimIds.Add(IdInLane(streamer, lane, taken));
+ unloadIds.Add(IdInLane(streamer, lane, taken));
+ survivorIds.Add(IdInLane(streamer, lane, taken));
+ }
+ }
+
+ streamer.Start();
+ lock (blockerIds)
+ {
+ foreach (uint id in blockerIds)
+ streamer.EnqueueLoad(id, LandblockStreamJobKind.LoadFar);
+ }
+ // Every worker is now inside its blocker build; everything below
+ // queues behind them, one lane each.
+ Assert.True(entered.Wait(TimeSpan.FromSeconds(5)));
+
+ foreach (uint id in victimIds)
+ streamer.EnqueueLoad(id, LandblockStreamJobKind.LoadFar);
+ foreach (uint id in unloadIds)
+ streamer.EnqueueUnload(id);
+
+ // Must supersede every load queued above on EVERY lane, while the
+ // queued unloads survive.
+ streamer.ClearPendingLoads();
+
+ foreach (uint id in survivorIds)
+ streamer.EnqueueLoad(id, LandblockStreamJobKind.LoadFar);
+
+ release.Set();
+
+ // Expected completions: 4 blocker Loaded + 4 Unloaded + 4 survivor
+ // Loaded. The victims produce nothing.
+ List results =
+ await DrainCountAsync(streamer, workerCount * 3);
+
+ var loadedResultIds = results
+ .OfType()
+ .Select(result => result.LandblockId)
+ .ToHashSet();
+ var unloadedResultIds = results
+ .OfType()
+ .Select(result => result.LandblockId)
+ .ToHashSet();
+
+ lock (blockerIds)
+ {
+ foreach (uint id in blockerIds)
+ Assert.Contains(id, loadedResultIds);
+ }
+ foreach (uint id in survivorIds)
+ Assert.Contains(id, loadedResultIds);
+ foreach (uint id in unloadIds)
+ Assert.Contains(id, unloadedResultIds);
+ foreach (uint id in victimIds)
+ {
+ Assert.DoesNotContain(id, loadedResultIds);
+ Assert.DoesNotContain(id, loadedIds);
+ }
+ }
+
+ [Fact]
+ public async Task NearTierJobs_RunBeforeQueuedFarJobs_WithinEachLane()
+ {
+ const int workerCount = 3;
+ using var release = new ManualResetEventSlim();
+ using var entered = new CountdownEvent(workerCount);
+ var executionOrder = new ConcurrentQueue();
+ var blockerIds = new HashSet();
+
+ using var streamer = new LandblockStreamer(
+ loadLandblock: (uint id, LandblockStreamJobKind _) =>
+ {
+ bool isBlocker;
+ lock (blockerIds)
+ isBlocker = blockerIds.Contains(id);
+ if (isBlocker)
+ {
+ entered.Signal();
+ release.Wait();
+ }
+ else
+ {
+ executionOrder.Enqueue(id);
+ }
+ return StubLandblock(id);
+ },
+ buildMeshOrNull: (_, _) => StubMesh(),
+ workerCount: workerCount);
+
+ var taken = new HashSet();
+ var farIds = new uint[workerCount];
+ var nearIds = new uint[workerCount];
+ lock (blockerIds)
+ {
+ for (int lane = 0; lane < workerCount; lane++)
+ {
+ blockerIds.Add(IdInLane(streamer, lane, taken));
+ farIds[lane] = IdInLane(streamer, lane, taken);
+ nearIds[lane] = IdInLane(streamer, lane, taken);
+ }
+ }
+
+ streamer.Start();
+ lock (blockerIds)
+ {
+ foreach (uint id in blockerIds)
+ streamer.EnqueueLoad(id, LandblockStreamJobKind.LoadFar);
+ }
+ Assert.True(entered.Wait(TimeSpan.FromSeconds(5)));
+
+ // Far first, near second — the near job must still run first once
+ // the lane's blocker completes.
+ for (int lane = 0; lane < workerCount; lane++)
+ {
+ streamer.EnqueueLoad(farIds[lane], LandblockStreamJobKind.LoadFar);
+ streamer.EnqueueLoad(nearIds[lane], LandblockStreamJobKind.LoadNear);
+ }
+
+ release.Set();
+
+ await DrainCountAsync(streamer, workerCount * 3);
+
+ var observed = executionOrder.ToList();
+ for (int lane = 0; lane < workerCount; lane++)
+ {
+ int nearIndex = observed.IndexOf(nearIds[lane]);
+ int farIndex = observed.IndexOf(farIds[lane]);
+ Assert.True(nearIndex >= 0 && farIndex >= 0);
+ Assert.True(
+ nearIndex < farIndex,
+ $"lane {lane}: near 0x{nearIds[lane]:X8} (index {nearIndex}) ran " +
+ $"after far 0x{farIds[lane]:X8} (index {farIndex}).");
+ }
+ }
+
+ [Fact]
+ public async Task SingleWorkerPool_ReproducesSerialGlobalOrdering()
+ {
+ // The degenerate pool of one must be today's (pre-#418) behavior
+ // exactly: one worker thread, global near-first execution order
+ // across DIFFERENT landblock ids.
+ var callOrder = new List();
+ var loaderThreads = new HashSet();
+
+ using var streamer = new LandblockStreamer(
+ loadLandblock: (uint id, LandblockStreamJobKind _) =>
+ {
+ callOrder.Add(id);
+ loaderThreads.Add(System.Environment.CurrentManagedThreadId);
+ return StubLandblock(id);
+ },
+ buildMeshOrNull: (_, _) => StubMesh(),
+ workerCount: 1);
+ Assert.Equal(1, streamer.WorkerCount);
+
+ streamer.EnqueueLoad(0xAAAAFFFFu, LandblockStreamJobKind.LoadFar);
+ streamer.EnqueueLoad(0xBBBBFFFFu, LandblockStreamJobKind.LoadFar);
+ streamer.EnqueueLoad(0xCCCCFFFFu, LandblockStreamJobKind.LoadNear);
+ streamer.Start();
+
+ List results = await DrainCountAsync(streamer, 3);
+
+ Assert.Equal(
+ new[] { 0xCCCCFFFFu, 0xAAAAFFFFu, 0xBBBBFFFFu },
+ callOrder);
+ Assert.Single(loaderThreads);
+ var first = Assert.IsType(results[0]);
+ Assert.Equal(0xCCCCFFFFu, first.LandblockId);
+ }
+
+ [Fact]
+ public async Task Dispose_JoinsEveryWorkerInThePool()
+ {
+ const int workerCount = 3;
+ using var release = new ManualResetEventSlim();
+ using var entered = new CountdownEvent(workerCount);
+ var loaderThreads = new ConcurrentDictionary();
+ var blockerIds = new HashSet();
+
+ var streamer = new LandblockStreamer(
+ loadLandblock: (uint id, LandblockStreamJobKind _) =>
+ {
+ loaderThreads.TryAdd(System.Environment.CurrentManagedThreadId, 0);
+ entered.Signal();
+ release.Wait();
+ return StubLandblock(id);
+ },
+ buildMeshOrNull: (_, _) => StubMesh(),
+ workerCount: workerCount);
+
+ try
+ {
+ var taken = new HashSet();
+ for (int lane = 0; lane < workerCount; lane++)
+ blockerIds.Add(IdInLane(streamer, lane, taken));
+
+ streamer.Start();
+ foreach (uint id in blockerIds)
+ streamer.EnqueueLoad(id, LandblockStreamJobKind.LoadFar);
+ Assert.True(entered.Wait(TimeSpan.FromSeconds(5)));
+ Assert.Equal(workerCount, loaderThreads.Count);
+
+ Task dispose = Task.Run(streamer.Dispose);
+ await Task.Delay(100);
+ // Dispose must be blocked on the still-building workers.
+ Assert.False(dispose.IsCompleted);
+
+ release.Set();
+ await dispose.WaitAsync(TimeSpan.FromSeconds(5));
+
+ Assert.Throws(
+ () => streamer.EnqueueLoad(0x1234FFFFu, LandblockStreamJobKind.LoadFar));
+ Assert.Throws(
+ () => streamer.EnqueueUnload(0x1234FFFFu));
+ Assert.Throws(streamer.ClearPendingLoads);
+ }
+ finally
+ {
+ release.Set();
+ streamer.Dispose();
+ }
+ }
+}
diff --git a/tests/AcDream.Core.Tests/Streaming/LandblockStreamerTests.cs b/tests/AcDream.Core.Tests/Streaming/LandblockStreamerTests.cs
index 122988a9..03a466f9 100644
--- a/tests/AcDream.Core.Tests/Streaming/LandblockStreamerTests.cs
+++ b/tests/AcDream.Core.Tests/Streaming/LandblockStreamerTests.cs
@@ -55,13 +55,19 @@ public class LandblockStreamerTests
System.Array.Empty(),
System.Array.Empty());
+ // workerCount: 1 — this test asserts a GLOBAL execution order across
+ // four DIFFERENT landblock ids, which only the serial degenerate pool
+ // guarantees (#418). The pool contract orders jobs per landblock id
+ // and prefers Near-tier per lane; the cross-lane variant lives in
+ // AcDream.App.Tests LandblockStreamerPoolTests.
using var streamer = new LandblockStreamer(
loadLandblock: (id, kind) =>
{
callOrder.Add((id, kind));
return new LoadedLandblock(id, new LandBlock(), System.Array.Empty());
},
- buildMeshOrNull: (_, _) => stubMesh);
+ buildMeshOrNull: (_, _) => stubMesh,
+ workerCount: 1);
streamer.EnqueueLoad(0xAAAAFFFFu, LandblockStreamJobKind.LoadFar);
streamer.EnqueueLoad(0xBBBBFFFFu, LandblockStreamJobKind.LoadFar);