File liveliness_manager.hpp
File List > astutedds > dcps > liveliness_manager.hpp
Go to the documentation of this file
//
// Copyright (c) 2026, Astute Systems PTY LTD
//
// This file is part of the Astute DDS developed by Astute Systems.
//
// See the commercial LICENSE file in the project root for full license details.
//
// @file liveliness_manager.hpp
// @brief DDS Liveliness QoS Policy Manager
//
// Implements AUTOMATIC, MANUAL_BY_PARTICIPANT, and MANUAL_BY_TOPIC liveliness policies.
// Reference: DDS DCPS specification 2.2.3.10
//
#ifndef ASTUTEDDS_DCPS_LIVELINESS_MANAGER_HPP
#define ASTUTEDDS_DCPS_LIVELINESS_MANAGER_HPP
#include <astutedds/dcps/dds_psm_types.hpp>
#include <astutedds/dcps/qos.hpp>
#include <astutedds/dcps/timer_service.hpp>
#include <astutedds/rtps/rtps_types.hpp>
#include <atomic>
#include <chrono>
#include <cstring>
#include <functional>
#include <map>
#include <memory>
#include <mutex>
#include <string>
#include <vector>
namespace astutedds::dcps
{
// Forward declarations
class DataWriter;
class DataReader;
class DomainParticipant;
struct LivelinessStatus
{
bool alive{true};
std::chrono::steady_clock::time_point last_asserted;
LivelinessQosPolicyKind kind{LivelinessQosPolicyKind::AUTOMATIC_LIVELINESS_QOS};
rtps::Duration_t lease_duration{rtps::Time_t::TIME_INFINITE()};
};
inline InstanceHandle_t liveliness_handle_from_guid(const rtps::GUID_t& g) noexcept
{
uint32_t h = 2166136261u;
for (uint8_t b : g.guidPrefix.value)
h = (h ^ b) * 16777619u;
for (uint8_t b : g.entityId.entityKey)
h = (h ^ b) * 16777619u;
h = (h ^ g.entityId.entityKind) * 16777619u;
// Keep it positive-looking (avoid HANDLE_NIL_VALUE = 0).
int32_t out = static_cast<int32_t>(h & 0x7FFFFFFF);
return out == 0 ? 1 : out;
}
using LivelinessChangedCallback = std::function<void(const LivelinessChangedStatus&)>;
using LivelinessLostCallback = std::function<void(const LivelinessLostStatus&)>;
using LivelinessChangedRoutingCallback =
std::function<void(const rtps::GUID_t&, const LivelinessChangedStatus&)>;
using LivelinessLostRoutingCallback =
std::function<void(const rtps::GUID_t&, const LivelinessLostStatus&)>;
using AssertionSenderCallback = std::function<void(const rtps::GUID_t&)>;
using ManualLivelinessSenderCallback = std::function<void(uint8_t /*kind*/)>;
struct TrackedWriter
{
rtps::GUID_t guid;
LivelinessQosPolicyKind kind;
std::chrono::steady_clock::time_point last_asserted;
std::chrono::milliseconds lease_duration;
bool alive{true};
bool is_local{false};
TimerId lease_timer_id{INVALID_TIMER_ID};
TimerId assertion_timer_id{INVALID_TIMER_ID};
};
class LivelinessManager
{
public:
explicit LivelinessManager(const rtps::GUID_t& participant_guid);
~LivelinessManager();
// Non-copyable
LivelinessManager(const LivelinessManager&) = delete;
LivelinessManager& operator=(const LivelinessManager&) = delete;
void start();
void stop();
void register_writer(const rtps::GUID_t& writer_guid, const LivelinessQosPolicy& qos);
void unregister_writer(const rtps::GUID_t& writer_guid);
void add_remote_writer(const rtps::GUID_t& writer_guid, const LivelinessQosPolicy& qos);
void remove_remote_writer(const rtps::GUID_t& writer_guid);
bool assert_liveliness(const rtps::GUID_t& writer_guid);
bool assert_liveliness_for_participant();
void on_data_received(const rtps::GUID_t& writer_guid);
LivelinessStatus get_status(const rtps::GUID_t& writer_guid) const;
LivelinessChangedStatus get_liveliness_changed_status() const;
void set_liveliness_changed_callback(LivelinessChangedCallback callback);
void set_liveliness_lost_callback(LivelinessLostCallback callback);
void set_liveliness_lost_routing_callback(LivelinessLostRoutingCallback cb);
void set_liveliness_changed_routing_callback(LivelinessChangedRoutingCallback cb);
void set_assertion_sender(AssertionSenderCallback callback);
void set_manual_liveliness_sender(ManualLivelinessSenderCallback callback);
void assert_remote_participant_liveliness(const rtps::GuidPrefix_t& sender_prefix,
uint8_t kind_byte);
LivelinessLostStatus get_liveliness_lost_status() const;
bool is_writer_alive(const rtps::GUID_t& writer_guid) const;
int32_t alive_count() const;
int32_t not_alive_count() const;
private:
static std::chrono::milliseconds duration_to_ms(const rtps::Duration_t& duration);
void notify_liveliness_changed(const rtps::GUID_t& writer_guid, bool alive);
void notify_liveliness_lost(const rtps::GUID_t& writer_guid);
void arm_writer_timers_locked_(TrackedWriter& w);
void cancel_writer_timers_locked_(TrackedWriter& w);
void on_local_lease_missed_(rtps::GUID_t guid);
void on_remote_lease_expired_(rtps::GUID_t guid);
void on_assertion_tick_(rtps::GUID_t guid);
rtps::GUID_t participant_guid_;
mutable std::mutex writers_mutex_;
std::map<rtps::GUID_t, TrackedWriter> tracked_writers_;
mutable std::mutex status_mutex_;
mutable LivelinessChangedStatus liveliness_changed_status_;
mutable LivelinessLostStatus liveliness_lost_status_;
mutable std::mutex callbacks_mutex_;
LivelinessChangedCallback liveliness_changed_callback_;
LivelinessLostCallback liveliness_lost_callback_;
LivelinessChangedRoutingCallback liveliness_changed_routing_callback_;
LivelinessLostRoutingCallback liveliness_lost_routing_callback_;
AssertionSenderCallback assertion_sender_;
ManualLivelinessSenderCallback manual_liveliness_sender_;
std::atomic<bool> started_{false};
TimerService timers_;
};
class WriterLiveliness
{
public:
WriterLiveliness(LivelinessManager& manager, const rtps::GUID_t& writer_guid, const LivelinessQosPolicy& qos);
~WriterLiveliness();
bool assert_liveliness();
void on_data_written();
LivelinessLostStatus get_liveliness_lost_status() const;
void set_listener(LivelinessLostCallback callback);
private:
LivelinessManager& manager_;
rtps::GUID_t writer_guid_;
LivelinessQosPolicy qos_;
LivelinessLostCallback callback_;
};
class ReaderLiveliness
{
public:
ReaderLiveliness(LivelinessManager& manager, const rtps::GUID_t& reader_guid, const LivelinessQosPolicy& qos);
~ReaderLiveliness() = default;
LivelinessChangedStatus get_liveliness_changed_status() const;
void set_listener(LivelinessChangedCallback callback);
bool is_matched_writer_alive(const rtps::GUID_t& writer_guid) const;
private:
LivelinessManager& manager_;
rtps::GUID_t reader_guid_;
LivelinessQosPolicy qos_;
};
} // namespace astutedds::dcps
#endif // ASTUTEDDS_DCPS_LIVELINESS_MANAGER_HPP