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;
};