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

Method SendTransferPacket

library/cpp/netliba/v12/udp_host.cpp:722–775  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

720 }
721
722 TUdpHost::ESentPacketResult TUdpHost::SendTransferPacket(TConnection* connection, TUdpOutTransfer* xfer, ui64 transferId) {
723 NHPTimer::STime tCopy = CurrentT;
724 float deltaT = (float)NHPTimer::GetTimePassed(&tCopy);
725 deltaT = ClampVal(deltaT, 0.0f, UdpTransferTimeout / 3);
726
727 bool isCanceled = false;
728 const int packetId = xfer->AckTracker.GetPacketToSend(deltaT, &isCanceled);
729 if (packetId == -1) {
730 if (isCanceled) {
731 if (xfer->TriedToSendAtLeastOnePacket) {
732 const ui8 err = FlushPackets(); // there must be FlushPackets before any Send.* function -
733 //it guarantee we have some space in packet buffer
734
735 if (err & TUdpHost::EFlashPacketResult::FPR_OUT_TRANSFERS_CHANGED) {
736 //also we must check _current_ transfer because FlushPackets() could remove it
737 if (!connection->IsSendTransferAlive(transferId)) {
738 return TUdpHost::ESentPacketResult::SPR_STOP_SENDING_TRANSFER;
739 }
740 }
741 SendCancelTransfer(S, connection, transferId, xfer->AckTos);
742 xfer->AckTracker.Congestion->ForceTimeAccount();
743 } else {
744 xfer->AckTracker.AckAll(); // isn't necessary actually, but just to be sure...
745 CanceledSend(TTransfer(connection, transferId));
746 }
747 }
748 return TUdpHost::ESentPacketResult::SPR_STOP_SENDING_TRANSFER;
749 }
750 const int dataSize = (packetId == xfer->PacketCount - 1) ? xfer->LastPacketSize : xfer->PacketSize;
751
752 // I intentionally use very specific comparison here to minimize possible backward compatibility issues
753 if (dataSize == UDP_SMALL_PACKET_SIZE && xfer->AckTracker.Congestion->GetMTU() == UDP_XSMALL_PACKET_SIZE) {
754 // Cerr << "SendTransferPacker: dataSize " << dataSize << " > mtu " << xfer->AckTracker.Congestion->GetMTU()
755 // << ", xfer " << ui64(xfer) << " marked as failed" << Endl;
756 FailedSend(TTransfer(connection, transferId));
757 return TUdpHost::ESentPacketResult::SPR_STOP_SENDING_TRANSFER;
758 }
759
760 const std::pair<char* const, ui8>& packetBuffer = GetPacketBuffer(dataSize + PACKET_HEADERS_SIZE, connection, transferId);
761
762 if (!packetBuffer.first) { //buffer overflow, or current transfer removed during packet flushing
763 if (packetBuffer.second & TUdpHost::EFlashPacketResult::FPR_OUT_TRANSFERS_CHANGED) {
764 //transfer removed
765 return TUdpHost::ESentPacketResult::SPR_STOP_SENDING_TRANSFER;
766 }
767 //xfer->AckTracker.AddToResend(packetId); no need here - FlushPackets adds packetId to resend
768 Y_ASSERT(packetBuffer.second & TUdpHost::EFlashPacketResult::FPR_OVERFLOW);
769 return TUdpHost::ESentPacketResult::SPR_OVERFLOW;
770 }
771
772 xfer->TriedToSendAtLeastOnePacket = true;
773 AddDataToPacketQueue(S, packetBuffer.first, connection, transferId, *xfer, packetId, dataSize);
774 return TUdpHost::ESentPacketResult::SPR_OK;
775 }
776
777 bool TUdpHost::SendCycle(ui8 prio) {
778 bool res1 = Connections.ForEachSendingConnection([prio, this](TSendingConnectionsList::iterator& sc) {

Callers

nothing calls this directly

Calls 8

SendCancelTransferFunction · 0.85
AddDataToPacketQueueFunction · 0.85
IsSendTransferAliveMethod · 0.80
TTransferClass · 0.70
GetPacketToSendMethod · 0.45
ForceTimeAccountMethod · 0.45
AckAllMethod · 0.45
GetMTUMethod · 0.45

Tested by

no test coverage detected