Add a 'size' to command_queue.

This adds a size (command_queue_st::size) to the ThreadSafe deque<Command> that was
AICurlPrivate::command_queue that is kept equal to the number of 'add'
commands in the deque minus the number of 'remove' commands in the
deque. The deque itself is now command_queue_st::commands.
This commit is contained in:
Aleric Inglewood
2013-04-05 20:24:48 +02:00
parent 07201a5cfe
commit f58bc148ba

View File

@@ -264,9 +264,15 @@ void Command::reset(void)
// //
// If at this point addRequest is called again, then it is detected that the ThreadSafeBufferedCurlEasyRequest is active. // If at this point addRequest is called again, then it is detected that the ThreadSafeBufferedCurlEasyRequest is active.
struct command_queue_st {
std::deque<Command> commands; // The commands
size_t size; // Number of add commands in the queue minus the number of remove commands.
};
// Multi-threaded queue for passing Command objects from the main-thread to the curl-thread. // Multi-threaded queue for passing Command objects from the main-thread to the curl-thread.
AIThreadSafeSimpleDC<std::deque<Command> > command_queue; AIThreadSafeSimpleDC<command_queue_st> command_queue; // Fills 'size' with zero, because it's a global.
typedef AIAccess<std::deque<Command> > command_queue_wat; typedef AIAccess<command_queue_st> command_queue_wat;
typedef AIAccess<command_queue_st> command_queue_rat;
AIThreadSafeDC<Command> command_being_processed; AIThreadSafeDC<Command> command_being_processed;
typedef AIWriteAccess<Command> command_being_processed_wat; typedef AIWriteAccess<Command> command_being_processed_wat;
@@ -1289,7 +1295,7 @@ void AICurlThread::process_commands(AICurlMultiHandle_wat const& multi_handle_w)
// Access command_queue, and move command to command_being_processed. // Access command_queue, and move command to command_being_processed.
{ {
command_queue_wat command_queue_w(command_queue); command_queue_wat command_queue_w(command_queue);
if (command_queue_w->empty()) if (command_queue_w->commands.empty())
{ {
mWakeUpFlagMutex.lock(); mWakeUpFlagMutex.lock();
mWakeUpFlag = false; mWakeUpFlag = false;
@@ -1297,8 +1303,22 @@ void AICurlThread::process_commands(AICurlMultiHandle_wat const& multi_handle_w)
break; break;
} }
// Move the next command from the queue into command_being_processed. // Move the next command from the queue into command_being_processed.
*command_being_processed_wat(command_being_processed) = command_queue_w->front(); command_st command;
command_queue_w->pop_front(); {
command_being_processed_wat command_being_processed_w(command_being_processed);
*command_being_processed_w = command_queue_w->commands.front();
command = command_being_processed_w->command();
}
// Update the size: the number netto number of pending requests in the command queue.
command_queue_w->commands.pop_front();
if (command == cmd_add)
{
command_queue_w->size--;
}
else if (command == cmd_remove)
{
command_queue_w->size++;
}
} }
// Access command_being_processed only. // Access command_being_processed only.
{ {
@@ -1929,7 +1949,8 @@ void clearCommandQueue(void)
{ {
// Clear the command queue now in order to avoid the global deinitialization order fiasco. // Clear the command queue now in order to avoid the global deinitialization order fiasco.
command_queue_wat command_queue_w(command_queue); command_queue_wat command_queue_w(command_queue);
command_queue_w->clear(); command_queue_w->commands.clear();
command_queue_w->size = 0;
} }
//----------------------------------------------------------------------------- //-----------------------------------------------------------------------------
@@ -2331,7 +2352,7 @@ void AICurlEasyRequest::addRequest(void)
// Find the last command added. // Find the last command added.
command_st cmd = cmd_none; command_st cmd = cmd_none;
for (std::deque<Command>::iterator iter = command_queue_w->begin(); iter != command_queue_w->end(); ++iter) for (std::deque<Command>::iterator iter = command_queue_w->commands.begin(); iter != command_queue_w->commands.end(); ++iter)
{ {
if (*iter == *this) if (*iter == *this)
{ {
@@ -2357,7 +2378,8 @@ void AICurlEasyRequest::addRequest(void)
} }
#endif #endif
// Add a command to add the new request to the multi session to the command queue. // Add a command to add the new request to the multi session to the command queue.
command_queue_w->push_back(Command(*this, cmd_add)); command_queue_w->commands.push_back(Command(*this, cmd_add));
command_queue_w->size++;
AICurlEasyRequest_wat(*get())->add_queued(); AICurlEasyRequest_wat(*get())->add_queued();
} }
// Something was added to the queue, wake up the thread to get it. // Something was added to the queue, wake up the thread to get it.
@@ -2381,7 +2403,7 @@ void AICurlEasyRequest::removeRequest(void)
// Find the last command added. // Find the last command added.
command_st cmd = cmd_none; command_st cmd = cmd_none;
for (std::deque<Command>::iterator iter = command_queue_w->begin(); iter != command_queue_w->end(); ++iter) for (std::deque<Command>::iterator iter = command_queue_w->commands.begin(); iter != command_queue_w->commands.end(); ++iter)
{ {
if (*iter == *this) if (*iter == *this)
{ {
@@ -2418,7 +2440,8 @@ void AICurlEasyRequest::removeRequest(void)
} }
#endif #endif
// Add a command to remove this request from the multi session to the command queue. // Add a command to remove this request from the multi session to the command queue.
command_queue_w->push_back(Command(*this, cmd_remove)); command_queue_w->commands.push_back(Command(*this, cmd_remove));
command_queue_w->size--;
// Suppress warning that would otherwise happen if the callbacks are revoked before the curl thread removed the request. // Suppress warning that would otherwise happen if the callbacks are revoked before the curl thread removed the request.
AICurlEasyRequest_wat(*get())->remove_queued(); AICurlEasyRequest_wat(*get())->remove_queued();
} }