Files
HawkeyeVision/framework/Inspectron.HawkEye/UDPB/UDPBReceiveBuffer.cs
2025-08-19 12:45:50 +02:00

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;
}
}
}
}