mirror of
https://github.com/pocoproject/poco.git
synced 2025-10-26 10:32:56 +01:00
#2787: add queueSize property to the AsyncChannel
This commit is contained in:
@@ -50,14 +50,14 @@ public:
|
|||||||
void setChannel(Channel::Ptr pChannel);
|
void setChannel(Channel::Ptr pChannel);
|
||||||
/// Connects the AsyncChannel to the given target channel.
|
/// Connects the AsyncChannel to the given target channel.
|
||||||
/// All messages will be forwarded to this channel.
|
/// All messages will be forwarded to this channel.
|
||||||
|
|
||||||
Channel::Ptr getChannel() const;
|
Channel::Ptr getChannel() const;
|
||||||
/// Returns the target channel.
|
/// Returns the target channel.
|
||||||
|
|
||||||
void open();
|
void open();
|
||||||
/// Opens the channel and creates the
|
/// Opens the channel and creates the
|
||||||
/// background logging thread.
|
/// background logging thread.
|
||||||
|
|
||||||
void close();
|
void close();
|
||||||
/// Closes the channel and stops the background
|
/// Closes the channel and stops the background
|
||||||
/// logging thread.
|
/// logging thread.
|
||||||
@@ -69,7 +69,7 @@ public:
|
|||||||
void setProperty(const std::string& name, const std::string& value);
|
void setProperty(const std::string& name, const std::string& value);
|
||||||
/// Sets or changes a configuration property.
|
/// Sets or changes a configuration property.
|
||||||
///
|
///
|
||||||
/// The "channel" property allows setting the target
|
/// The "channel" property allows setting the target
|
||||||
/// channel via the LoggingRegistry.
|
/// channel via the LoggingRegistry.
|
||||||
/// The "channel" property is set-only.
|
/// The "channel" property is set-only.
|
||||||
///
|
///
|
||||||
@@ -82,18 +82,34 @@ public:
|
|||||||
/// * highest
|
/// * highest
|
||||||
///
|
///
|
||||||
/// The "priority" property is set-only.
|
/// The "priority" property is set-only.
|
||||||
|
///
|
||||||
|
/// The "queueSize" property allows to limit the number
|
||||||
|
/// of messages in the queue. If the queue is full and
|
||||||
|
/// new messages are logged, these are dropped until the
|
||||||
|
/// queue has free capacity again. The number of dropped
|
||||||
|
/// messages is recorded, and a log message indicating the
|
||||||
|
/// number of dropped messages will be generated when the
|
||||||
|
/// queue has free capacity again.
|
||||||
|
/// In addition to an unsigned integer specifying the size,
|
||||||
|
/// this property can have the values "none" or "unlimited",
|
||||||
|
/// which disable the queue size limit. A size of 0 also
|
||||||
|
/// removes the limit.
|
||||||
|
///
|
||||||
|
/// The "queueSize" property is set-only.
|
||||||
|
|
||||||
protected:
|
protected:
|
||||||
~AsyncChannel();
|
~AsyncChannel();
|
||||||
void run();
|
void run();
|
||||||
void setPriority(const std::string& value);
|
void setPriority(const std::string& value);
|
||||||
|
|
||||||
private:
|
private:
|
||||||
Channel::Ptr _pChannel;
|
Channel::Ptr _pChannel;
|
||||||
Thread _thread;
|
Thread _thread;
|
||||||
FastMutex _threadMutex;
|
FastMutex _threadMutex;
|
||||||
FastMutex _channelMutex;
|
FastMutex _channelMutex;
|
||||||
NotificationQueue _queue;
|
NotificationQueue _queue;
|
||||||
|
std::size_t _queueSize = 0;
|
||||||
|
std::size_t _dropCount = 0;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -18,7 +18,10 @@
|
|||||||
#include "Poco/Formatter.h"
|
#include "Poco/Formatter.h"
|
||||||
#include "Poco/AutoPtr.h"
|
#include "Poco/AutoPtr.h"
|
||||||
#include "Poco/LoggingRegistry.h"
|
#include "Poco/LoggingRegistry.h"
|
||||||
|
#include "Poco/NumberParser.h"
|
||||||
#include "Poco/Exception.h"
|
#include "Poco/Exception.h"
|
||||||
|
#include "Poco/String.h"
|
||||||
|
#include "Poco/Format.h"
|
||||||
|
|
||||||
|
|
||||||
namespace Poco {
|
namespace Poco {
|
||||||
@@ -31,23 +34,23 @@ public:
|
|||||||
_msg(msg)
|
_msg(msg)
|
||||||
{
|
{
|
||||||
}
|
}
|
||||||
|
|
||||||
~MessageNotification()
|
~MessageNotification()
|
||||||
{
|
{
|
||||||
}
|
}
|
||||||
|
|
||||||
const Message& message() const
|
const Message& message() const
|
||||||
{
|
{
|
||||||
return _msg;
|
return _msg;
|
||||||
}
|
}
|
||||||
|
|
||||||
private:
|
private:
|
||||||
Message _msg;
|
Message _msg;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
||||||
AsyncChannel::AsyncChannel(Channel::Ptr pChannel, Thread::Priority prio):
|
AsyncChannel::AsyncChannel(Channel::Ptr pChannel, Thread::Priority prio):
|
||||||
_pChannel(pChannel),
|
_pChannel(pChannel),
|
||||||
_thread("AsyncChannel")
|
_thread("AsyncChannel")
|
||||||
{
|
{
|
||||||
_thread.setPriority(prio);
|
_thread.setPriority(prio);
|
||||||
@@ -70,7 +73,7 @@ AsyncChannel::~AsyncChannel()
|
|||||||
void AsyncChannel::setChannel(Channel::Ptr pChannel)
|
void AsyncChannel::setChannel(Channel::Ptr pChannel)
|
||||||
{
|
{
|
||||||
FastMutex::ScopedLock lock(_channelMutex);
|
FastMutex::ScopedLock lock(_channelMutex);
|
||||||
|
|
||||||
_pChannel = pChannel;
|
_pChannel = pChannel;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -94,10 +97,10 @@ void AsyncChannel::close()
|
|||||||
if (_thread.isRunning())
|
if (_thread.isRunning())
|
||||||
{
|
{
|
||||||
while (!_queue.empty()) Thread::sleep(100);
|
while (!_queue.empty()) Thread::sleep(100);
|
||||||
|
|
||||||
do
|
do
|
||||||
{
|
{
|
||||||
_queue.wakeUpAll();
|
_queue.wakeUpAll();
|
||||||
}
|
}
|
||||||
while (!_thread.tryJoin(100));
|
while (!_thread.tryJoin(100));
|
||||||
}
|
}
|
||||||
@@ -106,6 +109,18 @@ void AsyncChannel::close()
|
|||||||
|
|
||||||
void AsyncChannel::log(const Message& msg)
|
void AsyncChannel::log(const Message& msg)
|
||||||
{
|
{
|
||||||
|
if (_queueSize != 0 && _queue.size() >= _queueSize)
|
||||||
|
{
|
||||||
|
++_dropCount;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (_dropCount != 0)
|
||||||
|
{
|
||||||
|
_queue.enqueueNotification(new MessageNotification(Message(msg, Poco::format("Dropped %z messages.", _dropCount))));
|
||||||
|
_dropCount = 0;
|
||||||
|
}
|
||||||
|
|
||||||
open();
|
open();
|
||||||
|
|
||||||
_queue.enqueueNotification(new MessageNotification(msg));
|
_queue.enqueueNotification(new MessageNotification(msg));
|
||||||
@@ -115,11 +130,24 @@ void AsyncChannel::log(const Message& msg)
|
|||||||
void AsyncChannel::setProperty(const std::string& name, const std::string& value)
|
void AsyncChannel::setProperty(const std::string& name, const std::string& value)
|
||||||
{
|
{
|
||||||
if (name == "channel")
|
if (name == "channel")
|
||||||
|
{
|
||||||
setChannel(LoggingRegistry::defaultRegistry().channelForName(value));
|
setChannel(LoggingRegistry::defaultRegistry().channelForName(value));
|
||||||
|
}
|
||||||
else if (name == "priority")
|
else if (name == "priority")
|
||||||
|
{
|
||||||
setPriority(value);
|
setPriority(value);
|
||||||
|
}
|
||||||
|
else if (name == "queueSize")
|
||||||
|
{
|
||||||
|
if (Poco::icompare(value, "none") == 0 || Poco::icompare(value, "unlimited") == 0 || value.empty())
|
||||||
|
_queueSize = 0;
|
||||||
|
else
|
||||||
|
_queueSize = Poco::NumberParser::parseUnsigned(value);
|
||||||
|
}
|
||||||
else
|
else
|
||||||
|
{
|
||||||
Channel::setProperty(name, value);
|
Channel::setProperty(name, value);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -137,12 +165,12 @@ void AsyncChannel::run()
|
|||||||
nf = _queue.waitDequeueNotification();
|
nf = _queue.waitDequeueNotification();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
void AsyncChannel::setPriority(const std::string& value)
|
void AsyncChannel::setPriority(const std::string& value)
|
||||||
{
|
{
|
||||||
Thread::Priority prio = Thread::PRIO_NORMAL;
|
Thread::Priority prio = Thread::PRIO_NORMAL;
|
||||||
|
|
||||||
if (value == "lowest")
|
if (value == "lowest")
|
||||||
prio = Thread::PRIO_LOWEST;
|
prio = Thread::PRIO_LOWEST;
|
||||||
else if (value == "low")
|
else if (value == "low")
|
||||||
@@ -155,7 +183,7 @@ void AsyncChannel::setPriority(const std::string& value)
|
|||||||
prio = Thread::PRIO_HIGHEST;
|
prio = Thread::PRIO_HIGHEST;
|
||||||
else
|
else
|
||||||
throw InvalidArgumentException("thread priority", value);
|
throw InvalidArgumentException("thread priority", value);
|
||||||
|
|
||||||
_thread.setPriority(prio);
|
_thread.setPriority(prio);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user