NetSession.cs 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469
  1. using System;
  2. using System.Collections;
  3. using System.Collections.Generic;
  4. using System.Net.Sockets;
  5. using System.IO;
  6. using System.Threading;
  7. using CommonNetwork.Net;
  8. using CommonLang.IO;
  9. using CommonLang.ByteOrder;
  10. using CommonLang.Protocol;
  11. using System.Net;
  12. using CommonLang.Net;
  13. using CommonLang;
  14. using CommonLang.Log;
  15. namespace CommonNetwork.Sockets
  16. {
  17. public class NetSession : BaseNetSession
  18. {
  19. private Logger log = LoggerFactory.GetLogger("NetSession");
  20. private TcpClient mTCP = null;
  21. protected INetPackageCodec mCodec;
  22. private INetSessionListener mListener;
  23. private Queue<Object> mSendQueue = new Queue<Object>();
  24. private List<object> mSendingArray = new List<object>(11);
  25. private Thread mWriteThread;
  26. private Thread mReadThread;
  27. public NetSession()
  28. {
  29. }
  30. /// <summary>
  31. /// 判断当前网络是否已经连接
  32. /// </summary>
  33. /// <returns></returns>
  34. public override bool IsConnected
  35. {
  36. get
  37. {
  38. TcpClient tcp = mTCP;
  39. if (tcp != null)
  40. {
  41. return tcp.Connected;
  42. }
  43. return false;
  44. }
  45. }
  46. public override INetPackageCodec Codec
  47. {
  48. get { return mCodec; }
  49. }
  50. public override IPEndPoint RemoteAddress
  51. {
  52. get
  53. {
  54. if (mTCP != null) return mTCP.Client.RemoteEndPoint as IPEndPoint;
  55. return null;
  56. }
  57. }
  58. public TcpClient Session
  59. {
  60. get { return mTCP; }
  61. }
  62. public override bool Open(string url, INetPackageCodec codec, INetSessionListener listener)
  63. {
  64. bool ret = false;
  65. try
  66. {
  67. lock (this)
  68. {
  69. if (mTCP == null)
  70. {
  71. this.mCodec = codec;
  72. this.mListener = listener;
  73. this.mURL = url;
  74. string[] url_kv = url.Split(':');
  75. //this.mRemoteAddress = IPUtil.ToEndPoint(url_kv[0], int.Parse(url_kv[1]));
  76. //new IPEndPoint(IPAddress.Parse(url_kv[0]), int.Parse(url_kv[1]));
  77. //this.mLocalAddress = IPUtil.CreateLocalConnectorEndPoint();//new IPEndPoint(IPAddress.Parse("127.0.0.1"), 0);
  78. //this.mRecvQueue.Clear();
  79. lock (mSendQueue) this.mSendQueue.Clear();
  80. this.mReadThread = new Thread(new ThreadStart(this.runRead));
  81. this.mReadThread.IsBackground = true;
  82. this.mWriteThread = new Thread(new ThreadStart(this.runWrite));
  83. this.mWriteThread.IsBackground = true;
  84. // 建立SOCKET链接
  85. this.mTCP = new TcpClient();
  86. //this.mTCP.ReceiveTimeout = 5000;
  87. //this.mTCP.SendTimeout = 5000;
  88. this.send_buff = new MemoryStream(8192);
  89. this.recv_buff = new MemoryStream(8192);
  90. this.mTCP.NoDelay = true;
  91. this.mTCP.Connect(url_kv[0], int.Parse(url_kv[1]));
  92. //创建读写线程对象
  93. this.mReadThread.Start();
  94. this.mWriteThread.Start();
  95. ret = true;
  96. }
  97. }
  98. }
  99. catch (Exception err)
  100. {
  101. log.Error(err.Message, err);
  102. onException(new NetException("/n[Open:]" + URL +
  103. "/n[InnerException:]" + err.InnerException +
  104. "/n[Exception:]" + err.Message +
  105. "/n[Source:]" + err.Source +
  106. "/n[StackTrace:]" + err.StackTrace));
  107. }
  108. return ret;
  109. }
  110. public override bool Close()
  111. {
  112. bool ret = false;
  113. lock (this)
  114. {
  115. if (mTCP != null)
  116. {
  117. try
  118. {
  119. this.mTCP.Close();
  120. }
  121. catch (Exception err)
  122. {
  123. log.Error(err.Message, err);
  124. }
  125. this.mTCP = null;
  126. ret = true;
  127. }
  128. }
  129. if (ret)
  130. {
  131. lock (mSendQueue)
  132. {
  133. this.mSendQueue.Clear();
  134. }
  135. if (mReadThread != null)
  136. {
  137. try
  138. {
  139. this.mReadThread.Join(1000);
  140. }
  141. catch (Exception err)
  142. {
  143. log.Error(err.Message, err);
  144. }
  145. this.mReadThread = null;
  146. }
  147. lock (mSendQueue)
  148. {
  149. this.mSendQueue.Clear();
  150. }
  151. if (mWriteThread != null)
  152. {
  153. try
  154. {
  155. this.mWriteThread.Join(1000);
  156. }
  157. catch (Exception err)
  158. {
  159. log.Error(err.Message, err);
  160. }
  161. this.mWriteThread = null;
  162. }
  163. }
  164. return ret;
  165. }
  166. //-------------------------------------------------------------------------------------
  167. /// <summary>
  168. /// 发送一个消息,该方法将立即返回。
  169. /// </summary>
  170. /// <param name="data"></param>
  171. public override void Send(Object data)
  172. {
  173. if (mTCP != null)
  174. {
  175. lock (mSendQueue)
  176. {
  177. mSendQueue.Enqueue(data);
  178. // 通知写线程开始工作。
  179. Monitor.PulseAll(mSendQueue);
  180. }
  181. }
  182. }
  183. public override void SendResponse(IMessage rsponse, int requestMessageID)
  184. {
  185. rsponse.MessageID = requestMessageID;
  186. Send(rsponse);
  187. }
  188. //-------------------------------------------------------------------------------------
  189. override public string ToString()
  190. {
  191. return "Session[" + URL + "](" + GetHashCode() + ")";
  192. }
  193. //-------------------------------------------------------------------------------------
  194. private void onClose()
  195. {
  196. lock (mSendQueue)
  197. {
  198. mSendQueue.Clear();
  199. try
  200. {
  201. if (send_buff != null) send_buff.Dispose();
  202. if (recv_buff != null) recv_buff.Dispose();
  203. }
  204. catch (Exception err)
  205. {
  206. log.Error(err.Message, err);
  207. }
  208. send_buff = null;
  209. recv_buff = null;
  210. }
  211. mListener.sessionClosed(this);
  212. if (mOnSessionClosed != null)
  213. {
  214. mOnSessionClosed.Invoke(this);
  215. }
  216. }
  217. private void onOpen()
  218. {
  219. mListener.sessionOpened(this);
  220. if (mOnSessionOpened != null)
  221. {
  222. mOnSessionOpened.Invoke(this);
  223. }
  224. }
  225. private void onReceive(Object message)
  226. {
  227. try
  228. {
  229. mListener.messageReceived(this, message);
  230. if (mOnMessageReceived != null)
  231. {
  232. mOnMessageReceived.Invoke(this, message);
  233. }
  234. }
  235. catch (Exception err)
  236. {
  237. log.Error(err.Message, err);
  238. onException(err);
  239. }
  240. }
  241. private void onSent(Object message)
  242. {
  243. mListener.messageSent(this, message);
  244. if (mOnMessageSent != null)
  245. {
  246. mOnMessageSent.Invoke(this, message);
  247. }
  248. }
  249. private void onException(Exception err)
  250. {
  251. mListener.onError(this, err);
  252. if (mOnError != null)
  253. {
  254. mOnError.Invoke(this, err);
  255. }
  256. }
  257. //-------------------------------------------------------------------------------------
  258. private void innerClose()
  259. {
  260. lock (this)
  261. {
  262. if (mTCP != null)
  263. {
  264. try
  265. {
  266. this.mTCP.Close();
  267. }
  268. catch (Exception err)
  269. {
  270. log.Error(err.Message, err);
  271. }
  272. }
  273. }
  274. }
  275. private void runWrite()
  276. {
  277. TcpClient _socket = mTCP;
  278. try
  279. {
  280. onOpen();
  281. Stream output = _socket.GetStream();
  282. while (_socket.Connected)
  283. {
  284. lock (mSendQueue)
  285. {
  286. mSendingArray.Clear();
  287. if (mSendQueue.Count > 0)
  288. {
  289. mSendingArray.AddRange(mSendQueue);
  290. mSendQueue.Clear();
  291. }
  292. else
  293. {
  294. // 如果没有待传输消息,则等待输入
  295. Monitor.Wait(mSendQueue, 100);
  296. }
  297. }
  298. if (mSendingArray.Count > 0 && _socket.Connected)
  299. {
  300. for (int i = 0; i < mSendingArray.Count; i++)
  301. {
  302. object send_msg = mSendingArray[i];
  303. try
  304. {
  305. int sendBytes;
  306. if (doEncode(output, send_msg, out sendBytes))
  307. {
  308. output.Flush();
  309. mSendPacks++;
  310. onSent(send_msg);
  311. }
  312. mSendBytes += sendBytes;
  313. }
  314. catch (Exception err)
  315. {
  316. try
  317. {
  318. if (_socket.Connected)
  319. {
  320. log.Error(err.Message, err);
  321. innerClose();
  322. onException(new NetException(err.Message, err));
  323. }
  324. }
  325. catch (Exception err2)
  326. {
  327. log.Error(err.Message, err2);
  328. }
  329. break;
  330. }
  331. }
  332. }
  333. }
  334. }
  335. catch (Exception err)
  336. {
  337. log.Error(err.Message, err);
  338. onException(new NetException("/n[runWrite:]" + URL +
  339. "/n[InnerException:]" + err.InnerException +
  340. "/n[Exception:]" + err.Message +
  341. "/n[Source:]" + err.Source +
  342. "/n[StackTrace:]" + err.StackTrace, err));
  343. }
  344. finally
  345. {
  346. onClose();
  347. }
  348. }
  349. private void runRead()
  350. {
  351. TcpClient _socket = mTCP;
  352. try
  353. {
  354. Stream input = _socket.GetStream();
  355. while (_socket.Connected)
  356. {
  357. try
  358. {
  359. Object msg = null;
  360. int readBytes;
  361. if (doDecode(input, out msg, out readBytes))
  362. {
  363. mRecvPacks++;
  364. onReceive(msg);
  365. }
  366. mRecvBytes += readBytes;
  367. }
  368. catch (Exception err)
  369. {
  370. try
  371. {
  372. if (_socket.Connected)
  373. {
  374. log.Error(err.Message, err);
  375. innerClose();
  376. onException(new NetException(err.Message, err));
  377. }
  378. }
  379. catch (Exception err2)
  380. {
  381. log.Error(err.Message, err2);
  382. }
  383. break;
  384. }
  385. }
  386. }
  387. catch (Exception err)
  388. {
  389. log.Error(err.Message, err);
  390. onException(new NetException("runRead: " + err.Message));
  391. }
  392. }
  393. protected static readonly byte[] ZERO_BUFF = new byte[0];
  394. protected MemoryStream send_buff;
  395. protected MemoryStream recv_buff;
  396. protected virtual bool doEncode(Stream output, object send_msg, out int sendBytes)
  397. {
  398. sendBytes = 0;
  399. send_buff.Position = 0;
  400. if (Codec.doEncode(send_buff, send_msg))
  401. {
  402. byte[] h_body = send_buff.GetBuffer();
  403. int length = (int)send_buff.Position;
  404. LittleEdian.PutS32(output, length);
  405. output.Write(h_body, 0, length);
  406. sendBytes = length + 4;
  407. return true;
  408. }
  409. return false;
  410. }
  411. protected virtual bool doDecode(Stream input, out object msg, out int readBytes)
  412. {
  413. int length = LittleEdian.GetS32(input);
  414. if (length == -1)
  415. {
  416. msg = null;
  417. readBytes = 0;
  418. return false;
  419. }
  420. else if (length <= 0)
  421. {
  422. throw new Exception("received negative trunk size=" + length);
  423. }
  424. if (recv_buff.Capacity < length)
  425. {
  426. recv_buff.Capacity = length;
  427. }
  428. recv_buff.Position = 0;
  429. recv_buff.SetLength(length);
  430. byte[] buffer = recv_buff.GetBuffer();
  431. IOUtil.ReadToEnd(input, buffer, 0, length);
  432. readBytes = length + 4;
  433. recv_buff.Position = 0;
  434. if (Codec.doDecode(recv_buff, out msg))
  435. {
  436. if (recv_buff.Position != length)
  437. {
  438. throw new Exception(string.Format("can not decode full trunk size={0} type={1}", length, msg != null ? msg.GetType().FullName : msg));
  439. }
  440. return true;
  441. }
  442. return false;
  443. }
  444. }
  445. }