|
|
|
@ -24,6 +24,7 @@ namespace llarp
|
|
|
|
|
const ILinkMessage *msg,
|
|
|
|
|
SendStatusHandler callback)
|
|
|
|
|
{
|
|
|
|
|
const uint16_t priority = msg->Priority();
|
|
|
|
|
std::array< byte_t, MAX_LINK_MSG_SIZE > linkmsg_buffer;
|
|
|
|
|
llarp_buffer_t buf(linkmsg_buffer);
|
|
|
|
|
|
|
|
|
@ -40,7 +41,7 @@ namespace llarp
|
|
|
|
|
|
|
|
|
|
if(_linkManager->HasSessionTo(remote))
|
|
|
|
|
{
|
|
|
|
|
QueueOutboundMessage(remote, std::move(message), msg->pathid);
|
|
|
|
|
QueueOutboundMessage(remote, std::move(message), msg->pathid, priority);
|
|
|
|
|
return true;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@ -53,8 +54,9 @@ namespace llarp
|
|
|
|
|
pendingSessionMessageQueues.emplace(remote, MessageQueue());
|
|
|
|
|
|
|
|
|
|
MessageQueueEntry entry;
|
|
|
|
|
entry.message = message;
|
|
|
|
|
entry.router = remote;
|
|
|
|
|
entry.priority = priority;
|
|
|
|
|
entry.message = message;
|
|
|
|
|
entry.router = remote;
|
|
|
|
|
itr_pair.first->second.push(std::move(entry));
|
|
|
|
|
|
|
|
|
|
shouldCreateSession = itr_pair.second;
|
|
|
|
@ -232,13 +234,15 @@ namespace llarp
|
|
|
|
|
bool
|
|
|
|
|
OutboundMessageHandler::QueueOutboundMessage(const RouterID &remote,
|
|
|
|
|
Message &&msg,
|
|
|
|
|
const PathID_t &pathid)
|
|
|
|
|
const PathID_t &pathid,
|
|
|
|
|
uint16_t priority)
|
|
|
|
|
{
|
|
|
|
|
MessageQueueEntry entry;
|
|
|
|
|
entry.message = std::move(msg);
|
|
|
|
|
auto callback_copy = entry.message.second;
|
|
|
|
|
entry.router = remote;
|
|
|
|
|
entry.pathid = pathid;
|
|
|
|
|
entry.priority = priority;
|
|
|
|
|
if(outboundQueue.tryPushBack(std::move(entry))
|
|
|
|
|
!= llarp::thread::QueueReturn::Success)
|
|
|
|
|
{
|
|
|
|
@ -274,11 +278,13 @@ namespace llarp
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
MessageQueue &path_queue = itr_pair.first->second;
|
|
|
|
|
/*
|
|
|
|
|
if(path_queue.size() >= MAX_PATH_QUEUE_SIZE)
|
|
|
|
|
{
|
|
|
|
|
m_queueStats.dropped++;
|
|
|
|
|
path_queue.pop(); // head drop
|
|
|
|
|
}
|
|
|
|
|
*/
|
|
|
|
|
path_queue.push(std::move(entry));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
@ -310,10 +316,9 @@ namespace llarp
|
|
|
|
|
auto &non_routing_mq = outboundMessageQueues[zeroID];
|
|
|
|
|
while(not non_routing_mq.empty())
|
|
|
|
|
{
|
|
|
|
|
MessageQueueEntry entry = std::move(non_routing_mq.front());
|
|
|
|
|
non_routing_mq.pop();
|
|
|
|
|
|
|
|
|
|
const MessageQueueEntry &entry = non_routing_mq.top();
|
|
|
|
|
Send(entry.router, entry.message);
|
|
|
|
|
non_routing_mq.pop();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
size_t empty_count = 0;
|
|
|
|
@ -349,10 +354,11 @@ namespace llarp
|
|
|
|
|
auto &message_queue = outboundMessageQueues[pathid];
|
|
|
|
|
if(message_queue.size() > 0)
|
|
|
|
|
{
|
|
|
|
|
MessageQueueEntry entry = std::move(message_queue.front());
|
|
|
|
|
message_queue.pop();
|
|
|
|
|
const MessageQueueEntry &entry = message_queue.top();
|
|
|
|
|
|
|
|
|
|
Send(entry.router, entry.message);
|
|
|
|
|
message_queue.pop();
|
|
|
|
|
|
|
|
|
|
empty_count = 0;
|
|
|
|
|
sent_count++;
|
|
|
|
|
}
|
|
|
|
@ -395,8 +401,7 @@ namespace llarp
|
|
|
|
|
|
|
|
|
|
while(!movedMessages.empty())
|
|
|
|
|
{
|
|
|
|
|
MessageQueueEntry entry = std::move(movedMessages.front());
|
|
|
|
|
movedMessages.pop();
|
|
|
|
|
const MessageQueueEntry &entry = movedMessages.top();
|
|
|
|
|
|
|
|
|
|
if(status == SendStatus::Success)
|
|
|
|
|
{
|
|
|
|
@ -406,6 +411,7 @@ namespace llarp
|
|
|
|
|
{
|
|
|
|
|
DoCallback(entry.message.second, status);
|
|
|
|
|
}
|
|
|
|
|
movedMessages.pop();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|