TcpClientProvider.cs 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589
  1. using System;
  2. using System.Text;
  3. using System.Net;
  4. using System.Net.Sockets;
  5. using System.Threading;
  6. namespace IFramework.Net.Tcp
  7. {
  8. internal class TcpClientProvider : TcpSocket, IDisposable, ITcpClientProvider
  9. {
  10. #region variable
  11. private bool _isDisposed = false;
  12. private int bufferNumber = 8;
  13. private Encoding encoding = Encoding.UTF8;
  14. private int offsetNumber = 2;
  15. private ChannelProviderType channelProviderState = ChannelProviderType.Async;
  16. private LockParam lParam = new LockParam();
  17. private SocketTokenManager<SocketAsyncEventArgs> sendTokenManager = null;
  18. private SocketBufferManager sBufferManager = null;
  19. private AutoResetEvent mReset = new AutoResetEvent(false);
  20. #endregion
  21. #region properties
  22. /// <summary>
  23. /// 发送回调处理
  24. /// </summary>
  25. public OnSentHandler SentCallback { get; set; }
  26. /// <summary>
  27. /// 接收数据回调处理
  28. /// </summary>
  29. public OnReceivedHandler RecievedCallback { get; set; }
  30. /// <summary>
  31. /// 接受数据回调,返回缓冲区和偏移量
  32. /// </summary>
  33. public OnReceivedSegmentHandler ReceivedOffsetCallback { get; set; }
  34. /// <summary>
  35. /// 断开连接回调处理
  36. /// </summary>
  37. public OnDisconnectedHandler DisconnectedCallback { get; set; }
  38. /// <summary>
  39. /// 连接回调处理
  40. /// </summary>
  41. public OnConnectedHandler ConnectedCallback { get; set; }
  42. /// <summary>
  43. /// 是否连接状态
  44. /// </summary>
  45. public bool IsConnected
  46. {
  47. get { return isConnected; }
  48. }
  49. public int SendBufferPoolNumber { get { return sendTokenManager.Count; } }
  50. public ChannelProviderType ChannelProviderState
  51. {
  52. get { return channelProviderState; }
  53. }
  54. #endregion
  55. #region constructor
  56. public void Dispose()
  57. {
  58. Dispose(true);
  59. GC.SuppressFinalize(this);
  60. }
  61. protected virtual void Dispose(bool isDisposing)
  62. {
  63. if (_isDisposed) return;
  64. if (isDisposing)
  65. {
  66. DisposeSocketPool();
  67. SafeClose();
  68. _isDisposed = true;
  69. }
  70. }
  71. private void DisposeSocketPool()
  72. {
  73. sendTokenManager.Clear();
  74. if (sBufferManager != null)
  75. {
  76. sBufferManager.Clear();
  77. }
  78. }
  79. /// <summary>
  80. /// 构造
  81. /// </summary>
  82. /// <param name="chunkBufferSize">发送块缓冲区大小</param>
  83. /// <param name="bufferNumber">缓冲发送数</param>
  84. public TcpClientProvider(int chunkBufferSize = 4096, int bufferNumber = 8)
  85. :base(chunkBufferSize)
  86. {
  87. this.bufferNumber = bufferNumber;
  88. }
  89. #endregion
  90. #region public method
  91. /// <summary>
  92. /// 异步建立连接
  93. /// </summary>
  94. /// <param name="port"></param>
  95. /// <param name="ip"></param>
  96. public void Connect(int port, string ip)
  97. {
  98. try
  99. {
  100. if (!IsClose())
  101. {
  102. Close();
  103. }
  104. isConnected = false;
  105. channelProviderState = ChannelProviderType.Async;
  106. using (LockWait lwait = new LockWait(ref lParam))
  107. {
  108. CreatedConnectToBindArgs(port,ip);
  109. }
  110. }
  111. catch (Exception)
  112. {
  113. Close();
  114. throw;
  115. }
  116. }
  117. /// <summary>
  118. /// 异步等待连接返回结果
  119. /// </summary>
  120. /// <param name="port"></param>
  121. /// <param name="ip"></param>
  122. /// <returns></returns>
  123. public bool ConnectTo(int port,string ip)
  124. {
  125. try
  126. {
  127. if (!IsClose())
  128. {
  129. Close();
  130. }
  131. isConnected = false;
  132. channelProviderState = ChannelProviderType.AsyncWait;
  133. using (LockWait lwait = new LockWait(ref lParam))
  134. {
  135. CreatedConnectToBindArgs(port,ip);
  136. }
  137. mReset.WaitOne(connectioTimeout);
  138. isConnected = socket.Connected;
  139. return isConnected;
  140. }
  141. catch (Exception ex)
  142. {
  143. Close();
  144. throw ex;
  145. }
  146. }
  147. /// <summary>
  148. /// 同步连接
  149. /// </summary>
  150. /// <param name="port"></param>
  151. /// <param name="ip"></param>
  152. /// <returns></returns>
  153. public bool ConnectSync(int port, string ip)
  154. {
  155. if (!IsClose())
  156. {
  157. Close();
  158. }
  159. isConnected = false;
  160. channelProviderState = ChannelProviderType.Sync;
  161. int retry = 3;
  162. CreateTcpSocket(port, ip);
  163. //using (LockWait lwait = new LockWait(ref lParam))
  164. //{
  165. // CreatedConnectToBindArgs(port,ip);
  166. //}
  167. while (retry > 0)
  168. {
  169. try
  170. {
  171. --retry;
  172. socket.Connect(ipEndPoint);
  173. isConnected = true;
  174. return true;
  175. }
  176. catch (Exception ex)
  177. {
  178. Close();
  179. if (retry <= 0) throw ex;
  180. Thread.Sleep(1000);
  181. }
  182. }
  183. return false;
  184. }
  185. /// <summary>
  186. /// 根据偏移发送缓冲区数据
  187. /// </summary>
  188. /// <param name="buffer"></param>
  189. /// <param name="offset"></param>
  190. /// <param name="size"></param>
  191. public bool Send(SegmentOffset sendSegment, bool waiting = true)
  192. {
  193. try
  194. {
  195. if (IsClose())
  196. {
  197. Close();
  198. return false;
  199. }
  200. ArraySegment<byte>[] segItems = sBufferManager.BufferToSegments(sendSegment.buffer, sendSegment.offset, sendSegment.size);
  201. bool isWillEvent = true;
  202. foreach (var seg in segItems)
  203. {
  204. var tArgs = GetSocketAsyncFromSendPool(waiting);
  205. if (tArgs == null)
  206. {
  207. return false;
  208. }
  209. if (!sBufferManager.WriteBuffer(tArgs, seg.Array, seg.Offset, seg.Count))
  210. {
  211. sendTokenManager.Set(tArgs);
  212. throw new Exception(string.Format("发送缓冲区溢出...buffer block max size:{0}", sBufferManager.BlockSize));
  213. }
  214. if (tArgs.UserToken == null)
  215. ((SocketToken)tArgs.UserToken).TokenSocket = socket;
  216. if (IsClose())
  217. {
  218. Close();
  219. return false;
  220. }
  221. isWillEvent &= socket.SendAsync(tArgs);
  222. if (!isWillEvent)//can't trigger the io complated event to do
  223. {
  224. ProcessSentCallback(tArgs);
  225. }
  226. if (sendTokenManager.Count < (sendTokenManager.Capacity >> 2))
  227. Thread.Sleep(2);
  228. }
  229. return isWillEvent;
  230. }
  231. catch (Exception ex)
  232. {
  233. Close();
  234. throw ex;
  235. }
  236. }
  237. /// <summary>
  238. /// 发送文件
  239. /// </summary>
  240. /// <param name="filename"></param>
  241. public void SendFile(string filename)
  242. {
  243. socket.SendFile(filename);
  244. }
  245. /// <summary>
  246. /// 同步发送并接收数据,不设置receiveSegment 默认为只发数据
  247. /// </summary>
  248. /// <param name="buffer"></param>
  249. /// <param name="receiveBlock"></param>
  250. /// <param name="recAct"></param>
  251. /// <returns></returns>
  252. public int SendSync(SegmentOffset sendSegment,SegmentOffset receiveSegment)
  253. {
  254. if (channelProviderState != ChannelProviderType.Sync)
  255. {
  256. throw new Exception("需要使用同步连接...ConnectSync");
  257. }
  258. int sent = socket.Send(sendSegment.buffer, sendSegment.offset, sendSegment.size, SocketFlags.None);
  259. if (receiveSegment == null
  260. || receiveSegment.buffer == null
  261. || receiveSegment.size == 0) return sent;
  262. int cnt = socket.Receive(receiveSegment.buffer, receiveSegment.size, 0);
  263. return sent;
  264. }
  265. /// <summary>
  266. /// 同步接收数据
  267. /// </summary>
  268. /// <param name="receiveBlock"></param>
  269. /// <param name="receivedAction"></param>
  270. public void ReceiveSync(SegmentOffset receiveSegment, Action<SegmentOffset> receivedAction)
  271. {
  272. if (channelProviderState != ChannelProviderType.Sync)
  273. {
  274. throw new Exception("需要使用同步连接...ConnectSync");
  275. }
  276. int cnt = 0;
  277. do
  278. {
  279. if (socket.Connected == false) break;
  280. cnt = socket.Receive(receiveSegment.buffer, receiveSegment.size, 0);
  281. if (cnt <= 0) break;
  282. receivedAction(receiveSegment);
  283. } while (true);
  284. }
  285. /// <summary>
  286. /// 断开连接
  287. /// </summary>
  288. public void Disconnect()
  289. {
  290. Close();
  291. isConnected = false;
  292. }
  293. #endregion
  294. #region private method
  295. private void CreatedConnectToBindArgs(int port,string ip)
  296. {
  297. CreateTcpSocket(port,ip);
  298. //连接事件绑定
  299. var sArgs = new SocketAsyncEventArgs
  300. {
  301. RemoteEndPoint = ipEndPoint,
  302. UserToken = new SocketToken() { TokenSocket = socket }
  303. };
  304. sArgs.AcceptSocket = socket;
  305. sArgs.Completed += new EventHandler<SocketAsyncEventArgs>(IO_Completed);
  306. if (!socket.ConnectAsync(sArgs))
  307. {
  308. ProcessConnectCallback(sArgs);
  309. }
  310. }
  311. private void Close()
  312. {
  313. using (LockWait lwait = new LockWait(ref lParam))
  314. {
  315. DisposeSocketPool();
  316. SafeClose();
  317. isConnected = false;
  318. }
  319. }
  320. private bool IsClose()
  321. {
  322. return (IsConnected == false
  323. || socket == null
  324. || socket.Connected == false);
  325. }
  326. private SocketAsyncEventArgs GetSocketAsyncFromSendPool(bool waiting)
  327. {
  328. var tArgs = sendTokenManager.GetEmptyWait((retry) =>
  329. {
  330. return !IsClose();
  331. }, waiting);
  332. if (IsConnected == false) return null;
  333. if (tArgs == null)
  334. throw new Exception("发送缓冲池已用完,等待回收超时...");
  335. return tArgs;
  336. }
  337. private void InitializePool(int maxNumberOfConnections)
  338. {
  339. if(sendTokenManager!=null) sendTokenManager.Clear();
  340. if (sBufferManager != null) sBufferManager.Clear();
  341. int cnt = maxNumberOfConnections + offsetNumber;
  342. sendTokenManager = new SocketTokenManager<SocketAsyncEventArgs>(cnt);
  343. sBufferManager = new SocketBufferManager(cnt, receiveChunkSize);
  344. for (int i = 1; i <=cnt; ++i)
  345. {
  346. SocketAsyncEventArgs tArgs = new SocketAsyncEventArgs() {
  347. DisconnectReuseSocket=true
  348. };
  349. tArgs.Completed += IO_Completed;
  350. tArgs.UserToken = new SocketToken(i)
  351. {
  352. TokenSocket = socket,
  353. TokenId = i
  354. };
  355. sBufferManager.SetBuffer(tArgs);
  356. sendTokenManager.Set(tArgs);
  357. }
  358. }
  359. private void ProcessSentCallback(SocketAsyncEventArgs e)
  360. {
  361. try
  362. {
  363. if (e.SocketError == SocketError.Success)
  364. {
  365. if (SentCallback != null)
  366. {
  367. SocketToken sToken = e.UserToken as SocketToken;
  368. sToken.TokenIpEndPoint = (IPEndPoint)e.RemoteEndPoint;
  369. SentCallback(new SegmentToken(sToken, e.Buffer, e.Offset, e.BytesTransferred));
  370. }
  371. }
  372. else
  373. {
  374. ProcessDisconnectAsync(e);
  375. }
  376. }
  377. catch (Exception ex)
  378. {
  379. throw ex;
  380. }
  381. finally
  382. {
  383. sendTokenManager.Set(e);
  384. }
  385. }
  386. private void ProcessReceiveCallback(SocketAsyncEventArgs e)
  387. {
  388. if (e.BytesTransferred == 0
  389. || e.SocketError != SocketError.Success
  390. || e.AcceptSocket.Connected == false)
  391. {
  392. ProcessDisconnectAsync(e);
  393. return;
  394. }
  395. SocketToken sToken = e.UserToken as SocketToken;
  396. sToken.TokenIpEndPoint = (IPEndPoint)e.RemoteEndPoint;
  397. if (ReceivedOffsetCallback != null)
  398. ReceivedOffsetCallback(new SegmentToken(sToken, e.Buffer, e.Offset, e.BytesTransferred));
  399. if (RecievedCallback != null)
  400. {
  401. RecievedCallback(sToken, encoding.GetString(e.Buffer, e.Offset, e.BytesTransferred));
  402. }
  403. if (socket.Connected)
  404. {
  405. if (!e.AcceptSocket.ReceiveAsync(e))
  406. {
  407. ProcessReceiveCallback(e);
  408. }
  409. }
  410. }
  411. private void ProcessConnectCallback(SocketAsyncEventArgs e)
  412. {
  413. try
  414. {
  415. isConnected = (e.SocketError == SocketError.Success);
  416. if (isConnected)
  417. {
  418. using (LockWait lwait = new LockWait(ref lParam))
  419. {
  420. InitializePool(bufferNumber);
  421. }
  422. e.SetBuffer(receiveBuffer, 0, receiveChunkSize);
  423. if (ConnectedCallback != null)
  424. {
  425. SocketToken sToken = e.UserToken as SocketToken;
  426. sToken.TokenIpEndPoint = (IPEndPoint)e.RemoteEndPoint;
  427. ConnectedCallback(sToken, isConnected);
  428. }
  429. if (!e.AcceptSocket.ReceiveAsync(e))
  430. {
  431. ProcessReceiveCallback(e);
  432. }
  433. }
  434. else
  435. {
  436. ProcessDisconnectAsync(e);
  437. }
  438. if (channelProviderState == ChannelProviderType.AsyncWait)
  439. mReset.Set();
  440. }
  441. catch(Exception ex)
  442. {
  443. #if DEBUG
  444. Console.WriteLine(ex.TargetSite.Name + ex.Message);
  445. #endif
  446. }
  447. }
  448. private void ProcessDisconnectCallback(SocketAsyncEventArgs e)
  449. {
  450. try
  451. {
  452. isConnected = (e.SocketError == SocketError.Success);
  453. if (isConnected)
  454. {
  455. Close();
  456. }
  457. if (DisconnectedCallback != null)
  458. {
  459. SocketToken sToken = e.UserToken as SocketToken;
  460. sToken.TokenIpEndPoint = (IPEndPoint)e.RemoteEndPoint;
  461. DisconnectedCallback(sToken);
  462. }
  463. }
  464. catch (Exception ex)
  465. {
  466. throw ex;
  467. }
  468. }
  469. private void ProcessDisconnectAsync(SocketAsyncEventArgs e)
  470. {
  471. try
  472. {
  473. if (e.AcceptSocket == null) return;
  474. bool willRaiseEvent = false;
  475. if (e.AcceptSocket != null && e.AcceptSocket.Connected)
  476. willRaiseEvent = e.AcceptSocket.DisconnectAsync(e);
  477. if (!willRaiseEvent)
  478. {
  479. ProcessDisconnectCallback(e);
  480. }
  481. else
  482. {
  483. Close();
  484. }
  485. }
  486. catch (Exception ex)
  487. {
  488. #if DEBUG
  489. Console.WriteLine(ex.TargetSite.Name+ex.Message);
  490. #endif
  491. }
  492. }
  493. void IO_Completed(object sender, SocketAsyncEventArgs e)
  494. {
  495. switch (e.LastOperation)
  496. {
  497. case SocketAsyncOperation.Send:
  498. ProcessSentCallback(e);
  499. break;
  500. case SocketAsyncOperation.Receive:
  501. ProcessReceiveCallback(e);
  502. break;
  503. case SocketAsyncOperation.Connect:
  504. ProcessConnectCallback(e);
  505. break;
  506. case SocketAsyncOperation.Disconnect:
  507. ProcessDisconnectCallback(e);
  508. break;
  509. default:
  510. break;
  511. }
  512. }
  513. #endregion
  514. }
  515. }