| 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) { |
no test coverage detected