8#ifndef SEN_CORE_OBJ_DETAIL_EVENT_BUFFER_H
9#define SEN_CORE_OBJ_DETAIL_EVENT_BUFFER_H
45using SerializationFunc = std::function<void(
OutputStream&)>;
48struct SerializableEvent
50 MemberHash eventId = 0U;
51 TimeStamp creationTime;
52 SerializationFunc serializeFunc =
nullptr;
55 uint32_t serializedSize = 0U;
63class SerializableEventQueue
66 SEN_NOCOPY_NOMOVE(SerializableEventQueue)
69 SerializableEventQueue(std::size_t maxSize,
bool dropOldest): maxSize_(maxSize), dropOldest_(dropOldest) {}
71 ~SerializableEventQueue() =
default;
75 void push(SerializableEvent&& event);
76 void clear() { queue_.clear(); }
77 [[nodiscard]]
const std::list<SerializableEvent>& getContents() const noexcept {
return queue_; }
80 std::list<SerializableEvent> queue_;
82 std::size_t overflowCount_ = 0U;
94template <
typename... T>
95class EventBuffer final
97 SEN_MOVE_ONLY(EventBuffer)
103 EventBuffer() =
default;
108 [[nodiscard]] ConnectionGuard addConnection(Object* source, Callback callback, ConnId
id);
111 [[nodiscard]]
bool removeConnection(ConnId
id);
115 void produce(
Emit emissionMode,
117 TimeStamp creationTime,
120 bool addToTransportQueue,
121 SerializableEventQueue* transportQueue,
127 void dispatchFromStream(MemberHash eventId,
128 TimeStamp creationTime,
131 const Span<const uint8_t>& buffer,
132 Object* producer)
const;
135 void dispatch(MemberHash eventId,
136 TimeStamp creationTime,
139 bool addToTransportQueue,
140 SerializableEventQueue* transportQueue,
145 void immediateDispatch(MemberHash eventId,
146 TimeStamp creationTime,
155 using CallbackStorageType = std::shared_ptr<Callback>;
157 CallbackEntry(ConnId
id, CallbackStorageType callback): id_(id), callback_(std::move(callback)) {}
159 [[nodiscard]] ConnId getConnectionId() const noexcept {
return id_; }
160 [[nodiscard]]
const CallbackStorageType& getCallback() const noexcept {
return callback_; }
164 CallbackStorageType callback_;
169 mutable std::mutex callbacksMutex_;
170 std::vector<CallbackEntry> eventCallbacks_;
182inline void addVariantToList(
const T& val,
VarList& list)
185 VariantTraits<T>::valueToVariant(val, list.back());
194inline void SerializableEventQueue::push(SerializableEvent&& event)
198 queue_.push_back(std::move(event));
202 if (queue_.size() < maxSize_)
204 queue_.push_back(std::move(event));
211 queue_.push_back(std::move(event));
221template <
typename... T>
222inline EventBuffer<T...>::~EventBuffer()
226 for (
const auto& callbackEntry: eventCallbacks_)
228 if (
auto callbackLock = callbackEntry.getCallback()->lock(); callbackLock.isValid())
230 callbackLock.invalidate();
235template <
typename... T>
236inline ConnectionGuard EventBuffer<T...>::addConnection(Object* source, Callback callback, ConnId
id)
239 std::lock_guard lock(callbacksMutex_);
240 eventCallbacks_.emplace_back(
id, std::make_shared<
decltype(callback)>(std::move(callback)));
242 return {source,
id.get(), 0U,
true};
245template <
typename... T>
246inline bool EventBuffer<T...>::removeConnection(ConnId
id)
248 std::lock_guard lock(callbacksMutex_);
250 for (
auto itr = eventCallbacks_.begin(); itr != eventCallbacks_.end(); ++itr)
252 if (itr->getConnectionId() ==
id)
258 if (
auto callbackLock = itr->getCallback()->lock(); callbackLock.isValid())
260 callbackLock.invalidate();
263 eventCallbacks_.erase(itr);
271template <
typename... T>
272inline void EventBuffer<T...>::produce(
Emit emissionMode,
274 TimeStamp creationTime,
277 bool addToTransportQueue,
278 SerializableEventQueue* transportQueue,
285 immediateDispatch(eventId, creationTime, transportMode, producer, queue, args...);
290 [
this, eventId, creationTime, producerId, transportMode, addToTransportQueue, transportQueue, producer, args...]()
293 eventId, creationTime, producerId, transportMode, addToTransportQueue, transportQueue, producer, args...);
295 impl::cannotBeDropped(transportMode));
299template <
typename... T>
300inline void EventBuffer<T...>::dispatch(MemberHash eventId,
301 TimeStamp creationTime,
304 bool addToTransportQueue,
305 SerializableEventQueue* transportQueue,
309 EventInfo info {creationTime};
312 for (
const auto& callbackEntry: eventCallbacks_)
314 if (
auto callbackLock = callbackEntry.getCallback()->lock(); callbackLock.isValid())
316#if SEN_GCC_VERSION_CHECK_SMALLER(12, 4, 0)
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
324 auto work = [cb = callbackEntry.getCallback(), info, args...]() { cb->invoke(info, args...); };
326 callbackLock.pushAnswer(std::move(work), impl::cannotBeDropped(transportMode));
330 if (addToTransportQueue)
332 uint32_t serializedSize = 0U;
333 if constexpr (
sizeof...(T) != 0U)
335 serializedSize = (... + SerializationTraits<T>::serializedSize(args));
338 transportQueue->push({eventId,
343 (SerializationTraits<T>::write(out, args), ...);
350 producer->senImplEventEmitted(
355 result.reserve(sizeof...(T));
356 (addVariantToList<T>(args, result), ...);
362template <
typename... T>
363inline void EventBuffer<T...>::immediateDispatch(MemberHash eventId,
364 TimeStamp creationTime,
370 EventInfo info {creationTime};
372 for (
auto& callbackEntry: eventCallbacks_)
374 if (
auto callbackLock = callbackEntry.getCallback()->lock(); callbackLock.isValid())
376 if (callbackLock.isSameQueue(queue))
378 callbackEntry.getCallback()->invoke(info, args...);
382 auto work = [cb = callbackEntry.getCallback(), info, args...]() { cb->invoke(info, args...); };
383 callbackLock.pushAnswer(std::move(work), impl::cannotBeDropped(transportMode));
388 producer->senImplEventEmitted(
393 result.reserve(sizeof...(T));
394 (addVariantToList<T>(args, result), ...);
400template <
typename... T>
401inline void EventBuffer<T...>::dispatchFromStream(MemberHash eventId,
402 TimeStamp creationTime,
405 const Span<const uint8_t>& buffer,
406 Object* producer)
const
408 std::tuple<T...> argValues = {};
413 std::apply([&in](
auto&&... arg)
414 { ((SerializationTraits<std::remove_reference_t<
decltype(arg)>>::read(in, arg)), ...); },
419 &EventBuffer::dispatch,
420 std::tuple_cat(std::make_tuple(
this, eventId, creationTime, producerId, transportMode,
false,
nullptr, producer),
421 std::move(argValues)));
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