using Common.Attributes; using System.Collections.Concurrent; using System.Net.Sockets; using System.Text; namespace DeviceCommand.Flexible { /// /// 灵活型 TCP 通信类(免实例化,每次调用临时建立 TCP 连接)。 /// 支持字节/文本发送、定长字节读取、按结束符读取文本行四种操作, /// 内部使用按端点粒度的通信锁保证同一端点同一时刻只有一个 TCP 事务在执行, /// 不同端点之间互不阻塞,适用于偶发性、无需保持长连接的 TCP 设备通信场景。 /// [ADPCommand] public static class FTCP { // 按端点粒度的通信锁:同一端点同一时刻只有一个事务在执行,不同端点互不阻塞 private static readonly ConcurrentDictionary _commLocks = new(); /// 获取指定端点的通信锁(不存在则自动创建) private static SemaphoreSlim GetLock(string ipAddress, int port) => _commLocks.GetOrAdd($"{ipAddress}:{port}", _ => new SemaphoreSlim(1, 1)); #region Send /// /// 向指定 TCP 端点发送一段字节数据(发送完成后立即断开)。 /// /// 设备 IP 地址 /// TCP 端口号 /// 发送超时时间(毫秒) /// 要发送的字节数组 /// 异步取消令牌 public static async Task 发送字节数据(string ipAddress,int port,int sendTimeout, byte[] buffer, CancellationToken ct = default) { await GetLock(ipAddress, port).WaitAsync(ct); try { using var client = new TcpClient(); client.SendTimeout = sendTimeout; await client.ConnectAsync(ipAddress, port, ct); using NetworkStream stream = client.GetStream(); await stream.WriteAsync(buffer, 0, buffer.Length, ct).WaitAsync(TimeSpan.FromMilliseconds(sendTimeout)); } finally { GetLock(ipAddress, port).Release(); } } /// /// 向指定 TCP 端点发送一段文本(以 UTF8 编码为字节后发送,发送完成后立即断开)。 /// /// 设备 IP 地址 /// TCP 端口号 /// 发送超时时间(毫秒) /// 要发送的文本字符串 /// 异步取消令牌 public static Task 发送文本数据(string ipAddress, int port,int sendTimeout, string text,CancellationToken ct = default) { return 发送字节数据( ipAddress, port, sendTimeout, Encoding.UTF8.GetBytes(text),ct); } #endregion #region Read /// /// 连接指定 TCP 端点并读取指定长度的字节数据(连接关闭或读满为止)。 /// /// 设备 IP 地址 /// TCP 端口号 /// 接收超时时间(毫秒) /// 要读取的字节长度 /// 异步取消令牌 /// 实际读取到的字节数组(对端提前关闭时可能短于请求长度) public static async Task 读取字节数据(string ipAddress,int port,int receiveTimeout,int length,CancellationToken ct = default) { await GetLock(ipAddress, port).WaitAsync(ct); try { using var client = new TcpClient(); client.ReceiveTimeout = receiveTimeout; await client.ConnectAsync(ipAddress, port, ct); using NetworkStream stream = client.GetStream(); byte[] buffer = new byte[length]; int offset = 0; using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct); if (receiveTimeout > 0) cts.CancelAfter(receiveTimeout); while (offset < length) { int read = await stream.ReadAsync(buffer, offset, length - offset, cts.Token); if (read == 0) break; offset += read; } return buffer[..offset]; } finally { GetLock(ipAddress, port).Release(); } } /// /// 连接指定 TCP 端点并读取一行文本,读取到结束符为止(返回内容不含结束符)。 /// /// 设备 IP 地址 /// TCP 端口号 /// 接收超时时间(毫秒),超时未收到结束符抛出超时异常 /// 文本行结束符,默认换行符 "\n" /// 异步取消令牌 /// 去除结束符并去除首尾空白后的文本行 public static async Task 读取文本行( string ipAddress, int port, int receiveTimeout, string delimiter = "\n",CancellationToken ct = default) { await GetLock(ipAddress, port).WaitAsync(ct); try { using var client = new TcpClient(); client.ReceiveTimeout = receiveTimeout; await client.ConnectAsync(ipAddress, port, ct); using NetworkStream stream = client.GetStream(); var sb = new StringBuilder(); byte[] buffer = new byte[1024]; using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct); if (receiveTimeout > 0) cts.CancelAfter(receiveTimeout); while (!cts.Token.IsCancellationRequested) { int bytesRead = await stream.ReadAsync(buffer, 0, buffer.Length, cts.Token); if (bytesRead == 0) break; sb.Append(Encoding.UTF8.GetString(buffer, 0, bytesRead)); int index = sb.ToString().IndexOf(delimiter, StringComparison.Ordinal); if (index >= 0) return sb.ToString(0, index).Trim(); } throw new TimeoutException("读取超时或对端关闭"); } finally { GetLock(ipAddress, port).Release(); } } #endregion } }