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