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