TcpAgent.cpp 30 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306
  1. /*
  2. * Copyright: JessMA Open Source (ldcsaa@gmail.com)
  3. *
  4. * Author : Bruce Liang
  5. * Website : https://github.com/ldcsaa
  6. * Project : https://github.com/ldcsaa/HP-Socket
  7. * Blog : http://www.cnblogs.com/ldcsaa
  8. * Wiki : http://www.oschina.net/p/hp-socket
  9. * QQ Group : 44636872, 75375912
  10. *
  11. * Licensed under the Apache License, Version 2.0 (the "License");
  12. * you may not use this file except in compliance with the License.
  13. * You may obtain a copy of the License at
  14. *
  15. * http://www.apache.org/licenses/LICENSE-2.0
  16. *
  17. * Unless required by applicable law or agreed to in writing, software
  18. * distributed under the License is distributed on an "AS IS" BASIS,
  19. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  20. * See the License for the specific language governing permissions and
  21. * limitations under the License.
  22. */
  23. #include "TcpAgent.h"
  24. #include "./common/FileHelper.h"
  25. BOOL CTcpAgent::Start(LPCTSTR lpszBindAddress, BOOL bAsyncConnect)
  26. {
  27. if(!CheckParams() || !CheckStarting())
  28. return FALSE;
  29. PrepareStart();
  30. if(ParseBindAddress(lpszBindAddress))
  31. if(CreateWorkerThreads())
  32. {
  33. m_bAsyncConnect = bAsyncConnect;
  34. m_enState = SS_STARTED;
  35. return TRUE;
  36. }
  37. EXECUTE_RESTORE_ERROR(Stop());
  38. return FALSE;
  39. }
  40. void CTcpAgent::SetLastError(EnSocketError code, LPCSTR func, int ec)
  41. {
  42. m_enLastError = code;
  43. ::SetLastError(ec);
  44. }
  45. BOOL CTcpAgent::CheckParams()
  46. {
  47. if ((m_enSendPolicy >= SP_PACK && m_enSendPolicy <= SP_DIRECT) &&
  48. (m_enOnSendSyncPolicy >= OSSP_NONE && m_enOnSendSyncPolicy <= OSSP_RECEIVE) &&
  49. ((int)m_dwSyncConnectTimeout > 0) &&
  50. ((int)m_dwMaxConnectionCount > 0 && m_dwMaxConnectionCount <= MAX_CONNECTION_COUNT) &&
  51. ((int)m_dwWorkerThreadCount > 0 && m_dwWorkerThreadCount <= MAX_WORKER_THREAD_COUNT) &&
  52. ((int)m_dwSocketBufferSize >= MIN_SOCKET_BUFFER_SIZE) &&
  53. ((int)m_dwFreeSocketObjLockTime >= 1000) &&
  54. ((int)m_dwFreeSocketObjPool >= 0) &&
  55. ((int)m_dwFreeBufferObjPool >= 0) &&
  56. ((int)m_dwFreeSocketObjHold >= 0) &&
  57. ((int)m_dwFreeBufferObjHold >= 0) &&
  58. ((int)m_dwKeepAliveTime >= 1000 || m_dwKeepAliveTime == 0) &&
  59. ((int)m_dwKeepAliveInterval >= 1000 || m_dwKeepAliveInterval == 0) )
  60. return TRUE;
  61. SetLastError(SE_INVALID_PARAM, __FUNCTION__, ERROR_INVALID_PARAMETER);
  62. return FALSE;
  63. }
  64. void CTcpAgent::PrepareStart()
  65. {
  66. m_bfActiveSockets.Reset(m_dwMaxConnectionCount);
  67. m_lsFreeSocket.Reset(m_dwFreeSocketObjPool);
  68. m_bfObjPool.SetItemCapacity(m_dwSocketBufferSize);
  69. m_bfObjPool.SetPoolSize(m_dwFreeBufferObjPool);
  70. m_bfObjPool.SetPoolHold(m_dwFreeBufferObjHold);
  71. m_bfObjPool.Prepare();
  72. m_rcBuffers = make_unique<CBufferPtr[]>(m_dwWorkerThreadCount);
  73. for_each(m_rcBuffers.get(), m_rcBuffers.get() + m_dwWorkerThreadCount, [this](CBufferPtr& buff) {buff.Malloc(m_dwSocketBufferSize);});
  74. }
  75. BOOL CTcpAgent::CheckStarting()
  76. {
  77. CSpinLock locallock(m_csState);
  78. if(m_enState == SS_STOPPED)
  79. m_enState = SS_STARTING;
  80. else
  81. {
  82. SetLastError(SE_ILLEGAL_STATE, __FUNCTION__, ERROR_INVALID_STATE);
  83. return FALSE;
  84. }
  85. return TRUE;
  86. }
  87. BOOL CTcpAgent::CheckStoping()
  88. {
  89. if(m_enState != SS_STOPPED)
  90. {
  91. CSpinLock locallock(m_csState);
  92. if(HasStarted())
  93. {
  94. m_enState = SS_STOPPING;
  95. return TRUE;
  96. }
  97. }
  98. SetLastError(SE_ILLEGAL_STATE, __FUNCTION__, ERROR_INVALID_STATE);
  99. return FALSE;
  100. }
  101. BOOL CTcpAgent::ParseBindAddress(LPCTSTR lpszBindAddress)
  102. {
  103. if(::IsStrEmpty(lpszBindAddress))
  104. return TRUE;
  105. HP_SOCKADDR addr;
  106. if(::sockaddr_A_2_IN(lpszBindAddress, 0, addr))
  107. {
  108. SOCKET sock = socket(addr.family, SOCK_STREAM, IPPROTO_TCP);
  109. if(sock != INVALID_SOCKET)
  110. {
  111. if(::bind(sock, addr.Addr(), addr.AddrSize()) != SOCKET_ERROR)
  112. {
  113. addr.Copy(m_soAddr);
  114. return TRUE;
  115. }
  116. else
  117. SetLastError(SE_SOCKET_BIND, __FUNCTION__, ::WSAGetLastError());
  118. ::ManualCloseSocket(sock);
  119. }
  120. else
  121. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ::WSAGetLastError());
  122. }
  123. else
  124. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ::WSAGetLastError());
  125. return FALSE;
  126. }
  127. BOOL CTcpAgent::CreateWorkerThreads()
  128. {
  129. DWORD dwWorkerThreadCount = m_dwWorkerThreadCount
  130. #ifdef USE_EXTERNAL_GC
  131. + 1
  132. #endif
  133. ;
  134. if(!m_ioDispatcher.Start(this, DEFAULT_WORKER_MAX_EVENT_COUNT, dwWorkerThreadCount))
  135. {
  136. SetLastError(SE_WORKER_THREAD_CREATE, __FUNCTION__, ::WSAGetLastError());
  137. return FALSE;
  138. }
  139. #ifdef USE_EXTERNAL_GC
  140. m_fdGCTimer = m_ioDispatcher.AddTimer(m_dwWorkerThreadCount, GC_CHECK_INTERVAL, this);
  141. if(IS_INVALID_FD(m_fdGCTimer))
  142. {
  143. SetLastError(SE_GC_START, __FUNCTION__, ::WSAGetLastError());
  144. return FALSE;
  145. }
  146. #endif
  147. return TRUE;
  148. }
  149. BOOL CTcpAgent::Stop()
  150. {
  151. if(!CheckStoping())
  152. return FALSE;
  153. DisconnectClientSocket();
  154. WaitForClientSocketClose();
  155. WaitForWorkerThreadEnd();
  156. ReleaseClientSocket();
  157. FireShutdown();
  158. ReleaseFreeSocket();
  159. Reset();
  160. return TRUE;
  161. }
  162. void CTcpAgent::DisconnectClientSocket()
  163. {
  164. ::WaitFor(100);
  165. if(m_bfActiveSockets.Elements() == 0)
  166. return;
  167. TAgentSocketObjPtrPool::IndexSet indexes;
  168. m_bfActiveSockets.CopyIndexes(indexes);
  169. for(auto it = indexes.begin(), end = indexes.end(); it != end; ++it)
  170. Disconnect(*it);
  171. }
  172. void CTcpAgent::WaitForClientSocketClose()
  173. {
  174. while(m_bfActiveSockets.Elements() > 0)
  175. ::WaitFor(50);
  176. }
  177. void CTcpAgent::WaitForWorkerThreadEnd()
  178. {
  179. m_ioDispatcher.Stop();
  180. }
  181. void CTcpAgent::ReleaseClientSocket()
  182. {
  183. VERIFY(m_bfActiveSockets.IsEmpty());
  184. m_bfActiveSockets.Reset();
  185. }
  186. void CTcpAgent::ReleaseFreeSocket()
  187. {
  188. m_lsFreeSocket.Clear();
  189. #ifdef USE_EXTERNAL_GC
  190. if(IS_VALID_FD(m_fdGCTimer))
  191. {
  192. close(m_fdGCTimer);
  193. m_fdGCTimer = INVALID_FD;
  194. }
  195. #endif
  196. ReleaseGCSocketObj(TRUE);
  197. VERIFY(m_lsGCSocket.IsEmpty());
  198. }
  199. void CTcpAgent::Reset()
  200. {
  201. m_bfObjPool.Clear();
  202. m_phSocket.Reset();
  203. m_soAddr.Reset();
  204. m_rcBuffers = nullptr;
  205. m_enState = SS_STOPPED;
  206. m_evWait.SyncNotifyAll();
  207. }
  208. BOOL CTcpAgent::Connect(LPCTSTR lpszRemoteAddress, USHORT usPort, CONNID* pdwConnID, PVOID pExtra, USHORT usLocalPort, LPCTSTR lpszLocalAddress)
  209. {
  210. ASSERT(lpszRemoteAddress && usPort != 0);
  211. if(!HasStarted())
  212. {
  213. ::SetLastError(ERROR_INVALID_STATE);
  214. return FALSE;
  215. }
  216. if(!pdwConnID)
  217. pdwConnID = CreateLocalObject(CONNID);
  218. *pdwConnID = 0;
  219. HP_SOCKADDR addr;
  220. HP_SCOPE_HOST host(lpszRemoteAddress);
  221. SOCKET soClient = INVALID_SOCKET;
  222. DWORD result = CreateClientSocket(host.addr, usPort, lpszLocalAddress, usLocalPort, soClient, addr);
  223. if(result == NO_ERROR)
  224. {
  225. result = PrepareConnect(*pdwConnID, soClient);
  226. if(result == NO_ERROR)
  227. result = ConnectToServer(*pdwConnID, host.name, soClient, addr, pExtra);
  228. }
  229. if(result != NO_ERROR)
  230. {
  231. if(soClient != INVALID_SOCKET)
  232. ::ManualCloseSocket(soClient);
  233. ::SetLastError(result);
  234. }
  235. return (result == NO_ERROR);
  236. }
  237. int CTcpAgent::CreateClientSocket(LPCTSTR lpszRemoteAddress, USHORT usPort, LPCTSTR lpszLocalAddress, USHORT usLocalPort, SOCKET& soClient, HP_SOCKADDR& addr)
  238. {
  239. if(!::GetSockAddrByHostName(lpszRemoteAddress, usPort, addr))
  240. return ERROR_ADDRNOTAVAIL;
  241. HP_SOCKADDR* lpBindAddr = &m_soAddr;
  242. if(::IsStrNotEmpty(lpszLocalAddress))
  243. {
  244. lpBindAddr = CreateLocalObject(HP_SOCKADDR);
  245. if(!::sockaddr_A_2_IN(lpszLocalAddress, 0, *lpBindAddr))
  246. return ::WSAGetLastError();
  247. }
  248. BOOL bBind = lpBindAddr->IsSpecified();
  249. if(bBind && lpBindAddr->family != addr.family)
  250. return ERROR_AFNOSUPPORT;
  251. int result = NO_ERROR;
  252. soClient = socket(addr.family, SOCK_STREAM, IPPROTO_TCP);
  253. if(soClient == INVALID_SOCKET)
  254. result = ::WSAGetLastError();
  255. else
  256. {
  257. BOOL bOnOff = (m_dwKeepAliveTime > 0 && m_dwKeepAliveInterval > 0);
  258. VERIFY(IS_NO_ERROR(::SSO_KeepAliveVals(soClient, bOnOff, m_dwKeepAliveTime, m_dwKeepAliveInterval)));
  259. VERIFY(IS_NO_ERROR(::SSO_ReuseAddress(soClient, m_enReusePolicy)));
  260. VERIFY(IS_NO_ERROR(::SSO_NoDelay(soClient, m_bNoDelay)));
  261. if(bBind && usLocalPort == 0)
  262. {
  263. if(::bind(soClient, lpBindAddr->Addr(), lpBindAddr->AddrSize()) == SOCKET_ERROR)
  264. result = ::WSAGetLastError();
  265. }
  266. else if(usLocalPort != 0)
  267. {
  268. HP_SOCKADDR bindAddr = bBind ? *lpBindAddr : HP_SOCKADDR::AnyAddr(addr.family);
  269. bindAddr.SetPort(usLocalPort);
  270. if(::bind(soClient, bindAddr.Addr(), bindAddr.AddrSize()) == SOCKET_ERROR)
  271. result = ::WSAGetLastError();
  272. }
  273. }
  274. return result;
  275. }
  276. int CTcpAgent::PrepareConnect(CONNID& dwConnID, SOCKET soClient)
  277. {
  278. if(!m_bfActiveSockets.AcquireLock(dwConnID))
  279. return ERROR_CONNECTION_COUNT_LIMIT;
  280. if(TRIGGER(FirePrepareConnect(dwConnID, soClient)) == HR_ERROR)
  281. {
  282. VERIFY(m_bfActiveSockets.ReleaseLock(dwConnID, nullptr));
  283. return ENSURE_ERROR_CANCELLED;
  284. }
  285. return NO_ERROR;
  286. }
  287. int CTcpAgent::ConnectToServer(CONNID dwConnID, LPCTSTR lpszRemoteHostName, SOCKET& soClient, const HP_SOCKADDR& addr, PVOID pExtra)
  288. {
  289. TAgentSocketObj* pSocketObj = GetFreeSocketObj(dwConnID, soClient);
  290. AddClientSocketObj(dwConnID, pSocketObj, addr, lpszRemoteHostName, pExtra);
  291. int result = HAS_ERROR;
  292. VERIFY(::fcntl_SETFL(pSocketObj->socket, O_NOATIME | O_NONBLOCK | O_CLOEXEC));
  293. int rc = ::connect(pSocketObj->socket, addr.Addr(), addr.AddrSize());
  294. if(IS_NO_ERROR(rc) || IS_IO_PENDING_ERROR())
  295. {
  296. if(m_bAsyncConnect)
  297. {
  298. if(m_ioDispatcher.AddFD(pSocketObj->socket, EPOLLOUT, pSocketObj))
  299. result = NO_ERROR;
  300. }
  301. else
  302. {
  303. if(IS_HAS_ERROR(result))
  304. result = ::WaitForSocketWrite(pSocketObj->socket, m_dwSyncConnectTimeout);
  305. if(IS_NO_ERROR(result))
  306. {
  307. pSocketObj->SetConnected();
  308. if(TRIGGER(FireConnect(pSocketObj)) == HR_ERROR)
  309. result = ENSURE_ERROR_CANCELLED;
  310. else
  311. {
  312. UINT evts = (pSocketObj->IsPending() ? EPOLLOUT : 0) | (pSocketObj->IsPaused() ? 0 : EPOLLIN);
  313. if(!m_ioDispatcher.AddFD(pSocketObj->socket, evts | EPOLLRDHUP, pSocketObj))
  314. result = HAS_ERROR;
  315. }
  316. }
  317. }
  318. }
  319. if(result == HAS_ERROR)
  320. result = ::WSAGetLastError();
  321. if(result != NO_ERROR)
  322. {
  323. AddFreeSocketObj(pSocketObj, SCF_NONE);
  324. soClient = INVALID_SOCKET;
  325. }
  326. return result;
  327. }
  328. TAgentSocketObj* CTcpAgent::GetFreeSocketObj(CONNID dwConnID, SOCKET soClient)
  329. {
  330. DWORD dwIndex;
  331. TAgentSocketObj* pSocketObj = nullptr;
  332. if(m_lsFreeSocket.TryLock(&pSocketObj, dwIndex))
  333. {
  334. if(::GetTimeGap32(pSocketObj->freeTime) >= m_dwFreeSocketObjLockTime)
  335. VERIFY(m_lsFreeSocket.ReleaseLock(nullptr, dwIndex));
  336. else
  337. {
  338. VERIFY(m_lsFreeSocket.ReleaseLock(pSocketObj, dwIndex));
  339. pSocketObj = nullptr;
  340. }
  341. }
  342. if(!pSocketObj) pSocketObj = CreateSocketObj();
  343. pSocketObj->Reset(dwConnID, soClient);
  344. return pSocketObj;
  345. }
  346. TAgentSocketObj* CTcpAgent::CreateSocketObj()
  347. {
  348. return TAgentSocketObj::Construct(m_phSocket, m_bfObjPool);
  349. }
  350. void CTcpAgent::DeleteSocketObj(TAgentSocketObj* pSocketObj)
  351. {
  352. TAgentSocketObj::Destruct(pSocketObj);
  353. }
  354. void CTcpAgent::AddFreeSocketObj(TAgentSocketObj* pSocketObj, EnSocketCloseFlag enFlag, EnSocketOperation enOperation, int iErrorCode)
  355. {
  356. if(!InvalidSocketObj(pSocketObj))
  357. return;
  358. CloseClientSocketObj(pSocketObj, enFlag, enOperation, iErrorCode);
  359. m_bfActiveSockets.Remove(pSocketObj->connID);
  360. TAgentSocketObj::Release(pSocketObj);
  361. #ifndef USE_EXTERNAL_GC
  362. ReleaseGCSocketObj();
  363. #endif
  364. if(!m_lsFreeSocket.TryPut(pSocketObj))
  365. m_lsGCSocket.PushBack(pSocketObj);
  366. }
  367. void CTcpAgent::ReleaseGCSocketObj(BOOL bForce)
  368. {
  369. ::ReleaseGCObj(m_lsGCSocket, m_dwFreeSocketObjLockTime, bForce);
  370. }
  371. BOOL CTcpAgent::InvalidSocketObj(TAgentSocketObj* pSocketObj)
  372. {
  373. return TAgentSocketObj::InvalidSocketObj(pSocketObj);
  374. }
  375. void CTcpAgent::AddClientSocketObj(CONNID dwConnID, TAgentSocketObj* pSocketObj, const HP_SOCKADDR& remoteAddr, LPCTSTR lpszRemoteHostName, PVOID pExtra)
  376. {
  377. ASSERT(FindSocketObj(dwConnID) == nullptr);
  378. pSocketObj->connTime = ::TimeGetTime();
  379. pSocketObj->activeTime = pSocketObj->connTime;
  380. pSocketObj->host = lpszRemoteHostName;
  381. pSocketObj->extra = pExtra;
  382. pSocketObj->SetConnected(CST_CONNECTING);
  383. remoteAddr.Copy(pSocketObj->remoteAddr);
  384. VERIFY(m_bfActiveSockets.ReleaseLock(dwConnID, pSocketObj));
  385. }
  386. TAgentSocketObj* CTcpAgent::FindSocketObj(CONNID dwConnID)
  387. {
  388. TAgentSocketObj* pSocketObj = nullptr;
  389. if(m_bfActiveSockets.Get(dwConnID, &pSocketObj) != TAgentSocketObjPtrPool::GR_VALID)
  390. pSocketObj = nullptr;
  391. return pSocketObj;
  392. }
  393. void CTcpAgent::CloseClientSocketObj(TAgentSocketObj* pSocketObj, EnSocketCloseFlag enFlag, EnSocketOperation enOperation, int iErrorCode, int iShutdownFlag)
  394. {
  395. ASSERT(TAgentSocketObj::IsExist(pSocketObj));
  396. if(enFlag == SCF_CLOSE)
  397. FireClose(pSocketObj, SO_CLOSE, SE_OK);
  398. else if(enFlag == SCF_ERROR)
  399. FireClose(pSocketObj, enOperation, iErrorCode);
  400. SOCKET socket = pSocketObj->socket;
  401. pSocketObj->socket = INVALID_SOCKET;
  402. ::ManualCloseSocket(socket, iShutdownFlag);
  403. }
  404. BOOL CTcpAgent::GetLocalAddress(CONNID dwConnID, TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort)
  405. {
  406. ASSERT(lpszAddress != nullptr && iAddressLen > 0);
  407. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  408. if(!TAgentSocketObj::IsValid(pSocketObj))
  409. {
  410. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  411. return FALSE;
  412. }
  413. return ::GetSocketLocalAddress(pSocketObj->socket, lpszAddress, iAddressLen, usPort);
  414. }
  415. BOOL CTcpAgent::GetRemoteAddress(CONNID dwConnID, TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort)
  416. {
  417. ASSERT(lpszAddress != nullptr && iAddressLen > 0);
  418. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  419. if(!TAgentSocketObj::IsExist(pSocketObj))
  420. {
  421. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  422. return FALSE;
  423. }
  424. ADDRESS_FAMILY usFamily;
  425. return ::sockaddr_IN_2_A(pSocketObj->remoteAddr, usFamily, lpszAddress, iAddressLen, usPort);
  426. }
  427. BOOL CTcpAgent::GetRemoteHost(CONNID dwConnID, TCHAR lpszHost[], int& iHostLen, USHORT& usPort)
  428. {
  429. ASSERT(lpszHost != nullptr && iHostLen > 0);
  430. BOOL isOK = FALSE;
  431. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  432. if(!TAgentSocketObj::IsExist(pSocketObj))
  433. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  434. else
  435. {
  436. int iLen = pSocketObj->host.GetLength() + 1;
  437. if(iHostLen >= iLen)
  438. {
  439. memcpy(lpszHost, CA2CT((LPCSTR)pSocketObj->host), iLen * sizeof(TCHAR));
  440. usPort = pSocketObj->remoteAddr.Port();
  441. isOK = TRUE;
  442. }
  443. iHostLen = iLen;
  444. }
  445. return isOK;
  446. }
  447. BOOL CTcpAgent::GetRemoteHost(CONNID dwConnID, LPCSTR* lpszHost, USHORT* pusPort)
  448. {
  449. *lpszHost = nullptr;
  450. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  451. if(!TAgentSocketObj::IsExist(pSocketObj))
  452. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  453. else
  454. {
  455. *lpszHost = pSocketObj->host;
  456. if(pusPort)
  457. *pusPort = pSocketObj->remoteAddr.Port();
  458. }
  459. return (*lpszHost != nullptr && (*lpszHost)[0] != 0);
  460. }
  461. BOOL CTcpAgent::SetConnectionExtra(CONNID dwConnID, PVOID pExtra)
  462. {
  463. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  464. return SetConnectionExtra(pSocketObj, pExtra);
  465. }
  466. BOOL CTcpAgent::SetConnectionExtra(TAgentSocketObj* pSocketObj, PVOID pExtra)
  467. {
  468. if(!TAgentSocketObj::IsExist(pSocketObj))
  469. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  470. else
  471. {
  472. pSocketObj->extra = pExtra;
  473. return TRUE;
  474. }
  475. return FALSE;
  476. }
  477. BOOL CTcpAgent::GetConnectionExtra(CONNID dwConnID, PVOID* ppExtra)
  478. {
  479. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  480. return GetConnectionExtra(pSocketObj, ppExtra);
  481. }
  482. BOOL CTcpAgent::GetConnectionExtra(TAgentSocketObj* pSocketObj, PVOID* ppExtra)
  483. {
  484. ASSERT(ppExtra != nullptr);
  485. if(!TAgentSocketObj::IsExist(pSocketObj))
  486. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  487. else
  488. {
  489. *ppExtra = pSocketObj->extra;
  490. return TRUE;
  491. }
  492. return FALSE;
  493. }
  494. BOOL CTcpAgent::SetConnectionReserved(CONNID dwConnID, PVOID pReserved)
  495. {
  496. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  497. return SetConnectionReserved(pSocketObj, pReserved);
  498. }
  499. BOOL CTcpAgent::SetConnectionReserved(TAgentSocketObj* pSocketObj, PVOID pReserved)
  500. {
  501. if(!TAgentSocketObj::IsExist(pSocketObj))
  502. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  503. else
  504. {
  505. pSocketObj->reserved = pReserved;
  506. return TRUE;
  507. }
  508. return FALSE;
  509. }
  510. BOOL CTcpAgent::GetConnectionReserved(CONNID dwConnID, PVOID* ppReserved)
  511. {
  512. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  513. return GetConnectionReserved(pSocketObj, ppReserved);
  514. }
  515. BOOL CTcpAgent::GetConnectionReserved(TAgentSocketObj* pSocketObj, PVOID* ppReserved)
  516. {
  517. ASSERT(ppReserved != nullptr);
  518. if(!TAgentSocketObj::IsExist(pSocketObj))
  519. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  520. else
  521. {
  522. *ppReserved = pSocketObj->reserved;
  523. return TRUE;
  524. }
  525. return FALSE;
  526. }
  527. BOOL CTcpAgent::SetConnectionReserved2(CONNID dwConnID, PVOID pReserved2)
  528. {
  529. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  530. return SetConnectionReserved2(pSocketObj, pReserved2);
  531. }
  532. BOOL CTcpAgent::SetConnectionReserved2(TAgentSocketObj* pSocketObj, PVOID pReserved2)
  533. {
  534. if(!TAgentSocketObj::IsExist(pSocketObj))
  535. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  536. else
  537. {
  538. pSocketObj->reserved2 = pReserved2;
  539. return TRUE;
  540. }
  541. return FALSE;
  542. }
  543. BOOL CTcpAgent::GetConnectionReserved2(CONNID dwConnID, PVOID* ppReserved2)
  544. {
  545. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  546. return GetConnectionReserved2(pSocketObj, ppReserved2);
  547. }
  548. BOOL CTcpAgent::GetConnectionReserved2(TAgentSocketObj* pSocketObj, PVOID* ppReserved2)
  549. {
  550. ASSERT(ppReserved2 != nullptr);
  551. if(!TAgentSocketObj::IsExist(pSocketObj))
  552. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  553. else
  554. {
  555. *ppReserved2 = pSocketObj->reserved2;
  556. return TRUE;
  557. }
  558. return FALSE;
  559. }
  560. BOOL CTcpAgent::IsPauseReceive(CONNID dwConnID, BOOL& bPaused)
  561. {
  562. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  563. if(!TAgentSocketObj::IsValid(pSocketObj))
  564. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  565. else
  566. {
  567. bPaused = pSocketObj->paused;
  568. return TRUE;
  569. }
  570. return FALSE;
  571. }
  572. BOOL CTcpAgent::IsConnected(CONNID dwConnID)
  573. {
  574. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  575. if(!TAgentSocketObj::IsValid(pSocketObj))
  576. {
  577. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  578. return FALSE;
  579. }
  580. return pSocketObj->HasConnected();
  581. }
  582. BOOL CTcpAgent::GetPendingDataLength(CONNID dwConnID, int& iPending)
  583. {
  584. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  585. if(!TAgentSocketObj::IsValid(pSocketObj))
  586. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  587. else
  588. {
  589. iPending = pSocketObj->Pending();
  590. return TRUE;
  591. }
  592. return FALSE;
  593. }
  594. DWORD CTcpAgent::GetConnectionCount()
  595. {
  596. return m_bfActiveSockets.Elements();
  597. }
  598. BOOL CTcpAgent::GetAllConnectionIDs(CONNID pIDs[], DWORD& dwCount)
  599. {
  600. return m_bfActiveSockets.GetAllElementIndexes(pIDs, dwCount);
  601. }
  602. BOOL CTcpAgent::GetConnectPeriod(CONNID dwConnID, DWORD& dwPeriod)
  603. {
  604. BOOL isOK = TRUE;
  605. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  606. if(TAgentSocketObj::IsValid(pSocketObj))
  607. dwPeriod = ::GetTimeGap32(pSocketObj->connTime);
  608. else
  609. {
  610. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  611. isOK = FALSE;
  612. }
  613. return isOK;
  614. }
  615. BOOL CTcpAgent::GetSilencePeriod(CONNID dwConnID, DWORD& dwPeriod)
  616. {
  617. if(!m_bMarkSilence)
  618. return FALSE;
  619. BOOL isOK = TRUE;
  620. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  621. if(TAgentSocketObj::IsValid(pSocketObj))
  622. dwPeriod = ::GetTimeGap32(pSocketObj->activeTime);
  623. else
  624. {
  625. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  626. isOK = FALSE;
  627. }
  628. return isOK;
  629. }
  630. BOOL CTcpAgent::Disconnect(CONNID dwConnID, BOOL bForce)
  631. {
  632. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  633. if(!TAgentSocketObj::IsValid(pSocketObj))
  634. {
  635. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  636. return FALSE;
  637. }
  638. return m_ioDispatcher.SendCommandByFD(pSocketObj->socket, DISP_CMD_DISCONNECT, dwConnID, bForce);
  639. }
  640. BOOL CTcpAgent::DisconnectLongConnections(DWORD dwPeriod, BOOL bForce)
  641. {
  642. if(dwPeriod > MAX_CONNECTION_PERIOD)
  643. return FALSE;
  644. if(m_bfActiveSockets.Elements() == 0)
  645. return TRUE;
  646. DWORD now = ::TimeGetTime();
  647. TAgentSocketObjPtrPool::IndexSet indexes;
  648. m_bfActiveSockets.CopyIndexes(indexes);
  649. for(auto it = indexes.begin(), end = indexes.end(); it != end; ++it)
  650. {
  651. CONNID connID = *it;
  652. TAgentSocketObj* pSocketObj = FindSocketObj(connID);
  653. if(TAgentSocketObj::IsValid(pSocketObj) && (int)(now - pSocketObj->connTime) >= (int)dwPeriod)
  654. Disconnect(connID, bForce);
  655. }
  656. return TRUE;
  657. }
  658. BOOL CTcpAgent::DisconnectSilenceConnections(DWORD dwPeriod, BOOL bForce)
  659. {
  660. if(!m_bMarkSilence)
  661. return FALSE;
  662. if(dwPeriod > MAX_CONNECTION_PERIOD)
  663. return FALSE;
  664. if(m_bfActiveSockets.Elements() == 0)
  665. return TRUE;
  666. DWORD now = ::TimeGetTime();
  667. TAgentSocketObjPtrPool::IndexSet indexes;
  668. m_bfActiveSockets.CopyIndexes(indexes);
  669. for(auto it = indexes.begin(), end = indexes.end(); it != end; ++it)
  670. {
  671. CONNID connID = *it;
  672. TAgentSocketObj* pSocketObj = FindSocketObj(connID);
  673. if(TAgentSocketObj::IsValid(pSocketObj) && (int)(now - pSocketObj->activeTime) >= (int)dwPeriod)
  674. Disconnect(connID, bForce);
  675. }
  676. return TRUE;
  677. }
  678. BOOL CTcpAgent::PauseReceive(CONNID dwConnID, BOOL bPause)
  679. {
  680. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  681. if(!TAgentSocketObj::IsValid(pSocketObj))
  682. {
  683. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  684. return FALSE;
  685. }
  686. if(!pSocketObj->HasConnected())
  687. {
  688. ::SetLastError(ERROR_INVALID_STATE);
  689. return FALSE;
  690. }
  691. if(pSocketObj->paused == bPause)
  692. return TRUE;
  693. pSocketObj->paused = bPause;
  694. if(!bPause)
  695. return m_ioDispatcher.SendCommandByFD(pSocketObj->socket, DISP_CMD_UNPAUSE, pSocketObj->connID);
  696. return TRUE;
  697. }
  698. BOOL CTcpAgent::OnBeforeProcessIo(const TDispContext* pContext, PVOID pv, UINT events)
  699. {
  700. if(pv == this)
  701. {
  702. ReleaseGCSocketObj(FALSE);
  703. ::ReadTimer(m_fdGCTimer);
  704. return FALSE;
  705. }
  706. TAgentSocketObj* pSocketObj = (TAgentSocketObj*)(pv);
  707. if(!TAgentSocketObj::IsValid(pSocketObj))
  708. return FALSE;
  709. if(events & _EPOLL_ALL_ERROR_EVENTS)
  710. pSocketObj->SetConnected(FALSE);
  711. pSocketObj->Increment();
  712. if(!TAgentSocketObj::IsValid(pSocketObj))
  713. {
  714. pSocketObj->Decrement();
  715. return FALSE;
  716. }
  717. if(pSocketObj->IsConnecting())
  718. {
  719. HandleConnect(pContext, pSocketObj, events);
  720. pSocketObj->Decrement();
  721. return FALSE;
  722. }
  723. return TRUE;
  724. }
  725. VOID CTcpAgent::OnAfterProcessIo(const TDispContext* pContext, PVOID pv, UINT events, BOOL rs)
  726. {
  727. TAgentSocketObj* pSocketObj = (TAgentSocketObj*)(pv);
  728. if(TAgentSocketObj::IsValid(pSocketObj))
  729. {
  730. ASSERT(rs && !(events & (EPOLLERR | EPOLLHUP | EPOLLRDHUP)));
  731. UINT evts = (pSocketObj->IsPending() ? EPOLLOUT : 0) | (pSocketObj->IsPaused() ? 0 : EPOLLIN);
  732. m_ioDispatcher.ModFD(pSocketObj->socket, evts | EPOLLRDHUP, pSocketObj);
  733. }
  734. pSocketObj->Decrement();
  735. }
  736. VOID CTcpAgent::OnCommand(const TDispContext* pContext, TDispCommand* pCmd)
  737. {
  738. switch(pCmd->type)
  739. {
  740. case DISP_CMD_SEND:
  741. HandleCmdSend(pContext, (CONNID)(pCmd->wParam));
  742. break;
  743. case DISP_CMD_UNPAUSE:
  744. HandleCmdUnpause(pContext, (CONNID)(pCmd->wParam));
  745. break;
  746. case DISP_CMD_DISCONNECT:
  747. HandleCmdDisconnect(pContext, (CONNID)(pCmd->wParam), (BOOL)pCmd->lParam);
  748. break;
  749. }
  750. }
  751. VOID CTcpAgent::HandleCmdSend(const TDispContext* pContext, CONNID dwConnID)
  752. {
  753. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  754. if(TAgentSocketObj::IsValid(pSocketObj) && pSocketObj->IsPending())
  755. m_ioDispatcher.ProcessIo(pContext, pSocketObj, EPOLLOUT);
  756. }
  757. VOID CTcpAgent::HandleCmdUnpause(const TDispContext* pContext, CONNID dwConnID)
  758. {
  759. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  760. if(!TAgentSocketObj::IsValid(pSocketObj))
  761. return;
  762. if(BeforeUnpause(pSocketObj))
  763. m_ioDispatcher.ProcessIo(pContext, pSocketObj, EPOLLIN);
  764. else
  765. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_RECEIVE, ENSURE_ERROR_CANCELLED);
  766. }
  767. VOID CTcpAgent::HandleCmdDisconnect(const TDispContext* pContext, CONNID dwConnID, BOOL bForce)
  768. {
  769. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  770. if(TAgentSocketObj::IsValid(pSocketObj))
  771. m_ioDispatcher.ProcessIo(pContext, pSocketObj, EPOLLHUP);
  772. }
  773. BOOL CTcpAgent::OnReadyRead(const TDispContext* pContext, PVOID pv, UINT events)
  774. {
  775. return HandleReceive(pContext, (TAgentSocketObj*)pv, RETRIVE_EVENT_FLAG_H(events));
  776. }
  777. BOOL CTcpAgent::OnReadyWrite(const TDispContext* pContext, PVOID pv, UINT events)
  778. {
  779. return HandleSend(pContext, (TAgentSocketObj*)pv, RETRIVE_EVENT_FLAG_H(events));
  780. }
  781. BOOL CTcpAgent::OnHungUp(const TDispContext* pContext, PVOID pv, UINT events)
  782. {
  783. return HandleClose(pContext, (TAgentSocketObj*)pv, SCF_CLOSE, events);
  784. }
  785. BOOL CTcpAgent::OnError(const TDispContext* pContext, PVOID pv, UINT events)
  786. {
  787. return HandleClose(pContext, (TAgentSocketObj*)pv, SCF_ERROR, events);
  788. }
  789. VOID CTcpAgent::OnDispatchThreadStart(THR_ID tid)
  790. {
  791. OnWorkerThreadStart(tid);
  792. }
  793. VOID CTcpAgent::OnDispatchThreadEnd(THR_ID tid)
  794. {
  795. OnWorkerThreadEnd(tid);
  796. }
  797. BOOL CTcpAgent::HandleClose(const TDispContext* pContext, TAgentSocketObj* pSocketObj, EnSocketCloseFlag enFlag, UINT events)
  798. {
  799. EnSocketOperation enOperation = SO_CLOSE;
  800. if(events & _EPOLL_HUNGUP_EVENTS)
  801. enOperation = SO_CLOSE;
  802. else if(events & EPOLLIN)
  803. enOperation = SO_RECEIVE;
  804. else if(events & EPOLLOUT)
  805. enOperation = SO_SEND;
  806. int iErrorCode = 0;
  807. if(enFlag == SCF_ERROR)
  808. iErrorCode = ::SSO_GetError(pSocketObj->socket);
  809. AddFreeSocketObj(pSocketObj, enFlag, enOperation, iErrorCode);
  810. return TRUE;
  811. }
  812. BOOL CTcpAgent::HandleConnect(const TDispContext* pContext, TAgentSocketObj* pSocketObj, UINT events)
  813. {
  814. int code = ::SSO_GetError(pSocketObj->socket);
  815. if(!IS_NO_ERROR(code) || (events & _EPOLL_ERROR_EVENTS))
  816. {
  817. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_CONNECT, code);
  818. return FALSE;
  819. }
  820. if((events & (_EPOLL_HUNGUP_EVENTS | _EPOLL_READ_EVENTS)) || !(events & EPOLLOUT))
  821. {
  822. AddFreeSocketObj(pSocketObj, SCF_CLOSE, SO_CONNECT, SE_OK);
  823. return FALSE;
  824. }
  825. pSocketObj->SetConnected();
  826. if(TRIGGER(FireConnect(pSocketObj)) == HR_ERROR)
  827. {
  828. AddFreeSocketObj(pSocketObj, SCF_NONE);
  829. return FALSE;
  830. }
  831. UINT evts = (pSocketObj->IsPending() ? EPOLLOUT : 0) | (pSocketObj->IsPaused() ? 0 : EPOLLIN);
  832. if(!m_ioDispatcher.ModFD(pSocketObj->socket, evts | EPOLLRDHUP, pSocketObj))
  833. {
  834. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_CONNECT, ::WSAGetLastError());
  835. return FALSE;
  836. }
  837. return TRUE;
  838. }
  839. BOOL CTcpAgent::HandleReceive(const TDispContext* pContext, TAgentSocketObj* pSocketObj, int flag)
  840. {
  841. ASSERT(TAgentSocketObj::IsValid(pSocketObj));
  842. if(m_bMarkSilence) pSocketObj->activeTime = ::TimeGetTime();
  843. CBufferPtr& buffer = m_rcBuffers[pContext->GetIndex()];
  844. int reads = flag ? -1 : MAX_CONTINUE_READS;
  845. for(int i = 0; i < reads || reads < 0; i++)
  846. {
  847. if(pSocketObj->paused)
  848. break;
  849. int rc = (int)read(pSocketObj->socket, buffer.Ptr(), buffer.Size());
  850. if(rc > 0)
  851. {
  852. if(TRIGGER(FireReceive(pSocketObj, buffer.Ptr(), rc)) == HR_ERROR)
  853. {
  854. TRACE("<C-CNNID: %zu> OnReceive() event return 'HR_ERROR', connection will be closed !", pSocketObj->connID);
  855. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_RECEIVE, ENSURE_ERROR_CANCELLED);
  856. return FALSE;
  857. }
  858. }
  859. else if(rc == 0)
  860. {
  861. AddFreeSocketObj(pSocketObj, SCF_CLOSE, SO_RECEIVE, SE_OK);
  862. return FALSE;
  863. }
  864. else
  865. {
  866. ASSERT(rc == SOCKET_ERROR);
  867. int code = ::WSAGetLastError();
  868. if(code == ERROR_WOULDBLOCK)
  869. break;
  870. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_RECEIVE, code);
  871. return FALSE;
  872. }
  873. }
  874. return TRUE;
  875. }
  876. BOOL CTcpAgent::HandleSend(const TDispContext* pContext, TAgentSocketObj* pSocketObj, int flag)
  877. {
  878. ASSERT(TAgentSocketObj::IsValid(pSocketObj));
  879. if(!pSocketObj->IsPending())
  880. return TRUE;
  881. BOOL bBlocked = FALSE;
  882. int writes = flag ? -1 : MAX_CONTINUE_WRITES;
  883. TBufferObjList& sndBuff = pSocketObj->sndBuff;
  884. TItemPtr itPtr(sndBuff);
  885. for(int i = 0; i < writes || writes < 0; i++)
  886. {
  887. {
  888. CReentrantCriSecLock locallock(pSocketObj->csSend);
  889. itPtr = sndBuff.PopFront();
  890. }
  891. if(!itPtr.IsValid())
  892. break;
  893. ASSERT(!itPtr->IsEmpty());
  894. if(!SendItem(pSocketObj, itPtr, bBlocked))
  895. return FALSE;
  896. if(bBlocked)
  897. {
  898. ASSERT(!itPtr->IsEmpty());
  899. CReentrantCriSecLock locallock(pSocketObj->csSend);
  900. sndBuff.PushFront(itPtr.Detach());
  901. break;
  902. }
  903. }
  904. return TRUE;
  905. }
  906. BOOL CTcpAgent::SendItem(TAgentSocketObj* pSocketObj, TItem* pItem, BOOL& bBlocked)
  907. {
  908. while(!pItem->IsEmpty())
  909. {
  910. int rc = (int)write(pSocketObj->socket, pItem->Ptr(), pItem->Size());
  911. if(rc > 0)
  912. {
  913. if(TRIGGER(FireSend(pSocketObj, pItem->Ptr(), rc)) == HR_ERROR)
  914. {
  915. TRACE("<C-CNNID: %zu> OnSend() event should not return 'HR_ERROR' !!", pSocketObj->connID);
  916. ASSERT(FALSE);
  917. }
  918. pItem->Reduce(rc);
  919. }
  920. else if(rc == SOCKET_ERROR)
  921. {
  922. int code = ::WSAGetLastError();
  923. if(code == ERROR_WOULDBLOCK)
  924. {
  925. bBlocked = TRUE;
  926. break;
  927. }
  928. else
  929. {
  930. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_SEND, code);
  931. return FALSE;
  932. }
  933. }
  934. else
  935. ASSERT(FALSE);
  936. }
  937. return TRUE;
  938. }
  939. BOOL CTcpAgent::Send(CONNID dwConnID, const BYTE* pBuffer, int iLength, int iOffset)
  940. {
  941. ASSERT(pBuffer && iLength > 0);
  942. if(iOffset != 0) pBuffer += iOffset;
  943. WSABUF buffer;
  944. buffer.len = iLength;
  945. buffer.buf = (BYTE*)pBuffer;
  946. return SendPackets(dwConnID, &buffer, 1);
  947. }
  948. BOOL CTcpAgent::DoSendPackets(CONNID dwConnID, const WSABUF pBuffers[], int iCount)
  949. {
  950. ASSERT(pBuffers && iCount > 0);
  951. TAgentSocketObj* pSocketObj = FindSocketObj(dwConnID);
  952. if(!TAgentSocketObj::IsValid(pSocketObj))
  953. {
  954. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  955. return FALSE;
  956. }
  957. return DoSendPackets(pSocketObj, pBuffers, iCount);
  958. }
  959. BOOL CTcpAgent::DoSendPackets(TAgentSocketObj* pSocketObj, const WSABUF pBuffers[], int iCount)
  960. {
  961. ASSERT(pSocketObj && pBuffers && iCount > 0);
  962. int result = NO_ERROR;
  963. if(!pSocketObj->HasConnected())
  964. {
  965. ::SetLastError(ERROR_INVALID_STATE);
  966. return FALSE;
  967. }
  968. if(pBuffers && iCount > 0)
  969. {
  970. CLocalSafeCounter localcounter(*pSocketObj);
  971. CReentrantCriSecLock locallock(pSocketObj->csSend);
  972. if(TAgentSocketObj::IsValid(pSocketObj))
  973. result = SendInternal(pSocketObj, pBuffers, iCount);
  974. else
  975. result = ERROR_OBJECT_NOT_FOUND;
  976. }
  977. else
  978. result = ERROR_INVALID_PARAMETER;
  979. if(result != NO_ERROR)
  980. ::SetLastError(result);
  981. return (result == NO_ERROR);
  982. }
  983. int CTcpAgent::SendInternal(TAgentSocketObj* pSocketObj, const WSABUF pBuffers[], int iCount)
  984. {
  985. BOOL bPending = pSocketObj->IsPending();
  986. for(int i = 0; i < iCount; i++)
  987. {
  988. int iBufLen = pBuffers[i].len;
  989. if(iBufLen > 0)
  990. {
  991. BYTE* pBuffer = (BYTE*)pBuffers[i].buf;
  992. ASSERT(pBuffer);
  993. pSocketObj->sndBuff.Cat(pBuffer, iBufLen);
  994. ASSERT(pSocketObj->sndBuff.Length() > 0);
  995. }
  996. }
  997. if(!bPending && pSocketObj->IsPending())
  998. {
  999. if(!m_ioDispatcher.SendCommandByFD(pSocketObj->socket, DISP_CMD_SEND, pSocketObj->connID))
  1000. return ::GetLastError();
  1001. }
  1002. return NO_ERROR;
  1003. }
  1004. BOOL CTcpAgent::SendSmallFile(CONNID dwConnID, LPCTSTR lpszFileName, const LPWSABUF pHead, const LPWSABUF pTail)
  1005. {
  1006. CFile file;
  1007. CFileMapping fmap;
  1008. WSABUF szBuf[3];
  1009. HRESULT hr = ::MakeSmallFilePackage(lpszFileName, file, fmap, szBuf, pHead, pTail);
  1010. if(FAILED(hr))
  1011. {
  1012. ::SetLastError(hr);
  1013. return FALSE;
  1014. }
  1015. return SendPackets(dwConnID, szBuf, 3);
  1016. }