UDT协议实现分析——bind、listen与accept

UDT Server启动之后,基于UDT协议的UDP数据可靠传输才成为可能,因而接下来分析与UDT Server有关的几个主要API的实现,来了解下UDT Server是如何listening在特定UDP端口上的。主要有UDT::bind(),UDT::listen()和UDT::accept()等几个函数。

bind过程

通常UDT Server在创建UDT Socket之后,首先就要调用UDT::bind(),与一个特定的本地UDP端口地址进行绑定,以便可以在希望的端口上监听。这里来看一下UDT::bind()的实现:

1int CUDTUnited::bind(const UDTSOCKET u, const sockaddr* name, int namelen) { 2 CUDTSocket* s = locate(u); 3 if (NULL == s) 4 throw CUDTException(5, 4, 0); 5 6 CGuard cg(s->m_ControlLock); 7 8 // cannot bind a socket more than once 9 if (INIT != s->m_Status) 10 throw CUDTException(5, 0, 0); 11 12 // check the size of SOCKADDR structure 13 if (AF_INET == s->m_iIPversion) { 14 if (namelen != sizeof(sockaddr_in)) 15 throw CUDTException(5, 3, 0); 16 } else { 17 if (namelen != sizeof(sockaddr_in6)) 18 throw CUDTException(5, 3, 0); 19 } 20 21 s->m_pUDT->open(); 22 updateMux(s, name); 23 s->m_Status = OPENED; 24 25 // copy address information of local node 26 s->m_pUDT->m_pSndQueue->m_pChannel->getSockAddr(s->m_pSelfAddr); 27 28 return 0; 29} 30 31int CUDTUnited::bind(UDTSOCKET u, UDPSOCKET udpsock) { 32 CUDTSocket* s = locate(u); 33 if (NULL == s) 34 throw CUDTException(5, 4, 0); 35 36 CGuard cg(s->m_ControlLock); 37 38 // cannot bind a socket more than once 39 if (INIT != s->m_Status) 40 throw CUDTException(5, 0, 0); 41 42 sockaddr_in name4; 43 sockaddr_in6 name6; 44 sockaddr* name; 45 socklen_t namelen; 46 47 if (AF_INET == s->m_iIPversion) { 48 namelen = sizeof(sockaddr_in); 49 name = (sockaddr*) &name4; 50 } else { 51 namelen = sizeof(sockaddr_in6); 52 name = (sockaddr*) &name6; 53 } 54 55 if (-1 == ::getsockname(udpsock, name, &namelen)) 56 throw CUDTException(5, 3); 57 58 s->m_pUDT->open(); 59 updateMux(s, name, &udpsock); 60 s->m_Status = OPENED; 61 62 // copy address information of local node 63 s->m_pUDT->m_pSndQueue->m_pChannel->getSockAddr(s->m_pSelfAddr); 64 65 return 0; 66} 67 68 69 70int CUDT::bind(UDTSOCKET u, const sockaddr* name, int namelen) { 71 try { 72 return s_UDTUnited.bind(u, name, namelen); 73 } catch (CUDTException& e) { 74 s_UDTUnited.setError(new CUDTException(e)); 75 return ERROR; 76 } catch (bad_alloc&) { 77 s_UDTUnited.setError(new CUDTException(3, 2, 0)); 78 return ERROR; 79 } catch (...) { 80 s_UDTUnited.setError(new CUDTException(-1, 0, 0)); 81 return ERROR; 82 } 83} 84 85int CUDT::bind(UDTSOCKET u, UDPSOCKET udpsock) { 86 try { 87 return s_UDTUnited.bind(u, udpsock); 88 } catch (CUDTException& e) { 89 s_UDTUnited.setError(new CUDTException(e)); 90 return ERROR; 91 } catch (bad_alloc&) { 92 s_UDTUnited.setError(new CUDTException(3, 2, 0)); 93 return ERROR; 94 } catch (...) { 95 s_UDTUnited.setError(new CUDTException(-1, 0, 0)); 96 return ERROR; 97 } 98} 99 100 101int bind(UDTSOCKET u, const struct sockaddr* name, int namelen) { 102 return CUDT::bind(u, name, namelen); 103} 104 105int bind2(UDTSOCKET u, UDPSOCKET udpsock) { 106 return CUDT::bind(u, udpsock); 107}

UDT主要提供了两个bind接口,分别是UDT::bind()和,UDT::bind2()。UDT::bind()将一个UDT Socket与一个struct sockaddr对象描述的地址进行绑定,这需要UDT自己先创建相应的系统UDP socket,并将该系统UDP socket绑定到地址,然后把UDT Socket绑定到该系统UDP socket;UDT::bind2()则将一个UDT Socket直接与一个已经创建好的系统UDP socket进行绑定。

这两个API的实现结构与UDT::socket()的实现结构基本一致,一样是分为3层:UDT命名空间中提供了给应用程序调用的接口,可称为UDT API或User API;User API调用CUDT API,这一层主要用来做错误处理,也就是捕获动作实际执行过程中抛出的异常并保存起来,然后给应用程序使用;CUDT API调用CUDTUnited中API的实现。

这里主要来看CUDTUnited中bind()函数的实现。先来看CUDTUnited::bind(const UDTSOCKET u, const sockaddr* name, int namelen)函数的实现:

1. 调用CUDTUnited::locate(),根据SocketID,也就是UDT Socket handle在CUDTUnited的std::map<UDTSOCKET, CUDTSocket*> m_Sockets中找到对应的CUDTSocket结构(src/api.cpp):

1CUDTSocket* CUDTUnited::locate(const UDTSOCKET u) { 2 CGuard cg(m_ControlLock); 3 4 map<UDTSOCKET, CUDTSocket*>::iterator i = m_Sockets.find(u); 5 6 if ((i == m_Sockets.end()) || (i->second->m_Status == CLOSED)) 7 return NULL; 8 9 return i->second; 10}

若找不到,则直接返回;否则,继续执行。

2. 检查CUDTSocket对象的状态,如果当前的状态不为INIT,直接抛异常退出;否则,继续执行。

3. 根据本地IP地址的版本,检查绑定到的目标地址的长度的有效性。IP版本是在UDT Socket创建时指定的。如果无效,则直接抛异常退出;否则,继续执行。

4. 执行相应的CUDT的open()操作(src/core.cpp):

1void CUDT::open() { 2 CGuard cg(m_ConnectionLock); 3 4 // Initial sequence number, loss, acknowledgement, etc. 5 m_iPktSize = m_iMSS - 28; 6 m_iPayloadSize = m_iPktSize - CPacket::m_iPktHdrSize; 7 8 m_iEXPCount = 1; 9 m_iBandwidth = 1; 10 m_iDeliveryRate = 16; 11 m_iAckSeqNo = 0; 12 m_ullLastAckTime = 0; 13 14 // trace information 15 m_StartTime = CTimer::getTime(); 16 m_llSentTotal = m_llRecvTotal = m_iSndLossTotal = m_iRcvLossTotal = m_iRetransTotal = m_iSentACKTotal = 17 m_iRecvACKTotal = m_iSentNAKTotal = m_iRecvNAKTotal = 0; 18 m_LastSampleTime = CTimer::getTime(); 19 m_llTraceSent = m_llTraceRecv = m_iTraceSndLoss = m_iTraceRcvLoss = m_iTraceRetrans = m_iSentACK = m_iRecvACK = 20 m_iSentNAK = m_iRecvNAK = 0; 21 m_llSndDuration = m_llSndDurationTotal = 0; 22 23 // structures for queue 24 if (NULL == m_pSNode) 25 m_pSNode = new CSNode; 26 m_pSNode->m_pUDT = this; 27 m_pSNode->m_llTimeStamp = 1; 28 m_pSNode->m_iHeapLoc = -1; 29 30 if (NULL == m_pRNode) 31 m_pRNode = new CRNode; 32 m_pRNode->m_pUDT = this; 33 m_pRNode->m_llTimeStamp = 1; 34 m_pRNode->m_pPrev = m_pRNode->m_pNext = NULL; 35 m_pRNode->m_bOnList = false; 36 37 m_iRTT = 10 * m_iSYNInterval; 38 m_iRTTVar = m_iRTT >> 1; 39 m_ullCPUFrequency = CTimer::getCPUFrequency(); 40 41 // set up the timers 42 m_ullSYNInt = m_iSYNInterval * m_ullCPUFrequency; 43 44 // set minimum NAK and EXP timeout to 100ms 45 m_ullMinNakInt = 300000 * m_ullCPUFrequency; 46 m_ullMinExpInt = 300000 * m_ullCPUFrequency; 47 48 m_ullACKInt = m_ullSYNInt; 49 m_ullNAKInt = m_ullMinNakInt; 50 51 uint64_t currtime; 52 CTimer::rdtsc(currtime); 53 m_ullLastRspTime = currtime; 54 m_ullNextACKTime = currtime + m_ullSYNInt; 55 m_ullNextNAKTime = currtime + m_ullNAKInt; 56 57 m_iPktCount = 0; 58 m_iLightACKCount = 1; 59 60 m_ullTargetTime = 0; 61 m_ullTimeDiff = 0; 62 63 // Now UDT is opened. 64 m_bOpened = true; 65}

在这个函数中,主要还是对变量的初始化,后面会再结合UDT可靠传输的具体机制,来说明这些变量的具体含义。

5. 执行updateMux()函数更新UDT Socket的多路复用器的相关信息,后面我们会再来详细了解这个更新操作。

6. 将CUDTSocket对象的状态更新为OPENED。

7. 将发送队列的Channel的地址信息拷贝到本节点的s->m_pSelfAddr,m_pSelfAddrde对象的内存空间是在创建UDT Socket的CUDTUnited::newSocket()函数中分配的。

后面会再来解释UDT中Channel和多路复用器Multipexer的含义。

8. 返回0给调用者表示成功结束。

再来看CUDTUnited::bind(UDTSOCKET u, UDPSOCKET udpsock)函数将UDT Socket绑定到一个已经创建好的系统UDP socket的过程:

1. 调用CUDTUnited::locate(),根据SocketID,也就是UDT Socket handle在CUDTUnited的std::map<UDTSOCKET, CUDTSocket*> m_Sockets中找到对应的CUDTSocket结构。若找不到,则直接返回;否则,继续执行。

2. 检查CUDTSocket对象的状态,如果当前的状态不为INIT,直接抛异常退出;否则,继续执行。

3. 获取系统UDP socket的网络地址(含端口信息)。若获取失败则抛异常推出;否则,继续执行。

4. 执行相应的CUDT的open()操作对一些变量进行初始化。

5. 执行updateMux()函数更新UDT Socket的多路复用器的相关信息,后面我们会再来详细了解这个更新操作。

6. 将CUDTSocket对象的状态更新为OPENED。

7. 将发送队列的Channel的地址信息拷贝到本节点的s->m_pSelfAddr,m_pSelfAddrde对象的内存空间是在创建UDT Socket的CUDTUnited::newSocket()函数中分配的。

后面会再来解释UDT中Channel和多路复用器Multipexer的含义。

8. 返回0给调用者表示成功结束。m_MultiplexerLock

总体来说,bind操作使的UDT Socket状态机的状态由INIT状态,转换到了OPENED状态。

CUDTUnited的这两个bind()函数有如此多的重复逻辑,总让人觉得,是有方法做进一步的抽象,以消除重复的逻辑,并使这两个函数的实现都更加精简的。

UDT Socket与多路复用器的关联

bind()操作所做的最最重要的事大概就是将UDT Socket与多路复用器关联,也就是CUDTUnited::updateMux()函数的执行了。为了后面能够更清晰地说明更新多路复用器的操作过程,这里先说明一下UDT的多路复用器CMultiplexer、通道CChannel、发送队列CSndQueue和接收队列CRcvQueue的含义。

UDT中的通道CChannel是系统UDP socket的一个封装,它主要封装了系统UDP socket handle,IP版本号,socket地址的长度,发送缓冲区的大小及接收缓冲区的大小等信息,并提供了用于操作 系统UDP socket进行数据收发或属性设置等动作的函数。我们可以看一下这个class的定义(src/channel.h):

1class CChannel { 2 public: 3 CChannel(); 4 CChannel(int version); 5 ~CChannel(); 6 7 // Functionality: 8 // Open a UDP channel. 9 // Parameters: 10 // 0) [in] addr: The local address that UDP will use. 11 // Returned value: 12 // None. 13 void open(const sockaddr* addr = NULL); 14 15 // Functionality: 16 // Open a UDP channel based on an existing UDP socket. 17 // Parameters: 18 // 0) [in] udpsock: UDP socket descriptor. 19 // Returned value: 20 // None. 21 void open(UDPSOCKET udpsock); 22 23 // Functionality: 24 // Disconnect and close the UDP entity. 25 // Parameters: 26 // None. 27 // Returned value: 28 // None. 29 void close() const; 30 31 // Functionality: 32 // Get the UDP sending buffer size. 33 // Parameters: 34 // None. 35 // Returned value: 36 // Current UDP sending buffer size. 37 int getSndBufSize(); 38 39 // Functionality: 40 // Get the UDP receiving buffer size. 41 // Parameters: 42 // None. 43 // Returned value: 44 // Current UDP receiving buffer size. 45 int getRcvBufSize(); 46 47 // Functionality: 48 // Set the UDP sending buffer size. 49 // Parameters: 50 // 0) [in] size: expected UDP sending buffer size. 51 // Returned value: 52 // None. 53 void setSndBufSize(int size); 54 55 // Functionality: 56 // Set the UDP receiving buffer size. 57 // Parameters: 58 // 0) [in] size: expected UDP receiving buffer size. 59 // Returned value: 60 // None. 61 void setRcvBufSize(int size); 62 63 // Functionality: 64 // Query the socket address that the channel is using. 65 // Parameters: 66 // 0) [out] addr: pointer to store the returned socket address. 67 // Returned value: 68 // None. 69 void getSockAddr(sockaddr* addr) const; 70 71 // Functionality: 72 // Send a packet to the given address. 73 // Parameters: 74 // 0) [in] addr: pointer to the destination address. 75 // 1) [in] packet: reference to a CPacket entity. 76 // Returned value: 77 // Actual size of data sent. 78 int sendto(const sockaddr* addr, CPacket& packet) const; 79 80 // Functionality: 81 // Receive a packet from the channel and record the source address. 82 // Parameters: 83 // 0) [in] addr: pointer to the source address. 84 // 1) [in] packet: reference to a CPacket entity. 85 // Returned value: 86 // Actual size of data received. 87 int recvfrom(sockaddr* addr, CPacket& packet) const; 88 89 private: 90 void setUDPSockOpt(); 91 92 private: 93 int m_iIPversion; // IP version 94 int m_iSockAddrSize; // socket address structure size (pre-defined to avoid run-time test) 95 96 UDPSOCKET m_iSocket; // socket descriptor 97 98 int m_iSndBufSize; // UDP sending buffer size 99 int m_iRcvBufSize; // UDP receiving buffer size 100};

接收队列CRcvQueue在初始化时会起一个线程,该线程在被停掉前,会不断地由CChannel接收其它节点发送过来的UDP消息,可以将这个线程看做是listening在系统UDP 端口上的一个UDP Server。在接收到消息之后,该线程会根据消息的类型及目标 SocketID,把消息dispatch给不同的UDT Socket的CUDT对象。比如对于Handshake类型的消息就会dispatch给listening的UDT Socket的CUDT对象。后面我们研究具体的消息收发的时候再来仔细看这个类的设计。

发送队列CSndQueue,主要用于同步地向特定的目标发送一个UDT的Packet,或者在适当的时机异步地发送一些消息,它同样会在初始化是起一个线程,用来执行异步地发送任务。这个class是UDT做可靠传输的一个比较关键的class,后面我们研究具体的消息收发的时候再来仔细看这个类的设计。

UDT的多路复用器结构CMultiplexer将所有这些与特定的系统UDP socket相关联的CChannel,CRcvQueue,CSndQueue包在一起,并描述了这个系统UDP socket收发的数据的一些公有属性,有UDP 端口号,IP版本号,最大的包大小,引用计数,是否可复用,及用做哈希索引的ID等。可以看一下这个class的定义:

1struct CMultiplexer { 2 CSndQueue* m_pSndQueue; // The sending queue 3 CRcvQueue* m_pRcvQueue; // The receiving queue 4 CChannel* m_pChannel; // The UDP channel for sending and receiving 5 CTimer* m_pTimer; // The timer 6 7 int m_iPort; // The UDP port number of this multiplexer 8 int m_iIPversion; // IP version 9 int m_iMSS; // Maximum Segment Size 10 int m_iRefCount; // number of UDT instances that are associated with this multiplexer 11 bool m_bReusable; // if this one can be shared with others 12 13 int m_iID; // multiplexer ID 14 15 CMultiplexer() 16 : m_pSndQueue(NULL), 17 m_pRcvQueue(NULL), 18 m_pChannel(NULL), 19 m_pTimer(NULL), 20 m_iPort(0), 21 m_iIPversion(0), 22 m_iMSS(0), 23 m_iRefCount(0), 24 m_bReusable(true), 25 m_iID(0) { 26 } 27};

接着来看CUDTUnited::updateMux()函数的定义(src/api.cpp):

1void CUDTUnited::updateMux(CUDTSocket* s, const sockaddr* addr, const UDPSOCKET* udpsock) { 2 CGuard cg(m_ControlLock); 3 4 CMultiplexer m; 5 if ((s->m_pUDT->m_bReuseAddr) && (NULL != addr)) { 6 int port = (AF_INET == s->m_pUDT->m_iIPversion) ? 7 ntohs(((sockaddr_in*) addr)->sin_port) : ntohs(((sockaddr_in6*) addr)->sin6_port); 8 9 // find a reusable address 10 for (map<int, CMultiplexer>::iterator i = m_mMultiplexer.begin(); i != m_mMultiplexer.end(); ++i) { 11 if ((i->second.m_iIPversion == s->m_pUDT->m_iIPversion) && (i->second.m_iMSS == s->m_pUDT->m_iMSS) 12 && i->second.m_bReusable) { 13 if (i->second.m_iPort == port) { 14 // reuse the existing multiplexer 15 m = i->second; 16 break; 17 } 18 } 19 } 20 } 21 22 // a new multiplexer is needed 23 if (m.m_iID == 0) { 24 m.m_iMSS = s->m_pUDT->m_iMSS; 25 m.m_iIPversion = s->m_pUDT->m_iIPversion; 26 m.m_bReusable = s->m_pUDT->m_bReuseAddr; 27 m.m_iID = s->m_SocketID; 28 29 m.m_pChannel = new CChannel(s->m_pUDT->m_iIPversion); 30 m.m_pChannel->setSndBufSize(s->m_pUDT->m_iUDPSndBufSize); 31 m.m_pChannel->setRcvBufSize(s->m_pUDT->m_iUDPRcvBufSize); 32 33 try { 34 if (NULL != udpsock) 35 m.m_pChannel->open(*udpsock); 36 else 37 m.m_pChannel->open(addr); 38 } catch (CUDTException& e) { 39 m.m_pChannel->close(); 40 delete m.m_pChannel; 41 throw e; 42 } 43 44 sockaddr* sa = 45 (AF_INET == s->m_pUDT->m_iIPversion) ? (sockaddr*) new sockaddr_in : (sockaddr*) new sockaddr_in6; 46 m.m_pChannel->getSockAddr(sa); 47 m.m_iPort = (AF_INET == s->m_pUDT->m_iIPversion) ? 48 ntohs(((sockaddr_in*) sa)->sin_port) : ntohs(((sockaddr_in6*) sa)->sin6_port); 49 if (AF_INET == s->m_pUDT->m_iIPversion) 50 delete (sockaddr_in*) sa; 51 else 52 delete (sockaddr_in6*) sa; 53 54 m.m_pTimer = new CTimer; 55 56 m.m_pSndQueue = new CSndQueue; 57 m.m_pSndQueue->init(m.m_pChannel, m.m_pTimer); 58 m.m_pRcvQueue = new CRcvQueue; 59 m.m_pRcvQueue->init(32, s->m_pUDT->m_iPayloadSize, m.m_iIPversion, 1024, m.m_pChannel, m.m_pTimer); 60 61 m_mMultiplexer[m.m_iID] = m; 62 } 63 64 ++m.m_iRefCount; 65 s->m_pUDT->m_pSndQueue = m.m_pSndQueue; 66 s->m_pUDT->m_pRcvQueue = m.m_pRcvQueue; 67 s->m_iMuxID = m.m_iID; 68}

1. 这个函数首先会在已经创建的多路复用器的map中查找,看看是否存在 要与多路复用器关联的UDT Socket可用的多路复用器存在。对于一个UDT Socket来说,UDT Socket本身网络地址可复用,且某个多路复用器同时满足它的CChannel的UDP端口号与UDT Socket要bind的目标UDP端口号匹配,它的CChannel的IP地址版本及MSS与UDT Socket的IP地址版本及MSS匹配,它本身可复用,则该多路复用器就是该UDT Socket可用的多路复用器。

2. 若在前面的步骤中,没有找到可用的多路复用器,则创建一个。

根据UDT Socket的MSS值,IP版本号,及地址的可复用性来初始化CMultiplexer的对应值。设置CMultiplexer的ID为UDT Socket的SocketID。也就是说,某个CMultiplexer的ID就是与它关联的首个UDT Socket的SocketID。

创建CChannel,设置系统UDP socket发送缓冲区及接收缓冲区的大小。并执行CChannel的open()操作。在CChannel::open()中如果不是绑定的已经创建好的系统UDP socket的话,它会自行创建系统UDP socket,并绑定到目标端口上。

获取CChannel实际绑定的UDP端口号,赋值给m.m_iPort。

创建CTimer。

创建并初始化CSndQueue。

创建并初始化CRcvQueue。

将新建的CMultiplexer放进std::map<int, CMultiplexer> m_mMultiplexer中。

在CUDTUnited类定义中可以看到如下几行:

1private: 2 std::map<int, CMultiplexer> m_mMultiplexer; // UDP multiplexer 3 pthread_mutex_t m_MultiplexerLock;

原本设计似乎是要用m_MultiplexerLock来保证对m_mMultiplexer多线程的互斥访问的,但却没有一个地方有用到这个m_MultiplexerLock。不知是发现保护全无必要,还是有所遗漏?

3. 将UDT Socket与多路复用器关联起来,不管是找到的现成可用的,还是完全新创建的。这里可以看到所谓的将UDT Socket与多路复用器关联的含义,即是让CUDTSocket的CUDT对象m_pUDT的发送队列和接收队列指向CMultiplexer的发送队列和接收队列,设置CUDTSocket的多路复用器ID为CMultiplexer的ID m_iID,这样后面CUDTSocket和CUDT就可以使用发送队列CSndQueue和接收队列CRcvQueue进行数据的收发,并可在需要的时候找到相关的CMultiplexer对象了。

自此之后,CUDTSocket就有了可以用来收发数据的设施了。

总结一下UDT bind的主要过程。UDT bind过程中,做的最主要的事情就是,根据一个已经创建好的UDT Socket的一些信息及要绑定的本地UDP端口,找到或创建一个多路复用器CMultiplexer,将UDT Socket与该CMultiplexer关联,即设置CUDTSocket的多路复用器ID m_iMuxID为该CMultiplexer的ID,UDT Socket的发送队列指针和接收队列指针指向该CMultiplexer的发送队列和接收队列。后续UDT Socket就可以通过发送队列/接收队列及它们的CChannel进行数据的收发了。在这个过程中,UDT Socket状态机完成了状态由INIT到OPENED的转变。

listen过程

在UDT Server端,对UDT Socket执行了bind操作之后,就可以执行listen来等待其它节点的连接了。这里来看下UDT listen的过程(src/api.cpp):

1int CUDTUnited::listen(const UDTSOCKET u, int backlog) { 2 CUDTSocket* s = locate(u); 3 if (NULL == s) 4 throw CUDTException(5, 4, 0); 5 6 CGuard cg(s->m_ControlLock); 7 8 // do nothing if the socket is already listening 9 if (LISTENING == s->m_Status) 10 return 0; 11 12 // a socket can listen only if is in OPENED status 13 if (OPENED != s->m_Status) 14 throw CUDTException(5, 5, 0); 15 16 // listen is not supported in rendezvous connection setup 17 if (s->m_pUDT->m_bRendezvous) 18 throw CUDTException(5, 7, 0); 19 20 if (backlog <= 0) 21 throw CUDTException(5, 3, 0); 22 23 s->m_uiBackLog = backlog; 24 25 try { 26 s->m_pQueuedSockets = new set<UDTSOCKET>; 27 s->m_pAcceptSockets = new set<UDTSOCKET>; 28 } catch (...) { 29 delete s->m_pQueuedSockets; 30 delete s->m_pAcceptSockets; 31 throw CUDTException(3, 2, 0); 32 } 33 34 s->m_pUDT->listen(); 35 36 s->m_Status = LISTENING; 37 38 return 0; 39} 40 41 42int CUDT::listen(UDTSOCKET u, int backlog) { 43 try { 44 return s_UDTUnited.listen(u, backlog); 45 } catch (CUDTException& e) { 46 s_UDTUnited.setError(new CUDTException(e)); 47 return ERROR; 48 } catch (bad_alloc&) { 49 s_UDTUnited.setError(new CUDTException(3, 2, 0)); 50 return ERROR; 51 } catch (...) { 52 s_UDTUnited.setError(new CUDTException(-1, 0, 0)); 53 return ERROR; 54 } 55} 56 57 58int listen(UDTSOCKET u, int backlog) { 59 return CUDT::listen(u, backlog); 60}

这个API的实现同样分为3层,UDT命名空间提供的直接给应用程序调用User API层,CUDT API层用于做异常处理,CUDTUnited具体实现API的功能。这里直接来分析CUDTUnited::listen()函数:

1. 调用CUDTUnited::locate(),查找UDT Socket对应的CUDTSocket结构。若找不到,则抛出异常直接返回;否则,继续执行。

2. 检查CUDTSocket对象的状态,如果当前的状态为LISTENING,则说明UDT Socket已经处于监听状态了,直接返回;若当前状态不为OPENED,直接抛异常退出,否则,继续执行。这就限制了只有经过了bind操作的UDT Socket才能监听,也就是UDT Socket的状态只能由OPENED转为LISTENING。

3. 检查是否是rendezvous的UDT Socket,若是则抛出异常推出。这确保在监听的UDT Socket不能为rendezvous的。

4. 检查传入的backlog参数并进行设置。backlog参数用于指定Listening的UDT Socket同一时刻能够处理的最大的等待连接的请求数。Listening的UDT Socket在收到连接请求的Handshake消息后,经过几次来回确认,会创建新的UDT Socket以便于通过UDT::accept()函数返回给应用程序,用于与请求连接的发起方进行通信。backlog值用于限定,还没有通过accept()返回的新创建的UDT Socket的个数。

5. 创建两个UDTSOCKET的集合m_pQueuedSockets和m_pAcceptSockets,前者为Listening的UDT Socket的连接已经成功建立但还未通过UDT::accept()返回给应用程序的UDT Socket的集合;而后者则是已经通过UDT::accept()返回给应用程序的UDT Socket的集合。

6. 执行UDTSocket的CUDT的listen()操作,可以看一下CUDT listen动作的具体含义(src/core.cpp):

1void CUDT::listen() { 2 CGuard cg(m_ConnectionLock); 3 4 if (!m_bOpened) 5 throw CUDTException(5, 0, 0); 6 7 if (m_bConnecting || m_bConnected) 8 throw CUDTException(5, 2, 0); 9 10 // listen can be called more than once 11 if (m_bListening) 12 return; 13 14 // if there is already another socket listening on the same port 15 if (m_pRcvQueue->setListener(this) < 0) 16 throw CUDTException(5, 11, 0); 17 18 m_bListening = true; 19}

先是进行状态的合法性检查。

然后执行m_pRcvQueue->setListener(this),将本CUDT设置为接收队列的listener。

最后设置CUDT的状态m_bListening为true。

这里可以看出CUDTSocket与CUDT是表示UDT Socket的两层状态机,它们的状态之间有关联,但又有各自的描述方法。这样似乎大大增加了这个UDT Socket状态管理的复杂度了。

再来看一下CRcvQueue::setListener()(src/queue.cpp):

1int CRcvQueue::setListener(CUDT* u) { 2 CGuard lslock(m_LSLock); 3 4 if (NULL != m_pListener) 5 return -1; 6 7 m_pListener = u; 8 return 0; 9}

设置接收队列CRcvQueue的Listener。

7. 设置CUDTSocket的状态为LISTENING并返回。

可以看到对于UDT::listen()的调用,促使UDT Socket的状态由OPENED转换为了LISTENING。UDT::listen()主要的作用就是 为与UDT Socket关联的特定端口上的多路复用器CMultiplexer的接收队列CRcvQueue设置listener,这个动作最主要的意义在于消息的dispatch。我们知道CMultiplexer的接收队列CRcvQueue在创建、初始化时会起一个线程,不断地试图从网络接收UDP消息,在收到消息之后,将消息dispatch给不同的UDT Socket处理,其中的Handshake等消息,就会被dispatch给listener CUDT处理。后面在具体研究消息的收发时会再来详细研究这个过程。

accept过程

UDT Server端在对listening执行了UDT::listen()操作之后,就可以执行UDT::accept()操作来等待其它节点连接自己了。来看一下UDT::accept()的执行过程(src/api.cpp):

1UDTSOCKET CUDTUnited::accept(const UDTSOCKET listen, sockaddr* addr, int* addrlen) { 2 if ((NULL != addr) && (NULL == addrlen)) 3 throw CUDTException(5, 3, 0); 4 5 CUDTSocket* ls = locate(listen); 6 7 if (ls == NULL) 8 throw CUDTException(5, 4, 0); 9 10 // the "listen" socket must be in LISTENING status 11 if (LISTENING != ls->m_Status) 12 throw CUDTException(5, 6, 0); 13 14 // no "accept" in rendezvous connection setup 15 if (ls->m_pUDT->m_bRendezvous) 16 throw CUDTException(5, 7, 0); 17 18 UDTSOCKET u = CUDT::INVALID_SOCK; 19 bool accepted = false; 20 21 // !!only one conection can be set up each time!! 22#ifndef WIN32 23 while (!accepted) { 24 pthread_mutex_lock(&(ls->m_AcceptLock)); 25 26 if ((LISTENING != ls->m_Status) || ls->m_pUDT->m_bBroken) { 27 // This socket has been closed. 28 accepted = true; 29 } else if (ls->m_pQueuedSockets->size() > 0) { 30 u = *(ls->m_pQueuedSockets->begin()); 31 ls->m_pAcceptSockets->insert(ls->m_pAcceptSockets->end(), u); 32 ls->m_pQueuedSockets->erase(ls->m_pQueuedSockets->begin()); 33 accepted = true; 34 } else if (!ls->m_pUDT->m_bSynRecving) { 35 accepted = true; 36 } 37 38 if (!accepted && (LISTENING == ls->m_Status)) 39 pthread_cond_wait(&(ls->m_AcceptCond), &(ls->m_AcceptLock)); 40 41 if (ls->m_pQueuedSockets->empty()) 42 m_EPoll.update_events(listen, ls->m_pUDT->m_sPollID, UDT_EPOLL_IN, false); 43 44 pthread_mutex_unlock(&(ls->m_AcceptLock)); 45 } 46#else 47 while (!accepted) 48 { 49 WaitForSingleObject(ls->m_AcceptLock, INFINITE); 50 51 if (ls->m_pQueuedSockets->size() > 0) 52 { 53 u = *(ls->m_pQueuedSockets->begin()); 54 ls->m_pAcceptSockets->insert(ls->m_pAcceptSockets->end(), u); 55 ls->m_pQueuedSockets->erase(ls->m_pQueuedSockets->begin()); 56 57 accepted = true; 58 } 59 else if (!ls->m_pUDT->m_bSynRecving) 60 accepted = true; 61 62 ReleaseMutex(ls->m_AcceptLock); 63 64 if (!accepted & (LISTENING == ls->m_Status)) 65 WaitForSingleObject(ls->m_AcceptCond, INFINITE); 66 67 if ((LISTENING != ls->m_Status) || ls->m_pUDT->m_bBroken) 68 { 69 // Send signal to other threads that are waiting to accept. 70 SetEvent(ls->m_AcceptCond); 71 accepted = true; 72 } 73 74 if (ls->m_pQueuedSockets->empty()) 75 m_EPoll.update_events(listen, ls->m_pUDT->m_sPollID, UDT_EPOLL_IN, false); 76 } 77#endif 78 79 if (u == CUDT::INVALID_SOCK) { 80 // non-blocking receiving, no connection available 81 if (!ls->m_pUDT->m_bSynRecving) 82 throw CUDTException(6, 2, 0); 83 84 // listening socket is closed 85 throw CUDTException(5, 6, 0); 86 } 87 88 if ((addr != NULL) && (addrlen != NULL)) { 89 if (AF_INET == locate(u)->m_iIPversion) 90 *addrlen = sizeof(sockaddr_in); 91 else 92 *addrlen = sizeof(sockaddr_in6); 93 94 // copy address information of peer node 95 memcpy(addr, locate(u)->m_pPeerAddr, *addrlen); 96 } 97 98 return u; 99} 100 101 102 103UDTSOCKET CUDT::accept(UDTSOCKET u, sockaddr* addr, int* addrlen) { 104 try { 105 return s_UDTUnited.accept(u, addr, addrlen); 106 } catch (CUDTException& e) { 107 s_UDTUnited.setError(new CUDTException(e)); 108 return INVALID_SOCK; 109 } catch (...) { 110 s_UDTUnited.setError(new CUDTException(-1, 0, 0)); 111 return INVALID_SOCK; 112 } 113} 114 115 116UDTSOCKET accept(UDTSOCKET u, struct sockaddr* addr, int* addrlen) { 117 return CUDT::accept(u, addr, addrlen); 118}

这个API实现的3层结构与UDT::bind(),UDT::listen()一样,不再赘述。来看CUDTUnited::accept()的实现:

1. 调用CUDTUnited::locate(),查找UDT Socket对应的CUDTSocket结构。若找不到,则抛出异常直接返回;否则,继续执行。

2. 检查CUDTSocket对象的状态。可见CUDTUnited::accept()操作要求相应的UDT Socket必须处于LISTENING状态,且不能为Rendezvous模式。这个地方对于ls->m_pUDT->m_bRendezvous的检查似乎有些多余了,在CUDTUnited::listen()中可以看到,如果UDT Socket处于Rendezvous模式的话,根本就不可能完成状态由OPENED到LISTENING的转换,因而对于UDT Socket LISTENING状态的检查已经足够了。

3. 通过一个循环来等待其它节点的连接。这指的是,等待ls->m_pQueuedSockets中被放入为新的连接创建的UDT Socket。有新的连接时,CUDTUnited::accept()线程被唤醒,它会将UDT Socket从ls->m_pQueuedSockets中移到ls->m_pAcceptSockets,并准备将UDT Socket返回给调用者。当然CUDTUnited::accept()的等待过程结束的条件不只是有新连接进来,在Listening的UDT Socket被closed掉时,ls->m_pUDT->m_bBroken会被设置,UDT Socket的状态也可能会发生变化,此时等待过程会结束;或者UDT Socket处于同步接收状态,则无论是否有新连接,等待过程都会尽快结束。

4. 等待连接的过程意外退出,也就是在没有等到新连接进来的情况下等待过程就退出了的情况下,抛出异常退出。如果UDT Socket处于同步接收状态,抛出某个类型的异常,否则抛出另外一种类型的异常来表示UDT Socket被关闭了。这个地方的逻辑,向调用者展示的异常信息可能具有误导性,比如一个同步接收的UDT Socket被关闭了,向调用者展示的信息似乎仍然表明,UDT Socket是由于同步接收的问题而没有等到新连接进来才退出的。

5. 等到了新连接进来的情况下,将发起端的网络地址拷贝给调用者。

6. 将新UDT Socket的SocketID返回给调用者。

可以看到,UDT::accept()这个地方是一个典型的生产者-消费者模型。UDT::accept()是消费者,消费的对象是ls->m_pQueuedSockets中的UDT Socket。我们分析UDT::accept()函数的实现,只能看到这个关于生产-消费的故事的一半,另一半关于生产的故事则需要通过更仔细地分析CRcvQueue::worker()的执行来了解了。

总结一下这几个操作与Listening Socket状态变化之间的关系,如下图所示:

Done.

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )

KVM调整cpu和内存

一.修改kvm虚拟机的配置1、virsheditcentos7找到“memory”和“vcpu”标签,将<namecentos7</name<uuid2220a6d1a36a4fbb8523e078b3dfe795</uuid