From 57a3894234c03305d9d40a54d483e4fb79185e82 Mon Sep 17 00:00:00 2001 From: Pravus Date: Fri, 12 Jun 2026 03:00:29 +0200 Subject: [PATCH] opti: CRDT bridge concurrency improvements - outgoing message serialization moved outside the producers' lock via list double-buffering - WorldSyncCommandBuffer instance reused across batches (Renew) instead of per-tick allocation - idle fast path: empty sync buffers skip the world mutex acquisition and command buffer playback Co-Authored-By: Claude Fable 5 --- .../EngineAPIImplementation.cs | 19 +++-- .../OutgoingCRDTMessagesProvider.cs | 84 ++++++++++++------- .../WorldSynchronizer/CrdtEcsSynchronizer.cs | 40 +++++++-- .../ICRDTWorldSynchronizer.cs | 6 ++ .../IWorldSyncCommandBuffer.cs | 6 ++ .../Tests/CrdtWorldSynchronizerShould.cs | 29 +++++++ .../Tests/WorldSyncCommandBufferShould.cs | 51 +++++++++++ .../WorldSyncCommandBuffer.cs | 37 +++++++- 8 files changed, 226 insertions(+), 46 deletions(-) diff --git a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/JsModulesImplementation/EngineAPIImplementation.cs b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/JsModulesImplementation/EngineAPIImplementation.cs index 110ec1f41e0..2def62ef812 100644 --- a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/JsModulesImplementation/EngineAPIImplementation.cs +++ b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/JsModulesImplementation/EngineAPIImplementation.cs @@ -257,13 +257,22 @@ private void ApplySyncCommandBuffer(IWorldSyncCommandBuffer worldSyncBuffer) { try { - using MultiThreadSync.Scope mutex = multiThreadSync.GetScope(syncOwner); + if (worldSyncBuffer.IsEmpty) + { + // Nothing to apply to the World: skip acquiring the world mutex entirely, + // otherwise an idle scene blocks its runtime thread on the main thread every tick + crdtWorldSynchronizer.ReleaseSyncCommandBuffer(worldSyncBuffer); + } + else + { + using MultiThreadSync.Scope mutex = multiThreadSync.GetScope(syncOwner); - applyBufferSampler.Begin(); + applyBufferSampler.Begin(); - // Apply changes to the ECS World on the main thread - crdtWorldSynchronizer.ApplySyncCommandBuffer(worldSyncBuffer); - applyBufferSampler.End(); + // Apply changes to the ECS World on the main thread + crdtWorldSynchronizer.ApplySyncCommandBuffer(worldSyncBuffer); + applyBufferSampler.End(); + } // Allow system for which throttling is enabled to process once // If the scene is updated more frequently than Unity Loop the gate will be effectively open all the time diff --git a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/OutgoingMessages/OutgoingCRDTMessagesProvider.cs b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/OutgoingMessages/OutgoingCRDTMessagesProvider.cs index 59cb277ccd9..3fc3e2003d6 100644 --- a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/OutgoingMessages/OutgoingCRDTMessagesProvider.cs +++ b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/OutgoingMessages/OutgoingCRDTMessagesProvider.cs @@ -22,12 +22,24 @@ public class OutgoingCRDTMessagesProvider : IOutgoingCRDTMessagesProvider new (64, PoolConstants.SCENES_COUNT); internal readonly Dictionary lwwMessageIndices = INDICES_SHARED_POOL.Get(); - internal readonly List messages = MESSAGES_SHARED_POOL.Get(); + internal List messages = MESSAGES_SHARED_POOL.Get(); private readonly ISDKComponentsRegistry componentsRegistry; private readonly ICRDTProtocol crdtProtocol; private readonly ICRDTMemoryAllocator memoryAllocator; + /// + /// Guards and : a dedicated object is required + /// because is swapped with the spare list on serialization + /// + private readonly object writeLock = new (); + + /// + /// The second buffer is swapped with so the per-message serialization + /// runs outside the lock and does not stall the main-thread systems writing new messages + /// + private List spareMessages = MESSAGES_SHARED_POOL.Get(); + public OutgoingCRDTMessagesProvider(ISDKComponentsRegistry componentsRegistry, ICRDTProtocol crdtProtocol, ICRDTMemoryAllocator memoryAllocator) { this.componentsRegistry = componentsRegistry; @@ -39,6 +51,7 @@ public void Dispose() { INDICES_SHARED_POOL.Release(lwwMessageIndices); MESSAGES_SHARED_POOL.Release(messages); + MESSAGES_SHARED_POOL.Release(spareMessages); } public void AddDeleteMessage(CRDTEntity entity) where TMessage: class, IMessage @@ -72,7 +85,7 @@ public void AddPutMessage(TMessage message, CRDTEntity entity) where T private void AddLwwMessage(CRDTEntity entity, SDKComponentBridge componentBridge, in PendingMessage newMessage) { - lock (messages) + lock (writeLock) { var key = new OutgoingMessageKey(entity, componentBridge.Id); @@ -104,7 +117,7 @@ public TMessage AppendMessage(Action prepareMe var message = (TMessage)componentBridge.Pool.Rent(); prepareMessage(message, data); - lock (messages) { messages.Add(new PendingMessage(message, componentBridge, entity, CRDTMessageType.APPEND_COMPONENT, timestamp)); } + lock (writeLock) { messages.Add(new PendingMessage(message, componentBridge, entity, CRDTMessageType.APPEND_COMPONENT, timestamp)); } return message; } @@ -126,41 +139,50 @@ public OutgoingCRDTMessagesSyncBlock GetSerializationSyncBlock(Action processedMessages = OutgoingCRDTMessagesSyncBlock.MESSAGES_SHARED_POOL.Get(); - // While we do it we must synchronize - lock (messages) + List toSerialize; + + // Detach the pending batch under the lock: producers keep writing into the spare list + // while serialization (the expensive part) runs lock-free below + lock (writeLock) { - for (var i = 0; i < messages.Count; i++) - { - PendingMessage pendingMessage = messages[i]; + toSerialize = messages; + messages = spareMessages; + spareMessages = toSerialize; + lwwMessageIndices.Clear(); + } + + // Safe without the lock: this method is invoked from the scene runtime thread only (single consumer) + // so nothing else touches the detached list + for (var i = 0; i < toSerialize.Count; i++) + { + PendingMessage pendingMessage = toSerialize[i]; - actOnPendingMessage?.Invoke(pendingMessage); + actOnPendingMessage?.Invoke(pendingMessage); - IMemoryOwner memory; + IMemoryOwner memory; - switch (pendingMessage.MessageType) - { - case CRDTMessageType.PUT_COMPONENT: - memory = memoryAllocator.GetMemoryBuffer(pendingMessage.Message.CalculateSize()); - pendingMessage.Bridge.Serializer.SerializeInto(pendingMessage.Message, memory.Memory.Span); - pendingMessage.Bridge.Pool.Release(pendingMessage.Message); - processedMessages.Add(crdtProtocol.CreatePutMessage(pendingMessage.Entity, pendingMessage.Bridge.Id, memory)); - break; - case CRDTMessageType.APPEND_COMPONENT: - memory = memoryAllocator.GetMemoryBuffer(pendingMessage.Message.CalculateSize()); - pendingMessage.Bridge.Serializer.SerializeInto(pendingMessage.Message, memory.Memory.Span); - pendingMessage.Bridge.Pool.Release(pendingMessage.Message); - processedMessages.Add(crdtProtocol.CreateAppendMessage(pendingMessage.Entity, pendingMessage.Bridge.Id, pendingMessage.Timestamp, memory)); - break; - case CRDTMessageType.DELETE_COMPONENT: - processedMessages.Add(crdtProtocol.CreateDeleteMessage(pendingMessage.Entity, pendingMessage.Bridge.Id)); - break; - } + switch (pendingMessage.MessageType) + { + case CRDTMessageType.PUT_COMPONENT: + memory = memoryAllocator.GetMemoryBuffer(pendingMessage.Message.CalculateSize()); + pendingMessage.Bridge.Serializer.SerializeInto(pendingMessage.Message, memory.Memory.Span); + pendingMessage.Bridge.Pool.Release(pendingMessage.Message); + processedMessages.Add(crdtProtocol.CreatePutMessage(pendingMessage.Entity, pendingMessage.Bridge.Id, memory)); + break; + case CRDTMessageType.APPEND_COMPONENT: + memory = memoryAllocator.GetMemoryBuffer(pendingMessage.Message.CalculateSize()); + pendingMessage.Bridge.Serializer.SerializeInto(pendingMessage.Message, memory.Memory.Span); + pendingMessage.Bridge.Pool.Release(pendingMessage.Message); + processedMessages.Add(crdtProtocol.CreateAppendMessage(pendingMessage.Entity, pendingMessage.Bridge.Id, pendingMessage.Timestamp, memory)); + break; + case CRDTMessageType.DELETE_COMPONENT: + processedMessages.Add(crdtProtocol.CreateDeleteMessage(pendingMessage.Entity, pendingMessage.Bridge.Id)); + break; } - - messages.Clear(); - lwwMessageIndices.Clear(); } + toSerialize.Clear(); + // This list will be released on block.Dispose() return new OutgoingCRDTMessagesSyncBlock(processedMessages); } diff --git a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/CrdtEcsSynchronizer.cs b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/CrdtEcsSynchronizer.cs index acb6b50fbab..f565bf0cabc 100644 --- a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/CrdtEcsSynchronizer.cs +++ b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/CrdtEcsSynchronizer.cs @@ -28,6 +28,10 @@ public class CRDTWorldSynchronizer : ICRDTWorldSynchronizer // and it is not guaranteed as we use thread pools (in the most cases different threads are used for getting and applying command buffers) private readonly DCLSemaphoreSlim semaphore = new (1, 1); + // Single-slot reuse: the semaphore guarantees one outstanding buffer at a time + // so the same instance can be renewed instead of allocating a new one per batch (per scene tick) + private WorldSyncCommandBuffer reusableSyncBuffer; + private bool disposed; public IReadOnlyDictionary EntitiesMap => entitiesMap; @@ -68,6 +72,15 @@ public IWorldSyncCommandBuffer GetSyncCommandBuffer() throw new TimeoutException("Rent Wait Timeout: Couldn't rent command buffer"); #endif + WorldSyncCommandBuffer syncBuffer = reusableSyncBuffer; + + if (syncBuffer != null) + { + reusableSyncBuffer = null; + syncBuffer.Renew(); + return syncBuffer; + } + return new WorldSyncCommandBuffer(sdkComponentsRegistry, entityFactory, collectionsPool); } @@ -81,15 +94,28 @@ public void ApplySyncCommandBuffer(IWorldSyncCommandBuffer syncCommandBuffer) else syncCommandBuffer.Apply(world, reusableCommandBuffer, entitiesMap); } - finally - { + finally { FinalizeRent(syncCommandBuffer); } + } + + public void ReleaseSyncCommandBuffer(IWorldSyncCommandBuffer syncCommandBuffer) + { + try { syncCommandBuffer.Dispose(); } + finally { FinalizeRent(syncCommandBuffer); } + } + + private void FinalizeRent(IWorldSyncCommandBuffer syncCommandBuffer) + { + // Keep the disposed instance for the next rent; when the scene is disposed let it be GCed + // (the collections pool the buffer renews from is released too) + if (!disposed && syncCommandBuffer is WorldSyncCommandBuffer concreteBuffer) + reusableSyncBuffer = concreteBuffer; + #if !UNITY_WEBGL - // Pairs with semaphore.Wait in GetSyncCommandBuffer; must release on every path, - // including the disposed early-return and any exception thrown by Apply, - // otherwise the slot leaks and subsequent Rents time out forever. - semaphore.Release(); + // Pairs with semaphore.Wait in GetSyncCommandBuffer; must release on every path, + // including the disposed early-return and any exception thrown by Apply, + // otherwise the slot leaks and subsequent Rents time out forever. + semaphore.Release(); #endif - } } } } diff --git a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/ICRDTWorldSynchronizer.cs b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/ICRDTWorldSynchronizer.cs index 397400212b2..54bf29193ee 100644 --- a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/ICRDTWorldSynchronizer.cs +++ b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/ICRDTWorldSynchronizer.cs @@ -23,5 +23,11 @@ public interface ICRDTWorldSynchronizer : IDisposable /// /// void ApplySyncCommandBuffer(IWorldSyncCommandBuffer syncCommandBuffer); + + /// + /// Finalizes the command buffer without applying it (nothing to apply) and allows to rent it again. + /// Unlike it does not touch the World so no synchronization with the main thread is needed + /// + void ReleaseSyncCommandBuffer(IWorldSyncCommandBuffer syncCommandBuffer); } } diff --git a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/IWorldSyncCommandBuffer.cs b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/IWorldSyncCommandBuffer.cs index 49cd6ccb9cb..fde441de9be 100644 --- a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/IWorldSyncCommandBuffer.cs +++ b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/IWorldSyncCommandBuffer.cs @@ -14,6 +14,12 @@ namespace CrdtEcsBridge.WorldSynchronizer /// public interface IWorldSyncCommandBuffer : IDisposable { + /// + /// True if applying the buffer would not change the World at all. + /// Valid only after + /// + bool IsEmpty { get; } + /// /// Add messages in the order they are processed by CRDT /// This sync function is for ECS syncing and it heavily relies on the proper result from the CRDT Protocol diff --git a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/Tests/CrdtWorldSynchronizerShould.cs b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/Tests/CrdtWorldSynchronizerShould.cs index f33f5655a69..f453b81d317 100644 --- a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/Tests/CrdtWorldSynchronizerShould.cs +++ b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/Tests/CrdtWorldSynchronizerShould.cs @@ -34,5 +34,34 @@ public void ReleaseSyncBuffer() Assert.DoesNotThrow(() => crdtWorldSynchronizer.GetSyncCommandBuffer()); } + + [Test] + public void ReuseSyncBufferInstance() + { + //Arrange + IWorldSyncCommandBuffer first = crdtWorldSynchronizer.GetSyncCommandBuffer(); + first.FinalizeAndDeserialize(); + crdtWorldSynchronizer.ApplySyncCommandBuffer(first); + + //Act + IWorldSyncCommandBuffer second = crdtWorldSynchronizer.GetSyncCommandBuffer(); + + //Assert + Assert.AreSame(first, second); + } + + [Test] + public void AllowRentAfterReleaseWithoutApplying() + { + //Arrange + IWorldSyncCommandBuffer worldSyncCommandBuffer = crdtWorldSynchronizer.GetSyncCommandBuffer(); + worldSyncCommandBuffer.FinalizeAndDeserialize(); + + //Act + crdtWorldSynchronizer.ReleaseSyncCommandBuffer(worldSyncCommandBuffer); + + //Assert + Assert.DoesNotThrow(() => crdtWorldSynchronizer.GetSyncCommandBuffer()); + } } } diff --git a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/Tests/WorldSyncCommandBufferShould.cs b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/Tests/WorldSyncCommandBufferShould.cs index bb782feea54..e254c0a02ab 100644 --- a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/Tests/WorldSyncCommandBufferShould.cs +++ b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/Tests/WorldSyncCommandBufferShould.cs @@ -137,6 +137,57 @@ public void FailGracefullyOnUnknownComponent() Assert.AreEqual(CRDTReconciliationEffect.NoChanges, result); } + [Test] + public void ReportEmptyWhenThereAreNoChangesToApply() + { + //Act + worldSyncCommandBuffer.FinalizeAndDeserialize(); + + //Assert + Assert.IsTrue(worldSyncCommandBuffer.IsEmpty); + } + + [Test] + public void ReportNotEmptyWhenComponentChanged() + { + //Arrange + worldSyncCommandBuffer.SyncCRDTMessage(CreateTestMessage(), CRDTReconciliationEffect.ComponentAdded); + + //Act + worldSyncCommandBuffer.FinalizeAndDeserialize(); + + //Assert + Assert.IsFalse(worldSyncCommandBuffer.IsEmpty); + } + + [Test] + public void ReportNotEmptyWhenEntityDeleted() + { + //Arrange + var message = new CRDTMessage(CRDTMessageType.DELETE_ENTITY, ENTITY_ID, 0, 0, EmptyMemoryOwner.EMPTY); + worldSyncCommandBuffer.SyncCRDTMessage(message, CRDTReconciliationEffect.EntityDeleted); + + //Act + worldSyncCommandBuffer.FinalizeAndDeserialize(); + + //Assert + Assert.IsFalse(worldSyncCommandBuffer.IsEmpty); + } + + [Test] + public void ReportEmptyWhenAllChangesResolveToNoChanges() + { + //Arrange: Added followed by Deleted merges to NoChanges + worldSyncCommandBuffer.SyncCRDTMessage(CreateTestMessage(), CRDTReconciliationEffect.ComponentAdded); + worldSyncCommandBuffer.SyncCRDTMessage(CreateTestMessage(), CRDTReconciliationEffect.ComponentDeleted); + + //Act + worldSyncCommandBuffer.FinalizeAndDeserialize(); + + //Assert + Assert.IsTrue(worldSyncCommandBuffer.IsEmpty); + } + private static (CRDTMessage, CRDTReconciliationEffect effect, CRDTReconciliationEffect expected)[][] MessagesSource() { return new[] diff --git a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/WorldSyncCommandBuffer.cs b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/WorldSyncCommandBuffer.cs index 1f54ecd68ee..cf22a862e9d 100644 --- a/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/WorldSyncCommandBuffer.cs +++ b/Explorer/Assets/DCL/Infrastructure/CrdtEcsBridge/WorldSynchronizer/WorldSyncCommandBuffer.cs @@ -44,15 +44,27 @@ public class WorldSyncCommandBuffer : IWorldSyncCommandBuffer { (CRDTReconciliationEffect.ComponentDeleted, CRDTReconciliationEffect.ComponentDeleted), CRDTReconciliationEffect.ComponentDeleted }, }; - private readonly Dictionary> batchStates; - private readonly List deletedEntities; - private readonly ISDKComponentsRegistry sdkComponentsRegistry; private readonly ISceneEntityFactory entityFactory; private readonly WorldSyncCommandBufferCollectionsPool collectionsPool; + private Dictionary> batchStates; + private List deletedEntities; + private bool finalized; private bool deserialized; + private bool containsAnyChangesToApply; + + public bool IsEmpty + { + get + { + if (!deserialized) + throw new InvalidOperationException($"{nameof(FinalizeAndDeserialize)} must be called before {nameof(IsEmpty)}"); + + return deletedEntities.Count == 0 && !containsAnyChangesToApply; + } + } /// /// Can't contain a public ctor as should be instantiated within the assembly @@ -67,6 +79,23 @@ internal WorldSyncCommandBuffer(ISDKComponentsRegistry componentsRegistry, IScen this.collectionsPool = collectionsPool; } + /// + /// Prepares the disposed instance for reuse so a new buffer class is not allocated on every batch. + /// Must be called on the disposed instance only: the renting is serialized by the synchronizer's semaphore + /// + internal void Renew() + { + if (!finalized) + throw new InvalidOperationException($"{nameof(WorldSyncCommandBuffer)} can't be renewed while it is still in use"); + + batchStates = collectionsPool.GetMainDictionary(); + deletedEntities = collectionsPool.GetDeletedEntities(); + + finalized = false; + deserialized = false; + containsAnyChangesToApply = false; + } + public void Dispose() { if (finalized) return; @@ -197,6 +226,8 @@ public void FinalizeAndDeserialize() componentsBatch.Clear(); } + + containsAnyChangesToApply |= containsAnyChanges; } deserialized = true;