using System; using System.Collections.Concurrent; using System.ComponentModel; using System.Collections.Generic; using System.Linq; using System.Net; using System.Net.Sockets; using System.Text; using System.Threading.Tasks; using Newtonsoft.Json; using TeamAAS.Communication.Attributes; using TeamAAS.Communication.Base; using TeamAAS.Communication.Interfaces; using TeamAAS.Communication.Enums; namespace TeamAAS.Communication.Devices { [Communication("TCP 服务端", "基础通讯", "TCP 监听,接受客户端连接")] public class TcpServerCommunication : BindableCommunicationBase { [JsonIgnore] private TcpListener _listener; [JsonIgnore] private readonly object _sync = new object(); [JsonIgnore] private readonly List _clients = new List(); [JsonIgnore] private readonly ConcurrentQueue _received = new ConcurrentQueue(); [JsonIgnore] private System.Threading.Timer _reconnectTimer; [JsonIgnore] private volatile bool _hasStartedOnce; [JsonIgnore] private volatile bool _manualDisconnect; [JsonIgnore] private volatile bool _isReconnecting; private string _localIp = "0.0.0.0"; public string LocalIp { get { return _localIp; } set { SetProperty(ref _localIp, value); } } private int _port = 7790; public int Port { get { return _port; } set { SetProperty(ref _port, value); } } private Terminator _terminator = Terminator.None; [Category("III.数据格式"), DisplayName("结束符"), Description("发送时自动附加、接收时自动去除的结束符")] public Terminator Terminator { get { return _terminator; } set { SetProperty(ref _terminator, value); } } private DataEncoding _dataEncoding = DataEncoding.Default; [Category("III.数据格式"), DisplayName("编码格式"), Description("收发数据的编码格式")] public DataEncoding DataEncoding { get { return _dataEncoding; } set { SetProperty(ref _dataEncoding, value); } } [Browsable(false)] public override string EndpointUrl { get { return $"tcp://{LocalIp}:{Port}"; } set { if (!string.IsNullOrWhiteSpace(value) && value.StartsWith("tcp://")) { var uri = new Uri(value); if (!string.IsNullOrWhiteSpace(uri.Host)) LocalIp = uri.Host; if (uri.Port > 0) Port = uri.Port; } } } public override bool IsConnected => _listener != null; public override event Action ConnectChangedEvent; public override event Action DataReceivedEvent; public TcpServerCommunication() { } public override void Connect() { Disconnect(); IPAddress bindIp = IPAddress.Any; if (!string.IsNullOrWhiteSpace(LocalIp)) { try { bindIp = IPAddress.Parse(LocalIp.Trim()); } catch { bindIp = IPAddress.Any; } } _listener = new TcpListener(bindIp, Port); _listener.Start(); _manualDisconnect = false; _hasStartedOnce = true; StopReconnectTimer(); _ = Task.Run(AcceptLoop); ConnectChangedEvent?.Invoke(this, true); Notify(nameof(IsConnected)); } public override Task ConnectAsync() { return Task.Run(() => Connect()); } private async Task AcceptLoop() { try { while (_listener != null) { var client = await _listener.AcceptTcpClientAsync(); lock (_sync) _clients.Add(client); _ = Task.Run(() => ClientLoop(client)); } } catch { } finally { if (_hasStartedOnce && !_manualDisconnect) StartAutoReconnect(); } } private async Task ClientLoop(TcpClient client) { var buffer = new byte[4096]; var receiveBuffer = new StringBuilder(); try { var stream = client.GetStream(); while (client.Connected) { int n = await stream.ReadAsync(buffer, 0, buffer.Length); if (n <= 0) break; receiveBuffer.Append(GetEncoding().GetString(buffer, 0, n)); FlushReceivedMessages(receiveBuffer); } } catch { } finally { lock (_sync) _clients.Remove(client); client.Dispose(); } } /// /// 获取编码 /// public Encoding GetEncoding() { switch (DataEncoding) { case DataEncoding.ASCII: return Encoding.ASCII; case DataEncoding.UTF7: return Encoding.UTF7; case DataEncoding.UTF8: return Encoding.UTF8; case DataEncoding.UTF32: return Encoding.UTF32; case DataEncoding.Unicode: return Encoding.Unicode; case DataEncoding.BigEndianUnicode: return Encoding.BigEndianUnicode; case DataEncoding.GB2312: return Encoding.GetEncoding("gb2312"); default: return Encoding.Default; } } /// /// 获取结束符字符串 /// public string GetTerminatorString() { switch (Terminator) { case Terminator.CR: return "\r"; case Terminator.LF: return "\n"; case Terminator.CRLF: return "\r\n"; case Terminator.None: default: return string.Empty; } } /// /// 去除数据尾部结束符 /// public string TrimTerminator(string text) { if (string.IsNullOrEmpty(text)) return text; var term = GetTerminatorString(); if (string.IsNullOrEmpty(term)) return text; return text.EndsWith(term) ? text.Substring(0, text.Length - term.Length) : text; } /// /// 按结束符分帧:收到完整结束符才触发接收事件;无结束符时直接触发 /// private void FlushReceivedMessages(StringBuilder receiveBuffer) { var term = GetTerminatorString(); if (string.IsNullOrEmpty(term)) { if (receiveBuffer.Length > 0) { var text = receiveBuffer.ToString(); receiveBuffer.Clear(); _received.Enqueue(text); DataReceivedEvent?.Invoke(this, text); } return; } string content = receiveBuffer.ToString(); int idx; while ((idx = content.IndexOf(term, StringComparison.Ordinal)) >= 0) { var msg = content.Substring(0, idx); content = content.Substring(idx + term.Length); _received.Enqueue(msg); DataReceivedEvent?.Invoke(this, msg); } receiveBuffer.Clear(); receiveBuffer.Append(content); } /// /// 自动重连:监听曾成功启动,意外停止后每5秒尝试重新监听 /// private void StartAutoReconnect() { if (_reconnectTimer != null) return; if (_listener != null) { try { _listener.Stop(); } catch { } _listener = null; } lock (_sync) { foreach (var c in _clients) { try { c.Dispose(); } catch { } } _clients.Clear(); } ConnectChangedEvent?.Invoke(this, false); Notify(nameof(IsConnected)); _reconnectTimer = new System.Threading.Timer(ReconnectTimer_Elapsed, null, TimeSpan.FromSeconds(5), TimeSpan.FromSeconds(5)); } private void ReconnectTimer_Elapsed(object state) { if (_manualDisconnect) { StopReconnectTimer(); return; } if (_isReconnecting) return; _isReconnecting = true; try { if (_listener != null) { try { _listener.Stop(); } catch { } _listener = null; } IPAddress bindIp = IPAddress.Any; if (!string.IsNullOrWhiteSpace(LocalIp)) { try { bindIp = IPAddress.Parse(LocalIp.Trim()); } catch { bindIp = IPAddress.Any; } } var listener = new TcpListener(bindIp, Port); try { listener.Start(); _listener = listener; StopReconnectTimer(); _ = Task.Run(AcceptLoop); ConnectChangedEvent?.Invoke(this, true); Notify(nameof(IsConnected)); } catch { try { listener.Stop(); } catch { } } } finally { _isReconnecting = false; } } private void StopReconnectTimer() { var t = _reconnectTimer; _reconnectTimer = null; if (t != null) { try { t.Dispose(); } catch { } } } public void Send(string text) { var data = GetEncoding().GetBytes((text ?? string.Empty) + GetTerminatorString()); lock (_sync) { foreach (var c in _clients.ToList()) { try { c.GetStream().Write(data, 0, data.Length); } catch { } } } } public override object ReadValue(string address) { _received.TryDequeue(out var text); return text; } public override Task ReadValueAsync(string address) { return Task.FromResult(ReadValue(address)); } public override void WriteValue(string address, object value) { Send(value?.ToString() ?? string.Empty); } public override Task WriteValueAsync(string address, object value) { Send(value?.ToString() ?? string.Empty); return Task.CompletedTask; } public override void Disconnect() { StopReconnectTimer(); _manualDisconnect = true; if (_listener != null) { _listener.Stop(); _listener = null; } lock (_sync) { foreach (var c in _clients) { try { c.Dispose(); } catch { } } _clients.Clear(); } ConnectChangedEvent?.Invoke(this, false); Notify(nameof(IsConnected)); } public override void Dispose() { Disconnect(); } } }