using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Net;
using System.Threading;
using System.Threading.Tasks;
using ActorCore;
using MrWu.Debug;
using NetWorkMessage;
using Server.Core;
namespace Server.Net
{
///
/// 跨进程发送器
///
public class ProcessOuterSender : IProcessOuterSender
{
private TService service;
private uint RpcId;
public Actor Parent { get; private set; }
///
/// 所有的Session
///
private ConcurrentDictionary sessions = new ConcurrentDictionary();
///
/// 获取ActorId对应的地址
///
public Func GetActorIdAddress;
public const int TIMEOUT_TIME = 30 * 1000;
private ConcurrentDictionary requestCallBack =
new ConcurrentDictionary();
//每个网络
private Dictionary> rpcIdHashSets = new Dictionary>();
///
/// 断开链接的Session
///
private ConcurrentQueue disposedSessions = new ConcurrentQueue();
public ProcessOuterSender(Actor actor)
{
Parent = actor;
}
public void Start(string innerIp, int innerPort)
{
Start(new IPEndPoint(IPAddress.Parse(innerIp), innerPort));
}
///
/// 启动服务
///
///
public void Start(IPEndPoint ipEndPoint)
{
service = new TService(ipEndPoint, ServiceType.Actor);
this.service.AcceptCallBack = OnAccept;
this.service.ReadCallBack = OnRead;
this.service.ErrorCallBack = OnError;
}
public void Update()
{
service.Update();
NetDisconnectRpcResponse();
}
public void Send(ActorId fromActorId, ActorId actorId, IMessage messageObject)
{
this.SendInner(fromActorId, actorId, (MessageObject)messageObject);
}
public uint GetRpcId()
{
return ++RpcId;
}
private void NetDisconnectRpcResponse()
{
int dissCnt = disposedSessions.Count;
while (dissCnt -- > 0)
{
if (!disposedSessions.TryDequeue(out long sessionId))
{
break;
}
if (rpcIdHashSets.TryGetValue(sessionId, out HashSet rpcIds))
{
//最多处理100个
int rpcHandleCnt = Math.Min(rpcIds.Count, 100);
foreach (uint rpcId in rpcIds)
{
if (rpcHandleCnt-- <= 0)
{
break;
}
if (!requestCallBack.TryRemove(rpcId, out MessageSenderStruct action))
{
continue;
}
Debug.Info($"网络断开,回复RPC! {rpcId} {Thread.CurrentThread.ManagedThreadId}");
action.SetResult(MessageHelper.CreateResponse(action.RequestType, rpcId, NetErrorCode.ERR_ActorSessionDisconnect));
}
if (rpcIds.Count > 0)
{
disposedSessions.Enqueue(sessionId);
}
}
}
}
public async Task Call(ActorId fromActorId, ActorId actorId, IRequest request)
{
if (actorId == default)
{
throw new Exception($"actor id is 0: {request}");
}
uint rpcId = GetRpcId();
request.RpcId = rpcId;
IResponse response;
Type requestType = request.GetType();
MessageSenderStruct messageSenderStruct = new MessageSenderStruct(actorId, requestType);
if (!this.requestCallBack.TryAdd(rpcId, messageSenderStruct))
{
Debug.Error("ProcessOuterSender RpcId 竟然有重复的现象!!!");
response = MessageHelper.CreateResponse(requestType, rpcId, NetErrorCode.ERR_RpcIdAddFail);
return response;
}
long sessionInstanceId = SendInner(fromActorId, actorId, request as MessageObject);
//Debug.Info($"ProcessOuterSender Rpc Add:{rpcId},sessionInstanceId:{sessionInstanceId},threadId:{Thread.CurrentThread.ManagedThreadId}");
if (sessionInstanceId < 0)
{
Debug.Info("没发送成功!");
response = MessageHelper.CreateResponse(requestType, rpcId, NetErrorCode.ERR_NotFoundActorSession);
return response;
}
if (rpcIdHashSets.ContainsKey(sessionInstanceId))
{
rpcIdHashSets[sessionInstanceId].Add(rpcId);
}
else
{
rpcIdHashSets.Add(sessionInstanceId, new HashSet { rpcId });
}
async Task TimeOut()
{
await Parent.WaitAsync(TIMEOUT_TIME);
if (!requestCallBack.TryRemove(rpcId, out MessageSenderStruct action))
{
return;
}
Debug.Info("网络RPC发送超时");
IResponse response = MessageHelper.CreateResponse(requestType, rpcId, NetErrorCode.ERR_Timeout);
action.SetResult(response);
}
_ = TimeOut();
long beginTime = TimeInfo.Instance.ServerFrameTime();
response = await messageSenderStruct.Wait();
if (rpcIdHashSets.ContainsKey(sessionInstanceId))
{
rpcIdHashSets.Remove(rpcId);
}
long endTime = TimeInfo.Instance.ServerFrameTime();
long costTime = endTime - beginTime;
if (costTime > 200)
{
Debug.Warning($"actor rpc time > 200: {costTime} {requestType.FullName}");
}
Debug.Info("返回数据!");
return response;
}
private long SendInner(ActorId fromActorId, ActorId actorId, MessageObject messageObject)
{
if (actorId == default)
{
throw new Exception($"actor id is 0:{messageObject}");
}
//验证
ActorSession session = Get(actorId);
if (session == null)
{
Debug.Error($"没找到ActorSession:{actorId}");
return -1;
}
session.Send(fromActorId, actorId, messageObject);
return session.InstanceId;
}
#region 网络事件
private void OnAccept(long channelId, IPEndPoint ipEndPoint)
{
ActorSession actorSession = AddSession(channelId);
actorSession.RemoteAddress = ipEndPoint;
Debug.Log("OnAccept!");
}
private void OnRead(long channelId, MemoryBuffer memoryBuffer)
{
Debug.Log("OnRead xx");
ActorSession session = Get(channelId);
if (session == null)
{
Debug.Log("session is null");
return;
}
session.LastRecvTime = TimeInfo.Instance.ClientFrameTime();
(ActorId fromActorId, ActorId actorId, IMessage message) =
MessageSerializeHelper.ToActorMessage(service, memoryBuffer);
if (message == null)
{
Debug.Log("message is null!");
return;
}
if (message is IResponse response)
{
this.HandleIActorResponse(response);
return;
}
switch (message)
{
case IRequest:
{
_ = CallInner();
break;
async Task CallInner()
{
IRequest request = (IRequest)message;
uint rpdId = request.RpcId;
// 注意这里都不能抛异常,因为这里只是中转消息
IResponse res = await Parent.Call(actorId, request);
res.RpcId = rpdId;
//Debug.Info($"回复消息:{res.GetType()}");
Send(actorId, fromActorId, res);
((MessageObject)res).Dispose();
}
}
default:
Parent.Send(actorId, (MessageObject)message);
break;
}
}
private void OnError(long channelId, int error)
{
Debug.Info($"ProcessOuterSender {Thread.CurrentThread.ManagedThreadId} OnError:{error}");
ActorSession session = Get(channelId);
if (session == null)
{
return;
}
disposedSessions.Enqueue(session.InstanceId);
session.Error = error;
session.Dispose();
}
#endregion
public void Dispose()
{
Debug.Log("ProcessOuterSender Dispose!");
this.service.Dispose();
}
private ActorSession Get(ActorId actorId)
{
long channelId = actorId.GetChannelId();
if (this.sessions.TryGetValue(channelId, out ActorSession session))
{
return session;
}
IPEndPoint ipEndPoint = GetActorIdAddress(actorId);
if (ipEndPoint == null)
{
Debug.Error($"未找到 Actor 信息! {actorId}");
return null;
}
session = CreateInner(channelId, ipEndPoint);
return session;
}
private ActorSession Get(long channelId)
{
if (this.sessions.TryGetValue(channelId, out ActorSession session))
{
return session;
}
return null;
}
private void HandleIActorResponse(IResponse response)
{
Debug.Log($"ProcessOuterSender RpcId:{response.RpcId}");
if (!this.requestCallBack.TryRemove(response.RpcId, out MessageSenderStruct messageSenderStruct))
{
return;
}
Run(messageSenderStruct, response);
}
private void Run(MessageSenderStruct self, IResponse response)
{
self.SetResult(response);
}
private ActorSession CreateInner(long channelId, IPEndPoint ipEndPoint)
{
ActorSession session = AddSession(channelId);
session.RemoteAddress = ipEndPoint;
service.Create(channelId, session.RemoteAddress);
return session;
}
private ActorSession AddSession(long channelId)
{
ActorSession session = new ActorSession(channelId, service, SessionDisposed);
if (!this.sessions.TryAdd(channelId, session))
{
Debug.Error($"AddSession Session alreadyExists SessionId:{session.Id}");
return null;
}
return session;
}
private void SessionDisposed(long channelId)
{
if (!this.sessions.TryRemove(channelId, out ActorSession _))
{
Debug.Error($"RemoveSession Session not exists SessionId:{channelId}");
}
}
}
}