MCPcopy Create free account
hub / github.com/catboost/catboost / AddDataToPacketQueue

Function AddDataToPacketQueue

library/cpp/netliba/v12/udp_host_protocol.h:781–825  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

779 ///////////////////////////////////////////////////////////////////////////////
780
781 inline void AddDataToPacketQueue(TUdpSocket& s, char* packetBuffer,
782 TConnection* connection, const ui64 transferId, const TUdpOutTransfer& xfer,
783 const int packetId, const int dataSize) {
784 Y_ASSERT(xfer.PacketSize == UDP_PACKET_SIZE || xfer.PacketSize == UDP_SMALL_PACKET_SIZE || xfer.PacketSize == UDP_XSMALL_PACKET_SIZE);
785 Y_ASSERT(xfer.LastPacketSize < xfer.PacketSize);
786 Y_ASSERT(dataSize == xfer.PacketSize || dataSize == xfer.LastPacketSize);
787
788 EUdpCmd cmd = xfer.PacketSize == UDP_PACKET_SIZE ? DATA : DATA_SMALL;
789 TPosixSharedMemory* shm = xfer.Data->GetSharedData();
790
791 char* pktData = packetBuffer;
792 TOptionsVector TransferOptions;
793 if (packetId == 0) {
794 if (xfer.PacketPriority == PP_HIGH || xfer.PacketPriority == PP_SYSTEM) {
795 TransferOptions.TransferOpt.SetHighPriority(true);
796 }
797 if (shm) {
798 //fprintf(stderr, "TransferOptions.TransferOpt.SetSharedMemory sz: %i\n", (int)shm->GetSize());
799 //We assume size of one shared memory region can`t be more than size of one transfer which can`t be more than 1.8G
800 Y_ASSERT(shm->GetSizeT() <= std::numeric_limits<ui32>::max());
801 TransferOptions.TransferOpt.SetSharedMemory((ui32)shm->GetSizeT(), shm->GetId());
802 }
803 }
804 /*
805 Cerr << GetAddressAsString(connection->GetAddress()) << " Sending packet " << packetId << " of " << xfer.PacketCount
806 << " with xfer.PacketSize=" << xfer.PacketSize
807 << ", *conn=" << size_t(connection)
808 << ", conn.SmallMtuUseXs=" << connection->GetSmallMtuUseXs()
809 << ", dataSize=" << dataSize
810 << ", cmde=" << int(cmd)
811 << Endl;
812 */
813 WriteDataPacketHeader(&pktData, cmd, connection, transferId, packetId, xfer.AckTos, xfer.NetlibaColor, &TransferOptions);
814
815 // TODO: we can avoid memcpy here, but xfer.Data is a list of data chunks and we need whole packet...
816 // we may use iovec!
817 TBlockChainIterator dataReader(xfer.Data->GetChain());
818 dataReader.Seek((int)(packetId * xfer.PacketSize));
819 dataReader.Read(pktData, (int)dataSize);
820 pktData += dataSize;
821
822 // AddPacketToQueue YASSERTs buffer overflow
823 s.AddPacketToQueue(pktData - packetBuffer, {connection->GetWinsockAddress(), connection->GetWinsockMyAddress()},
824 xfer.DataTos, xfer.PacketSize);
825 }
826
827 inline bool ReadDataPacket(const EUdpCmd cmd, const char** pktData, const char* pktEnd, const int packetId,
828 TIntrusivePtr<TPosixSharedMemory>* shm, size_t* packetSize, const TOptionsVector& opt) {

Callers 1

SendTransferPacketMethod · 0.85

Calls 8

WriteDataPacketHeaderFunction · 0.85
GetSizeTMethod · 0.80
AddPacketToQueueMethod · 0.80
maxFunction · 0.50
GetSharedDataMethod · 0.45
GetIdMethod · 0.45
SeekMethod · 0.45
ReadMethod · 0.45

Tested by

no test coverage detected