forked from Kistler-Group/sdbus-cpp
feat: introduce sd-event integration
This commit is contained in:
committed by
Stanislav Angelovic
parent
8d24b2826e
commit
e95024192c
@@ -32,6 +32,9 @@
|
||||
#include <sdbus-c++/Error.h>
|
||||
#include "ScopeGuard.h"
|
||||
#include SDBUS_HEADER
|
||||
#ifndef SDBUS_basu // sd_event integration is not supported in basu-based sdbus-c++
|
||||
#include <systemd/sd-event.h>
|
||||
#endif
|
||||
#include <unistd.h>
|
||||
#include <poll.h>
|
||||
#include <sys/eventfd.h>
|
||||
@@ -274,6 +277,189 @@ void Connection::addMatchAsync(const std::string& match, message_handler callbac
|
||||
floatingMatchRules_.push_back(addMatchAsync(match, std::move(callback), std::move(installCallback)));
|
||||
}
|
||||
|
||||
void Connection::attachSdEventLoop(sd_event *event, int priority)
|
||||
{
|
||||
#ifndef SDBUS_basu
|
||||
auto pollData = getEventLoopPollData();
|
||||
|
||||
auto sdEvent = createSdEventSlot(event);
|
||||
auto sdTimeEventSource = createSdTimeEventSourceSlot(event, priority);
|
||||
auto sdIoEventSource = createSdIoEventSourceSlot(event, pollData.fd, priority);
|
||||
auto sdInternalEventSource = createSdInternalEventSourceSlot(event, pollData.eventFd, priority);
|
||||
|
||||
sdEvent_ = std::make_unique<SdEvent>(SdEvent{ std::move(sdEvent)
|
||||
, std::move(sdTimeEventSource)
|
||||
, std::move(sdIoEventSource)
|
||||
, std::move(sdInternalEventSource) });
|
||||
#else
|
||||
(void)event;
|
||||
(void)priority;
|
||||
SDBUS_THROW_ERROR("sd_event integration is not supported on this platform", EOPNOTSUPP);
|
||||
#endif
|
||||
}
|
||||
|
||||
void Connection::detachSdEventLoop()
|
||||
{
|
||||
sdEvent_.reset();
|
||||
}
|
||||
|
||||
sd_event *Connection::getSdEventLoop()
|
||||
{
|
||||
return sdEvent_ ? static_cast<sd_event*>(sdEvent_->sdEvent.get()) : nullptr;
|
||||
}
|
||||
|
||||
#ifndef SDBUS_basu
|
||||
|
||||
Slot Connection::createSdEventSlot(sd_event *event)
|
||||
{
|
||||
// Get default event if no event is provided by the caller
|
||||
if (event != nullptr)
|
||||
event = sd_event_ref(event);
|
||||
else
|
||||
(void)sd_event_default(&event);
|
||||
SDBUS_THROW_ERROR_IF(!event, "Invalid sd_event handle", EINVAL);
|
||||
|
||||
return Slot{event, [](void* event){ sd_event_unref((sd_event*)event); }};
|
||||
}
|
||||
|
||||
Slot Connection::createSdTimeEventSourceSlot(sd_event *event, int priority)
|
||||
{
|
||||
sd_event_source *timeEventSource{};
|
||||
auto r = sd_event_add_time(event, &timeEventSource, CLOCK_MONOTONIC, 0, 0, onSdTimerEvent, this);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to add timer event", -r);
|
||||
Slot sdTimeEventSource{timeEventSource, [](void* source){ deleteSdEventSource((sd_event_source*)source); }};
|
||||
|
||||
r = sd_event_source_set_priority(timeEventSource, priority);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set time event priority", -r);
|
||||
|
||||
r = sd_event_source_set_description(timeEventSource, "bus-time");
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set time event description", -r);
|
||||
|
||||
return sdTimeEventSource;
|
||||
}
|
||||
|
||||
Slot Connection::createSdIoEventSourceSlot(sd_event *event, int fd, int priority)
|
||||
{
|
||||
sd_event_source *ioEventSource{};
|
||||
auto r = sd_event_add_io(event, &ioEventSource, fd, 0, onSdIoEvent, this);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to add io event", -r);
|
||||
Slot sdIoEventSource{ioEventSource, [](void* source){ deleteSdEventSource((sd_event_source*)source); }};
|
||||
|
||||
r = sd_event_source_set_prepare(ioEventSource, onSdEventPrepare);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set prepare callback for IO event", -r);
|
||||
|
||||
r = sd_event_source_set_priority(ioEventSource, priority);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set priority for IO event", -r);
|
||||
|
||||
r = sd_event_source_set_description(ioEventSource, "bus-input");
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set priority for IO event", -r);
|
||||
|
||||
return sdIoEventSource;
|
||||
}
|
||||
|
||||
Slot Connection::createSdInternalEventSourceSlot(sd_event *event, int fd, int priority)
|
||||
{
|
||||
sd_event_source *internalEventSource{};
|
||||
auto r = sd_event_add_io(event, &internalEventSource, fd, 0, onSdInternalEvent, this);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to add internal event", -r);
|
||||
Slot sdInternalEventSource{internalEventSource, [](void* source){ deleteSdEventSource((sd_event_source*)source); }};
|
||||
|
||||
// sd-event loop calls prepare callbacks for all event sources, not just for the one that fired now.
|
||||
// So since onSdEventPrepare is already registered on ioEventSource, we don't need to duplicate it here.
|
||||
//r = sd_event_source_set_prepare(internalEventSource, onSdEventPrepare);
|
||||
//SDBUS_THROW_ERROR_IF(r < 0, "Failed to set prepare callback for internal event", -r);
|
||||
|
||||
r = sd_event_source_set_priority(internalEventSource, priority);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set priority for internal event", -r);
|
||||
|
||||
r = sd_event_source_set_description(internalEventSource, "internal-event");
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set priority for IO event", -r);
|
||||
|
||||
return sdInternalEventSource;
|
||||
}
|
||||
|
||||
int Connection::onSdTimerEvent(sd_event_source */*s*/, uint64_t /*usec*/, void *userdata)
|
||||
{
|
||||
auto connection = static_cast<Connection*>(userdata);
|
||||
assert(connection != nullptr);
|
||||
|
||||
(void)connection->processPendingEvent();
|
||||
|
||||
return 1;
|
||||
}
|
||||
|
||||
int Connection::onSdIoEvent(sd_event_source */*s*/, int /*fd*/, uint32_t /*revents*/, void *userdata)
|
||||
{
|
||||
auto connection = static_cast<Connection*>(userdata);
|
||||
assert(connection != nullptr);
|
||||
|
||||
(void)connection->processPendingEvent();
|
||||
|
||||
return 1;
|
||||
}
|
||||
|
||||
int Connection::onSdInternalEvent(sd_event_source */*s*/, int /*fd*/, uint32_t /*revents*/, void *userdata)
|
||||
{
|
||||
auto connection = static_cast<Connection*>(userdata);
|
||||
assert(connection != nullptr);
|
||||
|
||||
// It's not really necessary to processPendingEvent() here. We just clear the event fd.
|
||||
// The sd-event loop will before the next poll call prepare callbacks for all event sources,
|
||||
// including I/O bus fd. This will get up-to-date poll timeout, which will be zero if there
|
||||
// are pending D-Bus messages in the read queue, which will immediately wake up next poll
|
||||
// and go to onSdIoEvent() handler, which calls processPendingEvent(). Viola.
|
||||
// For external event loops that only have access to public sdbus-c++ API, processPendingEvent()
|
||||
// is the only option to clear event fd (it comes at a little extra cost but on the other hand
|
||||
// the solution is simpler for clients -- we don't provide an extra method for just clearing
|
||||
// the event fd. There is one method for both fd's -- and that's processPendingEvent().
|
||||
|
||||
// Kept here so that potential readers know what to do in their custom external event loops.
|
||||
//(void)connection->processPendingEvent();
|
||||
|
||||
connection->eventFd_.clear();
|
||||
|
||||
return 1;
|
||||
}
|
||||
|
||||
int Connection::onSdEventPrepare(sd_event_source */*s*/, void *userdata)
|
||||
{
|
||||
auto connection = static_cast<Connection*>(userdata);
|
||||
assert(connection != nullptr);
|
||||
|
||||
auto sdbusPollData = connection->getEventLoopPollData();
|
||||
|
||||
// Set poll events to watch out for on I/O fd
|
||||
auto* sdIoEventSource = static_cast<sd_event_source*>(connection->sdEvent_->sdIoEventSource.get());
|
||||
auto r = sd_event_source_set_io_events(sdIoEventSource, sdbusPollData.events);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set poll events for IO event source", -r);
|
||||
|
||||
// Set poll events to watch out for on internal event fd
|
||||
auto* sdInternalEventSource = static_cast<sd_event_source*>(connection->sdEvent_->sdInternalEventSource.get());
|
||||
r = sd_event_source_set_io_events(sdInternalEventSource, POLLIN);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set poll events for internal event source", -r);
|
||||
|
||||
// Set current timeout to the time event source (it may be zero if there are messages in the sd-bus queues to be processed)
|
||||
auto* sdTimeEventSource = static_cast<sd_event_source*>(connection->sdEvent_->sdTimeEventSource.get());
|
||||
r = sd_event_source_set_time(sdTimeEventSource, static_cast<uint64_t>(sdbusPollData.timeout.count()));
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to set timeout for time event source", -r);
|
||||
r = sd_event_source_set_enabled(sdTimeEventSource, SD_EVENT_ON);
|
||||
SDBUS_THROW_ERROR_IF(r < 0, "Failed to enable time event source", -r);
|
||||
|
||||
return 1;
|
||||
}
|
||||
|
||||
void Connection::deleteSdEventSource(sd_event_source *s)
|
||||
{
|
||||
#if LIBSYSTEMD_VERSION>=243
|
||||
sd_event_source_disable_unref(s);
|
||||
#else
|
||||
sd_event_source_set_enabled(s, SD_EVENT_OFF);
|
||||
sd_event_source_unref(s);
|
||||
#endif
|
||||
}
|
||||
|
||||
#endif // SDBUS_basu
|
||||
|
||||
Slot Connection::addObjectVTable( const std::string& objectPath
|
||||
, const std::string& interfaceName
|
||||
, const sd_bus_vtable* vtable
|
||||
|
||||
@@ -38,6 +38,8 @@
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
||||
struct sd_event_source;
|
||||
|
||||
namespace sdbus::internal {
|
||||
|
||||
class Connection final
|
||||
@@ -97,6 +99,10 @@ namespace sdbus::internal {
|
||||
[[nodiscard]] Slot addMatchAsync(const std::string& match, message_handler callback, message_handler installCallback) override;
|
||||
void addMatchAsync(const std::string& match, message_handler callback, message_handler installCallback, floating_slot_t) override;
|
||||
|
||||
void attachSdEventLoop(sd_event *event, int priority) override;
|
||||
void detachSdEventLoop() override;
|
||||
sd_event *getSdEventLoop() override;
|
||||
|
||||
const ISdBus& getSdBusInterface() const override;
|
||||
ISdBus& getSdBusInterface() override;
|
||||
|
||||
@@ -162,6 +168,19 @@ namespace sdbus::internal {
|
||||
static int sdbus_match_install_callback(sd_bus_message *sdbusMessage, void *userData, sd_bus_error *retError);
|
||||
|
||||
private:
|
||||
#ifndef SDBUS_basu // sd_event integration is not supported if instead of libsystemd we are based on basu
|
||||
Slot createSdEventSlot(sd_event *event);
|
||||
Slot createSdTimeEventSourceSlot(sd_event *event, int priority);
|
||||
Slot createSdIoEventSourceSlot(sd_event *event, int fd, int priority);
|
||||
Slot createSdInternalEventSourceSlot(sd_event *event, int fd, int priority);
|
||||
static void deleteSdEventSource(sd_event_source *s);
|
||||
|
||||
static int onSdTimerEvent(sd_event_source *s, uint64_t usec, void *userdata);
|
||||
static int onSdIoEvent(sd_event_source *s, int fd, uint32_t revents, void *userdata);
|
||||
static int onSdInternalEvent(sd_event_source *s, int fd, uint32_t revents, void *userdata);
|
||||
static int onSdEventPrepare(sd_event_source *s, void *userdata);
|
||||
#endif
|
||||
|
||||
struct EventFd
|
||||
{
|
||||
EventFd();
|
||||
@@ -180,6 +199,15 @@ namespace sdbus::internal {
|
||||
sd_bus_slot *slot;
|
||||
};
|
||||
|
||||
// sd-event integration
|
||||
struct SdEvent
|
||||
{
|
||||
Slot sdEvent;
|
||||
Slot sdTimeEventSource;
|
||||
Slot sdIoEventSource;
|
||||
Slot sdInternalEventSource;
|
||||
};
|
||||
|
||||
private:
|
||||
std::unique_ptr<ISdBus> sdbus_;
|
||||
BusPtr bus_;
|
||||
@@ -187,6 +215,7 @@ namespace sdbus::internal {
|
||||
EventFd loopExitFd_; // To wake up event loop I/O polling to exit
|
||||
EventFd eventFd_; // To wake up event loop I/O polling to re-enter poll with fresh PollData values
|
||||
std::vector<Slot> floatingMatchRules_;
|
||||
std::unique_ptr<SdEvent> sdEvent_; // Integration of systemd sd-event event loop implementation
|
||||
};
|
||||
|
||||
}
|
||||
|
||||
+5
-3
@@ -300,11 +300,13 @@ int SdBus::sd_bus_open_server(sd_bus **ret, int fd)
|
||||
|
||||
int SdBus::sd_bus_open_system_remote(sd_bus **ret, const char *host)
|
||||
{
|
||||
#ifdef SDBUS_basu
|
||||
#ifndef SDBUS_basu
|
||||
return ::sd_bus_open_system_remote(ret, host);
|
||||
#else
|
||||
(void)ret;
|
||||
(void)host;
|
||||
// https://git.sr.ht/~emersion/basu/commit/01d33b244eb6
|
||||
return -EOPNOTSUPP;
|
||||
#else
|
||||
return ::sd_bus_open_system_remote(ret, host);
|
||||
#endif
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user