using System;
using System.Collections.Generic;
using System.Diagnostics;
using AcDream.App.Rendering.Wb;
using AcDream.Core.Terrain;
using AcDream.Core.World;
namespace AcDream.App.Streaming;
///
/// Called once per frame from GameWindow.OnUpdate. Owns the
/// and uses delegates into
/// so tests can inject fakes. All work
/// happens on the render thread; the streamer itself is background.
///
///
/// Threading: not thread-safe. All calls must happen on the render thread.
///
///
public sealed class StreamingController
: IStreamingFrameBackend,
IWorldRevealStreamingScheduler
{
private readonly record struct DestinationReservation(
long RevealGeneration,
uint LandblockId,
int Radius);
private sealed class OriginRecenterRetirement
{
public bool RadiiConverged;
public bool GenerationAdvanced;
public bool PendingLoadsCleared;
public bool CompletionQueueCleared;
public bool RegionCleared;
public bool SpatialGenerationDetached;
public bool PreparationCommitted;
public (int X, int Y, bool IsSealedDungeon)? Destination;
public bool DestinationConfigured;
public bool DestinationLoadEnqueued;
}
private sealed class FullWindowRetirement
{
public bool GenerationAdvanced;
public bool PendingLoadsCleared;
public bool CompletionQueueCleared;
public bool RegionCleared;
public List? ResidentIds;
public IEnumerator? ResidentEnumerator;
public int RetirementCursor;
public bool PreparationCommitted;
}
private readonly Action _enqueueLoad;
private readonly Action _enqueueUnload;
private readonly ILandblockCompletionSource _completionSource;
private readonly Action? _clearPendingLoads;
private readonly LandblockPresentationPipeline _presentation;
private readonly Func
_isPublicationBlockedByRetirement;
private readonly StreamingWorkBudgetOptions _configuredWorkBudgetOptions;
private StreamingWorkBudget _workBudget;
private StreamingWorkMeter? _activeWorkMeter;
private StreamingWorkMeterSnapshot _lastWorkMeter;
private readonly StreamingCompletionQueue _completionQueue = new();
private long _nextCompletionSequence;
private int _legacyCompletionProfile = 4;
private DestinationReservation? _destinationReservation;
private long _lifetimeWorkOverruns;
private long _lifetimeOversizedProgress;
private double _maximumWorkFrameMilliseconds;
private string? _maximumWorkFrameStage;
private double _maximumWorkOperationMilliseconds;
private string? _maximumWorkOperationStage;
private readonly GpuWorldState _state;
private StreamingRegion? _region;
private RadiiReconfiguration? _pendingRadiiReconfiguration;
private bool _advancingRadiiReconfiguration;
private (int NearRadius, int FarRadius)? _deferredRadiiRequest;
private OriginRecenterRetirement? _originRecenterRetirement;
private FullWindowRetirement? _fullWindowRetirement;
private bool _advancingOriginRecenter;
// Hard streaming boundaries advance this token. Worker completions carry
// the generation captured at enqueue time, so an old overlapping load or
// unload cannot mutate the replacement window after a portal recenter.
private ulong _generation;
// True while streaming is collapsed to the single dungeon landblock the
// player stands in (the dungeon gate, #133 FPS). AC dungeons have NO
// adjacent landblocks — neighbors are unrelated ocean-grid dungeons that
// are never visible, so we stop loading the 25×25 window entirely.
private bool _collapsed;
// The dungeon landblock id we collapsed onto. Once collapsed we key the
// gate on this STABLE landblock, not the per-frame insideDungeon signal:
// CurrCell can momentarily resolve to null/outdoor mid-frame, and gating
// expand on that flicker thrashes collapse↔expand (reload storms + a light
// leak). We only expand when the observer actually moves to a different
// landblock (teleport/portal out).
private uint _collapsedCenter;
///
/// Near-tier radius (LBs from observer that load full detail: terrain +
/// scenery + entities). Runtime changes go through
/// .
///
public int NearRadius { get; private set; }
///
/// Far-tier radius (LBs from observer that load terrain only).
///
public int FarRadius { get; private set; }
///
/// 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 => _legacyCompletionProfile;
set
{
StreamingWorkBudgetOptions profile =
_configuredWorkBudgetOptions
.ScaleForLegacyCompletionCount(value);
_legacyCompletionProfile = value;
_workBudget = profile.ToBudget();
}
}
internal long ActiveRevealGeneration =>
_destinationReservation?.RevealGeneration ?? 0L;
internal uint DestinationLandblockId =>
_destinationReservation?.LandblockId ?? 0u;
internal int DestinationRadius =>
_destinationReservation?.Radius ?? 0;
internal void BeginDestinationReservation(
long revealGeneration,
uint destinationCell,
int requiredRenderRadius)
{
if (revealGeneration <= 0)
throw new ArgumentOutOfRangeException(nameof(revealGeneration));
if (destinationCell == 0u)
throw new ArgumentOutOfRangeException(nameof(destinationCell));
if (requiredRenderRadius < 0)
throw new ArgumentOutOfRangeException(nameof(requiredRenderRadius));
_destinationReservation = new DestinationReservation(
revealGeneration,
(destinationCell & 0xFFFF0000u) | 0xFFFFu,
requiredRenderRadius);
}
internal void EndDestinationReservation(long revealGeneration)
{
if (_destinationReservation is { } active
&& active.RevealGeneration == revealGeneration)
{
_destinationReservation = null;
}
}
void IWorldRevealStreamingScheduler.BeginDestinationReservation(
long revealGeneration,
uint destinationCell,
int requiredRenderRadius) =>
BeginDestinationReservation(
revealGeneration,
destinationCell,
requiredRenderRadius);
void IWorldRevealStreamingScheduler.EndDestinationReservation(
long revealGeneration) =>
EndDestinationReservation(revealGeneration);
// [FRAME-DIAG] (read by GameWindow's ACDREAM_WB_DIAG rollup): the standing
// 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
{
get
{
StreamingCompletionQueueSnapshot queued =
_completionQueue.CaptureSnapshot();
return new StreamingWorkDiagnostics(
_lastWorkMeter,
_lifetimeWorkOverruns,
_lifetimeOversizedProgress,
_maximumWorkFrameMilliseconds,
_maximumWorkFrameStage,
_maximumWorkOperationMilliseconds,
_maximumWorkOperationStage,
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;
///
/// True once every in-bounds landblock in the requested Chebyshev ring has
/// crossed the render-thread publication barrier. Worker completion and
/// world-state registration are not sufficient: all static GfxObj and
/// EnvCell shell meshes must have completed their render-thread upload.
/// Portal-space exit uses this alongside physics residency so the world
/// cannot be revealed while its render slots are still absent.
///
public bool IsRenderNeighborhoodResident(uint cellOrLandblockId, int radius)
{
if (radius < 0)
throw new ArgumentOutOfRangeException(nameof(radius));
// LandDefs::InboundValidCellId validates both map axes and the low-word
// class (outdoor cell, EnvCell, or canonical landblock sentinel).
if (!AcDream.Core.Physics.LandDefs.InboundValidCellId(cellOrLandblockId))
return false;
int cx = (int)((cellOrLandblockId >> 24) & 0xFFu);
int cy = (int)((cellOrLandblockId >> 16) & 0xFFu);
for (int dx = -radius; dx <= radius; dx++)
for (int dy = -radius; dy <= radius; dy++)
{
int nx = cx + dx;
int ny = cy + dy;
// Match PhysicsEngine.IsNeighborhoodTerrainResident: the outer
// 0xFF coordinate has no loadable neighbour beyond it.
if (nx < 0 || nx > 254 || ny < 0 || ny > 254)
continue;
uint canonical = ((uint)nx << 24) | ((uint)ny << 16) | 0xFFFFu;
if (!_state.IsNearTier(canonical) || !_state.IsRenderReady(canonical))
return false;
}
return true;
}
// [FRAME-DIAG]: full-window presentation retirements include explicit
// reloads and shared-origin recenter barriers. Both retire and later
// re-upload the complete resident window.
public int FullWindowRetirementCount { get; private set; }
public int LastFullWindowRetirementLandblockCount { get; private set; }
///
/// Internal compatibility seam for hermetic controller policy tests. Live
/// composition cannot supply presentation callbacks; it must use the
/// concrete pipeline constructor below.
///
internal StreamingController(
Action enqueueLoad,
Action enqueueUnload,
Func> drainCompletions,
Action applyTerrain,
GpuWorldState state,
int nearRadius,
int farRadius,
Action? removeTerrain = null,
Action? demoteNearLayer = null,
Action? clearPendingLoads = null,
Action? onLandblockLoaded = null,
Action? ensureEnvCellMeshes = null,
LandblockRetirementCoordinator? retirementCoordinator = null,
StreamingWorkBudgetOptions? workBudgetOptions = null)
: this(
enqueueLoad,
enqueueUnload,
drainCompletions,
state,
nearRadius,
farRadius,
new LandblockPresentationPipeline(
applyTerrain,
state,
onLandblockLoaded,
ensureEnvCellMeshes,
retirementCoordinator,
removeTerrain,
demoteNearLayer),
clearPendingLoads,
workBudgetOptions)
{
}
///
/// Concrete presentation path. Its pipeline is the sole presentation and
/// retirement owner, so no legacy publication delegate can be supplied or
/// silently ignored.
///
internal StreamingController(
Action enqueueLoad,
Action enqueueUnload,
Func> drainCompletions,
GpuWorldState state,
int nearRadius,
int farRadius,
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;
_completionSource = completionSource
?? throw new ArgumentNullException(nameof(completionSource));
_clearPendingLoads = clearPendingLoads;
_state = state;
_presentation = presentationPipeline
?? throw new ArgumentNullException(nameof(presentationPipeline));
_isPublicationBlockedByRetirement =
IsPublicationBlockedByRetirement;
_configuredWorkBudgetOptions =
workBudgetOptions ?? StreamingWorkBudgetOptions.Default;
_workBudget = _configuredWorkBudgetOptions.ToBudget();
if (!_presentation.MatchesState(_state))
{
throw new ArgumentException(
"The presentation pipeline must own the controller's world state.",
nameof(presentationPipeline));
}
NearRadius = nearRadius;
FarRadius = farRadius;
}
///
/// Reconciles an active outdoor streaming window to new quality radii
/// without replacing the worker/controller or blindly bootstrapping data
/// already published in . Pending jobs from the
/// old window are invalidated by generation; resident blocks are retained,
/// promoted, demoted, loaded, or unloaded from their actual current tier.
/// A sealed dungeon remains radius-zero until its normal exit edge while
/// remembering the new outdoor radii for that exit.
///
public void ReconfigureRadii(int nearRadius, int farRadius)
{
if (nearRadius < 0)
throw new ArgumentOutOfRangeException(nameof(nearRadius));
if (farRadius < nearRadius)
throw new ArgumentOutOfRangeException(
nameof(farRadius),
"Far radius must be greater than or equal to near radius.");
// A radius transaction may not admit work after an origin-recenter
// snapshot has been captured. Retain only the latest requested radii;
// the first Tick after origin commit applies them to the new frame.
if (_originRecenterRetirement is not null)
{
_deferredRadiiRequest = (nearRadius, farRadius);
return;
}
// Queue callbacks are injected seams and can synchronously reenter the
// controller. Last-request-wins deferral keeps the active mutation
// cursor stable; the request is applied after the current transaction
// converges instead of recursively replaying its active step.
if (_advancingRadiiReconfiguration)
{
_deferredRadiiRequest = (nearRadius, farRadius);
return;
}
// A prior callback may have failed after this controller accepted the
// request. Resume that exact transaction before considering another
// target; replacing it would orphan already-admitted generation work.
if (_pendingRadiiReconfiguration is not null)
AdvanceRadiiReconfiguration();
// This explicit call is newer than any request deferred by a prior
// failed attempt and therefore supersedes it.
_deferredRadiiRequest = null;
if (nearRadius == NearRadius && farRadius == FarRadius)
return;
// A staged old-window publication must either finish or fail before
// the desired-tier snapshot below is computed. Otherwise a landblock
// that becomes spatially resident during convergence is absent from
// the reconfiguration mutation ledger.
if (!ConvergePendingPublications())
{
_deferredRadiiRequest = (nearRadius, farRadius);
return;
}
if (_collapsed || _region is null)
{
NearRadius = nearRadius;
FarRadius = farRadius;
return;
}
var rebuilt = new StreamingRegion(
_region.CenterX,
_region.CenterY,
nearRadius,
farRadius);
TwoTierDiff bootstrap = rebuilt.ComputeFirstTickDiff();
var desiredNear = new HashSet(bootstrap.ToLoadNear);
var desiredFar = new HashSet(bootstrap.ToLoadFar);
var mutations = new List();
uint[] loaded = [.. _state.LoadedLandblockIds];
for (int i = 0; i < loaded.Length; i++)
{
uint id = loaded[i];
if (desiredNear.Contains(id))
{
if (!_state.IsNearTier(id))
mutations.Add(new RadiiMutation(
() => EnqueueLoad(id, LandblockStreamJobKind.PromoteToNear)));
}
else if (desiredFar.Contains(id))
{
if (_state.IsNearTier(id))
mutations.Add(new RadiiMutation(() => DemoteLandblock(id)));
}
else
{
mutations.Add(new RadiiMutation(() => EnqueueUnload(id)));
}
}
foreach (uint id in desiredNear)
if (!_state.IsLoaded(id))
mutations.Add(new RadiiMutation(
() => EnqueueLoad(id, LandblockStreamJobKind.LoadNear)));
foreach (uint id in desiredFar)
if (!_state.IsLoaded(id))
mutations.Add(new RadiiMutation(
() => EnqueueLoad(id, LandblockStreamJobKind.LoadFar)));
_pendingRadiiReconfiguration = new RadiiReconfiguration(
nearRadius,
farRadius,
rebuilt,
mutations);
AdvanceRadiiReconfiguration();
}
///
/// Compatibility constructor for deterministic tests and callers that do
/// not own an asynchronous worker. Production must use the generation-aware
/// overload above so hard-recenter results can be rejected by incarnation.
///
internal StreamingController(
Action enqueueLoad,
Action enqueueUnload,
Func> drainCompletions,
Action applyTerrain,
GpuWorldState state,
int nearRadius,
int farRadius,
Action? removeTerrain = null,
Action? demoteNearLayer = null,
Action? clearPendingLoads = null,
Action? onLandblockLoaded = null,
Action? ensureEnvCellMeshes = null,
LandblockRetirementCoordinator? retirementCoordinator = null,
StreamingWorkBudgetOptions? workBudgetOptions = null)
: this(
(id, kind, _) => enqueueLoad(id, kind),
(id, _) => enqueueUnload(id),
drainCompletions,
applyTerrain,
state,
nearRadius,
farRadius,
removeTerrain,
demoteNearLayer,
clearPendingLoads,
onLandblockLoaded,
ensureEnvCellMeshes,
retirementCoordinator,
workBudgetOptions)
{
}
private void EnqueueLoad(uint id, LandblockStreamJobKind kind) =>
_enqueueLoad(id, kind, _generation);
private void EnqueueUnload(uint id) => _enqueueUnload(id, _generation);
private void DemoteLandblock(uint id)
{
uint canonical = (id & 0xFFFF0000u) | 0xFFFFu;
_presentation.EnqueueNearLayerRetirement(canonical);
}
private bool AdvanceGeneration()
{
if (!ConvergePendingPublications())
return false;
_generation = unchecked(_generation + 1);
return true;
}
private bool ConvergePendingPublications(
bool preferDestination = false)
{
IReadOnlyList pending =
_presentation.GetPendingPublicationResults();
if (pending.Count == 0)
return true;
if (_activeWorkMeter is null)
{
// Settings/native callbacks may request policy changes between
// frames, but publication remains owned by the frame-scoped meter.
return false;
}
try
{
bool progressed = false;
bool destinationPending = false;
if (preferDestination)
{
for (int i = 0; i < pending.Count; i++)
{
if (!IsDestinationWork(pending[i].LandblockId))
continue;
destinationPending = true;
break;
}
}
// A partially prepared background landblock must not serialize the
// reveal-critical destination behind dictionary insertion order.
// Each publication owns independent receipts, so destination
// transactions can safely resume first without replaying or
// discarding the non-destination cursor.
for (int destinationPass = 1; destinationPass >= 0; destinationPass--)
{
bool requireDestination = destinationPass != 0;
for (int i = 0; i < pending.Count; i++)
{
LandblockStreamResult result = pending[i];
bool isDestination = IsDestinationWork(result.LandblockId);
if (destinationPending && !isDestination)
continue;
if (isDestination != requireDestination)
continue;
using StreamingWorkMeter.LaneScope lane =
_activeWorkMeter.EnterLane(
isDestination
? StreamingWorkLane.Destination
: StreamingWorkLane.NonDestination);
LandblockPublicationAdvance advance =
_presentation.ResumePublication(
result,
_activeWorkMeter,
ensureProgress: !progressed);
progressed |= advance.Progressed;
if (!advance.Completed)
return false;
}
}
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.
_completionQueue.RemoveResults(
pending,
result => !_presentation.HasPendingPublication(result));
}
}
///
/// Advance one frame. /
/// are landblock coordinates (0..255) of the current viewer — the camera
/// in offline mode, the server-sent player position in live.
///
/// Two-tier model (Phase A.5 T13):
///
/// - → enqueue LoadFar (terrain only, no entities)
/// - → enqueue LoadNear (terrain + entities)
/// - → enqueue PromoteToNear (entity layer for already-loaded terrain)
/// - → drop entities on render thread immediately (terrain stays)
/// - → enqueue full unload
///
///
public void Tick(int observerCx, int observerCy, bool insideDungeon = false)
{
if (_activeWorkMeter is not null)
throw new InvalidOperationException(
"StreamingController.Tick cannot be reentered.");
var meter = new StreamingWorkMeter(
_workBudget,
destinationReservationActive: _destinationReservation is not null);
_activeWorkMeter = meter;
try
{
if (_fullWindowRetirement is not null)
{
bool retirementWasPending =
_presentation.PendingRetirementCount != 0;
TryAdvanceFullWindowRetirement(meter);
AdvanceDestinationRetirementDependency(meter);
if (_presentation.UsesBudgetedRetirementSteps
|| retirementWasPending)
{
_presentation.AdvanceRetirements(meter);
}
if (_fullWindowRetirement is
{
PreparationCommitted: true,
}
&& _presentation.PendingRetirementCount == 0)
{
_fullWindowRetirement = null;
}
return;
}
if (_originRecenterRetirement is not null)
{
bool retirementWasPending =
_presentation.PendingRetirementCount != 0;
TryAdvanceOriginRecenterPreparation(meter);
AdvanceDestinationRetirementDependency(meter);
if (_presentation.UsesBudgetedRetirementSteps
|| retirementWasPending)
{
_presentation.AdvanceRetirements(meter);
}
return;
}
if (_pendingRadiiReconfiguration is not null)
AdvanceRadiiReconfiguration();
else if (_deferredRadiiRequest is { } deferred)
{
_deferredRadiiRequest = null;
ReconfigureRadii(deferred.NearRadius, deferred.FarRadius);
}
// 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.
bool retirementWasPendingAtFrameStart =
_presentation.PendingRetirementCount != 0;
AdvanceDestinationRetirementDependency(meter);
_presentation.AdvanceRetirements(meter);
bool destinationPublicationIncomplete =
_destinationReservation is not null
&& !IsRenderNeighborhoodResident(
DestinationLandblockId,
DestinationRadius);
if (!ConvergePendingPublications(
preferDestination: destinationPublicationIncomplete))
return;
uint centerId = StreamingRegion.EncodeLandblockId(observerCx, observerCy);
if (_collapsed)
{
// Hysteresis. Cases:
// - Still in the SAME dungeon landblock → hold (sweep stragglers).
// - In a DIFFERENT dungeon cell (multi-landblock dungeon / new dungeon)
// → re-collapse onto it.
// - CurrCell flickered null but the player hasn't gone anywhere: the
// observer landblock reverts to the position-derived value, which for a
// dungeon is only ever the ADJACENT off-by-one landblock (negative cell-
// local Y). Hold — never expand on an adjacent flicker.
// - Genuinely left to a DISTANT landblock (portal/teleport out, always far
// from the ocean-grid dungeon block) → expand.
if (insideDungeon && centerId != _collapsedCenter)
EnterDungeonCollapse(observerCx, observerCy, centerId);
else if (!insideDungeon && ChebyshevLandblocks(centerId, _collapsedCenter) > 1)
ExitDungeonExpand(observerCx, observerCy);
else
SweepCollapsed();
}
else if (insideDungeon)
{
EnterDungeonCollapse(observerCx, observerCy, centerId);
}
else
{
NormalTick(observerCx, observerCy);
}
DrainAndApply(
preferDestination: destinationPublicationIncomplete);
// Retirement is cleanup after immediate spatial withdrawal. Spend
// remaining capacity only after visible/destination publication
// work has had access to this frame's shared meter.
if (_presentation.UsesBudgetedRetirementSteps
&& !retirementWasPendingAtFrameStart)
_presentation.AdvanceRetirements(meter);
}
finally
{
meter.FinishFrame();
StreamingWorkMeterSnapshot snapshot = meter.Snapshot;
ObserveWorkLifetime(snapshot);
_lastWorkMeter = snapshot;
_activeWorkMeter = null;
}
}
private void AdvanceDestinationRetirementDependency(
StreamingWorkMeter meter)
{
if (_destinationReservation is null
|| !_presentation.IsRetirementPending(DestinationLandblockId))
{
return;
}
using StreamingWorkMeter.LaneScope lane =
meter.EnterLane(StreamingWorkLane.Destination);
_presentation.AdvancePriorityRetirement(
DestinationLandblockId,
meter);
}
private void ObserveWorkLifetime(StreamingWorkMeterSnapshot snapshot)
{
_lifetimeWorkOverruns = SaturatingAdd(
_lifetimeWorkOverruns,
snapshot.OverrunCount);
_lifetimeOversizedProgress = SaturatingAdd(
_lifetimeOversizedProgress,
snapshot.OversizedProgressCount);
if (_maximumWorkFrameStage is null
|| snapshot.ElapsedMilliseconds > _maximumWorkFrameMilliseconds)
{
_maximumWorkFrameMilliseconds = snapshot.ElapsedMilliseconds;
_maximumWorkFrameStage = snapshot.LastStage;
}
if (_maximumWorkOperationStage is null
|| snapshot.MaximumOperationMilliseconds
> _maximumWorkOperationMilliseconds)
{
_maximumWorkOperationMilliseconds =
snapshot.MaximumOperationMilliseconds;
_maximumWorkOperationStage = snapshot.MaximumOperationStage;
}
}
private static long SaturatingAdd(long left, int right) =>
left > long.MaxValue - right ? long.MaxValue : left + right;
private void AdvanceRadiiReconfiguration()
{
if (_advancingRadiiReconfiguration)
return;
RadiiReconfiguration pending = _pendingRadiiReconfiguration
?? throw new InvalidOperationException("No radii reconfiguration is pending.");
var failures = new List();
_advancingRadiiReconfiguration = true;
try
{
if (!pending.GenerationAdvanced)
{
if (!AdvanceGeneration())
return;
pending.GenerationAdvanced = true;
}
if (!pending.PendingLoadsCleared)
{
try
{
_clearPendingLoads?.Invoke();
pending.PendingLoadsCleared = true;
}
catch (Exception error)
{
if (error is StreamingMutationException { MutationCommitted: true })
pending.PendingLoadsCleared = true;
failures.Add(error);
}
}
// Do not admit new-generation work until cancellation of the old
// inbox is durable. Otherwise a successful retry of ClearPendingLoads
// would also erase mutations already marked complete by this ledger.
for (int i = 0; pending.PendingLoadsCleared && i < pending.Mutations.Count; i++)
{
RadiiMutation mutation = pending.Mutations[i];
if (mutation.Completed)
continue;
try
{
mutation.Apply();
mutation.Completed = true;
}
catch (Exception error)
{
if (error is StreamingMutationException { MutationCommitted: true })
mutation.Completed = true;
failures.Add(error);
}
}
bool converged = pending.PendingLoadsCleared
&& pending.Mutations.All(static mutation => mutation.Completed);
if (converged)
{
// Publish the new logical window only after every queue/retirement
// admission is durable. A later Tick can now trust Resident as the
// complete desired set and will never strand a hole.
pending.Region.MarkResidentFromBootstrap();
_region = pending.Region;
NearRadius = pending.NearRadius;
FarRadius = pending.FarRadius;
_completionQueue.Clear();
_pendingRadiiReconfiguration = null;
}
}
finally
{
_advancingRadiiReconfiguration = false;
}
if (failures.Count != 0)
{
throw new AggregateException(
"Streaming quality reconfiguration did not complete cleanly.",
failures);
}
if (_pendingRadiiReconfiguration is null
&& _deferredRadiiRequest is { } deferred)
{
_deferredRadiiRequest = null;
ReconfigureRadii(deferred.NearRadius, deferred.FarRadius);
}
}
private sealed class RadiiReconfiguration
{
public RadiiReconfiguration(
int nearRadius,
int farRadius,
StreamingRegion region,
List mutations)
{
NearRadius = nearRadius;
FarRadius = farRadius;
Region = region;
Mutations = mutations;
}
public int NearRadius { get; }
public int FarRadius { get; }
public StreamingRegion Region { get; }
public List Mutations { get; }
public bool GenerationAdvanced { get; set; }
public bool PendingLoadsCleared { get; set; }
}
private sealed class RadiiMutation
{
public RadiiMutation(Action apply) => Apply = apply;
public Action Apply { get; }
public bool Completed { get; set; }
}
///
/// #135: collapse to a single dungeon landblock IMMEDIATELY, before the first
/// has a chance to bootstrap the full 25×25 window. Called
/// from the login / teleport spawn path the instant the streaming center is
/// recentered onto a SEALED dungeon landblock.
///
/// The per-frame insideDungeon gate keys on the physics
/// CurrCell, which is only set once the player is PLACED — and placement
/// waits for the dungeon landblock to hydrate. So for the whole hydration window
/// (tens of seconds for a ~200-cell dungeon) the gate reads false and
/// would enqueue the ~24 unrelated ocean-grid neighbor
/// dungeons (+ ~19k entities each); the collapse then only mops them up after
/// placement. That mop-up is the 10→high FPS ramp users see at a dungeon login.
///
/// Pre-collapsing means the EXPENSIVE dungeon-neighbour window is never
/// enqueued. On teleport nothing is enqueued at all (this fires before the next
/// Tick recenters). On login a brief Holtburg outdoor window may be enqueued by the
/// frame-1 NormalTick (before the player's spawn arrives) and is immediately
/// cancelled by _clearPendingLoads here — cheap outdoor terrain, not the
/// ocean-grid dungeons, and a handful of already-dequeued loads get swept next
/// frame. Idempotent: a no-op when already collapsed onto this same landblock, so a
/// re-sent spawn or a same-frame double call costs nothing. Render-thread only,
/// same as .
///
public void InitializeKnownLoginCenter(int cx, int cy, bool isSealedDungeon)
{
// StreamingReadinessGate keeps the worker stopped until the player's
// real spawn supplies this center. There is therefore no guessed-center
// window to force-reload here. A sealed dungeon is the one special case:
// pin radius zero before the first Tick so no ocean-grid neighbours are
// ever enqueued.
if (isSealedDungeon)
PreCollapseToDungeon(cx, cy);
}
///
/// Pins streaming to a sealed dungeon before the first normal Tick.
///
public void PreCollapseToDungeon(int cx, int cy)
{
uint centerId = StreamingRegion.EncodeLandblockId(cx, cy);
if (_collapsed && _collapsedCenter == centerId) return;
EnterDungeonCollapse(cx, cy, centerId);
}
///
/// Outdoor / building-interior streaming — the original two-tier model.
///
private void NormalTick(int observerCx, int observerCy)
{
if (_region is null)
{
_region = new StreamingRegion(observerCx, observerCy, NearRadius, FarRadius);
var bootstrap = _region.ComputeFirstTickDiff();
EnqueueLoadsByRevealPriority(
bootstrap.ToLoadNear,
LandblockStreamJobKind.LoadNear);
foreach (var id in bootstrap.ToLoadFar) EnqueueLoad(id, LandblockStreamJobKind.LoadFar);
_region.MarkResidentFromBootstrap();
}
else if (_region.CenterX != observerCx || _region.CenterY != observerCy)
{
var diff = _region.RecenterTo(observerCx, observerCy);
EnqueueLoadsByRevealPriority(
diff.ToPromote,
LandblockStreamJobKind.PromoteToNear);
EnqueueLoadsByRevealPriority(
diff.ToLoadNear,
LandblockStreamJobKind.LoadNear);
foreach (var id in diff.ToLoadFar) EnqueueLoad(id, LandblockStreamJobKind.LoadFar);
foreach (var id in diff.ToDemote) DemoteLandblock(id);
foreach (var id in diff.ToUnload) EnqueueUnload(id);
}
}
///
/// Places the canonical reveal neighborhood at the front of the worker's
/// existing Near FIFO. Completion priority alone is too late: a
/// non-destination build can otherwise begin a multi-frame publication
/// before the destination result has even reached the render thread.
///
private void EnqueueLoadsByRevealPriority(
IReadOnlyList landblockIds,
LandblockStreamJobKind kind,
bool skipLoaded = false)
{
for (int destinationPass = 1; destinationPass >= 0; destinationPass--)
{
bool requireDestination = destinationPass != 0;
for (int i = 0; i < landblockIds.Count; i++)
{
uint id = landblockIds[i];
if ((!skipLoaded || !_state.IsLoaded(id))
&& IsDestinationWork(id) == requireDestination)
{
EnqueueLoad(id, kind);
}
}
}
}
///
/// Dungeon-entry edge: cancel the in-flight window load, unload every
/// resident neighbor, and pin streaming to the player's single dungeon
/// landblock. Retail-faithful — AC dungeons have no adjacent landblocks
/// (ACE LandblockManager.GetAdjacentIDs returns empty for a dungeon);
/// the 25×25 window was pulling in ~129 unrelated ocean-grid dungeons and
/// their thousands of emitters (#133 FPS). Unloading them also tears down
/// their lights, shrinking the static-light set toward retail's ≤40.
///
private void EnterDungeonCollapse(int cx, int cy, uint centerId)
{
bool logTransition = !_collapsed || _collapsedCenter != centerId;
if (logTransition)
Console.WriteLine($"streaming: dungeon collapse -> 0x{centerId:X8}");
if (!AdvanceGeneration())
return;
_collapsed = true;
_collapsedCenter = centerId;
_clearPendingLoads?.Invoke();
_completionQueue.Clear();
_region = null;
foreach (var id in _state.LoadedLandblockIds)
if (id != centerId) EnqueueUnload(id);
// Pin a radius-0 region so RecenterTo never re-expands while inside,
// and so the post-exit rebuild starts from a clean, consistent state.
_region = new StreamingRegion(cx, cy, 0, 0);
_region.MarkResidentFromBootstrap();
// The dungeon landblock itself must be (or become) loaded. If a prior
// ClearPendingLoads cancelled its queued load, re-enqueue it.
if (!_state.IsLoaded(centerId))
EnqueueLoad(centerId, LandblockStreamJobKind.LoadNear);
else if (!_state.IsNearTier(centerId))
EnqueueLoad(centerId, LandblockStreamJobKind.PromoteToNear);
}
///
/// While collapsed, unload any landblock that finished loading after the
/// collapse edge — a Load the worker had already dequeued before the
/// control job took
/// effect. At steady state only the dungeon landblock is resident, so this
/// is a no-op.
///
private void SweepCollapsed()
{
// Always preserve the true dungeon landblock (_collapsedCenter), never the
// per-frame observer landblock — a CurrCell flicker must not unload the dungeon.
foreach (var id in _state.LoadedLandblockIds)
if (id != _collapsedCenter) EnqueueUnload(id);
}
/// Chebyshev distance in landblock cells between two landblock ids.
private static int ChebyshevLandblocks(uint a, uint b)
{
int ax = (int)((a >> 24) & 0xFFu), ay = (int)((a >> 16) & 0xFFu);
int bx = (int)((b >> 24) & 0xFFu), by = (int)((b >> 16) & 0xFFu);
return Math.Max(Math.Abs(ax - bx), Math.Abs(ay - by));
}
///
/// True when belongs to the canonical reveal
/// generation's destination neighborhood.
///
private bool IsDestinationWork(uint id)
=> _destinationReservation is { } reservation
&& ChebyshevLandblocks(id, reservation.LandblockId)
<= reservation.Radius;
///
/// Dungeon-exit edge (portal to outdoors / teleport): rebuild the full
/// two-tier window at the new center and unload anything resident from the
/// collapsed state that falls outside it.
///
private void ExitDungeonExpand(int observerCx, int observerCy)
{
Console.WriteLine(
$"streaming: dungeon EXIT-expand -> ({observerCx},{observerCy}) " +
$"(was collapsed on 0x{_collapsedCenter:X8})");
if (!AdvanceGeneration())
return;
_collapsed = false;
var rebuilt = new StreamingRegion(observerCx, observerCy, NearRadius, FarRadius);
foreach (var id in _state.LoadedLandblockIds)
if (!rebuilt.Resident.Contains(id)) EnqueueUnload(id);
var boot = rebuilt.ComputeFirstTickDiff();
EnqueueLoadsByRevealPriority(
boot.ToLoadNear,
LandblockStreamJobKind.LoadNear,
skipLoaded: true);
foreach (var id in boot.ToLoadFar)
if (!_state.IsLoaded(id)) EnqueueLoad(id, LandblockStreamJobKind.LoadFar);
rebuilt.MarkResidentFromBootstrap();
_region = rebuilt;
}
///
/// Starts a reload of the current streaming window. The next and later
/// calls invalidate accepted work, capture and detach
/// stable residents, and converge render, physics, and static-lighting
/// owners under the shared frame budget before
/// bootstraps the window again. Hard login/portal transitions separately
/// publish their immediate observable edge through
/// .
/// Shared-origin teleports use instead so
/// every old-window owner retires while the old coordinate frame is active.
///
public void ForceReloadWindow()
{
BeginFullWindowRetirement();
}
///
/// Starts the old-window half of a shared-origin recenter. All accepted
/// worker work is invalidated and every resident landblock is detached
/// through the retryable presentation ledger while the old coordinate
/// frame is still active. remains blocked until
/// succeeds.
///
internal void BeginOriginRecenter()
{
_originRecenterRetirement ??= new OriginRecenterRetirement();
}
///
/// Reports whether every old resident has been detached and captured by
/// an exact retirement receipt. Deferred receipt cleanup may continue
/// after the shared origin changes; same-landblock publication remains
/// fenced by .
///
internal bool IsOriginRecenterRetirementComplete()
{
if (_originRecenterRetirement is null)
throw new InvalidOperationException(
"No streaming-origin recenter transaction is pending.");
return _originRecenterRetirement.PreparationCommitted;
}
///
/// Releases the streaming bootstrap gate after the composition root has
/// committed the new shared origin. Old detached receipts continue under
/// the shared frame meter and fence only a replacement publication with
/// the same canonical landblock key.
///
internal bool TryCommitOriginRecenter(
int destinationX,
int destinationY,
bool isSealedDungeon)
{
if (_advancingOriginRecenter)
return false;
_advancingOriginRecenter = true;
try
{
return TryCommitOriginRecenterCore(
destinationX,
destinationY,
isSealedDungeon);
}
finally
{
_advancingOriginRecenter = false;
}
}
private bool TryCommitOriginRecenterCore(
int destinationX,
int destinationY,
bool isSealedDungeon)
{
OriginRecenterRetirement transaction = _originRecenterRetirement
?? throw new InvalidOperationException(
"No streaming-origin recenter transaction is pending.");
if (!transaction.PreparationCommitted)
throw new InvalidOperationException(
"Streaming-origin retirement preparation has not completed.");
var destination = (destinationX, destinationY, isSealedDungeon);
if (transaction.Destination is { } retained && retained != destination)
{
throw new InvalidOperationException(
"A recenter transaction cannot commit two different destinations.");
}
transaction.Destination ??= destination;
if (!transaction.DestinationConfigured)
{
_collapsed = isSealedDungeon;
_collapsedCenter = isSealedDungeon
? StreamingRegion.EncodeLandblockId(destinationX, destinationY)
: 0u;
if (isSealedDungeon)
{
_region = new StreamingRegion(
destinationX,
destinationY,
nearRadius: 0,
farRadius: 0);
_region.MarkResidentFromBootstrap();
}
else
{
_region = null;
}
transaction.DestinationConfigured = true;
}
if (isSealedDungeon && !transaction.DestinationLoadEnqueued)
{
try
{
EnqueueLoad(
StreamingRegion.EncodeLandblockId(destinationX, destinationY),
LandblockStreamJobKind.LoadNear);
transaction.DestinationLoadEnqueued = true;
}
catch (StreamingMutationException error) when (error.MutationCommitted)
{
transaction.DestinationLoadEnqueued = true;
Console.WriteLine(
$"streaming: committed dungeon recenter enqueue reported failure: {error}");
return false;
}
catch (Exception error)
{
Console.WriteLine(
$"streaming: dungeon recenter enqueue will resume: {error}");
return false;
}
}
_originRecenterRetirement = null;
return true;
}
///
/// Releases a fully retired origin transaction at a session boundary
/// without bootstrapping a destination from the ending session.
///
internal bool TryCancelOriginRecenter()
{
if (_originRecenterRetirement is not { PreparationCommitted: true })
return false;
_collapsed = false;
_collapsedCenter = 0u;
_region = null;
_originRecenterRetirement = null;
return true;
}
private bool TryAdvanceOriginRecenterPreparation(StreamingWorkMeter meter)
{
ArgumentNullException.ThrowIfNull(meter);
OriginRecenterRetirement transaction = _originRecenterRetirement
?? throw new InvalidOperationException(
"No streaming-origin recenter transaction is pending.");
if (transaction.PreparationCommitted)
return true;
if (_advancingOriginRecenter)
return false;
_advancingOriginRecenter = true;
try
{
if (!transaction.RadiiConverged)
{
if (_pendingRadiiReconfiguration is not null
&& !TryRunStreamingWork(
meter,
new StreamingWorkCost(
EntityOperations: Math.Max(
1,
_pendingRadiiReconfiguration.Mutations.Count)),
"recenter-radii-convergence",
() =>
{
AdvanceRadiiReconfiguration();
return _pendingRadiiReconfiguration is null;
}))
{
return false;
}
if (_pendingRadiiReconfiguration is not null)
return false;
transaction.RadiiConverged = true;
}
if (!transaction.GenerationAdvanced)
{
// AdvanceGeneration meters any retained publication it must
// converge through the active frame meter. The token increment
// itself is constant-time and must not create a nested meter
// reservation around that retry.
if (!AdvanceGeneration())
return false;
transaction.GenerationAdvanced = true;
}
if (!transaction.PendingLoadsCleared)
{
bool ClearPendingLoads()
{
try
{
_clearPendingLoads?.Invoke();
transaction.PendingLoadsCleared = true;
return true;
}
catch (StreamingMutationException error) when (error.MutationCommitted)
{
transaction.PendingLoadsCleared = true;
Console.WriteLine(
$"streaming: committed pending-load clear reported failure: {error}");
return false;
}
}
if (!TryRunStreamingWork(
meter,
new StreamingWorkCost(EntityOperations: 1),
"recenter-clear-worker-inbox",
ClearPendingLoads))
return false;
}
if (!transaction.CompletionQueueCleared)
{
while (_completionQueue.Count != 0)
{
StreamingWorkAdmission admission = meter.TryReserve(
new StreamingWorkCost(EntityOperations: 1),
"recenter-release-completion");
if (admission == StreamingWorkAdmission.Yielded)
return false;
if (!_completionQueue.TryRemoveOne())
{
meter.Fail();
throw new InvalidOperationException(
"Completion queue count changed during recenter release.");
}
meter.Complete();
}
transaction.CompletionQueueCleared = true;
}
if (!transaction.RegionCleared)
{
if (!TryRunStreamingWork(
meter,
new StreamingWorkCost(EntityOperations: 1),
"recenter-clear-region",
() =>
{
_collapsed = false;
_region = null;
return true;
}))
{
return false;
}
transaction.RegionCleared = true;
}
if (!transaction.SpatialGenerationDetached)
{
StreamingWorkAdmission admission = meter.TryReserve(
new StreamingWorkCost(
EntityOperations:
_state.OriginRecenterSpatialOperationCount),
"recenter-detach-spatial-generation",
ensureProgress: true);
if (admission == StreamingWorkAdmission.Yielded)
return false;
try
{
GpuWorldRecenterRetirement detached =
_presentation.DetachAllForOriginRecenter();
transaction.SpatialGenerationDetached = true;
FullWindowRetirementCount++;
LastFullWindowRetirementLandblockCount =
detached.Landblocks.Count;
meter.Complete();
if (detached.ObserverFailure is not null)
{
Console.WriteLine(
"streaming: committed origin-recenter spatial " +
$"generation reported failure: {detached.ObserverFailure}");
return false;
}
}
catch
{
meter.Fail();
throw;
}
}
transaction.PreparationCommitted = true;
return true;
}
catch (Exception error)
{
Console.WriteLine(
$"streaming: origin-recenter preparation will resume: {error}");
return false;
}
finally
{
_advancingOriginRecenter = false;
}
}
private void BeginFullWindowRetirement()
{
_fullWindowRetirement ??= new FullWindowRetirement();
}
private bool TryAdvanceFullWindowRetirement(StreamingWorkMeter meter)
{
FullWindowRetirement transaction = _fullWindowRetirement
?? throw new InvalidOperationException(
"No full-window retirement transaction is pending.");
if (transaction.PreparationCommitted)
return true;
try
{
if (!transaction.GenerationAdvanced)
{
if (!AdvanceGeneration())
{
return false;
}
transaction.GenerationAdvanced = true;
}
if (!transaction.PendingLoadsCleared)
{
bool ClearPendingLoads()
{
try
{
_clearPendingLoads?.Invoke();
transaction.PendingLoadsCleared = true;
return true;
}
catch (StreamingMutationException error) when (error.MutationCommitted)
{
transaction.PendingLoadsCleared = true;
Console.WriteLine(
$"streaming: committed reload pending-load clear reported failure: {error}");
return false;
}
}
if (!TryRunStreamingWork(
meter,
new StreamingWorkCost(EntityOperations: 1),
"reload-clear-worker-inbox",
ClearPendingLoads))
{
return false;
}
}
if (!transaction.CompletionQueueCleared)
{
while (_completionQueue.Count != 0)
{
StreamingWorkAdmission admission = meter.TryReserve(
new StreamingWorkCost(EntityOperations: 1),
"reload-release-completion");
if (admission == StreamingWorkAdmission.Yielded)
return false;
if (!_completionQueue.TryRemoveOne())
{
meter.Fail();
throw new InvalidOperationException(
"Completion queue count changed during reload release.");
}
meter.Complete();
}
transaction.CompletionQueueCleared = true;
}
if (!transaction.RegionCleared)
{
if (!TryRunStreamingWork(
meter,
new StreamingWorkCost(EntityOperations: 1),
"reload-clear-region",
() =>
{
_collapsed = false;
_region = null;
return true;
}))
{
return false;
}
transaction.RegionCleared = true;
}
if (transaction.ResidentIds is null)
{
transaction.ResidentIds = [];
transaction.ResidentEnumerator =
_state.LoadedLandblockIds.GetEnumerator();
}
while (transaction.ResidentEnumerator is { } residentEnumerator)
{
StreamingWorkAdmission admission = meter.TryReserve(
new StreamingWorkCost(EntityOperations: 1),
"reload-capture-resident-id");
if (admission == StreamingWorkAdmission.Yielded)
return false;
bool moved;
try
{
moved = residentEnumerator.MoveNext();
if (moved)
transaction.ResidentIds.Add(residentEnumerator.Current);
else
{
residentEnumerator.Dispose();
transaction.ResidentEnumerator = null;
FullWindowRetirementCount++;
LastFullWindowRetirementLandblockCount =
transaction.ResidentIds.Count;
}
meter.Complete();
}
catch
{
meter.Fail();
throw;
}
}
List residentIds = transaction.ResidentIds
?? throw new InvalidOperationException(
"Full-window resident capture did not commit.");
while (transaction.RetirementCursor < residentIds.Count)
{
uint id = residentIds[transaction.RetirementCursor];
int entityCount = _state.TryGetLandblock(
id,
out LoadedLandblock? loaded)
? loaded!.Entities.Count
: 0;
StreamingWorkAdmission admission = meter.TryReserve(
new StreamingWorkCost(
EntityOperations: Math.Max(1, entityCount)),
$"reload-detach-0x{id:X8}");
if (admission == StreamingWorkAdmission.Yielded)
return false;
try
{
_presentation.EnqueueFullRetirement(id);
transaction.RetirementCursor++;
meter.Complete();
}
catch (Exception error)
{
if (!_state.IsLoaded(id))
transaction.RetirementCursor++;
meter.Fail();
Console.WriteLine(
$"streaming: full-window retirement for 0x{id:X8} " +
$"will resume: {error}");
return false;
}
}
transaction.PreparationCommitted = true;
return true;
}
catch (Exception error)
{
Console.WriteLine(
$"streaming: full-window preparation will resume: {error}");
return false;
}
}
///
/// 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(bool preferDestination = false)
{
StreamingWorkMeter meter = _activeWorkMeter
?? throw new InvalidOperationException(
"Completion scheduling requires an active frame meter.");
AdmitCompletions(meter);
bool destinationQueued =
preferDestination
&& _completionQueue.HasPriority(
StreamingCompletionPriority.Destination);
bool executed = false;
while (_completionQueue.TryPeekNext(
_isPublicationBlockedByRetirement,
out StreamingQueuedCompletion? completion,
destinationQueued
? StreamingCompletionPriority.Unload
: StreamingCompletionPriority.Far))
{
StreamingQueuedCompletion work = completion
?? throw new InvalidOperationException(
"The completion queue returned a null head.");
try
{
using StreamingWorkMeter.LaneScope lane = meter.EnterLane(
work.Priority == StreamingCompletionPriority.Destination
&& work.RevealGeneration == ActiveRevealGeneration
? StreamingWorkLane.Destination
: StreamingWorkLane.NonDestination);
LandblockPublicationAdvance advance = ApplyResult(
work,
meter,
ensureProgress: !executed);
executed |= advance.Progressed;
if (!advance.Completed)
break;
_completionQueue.RemoveHead(work);
}
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);
throw;
}
}
}
private static bool TryRunStreamingWork(
StreamingWorkMeter meter,
StreamingWorkCost cost,
string stage,
Func operation)
{
StreamingWorkAdmission admission = meter.TryReserve(cost, stage);
if (admission == StreamingWorkAdmission.Yielded)
return false;
try
{
bool completed = operation();
if (completed)
meter.Complete();
else
meter.Fail();
return completed;
}
catch
{
meter.Fail();
throw;
}
}
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);
StreamingCompletionPriority priority = stale
? StreamingCompletionPriority.Control
: ClassifyCompletion(result);
LandblockStreamCostEstimate estimate =
LandblockStreamResultCost.Estimate(result);
// A stale result belongs to a generation that was already
// cancelled. Reading it releases worker-outbox ownership; it does
// not admit payload into the current world. Charge only elapsed
// time so a large completed old window cannot consume the
// destination generation's completion quota for many frames.
StreamingWorkCost admissionCost = stale
? default
: new StreamingWorkCost(
CompletionAdmissions:
estimate.Work.CompletionAdmissions,
AdoptedCpuBytes:
estimate.Work.AdoptedCpuBytes);
using StreamingWorkMeter.LaneScope lane = meter.EnterLane(
priority == StreamingCompletionPriority.Destination
? StreamingWorkLane.Destination
: StreamingWorkLane.NonDestination);
StreamingWorkAdmission admission = meter.TryReserve(
admissionCost,
"completion-admission");
if (admission == StreamingWorkAdmission.Yielded)
return;
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,
priority,
priority == StreamingCompletionPriority.Destination
? ActiveRevealGeneration
: 0L,
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
&& IsDestinationWork(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 static LandblockPublicationAdvance RunSimpleResultOperation(
StreamingWorkMeter meter,
string stage,
bool ensureProgress,
Action operation)
{
StreamingWorkAdmission admission = meter.TryReserve(
default,
stage,
ensureProgress);
if (admission == StreamingWorkAdmission.Yielded)
return new LandblockPublicationAdvance(false, false);
try
{
operation();
meter.Complete();
return new LandblockPublicationAdvance(true, true);
}
catch
{
meter.Fail();
throw;
}
}
///
/// Apply a single with the full side-
/// effects: terrain upload, GPU state, and the re-hydration callback.
/// All priority queues route through this one publication path.
///
private LandblockPublicationAdvance ApplyResult(
StreamingQueuedCompletion work,
StreamingWorkMeter meter,
bool ensureProgress)
{
LandblockStreamResult result = work.Result;
if (_presentation.HasPendingPublication(result))
{
LandblockPublicationAdvance resumed =
_presentation.ResumePublication(
result,
meter,
ensureProgress);
if (resumed.Completed)
ReconcileCompletedPendingPublication(result.LandblockId);
return resumed;
}
if (IsStaleGeneration(result))
{
// A worker can finish one old load/unload after a hard recenter.
// Landblock id membership cannot distinguish overlapping windows;
// generation is the logical streaming incarnation boundary.
return RunSimpleResultOperation(
meter,
"execute-stale-generation",
ensureProgress,
static () => { });
}
if (result is LandblockStreamResult.Unloaded
&& _region?.TryGetDesiredTier(result.LandblockId, out _) == true)
{
// Normal recenter does not advance the hard-window generation. A
// rapid away->back can therefore re-own an id while its same-
// generation unload is already in the worker outbox. Current
// region ownership wins; the old unload must not tear it down.
return RunSimpleResultOperation(
meter,
"execute-reowned-unload",
ensureProgress,
static () => { });
}
if (result is LandblockStreamResult.Loaded
or LandblockStreamResult.Promoted)
{
if (_region is null
|| !_region.TryGetDesiredTier(result.LandblockId, out var desiredTier))
{
// A hard recenter/collapse can leave one worker job already in
// flight. Its completion belongs to the old region and must not
// resurrect a landblock the new StreamingRegion never owns.
return RunSimpleResultOperation(
meter,
"execute-undesired-load",
ensureProgress,
static () => { });
}
bool isNearCompletion = result is LandblockStreamResult.Promoted
or LandblockStreamResult.Loaded { Tier: LandblockStreamTier.Near };
if (isNearCompletion
&& _state.IsNearTier(result.LandblockId)
&& !_presentation.HasPendingPublication(result))
{
// Normal recenter can enqueue a second high-priority Near job
// while an earlier one is already in flight. The streamer
// supersedes Far/Unload work but intentionally does not dedupe
// high-priority jobs. Publication is idempotent at this owner:
// never append statics, replay defaults, or reapply the complete
// Near transaction to an already-Near landblock.
return RunSimpleResultOperation(
meter,
"execute-duplicate-near",
ensureProgress,
static () => { });
}
if (desiredTier == LandblockStreamTier.Far && isNearCompletion)
{
// Recenter can demote a Near job before its first completion.
// StreamingRegion already considers that landblock Far, so it
// deliberately emits no second LoadFar job. If no terrain is
// resident yet, publish the completed heightmap/mesh as a Far
// result while discarding its entity and EnvCell layers. This
// is the same payload a fresh LoadFar would have produced and
// closes the in-flight Near -> desired Far lifecycle without a
// duplicate worker read or a permanent terrain hole.
if (!_state.IsLoaded(result.LandblockId))
{
switch (result)
{
case LandblockStreamResult.Loaded loaded:
return _presentation.PublishAsFar(
loaded,
loaded.Build,
loaded.MeshData,
meter,
ensureProgress);
case LandblockStreamResult.Promoted promoted:
return _presentation.PublishAsFar(
promoted,
promoted.Build,
promoted.MeshData,
meter,
ensureProgress);
}
}
return RunSimpleResultOperation(
meter,
"execute-already-far",
ensureProgress,
static () => { });
}
}
switch (result)
{
case LandblockStreamResult.Loaded loaded:
if (loaded.Tier == LandblockStreamTier.Far
&& _state.IsNearTier(loaded.LandblockId))
{
// This Far completion was queued before a newer Near load
// or promotion. Applying it would erase the entity/cell
// layer that now owns the landblock.
return RunSimpleResultOperation(
meter,
"execute-stale-far-tier",
ensureProgress,
static () => { });
}
return _presentation.PublishLoaded(
loaded,
work.Estimate,
meter,
ensureProgress);
case LandblockStreamResult.Promoted promoted:
// PromoteToNear carries a complete build and mesh because the
// streamer deliberately lets it supersede a queued LoadFar. If
// that Far job never started, publish this as the real Near
// landblock; if the base is already resident, merge only the
// Near layer so existing live projections retain identity.
return _presentation.PublishPromoted(
promoted,
mergeIntoExistingLandblock:
_state.IsLoaded(promoted.LandblockId),
work.Estimate,
meter,
ensureProgress);
case LandblockStreamResult.Unloaded unloaded:
return RunSimpleResultOperation(
meter,
"execute-unload",
ensureProgress,
() => _presentation.EnqueueFullRetirement(
unloaded.LandblockId));
case LandblockStreamResult.Failed failed:
return RunSimpleResultOperation(
meter,
"execute-load-failure",
ensureProgress,
() => Console.WriteLine(
$"streaming: load failed for 0x{failed.LandblockId:X8}: {failed.Error}"));
case LandblockStreamResult.WorkerCrashed crashed:
return RunSimpleResultOperation(
meter,
"execute-worker-crash",
ensureProgress,
() => Console.WriteLine(
$"streaming: worker CRASHED: {crashed.Error}"));
default:
throw new InvalidOperationException(
$"Unsupported streaming result {result.GetType().Name}.");
}
}
private void ReconcileCompletedPendingPublication(uint landblockId)
{
if (_region is null
|| !_region.TryGetDesiredTier(landblockId, out LandblockStreamTier desiredTier))
{
if (_state.IsLoaded(landblockId))
_presentation.EnqueueFullRetirement(landblockId);
return;
}
if (desiredTier == LandblockStreamTier.Far
&& _state.IsNearTier(landblockId))
{
_presentation.EnqueueNearLayerRetirement(landblockId);
}
}
private bool IsStaleGeneration(LandblockStreamResult result) =>
result is not LandblockStreamResult.WorkerCrashed
&& result.Generation != _generation;
private bool IsPublicationBlockedByRetirement(LandblockStreamResult result) =>
result is LandblockStreamResult.Loaded or LandblockStreamResult.Promoted
&& _presentation.IsRetirementPending(result.LandblockId);
///
/// Returns the landblock id associated with .
/// For this is 0 by
/// convention (not tied to a specific landblock).
///
private static uint ResultLandblockId(LandblockStreamResult result) => result.LandblockId;
}