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();
}
}
}