1 #ifndef __AGGREGATORE_PULLSTREAMALIGNER__ 2 #define __AGGREGATORE_PULLSTREAMALIGNER__ 4 #include <aggregator/StreamAligner.hpp> 13 PullStreamBase() : has_data(
false) {}
14 virtual ~PullStreamBase() {}
15 virtual void pull() = 0;
16 virtual void push() = 0;
17 virtual void copyState(
const PullStreamBase& other ) = 0;
19 base::Time lastTime()
const {
return last_ts; }
20 bool hasData()
const {
return has_data; }
27 template <
class T>
class PullStream :
public PullStreamBase
30 typedef boost::function<bool (base::Time&, T&)> pull_callback_t;
32 PullStream( pull_callback_t pull_callback,
StreamAligner* sa,
size_t stream_index )
33 : stream_idx( stream_index ), sa( sa ), pull_callback( pull_callback ) {}
37 has_data = pull_callback( last_ts, last_data );
43 sa->
push( stream_idx, last_ts, last_data );
48 void copyState(
const PullStreamBase& other )
50 const PullStream<T> &pull_stream(
static_cast<const PullStream<T>&
>(other));
51 operator=( pull_stream );
58 pull_callback_t pull_callback;
62 static bool comparePullStreams(
const PullStreamBase* b1,
const PullStreamBase* b2 )
64 const base::Time &ts1( b1->lastTime() );
65 const base::Time &ts2( b2->lastTime() );
66 return ts1 < ts2 || !b2->hasData();
74 int idx = StreamAligner::registerStream<T>( callback, bufferSize, period, priority );
75 pull_streams.push_back(
new PullStream<T>( pull_callback,
this, idx ) );
83 if( !(*it)->hasData() )
89 if( first->hasData() )
pull_stream_vector pull_streams
Definition: PullStreamAligner.hpp:116
bool pull()
Definition: PullStreamAligner.hpp:79
void copyState(const PullStreamAligner &other)
Definition: PullStreamAligner.hpp:97
std::vector< PullStreamBase * > pull_stream_vector
Definition: PullStreamAligner.hpp:115
Definition: StreamAligner.hpp:18
~PullStreamAligner()
Definition: PullStreamAligner.hpp:108
int registerStream(typename PullStream< T >::pull_callback_t pull_callback, typename Stream< T >::callback_t callback, int bufferSize, base::Time period, int priority=-1)
Definition: PullStreamAligner.hpp:71
Definition: PullStreamAligner.hpp:8
void push(int idx, const base::Time &ts, const T &data)
Push new data into the stream.
Definition: StreamAligner.hpp:418
Definition: DetermineSampleTimestamp.hpp:8
Definition: StreamAligner.hpp:48
void copyState(const StreamAligner &other)
Definition: StreamAligner.hpp:252