ipcq 0.12.0
Loading...
Searching...
No Matches
writer.hpp
Go to the documentation of this file.
1/**
2 * @file
3 * @ingroup ipcq
4 * @brief Contains declarations for @c ipcq::BasicWriter.
5 * @copyright
6 * SPDX-FileCopyrightText: 2018-2023 European Southern Observatory (ESO)
7 *
8 * SPDX-License-Identifier: LGPL-3.0-only
9 */
10#ifndef IPCQ_WRITER_HPP_
11#define IPCQ_WRITER_HPP_
12
13#include "policies.hpp"
14#include "detail/shm.hpp"
15
16namespace ipcq {
17
18/**
19 * Notification policy.
20 *
21 * @ingroup ipcq
22 */
23enum class Notify {
24 /**
25 * Don't notify.
26 */
28 /**
29 * Notify all waiting readers.
30 */
32};
33
34/**
35 * Queue writer class.
36 *
37 * @tparam T Type of elements in queue.
38 * @tparam ConditionPolicy is the policy to use to signal readers when
39 * new elements have been appended to queue. ipcq provides @a ipcq::BoostConditionPolicy
40 * and @a ipcq::SpinConditionPolicy.
41 * @tparam ShmTraits Implementation details. Leave default.
42 *
43 * Basic Usage:
44 *
45 * struct MyTopicType {
46 * ...
47 * };
48 *
49 * // Create writer for the topic "topicname" with type MyTopicType.
50 * auto writer = ipcq::Writer<MyTopicType>("topicname", 100);
51 *
52 * while (!shouldStop) {
53 * MyTopicType sample;
54 * // Prepare sample
55 * ...
56 * auto err = writer.Write(sample, ipcq::Notify::NotifyAll);
57 * if (err) {
58 * // Handle error
59 * break;
60 * }
61 * }
62 * // Signal to readers that the queue is now closed.
63 * writer.Close();
64 *
65 * @sa ipcq::BoostConditionPolicy ipcq::SpinConditionPolicy.
66 * @ingroup ipcq
67 */
68template <class T, class ConditionPolicy, class ShmTraits=detail::BoostInterprocessTraits>
70 using ShmControl = typename detail::ShmControl<T, ConditionPolicy, ShmTraits>;
71 ShmControl m_ctl;
72
73public:
74 /**
75 * Remove shared memory object from the system.
76 *
77 * @note Internally this unlinks the file.
78 *
79 * @param topic_name Topic name, as provided to BasicWriter constructor.
80 * @return 0 on success and contents of `errno` on failure.
81 */
82 static std::error_code RemoveShmObject(char const* topic_name) noexcept {
84 detail::MakeShmObjectName(topic_name).c_str());
85 }
86
87 /**
88 * Creates a writer and allocates shared memory to have room for capacity Ts.
89 *
90 * @param topic_name Name of the topic to create, together with the template parameter @c T this
91 * forms the full topic for the queue.
92 * @param capacity The queue capacity in number of elements.
93 * @param mempolicy An optional memory policy that is used for the shared memory allocation.
94 */
95 explicit BasicWriter(char const* topic_name,
96 size_t capacity,
97 std::optional<numapp::MemPolicy> const& mempolicy = std::nullopt)
98 : m_ctl{topic_name, capacity, mempolicy} {
99 }
100
101 /**
102 * Get queue topic name.
103 *
104 * @returns topic name in queue.
105 */
106 [[nodiscard]] char const* TopicName() const noexcept {
107 return m_ctl.TopicName();
108 }
109
110 /**
111 * Query queue capacity, which is the number of elements the queue can hold.
112 *
113 * @return queue capacity in number of elements.
114 */
115 [[nodiscard]] constexpr size_t Capacity() const noexcept {
116 return m_ctl.Capacity();
117 }
118
119 /**
120 * Query number of elements in the queue.
121 *
122 * @return queue size in number of elements.
123 */
124 [[nodiscard]] constexpr size_t Size() noexcept {
125 return m_ctl.Size();
126 }
127
128 /**
129 * Query whether queue is closed or not.
130 *
131 * @returns true if queue is closed, false otherwise.
132 */
133 [[nodiscard]] constexpr bool IsClosed() const noexcept {
134 return m_ctl.IsClosed();
135 }
136
137 /**
138 * Close queue for further writing.
139 *
140 * Any waiting readers will be notified. Once existing data is read out subsequent reads will
141 * fail with `ipcq::Error::Closed`.
142 *
143 * @post Shared memory remains valid but new reader instances will not attach.
144 *
145 * @releases
146 *
147 * @sa IsClosed()
148 */
149 constexpr void Close() noexcept {
150 m_ctl.Close();
151 m_ctl.AtomicStoreState();
152 m_ctl.Signal();
153 }
154
155 /**
156 * Append single element to queue, possibly overwriting the front element.
157 *
158 * @releases
159 *
160 * @param data Element to copy into queue
161 * @param notify_readers Determines if readers will be notified (woken up if blocked) that new
162 * element is available.
163 * @return error value is queue is closed.
164 */
165 [[nodiscard]] std::error_code Write(T const& data, Notify notify_readers) noexcept {
166 return Write(&data, 1, notify_readers);
167 }
168
169
170 /**
171 * Remove single element from queue.
172 *
173 * @releases
174 *
175 * @pre @code Size() > 0 @endcode
176 */
177 constexpr void Pop() noexcept {
178 m_ctl.Pop();
179 // Make this change visible to all readers
180 m_ctl.AtomicStoreState();
181 }
182
183 /**
184 * Write @c count elements to queue.
185 *
186 * @pre @code count <= Capacity() @endcode
187 *
188 * @note Function will make room as necessary to write @c count elements.
189 *
190 * @releases
191 *
192 * @return Error::Closed if queue is closed.
193 */
194 [[nodiscard]] std::error_code
195 Write(const T* data, size_t count, Notify notify_readers) noexcept {
196 IPCQ_ASSERT(count <= m_ctl.Capacity());
197
198 // See if we need to pop elements from queue
199 size_t num_free = static_cast<size_t>(m_ctl.Capacity() - m_ctl.Size());
200 if (count > num_free) {
201 int num_to_pop = count - num_free;
202 while (num_to_pop--) {
203 m_ctl.Pop();
204 }
205 // Make this change visible to all readers to invalidate that will be written to
206 m_ctl.AtomicStoreState();
207 }
208
209 if (auto err = m_ctl.Push(data, count); err) {
210 return err;
211 }
212
213 // Make this change visible to all reader by performing release synchronization.
214 m_ctl.AtomicStoreState();
215
216 if (notify_readers == Notify::All) {
217 m_ctl.Signal();
218 }
219
220 return {};
221 }
222};
223
224/**
225 * Convenience alias with reasonable default condition policy.
226 *
227 * @ingroup ipcq
228 */
229template <class T>
231
232} // namespace ipcq
233#endif
Queue writer class.
Definition writer.hpp:69
constexpr bool IsClosed() const noexcept
Query whether queue is closed or not.
Definition writer.hpp:133
char const * TopicName() const noexcept
Get queue topic name.
Definition writer.hpp:106
constexpr void Pop() noexcept
Remove single element from queue.
Definition writer.hpp:177
constexpr size_t Capacity() const noexcept
Query queue capacity, which is the number of elements the queue can hold.
Definition writer.hpp:115
constexpr void Close() noexcept
Close queue for further writing.
Definition writer.hpp:149
std::error_code Write(const T *data, size_t count, Notify notify_readers) noexcept
Write count elements to queue.
Definition writer.hpp:195
BasicWriter(char const *topic_name, size_t capacity, std::optional< numapp::MemPolicy > const &mempolicy=std::nullopt)
Creates a writer and allocates shared memory to have room for capacity Ts.
Definition writer.hpp:95
constexpr size_t Size() noexcept
Query number of elements in the queue.
Definition writer.hpp:124
std::error_code Write(T const &data, Notify notify_readers) noexcept
Append single element to queue, possibly overwriting the front element.
Definition writer.hpp:165
static std::error_code RemoveShmObject(char const *topic_name) noexcept
Remove shared memory object from the system.
Definition writer.hpp:82
ShmControl provides the basic data model for the queue and low level operations on the shared memory ...
Definition shm.hpp:658
static std::error_code RemoveShmObject(char const *shm_object_name) noexcept
Attempts to remove shared memory object.
Definition shm.hpp:692
Notify
Notification policy.
Definition writer.hpp:23
BasicWriter< T, BoostConditionPolicy > Writer
Convenience alias with reasonable default condition policy.
Definition writer.hpp:230
@ None
Don't notify.
Definition writer.hpp:27
@ All
Notify all waiting readers.
Definition writer.hpp:31
std::string MakeShmObjectName(char const *topic_name)
Creates the shared memory object name from topic_name.
Definition shm.hpp:225
Declares ipcq policies.
Contains core shm implementation details.
#define IPCQ_ASSERT(x)
Definition shm.hpp:49