1 #ifndef __AGGREGATOR_HPP__ 2 #define __AGGREGATOR_HPP__ 4 #include <base/Time.hpp> 6 #include <base-logging/Logging.hpp> 8 #include <base/CircularBuffer.hpp> 10 #include <boost/function.hpp> 11 #include <boost/tuple/tuple.hpp> 14 #include <aggregator/StreamAlignerStatus.hpp> 24 StreamBase() : active(
true ) {}
25 virtual ~StreamBase() {}
26 virtual base::Time pop() = 0;
27 virtual bool hasData()
const = 0;
28 virtual int getPriority()
const = 0;
29 virtual base::Time latestTimeStamp()
const = 0;
30 virtual base::Time latestDataTime()
const = 0;
31 virtual base::Time earliestDataTime()
const = 0;
33 virtual void copyState(
const StreamBase& other ) = 0;
34 virtual void clear() = 0;
36 bool isActive()
const {
return active; }
37 void setActive(
bool active ) { this->active = active; }
39 friend std::ostream &
operator<<(std::ostream &stream,
const aggregator::StreamAligner::StreamBase &base);
48 template <
class T>
class Stream :
public StreamBase
51 typedef boost::function<void (const base::Time &ts, const T &value)>
callback_t;
54 typedef std::pair<base::Time,T>
item;
63 Stream( callback_t callback,
size_t bufferSize, base::Time period,
int priority,
const std::string &name )
64 : bufferSize( bufferSize ), callback(callback), period(period), lastTime(base::Time::fromSeconds(0)), priority(priority)
67 status.priority = priority;
70 buffer.set_capacity( bufferSize );
74 buffer.set_capacity( 20 );
76 status.buffer_size = buffer.capacity();
86 sample = buffer.front();
97 status.buffer_fill = buffer.size();
98 status.latest_data_time = latestDataTime();
99 status.earliest_data_time = earliestDataTime();
100 status.active = isActive();
114 void push(
const base::Time &ts,
const T &data )
118 status.samples_backward_in_time++;
130 status.samples_dropped_buffer_full++;
134 buffer.set_capacity(buffer.capacity() * 2);
135 status.buffer_size = buffer.capacity();
138 buffer.push_back( std::make_pair(ts, data) );
148 status.samples_processed++;
149 base::Time ts = buffer.front().first;
151 callback( ts, buffer.front().second );
156 throw std::runtime_error(
"pop() called on stream with no data.");
160 {
return !buffer.empty(); }
165 return buffer.front().first;
167 return lastTime + period;
178 return buffer.front().first;
184 lastTime = base::Time();
187 status.latest_sample_time = base::Time();
188 status.latest_data_time = base::Time();
189 status.samples_dropped_buffer_full = 0;
204 const base::Time &ts1( b1->latestTimeStamp() );
205 const base::Time &ts2( b2->latestTimeStamp() );
209 if(b1->hasData() && !b2->hasData())
212 if(!b1->hasData() && b2->hasData())
215 return b1->getPriority() < b2->getPriority();
240 : timeout(timeout), buffer_size_factor(2.0) {}
244 for(stream_vector::iterator it=streams.begin();it != streams.end();it++)
257 assert( streams.size() == other.
streams.size() );
258 for(
size_t i=0;i<streams.size();i++)
260 bool weGotStream = streams[i];
261 bool otherGotStream = other.
streams[i];
262 if(weGotStream != otherGotStream)
263 throw std::runtime_error(
"Stream setup of second stream aligner differs");
267 streams[i]->copyState( *other.
streams[i] );
296 throw std::runtime_error(
"invalid stream index.");
298 streams[idx]->setActive(
false );
310 throw std::runtime_error(
"invalid stream index.");
312 streams[idx]->setActive(
true );
323 throw std::runtime_error(
"invalid stream index.");
325 return streams[idx]->isActive();
338 throw std::runtime_error(
"invalid stream index.");
345 status.
streams[idx].active =
false;
371 if( period == base::Time() )
373 throw std::runtime_error(
"No buffer size provided for stream with unknown period.");
375 else if( period < base::Time() )
378 bufferSize = buffer_size_factor * std::ceil( timeout.toSeconds() / -period.toSeconds() );
379 period = base::Time();
383 bufferSize = buffer_size_factor * std::ceil( timeout.toSeconds() / period.toSeconds() );
387 if( bufferSize == 0 )
389 LOG_DEBUG_S <<
"dynamically allocating stream aligner buffer for stream: " << name;
392 StreamBase *newStream =
new Stream<T>(callback, bufferSize, period, priority, name);
395 for(
size_t i = 0; i < streams.size(); i++)
399 streams[i] = newStream;
405 streams.push_back( newStream );
407 return streams.size() - 1;
418 template <
class T>
void push(
int idx,
const base::Time &ts,
const T& data )
420 if( !streams.at(idx) )
421 throw std::runtime_error(
"invalid stream index.");
426 stream->status.samples_received++;
427 stream->status.latest_sample_time = ts;
432 stream->setActive(
true );
439 stream->status.samples_dropped_late_arriving++;
446 stream->
push( ts, data );
449 template <
class T>
bool getNextSample(
int idx, std::pair<base::Time,T> &sample)
const 451 if( !streams.at(idx) )
452 throw std::runtime_error(
"invalid stream index.");
478 if( streams.empty() )
485 for(stream_vector::iterator it=items.begin();it != items.end();it++)
491 if( (*it)->hasData() )
494 current_ts = (*it)->pop();
497 else if( (*it)->isActive() )
500 base::Time latestDataTime;
501 base::Time firstDataTime;
504 if(current_ts == base::Time())
507 for(stream_vector::iterator it2=items.begin();it2 != items.end();it2++)
510 if(*it2 && (*it2)->hasData())
512 if(latestDataTime < (*it2)->latestDataTime())
513 latestDataTime = (*it2)->latestDataTime();
515 if(firstDataTime == base::Time() || firstDataTime > (*it2)->earliestDataTime())
516 firstDataTime = (*it2)->earliestDataTime();
524 if(latestDataTime - firstDataTime < timeout)
542 for(
size_t i = 0; i < streams.size(); i++)
550 latest_ts = base::Time();
551 current_ts = base::Time();
586 if( !streams.at(idx) )
587 throw std::runtime_error(
"invalid stream index.");
589 return streams[idx]->getBufferStatus();
597 status.
time = base::Time::now();
601 for(
size_t i=0;i<streams.size();i++)
604 status.
streams[i] = streams[i]->getBufferStatus();
610 friend std::ostream &
operator<<(std::ostream &stream,
const aggregator::StreamAligner::StreamBase &base);
618 for(
size_t i=0;i<re.
streams.size();i++)
620 stream << i <<
":" << *re.
streams[i] << std::endl;
626 inline std::ostream &
operator<<(std::ostream &stream,
const aggregator::StreamAligner::StreamBase &base)
629 stream << base.getBufferStatus();
virtual void clear()
Definition: StreamAligner.hpp:182
bool step()
Definition: StreamAligner.hpp:476
std::vector< StreamBase * > stream_vector
Definition: StreamAligner.hpp:221
base::Time period
Definition: StreamAligner.hpp:58
void setTimeout(const base::Time &t)
Definition: StreamAligner.hpp:276
virtual void copyState(const StreamBase &other)
Definition: StreamAligner.hpp:104
virtual ~Stream()
Definition: StreamAligner.hpp:79
int registerStream(typename Stream< T >::callback_t callback, int bufferSize, base::Time period, int priority=-1, const std::string &name=std::string())
Definition: StreamAligner.hpp:367
bool getNextSample(item &sample) const
Definition: StreamAligner.hpp:81
bool isStreamActive(int idx) const
Definition: StreamAligner.hpp:320
virtual int getPriority() const
Definition: StreamAligner.hpp:90
base::Time latest_time
Definition: StreamAlignerStatus.hpp:101
base::Time getLatency() const
Definition: StreamAligner.hpp:567
virtual ~StreamAligner()
Definition: StreamAligner.hpp:242
base::Time time
Definition: StreamAlignerStatus.hpp:88
base::Time lastTime
Definition: StreamAligner.hpp:59
base::Time current_ts
Definition: StreamAligner.hpp:229
bool getNextSample(int idx, std::pair< base::Time, T > &sample) const
Definition: StreamAligner.hpp:449
const StreamAlignerStatus & getStatus() const
Definition: StreamAligner.hpp:595
double buffer_size_factor
Definition: StreamAligner.hpp:231
virtual const StreamStatus & getBufferStatus() const
Definition: StreamAligner.hpp:95
Definition: StreamAligner.hpp:18
base::Time getCurrentTime() const
Definition: StreamAligner.hpp:571
base::Time latest_ts
Definition: StreamAligner.hpp:226
base::Time timeout
Definition: StreamAligner.hpp:223
void disableStream(int idx)
Definition: StreamAligner.hpp:293
base::Time getTimeOut() const
Definition: StreamAligner.hpp:562
void clear()
Definition: StreamAligner.hpp:540
callback_t callback
Definition: StreamAligner.hpp:57
boost::circular_buffer< item > buffer
Definition: StreamAligner.hpp:55
std::string name
Definition: StreamAlignerStatus.hpp:92
void push(int idx, const base::Time &ts, const T &data)
Push new data into the stream.
Definition: StreamAligner.hpp:418
StreamAlignerStatus status
Definition: StreamAligner.hpp:236
void enableStream(int idx)
Definition: StreamAligner.hpp:307
size_t bufferSize
Definition: StreamAligner.hpp:56
friend std::ostream & operator<<(std::ostream &stream, const aggregator::StreamAligner::StreamBase &base)
Definition: StreamAligner.hpp:626
virtual base::Time earliestDataTime() const
Definition: StreamAligner.hpp:175
bool hasData() const
Definition: StreamAligner.hpp:159
Definition: StreamAlignerStatus.hpp:11
int getStreamSize() const
Definition: StreamAligner.hpp:579
static bool compareStreams(const StreamBase *b1, const StreamBase *b2)
Definition: StreamAligner.hpp:196
base::Time current_time
Definition: StreamAlignerStatus.hpp:98
std::vector< StreamStatus > streams
Definition: StreamAlignerStatus.hpp:111
boost::function< void(const base::Time &ts, const T &value)> callback_t
Definition: StreamAligner.hpp:51
StreamAligner(base::Time timeout=base::Time::fromSeconds(1))
Definition: StreamAligner.hpp:239
Definition: DetermineSampleTimestamp.hpp:8
std::pair< base::Time, T > item
Definition: StreamAligner.hpp:54
base::Time pop()
Definition: StreamAligner.hpp:144
stream_vector streams
Definition: StreamAligner.hpp:222
size_t samples_dropped_late_arriving
Definition: StreamAlignerStatus.hpp:108
Definition: StreamAligner.hpp:48
const StreamStatus & getBufferStatus(int idx) const
Definition: StreamAligner.hpp:584
base::Time latestTimeStamp() const
Definition: StreamAligner.hpp:162
int priority
Definition: StreamAligner.hpp:60
base::Time getLatestTime() const
Definition: StreamAligner.hpp:575
void push(const base::Time &ts, const T &data)
Definition: StreamAligner.hpp:114
virtual base::Time latestDataTime() const
Definition: StreamAligner.hpp:170
void copyState(const StreamAligner &other)
Definition: StreamAligner.hpp:252
void unregisterStream(int idx)
Definition: StreamAligner.hpp:334
Definition: StreamAlignerStatus.hpp:84
Stream(callback_t callback, size_t bufferSize, base::Time period, int priority, const std::string &name)
Definition: StreamAligner.hpp:63