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

Method FlushPackets

library/cpp/netliba/v12/udp_host.cpp:878–937  ·  view source on GitHub ↗

ui8 - error bitmap FPR_OVERFLOW - there were non-sent packets because of udp/kernel buffer overflow. FPR_OUT_TRANSFERS_CHANGED - some packets were not send because NO_ROUTE_TO_HOST Informs all transfers about problems ("failed transfer" or "add packet to resend"). Clears TUdpSocket packet queue.

Source from the content-addressed store, hash-verified

876 // Informs all transfers about problems ("failed transfer" or "add packet to resend").
877 // Clears TUdpSocket packet queue.
878 ui8 TUdpHost::FlushPackets() {
879 NewPacketsAfterLastFlush = 0;
880 //The logic is we must inform caller about errors, but it is not good idea
881 //to break flushing in case of NO_ROUTE_TO_HOST error. So, we will accumulate errors
882 //as bitmask and pass it to caller
883 ui8 errRet = TUdpHost::EFlashPacketResult::FPR_OK;
884 for (;;) {
885 size_t numSentPackets;
886 TVector<std::pair<char*, size_t>> failedPackets;
887 TUdpSocket::ESendError err = S.FlushPackets(&numSentPackets, &failedPackets);
888
889 if (err == TUdpSocket::SEND_OK) {
890 Y_ASSERT(S.IsPacketsQueueEmpty());
891 return errRet;
892
893 // most probably out of send buffer space (or something terrible has happened)
894 } else if (err == TUdpSocket::SEND_BUFFER_OVERFLOW) {
895 TVector<std::pair<char*, size_t>> notSentPackets;
896 S.GetPacketsInQueue(&notSentPackets);
897
898 for (size_t i = 0; i != notSentPackets.size(); ++i) {
899 const std::pair<char*, size_t>& notSentPacket = notSentPackets[i];
900 TTransfer transfer;
901 int packetId;
902 if (ParseDataPacketHeader(notSentPacket.first, notSentPacket.first + notSentPacket.second, &transfer, &packetId)) {
903 TConnection* connection = CheckedCast<TConnection*>(transfer.Connection.Get());
904 connection->GetSendQueue().Get(transfer.Id)->AckTracker.AddToResend(packetId);
905 //fprintf(stderr, "Failed send\n");
906 }
907 }
908 MaxWaitTime = 0;
909 S.ClearPacketsQueue(); // MUST be called only after flush (to simulate 100% flush in local connections)
910 Y_ASSERT(S.IsPacketsQueueEmpty());
911 errRet |= TUdpHost::EFlashPacketResult::FPR_OVERFLOW;
912 return errRet;
913
914 } else if (err == TUdpSocket::SEND_NO_ROUTE_TO_HOST || err == TUdpSocket::SEND_EINVAL) {
915 const char* errText = (err == TUdpSocket::SEND_NO_ROUTE_TO_HOST) ? "No route to host" : "Error in value";
916 for (size_t i = 0; i != failedPackets.size(); ++i) {
917 const std::pair<char*, size_t>& failedPacket = failedPackets[i];
918 TTransfer transfer;
919 int packetId;
920 if (ParseDataPacketHeader(failedPacket.first, failedPacket.first + failedPacket.second, &transfer, &packetId)) {
921 FailedSend(transfer);
922 fprintf(stderr, "%s, transfer: %" PRIu64 " failed, packetId: %" PRIi32 "\n",
923 errText, transfer.Id, packetId);
924 }
925 // packet is already dropped by socket queue
926 }
927 errRet |= TUdpHost::EFlashPacketResult::FPR_OUT_TRANSFERS_CHANGED;
928 continue;
929
930 } else {
931 Y_ASSERT(false);
932 break;
933 }
934 }
935 Y_ASSERT(false); // unreachable

Callers

nothing calls this directly

Calls 6

IsPacketsQueueEmptyMethod · 0.80
GetPacketsInQueueMethod · 0.80
ClearPacketsQueueMethod · 0.80
sizeMethod · 0.45
GetMethod · 0.45
AddToResendMethod · 0.45

Tested by

no test coverage detected