using System; using System.Collections.Generic; using System.IO; using System.Linq; using System.Net.WebSockets; using System.Text; using System.Threading; using System.Threading.Tasks; using XFEExtension.ArrayExtension; using XFEExtension.BufferExtension; using XFEExtension.DelegateExtension; using XFEExtension.FormatExtension; using XFEExtension.TaskExtension; namespace XFEExtension.CyberComm.XCCNetWork { class XCCNetWorkBase { /// /// 明文消息接收时触发 /// public EventHandler textMessageReceived; /// /// 二进制消息接收时触发 /// public EventHandler binaryMessageReceived; /// /// 异常消息接收时触发 /// public EventHandler exceptionMessageReceived; /// /// 连接关闭时触发 /// public EventHandler connectionClosed; /// /// 连接成功时触发 /// public EventHandler connected; } /// /// XCC客户端连接类型 /// public enum XCCClientType { /// /// 明文类型 /// TextMessageClient, /// /// 文件类型 /// FileTransportClient } /// /// XCC网络通讯 /// public class XCCNetWork { private readonly XCCNetWorkBase xCCNetWorkBase; /// /// XCC当前群组 /// public List Groups { get; set; } /// /// 明文消息接收时触发 /// public event EventHandler TextMessageReceived { add { xCCNetWorkBase.textMessageReceived += value; } remove { xCCNetWorkBase.textMessageReceived -= value; } } /// /// 二进制消息接收时触发 /// public event EventHandler BinaryMessageReceived { add { xCCNetWorkBase.binaryMessageReceived += value; } remove { xCCNetWorkBase.binaryMessageReceived -= value; } } /// /// 异常消息接收时触发 /// public event EventHandler ExceptionMessageReceived { add { xCCNetWorkBase.exceptionMessageReceived += value; } remove { xCCNetWorkBase.exceptionMessageReceived -= value; } } /// /// 连接关闭时触发 /// public event EventHandler ConnectionClosed { add { xCCNetWorkBase.connectionClosed += value; } remove { xCCNetWorkBase.connectionClosed -= value; } } /// /// 连接成功时触发 /// public event EventHandler Connected { add { xCCNetWorkBase.connected += value; } remove { xCCNetWorkBase.connected -= value; } } /// /// 创建XCC群组会话 /// /// 群组名称 /// 发送者 /// public XCCGroup CreateGroup(string groupId, string sender) { var group = new XCCGroupImpl(Guid.NewGuid().ToString(), groupId, sender, xCCNetWorkBase); Groups.Add(group); return group; } /// /// XCC网络通信会话 /// public XCCNetWork() { xCCNetWorkBase = new XCCNetWorkBase(); Groups = new List(); } } /// /// XCC群组 /// public abstract class XCCGroup { private event EndTaskTrigger UpdateTaskTrigger; private readonly XCCNetWorkBase workBase; private int reconnectTimes = -1; private bool readyToClose = false; #region 公有属性 /// /// 客户端标识名 /// public string Signature { get; } /// /// 群组ID /// public string GroupId { get; } /// /// 发送者 /// public string Sender { get; } /// /// 明文传输服务器是否连接 /// public bool TextMessageClientConnected { get; private set; } = false; /// /// 文件传输服务器是否连接 /// public bool FileTransportClientConnected { get; private set; } = false; /// /// WebSocket明文传输客户端 /// public ClientWebSocket TextMessageClientWebSocket { get; private set; } /// /// WebSocket文件传输客户端 /// public ClientWebSocket FileTransportClientWebSocket { get; private set; } #endregion #region 公有方法 /// /// 启动XCC会话 /// /// 是否自动重连 /// 最大重连次数,-1则为无限次 /// 重连尝试延迟 /// public async Task StartXCC(bool autoReconnect = true, int reconnectMaxTimes = -1, int reconnectTryDelay = 100) { var textMessageXCCTask = StartTextMessageXCC(autoReconnect, reconnectMaxTimes, reconnectTryDelay); var fileTransportXCCTask = StartFileTransportXCC(autoReconnect, reconnectMaxTimes, reconnectTryDelay); await Task.WhenAll(textMessageXCCTask, fileTransportXCCTask); } /// /// 启动XCC文本会话 /// /// 是否自动重连 /// 最大重连次数,-1则为无限次 /// 重连尝试延迟 /// public async Task StartTextMessageXCC(bool autoReconnect = true, int reconnectMaxTimes = -1, int reconnectTryDelay = 100) { XCCReconnect: TextMessageClientWebSocket = new ClientWebSocket(); Uri serverUri = new Uri("ws://xcc.api.xfegzs.com"); var base64GroupId = Convert.ToBase64String(Encoding.UTF8.GetBytes(GroupId)); var base64SenderId = Convert.ToBase64String(Encoding.UTF8.GetBytes(Sender)); TextMessageClientWebSocket.Options.SetRequestHeader("Group", base64GroupId); TextMessageClientWebSocket.Options.SetRequestHeader("Sender", base64SenderId); TextMessageClientWebSocket.Options.SetRequestHeader("Type", "Text"); TextMessageClientWebSocket.Options.SetRequestHeader("Signature", Signature); reconnectTimes++; try { if (TextMessageClientWebSocket.State != WebSocketState.Open) await TextMessageClientWebSocket.ConnectAsync(serverUri, CancellationToken.None); } catch (Exception ex) { if (readyToClose) { return; } if (TextMessageClientConnected == true) { workBase.connectionClosed?.Invoke(this, new XCCConnectionClosedEventArgsImpl(this, XCCClientType.TextMessageClient, TextMessageClientWebSocket, FileTransportClientWebSocket, false)); } TextMessageClientConnected = false; if (autoReconnect) { if (reconnectTimes <= reconnectMaxTimes || reconnectMaxTimes == -1) { Thread.Sleep(reconnectTryDelay); goto XCCReconnect; } } else { workBase.exceptionMessageReceived?.Invoke(this, new XCCExceptionMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.TextMessageClient, null, null, new XFECyberCommException("与XCC网络通讯明文服务器建立连接时发生异常", ex))); return; } } reconnectTimes = 0; TextMessageClientConnected = true; workBase.connected?.Invoke(this, new XCCConnectedEventArgsImpl(this, XCCClientType.TextMessageClient, TextMessageClientWebSocket, FileTransportClientWebSocket)); while (TextMessageClientWebSocket.State == WebSocketState.Open) { try { byte[] receiveBuffer = new byte[1024]; WebSocketReceiveResult receiveResult = await TextMessageClientWebSocket.ReceiveAsync(new ArraySegment(receiveBuffer), CancellationToken.None); var bufferList = new List(); bufferList.AddRange(receiveBuffer.Take(receiveResult.Count)); //ReceiveCompletedMessageByUsingWhile while (!receiveResult.EndOfMessage) { receiveResult = await TextMessageClientWebSocket.ReceiveAsync(new ArraySegment(receiveBuffer), CancellationToken.None); bufferList.AddRange(receiveBuffer.Take(receiveResult.Count)); } var receivedBinaryBuffer = bufferList.ToArray(); if (receiveResult.MessageType == WebSocketMessageType.Text) { try { var receivedMessage = Encoding.UTF8.GetString(receivedBinaryBuffer); var isHistory = receivedMessage.IndexOf("[XCCGetHistory]") == 0; if (isHistory) { receivedMessage = receivedMessage.Substring(15); } var unPackedMessage = receivedMessage.ToXFEArray(); var messageId = unPackedMessage[0]; var signature = unPackedMessage[1]; var message = unPackedMessage[2]; var senderName = unPackedMessage[3]; var sendTime = DateTime.Parse(unPackedMessage[4]); var messageType = XCCTextMessageType.Text; switch (signature) { case "[XCCTextMessage]": messageType = XCCTextMessageType.Text; break; case "[XCCImage]": messageType = XCCTextMessageType.Image; break; case "[XCCAudio]": messageType = XCCTextMessageType.Audio; break; case "[XCCVideo]": messageType = XCCTextMessageType.Video; break; default: break; } workBase.textMessageReceived?.Invoke(this, new XCCTextMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.TextMessageClient, messageId, messageType, message, senderName, sendTime, isHistory)); } catch (Exception ex) { workBase.exceptionMessageReceived?.Invoke(this, new XCCExceptionMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.TextMessageClient, null, null, new XFECyberCommException("接收XCC服务器消息时发生异常", ex))); } } else if (receiveResult.MessageType == WebSocketMessageType.Binary) { try { var xFEBuffer = XFEBuffer.ToXFEBuffer(receivedBinaryBuffer); var signature = Encoding.UTF8.GetString(xFEBuffer["Type"]); var messageId = Encoding.UTF8.GetString(xFEBuffer["ID"]); if (signature == "callback") UpdateTaskTrigger?.Invoke(true, messageId); } catch (Exception ex) { workBase.exceptionMessageReceived?.Invoke(this, new XCCExceptionMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.TextMessageClient, null, null, new XFECyberCommException("接收XCC服务器消息时发生异常", ex))); } } } catch (Exception ex) { try { await TextMessageClientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Close", CancellationToken.None); } catch { } if (TextMessageClientConnected == true) { workBase.connectionClosed?.Invoke(this, new XCCConnectionClosedEventArgsImpl(this, XCCClientType.TextMessageClient, TextMessageClientWebSocket, FileTransportClientWebSocket, false)); } TextMessageClientConnected = false; if (autoReconnect) { if (autoReconnect) { Thread.Sleep(reconnectTryDelay); if (reconnectTimes <= reconnectMaxTimes || reconnectMaxTimes == -1) goto XCCReconnect; } } else { workBase.exceptionMessageReceived?.Invoke(this, new XCCExceptionMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.TextMessageClient, null, null, new XFECyberCommException("与XCC网络通讯服务器建立连接时发生异常", ex))); return; } } } workBase.connectionClosed?.Invoke(this, new XCCConnectionClosedEventArgsImpl(this, XCCClientType.TextMessageClient, TextMessageClientWebSocket, FileTransportClientWebSocket, true)); } /// /// 启动XCC文件传输会话 /// /// 是否自动重连 /// 最大重连次数,-1则为无限次 /// 重连尝试延迟 /// public async Task StartFileTransportXCC(bool autoReconnect = true, int reconnectMaxTimes = -1, int reconnectTryDelay = 100) { XCCReconnect: FileTransportClientWebSocket = new ClientWebSocket(); Uri serverUri = new Uri("ws://xcc.api.xfegzs.com"); var base64GroupId = Convert.ToBase64String(Encoding.UTF8.GetBytes(GroupId)); var base64SenderId = Convert.ToBase64String(Encoding.UTF8.GetBytes(Sender)); FileTransportClientWebSocket.Options.SetRequestHeader("Group", base64GroupId); FileTransportClientWebSocket.Options.SetRequestHeader("Sender", base64SenderId); FileTransportClientWebSocket.Options.SetRequestHeader("Type", "File"); FileTransportClientWebSocket.Options.SetRequestHeader("Signature", Signature); reconnectTimes++; try { if (FileTransportClientWebSocket.State != WebSocketState.Open) await FileTransportClientWebSocket.ConnectAsync(serverUri, CancellationToken.None); } catch (Exception ex) { if (readyToClose) { return; } if (FileTransportClientConnected == true) { workBase.connectionClosed?.Invoke(this, new XCCConnectionClosedEventArgsImpl(this, XCCClientType.FileTransportClient, TextMessageClientWebSocket, FileTransportClientWebSocket, false)); } FileTransportClientConnected = false; if (autoReconnect) { if (reconnectTimes <= reconnectMaxTimes || reconnectMaxTimes == -1) { Thread.Sleep(reconnectTryDelay); goto XCCReconnect; } } else { workBase.exceptionMessageReceived?.Invoke(this, new XCCExceptionMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.FileTransportClient, null, null, new XFECyberCommException("与XCC网络通讯明文服务器建立连接时发生异常", ex))); return; } } reconnectTimes = 0; FileTransportClientConnected = true; workBase.connected?.Invoke(this, new XCCConnectedEventArgsImpl(this, XCCClientType.FileTransportClient, TextMessageClientWebSocket, FileTransportClientWebSocket)); while (FileTransportClientWebSocket.State == WebSocketState.Open) { try { byte[] receiveBuffer = new byte[1024]; WebSocketReceiveResult receiveResult = await FileTransportClientWebSocket.ReceiveAsync(new ArraySegment(receiveBuffer), CancellationToken.None); var bufferList = new List(); bufferList.AddRange(receiveBuffer.Take(receiveResult.Count)); //ReceiveCompletedMessageByUsingWhile while (!receiveResult.EndOfMessage) { receiveResult = await FileTransportClientWebSocket.ReceiveAsync(new ArraySegment(receiveBuffer), CancellationToken.None); bufferList.AddRange(receiveBuffer.Take(receiveResult.Count)); } var receivedBinaryBuffer = bufferList.ToArray(); if (receiveResult.MessageType == WebSocketMessageType.Binary) { try { var messageType = XCCBinaryMessageType.Binary; var xFEBuffer = XFEBuffer.ToXFEBuffer(receivedBinaryBuffer); var sender = Encoding.UTF8.GetString(xFEBuffer["Sender"]); var signature = Encoding.UTF8.GetString(xFEBuffer["Type"]); if (signature == "callback") return; var messageId = Encoding.UTF8.GetString(xFEBuffer["ID"]); bool isHistory = Encoding.UTF8.GetString(xFEBuffer["IsHistory"]) == "True"; var sendTime = DateTime.Parse(Encoding.UTF8.GetString(xFEBuffer["SendTime"])); byte[] unPackedBuffer = xFEBuffer[sender]; switch (signature) { case "text": messageType = XCCBinaryMessageType.Text; break; case "image": messageType = XCCBinaryMessageType.Image; break; case "audio": messageType = XCCBinaryMessageType.Audio; break; case "audio-buffer": messageType = XCCBinaryMessageType.AudioBuffer; break; case "video": messageType = XCCBinaryMessageType.Video; break; case "callback": UpdateTaskTrigger?.Invoke(true, messageId); continue; default: messageType = XCCBinaryMessageType.Binary; break; } workBase.binaryMessageReceived?.Invoke(this, new XCCBinaryMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.FileTransportClient, sender, messageId, unPackedBuffer, messageType, signature, sendTime, isHistory)); } catch (Exception ex) { workBase.exceptionMessageReceived?.Invoke(this, new XCCExceptionMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.FileTransportClient, null, null, new XFECyberCommException("接收XCC服务器消息时发生异常", ex))); } } } catch (Exception ex) { try { await FileTransportClientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Close", CancellationToken.None); } catch { } if (FileTransportClientConnected == true) { workBase.connectionClosed?.Invoke(this, new XCCConnectionClosedEventArgsImpl(this, XCCClientType.FileTransportClient, TextMessageClientWebSocket, FileTransportClientWebSocket, false)); } FileTransportClientConnected = false; if (autoReconnect) { if (autoReconnect) { Thread.Sleep(reconnectTryDelay); if (reconnectTimes <= reconnectMaxTimes || reconnectMaxTimes == -1) goto XCCReconnect; } } else { workBase.exceptionMessageReceived?.Invoke(this, new XCCExceptionMessageReceivedEventArgsImpl(this, TextMessageClientWebSocket, FileTransportClientWebSocket, XCCClientType.FileTransportClient, null, null, new XFECyberCommException("与XCC网络通讯服务器建立连接时发生异常", ex))); return; } } } workBase.connectionClosed?.Invoke(this, new XCCConnectionClosedEventArgsImpl(this, XCCClientType.FileTransportClient, TextMessageClientWebSocket, FileTransportClientWebSocket, true)); } /// /// 等待明文服务器和文件服务器均连接 /// /// public async Task WaitConnect() { await Task.Run(() => { while (!TextMessageClientConnected || !FileTransportClientConnected) { } }); } /// /// 发送文本消息 /// /// 待发送的文本 /// 最长超时时长 /// 服务器接收校验是否成功 public async Task SendTextMessage(string message, int timeout = 30000) { var messageId = Guid.NewGuid().ToString(); return await SendTextMessage(message, messageId, timeout); } /// /// 发送文本消息 /// /// 待发送的文本 /// 消息ID /// 最长超时时长 /// /// 服务器接收校验是否成功 public async Task SendTextMessage(string message, string messageId, int timeout) { try { byte[] sendBuffer = Encoding.UTF8.GetBytes(new string[] { messageId, "[XCCTextMessage]", message }.ToXFEString()); await TextMessageClientWebSocket.SendAsync(new ArraySegment(sendBuffer), WebSocketMessageType.Text, true, CancellationToken.None); var endTask = Task.Run(async () => { await Task.Delay(timeout); UpdateTaskTrigger?.Invoke(false, messageId); }); return await new XFEWaitTask(ref UpdateTaskTrigger, messageId); } catch (Exception ex) { throw new XFECyberCommException("客户端发送文本到服务器时出现异常", ex); } } /// /// 发送标准的文本消息 /// /// 发送者角色 /// 待发送的文本 /// /// 服务器接收校验是否成功 [Obsolete("发送者已统一,请使用SendTextMessage或SendBinaryTextMessage")] public async Task SendStandardTextMessage(string role, string message) { try { return await SendTextMessage(message); } catch (Exception ex) { throw new XFECyberCommException("客户端发送文本到服务器时出现异常", ex); } } /// /// 发送签名二进制消息 /// /// 二进制消息 /// 签名标识 /// 最长超时时长 /// 服务器接收校验是否成功 public async Task SendSignedBinaryMessage(byte[] message, string signature, int timeout = 10000) { var messageId = Guid.NewGuid().ToString(); return await SendSignedBinaryMessage(message, messageId, signature, timeout); } /// /// 发送签名二进制消息 /// /// 二进制消息 /// 消息ID /// 签名标识 /// 最长超时时长 /// 服务器接收校验是否成功 /// public async Task SendSignedBinaryMessage(byte[] message, string messageId, string signature, int timeout) { try { var xFEBuffer = new XFEBuffer(Sender, message, "Type", Encoding.UTF8.GetBytes(signature), "ID", Encoding.UTF8.GetBytes(messageId)); await FileTransportClientWebSocket.SendAsync(new ArraySegment(xFEBuffer.ToBuffer()), WebSocketMessageType.Binary, true, CancellationToken.None); var endTask = Task.Run(async () => { await Task.Delay(timeout); UpdateTaskTrigger?.Invoke(false, messageId); }); return await new XFEWaitTask(ref UpdateTaskTrigger, messageId); } catch (Exception ex) { throw new XFECyberCommException("客户端发送二进制数据到服务器时出现异常", ex); } } /// /// 发送二进制文本消息 /// /// 消息 /// 服务器接收校验是否成功 /// public async Task SendBinaryTextMessage(string message) { try { return await SendSignedBinaryMessage(Encoding.UTF8.GetBytes(message), "text"); } catch (Exception ex) { throw new XFECyberCommException("客户端发送文本到服务器时出现异常", ex); } } /// /// 发送默认标准的二进制消息 /// /// 待发送的二进制数据 /// 最长超时时长 /// /// 服务器接收校验是否成功 public async Task SendBinaryMessage(byte[] message, int timeout = 1000) { try { return await SendSignedBinaryMessage(message, "binary", timeout); } catch (Exception ex) { throw new XFECyberCommException("客户端发送二进制数据到服务器时出现异常", ex); } } /// /// 发送图片 /// /// 图片路径 /// 服务器接收校验是否成功 /// public async Task SendImage(string filePath) { try { return await SendSignedBinaryMessage(File.ReadAllBytes(filePath), "image", 60000); } catch (Exception ex) { throw new XFECyberCommException("客户端发送图片到服务器时出现异常", ex); } } /// /// 发送视频 /// /// 视频路径 /// 服务器接收校验是否成功 /// public async Task SendVideo(string filePath) { try { return await SendSignedBinaryMessage(File.ReadAllBytes(filePath), "video", 300000); } catch (Exception ex) { throw new XFECyberCommException("客户端发送视频到服务器时出现异常", ex); } } /// /// 发送音频 /// /// 音频路径 /// 服务器接收校验是否成功 /// public async Task SendAudio(string filePath) { try { return await SendSignedBinaryMessage(File.ReadAllBytes(filePath), "audio"); } catch (Exception ex) { throw new XFECyberCommException("客户端发送音频到服务器时出现异常", ex); } } /// /// 发送音频字节流(服务器不会缓存) /// /// 二进制音频流 /// 服务器接收校验是否成功 /// public async Task SendAudioBuffer(byte[] buffer) { try { return await SendSignedBinaryMessage(buffer, "audio-buffer"); } catch (Exception ex) { throw new XFECyberCommException("客户端发送音频到服务器时出现异常", ex); } } /// /// 获取历史记录 /// /// public async Task GetHistory() { try { var messageId = Guid.NewGuid().ToString(); byte[] sendBuffer = Encoding.UTF8.GetBytes(new string[] { messageId, "[XCCGetHistory]", "[XCCGetHistory]" }.ToXFEString()); await TextMessageClientWebSocket.SendAsync(new ArraySegment(sendBuffer), WebSocketMessageType.Text, true, CancellationToken.None); var endTask = Task.Run(async () => { await Task.Delay(5000); UpdateTaskTrigger?.Invoke(false, messageId); }); return await new XFEWaitTask(ref UpdateTaskTrigger, messageId); } catch (Exception ex) { throw new XFECyberCommException("客户端获取历史记录时出现异常", ex); } } /// /// 关闭XCC会话 /// /// /// public async Task CloseXCC() { try { readyToClose = true; await TextMessageClientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "客户端主动关闭连接", CancellationToken.None); await FileTransportClientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "客户端主动关闭连接", CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("客户端关闭连接时出现异常", ex); } } #endregion internal XCCGroup(string signature, string groupId, string sender, XCCNetWorkBase xCCNetWorkBase) { Signature = signature; GroupId = groupId; Sender = sender; workBase = xCCNetWorkBase; } } /// /// XCC文件类型 /// public enum XCCFileType { /// /// 图片 /// Image, /// /// 视频 /// Video, /// /// 音频 /// Audio } /// /// 消息接收触发器 /// /// /// 是否为历史消息 /// 消息 public delegate void MessageReceivedHandler(bool isHistory, T message); /// /// XCC消息接收器 /// public class XCCMessageReceiveHelper { private readonly Dictionary xCCFileDictionary = new Dictionary(); private readonly Dictionary> xCCMessageDictionary = new Dictionary>(); private bool loaded = false; /// /// 自动保存到本地 /// public bool AutoSaveInLocal { get; set; } /// /// 保存的根目录 /// public string SavePathRoot { get; set; } /// /// 接收到文件事件 /// public event MessageReceivedHandler FileReceived; /// /// 接收到文本事件 /// public event MessageReceivedHandler TextReceived; /// /// 错误发生事件 /// public event XFEEventHandler ExceptionOccurred; /// /// 接收到实时音频字节流事件 /// public event XFEEventHandler AudioBufferReceived; /// /// 从设置的根目录加载 /// /// public async Task Load() { await Task.Run(() => { if (Directory.Exists(SavePathRoot)) { foreach (var groupIdFullPath in Directory.EnumerateDirectories(SavePathRoot)) { var groupId = Path.GetFileName(groupIdFullPath); if (File.Exists($"{groupIdFullPath}/XFEMessage/XFEMessage.xfe")) { var xCCMessageList = new List(); foreach (var entry in new XFEMultiDictionary(File.ReadAllText($"{groupIdFullPath}/XFEMessage/XFEMessage.xfe"))) { var xCCMessage = XCCMessage.ConvertToXCCMessage(entry.Content, groupId); xCCMessageList.Add(xCCMessage); if (xCCMessage.MessageType == XCCTextMessageType.Text) TextReceived?.Invoke(true, xCCMessage); else FileReceived?.Invoke(true, LoadFile(xCCMessage)); } xCCMessageDictionary.Add(groupId, xCCMessageList); } } } else { Directory.CreateDirectory(SavePathRoot); } }); loaded = true; } /// /// 从设置的根目录的指定群组加载 /// /// 群组ID /// public async Task LoadGroup(string groupId) { await Task.Run(() => { if (File.Exists($"{SavePathRoot}/{groupId}/XFEMessage/XFEMessage.xfe")) { var xCCMessageList = new List(); foreach (var entry in new XFEMultiDictionary(File.ReadAllText($"{SavePathRoot}/{groupId}/XFEMessage/XFEMessage.xfe"))) { var xCCMessage = XCCMessage.ConvertToXCCMessage(entry.Content, groupId); xCCMessageList.Add(xCCMessage); if (xCCMessage.MessageType == XCCTextMessageType.Text) TextReceived?.Invoke(true, xCCMessage); else FileReceived?.Invoke(true, LoadFile(xCCMessage)); } xCCMessageDictionary.Add(groupId, xCCMessageList); } }); loaded = true; } /// /// 清理无用文件 /// /// public async Task ClearUselessFile() { if (!loaded) throw new XFEExtensionException("不能在加载完成前调用清理"); await Task.Run(() => { foreach (var groupId in xCCMessageDictionary.Keys) { foreach (var file in Directory.EnumerateFiles($"{SavePathRoot}/{groupId}")) { var filePath = $"{SavePathRoot}/{groupId}/{file}"; var messageId = Path.GetFileNameWithoutExtension(filePath); if (!xCCMessageDictionary.ContainsKey(groupId) || xCCMessageDictionary[groupId].Find(x => x.MessageId == messageId) is null) { File.Delete(filePath); } } } }); } private XCCFile LoadFile(XCCMessage xCCMessage) { var filePath = $"{SavePathRoot}/{xCCMessage.GroupId}/{xCCMessage.MessageId}.xfe"; byte[] fileBuffer = null; if (File.Exists(filePath)) fileBuffer = File.ReadAllBytes(filePath); XCCFile xCCFile; if (xCCFileDictionary.ContainsKey(xCCMessage.MessageId)) { if (!xCCFileDictionary[xCCMessage.MessageId].Loaded && fileBuffer != null) xCCFileDictionary[xCCMessage.MessageId].LoadFile(fileBuffer); return xCCFileDictionary[xCCMessage.MessageId]; } switch (xCCMessage.MessageType) { case XCCTextMessageType.Image: xCCFile = new XCCFile(xCCMessage.GroupId, xCCMessage.MessageId, XCCFileType.Image, xCCMessage.Sender, xCCMessage.SendTime, fileBuffer); break; case XCCTextMessageType.Audio: xCCFile = new XCCFile(xCCMessage.GroupId, xCCMessage.MessageId, XCCFileType.Audio, xCCMessage.Sender, xCCMessage.SendTime, fileBuffer); break; case XCCTextMessageType.Video: xCCFile = new XCCFile(xCCMessage.GroupId, xCCMessage.MessageId, XCCFileType.Video, xCCMessage.Sender, xCCMessage.SendTime, fileBuffer); break; default: return null; } xCCFileDictionary.Add(xCCMessage.MessageId, xCCFile); return xCCFile; } /// /// 获取文件 /// /// 消息ID /// public XCCFile GetFile(string messageId) { return xCCFileDictionary.ContainsKey(messageId) ? xCCFileDictionary[messageId] : null; } /// /// 添加文件 /// /// XCC文件实例 public void AddFile(XCCFile xCCFile) { xCCFileDictionary.Add(xCCFile.MessageId, xCCFile); if (AutoSaveInLocal && xCCFile.FileBuffer != null) SaveFile(xCCFile); } /// /// 保存文件 /// /// XCC文件实例 public void SaveFile(XCCFile xCCFile) { if (!Directory.Exists($"{SavePathRoot}/{xCCFile.GroupId}")) { Directory.CreateDirectory($"{SavePathRoot}/{xCCFile.GroupId}"); } File.WriteAllBytes($"{SavePathRoot}/{xCCFile.GroupId}/{xCCFile.MessageId}.xfe", xCCFile.FileBuffer); } /// /// 保存群组消息 /// /// public void SaveMessage(string groupId) { var filePath = $"{SavePathRoot}/{groupId}/XFEMessage"; if (!Directory.Exists(filePath)) { Directory.CreateDirectory(filePath); } var storageDictionary = new XFEDictionary(); foreach (var xCCMessage in xCCMessageDictionary[groupId]) { storageDictionary.Add(xCCMessage.MessageId, xCCMessage.ToString()); } File.WriteAllText(filePath + "/XFEMessage.xfe", storageDictionary.ToString()); } private void ReceiveFilePlaceHolder(XCCTextMessageReceivedEventArgs e, XCCFileType fileType) { var xCCFile = new XCCFile(e.GroupId, e.MessageId, fileType, e.Sender, e.SendTime); if (!xCCFileDictionary.ContainsKey(e.MessageId)) { xCCFileDictionary.Add(e.MessageId, xCCFile); } FileReceived?.Invoke(e.IsHistory, xCCFile); if (AutoSaveInLocal) SaveMessage(e.GroupId); } private void ReceiveTextMessage(object sender, XCCTextMessageReceivedEventArgs e) { var message = new XCCMessage(e.MessageId, e.MessageType, e.TextMessage, e.Sender, e.SendTime, e.GroupId); if (xCCMessageDictionary.ContainsKey(e.GroupId)) { if (xCCMessageDictionary[e.GroupId].Find(x => x.MessageId == e.MessageId) is null) { xCCMessageDictionary[e.GroupId].Add(message); } else { return; } } else { xCCMessageDictionary.Add(e.GroupId, new List() { message }); } switch (e.MessageType) { case XCCTextMessageType.Text: TextReceived?.Invoke(e.IsHistory, message); if (AutoSaveInLocal) SaveMessage(e.GroupId); break; case XCCTextMessageType.Image: ReceiveFilePlaceHolder(e, XCCFileType.Image); break; case XCCTextMessageType.Audio: ReceiveFilePlaceHolder(e, XCCFileType.Audio); break; case XCCTextMessageType.Video: ReceiveFilePlaceHolder(e, XCCFileType.Video); break; default: break; } } private void ReceiveBinaryMessage(object sender, XCCBinaryMessageReceivedEventArgs e) { if (e.MessageType == XCCBinaryMessageType.AudioBuffer) { AudioBufferReceived?.Invoke(e.BinaryMessage); return; } if (xCCFileDictionary.ContainsKey(e.MessageId)) { var xCCFile = xCCFileDictionary[e.MessageId]; if (!xCCFile.Loaded) { xCCFileDictionary[e.MessageId].LoadFile(e.BinaryMessage); if (AutoSaveInLocal) SaveFile(xCCFile); } } else { var fileType = XCCFileType.Image; switch (e.MessageType) { case XCCBinaryMessageType.Image: fileType = XCCFileType.Image; break; case XCCBinaryMessageType.Audio: fileType = XCCFileType.Audio; break; case XCCBinaryMessageType.AudioBuffer: break; case XCCBinaryMessageType.Video: fileType = XCCFileType.Video; break; default: break; } var xCCFile = new XCCFile(e.GroupId, e.MessageId, fileType, e.Sender, e.SendTime, e.BinaryMessage); xCCFileDictionary.Add(e.MessageId, xCCFile); if (!e.IsHistory) FileReceived?.Invoke(e.IsHistory, xCCFile); } } private void XCCNetWork_ExceptionMessageReceived(object sender, XCCExceptionMessageReceivedEventArgs e) { ExceptionOccurred?.Invoke(e.Exception); } /// /// XCC消息接收器 /// /// 保存根目录 /// 自动保存 public XCCMessageReceiveHelper(string savePathRoot, bool autoSaveInLocal = true) { AutoSaveInLocal = autoSaveInLocal; SavePathRoot = savePathRoot; } /// /// XCC消息接收器 /// /// 保存根目录 /// XCC网络通讯实例 /// 自动保存 public XCCMessageReceiveHelper(string savePathRoot, XCCNetWork xCCNetWork, bool autoSaveInLocal = true) { AutoSaveInLocal = autoSaveInLocal; SavePathRoot = savePathRoot; xCCNetWork.TextMessageReceived += ReceiveTextMessage; xCCNetWork.BinaryMessageReceived += ReceiveBinaryMessage; xCCNetWork.ExceptionMessageReceived += XCCNetWork_ExceptionMessageReceived; } } /// /// XCC文件 /// public class XCCFile { /// /// 文件加载完成时触发 /// public event XFEEventHandler FileLoaded; /// /// 群组ID /// public string GroupId { get; } /// /// 消息ID /// public string MessageId { get; } /// /// 发送者 /// public string Sender { get; } /// /// 发送时间 /// public DateTime? SendTime { get; } /// /// XCC文件类型 /// public XCCFileType FileType { get; } /// /// 是否已加载 /// public bool Loaded { get; private set; } /// /// 文件流 /// public byte[] FileBuffer { get; set; } /// /// 加载文件 /// /// 文件流 public void LoadFile(byte[] fileBuffer) { FileBuffer = fileBuffer; Loaded = true; FileLoaded?.Invoke(this); } /// /// XCC文件 /// /// 群组ID /// 文件消息ID /// 文件类型 /// 发送者 /// 发送时间 /// 文件的Buffer public XCCFile(string groupId, string messageId, XCCFileType fileType, string sender, DateTime? sendTime, byte[] fileBuffer = null) { GroupId = groupId; MessageId = messageId; FileType = fileType; Sender = sender; SendTime = sendTime; FileBuffer = fileBuffer; Loaded = fileBuffer != null; } } /// /// XCC消息 /// public class XCCMessage { /// /// 消息ID /// public string MessageId { get; } /// /// 消息类型 /// public XCCTextMessageType MessageType { get; } /// /// 消息内容 /// public string Message { get; } /// /// 发送者 /// public string Sender { get; } /// /// 发送时间 /// public DateTime SendTime { get; } /// /// 群组ID /// public string GroupId { get; } /// /// 封装为字符串 /// /// public override string ToString() { return new string[] { MessageId, MessageType.ToString(), Message, Sender, SendTime.ToString() }.ToXFEString(); } /// /// 将封装后的XCC消息字符串转换为XCC消息对象 /// /// /// /// public static XCCMessage ConvertToXCCMessage(string xCCMessageStringFormat, string groupId) { var unPackedMessage = xCCMessageStringFormat.ToXFEArray(); return new XCCMessage(unPackedMessage[0], (XCCTextMessageType)Enum.Parse(typeof(XCCTextMessageType), unPackedMessage[1]), unPackedMessage[2], unPackedMessage[3], DateTime.Parse(unPackedMessage[4]), groupId); } /// /// XCC消息 /// /// 消息ID /// 消息类型 /// 消息内容 /// 发送者 /// 发送时间 /// 群组ID public XCCMessage(string messageId, XCCTextMessageType messageType, string message, string sender, DateTime sendTime, string groupId) { MessageId = messageId; MessageType = messageType; Message = message; Sender = sender; SendTime = sendTime; GroupId = groupId; } } /// /// XCC网络通讯事件 /// public abstract class XCCMessageReceivedEventArgs : EventArgs { /// /// 触发事件的群组 /// public XCCGroup Group { get; } /// /// WebSocket明文传输客户端 /// public ClientWebSocket TextMessageClientWebSocket { get; private set; } /// /// WebSocket文件传输客户端 /// public ClientWebSocket FileTransportClientWebSocket { get; private set; } /// /// XCC服务器连接类型 /// public XCCClientType XCCClientType { get; } /// /// 消息ID /// public string MessageId { get; } /// /// 群组ID /// public string GroupId { get { return Group.GroupId; } } /// /// 发送者 /// public string Sender { get; } /// /// 回复文本消息 /// /// 待发送的文本 /// 发送进程 public async Task ReplyTextMessage(string message) { try { byte[] sendBuffer = Encoding.UTF8.GetBytes(message); await TextMessageClientWebSocket.SendAsync(new ArraySegment(sendBuffer), WebSocketMessageType.Text, true, CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("收到服务器端数据后客户端回复文本数据时出现异常", ex); } } /// /// 回复二进制消息 /// /// 二进制消息 /// public async Task ReplyBinaryMessage(byte[] message) { try { await FileTransportClientWebSocket.SendAsync(new ArraySegment(message), WebSocketMessageType.Binary, true, CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("收到服务器端数据后客户端回复文本数据时出现异常", ex); } } /// /// 关闭连接 /// /// public async Task Close() { try { if (TextMessageClientWebSocket.State == WebSocketState.Open) await TextMessageClientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "客户端主动关闭连接", CancellationToken.None); if (FileTransportClientWebSocket.State == WebSocketState.Open) await FileTransportClientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "客户端主动关闭连接", CancellationToken.None); } catch (Exception ex) { throw new XFECyberCommException("客户端关闭连接时出现异常", ex); } } internal XCCMessageReceivedEventArgs(XCCGroup group, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, XCCClientType xCCClientType, string sender, string messageId) { Group = group; TextMessageClientWebSocket = textMessageClientWebSocket; FileTransportClientWebSocket = fileTransportClientWebSocket; XCCClientType = xCCClientType; Sender = sender; MessageId = messageId; } } /// /// XFE网络通信明文返回消息类型 /// public enum XCCTextMessageType { /// /// 文本消息 /// Text, /// /// 图片消息 /// Image, /// /// 音频消息 /// Audio, /// /// 视频消息 /// Video } /// /// XCC网络通讯接收到明文消息事件 /// public abstract class XCCTextMessageReceivedEventArgs : XCCMessageReceivedEventArgs { /// /// 返回文本消息类型 /// public XCCTextMessageType MessageType { get; } /// /// 文本消息 /// public string TextMessage { get; } /// /// 发送时间 /// public DateTime SendTime { get; } /// /// 是否为历史消息 /// public bool IsHistory { get; } internal XCCTextMessageReceivedEventArgs(XCCGroup group, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, XCCClientType xCCClientType, string messageId, XCCTextMessageType messageType, string message, string sender, DateTime sendTime, bool isHistory) : base(group, textMessageClientWebSocket, fileTransportClientWebSocket, xCCClientType, sender, messageId) { MessageType = messageType; TextMessage = message; SendTime = sendTime; IsHistory = isHistory; } } /// /// XFE网络通信二进制返回消息类型 /// public enum XCCBinaryMessageType { /// /// 文本消息 /// Text, /// /// 二进制消息 /// Binary, /// /// 图片消息 /// Image, /// /// 音频消息 /// Audio, /// /// 实时音频 /// AudioBuffer, /// /// 视频消息 /// Video } /// /// XCC网络通讯接收到二进制消息事件 /// public abstract class XCCBinaryMessageReceivedEventArgs : XCCMessageReceivedEventArgs { /// /// 返回二进制消息类型 /// public XCCBinaryMessageType MessageType { get; } /// /// 消息签名 /// public string Signature { get; } /// /// 二进制消息 /// public byte[] BinaryMessage { get; } /// /// 发送时间 /// public DateTime SendTime { get; } /// /// 是否为历史消息 /// public bool IsHistory { get; } internal XCCBinaryMessageReceivedEventArgs(XCCGroup group, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, XCCClientType xCCClientType, string sender, string messageId, byte[] buffer, XCCBinaryMessageType messageType, string signature, DateTime sendTime, bool isHistory) : base(group, textMessageClientWebSocket, fileTransportClientWebSocket, xCCClientType, sender, messageId) { BinaryMessage = buffer; MessageType = messageType; Signature = signature; SendTime = sendTime; IsHistory = isHistory; } } /// /// XCC网络通讯期间异常事件 /// public abstract class XCCExceptionMessageReceivedEventArgs : XCCMessageReceivedEventArgs { /// /// 异常信息 /// public XFECyberCommException Exception { get; } internal XCCExceptionMessageReceivedEventArgs(XCCGroup group, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, XCCClientType xCCClientType, string sender, string messageId, XFECyberCommException exception) : base(group, textMessageClientWebSocket, fileTransportClientWebSocket, xCCClientType, sender, messageId) { Exception = exception; } } /// /// XCC会话关闭事件 /// public abstract class XCCConnectionClosedEventArgs : EventArgs { /// /// 触发事件的群组 /// public XCCGroup Group { get; } /// /// XCC服务器连接类型 /// public XCCClientType XCCClientType { get; } /// /// WebSocket明文传输客户端 /// public ClientWebSocket TextMessageClientWebSocket { get; private set; } /// /// WebSocket文件传输客户端 /// public ClientWebSocket FileTransportClientWebSocket { get; private set; } /// /// 是否正常关闭 /// public bool ClosedNormally { get; } internal XCCConnectionClosedEventArgs(XCCGroup group, XCCClientType xCCClientType, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, bool closedNormally) { Group = group; XCCClientType = xCCClientType; TextMessageClientWebSocket = textMessageClientWebSocket; FileTransportClientWebSocket = fileTransportClientWebSocket; ClosedNormally = closedNormally; } } /// /// XCC连接事件 /// public abstract class XCCConnectedEventArgs : EventArgs { /// /// 触发事件的群组 /// public XCCGroup Group { get; } /// /// XCC服务器连接类型 /// public XCCClientType XCCClientType { get; } /// /// WebSocket明文传输客户端 /// public ClientWebSocket TextMessageClientWebSocket { get; private set; } /// /// WebSocket文件传输客户端 /// public ClientWebSocket FileTransportClientWebSocket { get; private set; } internal XCCConnectedEventArgs(XCCGroup group, XCCClientType xCCClientType, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket) { Group = group; XCCClientType = xCCClientType; TextMessageClientWebSocket = textMessageClientWebSocket; FileTransportClientWebSocket = fileTransportClientWebSocket; } } class XCCGroupImpl : XCCGroup { internal XCCGroupImpl(string signature, string groupId, string sender, XCCNetWorkBase xCCNetWorkBase) : base(signature, groupId, sender, xCCNetWorkBase) { } } class XCCConnectionClosedEventArgsImpl : XCCConnectionClosedEventArgs { /// /// 使用群组、连接类型、传输客户端及关闭状态创建 XCC 会话关闭事件参数。 /// /// 触发关闭事件的 XCC 群组。 /// 关闭连接对应的 XCC 客户端类型。 /// 用于明文消息传输的 WebSocket 客户端。 /// 用于文件传输的 WebSocket 客户端。 /// 连接是否按正常流程关闭。 public XCCConnectionClosedEventArgsImpl(XCCGroup group, XCCClientType xCCClientType, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, bool closedNormally) : base(group, xCCClientType, textMessageClientWebSocket, fileTransportClientWebSocket, closedNormally) { } } class XCCConnectedEventArgsImpl : XCCConnectedEventArgs { /// /// 使用群组、连接类型及传输客户端创建 XCC 连接成功事件参数。 /// /// 触发连接事件的 XCC 群组。 /// 已建立连接的 XCC 客户端类型。 /// 用于明文消息传输的 WebSocket 客户端。 /// 用于文件传输的 WebSocket 客户端。 public XCCConnectedEventArgsImpl(XCCGroup group, XCCClientType xCCClientType, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket) : base(group, xCCClientType, textMessageClientWebSocket, fileTransportClientWebSocket) { } } class XCCTextMessageReceivedEventArgsImpl : XCCTextMessageReceivedEventArgs { internal XCCTextMessageReceivedEventArgsImpl(XCCGroup group, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, XCCClientType xCCClientType, string messageId, XCCTextMessageType messageType, string message, string sender, DateTime sendTime, bool isHistory) : base(group, textMessageClientWebSocket, fileTransportClientWebSocket, xCCClientType, messageId, messageType, message, sender, sendTime, isHistory) { } } class XCCBinaryMessageReceivedEventArgsImpl : XCCBinaryMessageReceivedEventArgs { internal XCCBinaryMessageReceivedEventArgsImpl(XCCGroup group, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, XCCClientType xCCClientType, string sender, string messageId, byte[] buffer, XCCBinaryMessageType messageType, string signature, DateTime sendTime, bool isHistory) : base(group, textMessageClientWebSocket, fileTransportClientWebSocket, xCCClientType, sender, messageId, buffer, messageType, signature, sendTime, isHistory) { } } class XCCExceptionMessageReceivedEventArgsImpl : XCCExceptionMessageReceivedEventArgs { internal XCCExceptionMessageReceivedEventArgsImpl(XCCGroup group, ClientWebSocket textMessageClientWebSocket, ClientWebSocket fileTransportClientWebSocket, XCCClientType xCCClientType, string sender, string messageId, XFECyberCommException exception) : base(group, textMessageClientWebSocket, fileTransportClientWebSocket, xCCClientType, sender, messageId, exception) { } } }