TCP Updates

This commit is contained in:
2022-11-08 09:46:47 -06:00
parent e9d19daf6a
commit ac68106fac
2 changed files with 128 additions and 112 deletions
@@ -0,0 +1,61 @@
using System.Collections.Generic;
namespace SharpOSC
{
public static class SlipFrame
{
static readonly byte END = 0xc0;
static readonly byte ESC = 0xdb;
static readonly byte ESC_END = 0xDC;
static readonly byte ESC_ESC = 0xDD;
public static List<byte[]> Decode(byte[] data)
{
List<byte[]> messages = new List<byte[]>();
List<byte> buffer = new List<byte>();
for (int i = 0; i < data.Length; i++)
{
if (data[i] == END && buffer.Count > 0)
{
messages.Add(buffer.ToArray());
buffer.Clear();
}
else if (data[i] != END)
{
buffer.Add(data[i]);
}
}
return messages;
}
public static byte[] Encode(byte[] data)
{
List<byte> slipData = new List<byte>();
byte[] esc_end = { ESC, ESC_END };
byte[] esc_esc = { ESC, ESC_ESC };
byte[] end = { END };
int length = data.Length;
for (int i = 0; i < length; i++)
{
if (data[i] == END)
{
slipData.AddRange(esc_end);
}
else if (data[i] == ESC)
{
slipData.AddRange(esc_esc);
}
else
{
slipData.Add(data[i]);
}
}
slipData.AddRange(end);
return slipData.ToArray();
}
}
}
@@ -5,6 +5,7 @@ using System.Text;
using System.Net.Sockets; using System.Net.Sockets;
using System.Threading; using System.Threading;
using Serilog; using Serilog;
using SuperSimpleTcp;
namespace SharpOSC namespace SharpOSC
{ {
@@ -19,6 +20,8 @@ namespace SharpOSC
public class TCPClient public class TCPClient
{ {
private ILogger _log = Log.Logger.ForContext<TCPClient>();
public int Port public int Port
{ {
get { return _port; } get { return _port; }
@@ -32,18 +35,10 @@ namespace SharpOSC
public delegate void MessageReceivedHandler(object source, MessageEventArgs args); public delegate void MessageReceivedHandler(object source, MessageEventArgs args);
public event MessageReceivedHandler MessageReceived; public event MessageReceivedHandler MessageReceived;
private Queue<OscPacket> SendQueue = new Queue<OscPacket>();
private Thread receivingThread;
private Thread sendThread;
string _address; string _address;
TcpClient client;
byte END = 0xc0; SimpleTcpClient tcpClient;
byte ESC = 0xdb;
byte ESC_END = 0xDC;
byte ESC_ESC = 0xDD;
public TCPClient(string address, int port) public TCPClient(string address, int port)
@@ -54,48 +49,81 @@ namespace SharpOSC
public bool Connect() public bool Connect()
{ {
try try
{ {
client = new TcpClient(Address, Port); tcpClient = new SimpleTcpClient(Address, Port);
receivingThread = new Thread(ReceiveLoop); tcpClient.Events.Connected += ClientConnected;
receivingThread.Start(); tcpClient.Events.DataReceived += DataReceived;
tcpClient.Events.Disconnected += ClientDisconneted;
sendThread = new Thread(SendLoop); tcpClient.Logger += TCPLog;
sendThread.Start(); tcpClient.Connect();
Log.Debug($"[tcpclient] connected to <{Address}:{Port}>");
return true; return true;
} }
catch (Exception e) catch (Exception e)
{ {
Log.Error(e.Message); _log.Error(e.Message);
return false; return false;
} }
} }
public void QueueForSending(OscPacket packet) private void TCPLog(string obj)
{ {
SendQueue.Enqueue(packet); _log.Verbose($"{obj}");
} }
private void SendLoop() private void ClientDisconneted(object sender, ConnectionEventArgs e)
{ {
while (client != null && client.Connected) _log.Verbose($"{e.IpPort} client disconnected: {e.Reason}");
Close();
}
private void DataReceived(object sender, DataReceivedEventArgs e)
{
_log.Verbose($"Raw Data Received contents: {Encoding.UTF8.GetString(e.Data.Array, 0, e.Data.Count)}");
_log.Verbose($"Raw Data Received size: {e.Data.Count}");
List<byte[]> messages = SlipFrame.Decode(e.Data.Array);
_log.Verbose($"Slip decoded {messages.Count} osc messages");
foreach (var message in messages)
{ {
if (SendQueue.Count > 0) _log.Verbose($"Raw message contents: {Encoding.UTF8.GetString(message, 0, message.Length)}");
try
{ {
OscPacket packet = SendQueue.Dequeue(); OscPacket packet = OscPacket.GetPacket(message);
Send(packet); OscMessage responseMessage = (OscMessage)packet;
if (packet == null)
{
_log.Error("packet is null");
}
if (responseMessage == null)
{
_log.Error("responeMessage is null");
}
_log.Debug($"OSC Message Received: {responseMessage.Address}");
OnMessageReceived(responseMessage);
_log.Debug($"After OnMessageReceived Event");
}
catch (Exception ex)
{
_log.Error($"Exception parsing OSC message: {ex.ToString()}");
} }
} }
} }
private void ClientConnected(object sender, ConnectionEventArgs e)
{
_log.Debug($"connected to <{Address}:{Port}>");
}
public void Send(byte[] message) public void Send(byte[] message)
{ {
byte[] slipData = SlipEncode(message); byte[] slipData = SlipFrame.Encode(message);
NetworkStream netStream = client.GetStream(); tcpClient.Send(slipData.ToArray());
netStream.Write(slipData.ToArray(), 0, slipData.ToArray().Length);
} }
public void Send(OscPacket packet) public void Send(OscPacket packet)
@@ -108,98 +136,25 @@ namespace SharpOSC
{ {
get get
{ {
if (client == null) if (tcpClient == null)
return false; return false;
else else
return client.Connected; return tcpClient.IsConnected;
} }
} }
public void ReceiveLoop()
{
while (client != null && client.Connected)
{
Receive();
}
//Log.Debug("[tcpclient] - ReceiveLoop has exited");
}
public void Receive()
{
Random random = new Random();
int num = random.Next(1000);
try
{
NetworkStream netStream = client.GetStream();
netStream.ReadTimeout = 250;
List<byte> responseData = new List<byte>();
if (netStream.CanRead)
{
//var watch = System.Diagnostics.Stopwatch.StartNew();
byte[] buffer = new byte[256];
int bytesRead = 0;
int reads = 0;
do
{
bytesRead = netStream.Read(buffer, 0, buffer.Length);
responseData.AddRange(buffer);
reads += 1;
Thread.Sleep(1);
//Log.Debug("Thread " + num + ": Bytes read: " + bytesRead + " - " + Encoding.UTF8.GetString(buffer));
} while (netStream.DataAvailable);
//Console.WriteLine("Raw TCP In: " + System.Text.Encoding.UTF8.GetString(responseData.ToArray()));
OscPacket packet = OscPacket.GetPacket(responseData.Skip(1).ToArray());
OscMessage responseMessage = (OscMessage)packet;
//watch.Stop();
//Console.WriteLine($"TCPCLient - message receive took {watch.ElapsedMilliseconds}ms and {reads} reads");
OnMessageReceived(responseMessage);
}
}
catch (Exception e)
{
//Console.WriteLine("TCPSENDER - Receive Exception: " + e.ToString());
}
}
public byte[] SlipEncode(byte[] data)
{
List<byte> slipData = new List<byte>();
byte[] esc_end = { ESC, ESC_END };
byte[] esc_esc = { ESC, ESC_ESC };
byte[] end = { END };
int length = data.Length;
for (int i = 0; i < length; i++)
{
if (data[i] == END)
{
slipData.AddRange(esc_end);
}
else if (data[i] == ESC)
{
slipData.AddRange(esc_esc);
}
else
{
slipData.Add(data[i]);
}
}
slipData.AddRange(end);
return slipData.ToArray();
}
public void Close() public void Close()
{ {
if (client != null) if (tcpClient != null)
{ {
if (client.Connected) tcpClient.Events.Connected -= ClientConnected;
tcpClient.Events.DataReceived -= DataReceived;
tcpClient.Events.Disconnected -= ClientDisconneted;
if (tcpClient.IsConnected)
{ {
Log.Debug($"[tcpClient] closing connection to {Address}"); _log.Debug($"closing connection to {Address}");
client.GetStream().Close(); tcpClient.Disconnect();
client.Close(); tcpClient.Dispose();
} }
} }
} }