Sen API
Sen Libraries
Loading...
Searching...
No Matches
event_buffer.h
Go to the documentation of this file.
1// === event_buffer.h ==================================================================================================
2// Sen Infrastructure
3// Released under the Apache License v2.0 (SPDX-License-Identifier Apache-2.0).
4// See the LICENSE.txt file for more information.
5// © Airbus SAS, Airbus Helicopters, and Airbus Defence and Space SAU/GmbH/SAS.
6// =====================================================================================================================
7
8#ifndef SEN_CORE_OBJ_DETAIL_EVENT_BUFFER_H
9#define SEN_CORE_OBJ_DETAIL_EVENT_BUFFER_H
10
11// sen
15#include "sen/core/base/span.h"
19#include "sen/core/meta/type.h"
20#include "sen/core/meta/var.h"
24#include "sen/core/obj/object.h"
25
26// std
27#include <cstddef>
28#include <cstdint>
29#include <functional>
30#include <list>
31#include <memory>
32#include <mutex>
33#include <tuple>
34#include <type_traits>
35#include <utility>
36#include <vector>
37
38namespace sen::impl
39{
40
41//--------------------------------------------------------------------------------------------------------------
42// SerializableEvent
43//--------------------------------------------------------------------------------------------------------------
44
45using SerializationFunc = std::function<void(OutputStream&)>;
46
48struct SerializableEvent
49{
50 MemberHash eventId = 0U;
51 TimeStamp creationTime;
52 SerializationFunc serializeFunc = nullptr;
53 ObjectId producerId;
55 uint32_t serializedSize = 0U;
56};
57
58//--------------------------------------------------------------------------------------------------------------
59// SerializableEventQueue
60//--------------------------------------------------------------------------------------------------------------
61
63class SerializableEventQueue
64{
65public:
66 SEN_NOCOPY_NOMOVE(SerializableEventQueue)
67
68public:
69 SerializableEventQueue(std::size_t maxSize, bool dropOldest): maxSize_(maxSize), dropOldest_(dropOldest) {}
70
71 ~SerializableEventQueue() = default;
72
73public:
74 void setOnDropped(std_util::move_only_function<void()>&& func) { onDropped_ = std::move(func); }
75 void push(SerializableEvent&& event);
76 void clear() { queue_.clear(); }
77 [[nodiscard]] const std::list<SerializableEvent>& getContents() const noexcept { return queue_; }
78
79private:
80 std::list<SerializableEvent> queue_;
81 std::size_t maxSize_;
82 std::size_t overflowCount_ = 0U;
83 bool dropOldest_;
84 std_util::move_only_function<void()> onDropped_ = []() {};
85};
86
87//--------------------------------------------------------------------------------------------------------------
88// EventBuffer
89//--------------------------------------------------------------------------------------------------------------
90
94template <typename... T>
95class EventBuffer final
96{
97 SEN_MOVE_ONLY(EventBuffer)
98
99public:
100 using Callback = EventCallback<T...>;
101
102public: // special members
103 EventBuffer() = default;
104 ~EventBuffer();
105
106public:
108 [[nodiscard]] ConnectionGuard addConnection(Object* source, Callback callback, ConnId id);
109
111 [[nodiscard]] bool removeConnection(ConnId id);
112
115 void produce(Emit emissionMode,
116 MemberHash eventId,
117 TimeStamp creationTime,
118 ObjectId producerId,
119 TransportMode transportMode,
120 bool addToTransportQueue,
121 SerializableEventQueue* transportQueue,
122 WorkQueue* queue,
123 Object* producer,
124 MaybeRef<T>... args) const;
125
127 void dispatchFromStream(MemberHash eventId,
128 TimeStamp creationTime,
129 ObjectId producerId,
130 TransportMode transportMode,
131 const Span<const uint8_t>& buffer,
132 Object* producer) const;
133
135 void dispatch(MemberHash eventId,
136 TimeStamp creationTime,
137 ObjectId producerId,
138 TransportMode transportMode,
139 bool addToTransportQueue,
140 SerializableEventQueue* transportQueue,
141 Object* producer,
142 MaybeRef<T>... args) const;
143
145 void immediateDispatch(MemberHash eventId,
146 TimeStamp creationTime,
147 TransportMode transportMode,
148 Object* producer,
149 WorkQueue* queue,
150 MaybeRef<T>... args) const;
151
152private:
153 struct CallbackEntry
154 {
155 using CallbackStorageType = std::shared_ptr<Callback>;
156
157 CallbackEntry(ConnId id, CallbackStorageType callback): id_(id), callback_(std::move(callback)) {}
158
159 [[nodiscard]] ConnId getConnectionId() const noexcept { return id_; }
160 [[nodiscard]] const CallbackStorageType& getCallback() const noexcept { return callback_; }
161
162 private:
163 ConnId id_;
164 CallbackStorageType callback_;
165 };
169 mutable std::mutex callbacksMutex_;
170 std::vector<CallbackEntry> eventCallbacks_;
171};
172
173//----------------------------------------------------------------------------------------------------------------------
174// Inline implementation
175//----------------------------------------------------------------------------------------------------------------------
176
177//--------------------------------------------------------------------------------------------------------------
178// Helpers
179//--------------------------------------------------------------------------------------------------------------
180
181template <typename T>
182inline void addVariantToList(const T& val, VarList& list)
183{
184 list.emplace_back();
185 VariantTraits<T>::valueToVariant(val, list.back());
186}
187
188[[nodiscard]] inline bool cannotBeDropped(TransportMode mode) noexcept { return mode == TransportMode::confirmed; }
189
190//--------------------------------------------------------------------------------------------------------------
191// SerializableEventQueue
192//--------------------------------------------------------------------------------------------------------------
193
194inline void SerializableEventQueue::push(SerializableEvent&& event)
195{
196 if (maxSize_ == 0)
197 {
198 queue_.push_back(std::move(event));
199 return;
200 }
201
202 if (queue_.size() < maxSize_)
203 {
204 queue_.push_back(std::move(event));
205 return;
206 }
207
208 if (dropOldest_)
209 {
210 queue_.pop_front();
211 queue_.push_back(std::move(event));
212 }
213
214 onDropped_();
215}
216
217//--------------------------------------------------------------------------------------------------------------
218// EventBuffer
219//--------------------------------------------------------------------------------------------------------------
220
221template <typename... T>
222inline EventBuffer<T...>::~EventBuffer()
223{
224 // A queued dispatch holds the callback by shared_ptr and can run after this buffer is gone, so
225 // cancel whatever is still attached, as removeConnection() does for an explicit disconnect.
226 for (const auto& callbackEntry: eventCallbacks_)
227 {
228 if (auto callbackLock = callbackEntry.getCallback()->lock(); callbackLock.isValid())
229 {
230 callbackLock.invalidate();
231 }
232 }
233}
234
235template <typename... T>
236inline ConnectionGuard EventBuffer<T...>::addConnection(Object* source, Callback callback, ConnId id)
237{
238 {
239 std::lock_guard lock(callbacksMutex_);
240 eventCallbacks_.emplace_back(id, std::make_shared<decltype(callback)>(std::move(callback)));
241 }
242 return {source, id.get(), 0U, true};
243}
244
245template <typename... T>
246inline bool EventBuffer<T...>::removeConnection(ConnId id)
247{
248 std::lock_guard lock(callbacksMutex_);
249
250 for (auto itr = eventCallbacks_.begin(); itr != eventCallbacks_.end(); ++itr)
251 {
252 if (itr->getConnectionId() == id)
253 {
254
255 // the callback might survive the erasure from this list when
256 // enqueued in some runner. Therefore, we need to explicitly cancel it,
257 // so that it doesn't get invoked from now on.
258 if (auto callbackLock = itr->getCallback()->lock(); callbackLock.isValid())
259 {
260 callbackLock.invalidate();
261 }
262
263 eventCallbacks_.erase(itr);
264 return true;
265 }
266 }
267
268 return false;
269}
270
271template <typename... T>
272inline void EventBuffer<T...>::produce(Emit emissionMode,
273 MemberHash eventId,
274 TimeStamp creationTime,
275 ObjectId producerId,
276 TransportMode transportMode,
277 bool addToTransportQueue,
278 SerializableEventQueue* transportQueue,
279 WorkQueue* queue,
280 Object* producer,
281 MaybeRef<T>... args) const
282{
283 if (emissionMode == Emit::now)
284 {
285 immediateDispatch(eventId, creationTime, transportMode, producer, queue, args...);
286 }
287 else
288 {
289 queue->push(
290 [this, eventId, creationTime, producerId, transportMode, addToTransportQueue, transportQueue, producer, args...]()
291 {
292 dispatch(
293 eventId, creationTime, producerId, transportMode, addToTransportQueue, transportQueue, producer, args...);
294 },
295 impl::cannotBeDropped(transportMode));
296 }
297}
298
299template <typename... T>
300inline void EventBuffer<T...>::dispatch(MemberHash eventId,
301 TimeStamp creationTime,
302 ObjectId producerId,
303 TransportMode transportMode,
304 bool addToTransportQueue,
305 SerializableEventQueue* transportQueue,
306 Object* producer,
307 MaybeRef<T>... args) const
308{
309 EventInfo info {creationTime};
310
311 // for (const auto& callback: callbacks_)
312 for (const auto& callbackEntry: eventCallbacks_)
313 {
314 if (auto callbackLock = callbackEntry.getCallback()->lock(); callbackLock.isValid())
315 {
316#if SEN_GCC_VERSION_CHECK_SMALLER(12, 4, 0)
317 // TODO (SEN-717): clean up with gcc12.4 on debian
318# pragma GCC diagnostic push
319# pragma GCC diagnostic ignored "-Wmaybe-uninitialized"
320 auto work = [cb = callbackEntry.getCallback(), info, args...]() { cb->invoke(info, args...); };
321# pragma GCC diagnostic pop
322#else
323 // Note: if modified, patch line below
324 auto work = [cb = callbackEntry.getCallback(), info, args...]() { cb->invoke(info, args...); };
325#endif
326 callbackLock.pushAnswer(std::move(work), impl::cannotBeDropped(transportMode));
327 }
328 }
329
330 if (addToTransportQueue)
331 {
332 uint32_t serializedSize = 0U;
333 if constexpr (sizeof...(T) != 0U)
334 {
335 serializedSize = (... + SerializationTraits<T>::serializedSize(args));
336 }
337
338 transportQueue->push({eventId,
339 creationTime,
340 [args...](OutputStream& out)
341 {
342 std::ignore = out;
343 (SerializationTraits<T>::write(out, args), ...);
344 },
345 producerId,
346 transportMode,
347 serializedSize});
348 }
349
350 producer->senImplEventEmitted(
351 eventId,
352 [args...]()
353 {
354 VarList result;
355 result.reserve(sizeof...(T));
356 (addVariantToList<T>(args, result), ...);
357 return result;
358 },
359 info);
360}
361
362template <typename... T>
363inline void EventBuffer<T...>::immediateDispatch(MemberHash eventId,
364 TimeStamp creationTime,
365 TransportMode transportMode,
366 Object* producer,
367 WorkQueue* queue,
368 MaybeRef<T>... args) const
369{
370 EventInfo info {creationTime};
371
372 for (auto& callbackEntry: eventCallbacks_)
373 {
374 if (auto callbackLock = callbackEntry.getCallback()->lock(); callbackLock.isValid())
375 {
376 if (callbackLock.isSameQueue(queue))
377 {
378 callbackEntry.getCallback()->invoke(info, args...);
379 }
380 else
381 {
382 auto work = [cb = callbackEntry.getCallback(), info, args...]() { cb->invoke(info, args...); };
383 callbackLock.pushAnswer(std::move(work), impl::cannotBeDropped(transportMode));
384 }
385 }
386 }
387
388 producer->senImplEventEmitted(
389 eventId,
390 [args...]()
391 {
392 VarList result;
393 result.reserve(sizeof...(T));
394 (addVariantToList<T>(args, result), ...);
395 return result;
396 },
397 info);
398}
399
400template <typename... T>
401inline void EventBuffer<T...>::dispatchFromStream(MemberHash eventId,
402 TimeStamp creationTime,
403 ObjectId producerId,
404 TransportMode transportMode,
405 const Span<const uint8_t>& buffer,
406 Object* producer) const
407{
408 std::tuple<T...> argValues = {};
409
410 InputStream in(buffer);
411
412 // deserialize into the tuple
413 std::apply([&in](auto&&... arg)
414 { ((SerializationTraits<std::remove_reference_t<decltype(arg)>>::read(in, arg)), ...); },
415 argValues);
416
417 // call emit with the tuple arguments
418 std::apply(
419 &EventBuffer::dispatch,
420 std::tuple_cat(std::make_tuple(this, eventId, creationTime, producerId, transportMode, false, nullptr, producer),
421 std::move(argValues)));
422}
423
424} // namespace sen::impl
425
426#endif // SEN_CORE_OBJ_DETAIL_EVENT_BUFFER_H
Here we define a set of template meta-programming helpers to let the compiler take some decisions bas...
InputStreamTemplate< LittleEndian > InputStream
Definition input_stream.h:84
Emit
How to emit an event.
Definition object.h:51
Callback< EventInfo, Args... > EventCallback
An event callback.
Definition callback.h:208
@ now
Directly when it happens.
Definition object.h:52
std::conditional_t< std::is_arithmetic_v< T >||shouldBePassedByValueV< T >, T, AddConstRef< T > > MaybeRef
returns 'const T&' or 'T' depending on the type
Definition class_helpers.h:46
std::vector< Var > VarList
A list of vars to represent sequences.
Definition var.h:104
TransportMode
How to transport information.
Definition type.h:56
@ multicast
Directed to all receivers, unreliable, unordered, no congestion control.
Definition type.h:58
@ confirmed
Directed to each receiver, reliable, ordered, with congestion control, relatively heavyweight.
Definition type.h:59
detail::MoveOnlyFunctionImpl< FwdArgs... > move_only_function
Definition move_only_function.h:256
OutputStreamTemplate< LittleEndian > OutputStream
Definition output_stream.h:64