| | | 1 | | using System; |
| | | 2 | | using UnityEngine; |
| | | 3 | | |
| | | 4 | | // The mailbox handles the communication latencies for agent-to-agent communication. The mailbox |
| | | 5 | | // owns the message queue and releases messages to be delivered after a set latency. |
| | | 6 | | public class Mailbox { |
| | 1 | 7 | | private static readonly Configs.LinkConfig _defaultLinkConfig = new Configs.LinkConfig { |
| | | 8 | | PacketDeliveryRatio = 1f, |
| | | 9 | | }; |
| | | 10 | | |
| | | 11 | | // Message queue. |
| | 2 | 12 | | private readonly PriorityQueue<PendingMessage> _messageQueue = |
| | | 13 | | new PriorityQueue<PendingMessage>(); |
| | | 14 | | |
| | | 15 | | // Enqueue a message to be delivered after a set latency. |
| | 2 | 16 | | public void Send(Message message) { |
| | 2 | 17 | | if (message == null) { |
| | 0 | 18 | | return; |
| | | 19 | | } |
| | | 20 | | |
| | 3 | 21 | | if (CommsManager.Instance.ContainsNode(message.Receiver)) { |
| | 1 | 22 | | Configs.LinkConfig config = GetLinkConfig(message); |
| | | 23 | | |
| | | 24 | | // TODO(Joseph0120): Set the packet delivery ratio to config.PacketDeliveryRatio. |
| | 1 | 25 | | float packetDeliveryRatio = Mathf.Clamp01(1); |
| | 1 | 26 | | if (UnityEngine.Random.value >= packetDeliveryRatio) { |
| | 0 | 27 | | return; |
| | | 28 | | } |
| | | 29 | | |
| | 1 | 30 | | float latency = config.LatencySeconds; |
| | 1 | 31 | | float jitter = Utilities.SampleStandardNormal() * config.LatencyStdSeconds; |
| | 1 | 32 | | float totalLatency = Math.Max(0f, latency + jitter); |
| | 1 | 33 | | float deliverAt = SimManager.Instance.ElapsedTime + totalLatency; |
| | | 34 | | |
| | | 35 | | // Deliver the message immediately if it is already due. |
| | 1 | 36 | | if (deliverAt <= SimManager.Instance.ElapsedTime) { |
| | 0 | 37 | | if (CommsManager.Instance.ContainsNode(message.Receiver)) { |
| | 0 | 38 | | message.Receiver.Receive(message); |
| | 0 | 39 | | } |
| | 1 | 40 | | } else { |
| | 1 | 41 | | var pendingMessage = new PendingMessage(message, deliverAt); |
| | 1 | 42 | | _messageQueue.Enqueue(pendingMessage, pendingMessage.DeliverAt); |
| | 1 | 43 | | } |
| | 1 | 44 | | } |
| | 2 | 45 | | } |
| | | 46 | | |
| | | 47 | | // Deliver should be called at every fixed update to check for due messages to be delivered. |
| | 3 | 48 | | public void Deliver() { |
| | 4 | 49 | | while (!_messageQueue.IsEmpty() && |
| | 1 | 50 | | _messageQueue.Peek().DeliverAt <= SimManager.Instance.ElapsedTime) { |
| | 1 | 51 | | PendingMessage pendingMessage = _messageQueue.Dequeue(); |
| | 2 | 52 | | if (CommsManager.Instance.ContainsNode(pendingMessage.Receiver)) { |
| | 1 | 53 | | pendingMessage.Receiver.Receive(pendingMessage.Message); |
| | 1 | 54 | | } |
| | 1 | 55 | | } |
| | 3 | 56 | | } |
| | | 57 | | |
| | | 58 | | // Clear all pending messages. |
| | 0 | 59 | | public void Clear() { |
| | 0 | 60 | | _messageQueue.Clear(); |
| | 0 | 61 | | } |
| | | 62 | | |
| | 1 | 63 | | private Configs.LinkConfig GetLinkConfig(Message message) { |
| | 1 | 64 | | Configs.CommunicationConfig communicationConfig = |
| | | 65 | | SimManager.Instance.SimulationConfig?.CommunicationConfig; |
| | 1 | 66 | | if (communicationConfig == null) { |
| | 0 | 67 | | return _defaultLinkConfig; |
| | | 68 | | } |
| | | 69 | | |
| | 1 | 70 | | Configs.AgentType senderType = message.Sender.EndpointType; |
| | 1 | 71 | | Configs.AgentType receiverType = message.Receiver.EndpointType; |
| | 3 | 72 | | foreach (Configs.LinkOverride linkOverride in communicationConfig.LinkOverrides) { |
| | 0 | 73 | | if (linkOverride.From == senderType && linkOverride.To == receiverType) { |
| | 0 | 74 | | return linkOverride.LinkConfig; |
| | | 75 | | } |
| | 0 | 76 | | } |
| | 1 | 77 | | return communicationConfig.LinkConfig ?? _defaultLinkConfig; |
| | 1 | 78 | | } |
| | | 79 | | } |