using System; using System.Collections.Generic; using System.Net; using System.Net.Sockets; using System.Reflection; using System.Linq; using UnityEngine; using UnityEngine.Events; /// /// TCP 通信的服务端 /// public class TcpServer { public event UnityAction connChangeEvent;//客户端连接状态变化事件 ProtocolByte protocol; Socket listenfd; Queue msgQueue = new Queue(); Queue closeConnIPQueue= new Queue();//断开消息记录队列 支线程触发的结束事件无法直接用于驱动Unity 变化,采用支线程存入断开的IP 主线程循环取出的方式处理 Queue newConnIPQueue = new Queue();//建立链接消息队列 支线程触发的链接事件无法直接用于驱动Unity 变化,采用支线程存入连接的IP 主线程循环取出的方式处理 public TcpServer() { } //玩家协议处理类 HandleServerMsg handlePlayerMsg = new HandleServerMsg(); int heartBeatTime = 20; //连接池 Conn[] conns; Dictionary connIndexDic = new Dictionary();//记录地址与链接池之间的关系 //最大挂起客户端数量 int maxConn = 40; //服务端启动调用方法 public void Start(string host, int port) { protocol = new ProtocolByte(); conns = new Conn[maxConn]; for (int i = 0; i < conns.Length; i++) { conns[i] = new Conn(); } listenfd = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); IPAddress dress = IPAddress.Parse(host); IPEndPoint endPoint = new IPEndPoint(dress, port); listenfd.Bind(endPoint); listenfd.Listen(maxConn); listenfd.BeginAccept(new AsyncCallback(AcceptCb), listenfd); } //循环提取数据并处理 public void ExtractingMsg() { int count = 0; while (msgQueue.Count > 0) { HandleMsg(msgQueue.Dequeue()); count++; } if (checkHeartTime > 0) { checkHeartTime -= Time.deltaTime; } else { checkHeartTime = 20; CheckBeat(); } while (closeConnIPQueue.Count > 0) { connChangeEvent?.Invoke(closeConnIPQueue.Dequeue(),false); } while (newConnIPQueue.Count>0) { connChangeEvent?.Invoke(newConnIPQueue.Dequeue(),true); } } //发送消息接口 public void Send(string ip, ProtocolBase proto) { if (connIndexDic.ContainsKey(ip)) Send(conns[connIndexDic[ip]], proto); } //广播消息接口 public void Broadcast(ProtocolBase proto) { for (int i = 0; i < conns.Length; i++) { if (conns[i] == null) continue; if (!conns[i].isUse) continue; // MessageBox.Show("发送广播消息至:" + conns[i].GetEndIPDress()); Send(conns[i], proto); } } public void Close() { for (int i = 0; i < conns.Length; i++) { Conn conn = conns[i]; if (conn == null) continue; if (!conn.isUse) continue; lock (conn) { conn.Close(); } } listenfd.Close(); } private float checkHeartTime = 20; //某端口断开时执行 private void ConnClose(string ip) { closeConnIPQueue.Enqueue(ip); } //组装消息并发送至固定conn端 void Send(Conn conn, ProtocolBase proto) { byte[] bytes = proto.Encode(); byte[] length = BitConverter.GetBytes(bytes.Length); byte[] sendBuff = length.Concat(bytes).ToArray(); try { conn.socket.BeginSend(sendBuff, 0, sendBuff.Length, SocketFlags.None, null, null); } catch (Exception e) { //MessageBox.Show("发送信息失败"); Debug.LogError("发送错误:" + e.Message); } } //获取空闲conn索引 private int GetIndex() { if (conns == null) return -1; for (int i = 0; i < conns.Length; i++) { if (conns[i] == null) { conns[i] = new Conn(); return i; } else if (!conns[i].isUse) return i; } return -1; } //检查链接的心跳 超时断开 private void CheckBeat() { //Debug.Log("心跳"); long timeNow = SysTime.GetTimeStamp(); for (int i = 0; i < conns.Length; i++) { Conn conn = conns[i]; if (conn == null) continue; if (!conn.isUse) continue; if (conn.lastTickTime < timeNow - heartBeatTime) { Debug.Log("[心跳引起断开连接]" + conn.GetEndIPDress()); lock (conn) conn.Close(); } } } //接收连接请求的回调 private void AcceptCb(IAsyncResult ar) { try { Socket socket = listenfd.EndAccept(ar); int index = GetIndex(); if (index < 0) { socket.Close(); //MessageBox.Show("无可用连接"); } else { Conn conn = conns[index]; conn.Init(socket, ConnClose); string adr = conn.GetEndIPDress(); connIndexDic[adr] = index; newConnIPQueue.Enqueue(adr); conn.socket.BeginReceive(conn.readBuffer, conn.bufferCount, conn.BuffRemain(), SocketFlags.None, ReceiveCb, conn); //Debug.Log("地址为:" + adr + " 的客户端连接成功" + "\r\n" + "连接占用挂起索引:" + index); } listenfd.BeginAccept(AcceptCb, null); } catch (Exception) { //MessageBox.Show("Accept失败"); } } //接收客户端数据的回调 private void ReceiveCb(IAsyncResult ar) { Conn conn = (Conn)ar.AsyncState; lock (conn) { try { int count = conn.socket.EndReceive(ar); if (count > Conn.BufferSize) { Debug.LogError($"数据包体长度{count}超出缓冲存储区域,需检查消息是否有误或修改Conn.BufferSize扩大缓冲区"); return; } if (count <= 0) { conn.Close(); return; } conn.bufferCount += count; ProcessData(conn); //处理完该条消息,继续接收回调 conn.socket.BeginReceive(conn.readBuffer, conn.bufferCount, conn.BuffRemain(), SocketFlags.None, ReceiveCb, conn); } catch (Exception ex) { string msg = ex.Message; conn.Close(); } } } //处理收到的信息 private void ProcessData(Conn conn) { //Debug.Log("处理一条消息"); if (conn.bufferCount < sizeof(Int32)) return; Array.Copy(conn.readBuffer, conn.lenBytes, sizeof(Int32)); conn.msgLength = BitConverter.ToInt32(conn.lenBytes, 0); if (conn.msgLength < 0) { Debug.Log($"数据消息长度解析失败,放弃本条消息"); conn.bufferCount = 0; return; } if (conn.msgLength + sizeof(Int32) > Conn.BufferSize) { Debug.LogError($"数据消息长度{conn.msgLength}超出缓冲存储区域,需重新设计传输数据结构或修改Conn.BufferSize扩大缓冲区"); conn.bufferCount = 0; return; } if (conn.bufferCount < conn.msgLength + sizeof(Int32)) { Debug.Log("消息分包,需等待后半部分数据"); return; } ProtocolByte proto = (ProtocolByte)protocol.Decode(conn.readBuffer, sizeof(Int32), conn.msgLength); ClientMsg clientMsg = new ClientMsg(conn.GetEndIPDress(), proto); msgQueue.Enqueue(clientMsg); // HandleMsg(conn,proto); //清除已处理的消息 int count = conn.bufferCount - sizeof(Int32) - conn.msgLength; Array.Copy(conn.readBuffer, sizeof(Int32) + conn.msgLength, conn.readBuffer, 0, count); conn.bufferCount = count; if (conn.bufferCount > 0) ProcessData(conn); } private void HandleMsg(ClientMsg msg) { string name = msg.protocolByte.GetProtocolName(); //Debug.Log("接到协议:" + name); string methodName = "Msg" + name; if (name == "HeatBeat") { if (connIndexDic.ContainsKey(msg.ip)) conns[connIndexDic[msg.ip]].lastTickTime = SysTime.GetTimeStamp(); } else { MethodInfo mm = handlePlayerMsg.GetType().GetMethod(methodName); if (mm == null) { return; } object[] obj = new object[] { msg }; mm.Invoke(handlePlayerMsg, obj); //try //{ //} //catch //{ } // MessageBox.Show("playerMsg已处理方法:" + methodName); } } } //客户端消息类 存储IP和消息 public class ClientMsg { public string ip; public ProtocolByte protocolByte; public ClientMsg(string ip,ProtocolByte protocolByte) { this.ip = ip; this.protocolByte = protocolByte; } }