cwSpscQueueTmpl.h/cpp : Initial commit.
This commit is contained in:
parent
c48ba08514
commit
b577d476d3
@ -34,6 +34,7 @@ namespace cw
|
||||
unsigned iter; // execution counter
|
||||
unsigned msgN; // count of msg's processed
|
||||
shared_t* share; // shared variables
|
||||
uint8_t prevId;
|
||||
} ctx_t;
|
||||
|
||||
|
||||
@ -57,7 +58,7 @@ namespace cw
|
||||
m->checksum += m->data[i];
|
||||
}
|
||||
|
||||
c->share->q->push(m);
|
||||
c->share->q->publish();
|
||||
|
||||
c->msgN++;
|
||||
|
||||
@ -84,7 +85,16 @@ namespace cw
|
||||
|
||||
if( curCheckSum != m->checksum )
|
||||
cwLogError(kOpFailRC,"Checksum mismatch.0x%x != 0x%x ",curCheckSum,m->checksum);
|
||||
else
|
||||
{
|
||||
uint8_t id = c->prevId + 1;
|
||||
if( id != m->data[0] )
|
||||
cwLogInfo("drop ");
|
||||
|
||||
c->prevId = m->data[0];
|
||||
}
|
||||
|
||||
c->msgN++;
|
||||
}
|
||||
|
||||
c->iter++;
|
||||
|
@ -10,31 +10,27 @@ namespace cw
|
||||
public:
|
||||
spScQueueTmpl( unsigned eleN )
|
||||
{
|
||||
_aV = mem::allocZ<T*>(eleN);
|
||||
_mem = mem::allocZ<T>(eleN);
|
||||
_wi.load(std::memory_order_release);
|
||||
_ri.load(std::memory_order_release);
|
||||
|
||||
for(unsigned i=0; i<eleN; ++i)
|
||||
_aV[i] = _mem + i;
|
||||
_aN = eleN;
|
||||
_aV = mem::allocZ<T>(eleN);
|
||||
_wi.store(0,std::memory_order_release);
|
||||
_ri.store(0,std::memory_order_release);
|
||||
}
|
||||
|
||||
virtual ~spScQueueTmpl()
|
||||
{
|
||||
mem::release(_aV);
|
||||
mem::release(_mem);
|
||||
}
|
||||
|
||||
T* get()
|
||||
{
|
||||
unsigned wi = _wi.load( std::memory_order_relaxed );
|
||||
return _aV[wi];
|
||||
return _aV + wi;
|
||||
}
|
||||
|
||||
rc_t push( T* v )
|
||||
rc_t publish()
|
||||
{
|
||||
unsigned ri = _ri.load( std::memory_order_acquire );
|
||||
unsigned wi = _wi.load( std::memory_order_relaxed );
|
||||
unsigned wi = _wi.load( std::memory_order_relaxed ); // _wi is written in this thread
|
||||
|
||||
// calc. the count of full elements
|
||||
unsigned n = wi >= ri ? wi-ri : (_aN-ri) + wi;
|
||||
@ -44,11 +40,9 @@ namespace cw
|
||||
if( n >= _aN-1 )
|
||||
return kBufTooSmallRC;
|
||||
|
||||
// store the new element
|
||||
_aV[wi] = v;
|
||||
|
||||
wi = (wi+1) % _aN;
|
||||
|
||||
// advance the write position
|
||||
_wi.store( wi, std::memory_order_release );
|
||||
|
||||
return kOkRC;
|
||||
@ -64,7 +58,7 @@ namespace cw
|
||||
if( n == 0 )
|
||||
return nullptr;
|
||||
|
||||
T* v = _aV[ri];
|
||||
T* v = _aV + ri;
|
||||
|
||||
ri = (ri+1) % _aN;
|
||||
|
||||
@ -78,8 +72,7 @@ namespace cw
|
||||
|
||||
private:
|
||||
unsigned _aN = 0;
|
||||
T** _aV = nullptr;
|
||||
T* _mem;
|
||||
T* _aV = nullptr;
|
||||
std::atomic<unsigned> _wi;
|
||||
std::atomic<unsigned> _ri;
|
||||
|
||||
|
Loading…
Reference in New Issue
Block a user