TcpAgent.cpp 29 KB

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