ModernUO/Projects/Server.Tests/Tests/Serialization/SerializationChunkSourceTests.cs
Kamron Batman 230c39851f
perf(saves): chunked worker handoff, LPT blob scheduling, heap pre-sizing
Replaces the per-entity round-robin ConcurrentQueue handoff between the game
loop and the serialization workers with pooled 4096-entity chunks published to
a single shared queue. The producer's per-entity cost drops from a synchronized
enqueue to a plain array store, and workers pulling whole chunks load-balance
dynamically: a worker busy with a thick entity simply takes fewer chunks.

Scheduling changes so indivisible multi-megabyte payloads no longer extend the
freeze tail:
- Persistence.SerializeAll pushes systems largest-first (LPT), using the
  previous save's SerializedLength (or the loaded file size on first save).
- Persistence self-payloads and entities whose previous size exceeds 1MB are
  published as dedicated single-entity chunks so they spread across workers
  instead of riding inside one shared chunk.
- Entity SerializedLength is stamped from the index at load so estimates exist
  on the first save after boot.

Worker heaps are pre-sized from the loaded save's .bin sizes (25% headroom) so
the first save doesn't pay copy-on-grow inside the freeze, and the drain loop
uses SpinWait backoff (never Sleep(1)) instead of hammering the queue head
while the producer works. Snapshot writing uses a 1MB FileStream buffer, and
per-worker entity/byte counts are logged at Debug for balance diagnostics.

Measured on a 24-core machine with a synthetic 10M-entity, 1.7GB world
(64/64/32MB system payloads, 24x 2MB thick entities) through the real
pipeline classes: freeze window 740ms -> 122-160ms steady state (~5-6x),
first save 490ms -> 330ms, steady-state allocations converge to zero, and
worker byte loads converge (previous max/min spread ~2x -> ~1.15x for the
small-entity stream).

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

222 lines
6.2 KiB
C#

using System.Collections.Generic;
using Xunit;
namespace Server.Tests;
public class SerializationChunkSourceTests
{
private class TestEntity : IGenericSerializable
{
public byte SerializedThread { get; set; }
public int SerializedPosition { get; set; }
public int SerializedLength { get; set; }
public int PayloadSize { get; init; } = 16;
public void Serialize(IGenericWriter writer)
{
for (var i = 0; i < PayloadSize; i++)
{
writer.Write((byte)(i & 0xFF));
}
}
}
private static List<TestEntity> Drain(SerializationChunkSource source)
{
var drained = new List<TestEntity>();
while (source.TryTake(out var chunk))
{
if (chunk.Single != null)
{
drained.Add((TestEntity)chunk.Single);
}
else
{
for (var i = 0; i < chunk.Count; i++)
{
drained.Add((TestEntity)chunk.Buffer[i]);
}
source.Return(chunk.Buffer, chunk.Count);
}
}
return drained;
}
[Fact]
public void PartialChunkIsNotVisibleUntilFlush()
{
var source = new SerializationChunkSource();
var entities = new List<TestEntity>();
for (var i = 0; i < 100; i++)
{
var e = new TestEntity();
entities.Add(e);
source.Push(e);
}
Assert.False(source.TryTake(out _));
source.Flush();
var drained = Drain(source);
Assert.Equal(entities, drained);
// Flush again should publish nothing
source.Flush();
Assert.False(source.TryTake(out _));
}
[Fact]
public void FullChunkPublishesWithoutFlush()
{
var source = new SerializationChunkSource();
for (var i = 0; i < 4096; i++)
{
source.Push(new TestEntity());
}
Assert.True(source.TryTake(out var chunk));
Assert.Null(chunk.Single);
Assert.Equal(4096, chunk.Count);
source.Return(chunk.Buffer, chunk.Count);
// Nothing partial left behind
source.Flush();
Assert.False(source.TryTake(out _));
}
[Fact]
public void PushSingleDoesNotDisturbPartialChunk()
{
var source = new SerializationChunkSource();
var small1 = new TestEntity();
var heavy = new TestEntity { SerializedLength = 2 * 1024 * 1024 };
var small2 = new TestEntity();
source.Push(small1);
source.PushSingle(heavy); // published immediately as a single, ahead of the partial chunk
source.Push(small2);
Assert.True(source.TryTake(out var chunk));
Assert.Same(heavy, chunk.Single);
Assert.Equal(1, chunk.Count);
Assert.False(source.TryTake(out _));
source.Flush();
var drained = Drain(source);
Assert.Equal([small1, small2], drained);
}
[Fact]
public void ReturnedBuffersAreClearedAndReused()
{
var source = new SerializationChunkSource();
for (var i = 0; i < 4096; i++)
{
source.Push(new TestEntity());
}
Assert.True(source.TryTake(out var chunk));
var buffer = chunk.Buffer;
source.Return(buffer, chunk.Count);
Assert.All(buffer, Assert.Null);
// Next fill rents the pooled buffer instead of allocating
source.Push(new TestEntity());
source.Flush();
Assert.True(source.TryTake(out var reused));
Assert.Same(buffer, reused.Buffer);
}
[Fact]
public void WorkersDrainAllEntitiesAndStampPositions()
{
var source = new SerializationChunkSource();
var workers = new SerializationThreadWorker[2];
for (var i = 0; i < workers.Length; i++)
{
workers[i] = new SerializationThreadWorker(i, source);
workers[i].AllocateHeap();
}
try
{
var entities = new List<TestEntity>();
for (var i = 0; i < 10_000; i++)
{
entities.Add(new TestEntity { PayloadSize = 16 + i % 64 });
}
// A heavy entity mid-stream gets a dedicated chunk
var heavy = new TestEntity { PayloadSize = 512 * 1024, SerializedLength = 4 * 1024 * 1024 };
entities.Insert(5000, heavy);
foreach (var worker in workers)
{
worker.Wake();
}
// Mirrors GenericEntityPersistence.Serialize: heavy check at the call site
foreach (var e in entities)
{
if (e.SerializedLength > SerializationChunkSource.HeavyEntityThreshold)
{
source.PushSingle(e);
}
else
{
source.Push(e);
}
}
// Mirrors World.PauseSerializationThreads
source.Flush();
foreach (var worker in workers)
{
worker.Sleep();
}
long totalEntities = 0;
long totalBytes = 0;
foreach (var worker in workers)
{
totalEntities += worker.EntitiesSerialized;
totalBytes += worker.BytesSerialized;
}
Assert.Equal(entities.Count, totalEntities);
long expectedBytes = 0;
foreach (var e in entities)
{
expectedBytes += e.PayloadSize;
// Every entity serialized exactly once with a consistent span on its worker's heap
Assert.Equal(e.PayloadSize, e.SerializedLength);
Assert.InRange(e.SerializedThread, (byte)0, (byte)(workers.Length - 1));
var heap = workers[e.SerializedThread].GetHeap(e.SerializedPosition, e.SerializedLength);
Assert.Equal(0, heap[0]);
Assert.Equal((e.SerializedLength - 1) & 0xFF, heap[^1]);
}
Assert.Equal(expectedBytes, totalBytes);
}
finally
{
foreach (var worker in workers)
{
worker.Exit();
}
}
}
}