127 lines
3.5 KiB
C#
127 lines
3.5 KiB
C#
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<byte[]> _receiveQueue;
|
|
|
|
public UDPBReceiveBuffer(UDPBSocket socket, ConcurrentQueue<byte[]> receiveQueue)
|
|
{
|
|
_socket = socket;
|
|
_receiveQueue = receiveQueue;
|
|
}
|
|
|
|
List<UDPBSequence> _sequences = new List<UDPBSequence>();
|
|
|
|
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;
|
|
|
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
} |