File types.hpp
File List > astutedds > rmw > types.hpp
Go to the documentation of this file
//
// Copyright (c) 2026, Astute Systems PTY LTD
//
// This file is part of the AstuteDDS RMW implementation.
//
// See the commercial LICENSE file in the project root for full license details.
//
#pragma once
#include <atomic>
#include <cstdint>
#include <map>
#include <mutex>
#include <new>
#include <string>
#include <utility>
#include <vector>
#include <fcntl.h>
#include <unistd.h>
#include <rmw/event_callback_type.h>
#include <rmw/qos_profiles.h>
#include <rmw/rmw.h>
#include <astutedds/c/astutedds.h>
#include <astutedds/rmw/identifier.hpp>
#include <astutedds/rmw/type_support.hpp>
// ---------------------------------------------------------------------------
// Context — one DomainParticipant per process
// ---------------------------------------------------------------------------
struct AstuteDDSContext
{
AstuteDDS_Participant participant { nullptr };
AstuteDDS_Publisher dds_pub { nullptr };
AstuteDDS_Subscriber dds_sub { nullptr };
std::mutex node_mtx;
std::vector<std::pair<std::string, std::string>> nodes;
std::mutex svc_mtx;
std::vector<std::string> services;
std::vector<std::string> clients;
struct NodeEndpoint
{
std::string node_name;
std::string node_namespace;
std::string dds_topic_name;
std::string type_name;
uint8_t gid[16] {};
rmw_qos_profile_t qos {};
};
std::mutex endpoint_mtx;
std::vector<NodeEndpoint> local_publishers;
std::vector<NodeEndpoint> local_subscriptions;
std::vector<NodeEndpoint> local_services;
std::vector<NodeEndpoint> local_clients;
std::mutex topic_mtx;
std::map<std::string, int> topic_refcount;
};
// ---------------------------------------------------------------------------
// Publisher / subscription
// ---------------------------------------------------------------------------
struct AstuteDDSPublisher
{
AstuteDDS_DataWriter writer { nullptr };
AstuteDDS_Topic topic { nullptr };
TypeSupportHandle type_support {};
std::string topic_name;
std::string type_name;
rmw_qos_profile_t qos {};
uint8_t gid[16] {};
std::mutex listener_mtx;
rmw_event_callback_t on_liveliness_lost_cb { nullptr };
const void * on_liveliness_lost_ud { nullptr };
rmw_event_callback_t on_offered_deadline_missed_cb { nullptr };
const void * on_offered_deadline_missed_ud { nullptr };
rmw_event_callback_t on_offered_incompatible_qos_cb { nullptr };
const void * on_offered_incompatible_qos_ud { nullptr };
rmw_event_callback_t on_publication_matched_cb { nullptr };
const void * on_publication_matched_ud { nullptr };
};
struct AstuteDDSSubscription
{
AstuteDDS_DataReader reader { nullptr };
AstuteDDS_Topic topic { nullptr };
TypeSupportHandle type_support {};
std::string topic_name;
std::string type_name;
rmw_qos_profile_t qos {};
uint8_t gid[16] {};
std::mutex listener_mtx;
rmw_event_callback_t on_new_message_cb { nullptr };
const void * on_new_message_user_data { nullptr };
rmw_event_callback_t on_subscription_matched_cb { nullptr };
const void * on_subscription_matched_ud { nullptr };
rmw_event_callback_t on_requested_incompatible_qos_cb { nullptr };
const void * on_requested_incompatible_qos_ud { nullptr };
rmw_event_callback_t on_liveliness_changed_cb { nullptr };
const void * on_liveliness_changed_ud { nullptr };
rmw_event_callback_t on_requested_deadline_missed_cb { nullptr };
const void * on_requested_deadline_missed_ud { nullptr };
rmw_event_callback_t on_message_lost_cb { nullptr };
const void * on_message_lost_ud { nullptr };
};
// ---------------------------------------------------------------------------
// Service / client (RPC via the AstuteDDS Requester/Replier primitives)
// ---------------------------------------------------------------------------
struct AstuteDDSClient
{
AstuteDDS_Requester requester { nullptr };
TypeSupportHandle request_ts {};
TypeSupportHandle response_ts {};
rmw_qos_profile_t qos {};
std::string service_name;
std::mutex cb_mtx;
rmw_event_callback_t on_new_response_cb { nullptr };
const void * on_new_response_user_data { nullptr };
};
struct AstuteDDSService
{
AstuteDDS_Replier replier { nullptr };
TypeSupportHandle request_ts {};
TypeSupportHandle response_ts {};
rmw_qos_profile_t qos {};
std::string service_name;
std::mutex cb_mtx;
rmw_event_callback_t on_new_request_cb { nullptr };
const void * on_new_request_user_data { nullptr };
};
// ---------------------------------------------------------------------------
// Guard condition — pipe-based for efficient poll()/select()
// ---------------------------------------------------------------------------
struct AstuteDDSGuardCondition
{
int pipe_read { -1 };
int pipe_write { -1 };
std::atomic<bool> triggered { false };
AstuteDDSGuardCondition()
{
int fds[2];
if (::pipe2(fds, O_NONBLOCK | O_CLOEXEC) == 0)
{
pipe_read = fds[0];
pipe_write = fds[1];
}
}
~AstuteDDSGuardCondition()
{
if (pipe_read >= 0) { ::close(pipe_read); }
if (pipe_write >= 0) { ::close(pipe_write); }
}
AstuteDDSGuardCondition(const AstuteDDSGuardCondition &) = delete;
AstuteDDSGuardCondition & operator=(const AstuteDDSGuardCondition &) = delete;
};
// ---------------------------------------------------------------------------
// Wait set + node
// ---------------------------------------------------------------------------
struct AstuteDDSWaitSet
{
std::vector<AstuteDDSSubscription *> subscriptions;
std::vector<AstuteDDSGuardCondition *> guard_conditions;
};
struct AstuteDDSNode
{
rmw_guard_condition_t * graph_guard_condition { nullptr };
AstuteDDSNode()
{
auto * gc_impl = new (std::nothrow) AstuteDDSGuardCondition();
graph_guard_condition = rmw_guard_condition_allocate();
if (graph_guard_condition)
{
graph_guard_condition->implementation_identifier = kIdentifier;
graph_guard_condition->data = gc_impl;
}
else
{
delete gc_impl;
}
}
~AstuteDDSNode()
{
if (graph_guard_condition)
{
auto * gc = reinterpret_cast<AstuteDDSGuardCondition *>(graph_guard_condition->data);
delete gc;
rmw_guard_condition_free(graph_guard_condition);
}
}
AstuteDDSNode(const AstuteDDSNode &) = delete;
AstuteDDSNode & operator=(const AstuteDDSNode &) = delete;
};