Endpoint.cs 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106
  1. using System;
  2. using System.Collections.Generic;
  3. using System.Runtime.CompilerServices;
  4. using Taobao.Top.Link.Channel;
  5. using Top.Api;
  6. namespace Taobao.Top.Link.Endpoints
  7. {
  8. // Abstract network model
  9. // https://docs.google.com/drawings/d/1PRfzMVNGE4NKkpD9A_-QlH2PV47MFumZX8LbCwhzpQg/edit
  10. public sealed class Endpoint
  11. {
  12. internal static int TIMEOUT = 10000; // connect timeout in 10 seconds
  13. private ITopLogger _log;
  14. // in/out endpoints
  15. private IList<EndpointProxy> _connected;
  16. private EndpointHandler _handler;
  17. /// <summary>get or set clientChannelSelector
  18. /// </summary>
  19. public IClientChannelSelector ChannelSelector { get; set; }
  20. /// <summary>message received, see RTT based
  21. /// </summary>
  22. public event EventHandler<EndpointContext> OnMessage;
  23. /// <summary>ack message received, see RTT based
  24. /// </summary>
  25. public event EventHandler<AckMessageArgs> OnAckMessage;
  26. /// <summary>get id
  27. /// </summary>
  28. public Identity Identity { get; private set; }
  29. public Endpoint(Identity identity) : this(Log.Instance, identity) { }
  30. public Endpoint(ITopLogger logger, Identity identity)
  31. {
  32. this._connected = new List<EndpointProxy>();
  33. this._log = logger;
  34. this.Identity = identity;
  35. this.ChannelSelector = new ClientChannelSharedSelector(logger);
  36. this._handler = new EndpointHandler(logger);
  37. this._handler.MessageHandler = ctx => OnMessage(this, ctx);
  38. this._handler.AckMessageHandler = (m, i) => OnAckMessage(this, new AckMessageArgs(m, i));
  39. }
  40. /// <summary>get connected endpoint by id
  41. /// </summary>
  42. /// <param name="target"></param>
  43. /// <returns></returns>
  44. [MethodImplAttribute(MethodImplOptions.Synchronized)]
  45. public EndpointProxy GetEndpoint(Identity target)
  46. {
  47. if (target.Equals(this.Identity))
  48. throw new LinkException(Text.E_ID_DUPLICATE);
  49. foreach (EndpointProxy p in this._connected)
  50. if (p.Identity != null && p.Identity.Equals(target))
  51. return p;
  52. return null;
  53. }
  54. /// <summary>connect endpoint
  55. /// </summary>
  56. /// <param name="target">target id</param>
  57. /// <param name="uri">target address</param>
  58. /// <returns></returns>
  59. public EndpointProxy GetEndpoint(Identity target, string uri)
  60. {
  61. return this.GetEndpoint(target, uri, null);
  62. }
  63. /// <summary>connect endpoint
  64. /// </summary>
  65. /// <param name="target">target id</param>
  66. /// <param name="uri">target address</param>
  67. /// <param name="extras">passed as connect message</param>
  68. /// <returns></returns>
  69. [MethodImplAttribute(MethodImplOptions.Synchronized)]
  70. public EndpointProxy GetEndpoint(Identity target, string uri, IDictionary<string, object> extras)
  71. {
  72. EndpointProxy e = this.GetEndpoint(target) ?? this.CreateProxy();
  73. e.Identity = target;
  74. // always clear, cached proxy will have broken channel
  75. e.Remove(uri);
  76. // always reget channel, make sure it's valid
  77. IClientChannel channel = this.ChannelSelector.GetChannel(new Uri(uri));
  78. // connect message
  79. Message msg = new Message();
  80. msg.MessageType = MessageType.CONNECT;
  81. IDictionary<string, object> content = new Dictionary<string, object>();
  82. this.Identity.Render(content);
  83. // pass extra data
  84. if (extras != null)
  85. foreach (var p in extras)
  86. content.Add(p);
  87. msg.Content = content;
  88. this._handler.SendAndWait(e, channel, msg, TIMEOUT);
  89. e.Add(channel);
  90. return e;
  91. }
  92. private EndpointProxy CreateProxy()
  93. {
  94. EndpointProxy e = new EndpointProxy(this._handler);
  95. this._connected.Add(e);
  96. return e;
  97. }
  98. }
  99. }