Server.cs 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197
  1. using System;
  2. using System.Collections.Generic;
  3. using System.Diagnostics.CodeAnalysis;
  4. using System.Linq;
  5. using System.Text;
  6. using System.Threading.Tasks;
  7. using NetMQ;
  8. namespace TcpEventBus
  9. {
  10. public class Server : ISocket
  11. {
  12. NetMQPoller poller = new NetMQPoller();
  13. public string IP { get; }
  14. public int Port { get; }
  15. public int Timeout { get; }
  16. public Server(string ip, int port,int timeout = 5000)
  17. {
  18. IP = ip;
  19. Port = port;
  20. IsConnect = false;
  21. Timeout = timeout;
  22. _server = new NetMQ.Sockets.SubscriberSocket($"@tcp://{IP}:{Port}");
  23. _server.ReceiveReady += _server_ReceiveReady;
  24. _server.Subscribe("");
  25. _client = new NetMQ.Sockets.PublisherSocket($"@tcp://{IP}:{Port - 1}");
  26. _client.Options.SendBuffer = 500 * 1024;
  27. _client.Options.SendHighWatermark = 10;
  28. poller.Add(_client);
  29. poller.Add(_server);
  30. }
  31. private void _server_ReceiveReady(object? sender, NetMQSocketEventArgs e)
  32. {
  33. OnData?.Invoke(e.Socket.ReceiveFrameBytes());
  34. }
  35. [AllowNull]
  36. private NetMQ.Sockets.SubscriberSocket _server;
  37. [AllowNull]
  38. private NetMQ.Sockets.PublisherSocket _client;
  39. public int ReceiveTimeout => 3000;
  40. [AllowNull]
  41. public Action OnConnected { get; set; }
  42. [AllowNull]
  43. public Action<ArraySegment<byte>> OnData { get; set; }
  44. [AllowNull]
  45. public Action OnDisconnected { get; set; }
  46. public bool IsConnect { get; private set; }
  47. public bool IsService => true;
  48. private object locker = new object();
  49. public void Disconnect()
  50. {
  51. if (System.Console.IsOutputRedirected)
  52. {
  53. System.Diagnostics.Debug.WriteLine($"{DateTime.Now}:{nameof(Disconnect)}");
  54. }
  55. else
  56. {
  57. Console.WriteLine($"{DateTime.Now}:{nameof(Disconnect)}");
  58. }
  59. poller.Stop();
  60. IsConnect = false;
  61. }
  62. public void Send(byte[] data)
  63. {
  64. lock (locker)
  65. {
  66. if (!IsConnect) return;
  67. try
  68. {
  69. _client.SendFrame(data);
  70. }
  71. catch
  72. {
  73. }
  74. }
  75. }
  76. public void Start()
  77. {
  78. if (System.Console.IsOutputRedirected)
  79. {
  80. System.Diagnostics.Debug.WriteLine($"{DateTime.Now}:{nameof(Start)}");
  81. }
  82. else
  83. {
  84. Console.WriteLine($"{DateTime.Now}:{nameof(Start)}");
  85. }
  86. Task.Run(() => poller.Run());
  87. IsConnect = true;
  88. }
  89. }
  90. public class Client : ISocket
  91. {
  92. NetMQ.NetMQPoller _poller = new NetMQPoller();
  93. public string IP { get; }
  94. public int Port { get; }
  95. private string ServerAddress => $"tcp://{IP}:{Port - 1}";
  96. private string ClientAddress => $"tcp://{IP}:{Port}";
  97. public int Timeout { get; }
  98. public Client(string ip, int port,int timeout = 5000)
  99. {
  100. IP = ip;
  101. Port = port;
  102. IsConnect = false;
  103. Timeout = timeout;
  104. _server = new NetMQ.Sockets.SubscriberSocket();
  105. _server.ReceiveReady += _server_ReceiveReady;
  106. _server.Subscribe("");
  107. _client = new NetMQ.Sockets.PublisherSocket();
  108. _client.Options.SendBuffer = 500 * 1024;
  109. _client.Options.SendHighWatermark = 10;
  110. _poller = new NetMQPoller();
  111. _poller.Add(_server);
  112. _poller.Add(_client);
  113. }
  114. private void _server_ReceiveReady(object? sender, NetMQSocketEventArgs e)
  115. {
  116. OnData?.Invoke(e.Socket.ReceiveFrameBytes());
  117. }
  118. [AllowNull]
  119. private NetMQ.Sockets.SubscriberSocket _server;
  120. [AllowNull]
  121. private NetMQ.Sockets.PublisherSocket _client;
  122. public int ReceiveTimeout => 3000;
  123. [AllowNull]
  124. public Action OnConnected { get; set; }
  125. [AllowNull]
  126. public Action<ArraySegment<byte>> OnData { get; set; }
  127. [AllowNull]
  128. public Action OnDisconnected { get; set; }
  129. public bool IsConnect { get; private set; }
  130. public bool IsService => false;
  131. private object locker = new object();
  132. public void Disconnect()
  133. {
  134. if (System.Console.IsOutputRedirected)
  135. {
  136. System.Diagnostics.Debug.WriteLine($"{DateTime.Now}:{nameof(Disconnect)}");
  137. }
  138. else
  139. {
  140. Console.WriteLine($"{DateTime.Now}:{nameof(Disconnect)}");
  141. }
  142. _server.Disconnect(ServerAddress);
  143. _client.Disconnect(ClientAddress);
  144. _poller.Stop();
  145. IsConnect = false;
  146. }
  147. public void Send(byte[] data)
  148. {
  149. lock (locker)
  150. {
  151. if (!IsConnect) return;
  152. try
  153. {
  154. _client.SendFrame(data);
  155. }
  156. catch
  157. {
  158. }
  159. }
  160. }
  161. public void Start()
  162. {
  163. _client.Connect(ClientAddress);
  164. _server.Connect(ServerAddress);
  165. Thread.Sleep(100);
  166. if (System.Console.IsOutputRedirected)
  167. {
  168. System.Diagnostics.Debug.WriteLine($"{DateTime.Now}:{nameof(Start)}");
  169. }
  170. else
  171. {
  172. Console.WriteLine($"{DateTime.Now}:{nameof(Start)}");
  173. }
  174. Task.Run(() => _poller.Run());
  175. IsConnect = true;
  176. }
  177. }
  178. }