ModernUO/Projects/Server/Serialization/SerializationChunkSource.cs
Kamron Batman 9acd701aaa
perf(saves): workers iterate entity dictionaries directly; main thread joins the drain
Removes the per-entity handoff from the freeze entirely. GenericEntityPersistence
publishes 4096-slot ranges over its dictionary's backing entries array, and workers
serialize occupied slots (value != null) directly via a ShadowEntry<TValue> struct
that mirrors the runtime's private Dictionary Entry layout. Safe because the
dictionary is frozen during Saving (mutations divert to the pending safety queues).

The layout is proven at startup before any code reads through it: validation
measures the true Entry stride via precise allocation accounting (so all shadow
reads are guaranteed in-bounds), then compares every key and value of a churned,
resized, freelist-exercised dictionary reading value slots as raw pointer bits
only - never materializing a managed reference until the layout is proven. If a
future runtime changes Dictionary internals, validation fails and saves fall back
to the enumerate-and-push path with a logged warning.

The main thread now joins the drain through an inline (threadless) worker after
publishing work, instead of idling while the thread workers finish - worth a full
worker share on the freeze and proportionally more on low-core hosts.

Measured through the real pipeline classes (24 cores, 10M entities, 1.7GB, dense
2-byte write profile): publish cost drops from ~55ms to ~0.1ms, the freeze is now
bound by pure serialize throughput at ~99ms steady state (vs ~740ms before this
branch, ~7.5x), and the first save after boot drops from ~468ms to ~139ms because
fine-grained ranges self-balance without needing size estimates.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-16 22:38:46 -07:00

161 lines
6.6 KiB
C#

/*************************************************************************
* ModernUO *
* Copyright 2019-2026 - ModernUO Development Team *
* Email: hi@modernuo.com *
* File: SerializationChunkSource.cs *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU General Public License as published by *
* the Free Software Foundation, either version 3 of the License, or *
* (at your option) any later version. *
* *
* You should have received a copy of the GNU General Public License *
* along with this program. If not, see <http://www.gnu.org/licenses/>. *
*************************************************************************/
using System;
using System.Collections.Concurrent;
using System.Runtime.CompilerServices;
using System.Runtime.InteropServices;
namespace Server;
/// <summary>
/// A range of backing-store slots that a serialization worker can serialize directly,
/// letting workers iterate a persistence's storage in parallel instead of the main thread
/// enumerating and handing off every entity. Implemented by
/// <see cref="GenericEntityPersistence{T}"/> over its dictionary's entries array.
/// </summary>
public interface ISlotRangeSource
{
/// <summary>
/// Serializes every occupied slot in [offset, offset + count) into the writer,
/// stamping each entity's thread/position/length. Returns the number serialized.
/// </summary>
int SerializeRange(BufferWriter writer, byte threadIndex, int offset, int count);
}
/// <summary>
/// Single-producer/multi-consumer handoff between the game loop and the serialization
/// thread workers during a world save. The producer batches entities into pooled chunks
/// so the per-entity cost is a plain array store instead of a synchronized enqueue, and
/// workers pull whole chunks so they naturally load-balance: a worker busy with a thick
/// entity simply takes fewer chunks.
/// Entities whose previous serialized size exceeds <see cref="HeavyEntityThreshold"/> are
/// published as dedicated single-entity chunks so multi-megabyte payloads spread across
/// workers instead of riding inside one chunk.
/// Persistences that support direct parallel iteration publish slot ranges instead of
/// filled chunks, removing the per-entity handoff from the freeze entirely.
/// </summary>
public sealed class SerializationChunkSource
{
// 4096 refs (32KB per chunk) keeps producer sync cost at one enqueue per 4096 entities
// while the drain tail stays sub-millisecond.
private const int ChunkCapacity = 4096;
/// <summary>
/// Entities whose previous <see cref="IGenericSerializable.SerializedLength"/> exceeds this
/// should be pushed with <see cref="PushSingle"/>. Callers do the check where the entity's
/// concrete type is known, so the size read is not an interface dispatch per entity.
/// </summary>
public const int HeavyEntityThreshold = 1024 * 1024; // 1MB
internal readonly struct Chunk
{
public readonly IGenericSerializable Single;
public readonly IGenericSerializable[] Buffer;
public readonly ISlotRangeSource Source;
public readonly int Offset;
public readonly int Count;
public Chunk(IGenericSerializable single)
{
Single = single;
Count = 1;
}
public Chunk(IGenericSerializable[] buffer, int count)
{
Buffer = buffer;
Count = count;
}
public Chunk(ISlotRangeSource source, int offset, int count)
{
Source = source;
Offset = offset;
Count = count;
}
}
private readonly ConcurrentQueue<Chunk> _chunks = new();
private readonly ConcurrentQueue<IGenericSerializable[]> _pool = new();
// Producer state - written only by the game loop thread.
private IGenericSerializable[] _current;
private int _count;
[MethodImpl(MethodImplOptions.AggressiveInlining)]
public void Push(IGenericSerializable entity)
{
var current = _current ??= Rent();
// Ref store skips the bounds and array-covariance checks. Safe by construction:
// _count is producer-thread-only and always < ChunkCapacity here (reset on publish),
// and the array's element type is exactly IGenericSerializable.
Unsafe.Add(ref MemoryMarshal.GetArrayDataReference(current), _count) = entity;
if (++_count == ChunkCapacity)
{
_chunks.Enqueue(new Chunk(current, ChunkCapacity));
_current = null;
_count = 0;
}
}
/// <summary>
/// Publishes the entity as a dedicated chunk regardless of its estimated size.
/// Used for persistence self-payloads, which can be large on the first save
/// before a <see cref="IGenericSerializable.SerializedLength"/> estimate exists.
/// </summary>
public void PushSingle(IGenericSerializable entity) => _chunks.Enqueue(new Chunk(entity));
/// <summary>
/// Publishes slot ranges covering [0, slotCount) of a directly-iterable persistence.
/// Workers claim ranges like any other chunk, so the per-entity handoff cost disappears
/// and load balancing is unchanged.
/// </summary>
public void PushSlotRanges(ISlotRangeSource source, int slotCount)
{
for (var offset = 0; offset < slotCount; offset += ChunkCapacity)
{
_chunks.Enqueue(new Chunk(source, offset, Math.Min(ChunkCapacity, slotCount - offset)));
}
}
/// <summary>
/// Publishes the partial chunk, if any. Must be called on the producer thread before
/// the workers are told to finish draining, or the tail of the stream is not serialized.
/// </summary>
public void Flush()
{
if (_count > 0)
{
_chunks.Enqueue(new Chunk(_current, _count));
_current = null;
_count = 0;
}
}
internal bool TryTake(out Chunk chunk) => _chunks.TryDequeue(out chunk);
internal void Return(IGenericSerializable[] buffer, int count)
{
// Clear so pooled chunks don't keep entities reachable between saves.
Array.Clear(buffer, 0, count);
_pool.Enqueue(buffer);
}
private IGenericSerializable[] Rent() =>
_pool.TryDequeue(out var buffer) ? buffer : new IGenericSerializable[ChunkCapacity];
}