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