Fixed 1 second delay for child processes.

This commit is contained in:
Gunnar Beutner 2013-02-10 01:35:40 +01:00
parent ee37e0cace
commit fc6df0ecbd
2 changed files with 89 additions and 66 deletions

View File

@ -24,7 +24,7 @@ using namespace icinga;
bool Process::m_ThreadCreated = false; bool Process::m_ThreadCreated = false;
boost::mutex Process::m_Mutex; boost::mutex Process::m_Mutex;
deque<Process::Ptr> Process::m_Tasks; deque<Process::Ptr> Process::m_Tasks;
condition_variable Process::m_TasksCV; int Process::m_TaskFd;
Process::Process(const String& command, const Dictionary::Ptr& environment) Process::Process(const String& command, const Dictionary::Ptr& environment)
: AsyncTask<Process, ProcessResult>(), m_Command(command), : AsyncTask<Process, ProcessResult>(), m_Command(command),
@ -33,7 +33,14 @@ Process::Process(const String& command, const Dictionary::Ptr& environment)
assert(Application::IsMainThread()); assert(Application::IsMainThread());
if (!m_ThreadCreated) { if (!m_ThreadCreated) {
thread t(&Process::WorkerThreadProc); int fds[2];
if (pipe(fds) < 0)
BOOST_THROW_EXCEPTION(PosixException("pipe() failed.", errno));
m_TaskFd = fds[1];
thread t(&Process::WorkerThreadProc, fds[0]);
t.detach(); t.detach();
m_ThreadCreated = true; m_ThreadCreated = true;
@ -42,87 +49,103 @@ Process::Process(const String& command, const Dictionary::Ptr& environment)
void Process::Run(void) void Process::Run(void)
{ {
boost::mutex::scoped_lock lock(m_Mutex); {
m_Tasks.push_back(GetSelf()); boost::mutex::scoped_lock lock(m_Mutex);
m_TasksCV.notify_all(); m_Tasks.push_back(GetSelf());
}
/**
* This little gem which is commonly known as the "self-pipe trick"
* takes care of waking up the select() call in the worker thread.
*/
if (write(m_TaskFd, "T", 1) < 0)
BOOST_THROW_EXCEPTION(PosixException("write() failed.", errno));
} }
void Process::WorkerThreadProc(void) void Process::WorkerThreadProc(int taskFd)
{ {
boost::mutex::scoped_lock lock(m_Mutex);
map<int, Process::Ptr> tasks; map<int, Process::Ptr> tasks;
for (;;) { for (;;) {
while (m_Tasks.empty() || tasks.size() >= MaxTasksPerThread) { map<int, Process::Ptr>::iterator it, prev;
lock.unlock();
map<int, Process::Ptr>::iterator it, prev;
#ifndef _MSC_VER #ifndef _MSC_VER
fd_set readfds; fd_set readfds;
int nfds = 0; int nfds = 0;
FD_ZERO(&readfds); FD_ZERO(&readfds);
FD_SET(taskFd, &readfds);
int fd; if (taskFd > nfds)
BOOST_FOREACH(tie(fd, tuples::ignore), tasks) { nfds = taskFd;
if (fd > nfds)
nfds = fd;
FD_SET(fd, &readfds); int fd;
} BOOST_FOREACH(tie(fd, tuples::ignore), tasks) {
if (fd > nfds)
nfds = fd;
timeval tv; FD_SET(fd, &readfds);
tv.tv_sec = 1;
tv.tv_usec = 0;
select(nfds + 1, &readfds, NULL, NULL, &tv);
#else /* _MSC_VER */
Utility::Sleep(1);
#endif /* _MSC_VER */
for (it = tasks.begin(); it != tasks.end(); ) {
int fd = it->first;
Process::Ptr task = it->second;
#ifndef _MSC_VER
if (!FD_ISSET(fd, &readfds)) {
it++;
continue;
}
#endif /* _MSC_VER */
if (!task->RunTask()) {
prev = it;
it++;
tasks.erase(prev);
Event::Post(boost::bind(&Process::FinishResult, task, task->m_Result));
} else {
it++;
}
}
lock.lock();
} }
while (!m_Tasks.empty() && tasks.size() < MaxTasksPerThread) { timeval tv;
Process::Ptr task = m_Tasks.front(); tv.tv_sec = 1;
m_Tasks.pop_front(); tv.tv_usec = 0;
select(nfds + 1, &readfds, NULL, NULL, &tv);
#else /* _MSC_VER */
Utility::Sleep(1);
#endif /* _MSC_VER */
lock.unlock(); if (FD_ISSET(taskFd, &readfds)) {
/* clear pipe */
char buffer[512];
int rc = read(taskFd, buffer, sizeof(buffer));
assert(rc >= 1);
try { while (tasks.size() < MaxTasksPerThread) {
task->InitTask(); Process::Ptr task;
int fd = task->GetFD(); {
if (fd >= 0) boost::mutex::scoped_lock lock(m_Mutex);
tasks[fd] = task;
} catch (...) { if (m_Tasks.empty())
Event::Post(boost::bind(&Process::FinishException, task, boost::current_exception())); break;
task = m_Tasks.front();
m_Tasks.pop_front();
}
try {
task->InitTask();
int fd = task->GetFD();
if (fd >= 0)
tasks[fd] = task;
} catch (...) {
Event::Post(boost::bind(&Process::FinishException, task, boost::current_exception()));
}
} }
}
lock.lock(); for (it = tasks.begin(); it != tasks.end(); ) {
int fd = it->first;
Process::Ptr task = it->second;
#ifndef _MSC_VER
if (!FD_ISSET(fd, &readfds)) {
it++;
continue;
}
#endif /* _MSC_VER */
if (!task->RunTask()) {
prev = it;
it++;
tasks.erase(prev);
Event::Post(boost::bind(&Process::FinishResult, task, task->m_Result));
} else {
it++;
}
} }
} }
} }

View File

@ -71,9 +71,9 @@ private:
static boost::mutex m_Mutex; static boost::mutex m_Mutex;
static deque<Process::Ptr> m_Tasks; static deque<Process::Ptr> m_Tasks;
static condition_variable m_TasksCV; static int m_TaskFd;
static void WorkerThreadProc(void); static void WorkerThreadProc(int taskFd);
void InitTask(void); void InitTask(void);
bool RunTask(void); bool RunTask(void);