using Model.Models;
using System;
using System.IO;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
namespace DeviceCommand.Base
{
public class Tcp : ITcp, IDisposable
{
public string IPAddress { get; set; } = "127.0.0.1";
public int Port { get; set; } = 502;
public int SendTimeout { get; set; } = 3000;
public int ReceiveTimeout { get; set; } = 3000;
private TcpClient _tcpClient;
public bool IsConnected => _tcpClient?.Connected ?? false;
protected readonly SemaphoreSlim _commLock = new(1, 1);
public Tcp()
{
_tcpClient = new TcpClient();
}
///
/// 通过 一次性配置 TCP 通信参数。
///
public Tcp(TcpConfig config) : this()
{
if (config == null) return;
ConfigureDevice(config.IPAddress, config.Port, config.SendTimeout, config.ReceiveTimeout);
}
public void ConfigureDevice(string ipAddress, int port, int sendTimeout = 3000, int receiveTimeout = 3000)
{
IPAddress = ipAddress;
Port = port;
SendTimeout = sendTimeout;
ReceiveTimeout = receiveTimeout;
}
public virtual async Task ConnectAsync(CancellationToken ct = default)
{
await _commLock.WaitAsync(ct);
try
{
return await ResetConnectionAsync(ct);
}
finally
{
_commLock.Release();
}
}
///
/// 核心内部方法:无锁状态下安全重置并重建连接
///
private async Task ResetConnectionAsync(CancellationToken ct)
{
if (_tcpClient != null)
{
_tcpClient.Close();
_tcpClient.Dispose();
}
_tcpClient = new TcpClient();
await _tcpClient.ConnectAsync(IPAddress, Port, ct).ConfigureAwait(false);
return true;
}
public virtual void Close()
{
if (_tcpClient?.Connected == true) _tcpClient.Close();
}
private async Task LoglessSendAsync(byte[] buffer, CancellationToken ct)
{
if (!IsConnected) throw new InvalidOperationException("TCP未连接。");
NetworkStream stream = _tcpClient.GetStream();
try
{
if (SendTimeout > 0)
{
// 使用 WaitAsync 控制发送超时
await stream.WriteAsync(buffer, 0, buffer.Length, ct)
.WaitAsync(TimeSpan.FromMilliseconds(SendTimeout), ct)
.ConfigureAwait(false);
}
else
{
await stream.WriteAsync(buffer, 0, buffer.Length, ct).ConfigureAwait(false);
}
}
catch (TimeoutException ex)
{
// 发送超时,安全起见直接断开重建
await ResetConnectionAsync(ct);
throw new TimeoutException($"TCP 发送数据超时({SendTimeout} ms),连接已重置以确保安全。", ex);
}
}
public async Task SendAsync(byte[] buffer, CancellationToken ct = default)
{
await _commLock.WaitAsync(ct);
try
{
await LoglessSendAsync(buffer, ct);
}
finally
{
_commLock.Release();
}
}
public async Task SendAsync(string str, CancellationToken ct = default)
{
await SendAsync(Encoding.UTF8.GetBytes(str), ct);
}
public async Task ReadAsync(int length, CancellationToken ct = default)
{
await _commLock.WaitAsync(ct);
try
{
if (!IsConnected) throw new InvalidOperationException("TCP未连接。");
NetworkStream stream = _tcpClient.GetStream();
byte[] buffer = new byte[length];
int offset = 0;
try
{
while (offset < length)
{
ct.ThrowIfCancellationRequested();
// 计算剩余超时时间(这里简单使用配置值,若要求极精准可加入不长计算)
int read = ReceiveTimeout > 0
? await stream.ReadAsync(buffer, offset, length - offset, ct)
.WaitAsync(TimeSpan.FromMilliseconds(ReceiveTimeout), ct)
.ConfigureAwait(false)
: await stream.ReadAsync(buffer, offset, length - offset, ct).ConfigureAwait(false);
if (read == 0) throw new IOException("远程主机已关闭连接");
offset += read;
}
return offset == 0 ? Array.Empty() : buffer[..offset];
}
catch (TimeoutException ex)
{
// 关键安全重置:WaitAsync 超时后抛弃老连接,防止未完成的 Read 污染后续 Buffer
await ResetConnectionAsync(ct);
throw new TimeoutException($"TCP 读取定长数据超时({ReceiveTimeout} ms),连接已重置避免数据错乱。", ex);
}
}
finally
{
_commLock.Release();
}
}
private async Task LoglessReadAsync(string delimiter, CancellationToken ct)
{
if (!IsConnected) throw new InvalidOperationException("TCP未连接。");
delimiter ??= "\n";
var sb = new StringBuilder();
byte[] buffer = new byte[1024];
NetworkStream stream = _tcpClient.GetStream();
try
{
while (true)
{
ct.ThrowIfCancellationRequested();
int bytesRead = ReceiveTimeout > 0
? await stream.ReadAsync(buffer, 0, buffer.Length, ct)
.WaitAsync(TimeSpan.FromMilliseconds(ReceiveTimeout), ct)
.ConfigureAwait(false)
: await stream.ReadAsync(buffer, 0, buffer.Length, ct).ConfigureAwait(false);
if (bytesRead == 0) throw new IOException("远程主机已关闭连接");
sb.Append(Encoding.UTF8.GetString(buffer, 0, bytesRead));
string currentText = sb.ToString();
int index = currentText.IndexOf(delimiter, StringComparison.Ordinal);
if (index >= 0)
{
return currentText.Substring(0, index).Trim();
}
}
}
catch (TimeoutException ex)
{
// 当 SCPI 仪表由于命令发错不响应时,WaitAsync 会立功并抛出 TimeoutException
// 此时我们在内部秒速“重置连接”,上层业务只会收到一个干净的超时报错,但由于锁和自动重连,下一次通信直接恢复!
await ResetConnectionAsync(ct);
throw new TimeoutException($"SCPI 读取超时(未收到结束符 '{delimiter}',等待:{ReceiveTimeout} ms),通信链路已自动刷新。", ex);
}
}
public async Task ReadAsync(string delimiter = "\n", CancellationToken ct = default)
{
await _commLock.WaitAsync(ct);
try
{
return await LoglessReadAsync(delimiter, ct);
}
finally
{
_commLock.Release();
}
}
public async Task WriteReadAsync(string command, string delimiter = "\n", CancellationToken ct = default)
{
await _commLock.WaitAsync(ct);
try
{
await LoglessSendAsync(Encoding.UTF8.GetBytes(command), ct);
return await LoglessReadAsync(delimiter, ct);
}
finally
{
_commLock.Release();
}
}
#region 扩展:读取所有可用的二进制网络字节 (无 SCPI 块头解析)
///
/// 无锁核心方法:发送命令并一次性读取所有回传的二进制原始数据包(不进行 # 协议头解析,专用于读取纯文件流如 PNG)
///
private async Task LoglessReadAllBytesAsync(string queryCommand, CancellationToken ct)
{
if (!IsConnected) throw new InvalidOperationException("TCP未连接。");
// 1. 发送查询命令(例如 :PRINt? PNG\n)
await LoglessSendAsync(Encoding.UTF8.GetBytes(queryCommand), ct);
NetworkStream stream = _tcpClient.GetStream();
// 2. 鼎阳示波器的截图数据大概在 100KB - 800KB 左右,使用动态内存流接收整个包
using (var ms = new MemoryStream())
{
byte[] buffer = new byte[8192]; // 8KB 缓冲区
try
{
// 先给设备短暂的反应时间,等待数据到达网络缓冲区
int delayCount = 0;
while (!_tcpClient.GetStream().DataAvailable && delayCount < 50)
{
await Task.Delay(10, ct);
delayCount++;
}
// 循环读取,直到网络流中没有更多数据
do
{
ct.ThrowIfCancellationRequested();
int read = ReceiveTimeout > 0
? await stream.ReadAsync(buffer, 0, buffer.Length, ct)
.WaitAsync(TimeSpan.FromMilliseconds(ReceiveTimeout), ct)
.ConfigureAwait(false)
: await stream.ReadAsync(buffer, 0, buffer.Length, ct).ConfigureAwait(false);
if (read == 0) break; // 远程流关闭
ms.Write(buffer, 0, read);
// 如果流里没有剩余数据了,退出读取(防止 ReadAsync 在没有数据时无限阻塞等待)
if (!stream.DataAvailable)
{
// 极短延时再确认一次,防止分包网络延迟引起的“假结束”
await Task.Delay(30, ct);
if (!stream.DataAvailable)
{
break;
}
}
} while (true);
return ms.ToArray();
}
catch (TimeoutException ex)
{
await ResetConnectionAsync(ct);
throw new TimeoutException($"读取二进制大包超时(等待:{ReceiveTimeout} ms),链路已重置。", ex);
}
}
}
///
/// 【公开方法】发送命令并读取设备回传的全部原始二进制字节数组(不带协议头解析,直接返回整个字节缓冲区)
///
public async Task ReadAllBytesAsync(string queryCommand, CancellationToken ct = default)
{
await _commLock.WaitAsync(ct);
try
{
return await LoglessReadAllBytesAsync(queryCommand, ct);
}
finally
{
_commLock.Release();
}
}
#endregion
public void Dispose()
{
_tcpClient?.Dispose();
_commLock?.Dispose();
}
}
}