File dcps.hpp

File List > astutedds > dcps > dcps.hpp

Go to the documentation of this file


#ifndef ASTUTEDDS_DCPS_DCPS_HPP
#define ASTUTEDDS_DCPS_DCPS_HPP

#include <astutedds/dcps/qos.hpp>
#include <astutedds/dcps/dds_psm_types.hpp>
#include <astutedds/dcps/builtin_topic_data.hpp>
#include <astutedds/dcps/builtin_readers.hpp>
#include <astutedds/dcps/liveliness_manager.hpp>
#include <astutedds/dcps/timer_service.hpp>
#include <astutedds/dcps/topic_description.hpp>
#include <astutedds/rtps/rtps_types.hpp>
#include <astutedds/rtps/udp_transport.hpp>
#include <array>
#include <memory>
#include <mutex>
#include <optional>
#include <vector>
#include <deque>
#include <string>
#include <functional>
#include <set>
#include <map>
#include <unordered_map>
#include <atomic>
#include <thread>
#include <chrono>
#include <condition_variable>

namespace astutedds::dcps
{

    // Forward declarations
    class DomainParticipant;
    class DomainParticipantFactory;
    class Topic;
    class Publisher;
    class Subscriber;
    class DataWriter;
    class DataReader;
    class ContentFilteredTopic;
    class ReadCondition;

    // ReturnCode_t is now defined as an unscoped enum in dds_psm_types.hpp
    // (included above). No duplicate definition here.

    // ------------------------------------------------------------------
    // Entity framework helpers (DDS §2.2.2.1 Entity)
    // ------------------------------------------------------------------
    // Every Entity (DomainParticipant, Topic, Publisher, Subscriber,
    // DataWriter, DataReader) exposes:
    //   - enable() / is_enabled()          §2.2.2.1.1.7
    //   - get_instance_handle()             §2.2.2.1.1.8
    //   - get_statuscondition()             §2.2.2.1.1.5
    //   - get_status_changes()              §2.2.2.1.1.6
    //
    // Each Entity is assigned a unique InstanceHandle_t at construction
    // via next_entity_instance_handle(); the value is non-NIL and stable
    // for the entity's lifetime.
    inline InstanceHandle_t next_entity_instance_handle() noexcept
    {
        static std::atomic<InstanceHandle_t> counter{HANDLE_NIL + 1};
        return counter.fetch_add(1, std::memory_order_relaxed);
    }

    class EntityStatusCondition
    {
    public:
        EntityStatusCondition() = default;

        void set_enabled_statuses(StatusMask mask) noexcept { enabled_statuses_ = mask; }
        StatusMask get_enabled_statuses() const noexcept { return enabled_statuses_; }

        void set_status_changes(StatusMask mask) noexcept { status_changes_ = mask; }
        StatusMask get_status_changes() const noexcept { return status_changes_; }

        bool trigger_value() const noexcept
        {
            return (enabled_statuses_ & status_changes_) != 0;
        }

    private:
        StatusMask enabled_statuses_{STATUS_MASK_ALL};
        StatusMask status_changes_{STATUS_MASK_NONE};
    };

    struct SampleInfo
    {
        bool valid_data{false};
        rtps::Time_t source_timestamp{};
        // DDS §2.2.2.5.5 — the wall-clock time at which the middleware
        // received the sample from the network (or observed it locally
        // for in-process writes).  Distinct from source_timestamp which
        // is the writer's asserted timestamp.
        rtps::Time_t reception_timestamp{};
        InstanceHandle_t instance_handle{HANDLE_NIL};
        InstanceHandle_t publication_handle{HANDLE_NIL};
        int32_t disposed_generation_count{0};
        int32_t no_writers_generation_count{0};
        int32_t sample_rank{0};
        int32_t generation_rank{0};
        int32_t absolute_generation_rank{0};

        // DDS PSM state fields
        SampleStateKind sample_state{NOT_READ_SAMPLE_STATE};
        ViewStateKind view_state{NEW_VIEW_STATE};
        InstanceStateKind instance_state{ALIVE_INSTANCE_STATE};
    };

    struct CacheChange
    {
        rtps::SequenceNumber_t sequenceNumber{};
        std::vector<uint8_t> serializedData;
        rtps::Time_t sourceTimestamp{};
        rtps::ChangeKind_t kind{rtps::ChangeKind_t::ALIVE};
    };

    struct ReceivedSample
    {
        rtps::SequenceNumber_t sequenceNumber{};
        std::vector<uint8_t> serializedData;
        rtps::Time_t sourceTimestamp{};
        // DDS §2.2.2.5.5 — wall-clock at which the middleware admitted
        // the sample into the reader cache.  Copied out to
        // SampleInfo.reception_timestamp on read/take.  Stamped by
        // DataReader::add_sample() when the value is not preset by an
        // upstream transport hook.
        rtps::Time_t receptionTimestamp{};
        rtps::GUID_t writerGuid{};
        bool valid{true};
        bool read{false};
        int32_t ownershipStrength{0};   // For EXCLUSIVE ownership filtering
        InstanceHandle_t instance_handle{HANDLE_NIL};  // Computed from @key fields
        // DDS §2.2.2.5.1.7 — carried through to SampleInfo.instance_state so
        // DISPOSE / UNREGISTER notifications flow from wire to application.
        InstanceStateKind instance_state{ALIVE_INSTANCE_STATE};
        // DDS §2.2.2.5.1.5 — snapshot of the reader's per-instance generation
        // counters at the moment this sample was admitted into the cache.
        // Copied out to SampleInfo on read/take.  Ranks (sample_rank,
        // generation_rank, absolute_generation_rank) are computed from
        // these snapshots per §2.2.2.5.1.6.
        int32_t disposed_generation_count{0};
        int32_t no_writers_generation_count{0};
        // Phase 3.d — RTPS 2.5 §8.7.6 coherent-set marker.  Non-zero
        // indicates this sample carries PID_COHERENT_SET on the wire and
        // belongs to the coherent set that started at this writer
        // sequence number.  0 means no coherent-set tag (regular sample).
        // end_of_coherent_set is set when the sample also carried
        // PID_END_COHERENT_SET (impl-specific closing marker emitted by
        // Publisher::end_coherent_changes for the last write in a set).
        uint64_t coherent_set_start_sn{0};
        bool     end_of_coherent_set{false};
    };

    class Topic : public TopicDescription
    {
    public:
        Topic(const std::string &name, const std::string &type_name, DomainParticipant *participant);
        ~Topic() override = default;

        // TopicDescription interface
        const char *get_name() const override { return name_.c_str(); }
        const char *get_type_name() const override { return type_name_.c_str(); }

        // Legacy accessors (used by shapes_demo)
        const std::string &name() const { return name_; }
        const std::string &type_name() const { return type_name_; }
        DomainParticipant *participant() const { return participant_; }

        void set_qos(const TopicQos &qos) { qos_ = qos; }
        const TopicQos &get_qos() const { return qos_; }

        // Entity framework (DDS §2.2.2.1)
        ReturnCode_t enable() noexcept { enabled_ = true; return ReturnCode_t::RETCODE_OK; }
        bool is_enabled() const noexcept { return enabled_; }
        InstanceHandle_t get_instance_handle() const noexcept { return instance_handle_; }
        StatusMask get_status_changes() const noexcept { return status_condition_.get_status_changes(); }
        EntityStatusCondition* get_statuscondition() noexcept { return &status_condition_; }

    private:
        std::string name_;
        std::string type_name_;
        DomainParticipant *participant_;
        TopicQos qos_;
        bool enabled_{false};
        InstanceHandle_t instance_handle_{next_entity_instance_handle()};
        EntityStatusCondition status_condition_{};
    };

    class DataWriter
    {
    public:
        DataWriter(Topic *topic, Publisher *publisher);
        virtual ~DataWriter() = default;
        DataWriter(const DataWriter &) = delete;
        DataWriter &operator=(const DataWriter &) = delete;

        bool write(const std::vector<uint8_t> &data);

        bool write(const std::vector<uint8_t> &xcdr1Data,
                   const std::vector<uint8_t> &xcdr2Data);

        uint64_t write_coherent(const std::vector<uint8_t> &xcdr1Data,
                                 const std::vector<uint8_t> &xcdr2Data,
                                 uint64_t in_start_sn,
                                 bool     end_of_set);

        uint64_t register_instance();

        uint64_t register_instance(const std::vector<uint8_t> &xcdr1Data);
        uint64_t register_instance(const std::vector<uint8_t> &xcdr1Data,
                                   const std::vector<uint8_t> &xcdr2Data);

        uint64_t lookup_instance(const std::vector<uint8_t> &xcdr1Data) const;

        bool unregister_instance(uint64_t handle);
        bool unregister_instance(const std::vector<uint8_t> &xcdr1Data);
        bool unregister_instance(const std::vector<uint8_t> &xcdr1Data,
                                 const std::vector<uint8_t> &xcdr2Data);

        bool dispose(uint64_t handle);
        bool dispose(const std::vector<uint8_t> &xcdr1Data);
        bool dispose(const std::vector<uint8_t> &xcdr1Data,
                     const std::vector<uint8_t> &xcdr2Data);

        //---------------------------------------------------------------------
        // DDS 1.4 §2.2.2.4.2.14 / .16 / .17 / .18 — `_w_timestamp` variants.
        // Every operation below mirrors its untimestamped sibling but pins
        // the sample's source_timestamp to @p source_timestamp instead of
        // deriving it from the current wall clock.  The timestamp is
        // recorded in the CacheChange for local delivery and, for the
        // dispose / unregister paths, propagated through the transport so
        // the wire INFO_TS matches.
        //---------------------------------------------------------------------
        bool write_w_timestamp(const std::vector<uint8_t> &data,
                               const DdsTime_t &source_timestamp);
        bool write_w_timestamp(const std::vector<uint8_t> &xcdr1Data,
                               const std::vector<uint8_t> &xcdr2Data,
                               const DdsTime_t &source_timestamp);

        uint64_t register_instance_w_timestamp(const std::vector<uint8_t> &xcdr1Data,
                                                const DdsTime_t &source_timestamp);
        uint64_t register_instance_w_timestamp(const std::vector<uint8_t> &xcdr1Data,
                                                const std::vector<uint8_t> &xcdr2Data,
                                                const DdsTime_t &source_timestamp);

        bool unregister_instance_w_timestamp(uint64_t handle,
                                              const DdsTime_t &source_timestamp);
        bool unregister_instance_w_timestamp(const std::vector<uint8_t> &xcdr1Data,
                                              const DdsTime_t &source_timestamp);
        bool unregister_instance_w_timestamp(const std::vector<uint8_t> &xcdr1Data,
                                              const std::vector<uint8_t> &xcdr2Data,
                                              const DdsTime_t &source_timestamp);

        bool dispose_w_timestamp(uint64_t handle,
                                  const DdsTime_t &source_timestamp);
        bool dispose_w_timestamp(const std::vector<uint8_t> &xcdr1Data,
                                  const DdsTime_t &source_timestamp);
        bool dispose_w_timestamp(const std::vector<uint8_t> &xcdr1Data,
                                  const std::vector<uint8_t> &xcdr2Data,
                                  const DdsTime_t &source_timestamp);

        ReturnCode_t get_key_value(std::vector<uint8_t> &key_holder,
                                    InstanceHandle_t handle) const;

        ReturnCode_t wait_for_acknowledgments(const Duration_t &max_wait)
        {
            if (qos_.reliability.kind != ReliabilityQosPolicyKind::RELIABLE_RELIABILITY_QOS)
                return ReturnCode_t::RETCODE_OK;

            // Target: all samples written so far have been acknowledged.
            // last_acked_sn_ is advanced by process_acknack() as readers confirm receipt.
            const uint64_t target =
                next_sequence_number_.low > 1 ? next_sequence_number_.low - 1 : 0;
            if (last_acked_sn_ >= target)
                return ReturnCode_t::RETCODE_OK;

            const auto duration = std::chrono::seconds(max_wait.sec)
                + std::chrono::nanoseconds(max_wait.nanosec);
            std::unique_lock<std::mutex> lock(mutex_);
            const bool acked = ack_cv_.wait_for(lock, duration,
                [this, target] { return last_acked_sn_ >= target; });
            return acked ? ReturnCode_t::RETCODE_OK : ReturnCode_t::RETCODE_TIMEOUT;
        }

        // Accessors
        Topic *topic() const { return topic_; }
        Topic *get_topic() const { return topic_; }  // PSM alias
        Publisher *publisher() const { return publisher_; }
        const rtps::GUID_t &guid() const { return guid_; }
        void set_guid(const rtps::GUID_t &guid) { guid_ = guid; }

        void set_qos(const DataWriterQos &qos);
        const DataWriterQos &get_qos() const { return qos_; }

        const std::deque<CacheChange> &history_cache() const { return history_cache_; }

        // Reliability support
        void add_matched_reader(const rtps::GUID_t &reader_guid);
        void remove_matched_reader(const rtps::GUID_t &reader_guid);
        void note_offered_incompatible_qos(QosPolicyId_t failed_policy_id);
        void process_acknack(const rtps::GUID_t &reader_guid,
                             const std::set<uint32_t> &requested_sns,
                             rtps::Count_t count);

        // Transport callback - set by DomainParticipant
        using SendCallback = std::function<void(const std::vector<uint8_t> &)>;
        void set_send_callback(SendCallback callback) { send_callback_ = callback; }

        // Dual-encoding transport callback (XCDR1 + XCDR2)
        using SendDualCallback = std::function<void(const std::vector<uint8_t> &,
                                                     const std::vector<uint8_t> &)>;
        void set_send_dual_callback(SendDualCallback callback) { send_dual_callback_ = callback; }

        // Phase 3.d — RTPS 2.5 §8.7.6 coherent-set send callback.
        // Publisher::create_datawriter binds this to
        // RtpsUdpTransport::publish_coherent().  @p in_coherent_set_start_sn
        // is 0 for the FIRST sample of a set (the transport allocates the
        // writer SN and stores it back via @p out_coherent_set_start_sn),
        // otherwise the previously-captured start SN is used verbatim.
        // @p end_of_coherent_set emits PID_END_COHERENT_SET on the last
        // sample so peer AstuteDDS readers flush the buffered set on
        // delivery of this sample.
        using SendCoherentCallback = std::function<void(
            const std::vector<uint8_t> &xcdr1Data,
            const std::vector<uint8_t> &xcdr2Data,
            uint64_t  in_coherent_set_start_sn,
            bool      end_of_coherent_set,
            uint64_t &out_coherent_set_start_sn)>;
        void set_send_coherent_callback(SendCoherentCallback callback)
        {
            send_coherent_callback_ = std::move(callback);
        }

        // Dispose transport callback — emits DATA(K) with PID_STATUS_INFO.
        // The uint32 parameter is the StatusInfo_t bit set (RTPS §9.6.3.9):
        // 0x00000001 = DISPOSED, 0x00000002 = UNREGISTERED,
        // 0x00000003 = DISPOSED | UNREGISTERED.  Set by
        // Publisher::create_datawriter to route via
        // DomainParticipant::send_dispose().  The trailing sec/nsec
        // override lets `_w_timestamp` variants (DDS §2.2.2.4.2.14/17)
        // pin the wire INFO_TS.  A negative seconds override means
        // "use current wall clock".
        using SendDisposeCallback =
            std::function<void(const std::vector<uint8_t> &xcdr1Data,
                                const std::vector<uint8_t> &xcdr2Data,
                                uint32_t statusInfoBits,
                                int64_t tsSecondsOverride,
                                uint32_t tsFractionOverride)>;
        void set_send_dispose_callback(SendDisposeCallback callback)
        {
            send_dispose_callback_ = std::move(callback);
        }

        // Intra-participant loopback callback — invoked from write() alongside
        // send_callback_ so a DataReader on the same DomainParticipant receives
        // this writer's samples without going through the RTPS/UDP path
        // (which filters out self-originated packets).  Set by
        // Publisher::create_datawriter to route via DomainParticipant::local_deliver().
        using LocalDeliverCallback =
            std::function<void(const std::vector<uint8_t> &data,
                                const rtps::GUID_t &writer_guid,
                                const rtps::SequenceNumber_t &sequence_number)>;
        void set_local_deliver_callback(LocalDeliverCallback callback)
        {
            local_deliver_callback_ = std::move(callback);
        }

        // Phase 3.d — coherent-set-aware local loopback used by write_coherent().
        using LocalDeliverCoherentCallback =
            std::function<void(const std::vector<uint8_t> &data,
                                const rtps::GUID_t &writer_guid,
                                const rtps::SequenceNumber_t &sequence_number,
                                uint64_t coherent_set_start_sn,
                                bool     end_of_coherent_set)>;
        void set_local_deliver_coherent_callback(LocalDeliverCoherentCallback callback)
        {
            local_deliver_coherent_callback_ = std::move(callback);
        }

        // Heartbeat callback — invoked by send_heartbeat(); set by DomainParticipant
        // to trigger an out-of-band HEARTBEAT submessage via the transport.
        using HeartbeatCallback = std::function<void()>;
        void set_heartbeat_callback(HeartbeatCallback callback) { heartbeat_callback_ = std::move(callback); }

        // DataWriter listener
        void set_listener(DataWriterListener *listener, StatusMask mask = STATUS_MASK_ALL)
        {
            listener_ = listener;
            listener_mask_ = mask;
        }
        DataWriterListener *get_listener() const { return listener_; }
        StatusMask get_listener_mask() const { return listener_mask_; }

        // Status accessors (DDS §2.2.4.1).  Each getter atomically returns the
        // current status snapshot and resets the corresponding `*_change` fields
        // to zero so subsequent reads only report deltas since the last call.
        void get_publication_matched_status(PublicationMatchedStatus &out);
        void get_offered_incompatible_qos_status(OfferedIncompatibleQosStatus &out);

        LivelinessLostStatus get_liveliness_lost_status() const;

        bool assert_liveliness();

        using LivelinessAssertHook = std::function<void(const rtps::GUID_t&)>;
        void set_liveliness_assert_hook(LivelinessAssertHook hook)
        {
            liveliness_assert_hook_ = std::move(hook);
        }
        using LivelinessLostStatusHook = std::function<LivelinessLostStatus()>;
        void set_liveliness_lost_hook(LivelinessLostStatusHook hook)
        {
            liveliness_lost_hook_ = std::move(hook);
        }

        OfferedDeadlineMissedStatus get_offered_deadline_missed_status();

        void set_deadline_timer_service(TimerService *ts)
        {
            deadline_timer_service_ = ts;
        }

        // Deadline monitoring (per DDS §2.2.3.15 — scoped to each keyed
        // instance).  start_deadline_timer() arms a default HANDLE_NIL bucket
        // so silent writers still report offered-deadline-missed; per-instance
        // buckets are created lazily by notify_write() as new keys are seen.
        void start_deadline_timer();
        void stop_deadline_timer();
        void notify_write(InstanceHandle_t handle);  // Reset/arm per-instance timer

        // Lifespan: remove expired samples from cache
        void remove_expired_samples();

        // Entity framework (DDS §2.2.2.1)
        ReturnCode_t enable() noexcept { enabled_ = true; return ReturnCode_t::RETCODE_OK; }
        bool is_enabled() const noexcept { return enabled_; }
        InstanceHandle_t get_instance_handle() const noexcept { return instance_handle_; }
        StatusMask get_status_changes() const noexcept { return status_condition_.get_status_changes(); }
        EntityStatusCondition* get_statuscondition() noexcept { return &status_condition_; }

        // Discovery / introspection (DDS §2.2.2.4.1.9-10)
        ReturnCode_t get_matched_subscriptions(std::vector<InstanceHandle_t>& handles) const;

        ReturnCode_t get_matched_subscription_data(SubscriptionBuiltinTopicData& data,
                                                    InstanceHandle_t handle) const;

    private:
        void send_heartbeat();
        void retransmit_samples(const std::set<uint32_t> &sequence_numbers);
        bool check_resource_limits();
        void apply_history_and_limits();
        // Fires from the TimerService worker when the deadline period elapses
        // without a write() for the given instance.  Bumps the offered-
        // deadline-missed counters (populating last_instance_handle),
        // dispatches the writer / participant listener, then reschedules.
        void on_deadline_missed_(InstanceHandle_t handle);

        Topic *topic_;
        Publisher *publisher_;
        rtps::GUID_t guid_;
        rtps::SequenceNumber_t next_sequence_number_{0, 1};
        std::deque<CacheChange> history_cache_;
        DataWriterQos qos_;
        // Entity framework state (DDS §2.2.2.1)
        bool enabled_{false};
        InstanceHandle_t instance_handle_{next_entity_instance_handle()};
        EntityStatusCondition status_condition_{};
        mutable std::mutex mutex_;
        SendCallback send_callback_;
        SendDualCallback send_dual_callback_;
        SendDisposeCallback send_dispose_callback_;
        HeartbeatCallback heartbeat_callback_;
        LocalDeliverCallback local_deliver_callback_;
        // Phase 3.d — RTPS §8.7.6 coherent-set send / loopback callbacks.
        SendCoherentCallback send_coherent_callback_;
        LocalDeliverCoherentCallback local_deliver_coherent_callback_;

        // §2.2.2.4.2.5 / §2.2.2.4.2.13 — per-instance key bytes cached from
        // register_instance(sample) or the most recent write() for that
        // handle.  dispose(handle) / unregister_instance(handle) look up
        // the payload here so they can emit a keyed DATA(K).  Both
        // encodings are stored so dispose_by_handle can pick the right
        // one for each remote reader.
        std::map<InstanceHandle_t, std::vector<uint8_t>> instance_key_payloads_xcdr1_;
        std::map<InstanceHandle_t, std::vector<uint8_t>> instance_key_payloads_xcdr2_;

        // Matched reader tracking (for wait_for_acknowledgments and status reporting)
        std::set<rtps::GUID_t> matched_readers_;

        // Acknowledgment tracking: last_acked_sn_ is set to the highest SN
        // confirmed by all matched readers; ack_cv_ is signaled on each update.
        uint64_t last_acked_sn_{0};
        std::condition_variable ack_cv_;

        // DataWriter listener
        DataWriterListener *listener_{nullptr};
        StatusMask listener_mask_{STATUS_MASK_ALL};

        // Persistent status counters (drained by getters per DDS §2.2.4.1).
        PublicationMatchedStatus     publication_matched_status_{};
        OfferedIncompatibleQosStatus offered_incompatible_qos_status_{};
        // §1.1b-B — accumulated across deadline-missed events until read.
        OfferedDeadlineMissedStatus  offered_deadline_missed_status_{};

        // §1.1b-B liveliness hooks — set by Publisher::create_datawriter to
        // route assertion and status queries through the participant's
        // LivelinessManager.  Left null when no participant is attached
        // (e.g. unit tests that construct DataWriters directly).
        LivelinessAssertHook     liveliness_assert_hook_;
        LivelinessLostStatusHook liveliness_lost_hook_;

        // §1.1b-B deadline watchdog on TimerService — replaces the old
        // dedicated deadline_thread_ / cv / running_ trio.  The
        // TimerService is owned by the participant and shared across
        // all writers.  §1.1b-B item 10 — one TimerId per keyed
        // instance (DDS §2.2.3.15); HANDLE_NIL acts as the default
        // bucket for unkeyed topics and for writers armed before any
        // write() has produced an instance handle.
        TimerService *deadline_timer_service_{nullptr};
        std::map<InstanceHandle_t, TimerId> deadline_timers_;

        // Lifespan tracking — each cache entry gets a creation timestamp
        std::deque<std::chrono::steady_clock::time_point> cache_timestamps_;
    };

    class ReadCondition
    {
    public:
        ReadCondition(DataReader *reader,
                      SampleStateMask sample_states,
                      ViewStateMask view_states,
                      InstanceStateMask instance_states)
            : reader_(reader),
              sample_state_mask_(sample_states),
              view_state_mask_(view_states),
              instance_state_mask_(instance_states) {}

        virtual ~ReadCondition() = default;
        ReadCondition(const ReadCondition &) = delete;
        ReadCondition &operator=(const ReadCondition &) = delete;

        SampleStateMask   get_sample_state_mask()   const noexcept { return sample_state_mask_; }
        ViewStateMask     get_view_state_mask()     const noexcept { return view_state_mask_; }
        InstanceStateMask get_instance_state_mask() const noexcept { return instance_state_mask_; }
        DataReader       *get_datareader()          const noexcept { return reader_; }

        bool trigger_value() const;

    private:
        DataReader       *reader_{nullptr};
        SampleStateMask   sample_state_mask_{ANY_SAMPLE_STATE};
        ViewStateMask     view_state_mask_{ANY_VIEW_STATE};
        InstanceStateMask instance_state_mask_{ANY_INSTANCE_STATE};
    };

    class DataReader
    {
    public:
        DataReader(Topic *topic, Subscriber *subscriber);
        virtual ~DataReader();
        DataReader(const DataReader &) = delete;
        DataReader &operator=(const DataReader &) = delete;

        std::vector<ReceivedSample> read(size_t max_samples = 0xFFFFFFFF);

        std::vector<ReceivedSample> take(size_t max_samples = 0xFFFFFFFF);

        ReturnCode_t read_next_sample(std::vector<uint8_t> &data, SampleInfo &info);

        ReturnCode_t take_next_sample(std::vector<uint8_t> &data, SampleInfo &info);

        void add_sample(const ReceivedSample &sample);

        ReturnCode_t get_key_value(std::vector<uint8_t> &key_holder, InstanceHandle_t handle);

        //---------------------------------------------------------------------
        // Phase 2.8 — DDS 1.4 §2.2.2.5.3.8 / §2.2.2.5.3.9 / §2.2.2.5.3.20
        // Classic-PSM loan-based read / take / return_loan.
        //
        // Semantics per spec:
        //   - If both `data_values` and `sample_infos` are empty on entry,
        //     the reader grants a loan.  The caller MUST invoke
        //     return_loan(data_values, sample_infos) before invoking any
        //     other read/take on the same pair.
        //   - If either sequence is non-empty on entry, no loan is granted
        //     (application-owned storage).
        //   - `max_samples` == -1 (LENGTH_UNLIMITED) means "no limit".
        //   - Populates each `SampleInfo` fully — including generation
        //     counts (Phase 2.9) and sample_rank / generation_rank per
        //     §2.2.2.5.1.6.
        //---------------------------------------------------------------------
        ReturnCode_t read(std::vector<std::vector<uint8_t>> &data_values,
                          std::vector<SampleInfo>            &sample_infos,
                          int32_t                             max_samples = -1,
                          SampleStateMask                     sample_states   = ANY_SAMPLE_STATE,
                          ViewStateMask                       view_states     = ANY_VIEW_STATE,
                          InstanceStateMask                   instance_states = ANY_INSTANCE_STATE);
        ReturnCode_t take(std::vector<std::vector<uint8_t>> &data_values,
                          std::vector<SampleInfo>            &sample_infos,
                          int32_t                             max_samples = -1,
                          SampleStateMask                     sample_states   = ANY_SAMPLE_STATE,
                          ViewStateMask                       view_states     = ANY_VIEW_STATE,
                          InstanceStateMask                   instance_states = ANY_INSTANCE_STATE);

        ReturnCode_t return_loan(std::vector<std::vector<uint8_t>> &data_values,
                                 std::vector<SampleInfo>            &sample_infos);

        // Accessors
        Topic *topic() const { return topic_; }
        TopicDescription *get_topicdescription() const;  // PSM: may be CFT
        Subscriber *subscriber() const { return subscriber_; }
        const rtps::GUID_t &guid() const { return guid_; }
        void set_guid(const rtps::GUID_t &guid) { guid_ = guid; }

        void set_qos(const DataReaderQos &qos);
        const DataReaderQos &get_qos() const { return qos_; }

        size_t unread_count() const;

        uint64_t latency_budget_exceeded_count() const;

        uint64_t autopurged_instances_count() const;

        size_t run_autopurge_sweep();

        // Reliability support
        void add_matched_writer(const rtps::GUID_t &writer_guid);
        void remove_matched_writer(const rtps::GUID_t &writer_guid);

        bool is_writer_matched(const rtps::GUID_t &writer_guid) const;
        void note_requested_incompatible_qos(QosPolicyId_t failed_policy_id);
        void process_heartbeat(const rtps::GUID_t &writer_guid,
                               const rtps::SequenceNumber_t &first_sn,
                               const rtps::SequenceNumber_t &last_sn,
                               rtps::Count_t count,
                               bool final);

        // Data received callback
        using DataCallback = std::function<void(const ReceivedSample &)>;
        void set_data_callback(DataCallback callback) { data_callback_ = callback; }

        // Listener support (DDS §2.1.4.3 — DataReaderListener dispatch)
        void set_listener(DataReaderListener* listener) { listener_ = listener; }
        DataReaderListener* get_listener() const { return listener_; }

        // Status accessors (DDS §2.2.4.1).  See DataWriter counterparts for
        // semantics; each call drains the `*_change` fields.
        void get_subscription_matched_status(SubscriptionMatchedStatus &out);
        void get_requested_incompatible_qos_status(RequestedIncompatibleQosStatus &out);

        LivelinessChangedStatus get_liveliness_changed_status() const;

        using LivelinessChangedStatusHook = std::function<LivelinessChangedStatus()>;
        void set_liveliness_changed_hook(LivelinessChangedStatusHook hook)
        {
            liveliness_changed_hook_ = std::move(hook);
        }

        // WaitSet notification — external condition variables to signal on data arrival
        using NotifyCallback = std::function<void()>;
        void add_notify_callback(NotifyCallback cb)
        {
            std::lock_guard<std::mutex> lock(mutex_);
            notify_callbacks_.push_back(std::move(cb));
        }

        // Destruction-time callback — external owners (e.g. the C API
        // eventfd side-map) register a hook that fires from ~DataReader()
        // so per-reader resources can be cleaned up deterministically
        // even when the reader is deleted through its parent Subscriber
        // rather than through the C API delete function.  Callbacks are
        // invoked in registration order with no lock held.
        using DestroyCallback = std::function<void()>;
        void add_destroy_callback(DestroyCallback cb)
        {
            std::lock_guard<std::mutex> lock(mutex_);
            destroy_callbacks_.push_back(std::move(cb));
        }

        // Topic description override (used when created via ContentFilteredTopic)
        void set_topicdescription(TopicDescription *td) { topicdescription_ = td; }

        RequestedDeadlineMissedStatus get_requested_deadline_missed_status();

        SampleLostStatus get_sample_lost_status();

        SampleRejectedStatus get_sample_rejected_status();

        void bump_sample_lost(int32_t count = 1);

        void set_deadline_timer_service(TimerService *ts)
        {
            deadline_timer_service_ = ts;
        }

        // Deadline monitoring (per DDS §2.2.3.15 — scoped to each keyed
        // instance).  start_deadline_timer() arms a default HANDLE_NIL bucket
        // for readers with no matched samples yet; per-instance buckets are
        // created lazily by notify_receive() as new keys arrive.
        void start_deadline_timer();
        void stop_deadline_timer();
        void notify_receive(InstanceHandle_t handle);  // Reset/arm per-instance timer

        // Entity framework (DDS §2.2.2.1)
        ReturnCode_t enable() noexcept { enabled_ = true; return ReturnCode_t::RETCODE_OK; }
        bool is_enabled() const noexcept { return enabled_; }
        InstanceHandle_t get_instance_handle() const noexcept { return instance_handle_; }
        StatusMask get_status_changes() const noexcept { return status_condition_.get_status_changes(); }
        EntityStatusCondition* get_entity_statuscondition() noexcept { return &status_condition_; }

        // Discovery / introspection (DDS §2.2.2.5.1.14-15)
        ReturnCode_t get_matched_publications(std::vector<InstanceHandle_t>& handles) const;

        ReturnCode_t get_matched_publication_data(PublicationBuiltinTopicData& data,
                                                    InstanceHandle_t handle) const;

        //---------------------------------------------------------------------
        // DDS 1.4 §2.2.2.5 — Phase 2.d DataReader completion.
        // ReadCondition factory, mask/condition-filtered read/take, per-
        // instance / next-instance read/take, lookup_instance, contained-
        // entity cleanup, historical-data wait.
        //---------------------------------------------------------------------

        ReadCondition *create_readcondition(SampleStateMask sample_states,
                                            ViewStateMask view_states,
                                            InstanceStateMask instance_states);

        ReturnCode_t delete_readcondition(ReadCondition *condition);

        //---------------------------------------------------------------------
        // Mask-filtered read/take variants (§2.2.2.5.3.8 / .9 / .16 / .17).
        // All variants share the same matcher: a sample is delivered iff
        // (sample.state & sample_state_mask) != 0
        //   && (sample.instance_state & instance_state_mask) != 0
        //   && the view_state is in @p view_state_mask.  view_state is
        // currently reported as NEW_VIEW_STATE for every sample (Phase
        // 2.d does not yet plumb per-instance view-state tracking).
        //---------------------------------------------------------------------

        std::vector<ReceivedSample> read(size_t max_samples,
                                         SampleStateMask sample_states,
                                         ViewStateMask view_states,
                                         InstanceStateMask instance_states);
        std::vector<ReceivedSample> take(size_t max_samples,
                                         SampleStateMask sample_states,
                                         ViewStateMask view_states,
                                         InstanceStateMask instance_states);

        std::vector<ReceivedSample> read_w_condition(ReadCondition *condition,
                                                     size_t max_samples = 0xFFFFFFFF);
        std::vector<ReceivedSample> take_w_condition(ReadCondition *condition,
                                                     size_t max_samples = 0xFFFFFFFF);

        std::vector<ReceivedSample> read_instance(InstanceHandle_t a_handle,
                                                   size_t max_samples = 0xFFFFFFFF,
                                                   SampleStateMask sample_states = ANY_SAMPLE_STATE,
                                                   ViewStateMask view_states = ANY_VIEW_STATE,
                                                   InstanceStateMask instance_states = ANY_INSTANCE_STATE);
        std::vector<ReceivedSample> take_instance(InstanceHandle_t a_handle,
                                                   size_t max_samples = 0xFFFFFFFF,
                                                   SampleStateMask sample_states = ANY_SAMPLE_STATE,
                                                   ViewStateMask view_states = ANY_VIEW_STATE,
                                                   InstanceStateMask instance_states = ANY_INSTANCE_STATE);
        std::vector<ReceivedSample> read_next_instance(InstanceHandle_t previous_handle,
                                                        size_t max_samples = 0xFFFFFFFF,
                                                        SampleStateMask sample_states = ANY_SAMPLE_STATE,
                                                        ViewStateMask view_states = ANY_VIEW_STATE,
                                                        InstanceStateMask instance_states = ANY_INSTANCE_STATE);
        std::vector<ReceivedSample> take_next_instance(InstanceHandle_t previous_handle,
                                                        size_t max_samples = 0xFFFFFFFF,
                                                        SampleStateMask sample_states = ANY_SAMPLE_STATE,
                                                        ViewStateMask view_states = ANY_VIEW_STATE,
                                                        InstanceStateMask instance_states = ANY_INSTANCE_STATE);

        std::vector<ReceivedSample> read_instance_w_condition(InstanceHandle_t a_handle,
                                                              ReadCondition *condition,
                                                              size_t max_samples = 0xFFFFFFFF);
        std::vector<ReceivedSample> take_instance_w_condition(InstanceHandle_t a_handle,
                                                              ReadCondition *condition,
                                                              size_t max_samples = 0xFFFFFFFF);
        std::vector<ReceivedSample> read_next_instance_w_condition(InstanceHandle_t previous_handle,
                                                                    ReadCondition *condition,
                                                                    size_t max_samples = 0xFFFFFFFF);
        std::vector<ReceivedSample> take_next_instance_w_condition(InstanceHandle_t previous_handle,
                                                                    ReadCondition *condition,
                                                                    size_t max_samples = 0xFFFFFFFF);

        InstanceHandle_t lookup_instance(const std::vector<uint8_t> &key_holder) const;

        ReturnCode_t delete_contained_entities();

        ReturnCode_t wait_for_historical_data(const Duration_t &max_wait);

        //---------------------------------------------------------------------
        // ISO C++ PSM §7.13.2 — Selector builder.
        //
        // A fluent builder that composes the mask triplet + max_samples
        // + optional instance handle into a single chainable expression
        // culminating in .read() / .take().  Purely a convenience layer
        // over the existing read()/take()/read_instance()/take_instance()
        // /read_next_instance()/take_next_instance() overloads.
        //
        // Example:
        //   auto samples = reader->select()
        //                      .max_samples(16)
        //                      .state(NOT_READ_SAMPLE_STATE,
        //                             ANY_VIEW_STATE,
        //                             ALIVE_INSTANCE_STATE)
        //                      .take();
        //---------------------------------------------------------------------
        class Selector
        {
        public:
            explicit Selector(DataReader *reader) noexcept : reader_(reader) {}

            Selector &max_samples(size_t n) noexcept
            {
                max_samples_ = n;
                return *this;
            }

            Selector &state(SampleStateMask sample_states, ViewStateMask view_states,
                            InstanceStateMask instance_states) noexcept
            {
                sample_states_   = sample_states;
                view_states_     = view_states;
                instance_states_ = instance_states;
                return *this;
            }

            Selector &sample_states(SampleStateMask m) noexcept
            {
                sample_states_ = m;
                return *this;
            }
            Selector &view_states(ViewStateMask m) noexcept
            {
                view_states_ = m;
                return *this;
            }
            Selector &instance_states(InstanceStateMask m) noexcept
            {
                instance_states_ = m;
                return *this;
            }

            Selector &instance(InstanceHandle_t h) noexcept
            {
                instance_       = h;
                mode_           = Mode::INSTANCE;
                return *this;
            }

            Selector &next_instance(InstanceHandle_t previous_handle) noexcept
            {
                instance_       = previous_handle;
                mode_           = Mode::NEXT_INSTANCE;
                return *this;
            }

            std::vector<ReceivedSample> read()
            {
                switch (mode_)
                {
                case Mode::ALL:
                    return reader_->read(max_samples_, sample_states_, view_states_,
                                         instance_states_);
                case Mode::INSTANCE:
                    return reader_->read_instance(instance_, max_samples_, sample_states_,
                                                  view_states_, instance_states_);
                case Mode::NEXT_INSTANCE:
                    return reader_->read_next_instance(instance_, max_samples_,
                                                       sample_states_, view_states_,
                                                       instance_states_);
                }
                return {};
            }

            std::vector<ReceivedSample> take()
            {
                switch (mode_)
                {
                case Mode::ALL:
                    return reader_->take(max_samples_, sample_states_, view_states_,
                                         instance_states_);
                case Mode::INSTANCE:
                    return reader_->take_instance(instance_, max_samples_, sample_states_,
                                                  view_states_, instance_states_);
                case Mode::NEXT_INSTANCE:
                    return reader_->take_next_instance(instance_, max_samples_,
                                                       sample_states_, view_states_,
                                                       instance_states_);
                }
                return {};
            }

        private:
            enum class Mode
            {
                ALL,
                INSTANCE,
                NEXT_INSTANCE
            };

            DataReader       *reader_{nullptr};
            size_t            max_samples_{0xFFFFFFFF};
            SampleStateMask   sample_states_{ANY_SAMPLE_STATE};
            ViewStateMask     view_states_{ANY_VIEW_STATE};
            InstanceStateMask instance_states_{ANY_INSTANCE_STATE};
            InstanceHandle_t  instance_{HANDLE_NIL};
            Mode              mode_{Mode::ALL};
        };

        Selector select() noexcept { return Selector(this); }

    protected:
        TopicDescription *topicdescription_{nullptr};  // non-null when created via CFT

    private:
        void send_acknack(const rtps::GUID_t &writer_guid);
        bool check_resource_limits() const;
        void apply_history_and_limits();
        // Phase 3.d — RTPS 2.5 §8.7.6 PID_COHERENT_SET receive path.
        // Body of the historical add_sample(): performs duplicate
        // detection, EXCLUSIVE-ownership filtering, TIME_BASED_FILTER,
        // history-cache insertion, listener dispatch.  Must be called
        // with mutex_ held.  The wrapping add_sample() takes mutex_
        // once, runs autopurge, applies the coherent-set gate, and
        // then either buffers the sample (mid-set) or forwards it
        // (and any pending set from the same writer) to this helper.
        void admit_sample_locked(const ReceivedSample &sample);
        // Phase 2.11.e — DDS §2.2.3.19 READER_DATA_LIFECYCLE.  Sweeps
        // instance_generation_ and evicts any NOT_ALIVE_* instance
        // whose autopurge_{disposed,nowriter}_samples_delay has
        // elapsed since state_change_time_ns.  Called with mutex_ held
        // (add_sample() / read() / take() etc. hold it) and returns
        // the number of instances purged this call.  A no-op when
        // both delays are INFINITE (default).
        size_t purge_expired_instances_locked();
        // §1.1b-B — fires from the TimerService worker when the reader
        // has not received a matching sample for the given instance within
        // its DEADLINE period.  Bumps the requested-deadline-missed counters
        // (populating last_instance_handle), dispatches the reader /
        // participant listener, then reschedules.
        void on_deadline_missed_(InstanceHandle_t handle);

        // Phase 2.9 (DDS §2.2.2.5.1.5 / §2.2.2.5.1.6 / §2.2.2.5.1.8).
        // Populates SampleInfo.disposed_generation_count,
        // no_writers_generation_count, sample_rank, generation_rank,
        // absolute_generation_rank, sample_state and view_state from
        // the given cached sample.  sample_rank is provided by the
        // caller (0 for the single-sample accessors).  Must be called
        // with mutex_ held.
        void populate_generation_and_ranks(const ReceivedSample &sample,
                                           SampleInfo &info,
                                           int32_t sample_rank) const;
        // Phase 2.9 — clears the per-instance view_new flag so that
        // subsequent reads report NOT_NEW_VIEW_STATE until the instance
        // becomes reborn.  Must be called with mutex_ held.
        void mark_instance_seen(InstanceHandle_t handle);

        Topic *topic_;
        Subscriber *subscriber_;
        rtps::GUID_t guid_;
        std::deque<ReceivedSample> history_cache_;
        DataReaderQos qos_;
        // Entity framework state (DDS §2.2.2.1)
        bool enabled_{false};
        InstanceHandle_t instance_handle_{next_entity_instance_handle()};
        EntityStatusCondition status_condition_{};
        mutable std::mutex mutex_;
        DataCallback data_callback_;
        DataReaderListener* listener_{nullptr};
        std::vector<NotifyCallback> notify_callbacks_;
        std::vector<DestroyCallback> destroy_callbacks_;

        // Persistent status counters (drained by getters per DDS §2.2.4.1).
        SubscriptionMatchedStatus      subscription_matched_status_{};
        RequestedIncompatibleQosStatus requested_incompatible_qos_status_{};
        // §1.1b-B — accumulated across deadline-missed events until read.
        RequestedDeadlineMissedStatus  requested_deadline_missed_status_{};
        // §1.1b-B — accumulated across sample-lost events until read.
        SampleLostStatus               sample_lost_status_{};
        // Phase 2.f — accumulated across sample-rejected events until read
        // (DDS §2.2.4.1.8 / §2.2.2.5.3.25).  Bumped when the reader's
        // cache is at RESOURCE_LIMITS.max_samples.
        SampleRejectedStatus           sample_rejected_status_{};

        // Phase 2.11.c — DDS §2.2.3.9 LATENCY_BUDGET diagnostic counter.
        // Bumped by add_sample() when the reader's latency_budget is
        // non-zero and the observed transit time (reception - source)
        // exceeds the budget.  Surfaced via
        // latency_budget_exceeded_count(); monotonic (not drained on
        // read).
        uint64_t latency_budget_exceeded_count_{0};

        // Phase 2.11.e — DDS §2.2.3.19 READER_DATA_LIFECYCLE diagnostic
        // counter.  Bumped once per instance that
        // purge_expired_instances_locked() evicts from
        // instance_generation_ (and whose cached samples are dropped)
        // because autopurge_{disposed,nowriter}_samples_delay has
        // elapsed.  Surfaced via autopurged_instances_count();
        // monotonic (not drained on read).
        uint64_t autopurged_instances_count_{0};

        // Phase 3.d — RTPS 2.5 §8.7.6 coherent-set receive buffer.
        // Samples carrying PID_COHERENT_SET are queued here per source
        // writer GUID until the set boundary is reached (either the
        // matching PID_END_COHERENT_SET marker on the last sample or
        // the arrival of a subsequent sample from the same writer with
        // a different / absent coherent-set tag).  On flush the whole
        // queue is admitted into history_cache_ under a single lock so
        // read() / take() cannot observe a partial set (DDS §2.2.3.6
        // PRESENTATION.coherent_access).
        std::map<rtps::GUID_t, std::vector<ReceivedSample>> coherent_pending_;

        // §1.1b-B liveliness hook — set by Subscriber::create_datareader.
        // Snapshot-copied from the participant's LivelinessManager on each
        // get_liveliness_changed_status() call.
        LivelinessChangedStatusHook liveliness_changed_hook_;

        // Matched writer tracking (for process_heartbeat and send_acknack)
        std::set<rtps::GUID_t> matched_writers_;

        // ACKNACK callback — invoked by send_acknack(); set by DomainParticipant
        // to deliver an ACKNACK via the transport for the given writer.
        using AcknackCallback = std::function<void(const rtps::GUID_t &)>;
        AcknackCallback acknack_callback_;

    public:
        void set_acknack_callback(AcknackCallback cb) { acknack_callback_ = std::move(cb); }

    private:

        // Time-based filter: track last accepted time per writer
        std::map<rtps::GUID_t, std::chrono::steady_clock::time_point> last_accepted_time_;

        // Ownership: per-instance tracking of current owner (EXCLUSIVE)
        struct InstanceOwnership {
            rtps::GUID_t ownerGuid{};
            int32_t ownerStrength{0};
        };
        std::map<int32_t, InstanceOwnership> instance_owners_;  // key = instance_handle hash

        // Instance registry: maps instance_handle → serialized key bytes (DDS §2.2.2.5.3)
        std::map<InstanceHandle_t, std::vector<uint8_t>> instance_key_data_;

        // DDS §2.2.2.5.1.5 — per-instance generation-count state.  For
        // every instance ever seen by this reader the middleware keeps
        // running counters that increment on NOT_ALIVE_DISPOSED→ALIVE
        // and NOT_ALIVE_NO_WRITERS→ALIVE transitions.  Snapshotted onto
        // each admitted ReceivedSample and copied out to SampleInfo on
        // read/take.  view_state is derived: NEW on first admit, becomes
        // NOT_NEW after the sample is delivered by read/take, and
        // returns to NEW on a re-alive transition.
        struct InstanceGenerationState
        {
            int32_t disposed_generation_count{0};
            int32_t no_writers_generation_count{0};
            InstanceStateKind current_state{ALIVE_INSTANCE_STATE};
            // True until at least one sample for this instance has been
            // surfaced by read()/take()/*next_sample() OR the instance
            // has been reborn since the last surface event.
            bool view_new{true};
            // Phase 2.11.e — DDS §2.2.3.19 READER_DATA_LIFECYCLE.
            // Wall-clock nanoseconds since epoch stamped when
            // current_state transitions into NOT_ALIVE_DISPOSED_INSTANCE_STATE
            // or NOT_ALIVE_NO_WRITERS_INSTANCE_STATE.  Zero while the
            // instance is ALIVE.  Used by purge_expired_instances_locked()
            // to decide whether autopurge_{disposed,nowriter}_samples_delay
            // has elapsed.
            int64_t state_change_time_ns{0};
        };
        std::map<InstanceHandle_t, InstanceGenerationState> instance_generation_;

        // Phase 2.8 (DDS §2.2.2.5.3.8 / .9 / .20) — outstanding loan
        // bookkeeping.  When a loan-granting read()/take() empties into
        // caller-supplied sequences, the (data_values, sample_infos)
        // address pair is inserted here so return_loan() can validate
        // that the pair was in fact granted by this reader.
        std::set<std::pair<const void *, const void *>> outstanding_loans_;

        // §1.1b-B deadline watchdog on TimerService — replaces the old
        // dedicated deadline_thread_ / cv / running_ trio.  §1.1b-B item 10 —
        // one TimerId per keyed instance (DDS §2.2.3.15 / §2.2.4.1.5
        // last_instance_handle); HANDLE_NIL acts as the default bucket for
        // unkeyed topics and for readers armed before any sample has arrived.
        TimerService *deadline_timer_service_{nullptr};
        std::map<InstanceHandle_t, TimerId> deadline_timers_;

        // §2.2.2.5.2 — ReadConditions owned by this reader.  Populated
        // by create_readcondition(), drained by delete_readcondition() /
        // delete_contained_entities() and the destructor.  Held under
        // mutex_ for iteration.
        std::vector<std::unique_ptr<ReadCondition>> read_conditions_;
    };

    class Publisher
    {
    public:
        Publisher(DomainParticipant *participant);
        ~Publisher() = default;
        Publisher(const Publisher &) = delete;
        Publisher &operator=(const Publisher &) = delete;

        DataWriter *create_datawriter(Topic *topic, const DataWriterQos &qos = DataWriterQos{});

        DataWriter *create_datawriter(Topic *topic, const DataWriterQos &qos,
                                      DataWriterListener *listener,
                                      StatusMask mask = STATUS_MASK_ALL)
        {
            DataWriter *writer = create_datawriter(topic, qos);
            if (writer && listener)
            {
                writer->set_listener(listener, mask);
            }
            return writer;
        }

        DataWriterQos default_datawriter_qos() const { return default_datawriter_qos_; }

        ReturnCode_t get_default_datawriter_qos(DataWriterQos &qos) const
        {
            qos = default_datawriter_qos_;
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t set_default_datawriter_qos(const DataWriterQos &qos)
        {
            default_datawriter_qos_ = qos;
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t copy_from_topic_qos(DataWriterQos &writer_qos,
                                         const TopicQos &topic_qos) const
        {
            writer_qos.durability          = topic_qos.durability;
            writer_qos.durability_service  = topic_qos.durability_service;
            writer_qos.deadline            = topic_qos.deadline;
            writer_qos.latency_budget      = topic_qos.latency_budget;
            writer_qos.liveliness          = topic_qos.liveliness;
            writer_qos.reliability         = topic_qos.reliability;
            writer_qos.destination_order   = topic_qos.destination_order;
            writer_qos.history             = topic_qos.history;
            writer_qos.resource_limits     = topic_qos.resource_limits;
            writer_qos.transport_priority  = topic_qos.transport_priority;
            writer_qos.lifespan            = topic_qos.lifespan;
            writer_qos.ownership           = topic_qos.ownership;
            writer_qos.representation      = topic_qos.representation;
            return ReturnCode_t::RETCODE_OK;
        }

        DataWriter *lookup_datawriter(const std::string &topic_name) const;

        bool delete_datawriter(DataWriter *writer);

        ReturnCode_t delete_contained_entities()
        {
            writers_.clear();
            coherent_buffer_.clear();
            coherent_active_ = false;
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t suspend_publications()
        {
            suspended_ = true;
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t resume_publications()
        {
            suspended_ = false;
            return ReturnCode_t::RETCODE_OK;
        }

        bool is_suspended() const noexcept { return suspended_; }

        ReturnCode_t wait_for_acknowledgments(const Duration_t &max_wait);

        ReturnCode_t begin_coherent_changes()
        {
            coherent_active_ = true;
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t end_coherent_changes();

        bool is_coherent_active() const { return coherent_active_; }

        void buffer_coherent_write(DataWriter *writer,
                                   const std::vector<uint8_t> &data);
        void buffer_coherent_write_dual(DataWriter *writer,
                                        const std::vector<uint8_t> &xcdr1,
                                        const std::vector<uint8_t> &xcdr2);

        DomainParticipant *participant() const { return participant_; }

        void set_qos(const PublisherQos &qos) { qos_ = qos; }
        const PublisherQos &get_qos() const { return qos_; }

        const std::vector<std::unique_ptr<DataWriter>> &writers() const { return writers_; }

        std::vector<DataWriter *> get_datawriters() const
        {
            std::vector<DataWriter *> result;
            result.reserve(writers_.size());
            for (const auto &w : writers_)
            {
                result.push_back(w.get());
            }
            return result;
        }

        // Entity framework (DDS §2.2.2.1)
        ReturnCode_t enable() noexcept { enabled_ = true; return ReturnCode_t::RETCODE_OK; }
        bool is_enabled() const noexcept { return enabled_; }
        InstanceHandle_t get_instance_handle() const noexcept { return instance_handle_; }
        StatusMask get_status_changes() const noexcept { return status_condition_.get_status_changes(); }
        EntityStatusCondition* get_statuscondition() noexcept { return &status_condition_; }

    private:
        DomainParticipant *participant_;
        PublisherQos qos_;
        std::vector<std::unique_ptr<DataWriter>> writers_;
        uint32_t next_writer_id_{1};

        // Entity framework state (DDS §2.2.2.1)
        bool enabled_{false};
        InstanceHandle_t instance_handle_{next_entity_instance_handle()};
        EntityStatusCondition status_condition_{};

        // §2.5 Publisher completion (DDS §2.2.2.4.1)
        DataWriterQos default_datawriter_qos_{};
        bool suspended_{false};

        // Coherent changes buffering
        bool coherent_active_{false};
        struct CoherentEntry {
            DataWriter *writer;
            std::vector<uint8_t> xcdr1;
            std::vector<uint8_t> xcdr2;
            bool dual;
        };
        std::vector<CoherentEntry> coherent_buffer_;
    };

    class Subscriber
    {
    public:
        Subscriber(DomainParticipant *participant);
        ~Subscriber() = default;
        Subscriber(const Subscriber &) = delete;
        Subscriber &operator=(const Subscriber &) = delete;

        DataReader *create_datareader(Topic *topic, const DataReaderQos &qos = DataReaderQos{});

        DataReader *create_datareader(TopicDescription *td, const DataReaderQos &qos,
                                      DataReaderListener *listener = nullptr,
                                      StatusMask mask = STATUS_MASK_ALL);

    private:
        DataReader *do_create_datareader(Topic *topic, const DataReaderQos &qos,
                                         ContentFilteredTopic *cft);

    public:

        DataReaderQos default_datareader_qos() const { return default_datareader_qos_; }

        ReturnCode_t get_default_datareader_qos(DataReaderQos &qos) const
        {
            qos = default_datareader_qos_;
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t set_default_datareader_qos(const DataReaderQos &qos)
        {
            default_datareader_qos_ = qos;
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t copy_from_topic_qos(DataReaderQos &reader_qos,
                                         const TopicQos &topic_qos) const
        {
            reader_qos.durability          = topic_qos.durability;
            reader_qos.deadline            = topic_qos.deadline;
            reader_qos.latency_budget      = topic_qos.latency_budget;
            reader_qos.liveliness          = topic_qos.liveliness;
            reader_qos.reliability         = topic_qos.reliability;
            reader_qos.destination_order   = topic_qos.destination_order;
            reader_qos.history             = topic_qos.history;
            reader_qos.resource_limits     = topic_qos.resource_limits;
            reader_qos.ownership           = topic_qos.ownership;
            reader_qos.representation      = topic_qos.representation;
            return ReturnCode_t::RETCODE_OK;
        }

        DataReader *lookup_datareader(const std::string &topic_name) const;

        bool delete_datareader(DataReader *reader);

        ReturnCode_t delete_contained_entities()
        {
            readers_.clear();
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t notify_datareaders();

        SampleLostStatus get_sample_lost_status();

        ReturnCode_t begin_access()
        {
            access_active_ = true;
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t end_access()
        {
            access_active_ = false;
            return ReturnCode_t::RETCODE_OK;
        }

        bool is_access_active() const { return access_active_; }

        DomainParticipant *participant() const { return participant_; }

        void set_qos(const SubscriberQos &qos) { qos_ = qos; }
        const SubscriberQos &get_qos() const { return qos_; }

        const std::vector<std::unique_ptr<DataReader>> &readers() const { return readers_; }

        std::vector<DataReader *> get_datareaders() const
        {
            std::vector<DataReader *> result;
            result.reserve(readers_.size());
            for (const auto &r : readers_)
            {
                result.push_back(r.get());
            }
            return result;
        }

        // Entity framework (DDS §2.2.2.1)
        ReturnCode_t enable() noexcept { enabled_ = true; return ReturnCode_t::RETCODE_OK; }
        bool is_enabled() const noexcept { return enabled_; }
        InstanceHandle_t get_instance_handle() const noexcept { return instance_handle_; }
        StatusMask get_status_changes() const noexcept { return status_condition_.get_status_changes(); }
        EntityStatusCondition* get_statuscondition() noexcept { return &status_condition_; }

    private:
        DomainParticipant *participant_;
        SubscriberQos qos_;
        std::vector<std::unique_ptr<DataReader>> readers_;
        uint32_t next_reader_id_{1};
        bool access_active_{false};

        // Entity framework state (DDS §2.2.2.1)
        bool enabled_{false};
        InstanceHandle_t instance_handle_{next_entity_instance_handle()};
        EntityStatusCondition status_condition_{};

        // §2.7 Subscriber completion (DDS §2.2.2.5.1)
        DataReaderQos default_datareader_qos_{};
    };

    // ── §1.4 Request / Reply (RPC) primitives ─────────────────────────
    //
    // The DDS-RPC pattern (OMG formal/2017-04-02) transports service calls
    // over two correlated DDS topics: `rq/<service>Request` and
    // `rr/<service>Reply`.  Every payload carries a 24-byte correlation
    // header prepended to the serialised user data:
    //
    //     bytes  0..15   16-byte requester GID   (identifies the caller)
    //     bytes 16..23   int64 sequence number   (little-endian)
    //
    // Requesters generate a stable GID at construction and auto-increment
    // the sequence number on every `send_request`.  Repliers echo the
    // caller's GID + sequence number on every `send_reply` so the
    // originating requester can correlate the response.  Requesters
    // filter incoming replies by GID so they do not see replies destined
    // for other requesters on the same service.

    struct RequesterQos
    {
        ReliabilityQosPolicy reliability{};   
        HistoryQosPolicy     history{};       
        DurabilityQosPolicy  durability{};    

        RequesterQos()
        {
            reliability.kind = ReliabilityQosPolicyKind::RELIABLE_RELIABILITY_QOS;
        }
    };
    using ReplierQos = RequesterQos;

    struct TakenRequest
    {
        std::vector<uint8_t>    payload;        
        std::array<uint8_t, 16> client_gid{};   
        int64_t                 sequence_number{0};
    };

    struct TakenReply
    {
        std::vector<uint8_t> payload;        
        int64_t              sequence_number{0};
    };

    class Requester
    {
    public:
        Requester(DomainParticipant *participant,
                  const std::string &service_name,
                  const std::string &request_type_name,
                  const std::string &reply_type_name,
                  const RequesterQos &qos);
        ~Requester();

        Requester(const Requester &) = delete;
        Requester &operator=(const Requester &) = delete;

        int64_t send_request(const std::vector<uint8_t> &payload);

        std::optional<TakenReply> take_reply();

        const std::array<uint8_t, 16> &gid() const { return gid_; }

        const std::string &service_name() const { return service_name_; }

        DataReader *reply_reader() const { return reader_; }
        DataWriter *request_writer() const { return writer_; }

    private:
        DomainParticipant       *participant_;
        std::string              service_name_;
        std::array<uint8_t, 16>  gid_{};
        std::atomic<int64_t>     next_seq_{1};

        // Owned indirectly through participant_'s topics_/publishers_/
        // subscribers_ vectors.  These raw pointers are non-owning.
        Topic       *req_topic_{nullptr};
        Topic       *rep_topic_{nullptr};
        Publisher   *publisher_{nullptr};
        Subscriber  *subscriber_{nullptr};
        DataWriter  *writer_{nullptr};
        DataReader  *reader_{nullptr};
    };

    class Replier
    {
    public:
        Replier(DomainParticipant *participant,
                const std::string &service_name,
                const std::string &request_type_name,
                const std::string &reply_type_name,
                const ReplierQos &qos);
        ~Replier();

        Replier(const Replier &) = delete;
        Replier &operator=(const Replier &) = delete;

        std::optional<TakenRequest> take_request();

        bool send_reply(const std::array<uint8_t, 16> &client_gid,
                        int64_t                         sequence_number,
                        const std::vector<uint8_t>     &payload);

        const std::string &service_name() const { return service_name_; }

        DataReader *request_reader() const { return reader_; }
        DataWriter *reply_writer() const { return writer_; }

    private:
        DomainParticipant *participant_;
        std::string        service_name_;

        Topic       *req_topic_{nullptr};
        Topic       *rep_topic_{nullptr};
        Publisher   *publisher_{nullptr};
        Subscriber  *subscriber_{nullptr};
        DataWriter  *writer_{nullptr};
        DataReader  *reader_{nullptr};
    };

    class DomainParticipant
    {
    public:
        DomainParticipant(uint32_t domain_id, uint32_t participant_id);
        ~DomainParticipant();
        DomainParticipant(const DomainParticipant &) = delete;
        DomainParticipant &operator=(const DomainParticipant &) = delete;

        bool enable();

        void stop();

        void set_participant_lease_duration_seconds(uint32_t sec);

        void set_spdp_announce_interval_ms(uint32_t ms);

        // ── Topic creation ──────────────────────────────────────────────

        Topic *create_topic(const std::string &topic_name,
                            const std::string &type_name,
                            const TopicQos &qos = TopicQos{});

        Topic *create_topic(const char *topic_name,
                            const char *type_name,
                            const TopicQos &qos,
                            TopicListener *listener,
                            StatusMask mask = STATUS_MASK_ALL)
        {
            (void)listener;
            (void)mask;
            return create_topic(std::string(topic_name), std::string(type_name), qos);
        }

        bool delete_topic(Topic *topic);

        // ── ContentFilteredTopic ────────────────────────────────────────

        ContentFilteredTopic *create_contentfilteredtopic(
            const char *name,
            Topic *related_topic,
            const char *filter_expression,
            const StringSeq &expression_parameters);

        bool delete_contentfilteredtopic(ContentFilteredTopic *cft);

        // ── Publisher / Subscriber creation ─────────────────────────────

        Publisher *create_publisher(const PublisherQos &qos = PublisherQos{});

        Publisher *create_publisher(const PublisherQos &qos,
                                    PublisherListener *listener,
                                    StatusMask mask = STATUS_MASK_ALL)
        {
            (void)listener;
            (void)mask;
            return create_publisher(qos);
        }

        bool delete_publisher(Publisher *publisher);

        Subscriber *create_subscriber(const SubscriberQos &qos = SubscriberQos{});

        Subscriber *create_subscriber(const SubscriberQos &qos,
                                      SubscriberListener *listener,
                                      StatusMask mask = STATUS_MASK_ALL)
        {
            (void)listener;
            (void)mask;
            return create_subscriber(qos);
        }

        bool delete_subscriber(Subscriber *subscriber);

        // ── §1.4 Request/Reply factories ────────────────────────────────

        Requester *create_requester(const std::string &service_name,
                                    const std::string &request_type_name,
                                    const std::string &reply_type_name,
                                    const RequesterQos &qos = RequesterQos{});

        bool delete_requester(Requester *requester);

        Replier *create_replier(const std::string &service_name,
                                const std::string &request_type_name,
                                const std::string &reply_type_name,
                                const ReplierQos &qos = ReplierQos{});

        bool delete_replier(Replier *replier);

        // ── Default QoS out-param accessors (PSM) ──────────────────────

        ReturnCode_t get_default_topic_qos(TopicQos &qos) const
        {
            qos = TopicQos{};
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t get_default_publisher_qos(PublisherQos &qos) const
        {
            qos = PublisherQos{};
            return ReturnCode_t::RETCODE_OK;
        }

        ReturnCode_t get_default_subscriber_qos(SubscriberQos &qos) const
        {
            qos = SubscriberQos{};
            return ReturnCode_t::RETCODE_OK;
        }

        // ── Cleanup ────────────────────────────────────────────────────

        ReturnCode_t delete_contained_entities();

        // ── Type Support Factory ───────────────────────────────────────

        using DataWriterFactory = std::function<DataWriter *(Topic *, Publisher *)>;
        using DataReaderFactory = std::function<DataReader *(Topic *, Subscriber *)>;

        void register_type_factory(const std::string &type_name,
                                   DataWriterFactory wf, DataReaderFactory rf);

        DataWriterFactory get_writer_factory(const std::string &type_name) const;

        DataReaderFactory get_reader_factory(const std::string &type_name) const;

        // ── Accessors ──────────────────────────────────────────────────

        uint32_t domain_id() const { return domain_id_; }
        uint32_t participant_id() const { return participant_id_; }
        const rtps::GuidPrefix_t &guid_prefix() const { return guid_prefix_; }

        rtps::RtpsUdpTransport *rtps_transport() const { return transport_.get(); }

        struct TopicSnapshot
        {
            std::string name;
            std::string type_name;
            TopicQos    qos;
        };
        std::vector<TopicSnapshot> topics_snapshot() const;

        void set_qos(const DomainParticipantQos &qos);
        const DomainParticipantQos &get_qos() const { return qos_; }

        bool is_enabled() const { return running_.load(); }

        // Entity framework (DDS §2.2.2.1) — DomainParticipant-specific overloads.
        // enable() and is_enabled() are declared above; get_instance_handle(),
        // get_status_changes() and get_statuscondition() below complete the
        // uniform Entity API required across all entity kinds.
        InstanceHandle_t get_instance_handle() const noexcept { return instance_handle_; }
        StatusMask get_status_changes() const noexcept { return status_condition_.get_status_changes(); }
        EntityStatusCondition* get_statuscondition() noexcept { return &status_condition_; }

        // Discovery / introspection (DDS §2.2.2.2.1.15–19)
        ReturnCode_t get_discovered_participants(std::vector<InstanceHandle_t>& participant_handles) const;

        ReturnCode_t get_discovered_participant_data(ParticipantBuiltinTopicData& data,
                                                      InstanceHandle_t handle) const;

        ReturnCode_t get_discovered_topics(std::vector<InstanceHandle_t>& topic_handles) const;

        ReturnCode_t get_discovered_topic_data(TopicBuiltinTopicData& data,
                                                InstanceHandle_t handle) const;

        BuiltinSubscriber* get_builtin_subscriber();

        bool contains_entity(InstanceHandle_t handle) const;

        // Listener management
        void set_listener(DomainParticipantListener *listener, StatusMask mask = STATUS_MASK_ALL)
        {
            listener_ = listener;
            listener_mask_ = mask;
        }
        DomainParticipantListener *get_listener() const { return listener_; }
        StatusMask get_listener_mask() const { return listener_mask_; }

        void send_data(const std::string &topic_name, const std::vector<uint8_t> &data);

        void send_data(const std::string &topic_name,
                       const std::vector<uint8_t> &xcdr1Data,
                       const std::vector<uint8_t> &xcdr2Data);

        void send_data_coherent(const std::string &topic_name,
                                const std::vector<uint8_t> &xcdr1Data,
                                const std::vector<uint8_t> &xcdr2Data,
                                uint64_t in_start_sn,
                                bool end_of_set,
                                uint64_t &out_start_sn);

        void send_dispose(const std::string &topic_name,
                          const std::vector<uint8_t> &xcdr1Data,
                          const std::vector<uint8_t> &xcdr2Data,
                          uint32_t statusInfoBits,
                          int64_t tsSecondsOverride = -1,
                          uint32_t tsFractionOverride = 0);

        using SendDataCallback = std::function<void(const std::string &topic_name, const std::vector<uint8_t> &data)>;
        void set_send_data_callback(SendDataCallback callback) { send_data_callback_ = callback; }

        using SubscribeCallback = std::function<void(const std::string &topic_name)>;
        void set_subscribe_callback(SubscribeCallback callback) { subscribe_callback_ = callback; }

        void register_writer(const std::string &topic_name,
                             const std::string &type_name,
                             const std::vector<std::string> &partitions = {},
                             const DataWriterQos &qos = {},
                             const std::vector<uint8_t> &topic_data = {},
                             const std::vector<uint8_t> &group_data = {});

        void register_reader(const std::string &topic_name,
                             const std::string &type_name,
                             const std::vector<std::string> &partitions = {},
                             const DataReaderQos &qos = {},
                             const std::vector<uint8_t> &contentFilterPayload = {},
                             const std::vector<uint8_t> &topic_data = {},
                             const std::vector<uint8_t> &group_data = {});

        void unregister_writer(const std::string &topic_name);

        void unregister_reader(const std::string &topic_name);

        void notify_subscription(const std::string &topic_name);

        void deliver_data(const std::string &topic_name, const std::vector<uint8_t> &data,
                          const rtps::GUID_t &writer_guid = rtps::GUID_t{},
                          int32_t ownershipStrength = 0,
                          const rtps::SequenceNumber_t &sequenceNumber = rtps::SequenceNumber_t{},
                          uint64_t coherent_set_start_sn = 0,
                          bool     end_of_coherent_set   = false);

        void local_deliver(const std::string &topic_name,
                            const std::vector<uint8_t> &data,
                            const rtps::GUID_t &writer_guid,
                            const rtps::SequenceNumber_t &sequenceNumber);

        void local_deliver_coherent(const std::string &topic_name,
                                     const std::vector<uint8_t> &data,
                                     const rtps::GUID_t &writer_guid,
                                     const rtps::SequenceNumber_t &sequenceNumber,
                                     uint64_t coherent_set_start_sn,
                                     bool     end_of_coherent_set);

        void deliver_dispose(const std::string &topic_name,
                             const std::vector<uint8_t> &keyPayload,
                             const rtps::GUID_t &writer_guid,
                             const rtps::SequenceNumber_t &sequenceNumber,
                             uint32_t statusInfoBits);

        void replay_history_to_reader(DataWriter *writer);

        using DataReceivedCallback = std::function<void(const std::string &topic, const ReceivedSample &)>;
        void set_data_received_callback(DataReceivedCallback callback) { data_received_callback_ = callback; }

        void match_reader_to_discovered_writers(DataReader *reader, const std::vector<std::string> &subPartitions);

        void match_writer_to_discovered_readers(DataWriter *writer, const std::vector<std::string> &pubPartitions);

        std::vector<std::pair<std::string, std::string>> get_discovered_topic_names_and_types() const;

        struct DiscoveredPublication
        {
            rtps::GUID_t endpointGuid;                    
            rtps::GuidPrefix_t participantGuid;           
            std::string topicName;
            std::string typeName;
            std::vector<std::string> partitions;
            std::vector<uint8_t> participantUserData;     
            std::string participantName;                  
            ReliabilityQosPolicy reliability;
            DurabilityQosPolicy durability;
            OwnershipQosPolicy ownership;
            int32_t ownershipStrength{0};
            DeadlineQosPolicy deadline;
            DataRepresentationQosPolicy data_representation;
            std::chrono::steady_clock::time_point lastSeen;
        };

        struct DiscoveredSubscription
        {
            rtps::GUID_t endpointGuid;
            rtps::GuidPrefix_t participantGuid;
            std::string topicName;
            std::string typeName;
            std::vector<std::string> partitions;
            std::vector<uint8_t> participantUserData;
            std::string participantName;
            ReliabilityQosPolicy reliability;
            DurabilityQosPolicy durability;
            OwnershipQosPolicy ownership;
            DeadlineQosPolicy deadline;
            DataRepresentationQosPolicy data_representation;
            std::chrono::steady_clock::time_point lastSeen;
        };

        std::vector<DiscoveredPublication> get_discovered_publications() const;

        std::vector<DiscoveredSubscription> get_discovered_subscriptions() const;

#ifdef ASTUTEDDS_COMMUNITY_EDITION
        // Community Edition endpoint slot management — one slot per DataWriter/DataReader.
        bool acquire_writer_slot();
        bool acquire_reader_slot();
        void release_writer_slot();
        void release_reader_slot();
#endif

    private:
        void discovery_loop();
        void generate_guid_prefix();
        rtps::GUID_t create_writer_guid(uint32_t writer_id);
        rtps::GUID_t create_reader_guid(uint32_t reader_id);

        uint32_t domain_id_;
        uint32_t participant_id_;
        rtps::GuidPrefix_t guid_prefix_;
        DomainParticipantQos qos_;

        // Entity framework state (DDS §2.2.2.1) — is_enabled() is derived
        // from running_ (started in enable()); instance_handle_ is stable
        // for the participant's lifetime.
        InstanceHandle_t instance_handle_{next_entity_instance_handle()};
        EntityStatusCondition status_condition_{};

        std::atomic<bool> running_{false};
        std::thread discovery_thread_;
        mutable std::mutex mutex_;

        std::vector<std::unique_ptr<Topic>> topics_;
        // §2.2.2.2.1.5 — create_topic() dedups by (name, type) and returns the
        // same native Topic* to every caller requesting that topic. Multiple
        // independent DataWriters/DataReaders (even from unrelated call sites
        // sharing this DomainParticipant) can therefore hold a raw pointer to
        // the same Topic. delete_topic() must only actually destroy the Topic
        // once every create_topic() caller has released it — otherwise the
        // first caller to tear down deletes the Topic out from under any
        // other still-open DataWriter/DataReader referencing it, causing a
        // use-after-free on the next write/close of that endpoint.
        std::unordered_map<Topic *, int> topic_ref_counts_;
        std::vector<std::unique_ptr<Publisher>> publishers_;
        std::vector<std::unique_ptr<Subscriber>> subscribers_;
        std::vector<std::unique_ptr<ContentFilteredTopic>> cfts_;
        std::vector<std::unique_ptr<Requester>> requesters_;
        std::vector<std::unique_ptr<Replier>> repliers_;

        // §2.2.2.2.1.13 built-in subscriber — lazily constructed on first
        // get_builtin_subscriber() call; owned for participant lifetime.
        std::unique_ptr<BuiltinSubscriber> builtin_subscriber_;

        SendDataCallback send_data_callback_;
        SubscribeCallback subscribe_callback_;
        DataReceivedCallback data_received_callback_;

        // Listener
        DomainParticipantListener *listener_{nullptr};
        StatusMask listener_mask_{0};

        // Type support factories (populated by register_type_factory)
        std::map<std::string, DataWriterFactory> writer_factories_;
        std::map<std::string, DataReaderFactory> reader_factories_;

        // Built-in RTPS transport
        std::unique_ptr<rtps::RtpsUdpTransport> transport_;

        // §1.1b-B liveliness watchdog — owns a TimerService and tracks
        // lease deadlines for local writers (lost) and remote writers
        // (changed).  Created in enable(), destroyed in stop().
        std::unique_ptr<LivelinessManager> liveliness_manager_;

        // §1.1b-B deadline watchdog — shared TimerService used by every
        // local DataWriter (offered deadline) and DataReader (requested
        // deadline).  Created in enable(), stopped in stop() *before*
        // publishers_/subscribers_ are cleared so no timer callback can
        // dereference a torn-down entity.
        std::unique_ptr<TimerService> deadline_timers_;

    public:
        LivelinessManager *liveliness_manager() { return liveliness_manager_.get(); }

        TimerService *deadline_timers() { return deadline_timers_.get(); }

    private:

        // Cached transport tuning values.  Applied at enable() time (before
        // transport init()) and also live-propagated if the setter is called
        // after enable().  See set_participant_lease_duration_seconds and
        // set_spdp_announce_interval_ms.
        uint32_t pending_lease_duration_sec_{100};
        uint32_t pending_spdp_announce_interval_ms_{3000};

        // Tracking sets for SEDP callback deduplication and data delivery gating.
        // matchedRemoteWriters_ contains GUIDs of remote writers whose QoS is
        // compatible with at least one local reader. deliver_data() only delivers
        // samples from writers in this set.
        std::set<rtps::GUID_t> matchedRemoteWriters_;
        std::set<rtps::GUID_t> incompatRemoteWriters_;
        std::set<rtps::GUID_t> matchedRemoteReaders_;
        std::set<rtps::GUID_t> incompatRemoteReaders_;

        // Pending TRANSIENT_LOCAL samples buffer.  Historical replay from a
        // TRANSIENT_LOCAL writer can reach us before the writer's SEDP
        // DATA(w) is processed on the reader side — the packets arrive out
        // of order on the network, and the writer emits the replay burst
        // synchronously on receipt of our SEDP DATA(r).  Without this
        // buffer, deliver_data() would drop those samples because the
        // writer is not yet in matchedRemoteWriters_.
        //
        // Bounded FIFO per writer GUID.  Drained on match (SEDP writer
        // discovered callback + retroactive match_reader_to_discovered_writers)
        // and purged on unmatch (writer_lost callback).
        struct PendingSample
        {
            std::string             topic_name;
            std::vector<uint8_t>    data;
            int32_t                 ownership_strength{0};
            rtps::SequenceNumber_t  sequence_number{};
            // Phase 3.d — carry RTPS §8.7.6 coherent-set markers through
            // the pending-samples queue so a late writer match still
            // preserves atomic-set semantics on drain.
            uint64_t                coherent_set_start_sn{0};
            bool                    end_of_coherent_set{false};
        };
        std::map<rtps::GUID_t, std::deque<PendingSample>> pending_samples_;
        // Per-writer cap to prevent unbounded memory growth if the writer is
        // never matched (e.g. incompatible QoS but discovered).  Larger than
        // any realistic TRANSIENT_LOCAL cache; drops oldest on overflow.
        static constexpr size_t kMaxPendingSamplesPerWriter = 1024;
        // Issue #89 — global cap on the number of *distinct* writer GUIDs
        // buffered in pending_samples_ at once.  Needed because
        // is_writer_participant_known_locked() (below) admits any writer
        // belonging to an SPDP-discovered remote participant, which is a
        // coarser/earlier gate than the per-topic SEDP writer table and
        // therefore accepts a wider (but still SPDP-authenticated) set of
        // writer GUIDs than before.  Bounds worst-case memory from a
        // misbehaving/compromised remote participant that announces many
        // writer entity IDs that never complete SEDP matching.
        static constexpr size_t kMaxPendingWriterEntries = 256;
        // Deliver a single sample to all locally matched readers on the
        // given topic.  Precondition: caller holds mutex_.  Used by both
        // deliver_data() and drain_pending_samples_locked().
        void deliver_sample_locked(const std::string &topic_name,
                                   const std::vector<uint8_t> &data,
                                   const rtps::GUID_t &writer_guid,
                                   int32_t ownership_strength,
                                   const rtps::SequenceNumber_t &sequence_number,
                                   uint64_t coherent_set_start_sn = 0,
                                   bool     end_of_coherent_set   = false);
        // Drain the pending-samples FIFO for @p writer_guid, delivering
        // each buffered sample to matched readers.  Called from the
        // writer-discovered / retroactive-match paths after
        // matchedRemoteWriters_.insert().  Precondition: caller holds mutex_.
        void drain_pending_samples_locked(const rtps::GUID_t &writer_guid);
        // Issue #89 — true if @p writer_guid's owning remote participant
        // has been discovered via SPDP, independent of whether that
        // specific writer entity's own SEDP DATA(w) has been processed
        // into get_discovered_writers() yet.  SPDP discovery of the
        // remote participant reliably precedes any directed traffic from
        // it (a writer can only unicast its TRANSIENT_LOCAL replay burst
        // to a reader's locator once mutual SPDP has completed), so this
        // closes the residual race where the per-topic SEDP writer table
        // lags the arrival of that writer's replay burst by a few hundred
        // microseconds — while still refusing to buffer for writer GUIDs
        // with no corresponding SPDP participant record at all (stray /
        // spoofed traffic). Precondition: caller holds mutex_.
        bool is_writer_participant_known_locked(const rtps::GUID_t &writer_guid) const;

#ifdef ASTUTEDDS_COMMUNITY_EDITION
        // Community Edition endpoint counters (thread-safe).
        std::atomic<uint32_t> ce_writer_count_{0};
        std::atomic<uint32_t> ce_reader_count_{0};
#endif
    };

    class DomainParticipantFactory
    {
    public:
        static DomainParticipantFactory &instance();

        static DomainParticipantFactory *get_instance()
        {
            return &instance();
        }

        DomainParticipant *create_participant(uint32_t domain_id,
                                              const DomainParticipantQos &qos = DomainParticipantQos{});

        DomainParticipant *create_participant(DomainId_t domain_id,
                                              const DomainParticipantQos &qos,
                                              DomainParticipantListener *listener,
                                              StatusMask mask = STATUS_MASK_ALL)
        {
            auto *dp = create_participant(static_cast<uint32_t>(domain_id), qos);
            if (dp && listener)
            {
                dp->set_listener(listener, mask);
            }
            return dp;
        }

        ReturnCode_t get_default_participant_qos(DomainParticipantQos &qos) const
        {
            qos = DomainParticipantQos{};
            return ReturnCode_t::RETCODE_OK;
        }

        bool delete_participant(DomainParticipant *participant);

        DomainParticipant *lookup_participant(uint32_t domain_id);

    private:
        DomainParticipantFactory() = default;

        std::mutex mutex_;
        std::vector<std::unique_ptr<DomainParticipant>> participants_;
        std::unordered_map<DomainParticipant*, uint32_t> ref_counts_;
        uint32_t next_participant_id_{0};
    };

} // namespace astutedds::dcps

// Include Cyclone-style namespace compatibility layer after class definitions
#include <astutedds/dcps/CorePolicy.hpp>

// Standard DDS namespace alias for compatibility
namespace DDS = astutedds::dcps;

#endif // ASTUTEDDS_DCPS_DCPS_HPP