AbstractCodec use fair_queue

This commit is contained in:
Michael Davidsaver
2015-12-14 16:59:55 -05:00
parent db86e47659
commit 730d30fe54
4 changed files with 120 additions and 190 deletions
+16 -17
View File
@@ -53,7 +53,7 @@ namespace epics {
//PROTECTED
_readMode(NORMAL), _version(0), _flags(0), _command(0), _payloadSize(0),
_remoteTransportSocketReceiveBufferSize(MAX_TCP_RECV), _totalBytesSent(0),
_blockingProcessQueue(false), _senderThread(0),
_senderThread(0),
_writeMode(PROCESS_SEND_QUEUE),
_writeOpReady(false),_lowLatency(false),
_socketBuffer(receiveBuffer),
@@ -98,7 +98,6 @@ namespace epics {
_maxSendPayloadSize =
_sendBuffer->getSize() - 2*PVA_MESSAGE_HEADER_SIZE;
_socketSendBufferSize = socketSendBufferSize;
_blockingProcessQueue = blockingProcessQueue;
}
@@ -851,7 +850,8 @@ namespace epics {
std::size_t senderProcessed = 0;
while (senderProcessed++ < MAX_MESSAGE_SEND)
{
TransportSender::shared_pointer sender = _sendQueue.take(-1);
TransportSender::shared_pointer sender;
_sendQueue.pop_front_try(sender);
if (sender.get() == 0)
{
// flush
@@ -860,19 +860,20 @@ namespace epics {
sendCompleted(); // do not schedule sending
if (_blockingProcessQueue) {
if (terminated()) // termination
if (terminated()) // termination
break;
sender = _sendQueue.take(0);
// termination (we want to process even if shutdown)
if (sender.get() == 0)
break;
}
else
return;
// termination (we want to process even if shutdown)
_sendQueue.pop_front(sender);
}
processSender(sender);
try{
processSender(sender);
}catch(...){
if (_sendBuffer->getPosition() > 0)
flush(true);
sendCompleted();
throw;
}
}
}
@@ -884,13 +885,13 @@ namespace epics {
void AbstractCodec::clearSendQueue()
{
_sendQueue.clean();
_sendQueue.clear();
}
void AbstractCodec::enqueueSendRequest(
TransportSender::shared_pointer const & sender) {
_sendQueue.put(sender);
_sendQueue.push_back(sender);
scheduleSend();
}
@@ -1066,8 +1067,6 @@ namespace epics {
// this is important to avoid cyclic refs (memory leak)
clearSendQueue();
_sendQueue.wakeup();
// post close
internalPostClose(true);
}
+5 -116
View File
@@ -112,120 +112,6 @@ namespace epics {
#endif
// TODO replace this queue with lock-free implementation
template<typename T>
class queue {
public:
queue(void) { }
//TODO
/*queue(queue const &T) = delete;
queue(queue &&T) = delete;
queue& operator=(const queue &T) = delete;
*/
~queue(void)
{
}
bool empty(void)
{
epics::pvData::Lock lock(_queueMutex);
return _queue.empty();
}
void clean()
{
epics::pvData::Lock lock(_queueMutex);
_queue.clear();
}
void wakeup()
{
if (!_wakeup.getAndSet(true))
{
_queueEvent.signal();
}
}
void put(T const & elem)
{
{
epics::pvData::Lock lock(_queueMutex);
_queue.push_back(elem);
}
_queueEvent.signal();
}
// TODO very sub-optimal (locks and empty() - pop() sequence; at least 2 locks!)
T take(int timeOut)
{
while (true)
{
bool isEmpty = empty();
if (isEmpty)
{
if (timeOut < 0) {
return T();
}
while (isEmpty)
{
if (timeOut == 0) {
_queueEvent.wait();
}
else {
_queueEvent.wait(timeOut);
}
isEmpty = empty();
if (isEmpty)
{
if (timeOut > 0) { // TODO spurious wakeup, but not critical
return T();
}
else // if (timeout == 0) cannot be negative
{
if (_wakeup.getAndSet(false)) {
return T();
}
}
}
}
}
else
{
epics::pvData::Lock lock(_queueMutex);
if (_queue.empty())
return T();
T sender = _queue.front();
_queue.pop_front();
return sender;
}
}
}
size_t size() {
epics::pvData::Lock lock(_queueMutex);
return _queue.size();
}
private:
std::deque<T> _queue;
epics::pvData::Event _queueEvent;
epics::pvData::Mutex _queueMutex;
AtomicValue<bool> _wakeup;
};
class epicsShareClass io_exception: public std::runtime_error {
public:
@@ -327,6 +213,10 @@ namespace epics {
char* /*deserializeTo*/,
std::size_t /*elementCount*/, std::size_t /*elementSize*/);
bool sendQueueEmpty() const {
return _sendQueue.empty();
}
protected:
virtual void sendBufferFull(int tries) = 0;
@@ -341,7 +231,6 @@ namespace epics {
int32_t _payloadSize; // TODO why not size_t?
epics::pvData::int32 _remoteTransportSocketReceiveBufferSize;
int64_t _totalBytesSent;
bool _blockingProcessQueue;
//TODO initialize union
osiSockAddr _sendTo;
epicsThreadId _senderThread;
@@ -352,7 +241,7 @@ namespace epics {
std::tr1::shared_ptr<epics::pvData::ByteBuffer> _socketBuffer;
std::tr1::shared_ptr<epics::pvData::ByteBuffer> _sendBuffer;
queue<TransportSender::shared_pointer> _sendQueue;
fair_queue<TransportSender> _sendQueue;
private:
+2 -1
View File
@@ -24,6 +24,7 @@
#include <pv/timer.h>
#include <pv/pvData.h>
#include <pv/sharedPtr.h>
#include <pv/fairQueue.h>
#ifdef remoteEpicsExportSharedSymbols
# define epicsExportSharedSymbols
@@ -142,7 +143,7 @@ namespace epics {
/**
* Interface defining transport sender (instance sending data over transport).
*/
class TransportSender : public Lockable {
class TransportSender : public Lockable, public fair_queue<TransportSender>::entry {
public:
POINTER_DEFINITIONS(TransportSender);