using System; using System.Collections.Generic; using System.Collections.Specialized; using System.Linq; using System.Net; using System.Net.WebSockets; using System.Text; using System.Threading; using System.Threading.Tasks; namespace XFEExtension.CyberComm { /// /// 传回的参数类型 /// public enum BackMessageType { /// /// 文本消息 /// Text, /// /// 二进制消息 /// Binary, /// /// 错误消息 /// Error } /// /// CyberComm客户端 /// public class CyberCommClient { private int reconnectTimes = -1; #region 公共属性 /// /// 指定的WS服务器URL /// public string ServerURL { get; set; } /// /// 是否已连接 /// public bool IsConnected { get; private set; } = false; /// /// 是否自动重连 /// public bool AutoReconnect { get; set; } /// /// 自动重连最大次数 /// public int ReconnectMaxTimes { get; set; } = -1; /// /// 自动重连尝试间隔 /// public int ReconnectTryDelay { get; set; } = 100; /// /// 是否自动接收完整消息 /// public bool AutoReceiveCompletedMessage { get; set; } /// /// 收到消息时触发 /// public event EventHandler MessageReceived; /// /// 连接关闭时触发 /// public event EventHandler ConnectionClosed; /// /// 连接成功时触发 /// public event EventHandler Connected; /// /// WebSocket客户端 /// public ClientWebSocket ClientWebSocket { get; private set; } #endregion #region 公有方法 /// /// 启动CyberComm客户端 /// /// public async Task StartCyberCommClient() { StartConnect: ClientWebSocket = new ClientWebSocket(); Uri serverUri = new Uri(ServerURL); reconnectTimes++; try { await ClientWebSocket.ConnectAsync(serverUri, CancellationToken.None); } catch (Exception ex) { if (IsConnected == true) { ConnectionClosed?.Invoke(this, EventArgs.Empty); } IsConnected = false; if (AutoReconnect) { if (reconnectTimes <= ReconnectMaxTimes || ReconnectMaxTimes == -1) { Thread.Sleep(ReconnectTryDelay); goto StartConnect; } } else { MessageReceived?.Invoke(this, new CyberCommClientEventArgsImpl(ClientWebSocket, new XFECyberCommException("连接服务器时发生异常", ex))); return; } } Connected?.Invoke(this, EventArgs.Empty); reconnectTimes = 0; while (ClientWebSocket.State == WebSocketState.Open) { try { byte[] receiveBuffer = new byte[1024]; WebSocketReceiveResult receiveResult = await ClientWebSocket.ReceiveAsync(new ArraySegment(receiveBuffer), CancellationToken.None); var bufferList = new List(); bufferList.AddRange(receiveBuffer.Take(receiveResult.Count)); //ReceiveCompletedMessageByUsingWhile if (AutoReceiveCompletedMessage) { while (!receiveResult.EndOfMessage) { receiveResult = await ClientWebSocket.ReceiveAsync(new ArraySegment(receiveBuffer), CancellationToken.None); bufferList.AddRange(receiveBuffer.Take(receiveResult.Count)); } } var receivedBinaryBuffer = bufferList.ToArray(); if (receiveResult.MessageType == WebSocketMessageType.Text) { var receivedMessage = Encoding.UTF8.GetString(receivedBinaryBuffer); MessageReceived?.Invoke(this, new CyberCommClientEventArgsImpl(ClientWebSocket, receivedMessage)); } if (receiveResult.MessageType == WebSocketMessageType.Binary) { MessageReceived?.Invoke(this, new CyberCommClientEventArgsImpl(ClientWebSocket, receivedBinaryBuffer)); } if (receiveResult.MessageType == WebSocketMessageType.Close) { await ClientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, string.Empty, CancellationToken.None); break; } } catch (Exception ex) { if (IsConnected == true) { ConnectionClosed?.Invoke(this, EventArgs.Empty); } IsConnected = false; if (AutoReconnect) { if (AutoReconnect) { Thread.Sleep(ReconnectTryDelay); if (reconnectTimes <= ReconnectMaxTimes || ReconnectMaxTimes == -1) goto StartConnect; } } else { MessageReceived?.Invoke(this, new CyberCommClientEventArgsImpl(ClientWebSocket, new XFECyberCommException("与服务器端通讯期间发生异常", ex))); return; } } } if (AutoReconnect) { goto StartConnect; } ConnectionClosed?.Invoke(this, EventArgs.Empty); } /// /// 发送文本消息 /// /// 待发送的文本 /// 发送进程 public async Task SendTextMessage(string message) { try { byte[] sendBuffer = Encoding.UTF8.GetBytes(message); await ClientWebSocket.SendAsync(new ArraySegment(sendBuffer), WebSocketMessageType.Text, true, CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("客户端发送文本到服务器时出现异常", ex); } } /// /// 发送二进制消息 /// /// 待发送的二进制数据 /// 发送进程 public async Task SendBinaryMessage(byte[] message) { try { await ClientWebSocket.SendAsync(new ArraySegment(message), WebSocketMessageType.Binary, true, CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("客户端发送二进制数据到服务器时出现异常", ex); } } /// /// 关闭CyberComm客户端 /// /// public async Task CloseCyberCommClient() { await ClientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, string.Empty, CancellationToken.None); } #endregion #region 构造函数 /// /// CyberComm客户端 /// /// WS服务器地址 /// 是否自动重连 /// 是否自动接收完整消息 public CyberCommClient(string serverURL, bool autoReconnect = true, bool autoReceiveCompletedMessage = true) { this.AutoReconnect = autoReconnect; this.AutoReceiveCompletedMessage = autoReceiveCompletedMessage; this.ServerURL = serverURL; } /// /// CyberComm客户端 /// public CyberCommClient() { } #endregion } /// /// CyberComm服务器 /// public class CyberCommServer { private readonly string serverURL; #region 公共属性 /// /// 是否自动接收完整消息 /// public bool AutoReceiveCompletedMessage { get; set; } /// /// 收到消息时触发 /// public event EventHandler MessageReceived; /// /// 客户端连接时触发 /// public event EventHandler ClientConnected; /// /// 服务器启动时触发 /// public event EventHandler ServerStarted; /// /// 连接关闭时触发 /// public event EventHandler ConnectionClosed; /// /// WebSocket服务器 /// public HttpListener WebSocketServer { get; } = new HttpListener(); #endregion #region 公有方法 /// /// 启动CyberComm服务器 /// /// public async Task StartCyberCommServer() { try { WebSocketServer.Prefixes.Add(serverURL); WebSocketServer.Start(); ServerStarted?.Invoke(this, EventArgs.Empty); while (true) { HttpListenerContext httpListenerContext = await WebSocketServer.GetContextAsync(); if (httpListenerContext.Request.IsWebSocketRequest) { HttpListenerWebSocketContext httpListenerWebSocketContext = await httpListenerContext.AcceptWebSocketAsync(null); string clientIP = httpListenerContext.Request.RemoteEndPoint.Address.ToString(); WebSocket webSocket = httpListenerWebSocketContext.WebSocket; NameValueCollection wsHeader = httpListenerWebSocketContext.Headers; ClientConnected?.Invoke(this, new CyberCommServerEventArgsImpl(webSocket, string.Empty, clientIP, wsHeader)); CyberCommClientConnected(webSocket, wsHeader, clientIP); } } } catch (Exception ex) { throw new XFECyberCommException("启动服务器时发生异常", ex); } } private async void CyberCommClientConnected(WebSocket webSocket, NameValueCollection wsHeader, string clientIP) { while (webSocket.State == WebSocketState.Open) { try { byte[] receiveBuffer = new byte[1024]; WebSocketReceiveResult receiveResult = await webSocket.ReceiveAsync(new ArraySegment(receiveBuffer), CancellationToken.None); //ReceiveCompletedMessageByUsingWhile var bufferList = new List(); if (AutoReceiveCompletedMessage) { bufferList.AddRange(receiveBuffer.Take(receiveResult.Count)); while (!receiveResult.EndOfMessage) { receiveResult = await webSocket.ReceiveAsync(new ArraySegment(receiveBuffer), CancellationToken.None); bufferList.AddRange(receiveBuffer.Take(receiveResult.Count)); } } var receivedBinaryBuffer = bufferList.ToArray(); if (receiveResult.MessageType == WebSocketMessageType.Text) { string receivedMessage = Encoding.UTF8.GetString(receivedBinaryBuffer); MessageReceived?.Invoke(this, new CyberCommServerEventArgsImpl(webSocket, receivedMessage, clientIP, wsHeader)); } if (receiveResult.MessageType == WebSocketMessageType.Binary) { MessageReceived?.Invoke(this, new CyberCommServerEventArgsImpl(webSocket, receivedBinaryBuffer, clientIP, wsHeader)); } if (receiveResult.MessageType == WebSocketMessageType.Close) { await webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Connection Closed", CancellationToken.None); break; } } catch (Exception ex) { MessageReceived?.Invoke(this, new CyberCommServerEventArgsImpl(webSocket, new XFECyberCommException("与客户端端通讯期间发生异常", ex), clientIP, wsHeader)); break; } } ConnectionClosed?.Invoke(this, new CyberCommServerEventArgsImpl(webSocket, string.Empty, clientIP, wsHeader)); webSocket.Dispose(); } #endregion #region 构造函数 /// /// CyberComm服务器,使用端口创建 /// /// 监听端口 /// 是否自动接收完整消息 public CyberCommServer(int listenPort, bool autoReceiveCompletedMessage = true) { this.serverURL = $"http://*:{listenPort}/"; this.AutoReceiveCompletedMessage = autoReceiveCompletedMessage; } /// /// CyberComm服务器,使用URL创建 /// /// 服务器URL /// 是否自动接收完整消息 public CyberCommServer(string serverURL, bool autoReceiveCompletedMessage = true) { this.serverURL = serverURL; this.AutoReceiveCompletedMessage = autoReceiveCompletedMessage; } #endregion } #region 事件参数及其实现类 /// /// CyberComm客户端事件参数 /// public abstract class CyberCommClientEventArgs : EventArgs { /// /// 消息类型 /// public BackMessageType MessageType { get; } /// /// 当前WebSocket /// public ClientWebSocket CurrentWebSocket { get; } /// /// 文本消息 /// public string TextMessage { get; } /// /// 异常消息(如果有的话) /// public XFECyberCommException Exception { get; } /// /// 二进制消息 /// public byte[] BinaryMessage { get; } /// /// 发送文本消息 /// /// 待发送的文本 /// /// 发送进程 public async Task ReplyMessage(string message) { try { byte[] sendBuffer = Encoding.UTF8.GetBytes(message); await CurrentWebSocket.SendAsync(new ArraySegment(sendBuffer), WebSocketMessageType.Text, true, CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("收到服务器端数据后客户端回复文本数据时出现异常", ex); } } /// /// 关闭连接 /// /// public async Task Close() { try { await CurrentWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Connection Closed", CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("关闭客户端连接时出现异常", ex); } } internal CyberCommClientEventArgs(ClientWebSocket clientWebSocket, string message) { CurrentWebSocket = clientWebSocket; TextMessage = message; MessageType = BackMessageType.Text; } internal CyberCommClientEventArgs(ClientWebSocket clientWebSocket, byte[] bytes) { CurrentWebSocket = clientWebSocket; BinaryMessage = bytes; MessageType = BackMessageType.Binary; } internal CyberCommClientEventArgs(ClientWebSocket clientWebSocket, XFECyberCommException ex) { CurrentWebSocket = clientWebSocket; Exception = ex; MessageType = BackMessageType.Error; } } /// /// CyberComm服务器事件参数 /// public abstract class CyberCommServerEventArgs { /// /// 消息类型 /// public BackMessageType MessageType { get; } /// /// 当前WebSocket /// public WebSocket CurrentWebSocket { get; } /// /// 客户端请求头 /// public NameValueCollection WSHeader { get; } /// /// 异常消息(如果有的话) /// public XFECyberCommException Exception { get; } /// /// 客户端IP地址 /// public string IpAddress { get; } /// /// 文本消息 /// public string TextMessage { get; } /// /// 二进制消息 /// public byte[] BinaryMessage { get; } /// /// 发送文本消息 /// /// 待发送的文本 /// /// 发送进程 public async Task ReplyMessage(string message) { try { byte[] sendBuffer = Encoding.UTF8.GetBytes(message); await CurrentWebSocket.SendAsync(new ArraySegment(sendBuffer), WebSocketMessageType.Text, true, CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("收到客户端数据后服务器端回复文本时出现异常", ex); } } /// /// 发送二进制消息 /// /// 二进制消息 /// /// public async Task ReplyBinaryMessage(byte[] bytes) { try { await CurrentWebSocket.SendAsync(new ArraySegment(bytes), WebSocketMessageType.Binary, true, CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("收到客户端数据后服务器端回复二进制数据时出现异常", ex); } } /// /// 关闭连接 /// /// public async Task Close() { try { await CurrentWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Connection Closed", CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("关闭服务器端连接时出现异常", ex); } } /// /// 强制关闭连接 /// /// public void ForceClose() { try { CurrentWebSocket.Abort(); } catch (Exception ex) { throw new XFECyberCommException("强制关闭服务器端连接时出现异常", ex); } } internal CyberCommServerEventArgs(WebSocket webSocket, string message, string ipAddress, NameValueCollection wsHeader) { CurrentWebSocket = webSocket; TextMessage = message; IpAddress = ipAddress; MessageType = BackMessageType.Text; WSHeader = wsHeader; } internal CyberCommServerEventArgs(WebSocket webSocket, byte[] bytes, string ipAddress, NameValueCollection wsHeader) { CurrentWebSocket = webSocket; BinaryMessage = bytes; IpAddress = ipAddress; MessageType = BackMessageType.Binary; WSHeader = wsHeader; } internal CyberCommServerEventArgs(WebSocket webSocket, XFECyberCommException ex, string ipAddress, NameValueCollection wsHeader) { CurrentWebSocket = webSocket; Exception = ex; IpAddress = ipAddress; MessageType = BackMessageType.Error; WSHeader = wsHeader; } } class CyberCommClientEventArgsImpl : CyberCommClientEventArgs { internal CyberCommClientEventArgsImpl(ClientWebSocket clientWebSocket, string message) : base(clientWebSocket, message) { } internal CyberCommClientEventArgsImpl(ClientWebSocket clientWebSocket, byte[] bytes) : base(clientWebSocket, bytes) { } internal CyberCommClientEventArgsImpl(ClientWebSocket clientWebSocket, XFECyberCommException ex) : base(clientWebSocket, ex) { } } class CyberCommServerEventArgsImpl : CyberCommServerEventArgs { internal CyberCommServerEventArgsImpl(WebSocket webSocket, string message, string ipAddress, NameValueCollection wsHeader) : base(webSocket, message, ipAddress, wsHeader) { } internal CyberCommServerEventArgsImpl(WebSocket webSocket, byte[] bytes, string ipAddress, NameValueCollection wsHeader) : base(webSocket, bytes, ipAddress, wsHeader) { } internal CyberCommServerEventArgsImpl(WebSocket webSocket, XFECyberCommException ex, string ipAddress, NameValueCollection wsHeader) : base(webSocket, ex, ipAddress, wsHeader) { } } #endregion /// /// CyberComm客户端群组 /// public class CyberCommGroup { private readonly List cyberCommList = new List(); /// /// 组ID /// public string GroupId { get; set; } /// /// 添加客户端 /// /// 客户端 public void Add(CyberCommServerEventArgs e) { cyberCommList.Add(e); } /// /// 移除指定的客户端 /// /// 客户端 public void Remove(WebSocket webSocket) { cyberCommList.Remove(cyberCommList.Find(x => x.CurrentWebSocket == webSocket)); } /// /// 移除指定索引的客户端 /// /// 客户单索引 public void RemoveAt(int index) { cyberCommList.RemoveAt(index); } /// /// 清空列表 /// public void Clear() { cyberCommList.Clear(); } /// /// 群组客户端数量 /// public int Count { get { return cyberCommList.Count; } } /// /// 索引器 /// /// 查找 /// 客户端 public CyberCommServerEventArgs this[Predicate findFunc] { get { return cyberCommList.Find(findFunc); } } /// /// 发送群组文本消息 /// /// 群发文本消息 public async Task SendGroupTextMessage(string message) { List tasks = new List(); foreach (CyberCommServerEventArgs cyberCommServerEventArgs in cyberCommList) { tasks.Add(cyberCommServerEventArgs.ReplyMessage(message)); } await Task.WhenAll(tasks); } /// /// 向指定的WS客户端发送群组文本消息 /// /// 群发文本消息 /// 指定的WS客户端 public async Task SendGroupTextMessage(string message, Func findFunc) { List tasks = new List(); foreach (CyberCommServerEventArgs cyberCommServerEventArgs in cyberCommList) { if (findFunc.Invoke(cyberCommServerEventArgs)) { tasks.Add(cyberCommServerEventArgs.ReplyMessage(message)); } } await Task.WhenAll(tasks); } /// /// 发送群组二进制消息 /// /// 群发二进制消息 public async Task SendGroupBinaryMessage(byte[] bytes) { List tasks = new List(); foreach (CyberCommServerEventArgs cyberCommServerEventArgs in cyberCommList) { tasks.Add(cyberCommServerEventArgs.ReplyBinaryMessage(bytes)); } await Task.WhenAll(tasks); } /// /// 向指定的WS客户端发送群组二进制消息 /// /// 群发二进制消息 /// 指定的WS客户端 public async Task SendGroupBinaryMessage(byte[] bytes, Func findFunc) { List tasks = new List(); foreach (CyberCommServerEventArgs cyberCommServerEventArgs in cyberCommList) { if (findFunc.Invoke(cyberCommServerEventArgs)) { tasks.Add(cyberCommServerEventArgs.ReplyBinaryMessage(bytes)); } } await Task.WhenAll(tasks); } /// /// 客户端群组 /// /// 群组ID public CyberCommGroup(string GroupId) { this.GroupId = GroupId; } /// /// 客户端群组 /// /// 群组ID /// 客户端群组 public CyberCommGroup(string GroupId, List cyberCommList) { this.GroupId = GroupId; this.cyberCommList = cyberCommList; } } /// /// 代签名的WebSocket /// public class SignedWebSocket { /// /// 签名 /// public string Signature { get; set; } /// /// 服务器 /// public WebSocket WebSocket { get; set; } /// /// 签名WebSocket /// /// 签名 /// WebSocket public SignedWebSocket(string signature, WebSocket webSocket) { Signature = signature; WebSocket = webSocket; } } /// /// 通信群组控制器 /// public class CyberCommGroupController { private readonly List commGroups; /// /// 刷新 /// public void Refresh() { for (int i = commGroups.Count - 1; i >= 0; i--) { if (commGroups[i].Count == 0) { commGroups.Remove(commGroups[i]); } } } /// /// 添加群组 /// /// public void AddGroup(CyberCommGroup commGroup) { commGroups.Add(commGroup); } /// /// 移除群组 /// /// public void RemoveGroup(CyberCommGroup commGroup) { commGroups.Remove(commGroup); } /// /// 移除指定索引的群组 /// /// public void RemoveGroupAt(int index) { commGroups.RemoveAt(index); } /// /// 清空群组 /// public void Clear() { commGroups.Clear(); } /// /// 群组数量 /// public int Count { get { return commGroups.Count; } } /// /// 索引器 /// /// /// public CyberCommGroup this[int index] { get { return commGroups[index]; } } /// /// 索引器 /// /// /// public CyberCommGroup this[string GroupId] { get { foreach (CyberCommGroup commGroup in commGroups) { if (commGroup.GroupId == GroupId) { return commGroup; } } return null; } } /// /// 通信群组控制器发送文本消息 /// /// 目标群组的ID /// 发送的文本消息 public async Task SendGroupTextMessage(string GroupId, string message) { CyberCommGroup commGroup = this[GroupId]; if (commGroup != null) { await commGroup.SendGroupTextMessage(message); } } /// /// 通信群组控制器发送二进制消息 /// /// 目标群组的ID /// 发送的二进制消息 public async Task SendGroupBinaryMessage(string GroupId, byte[] bytes) { CyberCommGroup commGroup = this[GroupId]; if (commGroup != null) { await commGroup.SendGroupBinaryMessage(bytes); } } /// /// 通信群组控制器 /// public CyberCommGroupController() { commGroups = new List(); } } }