From b8f6317fe135e4e46ab8c7b945c2b1041742dbe1 Mon Sep 17 00:00:00 2001 From: Erik Date: Fri, 24 Jul 2026 17:48:24 +0200 Subject: [PATCH] feat(streaming): enforce typed completion queues Replace the flat deferred list, priority scan, unload bypass, and count-only execution cap with exact destination/control/unload/Near/Far FIFOs behind one typed frame meter. Price worker results before adoption, retain exact retry identity, reject stale generations without payload retention, and publish queue pressure through lifecycle diagnostics. Tests: dotnet build AcDream.slnx -c Release --no-restore; dotnet test AcDream.slnx -c Release --no-restore (8138 passed, 5 skipped) --- docs/architecture/acdream-architecture.md | 15 +- docs/plans/2026-04-11-roadmap.md | 9 +- docs/plans/2026-05-12-milestones.md | 7 +- ...runtime-slice-e-cost-budgeted-streaming.md | 26 +- .../Composition/SessionPlayerComposition.cs | 2 +- .../Streaming/LandblockStreamer.cs | 49 +- .../Streaming/StreamingCompletionQueue.cs | 270 +++++++++ .../Streaming/StreamingController.cs | 562 +++++++++--------- .../Streaming/StreamingWorkBudget.cs | 29 +- .../Streaming/StreamingWorkBudgetOptions.cs | 43 ++ .../LandblockPresentationPipelineTests.cs | 12 +- .../StreamingCompletionQueueTests.cs | 207 +++++++ .../Streaming/StreamingWorkBudgetTests.cs | 301 +++++++++- .../Streaming/LandblockStreamerTests.cs | 28 + .../StreamingControllerPriorityApplyTests.cs | 37 +- 15 files changed, 1261 insertions(+), 336 deletions(-) create mode 100644 src/AcDream.App/Streaming/StreamingCompletionQueue.cs create mode 100644 tests/AcDream.App.Tests/Streaming/StreamingCompletionQueueTests.cs diff --git a/docs/architecture/acdream-architecture.md b/docs/architecture/acdream-architecture.md index fa41e8cc..e60a036c 100644 --- a/docs/architecture/acdream-architecture.md +++ b/docs/architecture/acdream-architecture.md @@ -712,11 +712,16 @@ partially hydrated cell membership or shell placement to the renderer. update-thread streaming cost. `StreamingWorkMeter` owns the frame-scoped admission ledger; `LandblockStreamResultCost` assigns deterministic charges to immutable completion arrays and logical retained entries without claiming to -measure CLR allocator overhead. During Slice E1 the ledger is observational: -the existing scheduler still executes its accepted work, while lifecycle -artifacts record would-yield limits, elapsed work, retained backlog bytes/age, -and pending publication/retirement counts. Enforcement belongs to the explicit -Slice E scheduler, not to presentation owners or the worker. +measure CLR allocator overhead. `LandblockStreamer` exposes one allocation-free +single-consumer peek/read source so the update thread prices a result before +adopting it. `StreamingController` is the sole scheduler: it admits through the +typed meter into reusable destination/control/unload/Near/Far FIFOs, preserves +exact reference identity and retry position, rejects stale generations, and +executes every class through the same meter. Priority changes order, never the +budget. Lifecycle artifacts publish worker and per-class backlog, retained +bytes/age, yields, oversizes, overruns, and pending publication/retirement +facts. Publication remains transaction-atomic until Slice E4; retirement +quiesce/cursors belong to E3. ### World-reveal readiness ownership diff --git a/docs/plans/2026-04-11-roadmap.md b/docs/plans/2026-04-11-roadmap.md index f8104a60..b2f2e2bb 100644 --- a/docs/plans/2026-04-11-roadmap.md +++ b/docs/plans/2026-04-11-roadmap.md @@ -1605,8 +1605,13 @@ port in any phase — no separate listing here. > and dense-Arwic physical routes passed; third-visit cache residence was > non-growing and no checkpoint retained staged or retiring bytes. [Residency > report](../research/2026-07-24-slice-d-unified-residency-report.md). Slice E -> cost-budgeted streaming/retirement is active. F/G and J retain their explicit -> approval gates. +> cost-budgeted streaming/retirement is active. E2 has replaced the flat +> deferred list, count-only throttle, priority hunt, and unload bypass with one +> typed-meter scheduler over explicit stable destination/control/unload/Near/ +> Far queues; production prices worker results before adoption and lifecycle +> artifacts expose exact backlog/retained-byte facts. E3 next adds immediate +> old-generation quiesce and cursored retirement. F/G and J retain their +> explicit approval gates. **Spec:** `docs/superpowers/specs/2026-07-05-modern-pipeline-design.md` (the umbrella design — read it first). **Goal:** smoothness first (no frame over diff --git a/docs/plans/2026-05-12-milestones.md b/docs/plans/2026-05-12-milestones.md index ff74caa0..a058c9d4 100644 --- a/docs/plans/2026-05-12-milestones.md +++ b/docs/plans/2026-05-12-milestones.md @@ -93,9 +93,10 @@ The separately authorized modern-runtime performance program has completed Slices A–D: corrected measurement, prepared-package bake/dedup, package-only production streaming, and typed unified residency. Its physical capped, uncapped, and dense connected gates pass without a visual-quality change. -Slice E cost-budgets portal publication and retirement next. Slices F/G remain -behind their explicit ECS approval gate, and Slice J remains behind its -gameplay-owner approval gate. +Slice E0–E2 now define and enforce typed completion admission through explicit +stable priority queues; E3 adds immediate old-generation quiesce and cursored +retirement next. Slices F/G remain behind their explicit ECS approval gate, +and Slice J remains behind its gameplay-owner approval gate. Slice 1's 18-cell favorite-spell overflow bar is user-accepted. Slice 2's status hand selected-object availability and existing diff --git a/docs/plans/2026-07-24-modern-runtime-slice-e-cost-budgeted-streaming.md b/docs/plans/2026-07-24-modern-runtime-slice-e-cost-budgeted-streaming.md index 37826461..964dc2b8 100644 --- a/docs/plans/2026-07-24-modern-runtime-slice-e-cost-budgeted-streaming.md +++ b/docs/plans/2026-07-24-modern-runtime-slice-e-cost-budgeted-streaming.md @@ -1,6 +1,6 @@ # Modern Runtime Slice E — Cost-budgeted streaming and retirement -**Status:** active — E0/E1 complete; E2 explicit admission queues next +**Status:** active — E0/E1/E2 complete; E3 quiesce and budgeted retirement next **Program:** `docs/plans/2026-07-24-modern-runtime-architecture.md` **Baseline:** Slice D closeout commit `66690805` **Behavior contract:** no visual-quality or retail-behavior reduction @@ -255,6 +255,30 @@ retry, stale-generation, and deferred-compaction tests remain green. - Keep current publication atomic until E4, but execute only work admitted by the typed meter. +**Complete 2026-07-24.** Production now consumes the worker channel through +one allocation-free `TryPeek`/`TryRead` source, prices the exact immutable +result before adoption, and admits it through the typed completion-count and +retained-CPU dimensions. Accepted work lives in reusable stable FIFOs for +destination, control, unload, Near, and Far classes. FIFO order is exact inside +each class, a blocked publication head does not reorder its tail, callback +appends cannot invalidate an active receipt, and exact-reference compaction +cannot conflate equal records or GUID/landblock reuse. + +`_deferredApply`, `MaxDrainIterations`, direct priority hunting, and the unload +budget bypass are gone. Destination and unload work win queue order but consume +the same execution meter. One indivisible publication may make progress after +admission work, but that privilege is granted only once per frame and every +oversize/overrun remains named. Stale generations consume a bounded admission +without retaining their payload. The old completion-count quality setting now +selects a complete scaled time/count/byte/operation profile rather than acting +as a hidden second throttle. + +Validation: 240 focused App streaming tests and 58 focused Core streamer/ +controller tests passed; the complete App suite passed 3,656 / 3 skipped, the +complete solution passed 8,138 / 5 skipped, and the Release solution build was +clean. Physical connected and visual gates are intentionally deferred until +E6; the active desktop was RDP and is not valid performance evidence. + ### E3 — Quiesce and budgeted retirement - Add destination-generation world quiesce. diff --git a/src/AcDream.App/Composition/SessionPlayerComposition.cs b/src/AcDream.App/Composition/SessionPlayerComposition.cs index d646c55b..f240d398 100644 --- a/src/AcDream.App/Composition/SessionPlayerComposition.cs +++ b/src/AcDream.App/Composition/SessionPlayerComposition.cs @@ -290,7 +290,7 @@ internal sealed class SessionPlayerCompositionPhase d.WorldOrigin.CenterX, d.WorldOrigin.CenterY))), enqueueUnload: streamerLease.Resource.EnqueueUnload, - drainCompletions: streamerLease.Resource.DrainCompletions, + completionSource: streamerLease.Resource, state: live.WorldState, nearRadius: nearRadius, farRadius: farRadius, diff --git a/src/AcDream.App/Streaming/LandblockStreamer.cs b/src/AcDream.App/Streaming/LandblockStreamer.cs index 84f2e9c8..234b90e1 100644 --- a/src/AcDream.App/Streaming/LandblockStreamer.cs +++ b/src/AcDream.App/Streaming/LandblockStreamer.cs @@ -43,7 +43,7 @@ namespace AcDream.App.Streaming; /// methods are thread-safe. /// /// -public sealed class LandblockStreamer : IDisposable +public sealed class LandblockStreamer : IDisposable, ILandblockCompletionSource { /// /// Default drain batch size. Tuned to cap GPU upload work the render @@ -61,6 +61,7 @@ public sealed class LandblockStreamer : IDisposable private readonly object _inboxGate = new(); private Thread? _worker; private Exception? _workerFailure; + private int _completionBacklog; private int _disposed; private readonly object _disposeGate = new(); private bool _disposeCompleted; @@ -276,11 +277,37 @@ public sealed class LandblockStreamer : IDisposable public IReadOnlyList DrainCompletions(int maxBatchSize = DefaultDrainBatchSize) { var batch = new List(maxBatchSize); - while (batch.Count < maxBatchSize && _outbox.Reader.TryRead(out var result)) + while (batch.Count < maxBatchSize && TryRead(out var result)) + { + if (result is null) + throw new InvalidOperationException( + "The completion channel returned a null result."); batch.Add(result); + } return batch; } + public int BacklogCount => Math.Max( + 0, + System.Threading.Volatile.Read(ref _completionBacklog)); + + public bool TryPeek(out LandblockStreamResult? result) => + _outbox.Reader.TryPeek(out result); + + public bool TryRead(out LandblockStreamResult? result) + { + if (!_outbox.Reader.TryRead(out result)) + return false; + System.Threading.Interlocked.Decrement(ref _completionBacklog); + return true; + } + + private void PublishResult(LandblockStreamResult result) + { + if (_outbox.Writer.TryWrite(result)) + System.Threading.Interlocked.Increment(ref _completionBacklog); + } + private void WorkerLoop() { var highPriority = new Queue(); @@ -337,7 +364,7 @@ public sealed class LandblockStreamer : IDisposable _workerFailure = ex; _inbox.Writer.TryComplete(ex); } - _outbox.Writer.TryWrite(new LandblockStreamResult.WorkerCrashed(ex.ToString())); + PublishResult(new LandblockStreamResult.WorkerCrashed(ex.ToString())); } finally { @@ -426,13 +453,13 @@ public sealed class LandblockStreamer : IDisposable var build = _loadLandblock(load.Request); if (build is null) { - _outbox.Writer.TryWrite(new LandblockStreamResult.Failed( + PublishResult(new LandblockStreamResult.Failed( load.LandblockId, "LandblockLoader.Load returned null", load.Generation)); break; } if (build.Origin != load.Origin) { - _outbox.Writer.TryWrite(new LandblockStreamResult.Failed( + PublishResult(new LandblockStreamResult.Failed( load.LandblockId, $"Landblock build origin {build.Origin} did not match request origin {load.Origin}", load.Generation)); @@ -444,18 +471,18 @@ public sealed class LandblockStreamer : IDisposable var promotedMesh = _buildMeshOrNull(load.LandblockId, lb); if (promotedMesh is null) { - _outbox.Writer.TryWrite(new LandblockStreamResult.Failed( + PublishResult(new LandblockStreamResult.Failed( load.LandblockId, "buildMeshOrNull returned null", load.Generation)); break; } - _outbox.Writer.TryWrite(new LandblockStreamResult.Promoted( + PublishResult(new LandblockStreamResult.Promoted( load.LandblockId, build, promotedMesh, load.Generation)); break; } var mesh = _buildMeshOrNull(load.LandblockId, lb); if (mesh is null) { - _outbox.Writer.TryWrite(new LandblockStreamResult.Failed( + PublishResult(new LandblockStreamResult.Failed( load.LandblockId, "buildMeshOrNull returned null", load.Generation)); break; } @@ -485,18 +512,18 @@ public sealed class LandblockStreamer : IDisposable PhysicsDatBundle.Empty); build = new LandblockBuild(lb, Origin: build.Origin); } - _outbox.Writer.TryWrite(new LandblockStreamResult.Loaded( + PublishResult(new LandblockStreamResult.Loaded( load.LandblockId, tier, build, mesh, load.Generation)); } catch (Exception ex) { - _outbox.Writer.TryWrite(new LandblockStreamResult.Failed( + PublishResult(new LandblockStreamResult.Failed( load.LandblockId, ex.ToString(), load.Generation)); } break; case LandblockStreamJob.Unload unload: - _outbox.Writer.TryWrite(new LandblockStreamResult.Unloaded( + PublishResult(new LandblockStreamResult.Unloaded( unload.LandblockId, unload.Generation)); break; diff --git a/src/AcDream.App/Streaming/StreamingCompletionQueue.cs b/src/AcDream.App/Streaming/StreamingCompletionQueue.cs new file mode 100644 index 00000000..f432b517 --- /dev/null +++ b/src/AcDream.App/Streaming/StreamingCompletionQueue.cs @@ -0,0 +1,270 @@ +using System.Diagnostics; + +namespace AcDream.App.Streaming; + +/// +/// Allocation-free single-consumer view of the streaming worker's completion +/// channel. Peek plus read lets the update thread price a payload before +/// adopting it into the scheduler. +/// +public interface ILandblockCompletionSource +{ + int BacklogCount { get; } + bool TryPeek(out LandblockStreamResult? result); + bool TryRead(out LandblockStreamResult? result); +} + +internal enum StreamingCompletionPriority : byte +{ + Destination = 0, + Control = 1, + Unload = 2, + Near = 3, + Far = 4, +} + +internal readonly record struct StreamingQueuedCompletion( + LandblockStreamResult Result, + LandblockStreamCostEstimate Estimate, + StreamingCompletionPriority Priority, + ulong Generation, + long Sequence, + long EnqueuedTimestamp) +{ + public StreamingWorkCost AdmissionCost => new( + CompletionAdmissions: Estimate.Work.CompletionAdmissions, + AdoptedCpuBytes: Estimate.Work.AdoptedCpuBytes); + + public StreamingWorkCost ExecutionCost => new( + EntityOperations: Estimate.Work.EntityOperations, + GpuUploadBytes: Estimate.Work.GpuUploadBytes, + GlRetireOperations: + Result is LandblockStreamResult.Unloaded ? 1 : 0); +} + +internal readonly record struct StreamingCompletionQueueSnapshot( + int Count, + long RetainedCpuBytes, + double OldestAgeMilliseconds, + int Destination, + int Control, + int Unload, + int Near, + int Far); + +/// +/// Render-thread-owned stable FIFO queues. Ordering is exact inside one +/// generation/priority; callbacks append and never invalidate an active node. +/// +internal sealed class StreamingCompletionQueue +{ + private readonly Queue[] _queues = + Enumerable.Range(0, Enum.GetValues().Length) + .Select(static _ => new Queue()) + .ToArray(); + private long _retainedCpuBytes; + private int _count; + + public int Count => _count; + public long RetainedCpuBytes => _retainedCpuBytes; + + public void Enqueue(StreamingQueuedCompletion completion) + { + _queues[(int)completion.Priority].Enqueue(completion); + _count++; + _retainedCpuBytes = SaturatingAdd( + _retainedCpuBytes, + completion.Estimate.Work.AdoptedCpuBytes); + } + + public bool TryPeekNext( + Func isBlocked, + out StreamingQueuedCompletion? completion) + { + ArgumentNullException.ThrowIfNull(isBlocked); + for (int priority = 0; priority < _queues.Length; priority++) + { + Queue queue = _queues[priority]; + if (queue.Count == 0) + continue; + StreamingQueuedCompletion head = queue.Peek(); + if (isBlocked(head.Result)) + continue; + + completion = head; + return true; + } + + completion = null; + return false; + } + + public void RemoveHead(StreamingQueuedCompletion completion) + { + Queue queue = + _queues[(int)completion.Priority]; + if (queue.Count == 0) + { + throw new InvalidOperationException( + "Streaming completion priority FIFO is empty."); + } + + StreamingQueuedCompletion head = queue.Peek(); + if (head.Sequence != completion.Sequence + || !ReferenceEquals(head.Result, completion.Result)) + { + throw new InvalidOperationException( + "Streaming completion is not the head of its priority FIFO."); + } + + queue.Dequeue(); + Release(completion); + } + + public int RemoveResults( + IReadOnlyList results, + Func shouldRemove) + { + ArgumentNullException.ThrowIfNull(results); + ArgumentNullException.ThrowIfNull(shouldRemove); + int removed = 0; + for (int priority = 0; priority < _queues.Length; priority++) + { + Queue queue = _queues[priority]; + int count = queue.Count; + for (int entry = 0; entry < count; entry++) + { + StreamingQueuedCompletion current = queue.Dequeue(); + bool matches = false; + for (int i = 0; i < results.Count; i++) + { + if (ReferenceEquals(current.Result, results[i])) + { + matches = true; + break; + } + } + + if (matches && shouldRemove(current.Result)) + { + Release(current); + removed++; + } + else + { + queue.Enqueue(current); + } + } + } + return removed; + } + + public void Clear() + { + for (int i = 0; i < _queues.Length; i++) + _queues[i].Clear(); + _count = 0; + _retainedCpuBytes = 0; + } + + public StreamingCompletionQueueSnapshot CaptureSnapshot() + { + long now = Stopwatch.GetTimestamp(); + long oldest = now; + bool found = false; + for (int i = 0; i < _queues.Length; i++) + { + foreach (StreamingQueuedCompletion completion in _queues[i]) + { + oldest = Math.Min(oldest, completion.EnqueuedTimestamp); + found = true; + } + } + + return new StreamingCompletionQueueSnapshot( + _count, + _retainedCpuBytes, + found + ? Math.Max(0L, now - oldest) * 1000.0 / Stopwatch.Frequency + : 0, + _queues[(int)StreamingCompletionPriority.Destination].Count, + _queues[(int)StreamingCompletionPriority.Control].Count, + _queues[(int)StreamingCompletionPriority.Unload].Count, + _queues[(int)StreamingCompletionPriority.Near].Count, + _queues[(int)StreamingCompletionPriority.Far].Count); + } + + private void Release(StreamingQueuedCompletion completion) + { + _count--; + _retainedCpuBytes = Math.Max( + 0, + _retainedCpuBytes - completion.Estimate.Work.AdoptedCpuBytes); + } + + private static long SaturatingAdd(long left, long right) => + left > long.MaxValue - right ? long.MaxValue : left + right; +} + +/// +/// Compatibility adapter for deterministic tests that still provide the +/// former batch-drain delegate. Production uses +/// directly. +/// +internal sealed class DelegateLandblockCompletionSource( + Func> drain) + : ILandblockCompletionSource +{ + private readonly Func> _drain = + drain ?? throw new ArgumentNullException(nameof(drain)); + private LandblockStreamResult? _peeked; + + public int BacklogCount => _peeked is null ? 0 : 1; + + public bool TryPeek(out LandblockStreamResult? result) + { + if (_peeked is null) + { + IReadOnlyList batch = _drain(1); + if (batch.Count > 1) + { + throw new InvalidOperationException( + "A single-result completion drain returned more than one result."); + } + if (batch.Count == 0) + { + result = null; + return false; + } + _peeked = batch[0]; + } + + result = _peeked; + return true; + } + + public bool TryRead(out LandblockStreamResult? result) + { + if (_peeked is not null) + { + result = _peeked; + _peeked = null; + return true; + } + + IReadOnlyList batch = _drain(1); + if (batch.Count > 1) + { + throw new InvalidOperationException( + "A single-result completion drain returned more than one result."); + } + if (batch.Count == 0) + { + result = null; + return false; + } + + result = batch[0]; + return true; + } +} diff --git a/src/AcDream.App/Streaming/StreamingController.cs b/src/AcDream.App/Streaming/StreamingController.cs index 38ae21ae..6dc08c8d 100644 --- a/src/AcDream.App/Streaming/StreamingController.cs +++ b/src/AcDream.App/Streaming/StreamingController.cs @@ -19,17 +19,12 @@ namespace AcDream.App.Streaming; /// public sealed class StreamingController : IStreamingFrameBackend { - private readonly record struct DeferredCompletion( - LandblockStreamResult Result, - StreamingWorkCost Cost, - long EnqueuedTimestamp); - private sealed class OriginRecenterRetirement { public bool RadiiConverged; public bool GenerationAdvanced; public bool PendingLoadsCleared; - public bool DeferredApplyCleared; + public bool CompletionQueueCleared; public bool RegionCleared; public List? ResidentIds; public int RetirementCursor; @@ -41,13 +36,18 @@ public sealed class StreamingController : IStreamingFrameBackend private readonly Action _enqueueLoad; private readonly Action _enqueueUnload; - private readonly Func> _drainCompletions; + private readonly ILandblockCompletionSource _completionSource; private readonly Action? _clearPendingLoads; private readonly LandblockPresentationPipeline _presentation; - private readonly StreamingWorkBudget _workBudget; + private readonly Func + _isPublicationBlockedByRetirement; + private readonly StreamingWorkBudgetOptions _configuredWorkBudgetOptions; + private StreamingWorkBudget _workBudget; private StreamingWorkMeter? _activeWorkMeter; private StreamingWorkMeterSnapshot _lastWorkMeter; - private long _deferredAdoptedCpuBytes; + private readonly StreamingCompletionQueue _completionQueue = new(); + private long _nextCompletionSequence; + private int _legacyCompletionProfile = 4; private readonly GpuWorldState _state; private StreamingRegion? _region; @@ -88,34 +88,28 @@ public sealed class StreamingController : IStreamingFrameBackend public int FarRadius { get; private set; } /// - /// Cap on completions drained per call. The cap is - /// the GPU upload budget for one frame: terrain mesh + per-entity GfxObj - /// sub-mesh uploads + texture uploads for one landblock take a few ms; - /// applying 25 of them in a single frame produces a memory spike - /// (observed: out-of-memory crash on the 5×5 first-frame load). - /// - /// - /// 4 is the original async-streamer value; it spreads a 5×5 first-frame - /// load over ~7 frames (~116ms at 60fps), which is below the human - /// perception threshold. Spawn races that previously dropped entities - /// while landblocks were in flight are now handled by - /// 's pending-spawn list, so spreading - /// completions doesn't lose any data. - /// + /// Compatibility quality selector for callers that still express the + /// streaming profile as a completion count. Four preserves the configured + /// High profile; other positive values scale every typed time/count/byte + /// dimension together. It is not an independent execution throttle. /// - public int MaxCompletionsPerFrame { get; set; } = 4; + public int MaxCompletionsPerFrame + { + get => _legacyCompletionProfile; + set + { + StreamingWorkBudgetOptions profile = + _configuredWorkBudgetOptions + .ScaleForLegacyCompletionCount(value); + _legacyCompletionProfile = value; + _workBudget = profile.ToBudget(); + } + } /// - /// #138: the teleport destination landblock id. When non-zero, - /// applies this landblock's - /// or - /// completion immediately - /// — even if it sits past the - /// position in the outbox — so the player can materialise at the - /// destination without waiting for every earlier-queued landblock to - /// drain first. Non-priority completions drained past it are buffered - /// in and applied over subsequent frames, - /// so no completions are lost and there is no GPU spike. + /// #138: the teleport destination landblock id. Loaded or promoted results + /// inside its priority ring enter the destination FIFO, ahead of ordinary + /// Near/Far work, but still consume the same typed frame budget. /// /// Set by GameWindow when a teleport starts; reset to 0 /// once the destination landblock has been applied (or when the @@ -124,32 +118,39 @@ public sealed class StreamingController : IStreamingFrameBackend public uint PriorityLandblockId { get; set; } /// - /// 2026-06-22: the radius (in landblocks, Chebyshev) around - /// that eager-applies - /// ahead of the per-frame budget. 0 (default) = only the single center landblock — the - /// original priority behaviour. A teleport sets this to the near ring so the player's - /// IMMEDIATE SURROUNDINGS (terrain + collision + scenery) are resident on arrival, not - /// just the one landblock they stand on. Without it, only the destination landblock - /// applies immediately and everything around it drains at , - /// so the player arrives to a near-empty world (the "Fort Tethana only one landblock - /// loaded" symptom) and their cell-walk can't root into neighbour cells yet (the - /// transient walk-through-walls). The far ring still drains at the budget. The eager - /// apply runs while the portal viewport replaces the world, so the GPU spike is hidden. + /// Radius in landblocks (Chebyshev distance) around + /// classified as destination work. + /// Zero means only the center. This changes queue priority, never budget + /// enforcement or atomic publication semantics. /// public int PriorityRadius { get; set; } // [FRAME-DIAG] (read by GameWindow's ACDREAM_WB_DIAG rollup): the standing - // deferred-LOAD backlog. A non-zero value during/after a teleport is the - // GPU-upload tail draining at MaxCompletionsPerFrame. Removable probe surface. - public int DeferredApplyBacklog => _deferredApply.Count; + // accepted completion backlog. A non-zero value during/after a teleport is + // the typed-budget publication tail. Removable compatibility probe surface. + public int DeferredApplyBacklog => _completionQueue.Count; public int PendingRetirementCount => _presentation.PendingRetirementCount; - public StreamingWorkDiagnostics WorkDiagnostics => new( - _lastWorkMeter, - _deferredApply.Count, - _deferredAdoptedCpuBytes, - OldestDeferredAgeMilliseconds(), - _presentation.PendingPublicationCount, - _presentation.PendingRetirementCount); + public StreamingWorkDiagnostics WorkDiagnostics + { + get + { + StreamingCompletionQueueSnapshot queued = + _completionQueue.CaptureSnapshot(); + return new StreamingWorkDiagnostics( + _lastWorkMeter, + queued.Count, + queued.RetainedCpuBytes, + queued.OldestAgeMilliseconds, + _presentation.PendingPublicationCount, + _presentation.PendingRetirementCount, + _completionSource.BacklogCount, + queued.Destination, + queued.Control, + queued.Unload, + queued.Near, + queued.Far); + } + } internal bool IsCollapsedToDungeon => _collapsed; /// @@ -195,9 +196,6 @@ public sealed class StreamingController : IStreamingFrameBackend public int FullWindowRetirementCount { get; private set; } public int LastFullWindowRetirementLandblockCount { get; private set; } - // Completions that were drained past a priority item get buffered here - // so they still apply over subsequent frames without loss. - private readonly List _deferredApply = new(); /// /// Internal compatibility seam for hermetic controller policy tests. Live /// composition cannot supply presentation callbacks; it must use the @@ -216,7 +214,8 @@ public sealed class StreamingController : IStreamingFrameBackend Action? clearPendingLoads = null, Action? onLandblockLoaded = null, Action? ensureEnvCellMeshes = null, - LandblockRetirementCoordinator? retirementCoordinator = null) + LandblockRetirementCoordinator? retirementCoordinator = null, + StreamingWorkBudgetOptions? workBudgetOptions = null) : this( enqueueLoad, enqueueUnload, @@ -232,7 +231,8 @@ public sealed class StreamingController : IStreamingFrameBackend retirementCoordinator, removeTerrain, demoteNearLayer), - clearPendingLoads) + clearPendingLoads, + workBudgetOptions) { } @@ -241,7 +241,7 @@ public sealed class StreamingController : IStreamingFrameBackend /// retirement owner, so no legacy publication delegate can be supplied or /// silently ignored. /// - public StreamingController( + internal StreamingController( Action enqueueLoad, Action enqueueUnload, Func> drainCompletions, @@ -251,16 +251,47 @@ public sealed class StreamingController : IStreamingFrameBackend LandblockPresentationPipeline presentationPipeline, Action? clearPendingLoads = null, StreamingWorkBudgetOptions? workBudgetOptions = null) + : this( + enqueueLoad, + enqueueUnload, + new DelegateLandblockCompletionSource(drainCompletions), + state, + nearRadius, + farRadius, + presentationPipeline, + clearPendingLoads, + workBudgetOptions) + { + } + + /// + /// Allocation-free production completion path. Peek/read is single-consumer + /// and lets the scheduler price a result before adopting it. + /// + public StreamingController( + Action enqueueLoad, + Action enqueueUnload, + ILandblockCompletionSource completionSource, + GpuWorldState state, + int nearRadius, + int farRadius, + LandblockPresentationPipeline presentationPipeline, + Action? clearPendingLoads = null, + StreamingWorkBudgetOptions? workBudgetOptions = null) { _enqueueLoad = enqueueLoad; _enqueueUnload = enqueueUnload; - _drainCompletions = drainCompletions; + _completionSource = completionSource + ?? throw new ArgumentNullException(nameof(completionSource)); _clearPendingLoads = clearPendingLoads; _state = state; _presentation = presentationPipeline ?? throw new ArgumentNullException(nameof(presentationPipeline)); - _workBudget = (workBudgetOptions ?? StreamingWorkBudgetOptions.Default) - .ToBudget(); + _isPublicationBlockedByRetirement = + IsPublicationBlockedByRetirement; + _configuredWorkBudgetOptions = + workBudgetOptions ?? StreamingWorkBudgetOptions.Default; + _workBudget = _configuredWorkBudgetOptions.ToBudget(); if (!_presentation.MatchesState(_state)) { throw new ArgumentException( @@ -400,7 +431,8 @@ public sealed class StreamingController : IStreamingFrameBackend Action? clearPendingLoads = null, Action? onLandblockLoaded = null, Action? ensureEnvCellMeshes = null, - LandblockRetirementCoordinator? retirementCoordinator = null) + LandblockRetirementCoordinator? retirementCoordinator = null, + StreamingWorkBudgetOptions? workBudgetOptions = null) : this( (id, kind, _) => enqueueLoad(id, kind), (id, _) => enqueueUnload(id), @@ -414,7 +446,8 @@ public sealed class StreamingController : IStreamingFrameBackend clearPendingLoads, onLandblockLoaded, ensureEnvCellMeshes, - retirementCoordinator) + retirementCoordinator, + workBudgetOptions) { } @@ -429,48 +462,47 @@ public sealed class StreamingController : IStreamingFrameBackend _presentation.BeginNearLayerRetirement(canonical); } - private void AdvanceGeneration() + private bool AdvanceGeneration() { - ConvergePendingPublications(); + if (!ConvergePendingPublications()) + return false; _generation = unchecked(_generation + 1); + return true; } - private void ConvergePendingPublications() + private bool ConvergePendingPublications() { IReadOnlyList pending = _presentation.GetPendingPublicationResults(); if (pending.Count == 0) - return; + return true; try { + bool progressed = false; for (int i = 0; i < pending.Count; i++) { LandblockStreamResult result = pending[i]; - ObserveResultOperation( + if (!TryObserveResultOperation( result, "publication-retry", - () => _presentation.ResumePublication(result)); + () => _presentation.ResumePublication(result), + ensureProgress: !progressed)) + { + return false; + } + progressed = true; } + return true; } finally { // A completed staged result can still be present in the deferred // FIFO that originally retained it. Remove only those exact // objects; unrelated accepted completions keep their order. - _deferredApply.RemoveAll(entry => - { - for (int i = 0; i < pending.Count; i++) - { - if (ReferenceEquals(entry.Result, pending[i])) - { - return !_presentation.HasPendingPublication( - entry.Result); - } - } - return false; - }); - RecalculateDeferredAdoptedBytes(); + _completionQueue.RemoveResults( + pending, + result => !_presentation.HasPendingPublication(result)); } } @@ -518,7 +550,8 @@ public sealed class StreamingController : IStreamingFrameBackend // Complete an admitted publication before this frame can mutate its // desired tier or spatial residence. A later recenter may then demote // or retire the fully known owner set through the normal ledger. - ConvergePendingPublications(); + if (!ConvergePendingPublications()) + return; uint centerId = StreamingRegion.EncodeLandblockId(observerCx, observerCy); @@ -574,7 +607,8 @@ public sealed class StreamingController : IStreamingFrameBackend if (!pending.GenerationAdvanced) { - AdvanceGeneration(); + if (!AdvanceGeneration()) + return; pending.GenerationAdvanced = true; } @@ -625,7 +659,7 @@ public sealed class StreamingController : IStreamingFrameBackend _region = pending.Region; NearRadius = pending.NearRadius; FarRadius = pending.FarRadius; - ClearDeferred(); + _completionQueue.Clear(); _pendingRadiiReconfiguration = null; } } @@ -762,11 +796,12 @@ public sealed class StreamingController : IStreamingFrameBackend bool logTransition = !_collapsed || _collapsedCenter != centerId; if (logTransition) Console.WriteLine($"streaming: dungeon collapse -> 0x{centerId:X8}"); - AdvanceGeneration(); + if (!AdvanceGeneration()) + return; _collapsed = true; _collapsedCenter = centerId; _clearPendingLoads?.Invoke(); - ClearDeferred(); + _completionQueue.Clear(); _region = null; foreach (var id in _state.LoadedLandblockIds) @@ -811,7 +846,7 @@ public sealed class StreamingController : IStreamingFrameBackend /// /// True when is the priority center or within /// landblocks of it (Chebyshev) — i.e. inside the teleport - /// near ring that eager-applies. With the default radius 0 + /// near ring that prioritizes. With the default radius 0 /// this reduces to an exact match on (the original /// single-landblock priority behaviour). /// @@ -829,7 +864,8 @@ public sealed class StreamingController : IStreamingFrameBackend Console.WriteLine( $"streaming: dungeon EXIT-expand -> ({observerCx},{observerCy}) " + $"(was collapsed on 0x{_collapsedCenter:X8})"); - AdvanceGeneration(); + if (!AdvanceGeneration()) + return; _collapsed = false; var rebuilt = new StreamingRegion(observerCx, observerCy, NearRadius, FarRadius); @@ -1027,7 +1063,8 @@ public sealed class StreamingController : IStreamingFrameBackend } if (!transaction.GenerationAdvanced) { - AdvanceGeneration(); + if (!AdvanceGeneration()) + return false; transaction.GenerationAdvanced = true; } if (!transaction.PendingLoadsCleared) @@ -1045,10 +1082,10 @@ public sealed class StreamingController : IStreamingFrameBackend return false; } } - if (!transaction.DeferredApplyCleared) + if (!transaction.CompletionQueueCleared) { - ClearDeferred(); - transaction.DeferredApplyCleared = true; + _completionQueue.Clear(); + transaction.CompletionQueueCleared = true; } if (!transaction.RegionCleared) { @@ -1103,10 +1140,11 @@ public sealed class StreamingController : IStreamingFrameBackend private void BeginFullWindowRetirement() { - AdvanceGeneration(); + if (!AdvanceGeneration()) + return; _collapsed = false; _clearPendingLoads?.Invoke(); - ClearDeferred(); + _completionQueue.Clear(); // Commit the old region boundary before any presentation owner is // advanced. New same-id publications are fenced by the presentation // pipeline's retryable retirement ledger. @@ -1120,164 +1158,171 @@ public sealed class StreamingController : IStreamingFrameBackend } /// - /// Apply streamed completions for this frame. LOADS (terrain mesh GPU uploads) are the - /// expensive part, so they are metered at to avoid a - /// GPU-upload spike; the overflow buffers in and drains over - /// subsequent frames. UNLOADS are cheap (they free GPU buffers — no upload) and are applied - /// IMMEDIATELY, never throttled: a teleport produces a whole window of unloads (~600), and - /// metering them at the load rate left the previous location's terrain resident for seconds - /// (rendering at its old world position as "floating terrain at the horizon"), and rapid - /// hops accumulated them faster than they cleared — a runaway resident count (951 observed - /// vs a 625 window) that also dragged FPS. Loads inside the teleport near ring - /// (, applied during portal travel) likewise bypass the budget so the - /// player materialises in a loaded world. + /// Adopts immutable worker results and advances explicit priority FIFOs + /// under the one typed frame meter. Destination and unload work win queue + /// order but no class bypasses time, bytes, entities, uploads, or retire + /// operation limits. /// private void DrainAndApply() { - // --- Step 1: drain the outbox in bounded chunks. Apply unloads + priority near-ring - // loads immediately; defer every other (budget-metered) load. Draining the whole - // outbox each frame (bounded by MaxDrainIterations) is what lets unloads flush - // promptly regardless of the load backlog — the throttle is on GPU UPLOADS, not on - // freeing them. The drain cap must NOT be gated behind the per-frame load budget - // (the prior version returned once the budget hit 0, stranding the outbox). - const int MaxDrainIterations = 64; // cap at 64 * MaxCompletionsPerFrame drained/frame - int iter = 0; - while (iter++ < MaxDrainIterations) - { - var chunk = _drainCompletions(MaxCompletionsPerFrame); - if (chunk.Count == 0) break; - for (int chunkIndex = 0; chunkIndex < chunk.Count; chunkIndex++) - { - LandblockStreamResult result = chunk[chunkIndex]; - // Reject obsolete hard-window work before it can occupy the - // metered deferred queue. ApplyResult repeats this invariant - // for direct callers and future drain paths. - if (IsStaleGeneration(result) - && !_presentation.HasPendingPublication(result)) - continue; - if (IsPublicationBlockedByRetirement(result)) - Defer(result); - else if (result is LandblockStreamResult.Unloaded - || IsWithinPriorityRing(ResultLandblockId(result))) - { - try - { - ApplyResultObserved( - result, - "priority-or-unload"); - } - catch - { - // The concrete presentation pipeline retains the exact - // unfinished stage. Preserve it and every untouched - // completion from this already-drained chunk in - // original FIFO order. Internal compatibility tests may - // use a callback pipeline that cannot retain a receipt; - // its failed result is consumed, but the untouched - // suffix must still survive. - if (_presentation.HasPendingPublication(result)) - Defer(result); - for (int suffix = chunkIndex + 1; suffix < chunk.Count; suffix++) - Defer(chunk[suffix]); - throw; - } - } - else - Defer(result); // a GPU-upload load — meter it in step 2 - } - } + StreamingWorkMeter meter = _activeWorkMeter + ?? throw new InvalidOperationException( + "Completion scheduling requires an active frame meter."); - // --- Step 2: apply the deferred LOAD backlog at the per-frame budget (FIFO, so - // earlier-queued landblocks win). Caps GPU upload per frame so a big diff doesn't - // spike. _deferredApply now only ever holds loads — unloads were applied in step 1. - int budget = MaxCompletionsPerFrame; - int applied = 0; - int read = 0; - int write = 0; - try + AdmitCompletions(meter); + bool executed = false; + while (_completionQueue.TryPeekNext( + _isPublicationBlockedByRetirement, + out StreamingQueuedCompletion? completion)) { - while (applied < budget && read < _deferredApply.Count) - { - DeferredCompletion entry = _deferredApply[read]; - LandblockStreamResult result = entry.Result; - if (IsPublicationBlockedByRetirement(result)) - { - if (write != read) - _deferredApply[write] = entry; - write++; - read++; - continue; - } + StreamingQueuedCompletion work = completion + ?? throw new InvalidOperationException( + "The completion queue returned a null head."); + StreamingWorkCost executionCost = + IsStaleGeneration(work.Result) + ? default + : work.ExecutionCost; + StreamingWorkAdmission admission = meter.TryReserve( + executionCost, + $"execute-{work.Priority}", + ensureProgress: !executed); + if (admission == StreamingWorkAdmission.Yielded) + break; - // Advance read only after the apply commits. If it throws, the - // finally block removes prior committed entries but retains this - // exact result and the untouched tail for retry. - try - { - ApplyResultObserved( - result, - "deferred-publication", - entry.Cost); - read++; - applied++; - } - catch - { - // A concrete suffix-stage failure keeps its pipeline - // receipt and must remain at this exact FIFO position. The - // internal compatibility callback cannot report its exact - // committed prefix and therefore removes its receipt; - // consume only that failed result so a later frame cannot - // replay it automatically. - if (!_presentation.HasPendingPublication(result)) - read++; - throw; - } + try + { + ApplyResult(work.Result); + _completionQueue.RemoveHead(work); + meter.Complete(); + executed = true; + } + catch + { + // Concrete publication failures retain their exact owner + // receipt and queue head. A compatibility callback without a + // receipt cannot safely replay, so consume only that result. + if (!_presentation.HasPendingPublication(work.Result)) + _completionQueue.RemoveHead(work); + meter.Fail(); + throw; } - } - finally - { - // One linear tail shift replaces budget-many RemoveAt shifts. The - // retained blocked prefix and the untouched FIFO tail keep their - // original relative order. - if (read > write) - RemoveDeferredRange(write, read - write); } } - private void ApplyResultObserved( - LandblockStreamResult result, - string stage, - StreamingWorkCost? knownCost = null) => - ObserveResultOperation( - result, - stage, - () => ApplyResult(result), - knownCost); + private void AdmitCompletions(StreamingWorkMeter meter) + { + while (_completionSource.TryPeek(out LandblockStreamResult? peeked)) + { + LandblockStreamResult result = peeked + ?? throw new InvalidOperationException( + "The completion source returned a null peek."); + bool stale = IsStaleGeneration(result) + && !_presentation.HasPendingPublication(result); + LandblockStreamCostEstimate estimate = + LandblockStreamResultCost.Estimate(result); + StreamingWorkCost admissionCost = stale + ? new StreamingWorkCost(CompletionAdmissions: 1) + : new StreamingWorkCost( + CompletionAdmissions: + estimate.Work.CompletionAdmissions, + AdoptedCpuBytes: + estimate.Work.AdoptedCpuBytes); + StreamingWorkAdmission admission = meter.TryReserve( + admissionCost, + "completion-admission"); + if (admission == StreamingWorkAdmission.Yielded) + return; - private void ObserveResultOperation( + try + { + if (!_completionSource.TryRead( + out LandblockStreamResult? consumed) + || !ReferenceEquals(result, consumed)) + { + throw new InvalidOperationException( + "The single-consumer completion source changed between peek and read."); + } + + if (!stale) + { + _completionQueue.Enqueue(new StreamingQueuedCompletion( + result, + estimate, + ClassifyCompletion(result), + result.Generation, + checked(_nextCompletionSequence++), + Stopwatch.GetTimestamp())); + } + meter.Complete(); + } + catch + { + meter.Fail(); + throw; + } + } + } + + private StreamingCompletionPriority ClassifyCompletion( + LandblockStreamResult result) + { + if (result is LandblockStreamResult.Loaded + or LandblockStreamResult.Promoted + && IsWithinPriorityRing(result.LandblockId)) + { + return StreamingCompletionPriority.Destination; + } + + return result switch + { + LandblockStreamResult.Failed + or LandblockStreamResult.WorkerCrashed => + StreamingCompletionPriority.Control, + LandblockStreamResult.Unloaded => + StreamingCompletionPriority.Unload, + LandblockStreamResult.Promoted + or LandblockStreamResult.Loaded + { + Tier: LandblockStreamTier.Near, + } => + StreamingCompletionPriority.Near, + _ => StreamingCompletionPriority.Far, + }; + } + + private bool TryObserveResultOperation( LandblockStreamResult result, string stage, Action operation, - StreamingWorkCost? knownCost = null) + StreamingWorkCost? knownCost = null, + bool ensureProgress = false) { StreamingWorkMeter? meter = _activeWorkMeter; if (meter is null) { operation(); - return; + return true; } - StreamingWorkCost cost = knownCost + StreamingWorkCost aggregate = knownCost ?? LandblockStreamResultCost.Estimate(result).Work; - StreamingWorkAdmission admission = meter.TryReserve(cost, stage); + StreamingWorkCost cost = new( + EntityOperations: aggregate.EntityOperations, + GpuUploadBytes: aggregate.GpuUploadBytes, + GlRetireOperations: + result is LandblockStreamResult.Unloaded ? 1 : 0); + StreamingWorkAdmission admission = meter.TryReserve( + cost, + stage, + ensureProgress); if (admission == StreamingWorkAdmission.Yielded) - meter.ForceReserveObserved(cost, stage); + return false; try { operation(); meter.Complete(); + return true; } catch { @@ -1286,69 +1331,10 @@ public sealed class StreamingController : IStreamingFrameBackend } } - private void Defer(LandblockStreamResult result) - { - StreamingWorkCost cost = LandblockStreamResultCost.Estimate(result).Work; - _deferredApply.Add(new DeferredCompletion( - result, - cost, - Stopwatch.GetTimestamp())); - _deferredAdoptedCpuBytes = SaturatingAdd( - _deferredAdoptedCpuBytes, - cost.AdoptedCpuBytes); - } - - private void ClearDeferred() - { - _deferredApply.Clear(); - _deferredAdoptedCpuBytes = 0; - } - - private void RemoveDeferredRange(int index, int count) - { - for (int i = index; i < index + count; i++) - { - _deferredAdoptedCpuBytes = Math.Max( - 0, - _deferredAdoptedCpuBytes - - _deferredApply[i].Cost.AdoptedCpuBytes); - } - _deferredApply.RemoveRange(index, count); - } - - private void RecalculateDeferredAdoptedBytes() - { - long total = 0; - for (int i = 0; i < _deferredApply.Count; i++) - { - total = SaturatingAdd( - total, - _deferredApply[i].Cost.AdoptedCpuBytes); - } - _deferredAdoptedCpuBytes = total; - } - - private double OldestDeferredAgeMilliseconds() - { - if (_deferredApply.Count == 0) - return 0; - - long oldest = _deferredApply[0].EnqueuedTimestamp; - for (int i = 1; i < _deferredApply.Count; i++) - oldest = Math.Min(oldest, _deferredApply[i].EnqueuedTimestamp); - return Math.Max(0L, Stopwatch.GetTimestamp() - oldest) - * 1000.0 - / Stopwatch.Frequency; - } - - private static long SaturatingAdd(long left, long right) => - left > long.MaxValue - right ? long.MaxValue : left + right; - /// /// Apply a single with the full side- /// effects: terrain upload, GPU state, and the re-hydration callback. - /// Extracted from the inline switch in the original DrainAndApply - /// so both the priority-hunt path and the normal drain path share it. + /// All priority queues route through this one publication path. /// private void ApplyResult(LandblockStreamResult result) { diff --git a/src/AcDream.App/Streaming/StreamingWorkBudget.cs b/src/AcDream.App/Streaming/StreamingWorkBudget.cs index 85cbba56..cfcd85eb 100644 --- a/src/AcDream.App/Streaming/StreamingWorkBudget.cs +++ b/src/AcDream.App/Streaming/StreamingWorkBudget.cs @@ -159,7 +159,13 @@ public readonly record struct StreamingWorkDiagnostics( long DeferredAdoptedCpuBytes, double OldestDeferredAgeMilliseconds, int PendingPublications, - int PendingRetirements); + int PendingRetirements, + int WorkerCompletionBacklog, + int DestinationBacklog, + int ControlBacklog, + int UnloadBacklog, + int NearBacklog, + int FarBacklog); /// /// Single-thread, frame-scoped admission meter. Callers reserve a known cost @@ -180,6 +186,7 @@ public sealed class StreamingWorkMeter private int _failures; private bool _frameOverrunRecorded; private bool _reservationActive; + private bool _ensuredProgressGranted; private string? _activeStage; private string? _lastStage; private StreamingWorkLimit _lastLimit; @@ -207,7 +214,8 @@ public sealed class StreamingWorkMeter public StreamingWorkAdmission TryReserve( StreamingWorkCost cost, - string stage) + string stage, + bool ensureProgress = false) { cost.Validate(); ArgumentException.ThrowIfNullOrWhiteSpace(stage); @@ -216,7 +224,11 @@ public sealed class StreamingWorkMeter "The prior streaming operation has not completed."); StreamingWorkLimit limit = FindLimit(cost); - if (limit != StreamingWorkLimit.None && _operations != 0) + bool grantEnsuredProgress = + ensureProgress && !_ensuredProgressGranted; + if (limit != StreamingWorkLimit.None + && _operations != 0 + && !grantEnsuredProgress) { _yields++; _lastStage = stage; @@ -228,14 +240,21 @@ public sealed class StreamingWorkMeter _operations++; _reservationActive = true; _activeStage = stage; - _lastStage = stage; - _lastLimit = limit; + if (ensureProgress) + _ensuredProgressGranted = true; if (limit != StreamingWorkLimit.None) { + _lastStage = stage; + _lastLimit = limit; _oversizedProgress++; return StreamingWorkAdmission.OversizedProgress; } + // Preserve the most recent constrained stage so later successful work + // in another budget dimension cannot erase the reason this frame + // yielded. With no constraint, retain the ordinary last-stage fact. + if (_lastLimit == StreamingWorkLimit.None) + _lastStage = stage; return StreamingWorkAdmission.Admitted; } diff --git a/src/AcDream.App/Streaming/StreamingWorkBudgetOptions.cs b/src/AcDream.App/Streaming/StreamingWorkBudgetOptions.cs index 894f6004..dcaef9f6 100644 --- a/src/AcDream.App/Streaming/StreamingWorkBudgetOptions.cs +++ b/src/AcDream.App/Streaming/StreamingWorkBudgetOptions.cs @@ -65,6 +65,31 @@ public sealed record StreamingWorkBudgetOptions( MaxGlRetireOperations, DestinationReserveFraction); + /// + /// Compatibility bridge for the pre-Slice-E quality setting. The former + /// landblock count now selects a complete work profile instead of acting + /// as a second hidden execution throttle. Four is the historical High + /// profile and therefore preserves this configured profile exactly. + /// + public StreamingWorkBudgetOptions ScaleForLegacyCompletionCount(int count) + { + if (count <= 0) + throw new ArgumentOutOfRangeException(nameof(count)); + + double scale = count / 4.0; + return this with + { + MaxUpdateMilliseconds = Math.Max( + 0.25, + MaxUpdateMilliseconds * scale), + MaxCompletionAdmissions = Scale(MaxCompletionAdmissions, scale), + MaxAdoptedCpuBytes = Scale(MaxAdoptedCpuBytes, scale), + MaxEntityOperations = Scale(MaxEntityOperations, scale), + MaxGpuUploadBytes = Scale(MaxGpuUploadBytes, scale), + MaxGlRetireOperations = Scale(MaxGlRetireOperations, scale), + }; + } + private static double ParsePositiveDouble(string? value, double fallback) => double.TryParse( value, @@ -118,4 +143,22 @@ public sealed record StreamingWorkBudgetOptions( return percent / 100f; } + + private static int Scale(int value, double scale) + { + if (scale >= int.MaxValue / (double)value) + return int.MaxValue; + return Math.Max( + 1, + (int)Math.Round(value * scale, MidpointRounding.AwayFromZero)); + } + + private static long Scale(long value, double scale) + { + if (scale >= long.MaxValue / (double)value) + return long.MaxValue; + return Math.Max( + 1L, + (long)Math.Round(value * scale, MidpointRounding.AwayFromZero)); + } } diff --git a/tests/AcDream.App.Tests/Streaming/LandblockPresentationPipelineTests.cs b/tests/AcDream.App.Tests/Streaming/LandblockPresentationPipelineTests.cs index a7871301..36c8ee85 100644 --- a/tests/AcDream.App.Tests/Streaming/LandblockPresentationPipelineTests.cs +++ b/tests/AcDream.App.Tests/Streaming/LandblockPresentationPipelineTests.cs @@ -1004,7 +1004,15 @@ public sealed class LandblockPresentationPipelineTests state, nearRadius: 0, farRadius: 2, - removeTerrain: id => calls.Add($"unload:{id:X8}")); + removeTerrain: id => calls.Add($"unload:{id:X8}"), + workBudgetOptions: new StreamingWorkBudgetOptions( + MaxUpdateMilliseconds: 100, + MaxCompletionAdmissions: 3, + MaxAdoptedCpuBytes: 1_000_000, + MaxEntityOperations: 1_000, + MaxGpuUploadBytes: 1_000_000, + MaxGlRetireOperations: 1_000, + DestinationReserveFraction: 0.75f)); controller.PriorityLandblockId = priorityId; Assert.Throws(() => controller.Tick(0x70, 0x70)); @@ -1021,8 +1029,8 @@ public sealed class LandblockPresentationPipelineTests Assert.Equal( [ $"publish:{priorityId:X8}", - $"publish:{laterId:X8}", $"unload:{unloadId:X8}", + $"publish:{laterId:X8}", ], calls); } diff --git a/tests/AcDream.App.Tests/Streaming/StreamingCompletionQueueTests.cs b/tests/AcDream.App.Tests/Streaming/StreamingCompletionQueueTests.cs new file mode 100644 index 00000000..7063d2d7 --- /dev/null +++ b/tests/AcDream.App.Tests/Streaming/StreamingCompletionQueueTests.cs @@ -0,0 +1,207 @@ +using System.Diagnostics; +using AcDream.App.Streaming; + +namespace AcDream.App.Tests.Streaming; + +public sealed class StreamingCompletionQueueTests +{ + [Fact] + public void PrioritySelectionIsStableFifoAndSkipsOnlyBlockedHead() + { + var queue = new StreamingCompletionQueue(); + StreamingQueuedCompletion farFirst = Completion( + new LandblockStreamResult.Failed(1, "far-first"), + StreamingCompletionPriority.Far, + sequence: 1); + StreamingQueuedCompletion nearBlocked = Completion( + new LandblockStreamResult.Failed(2, "near-blocked"), + StreamingCompletionPriority.Near, + sequence: 2); + StreamingQueuedCompletion nearTail = Completion( + new LandblockStreamResult.Failed(3, "near-tail"), + StreamingCompletionPriority.Near, + sequence: 3); + StreamingQueuedCompletion destinationFirst = Completion( + new LandblockStreamResult.Failed(4, "destination-first"), + StreamingCompletionPriority.Destination, + sequence: 4); + StreamingQueuedCompletion destinationTail = Completion( + new LandblockStreamResult.Failed(5, "destination-tail"), + StreamingCompletionPriority.Destination, + sequence: 5); + + queue.Enqueue(farFirst); + queue.Enqueue(nearBlocked); + queue.Enqueue(nearTail); + queue.Enqueue(destinationFirst); + queue.Enqueue(destinationTail); + + Assert.True(queue.TryPeekNext( + static _ => false, + out StreamingQueuedCompletion? selected)); + Assert.Equal(destinationFirst.Sequence, selected?.Sequence); + Assert.Same(destinationFirst.Result, selected?.Result); + queue.RemoveHead(destinationFirst); + + Assert.True(queue.TryPeekNext( + result => ReferenceEquals(result, destinationTail.Result) + || ReferenceEquals(result, nearBlocked.Result), + out selected)); + Assert.Equal(farFirst.Sequence, selected?.Sequence); + Assert.Same(farFirst.Result, selected?.Result); + + Assert.Throws(() => + queue.RemoveHead(nearTail)); + + Assert.True(queue.TryPeekNext( + static _ => false, + out selected)); + Assert.Equal(destinationTail.Sequence, selected?.Sequence); + Assert.Same(destinationTail.Result, selected?.Result); + queue.RemoveHead(destinationTail); + Assert.True(queue.TryPeekNext( + static _ => false, + out selected)); + Assert.Equal(nearBlocked.Sequence, selected?.Sequence); + Assert.Same(nearBlocked.Result, selected?.Result); + } + + [Fact] + public void RetainedBytesAndPriorityFactsTrackExactQueueLifetime() + { + var queue = new StreamingCompletionQueue(); + StreamingQueuedCompletion destination = Completion( + new LandblockStreamResult.Failed(10, "destination"), + StreamingCompletionPriority.Destination, + sequence: 1, + adoptedBytes: 12); + StreamingQueuedCompletion unload = Completion( + new LandblockStreamResult.Unloaded(11), + StreamingCompletionPriority.Unload, + sequence: 2, + adoptedBytes: 20); + StreamingQueuedCompletion near = Completion( + new LandblockStreamResult.Failed(12, "near"), + StreamingCompletionPriority.Near, + sequence: 3, + adoptedBytes: 30); + + queue.Enqueue(destination); + queue.Enqueue(unload); + queue.Enqueue(near); + + StreamingCompletionQueueSnapshot snapshot = queue.CaptureSnapshot(); + Assert.Equal(3, snapshot.Count); + Assert.Equal(62, snapshot.RetainedCpuBytes); + Assert.Equal(1, snapshot.Destination); + Assert.Equal(1, snapshot.Unload); + Assert.Equal(1, snapshot.Near); + Assert.True(snapshot.OldestAgeMilliseconds >= 0); + + queue.RemoveHead(destination); + snapshot = queue.CaptureSnapshot(); + Assert.Equal(2, snapshot.Count); + Assert.Equal(50, snapshot.RetainedCpuBytes); + + queue.Clear(); + snapshot = queue.CaptureSnapshot(); + Assert.Equal(0, snapshot.Count); + Assert.Equal(0, snapshot.RetainedCpuBytes); + Assert.Equal(0, snapshot.OldestAgeMilliseconds); + } + + [Fact] + public void ExactReferenceRemovalDoesNotConflateEqualRecords() + { + var queue = new StreamingCompletionQueue(); + var firstResult = + new LandblockStreamResult.Failed(20, "same", Generation: 7); + var equalButDistinctResult = + new LandblockStreamResult.Failed(20, "same", Generation: 7); + StreamingQueuedCompletion first = Completion( + firstResult, + StreamingCompletionPriority.Control, + sequence: 1, + adoptedBytes: 11); + StreamingQueuedCompletion second = Completion( + equalButDistinctResult, + StreamingCompletionPriority.Control, + sequence: 2, + adoptedBytes: 13); + queue.Enqueue(first); + queue.Enqueue(second); + + int removed = queue.RemoveResults( + [equalButDistinctResult], + static _ => true); + + Assert.Equal(1, removed); + Assert.Equal(1, queue.Count); + Assert.Equal(11, queue.RetainedCpuBytes); + Assert.True(queue.TryPeekNext( + static _ => false, + out StreamingQueuedCompletion? selected)); + Assert.Equal(first.Sequence, selected?.Sequence); + Assert.Same(first.Result, selected?.Result); + } + + [Fact] + public void CallbackAppendNeverInvalidatesActiveHead() + { + var queue = new StreamingCompletionQueue(); + StreamingQueuedCompletion active = Completion( + new LandblockStreamResult.Failed(30, "active"), + StreamingCompletionPriority.Near, + sequence: 1); + StreamingQueuedCompletion appended = Completion( + new LandblockStreamResult.Failed(31, "appended"), + StreamingCompletionPriority.Near, + sequence: 2); + queue.Enqueue(active); + + Assert.True(queue.TryPeekNext( + static _ => false, + out StreamingQueuedCompletion? selected)); + Assert.Equal(active.Sequence, selected?.Sequence); + Assert.Same(active.Result, selected?.Result); + + queue.Enqueue(appended); + queue.RemoveHead(active); + + Assert.True(queue.TryPeekNext( + static _ => false, + out selected)); + Assert.Equal(appended.Sequence, selected?.Sequence); + Assert.Same(appended.Result, selected?.Result); + } + + private static StreamingQueuedCompletion Completion( + LandblockStreamResult result, + StreamingCompletionPriority priority, + long sequence, + long adoptedBytes = 0) + { + var estimate = new LandblockStreamCostEstimate( + Work: new StreamingWorkCost( + CompletionAdmissions: 1, + AdoptedCpuBytes: adoptedBytes), + TerrainPayloadBytes: 0, + Entities: 0, + MeshReferences: 0, + VisibilityCells: 0, + EnvCellShells: 0, + PortalConnections: 0, + PortalPolygonVertices: 0, + PhysicsEnvCells: 0, + PhysicsEnvironments: 0, + PhysicsSetups: 0, + PhysicsGfxObjects: 0); + return new StreamingQueuedCompletion( + result, + estimate, + priority, + result.Generation, + sequence, + Stopwatch.GetTimestamp()); + } +} diff --git a/tests/AcDream.App.Tests/Streaming/StreamingWorkBudgetTests.cs b/tests/AcDream.App.Tests/Streaming/StreamingWorkBudgetTests.cs index 716f22fe..00128aae 100644 --- a/tests/AcDream.App.Tests/Streaming/StreamingWorkBudgetTests.cs +++ b/tests/AcDream.App.Tests/Streaming/StreamingWorkBudgetTests.cs @@ -83,6 +83,40 @@ public sealed class StreamingWorkBudgetTests snapshot.LastLimit); } + [Fact] + public void EnsureProgressCanCrossTheLimitOnlyOncePerFrame() + { + var meter = new StreamingWorkMeter( + Budget(cpuBytes: 1), + static () => 0, + timestampFrequency: 1_000); + + Assert.Equal( + StreamingWorkAdmission.Admitted, + meter.TryReserve( + new StreamingWorkCost(CompletionAdmissions: 1), + "admission")); + meter.Complete(); + Assert.Equal( + StreamingWorkAdmission.Admitted, + meter.TryReserve( + new StreamingWorkCost(AdoptedCpuBytes: 1), + "first-execution", + ensureProgress: true)); + meter.Complete(); + + Assert.Equal( + StreamingWorkAdmission.Yielded, + meter.TryReserve( + new StreamingWorkCost(AdoptedCpuBytes: 1), + "second-execution", + ensureProgress: true)); + Assert.Equal(1, meter.Snapshot.YieldCount); + Assert.Equal( + StreamingWorkLimit.AdoptedCpuBytes, + meter.Snapshot.LastLimit); + } + [Fact] public void MeterUsesInjectedMonotonicClockAndRecordsOverrunAndFailure() { @@ -129,6 +163,64 @@ public sealed class StreamingWorkBudgetTests timestampFrequency: 1_000)); } + [Theory] + [InlineData(2, 1.0, 2, 50, 4, 50, 2)] + [InlineData(4, 2.0, 4, 100, 8, 100, 4)] + [InlineData(6, 3.0, 6, 150, 12, 150, 6)] + public void LegacyCompletionSelectorScalesTheWholeProfile( + int count, + double milliseconds, + int admissions, + long cpuBytes, + int entityOperations, + long gpuBytes, + int retireOperations) + { + var options = new StreamingWorkBudgetOptions( + MaxUpdateMilliseconds: 2, + MaxCompletionAdmissions: 4, + MaxAdoptedCpuBytes: 100, + MaxEntityOperations: 8, + MaxGpuUploadBytes: 100, + MaxGlRetireOperations: 4, + DestinationReserveFraction: 0.75f); + + StreamingWorkBudgetOptions scaled = + options.ScaleForLegacyCompletionCount(count); + + Assert.Equal(milliseconds, scaled.MaxUpdateMilliseconds); + Assert.Equal(admissions, scaled.MaxCompletionAdmissions); + Assert.Equal(cpuBytes, scaled.MaxAdoptedCpuBytes); + Assert.Equal(entityOperations, scaled.MaxEntityOperations); + Assert.Equal(gpuBytes, scaled.MaxGpuUploadBytes); + Assert.Equal(retireOperations, scaled.MaxGlRetireOperations); + Assert.Equal(0.75f, scaled.DestinationReserveFraction); + } + + [Fact] + public void LegacyCompletionSelectorRejectsZeroAndSaturatesHugeProfiles() + { + var options = new StreamingWorkBudgetOptions( + MaxUpdateMilliseconds: 2, + MaxCompletionAdmissions: int.MaxValue, + MaxAdoptedCpuBytes: long.MaxValue, + MaxEntityOperations: int.MaxValue, + MaxGpuUploadBytes: long.MaxValue, + MaxGlRetireOperations: int.MaxValue, + DestinationReserveFraction: 0.75f); + + Assert.Throws(() => + options.ScaleForLegacyCompletionCount(0)); + + StreamingWorkBudgetOptions scaled = + options.ScaleForLegacyCompletionCount(int.MaxValue); + Assert.Equal(int.MaxValue, scaled.MaxCompletionAdmissions); + Assert.Equal(long.MaxValue, scaled.MaxAdoptedCpuBytes); + Assert.Equal(int.MaxValue, scaled.MaxEntityOperations); + Assert.Equal(long.MaxValue, scaled.MaxGpuUploadBytes); + Assert.Equal(int.MaxValue, scaled.MaxGlRetireOperations); + } + [Fact] public void CompletionCostChargesExactArraysAndDeterministicLogicalEntries() { @@ -239,10 +331,15 @@ public sealed class StreamingWorkBudgetTests applyTerrain: static (_, _) => { }, state: new GpuWorldState(), nearRadius: 1, - farRadius: 1) - { - MaxCompletionsPerFrame = 1, - }; + farRadius: 1, + workBudgetOptions: new StreamingWorkBudgetOptions( + MaxUpdateMilliseconds: 100, + MaxCompletionAdmissions: 2, + MaxAdoptedCpuBytes: 1_000, + MaxEntityOperations: 100, + MaxGpuUploadBytes: 44, + MaxGlRetireOperations: 100, + DestinationReserveFraction: 0.75f)); controller.Tick(32, 32); @@ -250,9 +347,10 @@ public sealed class StreamingWorkBudgetTests Assert.Equal(1, first.DeferredCompletions); Assert.Equal(44, first.DeferredAdoptedCpuBytes); Assert.True(first.OldestDeferredAgeMilliseconds >= 0); - Assert.Equal(1, first.LastFrame.Operations); - Assert.Equal(1, first.LastFrame.CompletedOperations); - Assert.Equal(44, first.LastFrame.Used.AdoptedCpuBytes); + Assert.Equal(3, first.LastFrame.Operations); + Assert.Equal(3, first.LastFrame.CompletedOperations); + Assert.Equal(1, first.LastFrame.YieldCount); + Assert.Equal(88, first.LastFrame.Used.AdoptedCpuBytes); controller.Tick(32, 32); @@ -263,6 +361,155 @@ public sealed class StreamingWorkBudgetTests Assert.Equal(1, second.LastFrame.CompletedOperations); } + [Fact] + public void ControllerBoundsProductionAdmissionByCountBeforeAdoption() + { + uint center = StreamingRegion.EncodeLandblockIdForTest(32, 32); + var source = new QueueCompletionSource( + Loaded(center), + Loaded(StreamingRegion.EncodeLandblockIdForTest(32, 33)), + Loaded(StreamingRegion.EncodeLandblockIdForTest(32, 34))); + int applied = 0; + StreamingController controller = Controller( + source, + () => applied++, + WorkOptions(admissions: 2, cpuBytes: 1_000)); + + controller.Tick(32, 32); + + Assert.Equal(2, applied); + Assert.Equal(1, source.BacklogCount); + Assert.Equal(1, controller.WorkDiagnostics.WorkerCompletionBacklog); + Assert.Equal(2, controller.WorkDiagnostics.LastFrame.Used.CompletionAdmissions); + Assert.Equal(0, controller.WorkDiagnostics.DeferredCompletions); + } + + [Fact] + public void ControllerBoundsProductionAdmissionByRetainedCpuBytes() + { + uint center = StreamingRegion.EncodeLandblockIdForTest(32, 32); + var source = new QueueCompletionSource( + Loaded(center), + Loaded(StreamingRegion.EncodeLandblockIdForTest(32, 33)), + Loaded(StreamingRegion.EncodeLandblockIdForTest(32, 34))); + int applied = 0; + StreamingController controller = Controller( + source, + () => applied++, + WorkOptions(admissions: 64, cpuBytes: 44)); + + controller.Tick(32, 32); + + Assert.Equal(1, applied); + Assert.Equal(2, source.BacklogCount); + Assert.Equal(44, controller.WorkDiagnostics.LastFrame.Used.AdoptedCpuBytes); + Assert.Equal( + StreamingWorkLimit.AdoptedCpuBytes, + controller.WorkDiagnostics.LastFrame.LastLimit); + } + + [Fact] + public void DestinationAndUnloadPriorityNeverBypassExecutionBudgets() + { + uint center = StreamingRegion.EncodeLandblockIdForTest(32, 32); + uint ordinary = StreamingRegion.EncodeLandblockIdForTest(33, 33); + var source = new QueueCompletionSource( + Loaded(ordinary), + Loaded(center), + new LandblockStreamResult.Unloaded(0x1111FFFFu), + new LandblockStreamResult.Unloaded(0x2222FFFFu)); + var applied = new List(); + StreamingController controller = Controller( + source, + id => applied.Add(id), + WorkOptions( + admissions: 4, + cpuBytes: 1_000, + gpuBytes: 44, + retireOperations: 1)); + controller.PriorityLandblockId = center; + + controller.Tick(32, 32); + + Assert.Equal([center], applied); + Assert.Equal(2, controller.WorkDiagnostics.DeferredCompletions); + Assert.Equal(1, controller.WorkDiagnostics.LastFrame.YieldCount); + Assert.Equal(44, controller.WorkDiagnostics.LastFrame.Used.GpuUploadBytes); + Assert.Equal(1, controller.WorkDiagnostics.LastFrame.Used.GlRetireOperations); + + controller.Tick(32, 32); + + Assert.Equal([center, ordinary], applied); + Assert.Equal(0, controller.WorkDiagnostics.DeferredCompletions); + Assert.Equal(1, controller.WorkDiagnostics.LastFrame.Used.GlRetireOperations); + } + + [Fact] + public void StaleGenerationResultsConsumeAdmissionButRetainNoPayload() + { + uint center = StreamingRegion.EncodeLandblockIdForTest(32, 32); + var source = new QueueCompletionSource( + Loaded(center, generation: 1), + Loaded(StreamingRegion.EncodeLandblockIdForTest(32, 33), generation: 1), + Loaded(center)); + int applied = 0; + StreamingController controller = Controller( + source, + () => applied++, + WorkOptions(admissions: 2, cpuBytes: 44)); + + controller.Tick(32, 32); + + Assert.Equal(0, applied); + Assert.Equal(1, source.BacklogCount); + Assert.Equal(0, controller.WorkDiagnostics.DeferredCompletions); + Assert.Equal(0, controller.WorkDiagnostics.DeferredAdoptedCpuBytes); + Assert.Equal(2, controller.WorkDiagnostics.LastFrame.Used.CompletionAdmissions); + Assert.Equal(0, controller.WorkDiagnostics.LastFrame.Used.AdoptedCpuBytes); + + controller.Tick(32, 32); + Assert.Equal(1, applied); + } + + private static StreamingController Controller( + ILandblockCompletionSource source, + Action applied, + StreamingWorkBudgetOptions options) => + Controller(source, _ => applied(), options); + + private static StreamingController Controller( + ILandblockCompletionSource source, + Action applied, + StreamingWorkBudgetOptions options) + { + var state = new GpuWorldState(); + var presentation = new LandblockPresentationPipeline( + (build, _) => applied(build.Landblock.LandblockId), + state); + return new StreamingController( + enqueueLoad: static (_, _, _) => { }, + enqueueUnload: static (_, _) => { }, + completionSource: source, + state, + nearRadius: 1, + farRadius: 1, + presentationPipeline: presentation, + workBudgetOptions: options); + } + + private static StreamingWorkBudgetOptions WorkOptions( + int admissions, + long cpuBytes, + long gpuBytes = 1_000, + int retireOperations = 100) => new( + MaxUpdateMilliseconds: 100, + MaxCompletionAdmissions: admissions, + MaxAdoptedCpuBytes: cpuBytes, + MaxEntityOperations: 100, + MaxGpuUploadBytes: gpuBytes, + MaxGlRetireOperations: retireOperations, + DestinationReserveFraction: 0.75f); + private static WorldEntity Entity(uint id, int meshRefs) { var refs = new MeshRef[meshRefs]; @@ -278,7 +525,9 @@ public sealed class StreamingWorkBudgetTests }; } - private static LandblockStreamResult.Loaded Loaded(uint landblockId) => new( + private static LandblockStreamResult.Loaded Loaded( + uint landblockId, + ulong generation = 0) => new( landblockId, LandblockStreamTier.Near, new LoadedLandblock( @@ -287,5 +536,39 @@ public sealed class StreamingWorkBudgetTests Array.Empty()), new LandblockMeshData( new TerrainVertex[1], - new uint[1])); + new uint[1]), + generation); + + private sealed class QueueCompletionSource( + params LandblockStreamResult[] results) + : ILandblockCompletionSource + { + private readonly Queue _results = new(results); + + public int BacklogCount => _results.Count; + + public bool TryPeek(out LandblockStreamResult? result) + { + if (_results.TryPeek(out LandblockStreamResult? peeked)) + { + result = peeked; + return true; + } + + result = null; + return false; + } + + public bool TryRead(out LandblockStreamResult? result) + { + if (_results.TryDequeue(out LandblockStreamResult? consumed)) + { + result = consumed; + return true; + } + + result = null; + return false; + } + } } diff --git a/tests/AcDream.Core.Tests/Streaming/LandblockStreamerTests.cs b/tests/AcDream.Core.Tests/Streaming/LandblockStreamerTests.cs index 454ce86a..122988a9 100644 --- a/tests/AcDream.Core.Tests/Streaming/LandblockStreamerTests.cs +++ b/tests/AcDream.Core.Tests/Streaming/LandblockStreamerTests.cs @@ -263,6 +263,34 @@ public class LandblockStreamerTests Assert.Equal(43ul, unloaded.Generation); } + [Fact] + public async Task CompletionSourcePeekPreservesExactResultAndBacklog() + { + using var streamer = new LandblockStreamer(loadLandblock: _ => null); + streamer.Start(); + streamer.EnqueueUnload(0xABCE0000u, generation: 44); + + for (int i = 0; + i < SpinMaxIterations && streamer.BacklogCount == 0; + i++) + { + await Task.Delay(SpinStepMs); + } + + Assert.Equal(1, streamer.BacklogCount); + Assert.True(streamer.TryPeek(out LandblockStreamResult? firstPeek)); + Assert.NotNull(firstPeek); + Assert.Equal(1, streamer.BacklogCount); + Assert.True(streamer.TryPeek(out LandblockStreamResult? secondPeek)); + Assert.Same(firstPeek, secondPeek); + Assert.Equal(1, streamer.BacklogCount); + + Assert.True(streamer.TryRead(out LandblockStreamResult? consumed)); + Assert.Same(firstPeek, consumed); + Assert.Equal(0, streamer.BacklogCount); + Assert.False(streamer.TryRead(out _)); + } + [Fact] public async Task Load_ExecutesLoaderOnWorkerThread() { diff --git a/tests/AcDream.Core.Tests/Streaming/StreamingControllerPriorityApplyTests.cs b/tests/AcDream.Core.Tests/Streaming/StreamingControllerPriorityApplyTests.cs index 20ff3a99..2a183ecc 100644 --- a/tests/AcDream.Core.Tests/Streaming/StreamingControllerPriorityApplyTests.cs +++ b/tests/AcDream.Core.Tests/Streaming/StreamingControllerPriorityApplyTests.cs @@ -18,7 +18,7 @@ public class StreamingControllerPriorityApplyTests generation); [Fact] - public void PriorityLandblock_isApplied_evenWhenBeyondPerFrameCap() + public void PriorityLandblock_WinsExecutionOrderAfterBoundedAdmission() { uint priority = StreamingRegion.EncodeLandblockIdForTest(169, 180); var outbox = new Queue(new LandblockStreamResult[] @@ -42,14 +42,18 @@ public class StreamingControllerPriorityApplyTests return batch; }, applyTerrain: (lb, _) => applied.Add(lb.LandblockId), - state: state, nearRadius: 4, farRadius: 12) + state: state, + nearRadius: 4, + farRadius: 12, + workBudgetOptions: GenerousBudget()) { MaxCompletionsPerFrame = 4 }; ctrl.PriorityLandblockId = priority; ctrl.Tick(169, 180); - Assert.Contains(priority, applied); // priority applied THIS tick - Assert.True(applied.Count <= 5); // did not blindly flush the whole outbox + Assert.NotEmpty(applied); + Assert.Equal(priority, applied[0]); + Assert.True(applied.Count <= 5); } [Fact] @@ -98,10 +102,8 @@ public class StreamingControllerPriorityApplyTests public void PriorityNeverArrives_noThrow_noLoss_noDoubleApply() { // While a PriorityLandblockId is set but the matching completion never - // arrives (e.g. the load failed), the hunt moves outbox completions into - // _deferredApply and they drain at the per-frame budget — same backpressure - // as the normal throttle, just relocated — until the caller clears - // PriorityLandblockId (the TAS MaxContinue safety net does this on timeout). + // arrives (e.g. the load failed), ordinary results still enter the typed + // Near/Far queues and advance under the same budget without loss. uint priority = StreamingRegion.EncodeLandblockIdForTest(169, 180); uint otherId0 = StreamingRegion.EncodeLandblockIdForTest(169, 181); uint otherId1 = StreamingRegion.EncodeLandblockIdForTest(169, 182); @@ -149,6 +151,15 @@ public class StreamingControllerPriorityApplyTests Assert.Equal(3, applied.Count); } + private static StreamingWorkBudgetOptions GenerousBudget() => new( + MaxUpdateMilliseconds: 100, + MaxCompletionAdmissions: 64, + MaxAdoptedCpuBytes: 16 * StreamingWorkBudgetOptions.MiB, + MaxEntityOperations: 10_000, + MaxGpuUploadBytes: 16 * StreamingWorkBudgetOptions.MiB, + MaxGlRetireOperations: 10_000, + DestinationReserveFraction: 0.75f); + [Fact] public void DeferredCompaction_ApplyFailureRetainsExactResultWithoutReplayingCommittedPrefix() { @@ -223,7 +234,15 @@ public class StreamingControllerPriorityApplyTests applyTerrain: (build, _) => applied.Add(build.LandblockId), state: state, nearRadius: 4, - farRadius: 12) + farRadius: 12, + workBudgetOptions: new StreamingWorkBudgetOptions( + MaxUpdateMilliseconds: 100, + MaxCompletionAdmissions: 4, + MaxAdoptedCpuBytes: 1_000_000, + MaxEntityOperations: 1_000, + MaxGpuUploadBytes: 1_000_000, + MaxGlRetireOperations: 1_000, + DestinationReserveFraction: 0.75f)) { MaxCompletionsPerFrame = 4, PriorityLandblockId = priority }; ctrl.Tick(169, 180);