39 void open(
const std::string& file_name,
40 std::shared_ptr<adtf_file::SampleFactory> sample_factory,
41 std::shared_ptr<adtf_file::StreamTypeFactory> stream_type_factory)
override;
43 std::vector<adtf_file::Stream>
getStreams()
const override;
50 std::shared_ptr<const adtf_file::StreamType>
getStreamTypeBefore(uint64_t item_index, uint16_t stream_id)
override;
51 void seekTo(uint64_t item_index)
override;
68 std::shared_ptr<adtf_file::Reader>
addFile(
const std::string& file_name,
69 std::shared_ptr<adtf_file::SampleFactory> sample_factory,
70 std::shared_ptr<adtf_file::StreamTypeFactory> stream_type_factory);
80 ReaderSequence() =
default;
81 ReaderSequence(
const ReaderSequence&) =
delete;
82 ReaderSequence(ReaderSequence&&) =
default;
84 void addReader(
const std::shared_ptr<Reader>& reader, std::optional<uint64_t> first_item_index)
86 if (!_elements.empty())
90 std::unordered_map<uint16_t, std::string> streams;
91 for (
const auto& stream : reader.getStreams())
93 streams[stream.stream_id] = stream.name;
98 if (get_streams(*reader) != get_streams(*_elements.front().reader))
100 throw std::runtime_error(
101 "Sequence files (split files) do not share the same streams with the same stream ids.");
107 Element element{reader, reader->getNextItem(), first_item_index};
108 _elements.insert(std::lower_bound(_elements.begin(), _elements.end(), element), element);
112 if (_elements.empty())
114 _empty_streams_cache = reader->getStreams();
119 _current = _elements.begin();
120 _empty_streams_cache.clear();
123 adtf_file::FileItem getNext()
125 while (_current != _elements.end())
127 if (!_current->item_cache.empty())
129 auto item = std::move(_current->item_cache.front());
130 _current->item_cache.pop_front();
136 auto item = _current->reader->getNextItem();
137 adjustTimeStamp(item);
143 if (_current != _elements.end())
145 if (
const auto next = std::next(_current); next != _elements.end())
147 _current_timestamp_upper_bound = next->first_item_timestamp;
151 _current_timestamp_upper_bound.reset();
155 _current->item_cache.insert(_current->item_cache.begin(),
156 _current->initial_stream_types.cbegin(),
157 _current->initial_stream_types.cend());
165 uint64_t getItemIndexForTimeStamp(std::chrono::nanoseconds time_stamp)
167 std::exception_ptr last_exception;
168 for (
auto element = _elements.begin(); element != _elements.end(); ++element)
170 if (
auto next = std::next(element); next != _elements.end())
172 if (next->first_item_timestamp < time_stamp)
178 else if (element->first_item_timestamp >= time_stamp)
181 _seek_cache.emplace_back(element, *element->first_item_index);
182 return _seek_cache.size() - 1;
185 const auto seekable_reader = std::dynamic_pointer_cast<adtf_file::SeekableReader>(element->reader);
188 const auto item_index = seekable_reader->getItemIndexForTimeStamp(time_stamp);
189 _seek_cache.emplace_back(element, item_index);
190 return _seek_cache.size() - 1;
194 last_exception = std::current_exception();
199 std::rethrow_exception(last_exception);
205 std::shared_ptr<const adtf_file::StreamType> getStreamTypeBefore(uint64_t item_index, uint16_t stream_id)
207 if (item_index >= _seek_cache.size())
209 throw std::logic_error(
"Invalid seek cache");
212 const auto cache = _seek_cache[item_index];
213 const auto seekable_reader = std::dynamic_pointer_cast<adtf_file::SeekableReader>(cache.first->reader);
214 return seekable_reader->getStreamTypeBefore(cache.second, stream_id);
217 void seekTo(uint64_t item_index)
219 if (item_index >= _seek_cache.size())
221 throw std::logic_error(
"Invalid seek cache");
224 const auto cache = _seek_cache[item_index];
227 const auto seekable_reader = std::dynamic_pointer_cast<adtf_file::SeekableReader>(cache.first->reader);
228 seekable_reader->seekTo(cache.second);
229 cache.first->item_cache.clear();
232 _current = cache.first;
234 if (
const auto next = std::next(_current); next != _elements.end())
236 _current_timestamp_upper_bound = next->first_item_timestamp;
240 _current_timestamp_upper_bound.reset();
243 for (
auto file = std::next(cache.first); file != _elements.end(); ++file)
245 const auto seekable_reader = std::dynamic_pointer_cast<adtf_file::SeekableReader>(file->reader);
246 seekable_reader->seekTo(*file->first_item_index);
247 file->item_cache.clear();
251 std::vector<adtf_file::Stream> getStreams()
const
253 if (_elements.empty())
255 return _empty_streams_cache;
258 auto streams = _elements.front().reader->getStreams();
260 for (
auto element = ++_elements.begin(); element != _elements.end(); ++element)
262 const auto next_parts = element->reader->getStreams();
263 size_t stream_index = 0;
264 for (
const auto& stream : next_parts)
268 streams[stream_index].item_count += stream.item_count + streams.size();
269 streams[stream_index].timestamp_of_last_item = stream.timestamp_of_last_item;
277 std::optional<uint64_t> getItemCount()
const
279 uint64_t item_count = 0;
280 for (
const auto& element : _elements)
282 if (
const auto reader_item_count = element.reader->getItemCount())
284 item_count += *reader_item_count + element.initial_stream_types.size();
291 if (!_elements.empty())
294 item_count -= _elements.front().initial_stream_types.size();
299 std::optional<double> getProgress()
const
301 if (_current != _elements.end())
304 100.0 * ::std::distance<Sequence::const_iterator>(_elements.begin(), _current) / _elements.size();
306 if (
const auto file_progress = _current->reader->getProgress())
308 progress += std::min(std::max(*file_progress, 0.0), 100.0) / _elements.size();
319 void adjustTimeStamp(adtf_file::FileItem& item)
321 if (_current_timestamp_upper_bound && item.
time_stamp > *_current_timestamp_upper_bound)
323 item.
time_stamp = *_current_timestamp_upper_bound;
329 Element(std::shared_ptr<Reader> reader,
330 adtf_file::FileItem first_item,
331 std::optional<uint64_t> first_item_index):
332 reader(std::move(reader)),
333 first_item_timestamp(first_item.time_stamp),
334 first_item_index(first_item_index),
335 item_cache({std::move(first_item)})
337 for (
const auto& stream : this->reader->getStreams())
339 initial_stream_types.push_back(
340 adtf_file::FileItem{stream.stream_id, first_item_timestamp, stream.initial_type});
344 std::shared_ptr<Reader> reader;
346 std::chrono::nanoseconds first_item_timestamp;
348 std::optional<uint64_t> first_item_index;
351 std::deque<adtf_file::FileItem> item_cache;
354 std::vector<adtf_file::FileItem> initial_stream_types;
356 operator const std::chrono::nanoseconds&()
const
358 return first_item_timestamp;
361 bool operator<(
const Element& rhs)
const
363 return first_item_timestamp < rhs.first_item_timestamp;
367 using Sequence = std::vector<Element>;
369 Sequence::iterator _current = _elements.end();
370 std::vector<adtf_file::Stream> _empty_streams_cache;
371 std::optional<std::chrono::nanoseconds> _current_timestamp_upper_bound;
373 std::vector<std::pair<Sequence::iterator, uint64_t>> _seek_cache;
378 ReaderSequence sequence;
379 std::unordered_map<uint16_t, uint16_t> stream_ids;
381 adtf_file::FileItem getNext()
383 auto item = sequence.getNext();
389 std::vector<std::pair<std::shared_ptr<Reader>, std::optional<uint64_t>>> _readers;
390 std::vector<StreamGroup> _stream_groups;
391 std::vector<adtf_file::Stream> _streams;
394 size_t stream_group_index;
397 std::vector<GroupStream> _group_streams;
402 adtf_file::FileItem file_item;
404 std::multimap<std::chrono::nanoseconds, QueueItem> _queue;
406 using stream_group_indices = std::vector<std::optional<uint64_t>>;
407 using stream_group_cache_item =
408 std::pair<std::chrono::nanoseconds, stream_group_indices>;
409 std::vector<stream_group_cache_item> _seek_cache;
411 bool _all_readers_seekable =
true;