EndpointProxy.cs 5.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154
  1. using System;
  2. using System.Collections.Generic;
  3. using System.Text;
  4. using Taobao.Top.Link.Channel;
  5. namespace Taobao.Top.Link.Endpoints
  6. {
  7. /// <summary>logic endpoint local proxy object
  8. /// </summary>
  9. public class EndpointProxy
  10. {
  11. private IList<IChannelSender> _senders;
  12. private Random _random;
  13. private EndpointHandler _handler;
  14. /// <summary>get id
  15. /// </summary>
  16. public Identity Identity { get; internal set; }
  17. /// <summary>known by both side, like a sessionId
  18. /// </summary>
  19. public string Token { get; internal set; }
  20. public EndpointProxy(EndpointHandler handler)
  21. {
  22. this._senders = new List<IChannelSender>();
  23. this._random = new Random();
  24. this._handler = handler;
  25. }
  26. internal void Add(IChannelSender sender)
  27. {
  28. if (!this._senders.Contains(sender))
  29. lock (this._senders)
  30. if (!this._senders.Contains(sender))
  31. this._senders.Add(sender);
  32. }
  33. internal void Remove(IChannelSender sender)
  34. {
  35. lock (this._senders)
  36. this._senders.Remove(sender);
  37. }
  38. internal void Remove(string uri)
  39. {
  40. lock (this._senders)
  41. {
  42. for (var i = 0; i < this._senders.Count; i++)
  43. {
  44. var channel = this._senders[i] as IClientChannel;
  45. if (channel == null)
  46. continue;
  47. if (!channel.Uri.ToString().Equals(uri))
  48. continue;
  49. this._senders.Remove(channel);
  50. i--;
  51. }
  52. }
  53. }
  54. internal void Close(string uri, string reason)
  55. {
  56. lock (this._senders)
  57. {
  58. for (var i = 0; i < this._senders.Count; i++)
  59. {
  60. var channel = this._senders[i] as IClientChannel;
  61. if (channel == null)
  62. continue;
  63. if (!channel.Uri.ToString().Equals(uri))
  64. continue;
  65. channel.Close(reason);
  66. this._senders.Remove(channel);
  67. i--;
  68. }
  69. }
  70. }
  71. /// <summary>check is there any sender can be used to send
  72. /// </summary>
  73. /// <returns></returns>
  74. public bool hasValidSender()
  75. {
  76. foreach (IClientChannel sender in this._senders)
  77. {
  78. if (sender.IsConnected)
  79. return true;
  80. }
  81. return false;
  82. }
  83. /// <summary>send message and wait reply
  84. /// </summary>
  85. /// <param name="message"></param>
  86. /// <returns></returns>
  87. public IDictionary<string, object> SendAndWait(IDictionary<string, object> message)
  88. {
  89. return this.SendAndWait(message, Endpoint.TIMEOUT);
  90. }
  91. /// <summary>send message and wait reply
  92. /// </summary>
  93. /// <param name="message"></param>
  94. /// <param name="timeout">timeout in milliseconds</param>
  95. /// <returns></returns>
  96. public IDictionary<string, object> SendAndWait(IDictionary<string, object> message, int timeout)
  97. {
  98. return this.SendAndWait(null, message, timeout);
  99. }
  100. /// <summary>send message and wait reply
  101. /// </summary>
  102. /// <param name="sender">use to send, must belong this proxy</param>
  103. /// <param name="message"></param>
  104. /// <param name="timeout">timeout in milliseconds</param>
  105. /// <returns></returns>
  106. internal IDictionary<string, object> SendAndWait(IChannelSender sender, IDictionary<string, object> message, int timeout)
  107. {
  108. return this._handler.SendAndWait(this,
  109. this.GetSender(sender),
  110. this.CreateMessage(message),
  111. timeout);
  112. }
  113. /// <summary>send message
  114. /// </summary>
  115. /// <param name="message"></param>
  116. public void Send(IDictionary<string, object> message)
  117. {
  118. this.Send(null, message);
  119. }
  120. /// <summary>send message
  121. /// </summary>
  122. /// <param name="sender">use to send, must belong this proxy</param>
  123. /// <param name="message"></param>
  124. internal void Send(IChannelSender sender, IDictionary<string, object> message)
  125. {
  126. this._handler.Send(this.CreateMessage(message), this.GetSender(sender));
  127. }
  128. private Message CreateMessage(IDictionary<string, object> message)
  129. {
  130. Message msg = new Message();
  131. msg.MessageType = MessageType.SEND;
  132. msg.Content = message;
  133. msg.Token = this.Token;
  134. return msg;
  135. }
  136. private IChannelSender GetSender(IChannelSender sender)
  137. {
  138. if (this._senders.Count == 0)
  139. throw new ChannelException(Text.E_NO_SENDER);
  140. if (this._senders.Contains(sender))
  141. return sender;
  142. return this._senders[this._random.Next(this._senders.Count)];
  143. }
  144. }
  145. }