using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using NetWorkMessage; using Server.Core; using MrWu.Debug; namespace ActorCore { public struct MessageInfo { /// /// 谁发来的 /// public ActorId FormActorId; /// /// 包体 /// public MessageObject MessageObject; } public class MailBox : IMailBox { private readonly Actor actor; private ActorId actorId => actor.ActorId; public MailBoxType MailBoxType => actor.MailBoxType; public int HandleMaxMessageCount => actor.HandleMaxMessageCount; private int actorTypeId => actor.ActorId.ActorTypeId; //这里可以改成普通字典,这里没有线程安全问题 private readonly ConcurrentDictionary requestCallback = new ConcurrentDictionary(); private readonly Dictionary actorMessageHandle = new Dictionary(); private readonly Dictionary actorRpcHandles = new Dictionary(); private readonly List list = new List(); public void RegisterMessageHandle(Type type, ActorHandle handle) { this.actorMessageHandle[type] = handle; } public void RegisterRPCHandle(Type type, ActorRpcHandle rpcHandle) { this.actorRpcHandles[type] = rpcHandle; } public MailBox(Actor actor) { this.actor = actor; } public void Update() { HandleActorMessage(); } /// /// 处理Actor消息 /// private void HandleActorMessage() { list.Clear(); ActorMessageQueue.Instance.Fetch(this.actorTypeId, HandleMaxMessageCount, list); // if (this.actorTypeId == ActorTypeId.NetInner) // { // Debug.Log($"当前数量:{list.Count}"); // } foreach (MessageInfo messageInfo in list) { //执行Actor消息 _ = HandleActorMessage(messageInfo.FormActorId,messageInfo.MessageObject); } } private async Task HandleActorMessage(ActorId fromActorId, MessageObject messageObject) { //Debug.Info($"HandleActorMessage {messageObject.GetType()}"); if (messageObject is IResponse response) { HandleIActorResponse(response); return; } if (MailBoxType == MailBoxType.OrderedMessage) { using (await actor.Wait(CoroutineLockType.MailBox, actorTypeId)) { try { //Debug.Log($"等待锁! {Thread.CurrentThread.ManagedThreadId}"); await _HandleActorMessage(fromActorId, messageObject); //Debug.Log($"结束锁! {Thread.CurrentThread.ManagedThreadId}"); } catch (Exception e) { Debug.Error("MailBox:" + e); throw e; } } } else { await _HandleActorMessage(fromActorId, messageObject); } } private async Task _HandleActorMessage(ActorId fromActorId,MessageObject messageObject) { switch (messageObject) { case IRequest request: await HandleIActorRequest(fromActorId,request); break; default: await HandleIActorMessage(fromActorId,messageObject); break; } //执行完就回收消息 - 如果消息不想被回收 那就不要用池 ReferencePool.Recycle(messageObject); } /// /// RPC 请求消息 /// /// > /// private async Task HandleIActorRequest(ActorId fromActorId,IRequest request) { //等待处理 Type type = request.GetType(); IResponse response = MessageHelper.CreateResponse(type,request.RpcId,0); if (!this.actorRpcHandles.TryGetValue(type,out ActorRpcHandle thisHandle)) { Debug.Error($"not found rpc handle:{type} ActorId:{this.actor.ActorId} {this.actorRpcHandles.Count}"); response.Error = NetErrorCode.ERR_NotFoundHandle; } uint rpcId = request.RpcId; if (response.Error == 0) { try { await thisHandle(fromActorId,request,response); } catch (Exception exception) { response.Error = NetErrorCode.ERR_RPCFail; response.Message = exception.ToString(); Debug.Log($"exception:{exception}"); } } response.RpcId = rpcId; Debug.Log($"Reply:{fromActorId} {response.GetType()}"); this.actor.Reply(fromActorId,response); } /// /// RPC 响应消息 /// /// private void HandleIActorResponse(IResponse response) { //Debug.Log($"HandleIActorResponse RpcId:{response.RpcId}"); if (!this.requestCallback.TryRemove(response.RpcId,out MessageSenderStruct actorMessageSender)) { Debug.Warning($"ActorResponse: {response.GetType()}"); return; } Run(actorMessageSender,response); } private void Run(MessageSenderStruct self,IResponse response) { self.SetResult(response); } /// /// 普通消息 /// /// /// private async Task HandleIActorMessage(ActorId fromActorId,MessageObject messageObject) { Type type = messageObject.GetType(); if (!this.actorMessageHandle.TryGetValue(type, out ActorHandle thisHandle)) { Debug.Error($"not found message handle:{type} ActorId:{this.actor.ActorId}"); return; } //todo 等待处理 //此处要加上MailType的判断,是否要等待上一个消息处理完成 //Debug.Log($"---- id:{this.ActorId.ActorTypeId}"); await thisHandle(fromActorId,messageObject); } public bool AddActorRpcResponse(uint rpcId, MessageSenderStruct messageSenderStruct) { return this.requestCallback.TryAdd(rpcId, messageSenderStruct); } public bool RemoveActorRpcResponse(uint rpcId,out MessageSenderStruct messageSenderStruct) { return this.requestCallback.TryRemove(rpcId, out messageSenderStruct); } public void Send(ActorId target, MessageObject messageObject) { if (target.Process == this.actorId.Process) { //发给内部Actor ActorMessageQueue.Instance.Send(this.actorId, target, messageObject); return; } //发给外部Actor NetInnerMessage netInnerMessage = NetInnerMessage.Create(); netInnerMessage.ActorId = target; netInnerMessage.MessageObject = messageObject; ActorMessageQueue.Instance.Send(this.actorId, new ActorId(this.actorId.Process, ActorTypeId.NetInner), netInnerMessage); } } }