ipcq 0.12.0
Loading...
Searching...
No Matches
reader.hpp
Go to the documentation of this file.
1/**
2 * @file
3 * @ingroup ipcq
4 * @brief Contains definition of ipcq::BasicReader.
5 * @copyright
6 * SPDX-FileCopyrightText: 2018-2023 European Southern Observatory (ESO)
7 *
8 * SPDX-License-Identifier: LGPL-3.0-only
9 */
10#ifndef IPCQ_READER_HPP_
11#define IPCQ_READER_HPP_
12#include <chrono>
13#include <optional>
14#include <thread>
15
16#include "detail/shm.hpp"
17#include "error.hpp"
18#include "policies.hpp"
19
20namespace ipcq {
21namespace detail {
22
23template <class>
24inline constexpr bool always_false_v = false; // NOLINT
25
26template <bool nothrow, class State>
28template <class State>
29struct ReaderRollback<true, State> {};
30template <class State>
31struct ReaderRollback<false, State> {
32 State state;
33};
34
35} // namespace detail
36
37template <class T, class ConditionPolicy, class ShmTraits>
38class BasicReader;
39
40/**
41 * Queue BasicReader class.
42 *
43 * @tparam T Type of elements in queue.
44 * @tparam ConditionPolicy Policy used to determine how readers are notified when
45 * queue has new elements appeneded to it.
46 * @tparam ShmTraits Implementation details. Leave default.
47 *
48 * Basic Usage:
49 *
50 * struct MyTopicType {
51 * ...
52 * };
53 * // Create reader for the topic "topicname" with type MyTopicType.
54 * // If the Writer is not yet created, this function will retry until it gives up
55 * // after the given timeout.
56 * auto reader = ipcq::Reader<MyTopicType>::MakeReader("topicname", 30s);
57 *
58 * while (true) {
59 * MyTopicType sample;
60 * // Read sample or timeout after 2s
61 * auto err = reader.Read(ipcq::OutputIteratorAdapter(&sample), 2s);
62 * if (err) {
63 * // Handle error
64 * break;
65 * }
66 * // Process sample
67 * ...
68 * }
69 *
70 * @ingroup ipcq
71 */
72template <class T, class ConditionPolicy, class ShmTraits = detail::BoostInterprocessTraits>
74 using ShmControl = typename detail::ShmControl<T, ConditionPolicy, ShmTraits>;
75 using ConstIterator = typename ShmControl::ConstIterator;
76 using CounterType = typename ShmControl::CounterType;
77
78 /**
79 * Local state of what has been copied.
80 */
81 ShmControl m_ctl;
82 struct State {
83 std::optional<ConstIterator> last_element; //< Iterator to next element to read from queue
84 CounterType next_counter;
85 };
86 State m_state;
87
88 /**
89 * Wait for iterator to next element to become valid
90 *
91 * Performs acquire-synchronization.
92 *
93 * @returns Error::Closed if queue is closed and element to read is not valid.
94 * @returns Error::Timeout if writer did not signal us before `timeout` time elapsed and
95 * element iterator become valid.
96 */
97 template <class Rep, class Period>
98 [[nodiscard]] std::pair<std::error_code, ConstIterator>
99 AwaitNextIterator(std::chrono::duration<Rep, Period> timeout) noexcept {
100 m_ctl.AtomicLoadState();
101
102 ConstIterator it{};
103 if (m_state.last_element) {
104 it = *m_state.last_element;
105 ++it; // This might advance to End()
106 } else {
107 it = m_ctl.Begin();
108 }
109
110 if (it == m_ctl.End()) {
111 // This may also happen if queue has 1 element or more generally if writer catches up to
112 // reader such that when m_last_element was incremented, it compares equally to "one
113 // past" last valid element.
114 if (m_ctl.IsClosed()) {
115 return {make_error_code(Error::Closed), {}};
116 }
117
118 // Predicate returns false if the waiting should be continued, true if waiting should be
119 // stopped.
120 //
121 // @todo: Optimize by only using a predicate if spurious wakeups are possible?
122 // @todo: Add prefetch hints?
123 if (!m_ctl.AwaitSignal(timeout, [this, &it]() mutable {
124 // Guard against spurious wakeups while waiting for condition
125 m_ctl.AtomicLoadState();
126 auto no_wait_more = m_ctl.IsClosed();
127
128 if (m_ctl.IsEmpty()) {
129 return no_wait_more;
130 }
131 // if the queue has a single element then this doesnt work
132 if (m_state.last_element) {
133 auto test_it = *m_state.last_element;
134 ++test_it; // note: may iterate to m_ctl.End()
135 if (test_it < m_ctl.Begin() || test_it >= m_ctl.End()) {
136 // invalid iterator
137 return no_wait_more;
138 }
139 it = test_it;
140 return true;
141 } else {
142 it = m_ctl.Begin();
143 return true;
144 }
145 })) {
146 // Timeout occured
147 return {make_error_code(Error::Timeout), {}};
148 } else {
149 // Not timeout, but could also be because Writer woke us up due to queue closing.
150 // If queue is closed, only report it as an error if iterator is invalid
151 if (it == m_ctl.End()) {
152 return {make_error_code(Error::Closed), it};
153 }
154 }
155 } else if (it < m_ctl.Begin() || it >= m_ctl.End()) {
157 }
158
159 // The predicate checks everything, so it should be good now!
160 IPCQ_ASSERT(it >= m_ctl.Begin() && it < m_ctl.End());
161 return {{}, it};
162 }
163
164public:
165 /**
166 * @name Construction
167 * @{
168 */
169 /**
170 * Create and connect reader to existing queue.
171 *
172 * @param topic_name Name of topic to associate reader with.
173 *
174 * @throws boost::interprocess::interprocess_exception if attaching to shared memory fails.
175 * @throws std::system_error with
176 * ipcq::Error::TypeMismatch if T or ConditionPolicy does not match types in queue.
177 * ipcq::Error::Closed if queue is closed, not yet ready or writer process does not exist.
178 *
179 * @sa BasicReader::MakeReader
180 */
181 explicit BasicReader(char const* topic_name) : m_ctl(topic_name), m_state{{}, {0}} {
182 // @todo: Decide if Reset() should be called automatically and immediately.
183 }
184
185 /**
186 * Factory function that retries to create a reader until success or timeout.
187 */
188 template <class Rep, class Period>
189 [[nodiscard]] static BasicReader
190 MakeReader(char const* topic_name, std::chrono::duration<Rep, Period> timeout) {
191 using namespace std::chrono;
192 auto expiry = high_resolution_clock::now() + timeout;
193
194 do {
195 try {
196 return BasicReader(topic_name);
197 } catch (boost::interprocess::interprocess_exception const& e) {
198 std::this_thread::sleep_for(500ms);
199 }
200 } while (high_resolution_clock::now() < expiry);
201 throw std::system_error(make_error_code(Error::Timeout),
202 "Timeout waiting for ipcq shared memory queue to be created");
203 }
204 /** @} */
205 /**
206 * @name Accessors
207 * @{
208 */
209 /**
210 * Get queue topic name.
211 *
212 * @returns topic name in queue.
213 */
214 [[nodiscard]] char const* TopicName() const noexcept {
215 return m_ctl.TopicName();
216 }
217
218 /**
219 * Query whether queue is closed or not.
220 *
221 * @returns true if queue is closed, false otherwise.
222 */
223 [[nodiscard]] constexpr bool IsClosed() const noexcept {
224 return m_ctl.IsClosed();
225 }
226
227 /**
228 * Query queue capacity, which is the number of elements the queue can hold.
229 *
230 * @return queue capacity in number of elements.
231 */
232 [[nodiscard]] constexpr size_t Capacity() const noexcept {
233 return m_ctl.Capacity();
234 }
235
236 /**
237 * Query number of elements in the queue.
238 *
239 * @return queue size in number of elements.
240 *
241 * @sa BasicReader::Synchronize
242 */
243 [[nodiscard]] constexpr size_t Size() const noexcept {
244 return m_ctl.Size();
245 }
246
247 /**
248 * Query number of elements that is currently available for reading. In other words it's the
249 * number of unread elements.
250 *
251 * @note This does not synchronize with shared memory.
252 *
253 * This can e.g. be used to calculate how near the reader is to be overwritten as the following
254 * approaches zero:
255 * @code
256 * reader.Size() - reader.NumAvailable()
257 * @endcode
258 *
259 * @return number of currently available elements.
260 */
261 [[nodiscard]] constexpr size_t NumAvailable() const noexcept {
262 if (m_state.last_element) {
263 return std::distance(*m_state.last_element, m_ctl.End()) -
264 1; // - 1 since End is one-past
265 } else {
266 return m_ctl.Size();
267 }
268 }
269 /** @} */
270
271 /**
272 * @name Read operations
273 * @{
274 */
275
276 /**
277 * Reads available data if any, or waits for notification from writer, and then reads available
278 * data from queue.
279 *
280 * @note Method is not guaranteed to read all the requested @c count elements.
281 *
282 * Reads a maximum of @a count elements from the queue synchronously by invoking the provided
283 * @c Operation @c op once with each element. The reader advances the next element to read
284 * automatically.
285 *
286 * If @a any data is already available then no waiting will occur.
287 * If no data is available then it will wait for writer notification and then
288 * read a maximum of `count` elements. If no notification occurs before specified @c timeout
289 * duration no reads will occur and @c ipcq::Error::Timeout is returned.
290 *
291 * @warning Data read from the queue must be considered invalid until the return value error
292 * code is verified to be 0.
293 *
294 * @acquires
295 *
296 * @param op Operation function used to operate on read element. Is invoked once for each
297 * element. Adapters to standard containers can be used as operation c.f. @c
298 * ipcq::OutputIteratorAdapter and @c ipcq::BackInserter.
299 * @param count Maximum number of elements to read from queue.
300 * @param timeout Maximum time to wait for any data, if no data is not already available to
301 * read. Value range and precision is guaranteed to 1 microsecond in range 1 microsecond to 256
302 * hours.
303 *
304 * @tparam Operation callable with requirements: @c `std::is_invocable_v<Operation, T
305 * const&>`, e.g. a callable with signature @c `void Operation(T const& element)`.
306 * or @c `void Operation(T const& element) noexcept`.
307 *
308 * @returns a pair with error code and number of read elements.
309 * @returns Error::InconsistentState if memory consistency errors occurred (see Reset() for how
310 * to recover).
311 * @returns Error::Timeout if no data became available before timing out.
312 *
313 * @throw Exceptions originating from @a op. Function provides strong exception guarantee for
314 * side-effects in the class itself but not in provided read-operation @c `op`. User may attempt
315 * to read the same elements again by issuing a new call to @a Read..
316 *
317 * @sa ipcq::OutputIteratorAdapter ipcq::BackInserter ipcq::BasicReader::Reset
318 */
319 template <class Operation, class Rep, class Period>
320 [[nodiscard]] std::pair<std::error_code, size_t>
321 Read(Operation&& op, size_t count, std::chrono::duration<Rep, Period> timeout) noexcept(
322 std::is_nothrow_invocable_v<Operation, T const&>) {
323 return ReadExt(std::forward<Operation>(op), count, timeout);
324 }
325
326 /**
327 * Skip up to @c count elements from queue.
328 *
329 * @note It is functionally equivalent as @c Read with a no-op @c op.
330 *
331 * @acquires
332 *
333 * @returns error, count pair where count is the number of elements skipped.
334 * @sa BasicReader::Read
335 */
336 template <class Rep, class Period>
337 [[nodiscard]] std::pair<std::error_code, size_t>
338 Skip(size_t count, std::chrono::duration<Rep, Period> timeout) noexcept {
339 return ReadExt([](T const&) noexcept {}, count, timeout);
340 }
341
342 /** @} */
343
344 /**
345 * @name Modifiers
346 * @{
347 */
348 /**
349 * Synchronizes internal state from shared memory.
350 *
351 * This is mainly useful to explicitly synchronize state to update non-synchronizing methods
352 * like:
353 * - BasicReader::Size()
354 * - BasicReader::NumAvailable()
355 *
356 * @note Synchronize should not be invoked for BasicReader::Read as it performs synchronization
357 * as a side-effect of reading.
358 *
359 * @acquires
360 */
361 constexpr void Synchronize() noexcept {
362 m_ctl.AtomicLoadState();
363 }
364
365 /**
366 * Reset internal reader state to synchronize with writer and recover from
367 * ipcq::Error::InconsistentState.
368 *
369 * To reset successfully the queue cannot be empty, as the reader must get state information
370 * from the queue to know what to expect next. In this case it returns ipcq::error::WouldBlock
371 * instead of waiting for data. Caller can then make the choice to assume the queue is in its
372 * initial state. If it was not the next read operation will fail with
373 * ipcq::Error::InconsistentState.
374 *
375 * @post On success the next element to read will be the next element to be written.
376 *
377 * @acquires
378 *
379 * @param keep Number of samples to keep as *unread*. Default is to keep 0 which is to say the
380 * next element to read is the next written. @c keep is automatically truncated to number of
381 * available elements.
382 *
383 * @returns 0 On success.
384 * @returns ipcq::Error::Closed if queue is closed and @c keep is 0 (i.e. reading will not be
385 * possible after Reset()).
386 * @returns ipcq::Error::WouldBlock if queue is empty.
387 */
388 [[nodiscard]] std::error_code Reset(std::size_t keep = 0u) noexcept {
389 m_ctl.AtomicLoadState();
390
391 if (keep == 0u && m_ctl.IsClosed()) {
393 }
394
395 if (m_ctl.IsEmpty()) {
396 // Note that empty does not mean that next element is 0. It could have been filled,
397 // and then emptied.
399 }
400 auto const distance_end = m_ctl.Size();
401 auto const distance_newest = distance_end - 1u;
402 keep = keep > distance_end ? distance_end : keep;
403
404 // Update m_last_element to the currently last element - keep.
405 auto it = m_ctl.Begin();
406 // The expected counter is always derived from last element to reduce likelihood of reading
407 // the counter value mid-update.
408 m_state.next_counter =
409 (it + distance_newest)->counter.load(std::memory_order_relaxed) - keep + 1u;
410 if (keep > distance_newest) {
411 // Set no last element so when waiting for next value Begin() is used.
412 m_state.last_element.reset();
413 } else {
414 it += distance_newest - keep;
415 m_state.last_element = it;
416 }
417
418 return {};
419 }
420 /** @} */
421
422private:
423 template <class Operation, class Rep, class Period>
424 [[nodiscard]] std::pair<std::error_code, size_t>
425 ReadExt(Operation op, size_t count, std::chrono::duration<Rep, Period> timeout) noexcept(
426 std::is_nothrow_invocable_v<Operation, T const&>) {
427 constexpr auto nothrow = std::is_nothrow_invocable_v<Operation, T const&>;
428 // Rollback is only used when operation may throw.
429 [[maybe_unused]] detail::ReaderRollback<nothrow, State> rollback;
430 if constexpr (!nothrow) {
431 rollback.state = m_state;
432 }
433
434 auto num_elements_copied{0u};
435
436 auto [err, it] = AwaitNextIterator(timeout);
437 if (err) {
438 return {err, num_elements_copied};
439 }
440 IPCQ_ASSERT(&*it != nullptr);
441
442 for (auto count_it = 0u; count_it < count; ++count_it) {
443 if (it == m_ctl.End()) {
444 return {{}, num_elements_copied};
445 }
446 if constexpr (nothrow) {
447 op(it->element);
448 } else if constexpr (std::is_invocable_v<Operation, T const&>) {
449 try {
450 op(it->element);
451 } catch (...) {
452 // Roll back and rethrow
453 m_state = rollback.state;
454 throw;
455 }
456 } else {
457 static_assert(detail::always_false_v<Operation>,
458 "Operation must be invocable with `T const&`");
459 }
460
461 // Load counter value atomically
462 auto counter = it->counter.load(std::memory_order_acquire);
463 if (counter != m_state.next_counter) {
464 // This may occur if reader missed an entire cycle of the queue such that the
465 // iterator position is valid, but only because it was completely missed, for
466 // example.
467 return {make_error_code(Error::InconsistentState), num_elements_copied};
468 }
469 // Success, iterate to next element (but only if queue has > 1 elements
470 if (Capacity() > 1) {
471 m_state.last_element = it;
472 }
473
474 ++num_elements_copied;
475 ++m_state.next_counter;
476 ++it;
477 } // end for
478
479 return {{}, num_elements_copied};
480 }
481};
482
483/**
484 * Convenience alias with reasonable default condition policy.
485 *
486 * @ingroup ipcq
487 */
488template <class T>
490
491} // namespace ipcq
492
493#endif
Queue BasicReader class.
Definition reader.hpp:73
BasicReader(char const *topic_name)
Create and connect reader to existing queue.
Definition reader.hpp:181
constexpr bool IsClosed() const noexcept
Query whether queue is closed or not.
Definition reader.hpp:223
std::pair< std::error_code, size_t > Skip(size_t count, std::chrono::duration< Rep, Period > timeout) noexcept
Skip up to count elements from queue.
Definition reader.hpp:338
constexpr void Synchronize() noexcept
Synchronizes internal state from shared memory.
Definition reader.hpp:361
constexpr size_t Capacity() const noexcept
Query queue capacity, which is the number of elements the queue can hold.
Definition reader.hpp:232
constexpr size_t Size() const noexcept
Query number of elements in the queue.
Definition reader.hpp:243
static BasicReader MakeReader(char const *topic_name, std::chrono::duration< Rep, Period > timeout)
Factory function that retries to create a reader until success or timeout.
Definition reader.hpp:190
std::pair< std::error_code, size_t > Read(Operation &&op, size_t count, std::chrono::duration< Rep, Period > timeout) noexcept(std::is_nothrow_invocable_v< Operation, T const & >)
Reads available data if any, or waits for notification from writer, and then reads available data fro...
Definition reader.hpp:321
constexpr size_t NumAvailable() const noexcept
Query number of elements that is currently available for reading.
Definition reader.hpp:261
std::error_code Reset(std::size_t keep=0u) noexcept
Reset internal reader state to synchronize with writer and recover from ipcq::Error::InconsistentStat...
Definition reader.hpp:388
char const * TopicName() const noexcept
Get queue topic name.
Definition reader.hpp:214
ShmControl provides the basic data model for the queue and low level operations on the shared memory ...
Definition shm.hpp:658
Contains declarations for ipcq::Error and related.
BasicReader< T, detail::BoostConditionPolicy > Reader
Convenience alias with reasonable default condition policy.
Definition reader.hpp:489
std::error_code make_error_code(Error e)
Function overload to create std::error_code from ipcq::Error.
Definition error.hpp:120
@ Closed
Queue is closed or otherwise not available for reading.
Definition error.hpp:33
@ InconsistentState
Reader is in an inconsistent state w.r.t.
Definition error.hpp:44
@ Timeout
Operation timed out.
Definition error.hpp:37
@ WouldBlock
Non-blocking operation would block.
Definition error.hpp:29
constexpr bool always_false_v
Definition reader.hpp:24
Declares ipcq policies.
Contains core shm implementation details.
#define IPCQ_ASSERT(x)
Definition shm.hpp:49