75 using ConstIterator =
typename ShmControl::ConstIterator;
76 using CounterType =
typename ShmControl::CounterType;
83 std::optional<ConstIterator> last_element;
84 CounterType next_counter;
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();
103 if (m_state.last_element) {
104 it = *m_state.last_element;
110 if (it == m_ctl.End()) {
114 if (m_ctl.IsClosed()) {
123 if (!m_ctl.AwaitSignal(timeout, [
this, &it]()
mutable {
125 m_ctl.AtomicLoadState();
126 auto no_wait_more = m_ctl.IsClosed();
128 if (m_ctl.IsEmpty()) {
132 if (m_state.last_element) {
133 auto test_it = *m_state.last_element;
135 if (test_it < m_ctl.Begin() || test_it >= m_ctl.End()) {
151 if (it == m_ctl.End()) {
155 }
else if (it < m_ctl.Begin() || it >= m_ctl.End()) {
160 IPCQ_ASSERT(it >= m_ctl.Begin() && it < m_ctl.End());
181 explicit BasicReader(
char const* topic_name) : m_ctl(topic_name), m_state{{}, {0}} {
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;
197 }
catch (boost::interprocess::interprocess_exception
const& e) {
198 std::this_thread::sleep_for(500ms);
200 }
while (high_resolution_clock::now() < expiry);
202 "Timeout waiting for ipcq shared memory queue to be created");
215 return m_ctl.TopicName();
223 [[nodiscard]]
constexpr bool IsClosed() const noexcept {
224 return m_ctl.IsClosed();
232 [[nodiscard]]
constexpr size_t Capacity() const noexcept {
233 return m_ctl.Capacity();
243 [[nodiscard]]
constexpr size_t Size() const noexcept {
262 if (m_state.last_element) {
263 return std::distance(*m_state.last_element, m_ctl.End()) -
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);
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);
362 m_ctl.AtomicLoadState();
388 [[nodiscard]] std::error_code
Reset(std::size_t keep = 0u)
noexcept {
389 m_ctl.AtomicLoadState();
391 if (keep == 0u && m_ctl.IsClosed()) {
395 if (m_ctl.IsEmpty()) {
400 auto const distance_end = m_ctl.Size();
401 auto const distance_newest = distance_end - 1u;
402 keep = keep > distance_end ? distance_end : keep;
405 auto it = m_ctl.Begin();
408 m_state.next_counter =
409 (it + distance_newest)->counter.load(std::memory_order_relaxed) - keep + 1u;
410 if (keep > distance_newest) {
412 m_state.last_element.reset();
414 it += distance_newest - keep;
415 m_state.last_element = it;
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&>;
430 if constexpr (!nothrow) {
431 rollback.state = m_state;
434 auto num_elements_copied{0u};
436 auto [err, it] = AwaitNextIterator(timeout);
438 return {err, num_elements_copied};
442 for (
auto count_it = 0u; count_it < count; ++count_it) {
443 if (it == m_ctl.End()) {
444 return {{}, num_elements_copied};
446 if constexpr (nothrow) {
448 }
else if constexpr (std::is_invocable_v<Operation, T const&>) {
453 m_state = rollback.state;
457 static_assert(detail::always_false_v<Operation>,
458 "Operation must be invocable with `T const&`");
462 auto counter = it->counter.load(std::memory_order_acquire);
463 if (counter != m_state.next_counter) {
467 return {
make_error_code(Error::InconsistentState), num_elements_copied};
470 if (Capacity() > 1) {
471 m_state.last_element = it;
474 ++num_elements_copied;
475 ++m_state.next_counter;
479 return {{}, num_elements_copied};
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.
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.
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...