From 39967e78bd2a752a29e780d52baf8aae79f08e30 Mon Sep 17 00:00:00 2001 From: Erik Date: Mon, 17 Aug 2026 19:34:20 +0200 Subject: [PATCH] perf #418: parallelize landblock builds across a striped worker pool Login publishes the 25x25 window at a flat 32 blocks/s (~27 s in the tunnel). The reveal-timing probe A/B (695a27b4) showed the consumer budget env ceilings change nothing, which was read as producer-limited: one "acdream.streaming.worker" thread, ~31 ms/block. This replaces the single worker with min(ProcessorCount-2, 8) workers, floor 1. Design: striped/affinity dispatch. Each worker owns one unbounded lane channel plus its own high/low priority queues; jobs route to lane = ((id >> 16) * 2654435761) % N (the low word of a landblock id is constant, so the id is mixed before reduction). Striping was chosen over a shared queue + in-flight conflict tracker because it preserves the per-landblock contract structurally rather than by bookkeeping: every job for one id lives on one lane, so per-id enqueue order IS execution and completion-arrival order, and the same-landblock supersede rules (PromoteToNear removes queued LoadFar/Unload) keep seeing every queued job for that id. Contract, point by point: - Per-landblock ordering: same id -> same lane -> serial FIFO. - ClearLoads: broadcast to every lane inside the same _inboxGate lock that serializes enqueues, so any load enqueued before ClearPendingLoads() returns sits ahead of its lane's ClearLoads copy in that lane's FIFO and is dropped at read time, exactly like the single-thread path. Already-dequeued builds still complete (now up to one per worker instead of one total); StreamingController's SweepCollapsed already unloads those uniformly. - Priority: per-lane high/low split unchanged. Cross-lane, priority is not globally ordered (a lane cannot run another lane's job), which the contract permits; near-tier jobs hash-spread across lanes and are preferred within each. - Outbox: SingleWriter flipped to false; nothing assumed single-writer (PublishResult already used TryWrite + an Interlocked backlog, and the consumer's peek->read head-stability holds because only the single reader ever moves the head). Cross-landblock arrival order was verified arbitrary-tolerant before relying on it: StreamingController.AdmitCompletions classifies each result independently into per-priority FIFOs (generation staleness + per-landblock retirement blocking); per-landblock arrival order is preserved by striping. - Crash surface: per-worker. The first real crash publishes WorkerCrashed (prefixed "worker N:" in pools > 1), sets _workerFailure, completes every lane, and cancels the pool (a crash still ends all processing, as before); siblings that merely observe the closed lanes (ChannelClosedException) exit quietly instead of reporting spurious crashes; the outbox completes only when the LAST worker exits so no in-flight completions are dropped. - Disposal: joins every worker under the same _disposeGate; Start stays idempotent and dispose-serialized. Thread-safety audit of the production build closures (SessionPlayerComposition), per shared object: - DatCollection (every read in LandblockBuildFactory.BuildLocked: LandblockLoader.Load, SceneryGenerator.Generate, SetupMesh.Flatten, CellMesh.Build, GfxObjBounds.Get, GfxObjDegradeResolver): NOT thread-safe; already serialized under the shared _datLock, which BuildLocked holds for the whole read transaction. Unchanged; the probe run measured hold 0-13 ms / wait <= 12 ms during the login window, so the lock is not the new bottleneck and the build was NOT serialized beyond it. - PakPreparedAssetSource / PakReader (BuildPreparedCollisionClosure, outside the lock): immutable TOC array + read-only MemoryMappedViewAccessor random-access reads + ConcurrentDictionary verdict caches - safe for N concurrent readers (Slice I3 design; the headless SharedPreparedCollisionCache wrapper is fully lock-protected). - LandblockMesh.Build (outside the lock): pure math over the dat record + the composition-time height table + the immutable TerrainBlendingContext record; the shared SurfaceCache is a ConcurrentDictionary and BuildSurface is deterministic, so its lookup-or-build race is last-write-wins-benign (the code already documented exactly this). - PhysicsDiagnostics probe statics: read-only bools + thread-safe Console writes. MEASURED OUTCOME (gate 4): the timing acceptance did NOT pass, and per the task contract that is reported, not tuned around. With 8 workers on this 16-core machine all 625 builds complete in ~203 ms (ACDREAM_PROBE_TELEPORT BUILD lines t=3475390..3475593) - the producer is off the critical path - but loaded= still advances at exactly +32/1000 ms and SUMMARY totalMs measured 27395 and 27503 across two runs (baseline 26728). The 32/s pacer is in the consumer admission/publication path and is not governed by the StreamingWorkBudgetOptions env ceilings. #418 stays IN-PROGRESS on the consumer side; see docs/ISSUES.md for the evidence chain. Tests: per-landblock ordering under 4-worker contention, cross-lane ClearLoads drop, per-lane near-before-far preference, pool-of-1 serial equivalence, disposal joining every worker, lane-spread guard, and worker-count validation (LandblockStreamerPoolTests). Two existing tests asserted a GLOBAL cross-landblock execution order - a serial implementation detail, not the contract - and now pin workerCount: 1 with justification comments (LoadNear_OvertakesQueuedFarLoads, TwoQueuedLoads_RetainTheirDistinctOriginAndGeneration). Gates: Release build 0 errors; App suite 5575 passed / 3 skipped (5568 + 7 new); Runtime suite 1756/0. Co-Authored-By: Claude Fable 5 --- docs/ISSUES.md | 31 ++ .../Streaming/LandblockStreamer.cs | 226 ++++++--- .../Streaming/LandblockBuildOriginTests.cs | 6 +- .../Streaming/LandblockStreamerPoolTests.cs | 431 ++++++++++++++++++ .../Streaming/LandblockStreamerTests.cs | 8 +- 5 files changed, 643 insertions(+), 59 deletions(-) create mode 100644 tests/AcDream.App.Tests/Streaming/LandblockStreamerPoolTests.cs 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);