Skip to content

Commit

Permalink
Log and queue added
Browse files Browse the repository at this point in the history
  • Loading branch information
orignal committed Dec 10, 2013
1 parent d07f5d0 commit 465075d
Show file tree
Hide file tree
Showing 3 changed files with 149 additions and 0 deletions.
3 changes: 3 additions & 0 deletions Log.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
#include "Log.h"

i2p::util::MsgQueue<LogMsg> g_Log;
45 changes: 45 additions & 0 deletions Log.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
#ifndef LOG_H__
#define LOG_H__

#include <iostream>
#include <sstream>
#include "Queue.h"

struct LogMsg
{
std::stringstream s;
std::ostream& output;

LogMsg (std::ostream& o = std::cout): output (o) {};

void Process ()
{
output << s.str ();
}
};

extern i2p::util::MsgQueue<LogMsg> g_Log;

template<typename TValue>
void LogPrint (std::stringstream& s, TValue arg)
{
s << arg;
}

template<typename TValue, typename... TArgs>
void LogPrint (std::stringstream& s, TValue arg, TArgs... args)
{
LogPrint (s, arg);
LogPrint (s, args...);
}

template<typename... TArgs>
void LogPrint (TArgs... args)
{
LogMsg * msg = new LogMsg ();
LogPrint (msg->s, args...);
msg->s << std::endl;
g_Log.Put (msg);
}

#endif
101 changes: 101 additions & 0 deletions Queue.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
#ifndef QUEUE_H__
#define QUEUE_H__

#include <queue>
#include <mutex>
#include <thread>
#include <condition_variable>

namespace i2p
{
namespace util
{
template<typename Element>
class Queue
{
public:

void Put (Element * e)
{
std::unique_lock<std::mutex> l(m_QueueMutex);
m_Queue.push (e);
m_NonEmpty.notify_one ();
}

Element * GetNext ()
{
std::unique_lock<std::mutex> l(m_QueueMutex);
Element * el = GetNonThreadSafe ();
if (!el)
{
m_NonEmpty.wait (l);
el = GetNonThreadSafe ();
}
return el;
}

Element * GetNextWithTimeout (int usec)
{
std::unique_lock<std::mutex> l(m_QueueMutex);
Element * el = GetNonThreadSafe ();
if (!el)
{
m_NonEmpty.wait_for (l, std::chrono::milliseconds (usec));
el = GetNonThreadSafe ();
}
return el;
}

void WakeUp () { m_NonEmpty.notify_one (); };

Element * Get ()
{
std::unique_lock<std::mutex> l(m_QueueMutex);
return GetNonThreadSafe ();
}

private:

Element * GetNonThreadSafe ()
{
if (!m_Queue.empty ())
{
Element * el = m_Queue.front ();
m_Queue.pop ();
return el;
}
return nullptr;
}

private:

std::queue<Element *> m_Queue;
std::mutex m_QueueMutex;
std::condition_variable m_NonEmpty;
};

template<class Msg>
class MsgQueue: public Queue<Msg>
{
public:

MsgQueue (): m_Thread (std::bind (&MsgQueue<Msg>::Run, this)) {};

private:
void Run ()
{
Msg * msg = nullptr;
while ((msg = Queue<Msg>::GetNext ()) != nullptr)
{
msg->Process ();
delete msg;
}
}

private:
std::thread m_Thread;
};
}
}

#endif

0 comments on commit 465075d

Please sign in to comment.