|
|
@ -28,8 +28,12 @@ namespace stream |
|
|
|
Stream::~Stream () |
|
|
|
Stream::~Stream () |
|
|
|
{ |
|
|
|
{ |
|
|
|
m_ReceiveTimer.cancel (); |
|
|
|
m_ReceiveTimer.cancel (); |
|
|
|
while (auto packet = m_ReceiveQueue.Get ()) |
|
|
|
while (!m_ReceiveQueue.empty ()) |
|
|
|
|
|
|
|
{ |
|
|
|
|
|
|
|
auto packet = m_ReceiveQueue.front (); |
|
|
|
|
|
|
|
m_ReceiveQueue.pop (); |
|
|
|
delete packet; |
|
|
|
delete packet; |
|
|
|
|
|
|
|
} |
|
|
|
for (auto it: m_SavedPackets) |
|
|
|
for (auto it: m_SavedPackets) |
|
|
|
delete it; |
|
|
|
delete it; |
|
|
|
} |
|
|
|
} |
|
|
@ -118,7 +122,7 @@ namespace stream |
|
|
|
packet->offset = packet->GetPayload () - packet->buf; |
|
|
|
packet->offset = packet->GetPayload () - packet->buf; |
|
|
|
if (packet->GetLength () > 0) |
|
|
|
if (packet->GetLength () > 0) |
|
|
|
{ |
|
|
|
{ |
|
|
|
m_ReceiveQueue.Put (packet); |
|
|
|
m_ReceiveQueue.push (packet); |
|
|
|
m_ReceiveTimer.cancel (); |
|
|
|
m_ReceiveTimer.cancel (); |
|
|
|
} |
|
|
|
} |
|
|
|
else |
|
|
|
else |
|
|
@ -131,7 +135,6 @@ namespace stream |
|
|
|
LogPrint ("Closed"); |
|
|
|
LogPrint ("Closed"); |
|
|
|
SendQuickAck (); // send ack for close explicitly?
|
|
|
|
SendQuickAck (); // send ack for close explicitly?
|
|
|
|
m_IsOpen = false; |
|
|
|
m_IsOpen = false; |
|
|
|
m_ReceiveQueue.WakeUp (); |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
@ -239,31 +242,25 @@ namespace stream |
|
|
|
|
|
|
|
|
|
|
|
if (SendPacket (packet, size)) |
|
|
|
if (SendPacket (packet, size)) |
|
|
|
LogPrint ("FIN sent"); |
|
|
|
LogPrint ("FIN sent"); |
|
|
|
m_ReceiveQueue.WakeUp (); |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
size_t Stream::ConcatenatePackets (uint8_t * buf, size_t len) |
|
|
|
size_t Stream::ConcatenatePackets (uint8_t * buf, size_t len) |
|
|
|
{ |
|
|
|
{ |
|
|
|
size_t pos = 0; |
|
|
|
size_t pos = 0; |
|
|
|
while (pos < len) |
|
|
|
while (pos < len && !m_ReceiveQueue.empty ()) |
|
|
|
{ |
|
|
|
|
|
|
|
Packet * packet = m_ReceiveQueue.Peek (); |
|
|
|
|
|
|
|
if (packet) |
|
|
|
|
|
|
|
{ |
|
|
|
{ |
|
|
|
|
|
|
|
Packet * packet = m_ReceiveQueue.front (); |
|
|
|
size_t l = std::min (packet->GetLength (), len - pos); |
|
|
|
size_t l = std::min (packet->GetLength (), len - pos); |
|
|
|
memcpy (buf + pos, packet->GetBuffer (), l); |
|
|
|
memcpy (buf + pos, packet->GetBuffer (), l); |
|
|
|
pos += l; |
|
|
|
pos += l; |
|
|
|
packet->offset += l; |
|
|
|
packet->offset += l; |
|
|
|
if (!packet->GetLength ()) |
|
|
|
if (!packet->GetLength ()) |
|
|
|
{ |
|
|
|
{ |
|
|
|
m_ReceiveQueue.Get (); |
|
|
|
m_ReceiveQueue.pop (); |
|
|
|
delete packet; |
|
|
|
delete packet; |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
else // no more data available
|
|
|
|
|
|
|
|
break; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
return pos; |
|
|
|
return pos; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|