MCPcopy Create free account
hub / github.com/PurpleI2P/i2pd / SendBuffer

Method SendBuffer

libi2pd/Streaming.cpp:947–1132  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

945 }
946
947 void Stream::SendBuffer ()
948 {
949 if (m_RemoteLeaseSet) // don't scheudle send for first SYN for incoming stream
950 ScheduleSend ();
951 auto ts = i2p::util::GetMillisecondsSinceEpoch ();
952 int numMsgs = m_WindowSize - m_SentPackets.size ();
953 if (numMsgs <= 0 || !m_IsSendTime) // window is full
954 {
955 m_LastSendTime = ts;
956 return;
957 }
958 else if (numMsgs > m_NumPacketsToSend)
959 numMsgs = m_NumPacketsToSend;
960
961 if (!m_RemoteLeaseSet) m_RemoteLeaseSet = m_LocalDestination.GetOwner ()->FindLeaseSet (m_RemoteIdentity->GetIdentHash ());
962 if (m_RemoteLeaseSet)
963 {
964 if (!m_RoutingSession)
965 m_RoutingSession = m_LocalDestination.GetOwner ()->GetRoutingSession (m_RemoteLeaseSet, true, false);
966 }
967
968 if (m_RoutingSession)
969 {
970 m_IsJavaClient = m_RoutingSession->IsWithJava ();
971 if (m_IsJavaClient) m_MaxWindowSize = 64;
972 int numSentPackets = m_RoutingSession->NumSentPackets ();
973 int numPacketsToSend = m_MaxWindowSize - numSentPackets;
974 if (numPacketsToSend <= 0) // shared window is full
975 {
976 if (m_LastReceivedSequenceNumber <= 0 && m_SequenceNumber == 0)
977 {
978 LogPrint (eLogWarning, "Streaming: limit of unacknowledged packets has been reached, terminate, rSID=", m_RecvStreamID, ", sSID=", m_SendStreamID);
979 m_Status = eStreamStatusReset;
980 Close ();
981 return;
982 }
983 m_LastSendTime = ts;
984 return;
985 }
986 else if (numMsgs > numPacketsToSend)
987 numMsgs = numPacketsToSend;
988 }
989 bool isNoAck = m_LastReceivedSequenceNumber < 0; // first packet
990 std::vector<Packet *> packets;
991 while ((m_Status == eStreamStatusNew) || (IsEstablished () && !m_SendBuffer.IsEmpty () && numMsgs > 0))
992 {
993 Packet * p = m_LocalDestination.NewPacket ();
994 uint8_t * packet = p->GetBuffer ();
995 // TODO: implement setters
996 size_t size = 0;
997 htobe32buf (packet + size, m_SendStreamID);
998 size += 4; // sendStreamID
999 htobe32buf (packet + size, m_RecvStreamID);
1000 size += 4; // receiveStreamID
1001 htobe32buf (packet + size, m_SequenceNumber++);
1002 size += 4; // sequenceNum
1003 if (isNoAck)
1004 htobuf32 (packet + size, 0);

Callers 4

SendTunnelDataMsgsMethod · 0.45
SendTunnelDataMsgsToMethod · 0.45
FlushTunnelDataMsgsMethod · 0.45
AsyncSendMethod · 0.45

Calls 15

LogPrintFunction · 0.85
htobe32bufFunction · 0.85
htobuf32Function · 0.85
htobe16bufFunction · 0.85
htobuf16Function · 0.85
GetRoutingSessionMethod · 0.80
IsWithJavaMethod · 0.80
NumSentPacketsMethod · 0.80
NewPacketMethod · 0.80
IsOfflineSignatureMethod · 0.80
push_backMethod · 0.80

Tested by

no test coverage detected