using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; namespace Inspectron.HawkEye.UDPB { public class UDPBReceiveBuffer { private readonly UDPBSocket _socket; private readonly ConcurrentQueue _receiveQueue; public UDPBReceiveBuffer(UDPBSocket socket, ConcurrentQueue receiveQueue) { _socket = socket; _receiveQueue = receiveQueue; } List _sequences = new List(); public byte[] LastReconstructedBuffer { get; private set; } public bool IsSequenceFinished() { var res = _sequences.OrderByDescending(x=>x.SequenceId).FirstOrDefault(x => x.IsFinished); if (res != null) { LastReconstructedBuffer = res.LastReconstructedBuffer; _sequences.Remove(res); return true; } return false; } public UDPBSequence StartSequence(uint sequence) { while (_sequences.Count > 5) { _sequences.Remove(_sequences.OrderBy(x => x.SequenceId).First()); } var exists = _sequences.FirstOrDefault(x => x.SequenceId == sequence); if (exists != null) { Console.WriteLine($"Sequence exists {sequence}"); return exists; } var res = new UDPBSequence(sequence, _socket, _receiveQueue); _sequences.Add(res); return res; } public void AddPacketBytes(byte[] data, int offset, int len) { var packet = new UDPBPacket(); packet.Deserialize(data, offset, len); if (packet.Type == UDPBPacket.EPacketType.Data) { var seqId = packet.SequenceId; var sequence = _sequences.FirstOrDefault(x => x.SequenceId == seqId); if (sequence != null) { sequence.AddData(packet); } else { var s=StartSequence(seqId); s.AddData(packet); } }else if (packet.Type == UDPBPacket.EPacketType.Ok) { var nb = new byte[len]; Array.Copy(data,offset,nb,0,len); _receiveQueue.Enqueue(nb); } else { AddCommand(packet); } } private void AddCommand(UDPBPacket commandPacket) { switch (commandPacket.Type) { case UDPBPacket.EPacketType.CloseSequence: bool finishSuccess; int dataSize; var sequence = _sequences.FirstOrDefault(x => x.SequenceId == commandPacket.SequenceId); if (sequence == null) return; finishSuccess = sequence.FinishSequence(commandPacket.MessageId, out var missingPacketIds, out dataSize); //inform error break; case UDPBPacket.EPacketType.BeginSequence: StartSequence(commandPacket.SequenceId); break; } } } }