From 003491472f08520bf9905abe54d715b894545902 Mon Sep 17 00:00:00 2001 From: Kamron Batman <3953314+kamronbatman@users.noreply.github.com> Date: Sun, 6 Sep 2026 13:46:14 -0700 Subject: [PATCH] =?UTF-8?q?fix:=20world=20save=20failure=20handling=20?= =?UTF-8?q?=E2=80=94=20no=20partial=20snapshots,=20no=20worker=20crash,=20?= =?UTF-8?q?staged=20publication=20(#2612)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three pre-existing world-save failure-handling bugs, found while reviewing the delta-saves engine (#2610), split out so they land and get adopted on their own before the larger change. #2610 will be rebased on top once this merges. Tracked in project Delta Saves (#7). ## 1. A snapshot missing a segment was published `GenericEntityPersistence.WriteSnapshot` logged a segment (or self-payload) write error and carried on. The save was published without those records, and every entity in the failed segment was deleted at the next load. The error now propagates: `WriteFiles` never moves a partial snapshot over `Saves/`, the previous save stays authoritative, and staff are told. ## 2. A serializer exception killed the process, or published a partial freeze An entity whose `Serialize` threw on a worker thread was an unhandled exception on that thread, which terminates the process. On the inline (main-thread) drain it was caught, but the snapshot was still queued for writing with whatever the workers had produced. Workers now record the first exception and keep draining so the queue empties and the wake/pause handshake completes; after every worker paused, `Snapshot` treats a recorded error (or a failure of the drain itself) as a failed save: nothing is written, the previous save stays, staff are told, and the world resumes. A `WorldSave` handler throwing is logged but does not invalidate the serialized snapshot. ## 3. Publication was not transactional Publishing moved the previous `Saves/` away (`AutoArchive` in `WorldSavePostSnapshot`) and only then moved the new files in. A subscriber throwing after the archive, or the final move failing, left nothing at `Saves/` until the next save; a crash in that window lost the newest save at the next boot. The snapshot now moves to `Saves.next` as soon as its files are complete; only then do subscribers archive the previous save (same event, same `OldSavePath`, so `AutoArchive` and shard subscribers are unchanged), and the staged directory is renamed into place, which is atomic on one volume (file-by-file move across volumes). A staged directory is by construction a complete save newer than `Saves/`, so an interrupted publish is finished at the next boot (before load) or before the next save: whatever sits at `Saves/` is set aside as `Saves.previous-` (never deleted) and the staged save is published, with a warning naming the directory to archive or delete by hand. ## Tests Server.Tests 860 / UOContent.Tests green. New: a throwing serializer is recorded on the worker and the next drain starts clean; a segment that cannot be indexed fails `WriteSnapshot` instead of dropping records; staged-save recovery with and without an existing `Saves/`, and as a no-op. --- .../GenericEntityPersistenceRoundTripTests.cs | 150 +++++++++++++++++- .../Serialization/StagedSavePublishTests.cs | 83 ++++++++++ .../Serialization/GenericEntityPersistence.cs | 7 +- .../SerializationThreadWorker.cs | 27 ++++ Projects/Server/World/World.cs | 144 +++++++++++++++-- 5 files changed, 390 insertions(+), 21 deletions(-) create mode 100644 Projects/Server.Tests/Tests/Serialization/StagedSavePublishTests.cs diff --git a/Projects/Server.Tests/Tests/Serialization/GenericEntityPersistenceRoundTripTests.cs b/Projects/Server.Tests/Tests/Serialization/GenericEntityPersistenceRoundTripTests.cs index f2114a427..48ba16e55 100644 --- a/Projects/Server.Tests/Tests/Serialization/GenericEntityPersistenceRoundTripTests.cs +++ b/Projects/Server.Tests/Tests/Serialization/GenericEntityPersistenceRoundTripTests.cs @@ -20,7 +20,7 @@ internal class RoundTripEntity : ISerializable { } - public void Serialize(IGenericWriter writer) + public virtual void Serialize(IGenericWriter writer) { writer.Write(Value); writer.Write(Name); @@ -33,6 +33,15 @@ internal class RoundTripEntity : ISerializable } } +internal class ThrowingRoundTripEntity : RoundTripEntity +{ + public ThrowingRoundTripEntity(Serial serial) : base(serial) + { + } + + public override void Serialize(IGenericWriter writer) => throw new InvalidOperationException("broken serializer"); +} + [Collection("Sequential Server Tests")] public class GenericEntityPersistenceRoundTripTests { @@ -161,4 +170,143 @@ public class GenericEntityPersistenceRoundTripTests Directory.Delete(dir, true); } } + + /// + /// A serializer throwing on a worker used to be an unhandled exception on that thread. + /// The worker records it, finishes the drain so the handshake completes, and the loop + /// fails the save. + /// + [Fact] + public void SerializerException_IsRecordedOnTheWorker_AndTheDrainCompletes() + { + 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(); + } + + var persistence = new RoundTripPersistence(2002); + + try + { + for (var i = 1; i <= 100; i++) + { + var serial = (Serial)(uint)i; + persistence.EntitiesBySerial[serial] = i == 50 + ? new ThrowingRoundTripEntity(serial) + : new RoundTripEntity(serial) { Value = i, Name = $"entity-{i}" }; + } + + foreach (var worker in workers) + { + worker.Wake(); + } + + source.SetOwner(persistence); + Assert.True(persistence.TrySnapshotEntries(out var slotCount)); + source.PushSlotRanges(persistence, slotCount); + source.Flush(); + + foreach (var worker in workers) + { + worker.Sleep(); + } + + Exception error = null; + foreach (var worker in workers) + { + error ??= worker.Error; + } + + Assert.IsType(error); + persistence.PostWorldSave(); + + // The next drain starts clean. + foreach (var worker in workers) + { + worker.Wake(); + } + + foreach (var worker in workers) + { + worker.Sleep(); + Assert.Null(worker.Error); + } + } + finally + { + persistence.Unregister(); + + foreach (var worker in workers) + { + worker.Exit(); + } + } + } + + /// + /// A segment that cannot be written used to be logged and dropped, publishing a save + /// without those entities (the loader then deletes them). It now fails the save. + /// + [Fact] + public void WriteSnapshot_FailsTheSave_InsteadOfDroppingASegment() + { + 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(); + } + + var previousWorkers = World._threadWorkers; + World._threadWorkers = workers; + + var persistence = new RoundTripPersistence(2003); + var dir = Path.Combine(Path.GetTempPath(), $"muo-segmentfail-{Guid.NewGuid():N}"); + Directory.CreateDirectory(dir); + + try + { + for (var i = 1; i <= 100; i++) + { + var serial = (Serial)(uint)i; + persistence.EntitiesBySerial[serial] = new RoundTripEntity(serial) { Value = i, Name = $"entity-{i}" }; + } + + // Deliberately not registered: the writer cannot index the type. + + foreach (var worker in workers) + { + worker.Wake(); + } + + source.SetOwner(persistence); + Assert.True(persistence.TrySnapshotEntries(out var slotCount)); + source.PushSlotRanges(persistence, slotCount); + source.Flush(); + + foreach (var worker in workers) + { + worker.Sleep(); + } + + Assert.Throws(() => persistence.WriteSnapshot(dir)); + persistence.PostWorldSave(); + } + finally + { + persistence.Unregister(); + + foreach (var worker in workers) + { + worker.Exit(); + } + + World._threadWorkers = previousWorkers; + Directory.Delete(dir, true); + } + } } diff --git a/Projects/Server.Tests/Tests/Serialization/StagedSavePublishTests.cs b/Projects/Server.Tests/Tests/Serialization/StagedSavePublishTests.cs new file mode 100644 index 000000000..2b1b3aff0 --- /dev/null +++ b/Projects/Server.Tests/Tests/Serialization/StagedSavePublishTests.cs @@ -0,0 +1,83 @@ +using System; +using System.IO; +using Xunit; + +namespace Server.Tests; + +/// +/// The publish protocol: a complete snapshot is staged next to Saves/ before the previous +/// save is touched, and an interrupted publish is finished at the next boot or save. +/// +[Collection("Sequential Server Tests")] +public class StagedSavePublishTests : IDisposable +{ + private readonly string _root = Path.Combine(Path.GetTempPath(), $"muo-staged-{Guid.NewGuid():N}"); + private readonly string _previousSavePath; + + public StagedSavePublishTests() + { + Directory.CreateDirectory(_root); + _previousSavePath = World.SavePath; + World.SetSavePathForTest(Path.Combine(_root, "Saves")); + } + + public void Dispose() + { + World.SetSavePathForTest(_previousSavePath); + + try + { + Directory.Delete(_root, true); + } + catch + { + // best effort + } + } + + private static void WriteMarker(string dir, string name) + { + Directory.CreateDirectory(dir); + File.WriteAllText(Path.Combine(dir, "marker.txt"), name); + } + + private static string ReadMarker(string dir) => File.ReadAllText(Path.Combine(dir, "marker.txt")); + + [Fact] + public void NothingStaged_RecoveryIsANoOp() + { + WriteMarker(World.SavePath, "current"); + + World.RecoverStagedSave(); + + Assert.Equal("current", ReadMarker(World.SavePath)); + Assert.Single(Directory.GetDirectories(_root)); + } + + [Fact] + public void StagedSave_ReplacesSaves_AndKeepsThePreviousOne() + { + WriteMarker(World.SavePath, "old"); + WriteMarker(World.StagedSavePath, "new"); + + World.RecoverStagedSave(); + + Assert.Equal("new", ReadMarker(World.SavePath)); + Assert.False(Directory.Exists(World.StagedSavePath)); + + var aside = Array.FindAll(Directory.GetDirectories(_root), d => d.Contains(".previous-")); + Assert.Single(aside); + Assert.Equal("old", ReadMarker(aside[0])); + } + + [Fact] + public void StagedSave_WithNoSaves_IsPublished() + { + WriteMarker(World.StagedSavePath, "new"); + + World.RecoverStagedSave(); + + Assert.Equal("new", ReadMarker(World.SavePath)); + Assert.Single(Directory.GetDirectories(_root)); + } +} diff --git a/Projects/Server/Serialization/GenericEntityPersistence.cs b/Projects/Server/Serialization/GenericEntityPersistence.cs index 9a1490504..70e9c6aab 100644 --- a/Projects/Server/Serialization/GenericEntityPersistence.cs +++ b/Projects/Server/Serialization/GenericEntityPersistence.cs @@ -158,6 +158,7 @@ public class GenericEntityPersistence : GenericPersistence, IGenericEntityPer _selfPosition, _selfLength ); + throw; } binPosition += _selfLength; @@ -205,6 +206,7 @@ public class GenericEntityPersistence : GenericPersistence, IGenericEntityPer } catch (Exception error) { + // Never publish a partial snapshot: entities missing from the idx are deleted on load. logger.Error( error, "Error writing segment: (Thread: {Thread} - {Start}, {Records} records)", @@ -212,6 +214,7 @@ public class GenericEntityPersistence : GenericPersistence, IGenericEntityPer segment.HeapStart, segment.RecordCount ); + throw; } } } @@ -318,9 +321,7 @@ public class GenericEntityPersistence : GenericPersistence, IGenericEntityPer private ushort GetTypeIndex(T entity) { // Every path into EntitiesBySerial registers the type first, so this cannot fire. - // If it ever does, the segment-level catch in WriteSnapshot logs it and moves on — - // the failed segment's records are dropped from the idx while binPosition rewinds, - // so treat any occurrence as a serious bug in an insertion path, not a bad entity. + // If it does, the save fails; treat it as a bug in an insertion path, not a bad entity. if (!_typeIndexes.TryGetValue(entity.GetType(), out var typeIndex)) { throw new InvalidOperationException( diff --git a/Projects/Server/Serialization/SerializationThreadWorker.cs b/Projects/Server/Serialization/SerializationThreadWorker.cs index 5a9b5ed1c..74973a80c 100644 --- a/Projects/Server/Serialization/SerializationThreadWorker.cs +++ b/Projects/Server/Serialization/SerializationThreadWorker.cs @@ -78,6 +78,12 @@ public class SerializationThreadWorker internal List Lengths => _lengths; internal List BufferEntities => _bufferEntities; + /// + /// First serializer exception during the drain. The drain continues so the handshake + /// completes; the loop fails the save once every worker has paused. + /// + public Exception Error { get; private set; } + /// /// Releases the write logs after the snapshot is written so serialized entity /// references don't linger between saves. Capacity is retained: the logs regrow to @@ -154,6 +160,25 @@ public class SerializationThreadWorker public ReadOnlySpan GetHeap(int start, int length) => _heap.AsSpan(start, length); private long ProcessChunk(in SerializationChunkSource.Chunk chunk, BufferWriter writer) + { + try + { + return ProcessChunkCore(in chunk, writer); + } + catch (Exception ex) + { + Error ??= ex; + + if (chunk.Buffer != null) + { + _chunkSource.Return(chunk.Buffer, chunk.Count); + } + + return 0; + } + } + + private long ProcessChunkCore(in SerializationChunkSource.Chunk chunk, BufferWriter writer) { if (chunk.Single != null) { @@ -213,6 +238,7 @@ public class SerializationThreadWorker public void DrainInline() { ReleaseWriteLogs(); + Error = null; var writer = new BufferWriter(_heap, true); var entities = 0L; @@ -238,6 +264,7 @@ public class SerializationThreadWorker while (worker._startEvent.WaitOne()) { worker.ReleaseWriteLogs(); + worker.Error = null; var writer = new BufferWriter(worker._heap, true); var entities = 0L; diff --git a/Projects/Server/World/World.cs b/Projects/Server/World/World.cs index c00fc85f6..7ce1b52f2 100644 --- a/Projects/Server/World/World.cs +++ b/Projects/Server/World/World.cs @@ -122,6 +122,8 @@ public static class World UseMultiThreadedSaves = ServerConfiguration.GetOrUpdateSetting("world.useMultithreadedSaves", true); } + internal static void SetSavePathForTest(string path) => SavePath = path; + [MethodImpl(MethodImplOptions.AggressiveInlining)] public static void WaitForWriteCompletion() => _diskWriteHandle.WaitOne(); @@ -195,6 +197,8 @@ public static class World logger.Information("Loading world"); var watch = Stopwatch.StartNew(); + RecoverStagedSave(); + Persistence.Load(SavePath); EventSink.InvokeWorldLoad(); @@ -271,6 +275,8 @@ public static class World { try { + RecoverStagedSave(); + // Allocate the heaps for the GC foreach (var worker in _threadWorkers) { @@ -322,22 +328,49 @@ public static class World } Persistence.SerializeAll(); - PauseSerializationThreads(); - LogWorkerBalance(); - - EventSink.InvokeWorldSave(); } catch (Exception ex) { exception = ex; } - WorldState = WorldState.WritingSave; - ThreadPool.QueueUserWorkItem(WriteFiles, snapshotPath); + // Always join the workers; any serializer exception fails the save. + try + { + PauseSerializationThreads(); + } + catch (Exception ex) + { + exception ??= ex; + } + + for (var i = 0; i < _threadWorkers.Length; i++) + { + exception ??= _threadWorkers[i].Error; + } + + if (exception == null) + { + LogWorkerBalance(); + + try + { + EventSink.InvokeWorldSave(); + } + catch (Exception ex) + { + logger.Error(ex, "A WorldSave handler failed"); + Persistence.TraceException(ex); + } + } + watch.Stop(); if (exception == null) { + WorldState = WorldState.WritingSave; + ThreadPool.QueueUserWorkItem(WriteFiles, snapshotPath); + var duration = watch.Elapsed.TotalSeconds; logger.Information("Saving world {Status} ({Duration:F2} seconds)", "done", duration); @@ -349,6 +382,9 @@ public static class World Persistence.TraceException(exception); BroadcastStaff(0x35, true, "World save failed! Check the logs!"); + + _diskWriteHandle.Set(); + FinishWorldSave(); } } @@ -361,17 +397,7 @@ public static class World logger.Information("Writing world save snapshot"); Persistence.WriteSnapshotAll(snapshotPath); - - try - { - EventSink.InvokeWorldSavePostSnapshot(SavePath, snapshotPath); - PathUtility.MoveDirectoryContents(snapshotPath, SavePath); - Directory.SetLastWriteTimeUtc(SavePath, Core.Now); - } - catch (Exception ex) - { - Persistence.TraceException(ex); - } + PublishSnapshot(snapshotPath); watch.Stop(); logger.Information("Writing world save snapshot {Status} ({Duration:F2} seconds)", "done", watch.Elapsed.TotalSeconds); @@ -388,6 +414,90 @@ public static class World Core.LoopContext.Post(FinishWorldSave); } + /// + /// A complete snapshot is staged here before the previous save is touched; a staged + /// directory is always a complete save newer than . + /// + internal static string StagedSavePath => SavePath + ".next"; + + // Stage, let subscribers archive the previous save, then rename the staged save into place. + private static void PublishSnapshot(string snapshotPath) + { + var staging = StagedSavePath; + + if (Directory.Exists(staging)) + { + SetAside(staging, "unpublished"); + } + + MoveDirectory(snapshotPath, staging); + PublishStagedSave(archive: true); + } + + private static void PublishStagedSave(bool archive) + { + var staging = StagedSavePath; + + if (archive) + { + EventSink.InvokeWorldSavePostSnapshot(SavePath, staging); + } + + if (Directory.Exists(SavePath)) + { + SetAside(SavePath, "previous"); + } + + MoveDirectory(staging, SavePath); + Directory.SetLastWriteTimeUtc(SavePath, Core.Now); + } + + /// + /// Finishes an interrupted publish. Runs at boot (before load) and before every save; + /// whatever is at Saves/ is set aside, never deleted. + /// + internal static void RecoverStagedSave() + { + var staging = StagedSavePath; + + if (!Directory.Exists(staging)) + { + return; + } + + logger.Warning( + "A complete world save was staged at {Staging} but never published; publishing it now.", + staging + ); + + PublishStagedSave(archive: false); + } + + private static void SetAside(string path, string reason) + { + var aside = $"{path}.{reason}-{Core.Now:yyyy-MM-dd-HH-mm-ss-fff}"; + MoveDirectory(path, aside); + logger.Warning("Set aside {Path} as {Aside}; delete or archive it by hand.", path, aside); + } + + // Atomic rename on one volume, file-by-file move otherwise. + private static void MoveDirectory(string source, string destination) + { + try + { + Directory.Move(source, destination); + } + catch (IOException) + { + if (Directory.Exists(destination)) + { + throw; + } + + PathUtility.MoveDirectoryContents(source, destination); + } + } + private static void FinishWorldSave() { WorldState = WorldState.Running;