82 using OnRemoveCallback = std::function<
void(
const Function&)>;
88 data_(ListenersData{{}, {}})
112 data_(ListenersData{{}, std::move(on_listener_removal)})
126 template <
typename UpdaterFunc>
129 std::string_view name,
131 UpdaterFunc&& updater
133 const std::shared_lock lock(event_mutex_);
134 std::forward<UpdaterFunc>(updater)();
135 return DoAddListener(id, name, std::move(func));
139 template <
typename Class,
typename UpdaterFunc>
142 std::string_view name,
143 void (Class::*func)(Args...),
144 UpdaterFunc&& updater
146 return DoUpdateAndListen(
149 [obj, func](Args... args) { (obj->*func)(args...); },
150 std::forward<UpdaterFunc>(updater)
163 template <
typename UpdaterFunc>
165 utils::ResourceScopeStorage& scopes,
167 std::string_view name,
169 UpdaterFunc&& updater
171 auto state = std::make_unique<ResourceScopeState>();
174 const std::shared_lock lock(event_mutex_);
176 state->scope = DoAddListener(id, name, [&state_ref = *state](Args...) { state_ref.event_skipped =
true; });
177 auto data = data_.Lock();
178 state->listener = data->listeners.at(id);
183 constexpr utils::ResourceScopeStorage::Priority kPriority{-1};
185 utils::impl::InternalTag{},
187 [state = std::move(state), func = std::move(func), updater = std::forward<UpdaterFunc>(updater)]()
mutable {
188 const std::shared_lock sema_lock(state->listener->sema);
190 if (state->event_skipped) {
193 state->listener->callback = std::move(func);
194 return std::move(state->scope);
200 template <
typename Class,
typename UpdaterFunc>
202 utils::ResourceScopeStorage& scopes,
204 std::string_view name,
205 void (Class::*func)(Args...),
206 UpdaterFunc&& updater
208 DoUpdateAndListenScoped(
212 [obj, func](Args... args) { (obj->*func)(args...); },
213 std::forward<UpdaterFunc>(updater)
224 std::shared_ptr<
const Listener> listener;
227 std::vector<Task> tasks;
237 std::shared_lock<engine::SharedMutex> tmp_lock{event_mutex_, std::adopt_lock};
246 auto data = data_.Lock();
247 auto& listeners = data->listeners;
248 tasks.reserve(listeners.size());
250 for (
const auto& [_, listener] : listeners) {
251 tasks.push_back(Task{
255 [&, &callback = listener->callback, sema_lock = std::shared_lock(listener->sema)] {
264 for (
auto& task : tasks) {
265 impl::WaitForTask(task.listener->name, task.task);
270 const std::string&
Name()
const noexcept {
return name_; }
273 struct Listener
final {
275 mutable engine::Semaphore sema;
278 mutable Function callback;
279 std::string task_name;
281 Listener(std::string name, Function callback, std::string task_name)
283 name(std::move(name)),
284 callback(std::move(callback)),
285 task_name(std::move(task_name))
289 struct ResourceScopeState {
290 bool event_skipped{
false};
291 std::shared_ptr<
const Listener> listener;
293 AsyncEventSubscriberScope scope;
296 struct ListenersData
final {
297 std::unordered_map<FunctionId, std::shared_ptr<
const Listener>, FunctionId::Hash> listeners;
298 OnRemoveCallback on_listener_removal;
301 void RemoveListener(FunctionId id, [[maybe_unused]] UnsubscribingKind kind)
noexcept final {
302 const engine::TaskCancellationBlocker blocker;
303 const std::shared_lock lock(event_mutex_);
304 std::shared_ptr<
const Listener> listener;
305 OnRemoveCallback on_listener_removal;
308 auto data = data_.Lock();
309 auto& listeners = data->listeners;
310 const auto iter = listeners.find(id);
312 if (iter == listeners.end()) {
313 impl::ReportNotSubscribed(
Name());
317 listener = iter->second;
319 on_listener_removal = data->on_listener_removal;
321 listeners.erase(iter);
325 (
void)std::shared_lock(listener->sema);
330 if constexpr (impl::kCheckSubscriptionUB) {
332 impl::CheckDataUsedByCallbackHasNotBeenDestroyedBeforeUnsubscribing(
341 AsyncEventSubscriberScope DoAddListener(FunctionId id, std::string_view name, Function&& func)
final {
344 auto data = data_.Lock();
345 auto& listeners = data->listeners;
346 auto task_name = impl::MakeAsyncChannelName(name_, name);
347 const auto [iterator, success] = listeners.emplace(
349 std::make_shared<
const Listener>(std::string{name}, std::move(func), std::move(task_name))
352 impl::ReportAlreadySubscribed(
Name(), name);
354 return AsyncEventSubscriberScope(
utils::impl::InternalTag{}, *
this, id);
357 const std::string name_;
366 mutable engine::SharedMutex event_mutex_;