userver: concurrent::AsyncEventChannel< Args > Class Template Reference
Loading...
Searching...
No Matches
concurrent::AsyncEventChannel< Args > Class Template Reference

#include <userver/concurrent/async_event_channel.hpp>

Detailed Description

template<typename... Args>
class concurrent::AsyncEventChannel< Args >

AsyncEventChannel is an in-process pub-sub with strict FIFO serialization, i.e. only after the event was processed a new event may appear for processing, same listener is never called concurrently.

Example usage:

enum class WeatherKind { kSunny, kRainy };
class WeatherStorage final {
public:
explicit WeatherStorage(WeatherKind value)
: value_(value),
channel_("weather")
{}
WeatherKind Get() const { return value_.load(); }
concurrent::AsyncEventSource<WeatherKind>& GetSource() { return channel_; }
template <typename Class>
void UpdateAndListen(
utils::ResourceScopeStorage& scopes,
Class* obj,
std::string_view name,
void (Class::*func)(WeatherKind)
) {
channel_.DoUpdateAndListenScoped(scopes, obj, name, func, [this, obj, func] { (obj->*func)(Get()); });
}
template <typename Class>
concurrent::AsyncEventSubscriberScope UpdateAndListen(
Class* obj,
std::string_view name,
void (Class::*func)(WeatherKind)
) {
return channel_.DoUpdateAndListen(obj, name, func, [&] { (obj->*func)(Get()); });
}
void Set(WeatherKind value) {
value_.store(value);
channel_.SendEvent(value);
}
private:
std::atomic<WeatherKind> value_;
concurrent::AsyncEventChannel<WeatherKind> channel_;
};
enum class CoatKind { kJacket, kRaincoat };
class CoatStorage final {
public:
explicit CoatStorage(utils::ResourceScopeStorage& scopes, WeatherStorage& weather_storage) {
weather_storage.UpdateAndListen(scopes, this, "coats", &CoatStorage::OnWeatherUpdate);
}
CoatKind Get() const { return value_.load(); }
private:
void OnWeatherUpdate(WeatherKind weather) { value_.store(ComputeCoat(weather)); }
static CoatKind ComputeCoat(WeatherKind weather);
std::atomic<CoatKind> value_{};
};
UTEST(AsyncEventChannel, UpdateAndListenSample) {
WeatherStorage weather_storage(WeatherKind::kSunny);
const utils::WithResourceScopes<CoatStorage> coat_storage(std::in_place, weather_storage);
EXPECT_EQ(coat_storage->Get(), CoatKind::kJacket);
weather_storage.Set(WeatherKind::kRainy);
EXPECT_EQ(coat_storage->Get(), CoatKind::kRaincoat);
}

Definition at line 79 of file async_event_channel.hpp.

Inheritance diagram for concurrent::AsyncEventChannel< Args >:

Public Types

using Function = typename AsyncEventSource<Args...>::Function
using OnRemoveCallback = std::function<void(const Function&)>

Public Member Functions

 AsyncEventChannel (std::string_view name)
 The primary constructor.
 AsyncEventChannel (std::string_view name, OnRemoveCallback on_listener_removal)
 The constructor with AsyncEventSubscriberScope usage checking.
template<typename UpdaterFunc>
AsyncEventSubscriberScope DoUpdateAndListen (FunctionId id, std::string_view name, Function &&func, UpdaterFunc &&updater)
 For use in UpdateAndListen of specific event channels.
template<typename Class, typename UpdaterFunc>
AsyncEventSubscriberScope DoUpdateAndListen (Class *obj, std::string_view name, void(Class::*func)(Args...), UpdaterFunc &&updater)
 This is an overloaded member function, provided for convenience. It differs from the above function only in what argument(s) it accepts.
template<typename UpdaterFunc>
void DoUpdateAndListenScoped (utils::ResourceScopeStorage &scopes, FunctionId id, std::string_view name, Function &&func, UpdaterFunc &&updater)
 Like DoUpdateAndListen, but binds the subscription to scopes.
template<typename Class, typename UpdaterFunc>
void DoUpdateAndListenScoped (utils::ResourceScopeStorage &scopes, Class *obj, std::string_view name, void(Class::*func)(Args...), UpdaterFunc &&updater)
 This is an overloaded member function, provided for convenience. It differs from the above function only in what argument(s) it accepts.
void SendEvent (Args... args) const
const std::string & Name () const noexcept
AsyncEventSubscriberScope AddListener (Class *obj, std::string_view name, void(Class::*func)(Args...))
 Subscribes to updates from this event source.

Member Typedef Documentation

◆ Function

template<typename... Args>
using concurrent::AsyncEventChannel< Args >::Function = typename AsyncEventSource<Args...>::Function

Definition at line 81 of file async_event_channel.hpp.

◆ OnRemoveCallback

template<typename... Args>
using concurrent::AsyncEventChannel< Args >::OnRemoveCallback = std::function<void(const Function&)>

Definition at line 82 of file async_event_channel.hpp.

Constructor & Destructor Documentation

◆ AsyncEventChannel() [1/2]

template<typename... Args>
concurrent::AsyncEventChannel< Args >::AsyncEventChannel ( std::string_view name)
inlineexplicit

The primary constructor.

Parameters
nameused for diagnostic purposes and is also accessible with Name

Definition at line 86 of file async_event_channel.hpp.

◆ AsyncEventChannel() [2/2]

template<typename... Args>
concurrent::AsyncEventChannel< Args >::AsyncEventChannel ( std::string_view name,
OnRemoveCallback on_listener_removal )
inline

The constructor with AsyncEventSubscriberScope usage checking.

The constructor with a callback that is called on listener removal, both on Unsubscribe and on automatic teardown. The callback takes a reference to Function as input. This is useful for checking the lifetime of data captured by the listener update function.

Note
Works only in debug mode.
Warning
Data captured by on_listener_removal function must be valid until the AsyncEventChannel object is completely destroyed.

Example usage:

auto on_remove = [](std::function<void(int)> func) { func(1); };
concurrent::AsyncEventChannel<int> channel("channel", on_remove);
int value = 0;
{
channel.AddListener(concurrent::FunctionId(&sub), "sub", [&value](int new_value) { value = new_value; });
sub.Unsubscribe();
}
if constexpr (concurrent::impl::kCheckSubscriptionUB) {
EXPECT_EQ(value, 1);
} else {
EXPECT_EQ(value, 0);
}
Parameters
nameused for diagnostic purposes and is also accessible with Name
on_listener_removalthe callback used for check
See also
impl::CheckDataUsedByCallbackHasNotBeenDestroyedBeforeUnsubscribing

Definition at line 110 of file async_event_channel.hpp.

Member Function Documentation

◆ AddListener()

AsyncEventSubscriberScope concurrent::AsyncEventSource< Args >::AddListener ( Class * obj,
std::string_view name,
void(Class::* func )(Args...) )
inlineinherited

Subscribes to updates from this event source.

The listener won't be called immediately. To process the current value and then listen to updates, use UpdateAndListen of specific event channels.

Warning
Listeners should not be added or removed while processing the event inside another listener.

Example usage:

UTEST(AsyncEventChannel, AddListenerSample) {
WeatherStorage weather_storage(WeatherKind::kSunny);
std::vector<WeatherKind> recorded_weather;
weather_storage.GetSource()
.AddListener(concurrent::FunctionId(&recorder), "recorder", [&](WeatherKind weather) {
recorded_weather.push_back(weather);
});
weather_storage.Set(WeatherKind::kRainy);
weather_storage.Set(WeatherKind::kSunny);
weather_storage.Set(WeatherKind::kSunny);
EXPECT_EQ(recorded_weather, (std::vector{WeatherKind::kRainy, WeatherKind::kSunny, WeatherKind::kSunny}));
recorder.Unsubscribe();
}
Parameters
objthe subscriber, which is the owner of the listener method, and is also used as the unique identifier of the subscription for this AsyncEventSource
namethe name of the subscriber, for diagnostic purposes
functhe listener method, usually called On<DataName>Update, e.g. OnConfigUpdate or OnCacheUpdate
Returns
a AsyncEventSubscriberScope controlling the subscription. Store the scope as a member and call Unsubscribe explicitly.

Definition at line 132 of file async_event_source.hpp.

◆ DoUpdateAndListen() [1/2]

template<typename... Args>
template<typename Class, typename UpdaterFunc>
AsyncEventSubscriberScope concurrent::AsyncEventChannel< Args >::DoUpdateAndListen ( Class * obj,
std::string_view name,
void(Class::* func )(Args...),
UpdaterFunc && updater )
inline

This is an overloaded member function, provided for convenience. It differs from the above function only in what argument(s) it accepts.

Definition at line 140 of file async_event_channel.hpp.

◆ DoUpdateAndListen() [2/2]

template<typename... Args>
template<typename UpdaterFunc>
AsyncEventSubscriberScope concurrent::AsyncEventChannel< Args >::DoUpdateAndListen ( FunctionId id,
std::string_view name,
Function && func,
UpdaterFunc && updater )
inline

For use in UpdateAndListen of specific event channels.

Atomically calls updater, which should invoke func with the previously sent event, and subscribes to new events as if using AddListener.

Parameters
idthe subscriber class instance, see also a simpler DoUpdateAndListen overload below
namethe name of the subscriber
functhe callback that is called on each update
updaterthe initial () -> void callback that should call func with the current value
See also
AsyncEventSource::AddListener

Definition at line 127 of file async_event_channel.hpp.

◆ DoUpdateAndListenScoped() [1/2]

template<typename... Args>
template<typename Class, typename UpdaterFunc>
void concurrent::AsyncEventChannel< Args >::DoUpdateAndListenScoped ( utils::ResourceScopeStorage & scopes,
Class * obj,
std::string_view name,
void(Class::* func )(Args...),
UpdaterFunc && updater )
inline

This is an overloaded member function, provided for convenience. It differs from the above function only in what argument(s) it accepts.

Definition at line 201 of file async_event_channel.hpp.

◆ DoUpdateAndListenScoped() [2/2]

template<typename... Args>
template<typename UpdaterFunc>
void concurrent::AsyncEventChannel< Args >::DoUpdateAndListenScoped ( utils::ResourceScopeStorage & scopes,
FunctionId id,
std::string_view name,
Function && func,
UpdaterFunc && updater )
inline

Like DoUpdateAndListen, but binds the subscription to scopes.

Synchronously calls updater and subscribes with a stub that only records whether an event arrived during construction. When the scope is entered, updater is called again if an event was skipped, and the stub is replaced with func. Unsubscribe runs in utils::ResourceScopeStorage::BeforeDestruction.

Warning
updater is not only invoked inline. It is stored and may run again after this call returns, when the scope is entered. Do not capture locals by reference: copy the pointers and values the updater needs (for example obj and func) and capture this explicitly.

Definition at line 164 of file async_event_channel.hpp.

◆ Name()

template<typename... Args>
const std::string & concurrent::AsyncEventChannel< Args >::Name ( ) const
inlinenoexcept
Returns
the name of this event channel

Definition at line 270 of file async_event_channel.hpp.

◆ SendEvent()

template<typename... Args>
void concurrent::AsyncEventChannel< Args >::SendEvent ( Args... args) const
inline

Send the next event and wait until all the listeners process it.

Strict FIFO serialization is guaranteed, i.e. only after this event is processed a new event may be delivered for the subscribers, same listener/subscriber is never called concurrently.

Definition at line 222 of file async_event_channel.hpp.


The documentation for this class was generated from the following file: