diff --git a/MLAPI-Editor/NetworkingManagerEditor.cs b/MLAPI-Editor/NetworkingManagerEditor.cs index 72af7da169..1de97dc77a 100644 --- a/MLAPI-Editor/NetworkingManagerEditor.cs +++ b/MLAPI-Editor/NetworkingManagerEditor.cs @@ -42,6 +42,8 @@ public class NetworkingManagerEditor : Editor private SerializedProperty networkIdRecycleDelayProperty; private SerializedProperty rpcHashSizeProperty; private SerializedProperty loadSceneTimeOutProperty; + private SerializedProperty enableMessageBufferingProperty; + private SerializedProperty messageBufferTimeoutProperty; private SerializedProperty enableEncryptionProperty; private SerializedProperty signKeyExchangeProperty; private SerializedProperty serverBase64PfxCertificateProperty; @@ -121,6 +123,8 @@ private void Init() networkIdRecycleDelayProperty = networkConfigProperty.FindPropertyRelative("NetworkIdRecycleDelay"); rpcHashSizeProperty = networkConfigProperty.FindPropertyRelative("RpcHashSize"); loadSceneTimeOutProperty = networkConfigProperty.FindPropertyRelative("LoadSceneTimeOut"); + enableMessageBufferingProperty = networkConfigProperty.FindPropertyRelative("EnableMessageBuffering"); + messageBufferTimeoutProperty = networkConfigProperty.FindPropertyRelative("MessageBufferTimeout"); enableEncryptionProperty = networkConfigProperty.FindPropertyRelative("EnableEncryption"); signKeyExchangeProperty = networkConfigProperty.FindPropertyRelative("SignKeyExchange"); serverBase64PfxCertificateProperty = networkConfigProperty.FindPropertyRelative("ServerBase64PfxCertificate"); @@ -160,6 +164,8 @@ private void CheckNullProperties() networkIdRecycleDelayProperty = networkConfigProperty.FindPropertyRelative("NetworkIdRecycleDelay"); rpcHashSizeProperty = networkConfigProperty.FindPropertyRelative("RpcHashSize"); loadSceneTimeOutProperty = networkConfigProperty.FindPropertyRelative("LoadSceneTimeOut"); + enableMessageBufferingProperty = networkConfigProperty.FindPropertyRelative("EnableMessageBuffering"); + messageBufferTimeoutProperty = networkConfigProperty.FindPropertyRelative("MessageBufferTimeout"); enableEncryptionProperty = networkConfigProperty.FindPropertyRelative("EnableEncryption"); signKeyExchangeProperty = networkConfigProperty.FindPropertyRelative("SignKeyExchange"); serverBase64PfxCertificateProperty = networkConfigProperty.FindPropertyRelative("ServerBase64PfxCertificate"); @@ -345,6 +351,13 @@ public override void OnInspectorGUI() EditorGUILayout.PropertyField(networkIdRecycleDelayProperty); } + EditorGUILayout.PropertyField(enableMessageBufferingProperty); + + using (new EditorGUI.DisabledScope(!networkingManager.NetworkConfig.EnableMessageBuffering)) + { + EditorGUILayout.PropertyField(messageBufferTimeoutProperty); + } + EditorGUILayout.LabelField("Bandwidth", EditorStyles.boldLabel); EditorGUILayout.PropertyField(rpcHashSizeProperty); diff --git a/MLAPI/Configuration/NetworkConfig.cs b/MLAPI/Configuration/NetworkConfig.cs index 42f079a373..47163ee64c 100644 --- a/MLAPI/Configuration/NetworkConfig.cs +++ b/MLAPI/Configuration/NetworkConfig.cs @@ -166,6 +166,16 @@ public class NetworkConfig [Tooltip("The amount of seconds to wait for all clients to load a requested scene")] public int LoadSceneTimeOut = 120; /// + /// Whether or not message buffering should be enabled. This will resolve most out of order messages during spawn. + /// + [Tooltip("Whether or not message buffering should be enabled. This will resolve most out of order messages during spawn")] + public bool EnableMessageBuffering = true; + /// + /// The amount of time a message should be buffered for without being consumed. If it is not consumed within this time, it will be dropped. + /// + [Tooltip("The amount of time a message should be buffered for without being consumed. If it is not consumed within this time, it will be dropped")] + public float MessageBufferTimeout = 20f; + /// /// Whether or not to enable the ECDHE key exchange to allow for encryption and authentication of messages /// [Tooltip("Whether or not to enable the ECDHE key exchange to allow for encryption and authentication of messages")] diff --git a/MLAPI/Core/NetworkingManager.cs b/MLAPI/Core/NetworkingManager.cs index 3db86abdd0..f1a40d3eb5 100644 --- a/MLAPI/Core/NetworkingManager.cs +++ b/MLAPI/Core/NetworkingManager.cs @@ -25,6 +25,7 @@ using static MLAPI.Messaging.CustomMessagingManager; using MLAPI.Exceptions; using MLAPI.Transports.Tasks; +using MLAPI.Messaging.Buffering; namespace MLAPI { @@ -695,6 +696,11 @@ private void Update() NetworkedObject.NetworkedBehaviourUpdate(); } + if (!IsServer && NetworkConfig.EnableMessageBuffering) + { + BufferManager.CleanBuffer(); + } + if (IsServer) { lastEventTickTime = NetworkTime; @@ -881,7 +887,7 @@ private void HandleRawTransportPoll(NetEventType eventType, ulong clientId, stri case NetEventType.Data: if (LogHelper.CurrentLogLevel <= LogLevel.Developer) LogHelper.LogInfo($"Incoming Data From {clientId} : {payload.Count} bytes"); - HandleIncomingData(clientId, channelName, payload, receiveTime); + HandleIncomingData(clientId, channelName, payload, receiveTime, true); break; case NetEventType.Disconnect: NetworkProfiler.StartEvent(TickType.Receive, 0, "NONE", "TRANSPORT_DISCONNECT"); @@ -905,7 +911,7 @@ private void HandleRawTransportPoll(NetEventType eventType, ulong clientId, stri private readonly BitStream inputStreamWrapper = new BitStream(new byte[0]); - private void HandleIncomingData(ulong clientId, string channelName, ArraySegment data, float receiveTime) + internal void HandleIncomingData(ulong clientId, string channelName, ArraySegment data, float receiveTime, bool allowBuffer) { if (LogHelper.CurrentLogLevel <= LogLevel.Developer) LogHelper.LogInfo("Unwrapping Data Header"); @@ -939,7 +945,31 @@ private void HandleIncomingData(ulong clientId, string channelName, ArraySegment return; } + + void bufferCallback(ulong networkId) + { + if (!allowBuffer) + { + // This is to prevent recursive buffering + if (LogHelper.CurrentLogLevel <= LogLevel.Error) LogHelper.LogError("A message of type " + MLAPIConstants.MESSAGE_NAMES[messageType] + " was recursivley buffered. It has been dropped."); + return; + } + + if (!NetworkConfig.EnableMessageBuffering) + { + throw new InvalidOperationException("Cannot buffer with buffering disabled."); + } + + if (IsServer) + { + throw new InvalidOperationException("Cannot buffer on server."); + } + + BufferManager.BufferMessageForNetworkId(networkId, clientId, channelName, receiveTime, data); + } + #region INTERNAL MESSAGE + switch (messageType) { case MLAPIConstants.MLAPI_CONNECTION_REQUEST: @@ -955,7 +985,7 @@ private void HandleIncomingData(ulong clientId, string channelName, ArraySegment if (IsClient) InternalMessageHandler.HandleDestroyObject(clientId, messageStream); break; case MLAPIConstants.MLAPI_SWITCH_SCENE: - if (IsClient && NetworkConfig.EnableSceneManagement) InternalMessageHandler.HandleSwitchScene(clientId, messageStream); + if (IsClient) InternalMessageHandler.HandleSwitchScene(clientId, messageStream); break; case MLAPIConstants.MLAPI_CHANGE_OWNER: if (IsClient) InternalMessageHandler.HandleChangeOwner(clientId, messageStream); @@ -970,10 +1000,10 @@ private void HandleIncomingData(ulong clientId, string channelName, ArraySegment if (IsClient) InternalMessageHandler.HandleTimeSync(clientId, messageStream, receiveTime); break; case MLAPIConstants.MLAPI_NETWORKED_VAR_DELTA: - InternalMessageHandler.HandleNetworkedVarDelta(clientId, messageStream); + InternalMessageHandler.HandleNetworkedVarDelta(clientId, messageStream, bufferCallback); break; case MLAPIConstants.MLAPI_NETWORKED_VAR_UPDATE: - InternalMessageHandler.HandleNetworkedVarUpdate(clientId, messageStream); + InternalMessageHandler.HandleNetworkedVarUpdate(clientId, messageStream, bufferCallback); break; case MLAPIConstants.MLAPI_SERVER_RPC: if (IsServer) InternalMessageHandler.HandleServerRPC(clientId, messageStream); @@ -985,10 +1015,10 @@ private void HandleIncomingData(ulong clientId, string channelName, ArraySegment if (IsClient) InternalMessageHandler.HandleServerRPCResponse(clientId, messageStream); break; case MLAPIConstants.MLAPI_CLIENT_RPC: - if (IsClient) InternalMessageHandler.HandleClientRPC(clientId, messageStream); + if (IsClient) InternalMessageHandler.HandleClientRPC(clientId, messageStream, bufferCallback); break; case MLAPIConstants.MLAPI_CLIENT_RPC_REQUEST: - if (IsClient) InternalMessageHandler.HandleClientRPCRequest(clientId, messageStream, channelName, security); + if (IsClient) InternalMessageHandler.HandleClientRPCRequest(clientId, messageStream, channelName, security, bufferCallback); break; case MLAPIConstants.MLAPI_CLIENT_RPC_RESPONSE: if (IsServer) InternalMessageHandler.HandleClientRPCResponse(clientId, messageStream); diff --git a/MLAPI/Messaging/Buffering/BufferManager.cs b/MLAPI/Messaging/Buffering/BufferManager.cs new file mode 100644 index 0000000000..0bab140d41 --- /dev/null +++ b/MLAPI/Messaging/Buffering/BufferManager.cs @@ -0,0 +1,92 @@ +using System; +using System.Collections.Generic; +using MLAPI.Serialization.Pooled; +using UnityEngine; + +namespace MLAPI.Messaging.Buffering +{ + internal static class BufferManager + { + private static readonly Dictionary> bufferQueues = new Dictionary>(); + + internal struct BufferedMessage + { + internal ulong sender; + internal string channelName; + internal PooledBitStream payload; + internal float receiveTime; + internal float bufferTime; + } + + internal static Queue ConsumeBuffersForNetworkId(ulong networkId) + { + if (bufferQueues.ContainsKey(networkId)) + { + Queue message = bufferQueues[networkId]; + + bufferQueues.Remove(networkId); + + return message; + } + else + { + return null; + } + } + + internal static void RecycleConsumedBufferedMessage(BufferedMessage message) + { + message.payload.Dispose(); + } + + internal static void BufferMessageForNetworkId(ulong networkId, ulong sender, string channelName, float receiveTime, ArraySegment payload) + { + if (!bufferQueues.ContainsKey(networkId)) + { + bufferQueues.Add(networkId, new Queue()); + } + + Queue queue = bufferQueues[networkId]; + + PooledBitStream payloadStream = PooledBitStream.Get(); + + payloadStream.Write(payload.Array, payload.Offset, payload.Count); + payloadStream.Position = 0; + + queue.Enqueue(new BufferedMessage() + { + bufferTime = Time.realtimeSinceStartup, + channelName = channelName, + payload = payloadStream, + receiveTime = receiveTime, + sender = sender + }); + } + + private static readonly List _keysToDestroy = new List(); + internal static void CleanBuffer() + { + foreach (KeyValuePair> pair in bufferQueues) + { + while (pair.Value.Count > 0 && Time.realtimeSinceStartup - pair.Value.Peek().bufferTime >= NetworkingManager.Singleton.NetworkConfig.MessageBufferTimeout) + { + BufferedMessage message = pair.Value.Dequeue(); + + RecycleConsumedBufferedMessage(message); + } + + if (pair.Value.Count == 0) + { + _keysToDestroy.Add(pair.Key); + } + } + + for (int i = 0; i < _keysToDestroy.Count; i++) + { + bufferQueues.Remove(_keysToDestroy[i]); + } + + _keysToDestroy.Clear(); + } + } +} diff --git a/MLAPI/Messaging/InternalMessageHandler.cs b/MLAPI/Messaging/InternalMessageHandler.cs index 4224ef11eb..30b657d357 100644 --- a/MLAPI/Messaging/InternalMessageHandler.cs +++ b/MLAPI/Messaging/InternalMessageHandler.cs @@ -14,6 +14,8 @@ using UnityEngine; using UnityEngine.Events; using UnityEngine.SceneManagement; +using System.Collections.Generic; +using MLAPI.Messaging.Buffering; namespace MLAPI.Messaging { @@ -374,6 +376,21 @@ internal static void HandleAddObject(ulong clientId, Stream stream) NetworkedObject netObject = SpawnManager.CreateLocalNetworkedObject(softSync, instanceId, prefabHash, parentNetworkId, pos, rot); SpawnManager.SpawnNetworkedObjectLocally(netObject, networkId, softSync, isPlayerObject, ownerId, stream, hasPayload, payLoadLength, true, false); + + Queue bufferQueue = BufferManager.ConsumeBuffersForNetworkId(networkId); + + // Apply buffered messages + if (bufferQueue != null) + { + while (bufferQueue.Count > 0) + { + BufferManager.BufferedMessage message = bufferQueue.Dequeue(); + + NetworkingManager.Singleton.HandleIncomingData(message.sender, message.channelName, new ArraySegment(message.payload.GetBuffer(), (int)message.payload.Position, (int)message.payload.Length), message.receiveTime, false); + + BufferManager.RecycleConsumedBufferedMessage(message); + } + } } } @@ -465,7 +482,7 @@ internal static void HandleTimeSync(ulong clientId, Stream stream, float receive } } - internal static void HandleNetworkedVarDelta(ulong clientId, Stream stream) + internal static void HandleNetworkedVarDelta(ulong clientId, Stream stream, Action bufferCallback) { if (!NetworkingManager.Singleton.NetworkConfig.EnableNetworkedVar) { @@ -481,22 +498,29 @@ internal static void HandleNetworkedVarDelta(ulong clientId, Stream stream) if (SpawnManager.SpawnedObjects.ContainsKey(networkId)) { NetworkedBehaviour instance = SpawnManager.SpawnedObjects[networkId].GetBehaviourAtOrderIndex(orderIndex); + if (instance == null) { - if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVar message recieved for a non existant behaviour"); - return; + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVarDelta message recieved for a non existant behaviour. NetworkId: " + networkId + ", behaviourIndex: " + orderIndex); } - NetworkedBehaviour.HandleNetworkedVarDeltas(instance.networkedVarFields, stream, clientId, instance); + else + { + NetworkedBehaviour.HandleNetworkedVarDeltas(instance.networkedVarFields, stream, clientId, instance); + } + } + else if (NetworkingManager.Singleton.IsServer || !NetworkingManager.Singleton.NetworkConfig.EnableMessageBuffering) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVarDelta message recieved for a non existant object with id: " + networkId + ". This delta was lost."); } else { - if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVar message recieved for a non existant object with id: " + networkId); - return; + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVarDelta message recieved for a non existant object with id: " + networkId + ". This delta will be buffered and might be recovered."); + bufferCallback(networkId); } } } - internal static void HandleNetworkedVarUpdate(ulong clientId, Stream stream) + internal static void HandleNetworkedVarUpdate(ulong clientId, Stream stream, Action bufferCallback) { if (!NetworkingManager.Singleton.NetworkConfig.EnableNetworkedVar) { @@ -512,17 +536,24 @@ internal static void HandleNetworkedVarUpdate(ulong clientId, Stream stream) if (SpawnManager.SpawnedObjects.ContainsKey(networkId)) { NetworkedBehaviour instance = SpawnManager.SpawnedObjects[networkId].GetBehaviourAtOrderIndex(orderIndex); + if (instance == null) { - if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVar message recieved for a non existant behaviour"); - return; + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVarUpdate message recieved for a non existant behaviour. NetworkId: " + networkId + ", behaviourIndex: " + orderIndex); + } + else + { + NetworkedBehaviour.HandleNetworkedVarUpdate(instance.networkedVarFields, stream, clientId, instance); } - NetworkedBehaviour.HandleNetworkedVarUpdate(instance.networkedVarFields, stream, clientId, instance); + } + else if (NetworkingManager.Singleton.IsServer || !NetworkingManager.Singleton.NetworkConfig.EnableMessageBuffering) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVarUpdate message recieved for a non existant object with id: " + networkId + ". This delta was lost."); } else { - if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVar message recieved for a non existant object with id: " + networkId); - return; + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("NetworkedVarUpdate message recieved for a non existant object with id: " + networkId + ". This delta will be buffered and might be recovered."); + bufferCallback(networkId); } } } @@ -563,11 +594,20 @@ internal static void HandleServerRPC(ulong clientId, Stream stream) if (SpawnManager.SpawnedObjects.ContainsKey(networkId)) { NetworkedBehaviour behaviour = SpawnManager.SpawnedObjects[networkId].GetBehaviourAtOrderIndex(behaviourId); - if (behaviour != null) + + if (behaviour == null) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ServerRPC message recieved for a non existant behaviour. NetworkId: " + networkId + ", behaviourIndex: " + behaviourId); + } + else { behaviour.OnRemoteServerRPC(hash, clientId, stream); } } + else if (NetworkingManager.Singleton.IsServer || !NetworkingManager.Singleton.NetworkConfig.EnableMessageBuffering) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ServerRPC message recieved for a non existant object with id: " + networkId + ". This message is lost."); + } } } @@ -583,7 +623,12 @@ internal static void HandleServerRPCRequest(ulong clientId, Stream stream, strin if (SpawnManager.SpawnedObjects.ContainsKey(networkId)) { NetworkedBehaviour behaviour = SpawnManager.SpawnedObjects[networkId].GetBehaviourAtOrderIndex(behaviourId); - if (behaviour != null) + + if (behaviour == null) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ServerRPCRequest message recieved for a non existant behaviour. NetworkId: " + networkId + ", behaviourIndex: " + behaviourId); + } + else { object result = behaviour.OnRemoteServerRPC(hash, clientId, stream); @@ -599,6 +644,10 @@ internal static void HandleServerRPCRequest(ulong clientId, Stream stream, strin } } } + else + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ServerRPCRequest message recieved for a non existant object with id: " + networkId + ". This message is lost."); + } } } @@ -618,10 +667,14 @@ internal static void HandleServerRPCResponse(ulong clientId, Stream stream) responseBase.Result = reader.ReadObjectPacked(responseBase.Type); responseBase.IsSuccessful = true; } + else + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ServerRPCResponse message recieved for a non existant responseId: " + responseId + ". This response is lost."); + } } } - - internal static void HandleClientRPC(ulong clientId, Stream stream) + + internal static void HandleClientRPC(ulong clientId, Stream stream, Action bufferCallback) { using (PooledBitReader reader = PooledBitReader.Get(stream)) { @@ -632,15 +685,29 @@ internal static void HandleClientRPC(ulong clientId, Stream stream) if (SpawnManager.SpawnedObjects.ContainsKey(networkId)) { NetworkedBehaviour behaviour = SpawnManager.SpawnedObjects[networkId].GetBehaviourAtOrderIndex(behaviourId); - if (behaviour != null) + + if (behaviour == null) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ClientRPC message recieved for a non existant behaviour. NetworkId: " + networkId + ", behaviourIndex: " + behaviourId); + } + else { behaviour.OnRemoteClientRPC(hash, clientId, stream); } } + else if (NetworkingManager.Singleton.IsServer || !NetworkingManager.Singleton.NetworkConfig.EnableMessageBuffering) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ClientRPC message recieved for a non existant object with id: " + networkId + ". This message is lost."); + } + else + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ClientRPC message recieved for a non existant object with id: " + networkId + ". This message will be buffered and might be recovered."); + bufferCallback(networkId); + } } } - - internal static void HandleClientRPCRequest(ulong clientId, Stream stream, string channelName, SecuritySendFlags security) + + internal static void HandleClientRPCRequest(ulong clientId, Stream stream, string channelName, SecuritySendFlags security, Action bufferCallback) { using (PooledBitReader reader = PooledBitReader.Get(stream)) { @@ -652,7 +719,12 @@ internal static void HandleClientRPCRequest(ulong clientId, Stream stream, strin if (SpawnManager.SpawnedObjects.ContainsKey(networkId)) { NetworkedBehaviour behaviour = SpawnManager.SpawnedObjects[networkId].GetBehaviourAtOrderIndex(behaviourId); - if (behaviour != null) + + if (behaviour == null) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ClientRPCRequest message recieved for a non existant behaviour. NetworkId: " + networkId + ", behaviourIndex: " + behaviourId); + } + else { object result = behaviour.OnRemoteClientRPC(hash, clientId, stream); @@ -668,6 +740,15 @@ internal static void HandleClientRPCRequest(ulong clientId, Stream stream, strin } } } + else if (NetworkingManager.Singleton.IsServer || !NetworkingManager.Singleton.NetworkConfig.EnableMessageBuffering) + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ClientRPCRequest message recieved for a non existant object with id: " + networkId + ". This message is lost."); + } + else + { + if (LogHelper.CurrentLogLevel <= LogLevel.Normal) LogHelper.LogWarning("ClientRPCRequest message recieved for a non existant object with id: " + networkId + ". This message will be buffered and might be recovered."); + bufferCallback(networkId); + } } } diff --git a/MLAPI/SceneManagement/NetworkSceneManager.cs b/MLAPI/SceneManagement/NetworkSceneManager.cs index 5456cbac34..83a4c04d0a 100644 --- a/MLAPI/SceneManagement/NetworkSceneManager.cs +++ b/MLAPI/SceneManagement/NetworkSceneManager.cs @@ -10,6 +10,7 @@ using MLAPI.Spawning; using UnityEngine; using UnityEngine.SceneManagement; +using MLAPI.Messaging.Buffering; namespace MLAPI.SceneManagement { @@ -361,6 +362,21 @@ private static void OnSceneUnloadClient(Guid switchSceneGuid, Stream objectStrea NetworkedObject networkedObject = SpawnManager.CreateLocalNetworkedObject(false, 0, prefabHash, parentNetworkId, position, rotation); SpawnManager.SpawnNetworkedObjectLocally(networkedObject, networkId, true, isPlayerObject, owner, objectStream, false, 0, true, false); + + Queue bufferQueue = BufferManager.ConsumeBuffersForNetworkId(networkId); + + // Apply buffered messages + if (bufferQueue != null) + { + while (bufferQueue.Count > 0) + { + BufferManager.BufferedMessage message = bufferQueue.Dequeue(); + + NetworkingManager.Singleton.HandleIncomingData(message.sender, message.channelName, new ArraySegment(message.payload.GetBuffer(), (int)message.payload.Position, (int)message.payload.Length), message.receiveTime, false); + + BufferManager.RecycleConsumedBufferedMessage(message); + } + } } } } @@ -391,6 +407,21 @@ private static void OnSceneUnloadClient(Guid switchSceneGuid, Stream objectStrea NetworkedObject networkedObject = SpawnManager.CreateLocalNetworkedObject(true, instanceId, 0, parentNetworkId, null, null); SpawnManager.SpawnNetworkedObjectLocally(networkedObject, networkId, true, isPlayerObject, owner, objectStream, false, 0, true, false); + + Queue bufferQueue = BufferManager.ConsumeBuffersForNetworkId(networkId); + + // Apply buffered messages + if (bufferQueue != null) + { + while (bufferQueue.Count > 0) + { + BufferManager.BufferedMessage message = bufferQueue.Dequeue(); + + NetworkingManager.Singleton.HandleIncomingData(message.sender, message.channelName, new ArraySegment(message.payload.GetBuffer(), (int)message.payload.Position, (int)message.payload.Length), message.receiveTime, false); + + BufferManager.RecycleConsumedBufferedMessage(message); + } + } } } }