using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.IO;
using System.Net;
using GameMessage;
using K4os.Compression.LZ4;
using MessagePack;
using MrWu.Debug;
using NetWorkMessage;
using Server.Core;
using UnityGame;
namespace Server.Net
{
///
/// 用户链接管理 这个是管理新的用户连接 子线程
///
public class UserNetSessionManager : IDisposable
{
private TService Service { get; set; }
///
/// 60s 超时断开 25秒一个心跳包
///
private long TimeOut = 60 * 1000;
///
/// 所有Session
///
private Dictionary Sessions = new Dictionary();
private readonly Queue SessionIds = new Queue();
public int ConnectCnt;
public void Start(string innerIp, int innerPort)
{
Start(new IPEndPoint(IPAddress.Parse(innerIp), innerPort));
}
public void Start(IPEndPoint ipEndPoint)
{
Service = new TService(ipEndPoint, ServiceType.Router);
Debug.Log($"启动监听:{ipEndPoint}");
this.Service.AcceptCallBack = OnAccept;
this.Service.ReadCallBack = OnRead;
this.Service.ErrorCallBack = OnError;
}
public void Log()
{
Debug.ImportantLog($"当前用户链接数:{Sessions.Count}");
Service.Log();
}
public void Update()
{
CheckTimeOut();
Service.Update();
}
public void Dispose()
{
Service.Dispose();
}
private void SendMessage(long id, MessageObject message)
{
(uint opcode, MemoryBuffer memoryBuffer) =
MessageSerializeHelper.ToRouterMemoryBuffer(Service, message);
memoryBuffer.Seek(0, SeekOrigin.Begin);
Service.Send(id, memoryBuffer);
message.Dispose();
}
#region 网络事件
private void OnAccept(long channelId, IPEndPoint ipEndPoint)
{
Debug.Log($"User OnAccept! {System.Threading.Thread.CurrentThread.ManagedThreadId}");
//userSessionEvents.Enqueue(UserSessionEvent.Create(channelId, ipEndPoint));
UserNetSession userSession = new UserNetSession(channelId, Service, SessionDisposed);
userSession.LastRecvTime = TimeInfo.Instance.ClientFrameTime();
if (this.Sessions.ContainsKey(channelId))
{
Debug.Error($"AddSession UserSession alreadyExists SessionId:{userSession.Id}");
return;
}
Debug.Info($"用户连接:{System.Threading.Thread.CurrentThread.ManagedThreadId}");
this.Sessions.Add(channelId, userSession);
ConnectCnt++;
userSession.RemoteAddress = ipEndPoint;
SessionIds.Enqueue(channelId);
}
public void ForceDisConnect(long id,long instanceId,int finCode)
{
UserNetSession session = Get(id);
if (session == null || session.InstanceId != instanceId)
{
Debug.Info("强制玩家断开时,链接已经不在!");
return;
}
session.ForceDisConnectCode = finCode;
if (finCode == 0)
{
TChannel channel = Service.Get(id);
channel.Dispose();
}
else
{
UserSessionFin(session, false);
}
}
///
/// 用户断开链接
///
///
/// 是否是主动断开
private void UserSessionFin(UserNetSession session, bool initiative)
{
if (session == null)
{
Debug.Error($"断开时找不到session!");
return;
}
long id = session.Id;
if (session.InstanceId > 0)
{
UserSessionManager.Instance.AddUserSessionEvt(UserSessionEvt.Create(session.Id, session.InstanceId, UserSessionEvtType.Fin));
}
//Debug.Info($"玩家断开:{id} {initiative} {session.State}");
if (!initiative)
{
//已经处理过异常断开
if (session.IsDisConnected)
{
return;
}
session.IsDisConnected = true;
}
if (session.State == UserSessionState.Using)
{
{
//置为空
session.InstanceId = 0;
//Debug.Info($"发送断开原因! {session.ForceDisConnectCode}");
session.State = UserSessionState.Idle;
//发送确认断开
RouterFinAckMessage ackMessage = RouterFinAckMessage.Create(session.ForceDisConnectCode);
SendMessage(session.Id, ackMessage);
ackMessage.Dispose();
}
}
if (!initiative) //只要不是主动断开,那就是真断开了
{
session.State = UserSessionState.Fin;
ConnectCnt--;
if (!this.Sessions.Remove(id))
{
Debug.Error($"RemoveSession UserSession not exists SessionId:{id}");
}
}
else
{
}
}
//单线程
private void OnRead(long channelId, MemoryBuffer memoryBuffer)
{
// 解包,把包数据发给主线程处理
UserNetSession session = Get(channelId);
if (session == null)
{
Debug.Error($"收到包,未找到链接,链接可能已经断开!{channelId}");
return;
}
session.LastRecvTime = TimeInfo.Instance.ClientFrameTime();
//Debug.Log($"收到包:{session.LastRecvTime} {System.Threading.Thread.CurrentThread.ManagedThreadId}");
(uint opcode, IMessage message) = MessageSerializeHelper.ToRouterMessage(Service, memoryBuffer);
Service.Recycle(memoryBuffer);
if (message == null)
{
Debug.Error("message is null!");
return;
}
switch (opcode)
{
case MessageOpcode.Client2ServerMessage:
OnOuterMessage(session, message);
break;
case MessageOpcode.RouterHeart:
OnRouterHeartMessage(session, message);
break;
//请求连接
case MessageOpcode.RouterSyn:
OnRouterSynMessage(session, message);
break;
// case MessageOpcode.RouterAck:
// OnRouterAckMessage(session, message);
// break;
case MessageOpcode.RouterFin:
OnRouterFinMessage(session, message);
break;
// case MessageOpcode.RouterFinAck:
// OnRouterFinAckMessage(session, message);
// break;
default:
Debug.Warning($"收到未知的包类型:{opcode}");
break;
}
}
private void OnError(long channelId, int error)
{
Debug.Info($"User OnError :{error}");
UserNetSession userSession = Get(channelId);
if (userSession == null)
{
return;
}
userSession.Error = error;
userSession.Dispose();
}
#endregion
private UserNetSession Get(long channelId)
{
if (this.Sessions.TryGetValue(channelId, out UserNetSession userSession))
{
return userSession;
}
return null;
}
//销毁事件 这个是多线程的
private void SessionDisposed(long channelId)
{
UserNetSession session = Get(channelId);
UserSessionFin(session, false);
}
//[长度][版本][压缩位][Actor][opcode]
private const int PackIndex = 4 + 1 + 1 + 3 + 4;
private void OnOuterMessage(UserNetSession userSession, IMessage message)
{
if (userSession.ForceDisConnectCode > 0)
{
Debug.Info("这个Session 已经被强制下线!");
return;
}
if (message is Client2ServerMessage outerMessage)
{
userSession.Ver = outerMessage.ver;
//解包
try
{
int opcode = outerMessage.OpCode;
Type messageType = MessageManager.GetType(opcode);
//Debug.Log($"收到包 opcode:{opcode} ver:{outerMessage.ver} compress:{outerMessage.Compress} length:{outerMessage.SourceData.Length}");
ReadOnlyMemory sourceData = null;
if (outerMessage.Compress == 1) //压缩了,需要解压
{
//原始数据长度
//int originLength = BitConverter.ToInt32(outerMessage.SourceData, PackIndex);
int inputLength = outerMessage.SourceData.Length - PackIndex - 4;
byte[] packSourceData =
LZ4Pickler.Unpickle(outerMessage.SourceData, PackIndex + 4, inputLength);
sourceData = new ReadOnlyMemory(packSourceData);
}
else
{
sourceData = new ReadOnlyMemory(outerMessage.SourceData, PackIndex,
outerMessage.SourceData.Length - PackIndex);
}
MessageData messageData = (MessageData)MessagePackSerializer.Deserialize(messageType, sourceData);
//主现成去处理
UserSessionManager.Instance.AddUserSessionEvt(UserSessionEvt.Create(userSession.Id, userSession.InstanceId,GamePacket.Create(opcode, messageData),outerMessage.RemoteIpAddress));
}
catch (Exception e)
{
Debug.Error($"解包报错:{e.Message}");
}
}
}
///
/// 心跳包
///
///
///
private void OnRouterHeartMessage(UserNetSession userSession, IMessage message)
{
if (message is RouterHeart routerHeart)
{
SendMessage(userSession.Id, routerHeart);
}
}
///
/// 链接
///
///
///
private void OnRouterSynMessage(UserNetSession userSession, IMessage message)
{
if (message is RouterSynMessage routerSyn)
{
if (userSession.State != UserSessionState.Idle)
{
Debug.Error("链接正在使用,不能确认连接!");
return;
}
userSession.ConnectAck();
Debug.Log("发送确认连接!");
//实例ID
userSession.State = UserSessionState.Using;
RouterAckMessage ackMessage = RouterAckMessage.Create();
SendMessage(userSession.Id, ackMessage);
//链接完成的事件
UserSessionManager.Instance.AddUserSessionEvt(UserSessionEvt.Create(userSession.Id, userSession.InstanceId, UserSessionEvtType.Connected));
}
}
///
/// 断开
///
///
///
private void OnRouterFinMessage(UserNetSession userSession, IMessage message)
{
if (message is RouterFinMessage routerFin)
{
UserSessionFin(userSession, userSession.ForceDisConnectCode <= 0);
Debug.Info("用户主动断开连接!");
}
}
private void CheckTimeOut()
{
long timeNow = TimeInfo.Instance.FrameTime;
const int MaxCheckNum = 5;
int n = SessionIds.Count;
if (n > MaxCheckNum)
{
n = MaxCheckNum;
}
for (int i = 0; i < n; i++)
{
if (this.Sessions.Count <= 0)
{
break;
}
long sessionId = this.SessionIds.Dequeue();
UserNetSession session = Get(sessionId);
if (session == null)
{
continue;
}
if (timeNow - session.LastRecvTime > TimeOut)
{
Debug.Info($"超时断线:{session.IsDisConnected} {session.LastRecvTime} {timeNow}");
session.Dispose();
continue;
}
this.SessionIds.Enqueue(sessionId);
}
}
public void Send2Router(long id,long instanceId, MessageObject message)
{
UserNetSession userNetSession = Get(id);
if (userNetSession == null || userNetSession.InstanceId != instanceId)
{
Debug.Info("用户已经断开链接,不发送");
return;
}
(uint opcode, MemoryBuffer memoryBuffer) =
MessageSerializeHelper.ToRouterMemoryBuffer(Service, message);
memoryBuffer.Seek(0, SeekOrigin.Begin);
Service.Send(id, memoryBuffer);
message.Dispose();
}
}
}