EndpointHandler.cs 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177
  1. using System;
  2. using System.Collections.Generic;
  3. using System.IO;
  4. using System.Text;
  5. using Taobao.Top.Link.Channel;
  6. using Top.Api;
  7. namespace Taobao.Top.Link.Endpoints
  8. {
  9. /// <summary>deal with protocol/callback/send
  10. /// </summary>
  11. public class EndpointHandler
  12. {
  13. private ITopLogger _log;
  14. private int _flag;
  15. //all connect in/out endpoints
  16. private IDictionary<string, Identity> _idByToken;
  17. private IDictionary<int, SendCallback> _callbacks;
  18. private EventHandler<ChannelContext> _onMessage;
  19. public Action<EndpointContext> MessageHandler { get; set; }
  20. public OnAckMessage AckMessageHandler { get; set; }
  21. public EndpointHandler(ITopLogger logger)
  22. {
  23. this._log = logger;
  24. this._idByToken = new Dictionary<string, Identity>();
  25. this._callbacks = new Dictionary<int, SendCallback>();
  26. this._onMessage = new EventHandler<ChannelContext>(this.OnMessage);
  27. }
  28. public void Send(Message message, IChannelSender sender)
  29. {
  30. this.Send(message, sender, null);
  31. }
  32. public void Send(Message message, IChannelSender sender, SendCallback callback)
  33. {
  34. if (callback != null)
  35. {
  36. message.Flag = System.Threading.Interlocked.Increment(ref this._flag);
  37. this._callbacks.Add(message.Flag, callback);
  38. }
  39. using (var s = new MemoryStream())
  40. {
  41. MessageIO.WriteMessage(s, message);
  42. this.GetChannel(sender).Send(s.ToArray());
  43. }
  44. }
  45. internal IDictionary<string, object> SendAndWait(EndpointProxy e
  46. , IChannelSender sender
  47. , Message message
  48. , int timeout)
  49. {
  50. SendCallback callback = new SendCallback(e);
  51. this.Send(message, sender, callback);
  52. callback.WaitReturn(timeout);
  53. if (callback.Error != null)
  54. throw callback.Error;
  55. return callback.Return;
  56. }
  57. private IClientChannel GetChannel(IChannelSender sender)
  58. {
  59. var channel = sender as IClientChannel;
  60. if (channel.OnMessage == null)
  61. channel.OnMessage = this._onMessage;
  62. return channel;
  63. }
  64. private void OnMessage(object sender, ChannelContext ctx)
  65. {
  66. Message msg = MessageIO.ReadMessage(new MemoryStream((byte[])ctx.Message));
  67. SendCallback callback = this._callbacks.ContainsKey(msg.Flag)
  68. ? this._callbacks[msg.Flag]
  69. : null;
  70. if (msg.MessageType == MessageType.CONNECTACK)
  71. {
  72. this.HandleConnectAck(callback, msg);
  73. return;
  74. }
  75. Identity msgFrom = msg.Token != null && this._idByToken.ContainsKey(msg.Token)
  76. ? this._idByToken[msg.Token]
  77. : null;
  78. // must CONNECT/CONNECTACK for got token before SEND
  79. if (msgFrom == null)
  80. {
  81. var error = new LinkException(Text.E_UNKNOWN_MSG_FROM);
  82. if (callback == null)
  83. throw error;
  84. callback.Error = error;
  85. return;
  86. }
  87. #region raise callback of client
  88. if (callback != null)
  89. {
  90. this.HandleCallback(callback, msg, msgFrom);
  91. return;
  92. }
  93. else if (this.IsError(msg))
  94. {
  95. this._log.Error(Text.E_GOT_ERROR, msg.StatusCode, msg.StatusPhase);
  96. return;
  97. }
  98. #endregion
  99. #region raise event
  100. if (msg.MessageType == MessageType.SENDACK)
  101. {
  102. if (this.AckMessageHandler != null)
  103. this.AckMessageHandler(msg.Content, msgFrom);
  104. return;
  105. }
  106. if (this.MessageHandler == null)
  107. return;
  108. EndpointContext endpointContext = new EndpointContext(ctx, this, msgFrom, msg.Flag, msg.Token);
  109. endpointContext.Message = msg.Content;
  110. try
  111. {
  112. this.MessageHandler(endpointContext);
  113. }
  114. catch (Exception e)
  115. {
  116. // onMessage error should be reply to client
  117. if (e is LinkException)
  118. endpointContext.Error(
  119. ((LinkException)e).ErrorCode,
  120. ((LinkException)e).Message);
  121. else
  122. endpointContext.Error(0, e.Message);
  123. }
  124. #endregion
  125. }
  126. private void HandleConnectAck(SendCallback callback, Message msg)
  127. {
  128. if (callback == null)
  129. throw new LinkException(Text.E_NO_CALLBACK);
  130. if (this.IsError(msg))
  131. callback.Error = new LinkException(msg.StatusCode, msg.StatusPhase);
  132. else
  133. {
  134. callback.Return = null;
  135. // set token for proxy for sending message next time
  136. callback.Target.Token = msg.Token;
  137. // store token from target endpoint for receiving it's message
  138. // next time
  139. if (this._idByToken.ContainsKey(msg.Token))
  140. this._idByToken[msg.Token] = callback.Target.Identity;
  141. else
  142. this._idByToken.Add(msg.Token, callback.Target.Identity);
  143. this._log.Info(Text.E_CONNECT_SUCCESS, callback.Target.Identity, msg.Token);
  144. }
  145. }
  146. private void HandleCallback(SendCallback callback, Message msg, Identity msgFrom)
  147. {
  148. if (!callback.Target.Identity.Equals(msgFrom))
  149. {
  150. this._log.Warn(Text.E_IDENTITY_NOT_MATCH_WITH_CALLBACK, msgFrom, callback.Target.Identity);
  151. return;
  152. }
  153. if (this.IsError(msg))
  154. callback.Error = new LinkException(msg.StatusCode, msg.StatusPhase);
  155. else
  156. callback.Return = msg.Content;
  157. }
  158. private bool IsError(Message msg)
  159. {
  160. return msg.StatusCode > 0 || !string.IsNullOrEmpty(msg.StatusPhase);
  161. }
  162. public delegate void OnAckMessage(IDictionary<string, object> message, Identity messageFrom);
  163. }
  164. }