Files
hjha-server/ServerCore/NetWork/ProcessOuterSender.cs
xiaoou e9616125ce feat: initial commit - HJHA game server full source
6 major server modules (PdkFriendServer/GlobalSever/ServerCore/GameModule/GameNetModule) +
game logic (GameFix/GameDAL/ServerData) +
network layer (NetWorkMessage) +
data layer (ObjectModel) +
utilities (MrWu/Core/Config/CloudAPI/dll) +
adapters (zyxAdapter/base)

.NET 8.0 C# solution, 16 projects, 958 source files
2026-07-07 12:02:15 +08:00

376 lines
12 KiB
C#
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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
{
/// <summary>
/// 跨进程发送器
/// </summary>
public class ProcessOuterSender : IProcessOuterSender
{
private TService service;
private uint RpcId;
public Actor Parent { get; private set; }
/// <summary>
/// 所有的Session
/// </summary>
private ConcurrentDictionary<long, ActorSession> sessions = new ConcurrentDictionary<long, ActorSession>();
/// <summary>
/// 获取ActorId对应的地址
/// </summary>
public Func<ActorId, IPEndPoint> GetActorIdAddress;
public const int TIMEOUT_TIME = 30 * 1000;
private ConcurrentDictionary<uint, MessageSenderStruct> requestCallBack =
new ConcurrentDictionary<uint, MessageSenderStruct>();
//每个网络
private Dictionary<long,HashSet<uint>> rpcIdHashSets = new Dictionary<long, HashSet<uint>>();
/// <summary>
/// 断开链接的Session
/// </summary>
private ConcurrentQueue<long> disposedSessions = new ConcurrentQueue<long>();
public ProcessOuterSender(Actor actor)
{
Parent = actor;
}
public void Start(string innerIp, int innerPort)
{
Start(new IPEndPoint(IPAddress.Parse(innerIp), innerPort));
}
/// <summary>
/// 启动服务
/// </summary>
/// <param name="ipEndPoint"></param>
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<uint> 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<IResponse> 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<uint> { 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}");
}
}
}
}