Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions MLAPI-Editor/NetworkingManagerEditor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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);

Expand Down
10 changes: 10 additions & 0 deletions MLAPI/Configuration/NetworkConfig.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
/// <summary>
/// Whether or not message buffering should be enabled. This will resolve most out of order messages during spawn.
/// </summary>
[Tooltip("Whether or not message buffering should be enabled. This will resolve most out of order messages during spawn")]
public bool EnableMessageBuffering = true;
/// <summary>
/// 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.
/// </summary>
[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;
/// <summary>
/// Whether or not to enable the ECDHE key exchange to allow for encryption and authentication of messages
/// </summary>
[Tooltip("Whether or not to enable the ECDHE key exchange to allow for encryption and authentication of messages")]
Expand Down
44 changes: 37 additions & 7 deletions MLAPI/Core/NetworkingManager.cs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
using static MLAPI.Messaging.CustomMessagingManager;
using MLAPI.Exceptions;
using MLAPI.Transports.Tasks;
using MLAPI.Messaging.Buffering;

namespace MLAPI
{
Expand Down Expand Up @@ -695,6 +696,11 @@ private void Update()
NetworkedObject.NetworkedBehaviourUpdate();
}

if (!IsServer && NetworkConfig.EnableMessageBuffering)
{
BufferManager.CleanBuffer();
}

if (IsServer)
{
lastEventTickTime = NetworkTime;
Expand Down Expand Up @@ -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");
Expand All @@ -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<byte> data, float receiveTime)
internal void HandleIncomingData(ulong clientId, string channelName, ArraySegment<byte> data, float receiveTime, bool allowBuffer)
{
if (LogHelper.CurrentLogLevel <= LogLevel.Developer) LogHelper.LogInfo("Unwrapping Data Header");

Expand Down Expand Up @@ -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:
Expand All @@ -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);
Expand All @@ -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);
Expand All @@ -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);
Expand Down
92 changes: 92 additions & 0 deletions MLAPI/Messaging/Buffering/BufferManager.cs
Original file line number Diff line number Diff line change
@@ -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<ulong, Queue<BufferedMessage>> bufferQueues = new Dictionary<ulong, Queue<BufferedMessage>>();

internal struct BufferedMessage
{
internal ulong sender;
internal string channelName;
internal PooledBitStream payload;
internal float receiveTime;
internal float bufferTime;
}

internal static Queue<BufferedMessage> ConsumeBuffersForNetworkId(ulong networkId)
{
if (bufferQueues.ContainsKey(networkId))
{
Queue<BufferedMessage> 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<byte> payload)
{
if (!bufferQueues.ContainsKey(networkId))
{
bufferQueues.Add(networkId, new Queue<BufferedMessage>());
}

Queue<BufferedMessage> 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<ulong> _keysToDestroy = new List<ulong>();
internal static void CleanBuffer()
{
foreach (KeyValuePair<ulong, Queue<BufferedMessage>> 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();
}
}
}
Loading