aggregator
StreamAligner.hpp
Go to the documentation of this file.
1 #ifndef __AGGREGATOR_HPP__
2 #define __AGGREGATOR_HPP__
3 
4 #include <base/Time.hpp>
5 #include <cmath>
6 #include <base-logging/Logging.hpp>
7 #include <vector>
8 #include <base/CircularBuffer.hpp>
9 #include <algorithm>
10 #include <boost/function.hpp>
11 #include <boost/tuple/tuple.hpp>
12 #include <stdexcept>
13 #include <iostream>
14 #include <aggregator/StreamAlignerStatus.hpp>
15 
16 namespace aggregator {
17 
19  {
20  class StreamBase
21  {
22  friend class StreamAligner;
23  public:
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;
32  virtual const StreamStatus &getBufferStatus() const = 0;
33  virtual void copyState( const StreamBase& other ) = 0;
34  virtual void clear() = 0;
35 
36  bool isActive() const { return active; }
37  void setActive( bool active ) { this->active = active; }
38 
39  friend std::ostream &operator<<(std::ostream &stream, const aggregator::StreamAligner::StreamBase &base);
40 
41  protected:
42  mutable StreamStatus status;
44  bool active;
45  };
46 
47  public:
48  template <class T> class Stream : public StreamBase
49  {
50  public:
51  typedef boost::function<void (const base::Time &ts, const T &value)> callback_t;
52 
53  protected:
54  typedef std::pair<base::Time,T> item;
55  boost::circular_buffer<item> buffer;
56  size_t bufferSize;
58  base::Time period;
59  base::Time lastTime;
60  int priority;
61 
62  public:
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)
65  {
66  status.name = name;
67  status.priority = priority;
68 
69  if (bufferSize > 0)
70  buffer.set_capacity( bufferSize );
71  else
72  {
73  // initial size, will be reallocated at runtime
74  buffer.set_capacity( 20 );
75  }
76  status.buffer_size = buffer.capacity();
77  }
78 
79  virtual ~Stream() {};
80 
81  bool getNextSample(item &sample) const
82  {
83  if(buffer.empty())
84  return false;
85 
86  sample = buffer.front();
87  return true;
88  }
89 
90  virtual int getPriority() const
91  {
92  return priority;
93  }
94 
95  virtual const StreamStatus &getBufferStatus() const
96  {
97  status.buffer_fill = buffer.size();
98  status.latest_data_time = latestDataTime();
99  status.earliest_data_time = earliestDataTime();
100  status.active = isActive();
101  return status;
102  }
103 
104  virtual void copyState( const StreamBase& other )
105  {
106  const Stream<T> &stream(dynamic_cast<const Stream<T>& >(other));
107 
108  lastTime = stream.lastTime;
109  buffer = stream.buffer;
110  bufferSize = stream.bufferSize;
111  status = stream.status;
112  }
113 
114  void push(const base::Time &ts, const T &data )
115  {
116  if(ts < lastTime)
117  {
118  status.samples_backward_in_time++;
119  return;
120  }
121 
122  lastTime = ts;
123 
124  if (buffer.full())
125  {
126  if (bufferSize > 0)
127  {
128  // if the buffer is full, just use the behaviour of the circular
129  // buffer: discard old data.
130  status.samples_dropped_buffer_full++;
131  }
132  else
133  {
134  buffer.set_capacity(buffer.capacity() * 2);
135  status.buffer_size = buffer.capacity();
136  }
137  }
138  buffer.push_back( std::make_pair(ts, data) );
139  }
140 
144  base::Time pop()
145  {
146  if( hasData() )
147  {
148  status.samples_processed++;
149  base::Time ts = buffer.front().first;
150  if(callback)
151  callback( ts, buffer.front().second );
152  buffer.pop_front();
153  return ts;
154  }
155 
156  throw std::runtime_error("pop() called on stream with no data.");
157  }
158 
159  bool hasData() const
160  { return !buffer.empty(); }
161 
162  base::Time latestTimeStamp() const
163  {
164  if( hasData() )
165  return buffer.front().first;
166  else
167  return lastTime + period;
168  }
169 
170  virtual base::Time latestDataTime() const
171  {
172  return lastTime;
173  }
174 
175  virtual base::Time earliestDataTime() const
176  {
177  if( hasData() )
178  return buffer.front().first;
179  return base::Time();
180  }
181 
182  virtual void clear()
183  {
184  lastTime = base::Time();
185  buffer.clear();
186 
187  status.latest_sample_time = base::Time();
188  status.latest_data_time = base::Time();
189  status.samples_dropped_buffer_full = 0;
190  status.samples_dropped_late_arriving = 0;
191  status.buffer_fill = 0;
192  status.active = true;
193  };
194  };
195 
196  static bool compareStreams( const StreamBase* b1, const StreamBase* b2 )
197  {
198  if(!b1)
199  return false;
200 
201  if(!b2)
202  return true;
203 
204  const base::Time &ts1( b1->latestTimeStamp() );
205  const base::Time &ts2( b2->latestTimeStamp() );
206 
207  if(ts1 == ts2)
208  {
209  if(b1->hasData() && !b2->hasData())
210  return true;
211 
212  if(!b1->hasData() && b2->hasData())
213  return false;
214 
215  return b1->getPriority() < b2->getPriority();
216  }
217 
218  return ts1 < ts2;
219  }
220 
221  typedef std::vector<StreamBase*> stream_vector;
223  base::Time timeout;
224 
226  base::Time latest_ts;
227 
229  base::Time current_ts;
230 
232 
237 
238  public:
239  explicit StreamAligner(base::Time timeout = base::Time::fromSeconds(1))
241 
242  virtual ~StreamAligner()
243  {
244  for(stream_vector::iterator it=streams.begin();it != streams.end();it++)
245  delete *it;
246  }
247 
252  void copyState(const StreamAligner& other)
253  {
254  latest_ts = other.latest_ts;
255  current_ts = other.current_ts;
256 
257  assert( streams.size() == other.streams.size() );
258  for(size_t i=0;i<streams.size();i++)
259  {
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");
264 
265  if(streams[i])
266  {
267  streams[i]->copyState( *other.streams[i] );
268  }
269  }
270  }
271 
276  void setTimeout(const base::Time &t )
277  {
278  timeout = t;
279  }
280 
293  void disableStream( int idx )
294  {
295  if(!streams[idx])
296  throw std::runtime_error("invalid stream index.");
297 
298  streams[idx]->setActive( false );
299  }
300 
307  void enableStream( int idx )
308  {
309  if(!streams[idx])
310  throw std::runtime_error("invalid stream index.");
311 
312  streams[idx]->setActive( true );
313  }
314 
320  bool isStreamActive( int idx ) const
321  {
322  if(!streams[idx])
323  throw std::runtime_error("invalid stream index.");
324 
325  return streams[idx]->isActive();
326  }
327 
334  void unregisterStream(int idx)
335  {
336  if(!streams[idx])
337  {
338  throw std::runtime_error("invalid stream index.");
339  }
340 
341  delete streams[idx];
342 
343  streams[idx] = 0;
344 
345  status.streams[idx].active = false;
346  }
347 
367  template <class T> int registerStream( typename Stream<T>::callback_t callback, int bufferSize, base::Time period, int priority = -1, const std::string &name = std::string())
368  {
369  if( bufferSize < 0 )
370  {
371  if( period == base::Time() )
372  {
373  throw std::runtime_error("No buffer size provided for stream with unknown period.");
374  }
375  else if( period < base::Time() )
376  {
377  // for a negative period, just calculate the buffer size, but don't set any lookahead.
378  bufferSize = buffer_size_factor * std::ceil( timeout.toSeconds() / -period.toSeconds() );
379  period = base::Time();
380  }
381  else
382  {
383  bufferSize = buffer_size_factor * std::ceil( timeout.toSeconds() / period.toSeconds() );
384  }
385  }
386 
387  if( bufferSize == 0 )
388  {
389  LOG_DEBUG_S << "dynamically allocating stream aligner buffer for stream: " << name;
390  }
391 
392  StreamBase *newStream = new Stream<T>(callback, bufferSize, period, priority, name);
393 
394  //check if there is a free slot from a previous deleted stream
395  for(size_t i = 0; i < streams.size(); i++)
396  {
397  if(!streams[i])
398  {
399  streams[i] = newStream;
400  status.streams[i] = StreamStatus();
401  return i;
402  }
403  }
404 
405  streams.push_back( newStream );
406  status.streams.push_back(StreamStatus());
407  return streams.size() - 1;
408  }
409 
418  template <class T> void push( int idx, const base::Time &ts, const T& data )
419  {
420  if( !streams.at(idx) )
421  throw std::runtime_error("invalid stream index.");
422 
423  Stream<T>* stream = dynamic_cast<Stream<T>*>(streams[idx]);
424  assert( stream );
425 
426  stream->status.samples_received++;
427  stream->status.latest_sample_time = ts;
428 
429  // mark stream as active, since it is receiving data items will
430  // have no effect on an already active stream, but enables
431  // streams which have been marked passive before.
432  stream->setActive( true );
433 
434  //any sample, that is older than the last replayed sample
435  //will never be played back and gets dropped by default
436  if(ts < current_ts)
437  {
439  stream->status.samples_dropped_late_arriving++;
440  return;
441  }
442 
443  if( ts > latest_ts )
444  latest_ts = ts;
445 
446  stream->push( ts, data );
447  }
448 
449  template <class T> bool getNextSample( int idx, std::pair<base::Time,T> &sample) const
450  {
451  if( !streams.at(idx) )
452  throw std::runtime_error("invalid stream index.");
453 
454  Stream<T>* stream = dynamic_cast<Stream<T>*>(streams[idx]);
455  assert( stream );
456 
457  return stream->getNextSample(sample);
458  }
459 
476  bool step()
477  {
478  if( streams.empty() )
479  return false;
480 
481  // copy streams vector and sort it by next ts
482  stream_vector items = streams;
483  std::sort( items.begin(), items.end(), &compareStreams );
484 
485  for(stream_vector::iterator it=items.begin();it != items.end();it++)
486  {
487  //first stream is unregistered no data there
488  if(!*it)
489  return false;
490 
491  if( (*it)->hasData() )
492  {
493  // if stream has current data, pop that data
494  current_ts = (*it)->pop();
495  return true;
496  }
497  else if( (*it)->isActive() )
498  {
499 
500  base::Time latestDataTime;
501  base::Time firstDataTime;
502 
503  //initalization case
504  if(current_ts == base::Time())
505  {
506  //check if one stream timed out
507  for(stream_vector::iterator it2=items.begin();it2 != items.end();it2++)
508  {
509 
510  if(*it2 && (*it2)->hasData())
511  {
512  if(latestDataTime < (*it2)->latestDataTime())
513  latestDataTime = (*it2)->latestDataTime();
514 
515  if(firstDataTime == base::Time() || firstDataTime > (*it2)->earliestDataTime())
516  firstDataTime = (*it2)->earliestDataTime();
517  }
518  }
519  } else {
520  latestDataTime = latest_ts;
521  firstDataTime = current_ts;
522  }
523 
524  if(latestDataTime - firstDataTime < timeout)
525  {
526  // if there is no data, but the expected data has
527  // not run out yet, wait for it.
528  return false;
529  }
530  }
531  }
532  return false;
533  }
534 
540  void clear()
541  {
542  for(size_t i = 0; i < streams.size(); i++)
543  {
544  if(streams[i])
545  {
546  streams[i]->clear();
547  }
548  }
549 
550  latest_ts = base::Time();
551  current_ts = base::Time();
552 
553  status.current_time = base::Time();
554  status.latest_time = base::Time();
556  }
557 
562  base::Time getTimeOut() const { return timeout; };
563 
567  base::Time getLatency() const { return latest_ts - current_ts; };
568 
571  base::Time getCurrentTime() const { return current_ts; };
572 
575  base::Time getLatestTime() const { return latest_ts; }
576 
579  int getStreamSize() const { return streams.size(); }
580 
584  const StreamStatus &getBufferStatus(int idx) const
585  {
586  if( !streams.at(idx) )
587  throw std::runtime_error("invalid stream index.");
588 
589  return streams[idx]->getBufferStatus();
590  }
591 
596  {
597  status.time = base::Time::now();
600 
601  for(size_t i=0;i<streams.size();i++)
602  {
603  if(streams[i])
604  status.streams[i] = streams[i]->getBufferStatus();
605  }
606 
607  return status;
608  }
609 
610  friend std::ostream &operator<<(std::ostream &stream, const aggregator::StreamAligner::StreamBase &base);
611  friend std::ostream &operator<<(std::ostream &stream, const aggregator::StreamAligner &re);
612  };
613 
614  inline std::ostream &operator<<(std::ostream &stream, const aggregator::StreamAligner &re)
615  {
616  using ::operator <<;
617  stream << "current time: " << re.getCurrentTime() << " latest time:" << re.getLatestTime() << " latency: " << re.getLatency() << std::endl;
618  for(size_t i=0;i<re.streams.size();i++)
619  {
620  stream << i << ":" << *re.streams[i] << std::endl;
621  }
622 
623  return stream;
624  }
625 
626  inline std::ostream &operator<<(std::ostream &stream, const aggregator::StreamAligner::StreamBase &base)
627  {
628  using ::operator <<;
629  stream << base.getBufferStatus();
630  return stream;
631  }
632 }
633 
634 #endif
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
size_t buffer_fill
Definition: StreamAlignerStatus.hpp:16
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
std::ostream & operator<<(std::ostream &stream, const aggregator::StreamAligner &re)
Definition: StreamAligner.hpp:614
void clear()
Definition: StreamAligner.hpp:540
callback_t callback
Definition: StreamAligner.hpp:57
boost::circular_buffer< item > buffer
Definition: StreamAligner.hpp:55
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
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