/************************************************************************* * ModernUO * * Copyright 2019-2026 - ModernUO Development Team * * Email: hi@modernuo.com * * File: SerializationThreadWorker.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 . * *************************************************************************/ using System; using System.Runtime.CompilerServices; using System.Threading; namespace Server; public class SerializationThreadWorker { private const int MinHeapSize = 1024 * 1024; // 1MB private readonly int _index; private readonly Thread _thread; private readonly AutoResetEvent _startEvent; // Main thread tells the thread to start working private readonly AutoResetEvent _stopEvent; // Main thread waits for the worker finish draining private readonly SerializationChunkSource _chunkSource; private readonly int _heapSizeHint; private bool _pause; private bool _exit; private bool _exited; private byte[] _heap; private long _entitiesSerialized; private long _bytesSerialized; public SerializationThreadWorker(int index, SerializationChunkSource chunkSource, int heapSizeHint = 0) : this(index, chunkSource, heapSizeHint, inline: false) { } private SerializationThreadWorker(int index, SerializationChunkSource chunkSource, int heapSizeHint, bool inline) { _index = index; _chunkSource = chunkSource; _heapSizeHint = heapSizeHint; if (!inline) { _startEvent = new AutoResetEvent(false); _stopEvent = new AutoResetEvent(false); _thread = new Thread(Execute); _thread.Start(this); } } /// /// Creates a worker with no thread of its own. The owner drains chunks inline via /// — used by the main thread to join the drain instead of /// idling while the thread workers finish. /// public static SerializationThreadWorker CreateInline(int index, SerializationChunkSource chunkSource, int heapSizeHint = 0) => new(index, chunkSource, heapSizeHint, inline: true); // Stats from the most recent save, for diagnosing load balance. public long EntitiesSerialized => _entitiesSerialized; public long BytesSerialized => _bytesSerialized; public void Wake() { _startEvent.Set(); } public void Sleep() { Volatile.Write(ref _pause, true); _stopEvent.WaitOne(); } public void Exit() { if (_exited) { return; } _exited = true; _exit = true; Wake(); Sleep(); } // Sized from the previous world load so the first save doesn't pay copy-on-grow during the freeze. public void AllocateHeap() => _heap ??= GC.AllocateUninitializedArray(Math.Max(MinHeapSize, _heapSizeHint)); [MethodImpl(MethodImplOptions.AggressiveInlining)] public ReadOnlySpan GetHeap(int start, int length) => _heap.AsSpan(start, length); [MethodImpl(MethodImplOptions.AggressiveInlining)] internal static void Serialize(IGenericSerializable e, BufferWriter writer, byte threadIndex) { e.SerializedThread = threadIndex; var start = e.SerializedPosition = (int)writer.Position; e.Serialize(writer); e.SerializedLength = (int)(writer.Position - start); } private static long ProcessChunk( in SerializationChunkSource.Chunk chunk, SerializationChunkSource chunkSource, BufferWriter writer, byte threadIndex ) { if (chunk.Single != null) { Serialize(chunk.Single, writer, threadIndex); return 1; } if (chunk.Source != null) { return chunk.Source.SerializeRange(writer, threadIndex, chunk.Offset, chunk.Count); } var buffer = chunk.Buffer; var count = chunk.Count; for (var i = 0; i < count; i++) { Serialize(buffer[i], writer, threadIndex); } chunkSource.Return(buffer, count); return count; } /// /// Drains chunks on the calling thread until the queue is empty, then returns. /// Only valid on inline workers; the main thread calls this after publishing all work /// so it contributes drain throughput instead of idling. /// public void DrainInline() { var writer = new BufferWriter(_heap, true, World.SerializedTypes); var threadIndex = (byte)_index; var entities = 0L; while (_chunkSource.TryTake(out var chunk)) { entities += ProcessChunk(in chunk, _chunkSource, writer, threadIndex); } _heap = writer.Buffer; _entitiesSerialized = entities; _bytesSerialized = writer.Position; writer.Close(); } private static void Execute(object obj) { var worker = (SerializationThreadWorker)obj; var threadIndex = (byte)worker._index; var chunkSource = worker._chunkSource; var serializedTypes = World.SerializedTypes; while (worker._startEvent.WaitOne()) { var writer = new BufferWriter(worker._heap, true, serializedTypes); var entities = 0L; var spinner = new SpinWait(); while (true) { var pauseRequested = Volatile.Read(ref worker._pause); if (chunkSource.TryTake(out var chunk)) { spinner.Reset(); entities += ProcessChunk(in chunk, chunkSource, writer, threadIndex); } else if (pauseRequested) // Break when finished { break; } else { // Idle backoff instead of hammering the queue head while the producer works. // sleep1Threshold: -1 keeps escalation at Yield/Sleep(0) and never Sleep(1), // avoiding timer-resolution stalls at the end of the drain. spinner.SpinOnce(-1); } } worker._heap = writer.Buffer; worker._entitiesSerialized = entities; worker._bytesSerialized = writer.Position; writer.Close(); worker._stopEvent.Set(); // Allow the main thread to continue now that we are finished worker._pause = false; if (Core.Closing || worker._exit) { return; } } } }