diff --git a/azure-pipelines.yml b/azure-pipelines.yml index 87a563dc76..b427c484e3 100644 --- a/azure-pipelines.yml +++ b/azure-pipelines.yml @@ -29,6 +29,7 @@ jobs: build_type: Release sanitizer: undefined cc: gcc + cxx: g++ 'Ubuntu 24.04 (Debug, x86_64, Iceoryx)': image: ubuntu-24.04 # No address sanitizer because of this in test run: @@ -39,6 +40,7 @@ jobs: sanitizer: undefined iceoryx: on cc: gcc + cxx: g++ coverage: on 'Ubuntu 24.04 (Debug, x86_64, Iceoryx2)': image: ubuntu-24.04 @@ -50,6 +52,7 @@ jobs: sanitizer: undefined iceoryx2: on cc: gcc + cxx: g++ coverage: on 'Ubuntu 24.04 (Release, x86_64)': image: ubuntu-24.04 @@ -67,9 +70,11 @@ jobs: topic_discovery: off idlc_xtests: off # temporary disabled because of passing -t option to idlc in this test for recursive types cc: gcc-12 + cxx: g++-12 'Ubuntu 24.04 with GCC 12 (Debug, x86_64, no tests)': image: ubuntu-24.04 cc: gcc-12 + cxx: g++-12 testing: off idlc_xtests: off 'Ubuntu 24.04 with Clang (Debug, x86_64)': @@ -77,11 +82,13 @@ jobs: analyzer: on sanitizer: address,undefined cc: clang + cxx: clang++ 'Ubuntu 24.04 with Clang (Debug, x86_64, no security)': image: ubuntu-24.04 sanitizer: address,undefined security: off cc: clang + cxx: clang++ 'Ubuntu 24.04 with Clang (Release, x86_64, no topic discovery)': image: ubuntu-24.04 build_type: Release @@ -89,20 +96,24 @@ jobs: topic_discovery: off idlc_xtests: off # temporary disabled because of passing -t option to idlc in this test for recursive types cc: clang + cxx: clang++ 'macOS 14 with Clang (Debug, x86_64)': image: macos-14 sanitizer: address,undefined deadline_update_skip: on cc: clang + cxx: clang++ coverage: on 'macOS 14 with Clang (Release, x86_64)': image: macos-14 build_type: Release sanitizer: undefined cc: clang + cxx: clang++ 'macOS 14 with GCC 14 (Debug, analyzer, x86_64)': image: macos-14 cc: gcc-14 + cxx: g++-14 analyzer: on # 32-bit Windows: without SSL/security because Chocolateley only provides 64-bit OpenSSL 'Windows 2025 with Visual Studio 2022 (Debug, x86, no security)': diff --git a/examples/CMakeLists.txt b/examples/CMakeLists.txt index 1c5d5e6029..54a4f148a0 100644 --- a/examples/CMakeLists.txt +++ b/examples/CMakeLists.txt @@ -76,6 +76,7 @@ if (ENABLE_TOPIC_DISCOVERY) endif () add_subdirectory(helloworld) +add_subdirectory(recorder_and_replayer) add_subdirectory(roundtrip) add_subdirectory(throughput) if (ENABLE_TOPIC_DISCOVERY) diff --git a/examples/recorder_and_replayer/CMakeLists.txt b/examples/recorder_and_replayer/CMakeLists.txt new file mode 100644 index 0000000000..9103870325 --- /dev/null +++ b/examples/recorder_and_replayer/CMakeLists.txt @@ -0,0 +1,99 @@ +# Copyright(c) 2023 ZettaScale Technology and others +# +# This program and the accompanying materials are made available under the +# terms of the Eclipse Public License v. 2.0 which is available at +# http://www.eclipse.org/legal/epl-2.0, or the Eclipse Distribution License +# v. 1.0 which is available at +# http://www.eclipse.org/org/documents/edl-v10.php. +# +# SPDX-License-Identifier: EPL-2.0 OR BSD-3-Clause + +cmake_minimum_required(VERSION 3.16) +project(recorder_replayer LANGUAGES C CXX) + +set(CMAKE_CXX_STANDARD 17) +set(CMAKE_CXX_STANDARD_REQUIRED) +include(FetchContent) + +# Prevent subprojects (like mcap_builder) from attempting to find lz4/zstd +set(CMAKE_DISABLE_FIND_PACKAGE_lz4 TRUE) +set(CMAKE_DISABLE_FIND_PACKAGE_zstd TRUE) +set(CMAKE_DISABLE_FIND_PACKAGE_LZ4 TRUE) +set(CMAKE_DISABLE_FIND_PACKAGE_ZSTD TRUE) + +set(BUILD_SHARED_LIBS OFF CACHE BOOL "" FORCE) +FetchContent_Declare( + mcap_builder + GIT_REPOSITORY https://github.com/olympus-robotics/mcap_builder.git + GIT_TAG main +) +FetchContent_MakeAvailable(mcap_builder) + +if(NOT TARGET CycloneDDS::ddsc) + find_package(CycloneDDS REQUIRED) +endif() + +add_library(dds_topic_descriptor_serde STATIC dds_topic_descriptor_serde.c) +if(TARGET CycloneDDS::ddsc) + # dds_topic_descriptor_serde need dds_stream_countops() + target_include_directories(dds_topic_descriptor_serde + PRIVATE + "${CMAKE_CURRENT_SOURCE_DIR}/../../src/core/cdr/include") +endif() +target_include_directories(dds_topic_descriptor_serde PUBLIC .) + +target_link_libraries(dds_topic_descriptor_serde PRIVATE CycloneDDS::ddsc) + +add_executable(recorder dds_recorder.cpp) + +if(TARGET CycloneDDS::ddsc) + target_include_directories(recorder + PRIVATE + "${CMAKE_CURRENT_SOURCE_DIR}/../../src/core/cdr/include" + "${CMAKE_CURRENT_SOURCE_DIR}/../../src/core/ddsi/include") +endif() +target_link_libraries(recorder CycloneDDS::ddsc dds_topic_descriptor_serde mcap) + + +# add_executable(replayer dds_replayer.cpp) + +# if(TARGET CycloneDDS::ddsc) +# target_include_directories(replayer +# PRIVATE +# "${CMAKE_CURRENT_SOURCE_DIR}/../../src/core/cdr/include" +# "${CMAKE_CURRENT_SOURCE_DIR}/../../src/core/ddsi/include") +# endif() +# target_link_libraries(replayer CycloneDDS::ddsc dds_topic_descriptor_serde mcap) + +# force to ignore many compile warnings in mcap... +if(CMAKE_CXX_COMPILER_ID MATCHES "Clang" OR CMAKE_CXX_COMPILER_ID MATCHES "AppleClang") + set(disable_flags + "-Wno-documentation" + "-Wno-error=documentation" + "-Wno-sign-conversion" + "-Wno-error=sign-conversion" + "-Wno-implicit-int-conversion" + "-Wno-error=implicit-int-conversion" + ) + target_compile_options(mcap PRIVATE ${disable_flags}) + target_compile_options(recorder PRIVATE ${disable_flags}) + # target_compile_options(replayer PRIVATE ${disable_flags}) +endif() + +if(MSVC) + target_compile_options(mcap PRIVATE /wd4251 /WX-) + target_compile_options(recorder PRIVATE /wd4251 /WX-) +endif() + +if(CMAKE_C_COMPILER_ID STREQUAL "GNU") + if(ANALYZER STREQUAL "on") + if(CMAKE_C_COMPILER_VERSION VERSION_GREATER_EQUAL "12") + set(disable_flags + "-Wno-error=analyzer-malloc-leak" + "-Wno-analyzer-malloc-leak" + ) + target_compile_options(mcap PRIVATE ${disable_flags}) + target_compile_options(recorder PRIVATE ${disable_flags}) + endif() + endif() +endif() diff --git a/examples/recorder_and_replayer/dds_recorder.cpp b/examples/recorder_and_replayer/dds_recorder.cpp new file mode 100644 index 0000000000..22d4dee2e6 --- /dev/null +++ b/examples/recorder_and_replayer/dds_recorder.cpp @@ -0,0 +1,367 @@ +/* + * Copyright(c) 2025 ZettaScale Technology and others + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v. 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Eclipse Distribution License + * v. 1.0 which is available at + * http://www.eclipse.org/org/documents/edl-v10.php. + * + * SPDX-License-Identifier: EPL-2.0 OR BSD-3-Clause + */ + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "dds/dds.h" +#include "mcap/writer.hpp" + +// most part is copied from dynsub.c, but ignore xtypeobj. +// we do not have a reliable way with c api to obtain idl at runtime, so we have to +// offer it as input. + +#include "dds/ddsrt/threads.h" +#include "dds/ddsi/ddsi_serdata.h" +#include "dds_topic_descriptor_serde.h" + +// For convenience, the DDS participant is global +static dds_entity_t participant; + +static dds_entity_t termcond; + +// Helper function to wait for a DCPSPublication/DCPSSubscription to show up with the desired topic name, +// and its descriptor +static dds_return_t get_topic_and_desc (const char *topic_name, dds_duration_t timeout, dds_entity_t *topic, + dds_topic_descriptor_t **descriptor) +{ + *descriptor = NULL; + const dds_entity_t waitset = dds_create_waitset (participant); + // we only care publisher as we are recorder + const dds_entity_t dcpspublication_reader = dds_create_reader (participant, DDS_BUILTIN_TOPIC_DCPSPUBLICATION, NULL, NULL); + const dds_entity_t dcpspublication_readcond = dds_create_readcondition (dcpspublication_reader, DDS_ANY_STATE); + (void)dds_waitset_attach (waitset, dcpspublication_readcond, dcpspublication_reader); + const dds_time_t abstimeout = (timeout == DDS_INFINITY) ? DDS_NEVER : dds_time () + timeout; + dds_return_t ret = DDS_RETCODE_OK; + dds_attach_t triggered_reader_x; + while (dds_waitset_wait_until (waitset, &triggered_reader_x, 1, abstimeout) > 0) + { + void *epraw = NULL; + dds_sample_info_t si; + dds_entity_t triggered_reader = (dds_entity_t)triggered_reader_x; + if (dds_take (triggered_reader, &epraw, &si, 1, 1) <= 0) + continue; + dds_builtintopic_endpoint_t *ep = (dds_builtintopic_endpoint_t *)epraw; + const dds_typeinfo_t *typeinfo = NULL; + + // we need the topic descriptor so we do not use dds_find_topic + if (ep->topic_name != NULL && strcmp (ep->topic_name, topic_name) == 0 && + dds_builtintopic_get_endpoint_type_info (ep, &typeinfo) == 0 && typeinfo) + { + if ((ret = dds_create_topic_descriptor (DDS_FIND_SCOPE_GLOBAL, participant, typeinfo, DDS_SECS (10), descriptor)) < 0) + { + fprintf (stderr, "dds_create_topic_descriptor: %s\n", dds_strretcode (ret)); + dds_return_loan (triggered_reader, &epraw, 1); + goto error; + } + dds_qset_data_representation (ep->qos, 0, NULL); + if ((*topic = dds_create_topic (participant, *descriptor, ep->topic_name, ep->qos, NULL)) < 0) + { + fprintf (stderr, "dds_create_topic_descriptor: %s (be sure to enable topic discovery in the configuration)\n", + dds_strretcode (*topic)); + dds_delete_topic_descriptor (*descriptor); + dds_return_loan (triggered_reader, &epraw, 1); + goto error; + } + } + dds_return_loan (triggered_reader, &epraw, 1); + } + +error: + dds_delete (dcpspublication_reader); + dds_delete (waitset); + return (*descriptor != NULL) ? DDS_RETCODE_OK : DDS_RETCODE_TIMEOUT; +} + +static bool prepare_mcap (const char *topic_name, const char *type_name, const char *idl_file, const dds_topic_descriptor_t *desc, + mcap::McapWriter &writer, mcap::Channel &channel) +{ + auto opts = mcap::McapWriterOptions (""); + std::string outPath = std::string (topic_name) + ".mcap"; + // default chunking is fine; keep CRCs on for safety + auto status = writer.open (outPath, opts); + if (!status.ok ()) + { + std::cerr << "Failed to open MCAP: " << status.message << "\n"; + return false; + } + + // as it's not very easy to obtain idl text in runtime, we have to offer it manually... + // if it's not presented, tools such as foxglove/lichtblick won't be able to use it. + // but we still can replay with our dds_replayer. + // + // Note from mcap documents: + // the idl text must be the text of a single, self-contained OMG IDL source file. + // That is, all referenced type definitions must be present, and there must be no + // preprocessor directives, i.e. #include "another.idl". + + mcap::Schema schema; + schema.id = 0; + if (idl_file != NULL && idl_file[0] != '\0' && type_name != NULL && type_name[0] != '\0') + { + std::vector idl_text; + FILE *fp = fopen (idl_file, "rb"); + if (fp) + { + if (fseek (fp, 0, SEEK_END) == 0) + { + long len = ftell (fp); + if (len > 0) + { + idl_text.resize (static_cast (len)); + rewind (fp); + size_t n = fread (idl_text.data (), 1, idl_text.size (), fp); + idl_text.resize (n); + } + } + fclose (fp); + } else + { + std::cerr << "Warning: failed to open IDL file: " << idl_file << "\n"; + return false; + } + + schema.name = type_name; + schema.encoding = "omgidl"; + schema.data = std::move (idl_text); + writer.addSchema (schema); + } + + // fill channel(i.e. topic in mcap) + channel.topic = topic_name; + channel.messageEncoding = "cdr"; + channel.schemaId = schema.id; + + // we need to save the desc in the metadata of mcap file, + // so we can use it to build the topic when we replay it. + const size_t buffer_sz = dds_topic_descriptor_serialized_size (desc); + void *buffer = (void *)malloc (buffer_sz); + dds_topic_descriptor_serialize (desc, buffer, buffer_sz); + channel.metadata.emplace ("topic_descriptor", std::string (reinterpret_cast (buffer), buffer_sz)); + free (buffer); + writer.addChannel (channel); + return true; +} + +static bool write_to_mcap (const dds_entity_t reader, mcap::McapWriter *writer, mcap::Channel *channel) +{ + const size_t MAX_SAMPLE_SIZE = 10; + + struct ddsi_serdata *sd[MAX_SAMPLE_SIZE] = {NULL}; + dds_sample_info_t si[MAX_SAMPLE_SIZE]; + dds_return_t n = dds_takecdr (reader, sd, MAX_SAMPLE_SIZE, si, 0); + if (n < 0) + { + fprintf (stderr, "dds_takecdr: %s\n", dds_strretcode (n)); + return false; + } + + if (n != 0) + { + // save raw CDR into mcap + for (int32_t i = 0; i < n; i++) + { + if (si[i].valid_data) + { + size_t size = ddsi_serdata_size (sd[i]); + std::vector payload (size); + ddsi_serdata_to_ser (sd[i], 0, size, payload.data ()); // Extract the data from the buffer + ddsi_serdata_unref (sd[i]); + + // Build MCAP message + mcap::Message msg; + msg.channelId = channel->id; + const uint64_t ts = + (uint64_t)std::chrono::duration_cast (std::chrono::system_clock::now ().time_since_epoch ()).count (); + msg.logTime = ts; + msg.publishTime = si->source_timestamp; + msg.data = payload.data (); + msg.dataSize = payload.size (); + + const auto st = writer->write (msg); + if (st.ok ()) + { + printf ("one msg written to mcap..\n"); + } else + { + std::cerr << "MCAP write failed: " << st.message << "\n"; + return false; + } + } + } + } + + return true; +} +#if !DDSRT_WITH_FREERTOS && !__ZEPHYR__ +static void signal_handler (int sig) +{ + (void)sig; + dds_set_guardcondition (termcond, true); +} +#endif + +#if !_WIN32 && !DDSRT_WITH_FREERTOS && !__ZEPHYR__ +static uint32_t sigthread (void *varg) +{ + sigset_t *set = (sigset_t *)varg; + int sig; + if (sigwait (set, &sig) == 0) + signal_handler (sig); + return 0; +} +#endif + +int main (int argc, char **argv) +{ + dds_return_t ret = 0; + dds_entity_t topic = 0; + const char *topic_name = NULL; + const char *idl_file = NULL; + const char *type_name = NULL; + + if (argc == 2) + { + topic_name = argv[1]; + } else if (argc == 6 && strcmp (argv[2], "-i") == 0 && strcmp (argv[4], "-t") == 0) + { + topic_name = argv[1]; + idl_file = argv[3]; + type_name = argv[5]; + } else + { + fprintf (stderr, + "Usage:\n" + "%s [-i -t ]\n" + "\n" + "For example:\n" + "./bin/recorder HelloWorldData_Msg -i ../examples/helloworld/HelloWorldData.idl -t HelloWorldData::Msg\n" + "\n" + "Note:\n" + " - To stop recording, just press ctrl+c.\n" + " - The IDL should contain the complete type definition of .\n" + " - The IDL must not contain '#include' directives.\n" + " - if idl is not provided, the output mcap won't be able be visualized by tools like foxglove/lichtblick,\n" + " but it always can be replayed by our dds_replayer.\n", + argv[0]); + return 1; + } + + participant = dds_create_participant (DDS_DOMAIN_DEFAULT, NULL, NULL); + if (participant < 0) + { + fprintf (stderr, "dds_create_participant: %s\n", dds_strretcode (participant)); + return 1; + } + + // get a topic and topic desc ... + dds_topic_descriptor_t *desc; + if ((ret = get_topic_and_desc (topic_name, DDS_SECS (10), &topic, &desc)) < 0) + { + fprintf (stderr, "get_topic_and_desc: %s\n", dds_strretcode (ret)); + dds_delete (participant); + return 1; + } + printf ("found topic and desc..\n"); + + // prepare mcap writer + mcap::McapWriter writer; + mcap::Channel channel; + if (!prepare_mcap (topic_name, type_name, idl_file, desc, writer, channel)) + { + fprintf (stderr, "prepare mcap failed\n"); + dds_delete (participant); + return 1; + } + dds_delete_topic_descriptor (desc); + printf ("mcap is ready\n"); + + + // ... given those, we can create a reader just like we do normally ... + const dds_entity_t reader = dds_create_reader (participant, topic, NULL, NULL); + // ... and create a waitset that allows us to wait for any incoming data ... + const dds_entity_t waitset = dds_create_waitset (participant); + const dds_entity_t readcond = dds_create_readcondition (reader, DDS_ANY_STATE); + (void)dds_waitset_attach (waitset, readcond, 0); + + termcond = dds_create_guardcondition (participant); + (void)dds_waitset_attach (waitset, termcond, 0); + +#ifdef _WIN32 + signal (SIGINT, signal_handler); +#elif !DDSRT_WITH_FREERTOS && !__ZEPHYR__ + ddsrt_thread_t sigtid; + sigset_t sigset, osigset; + sigemptyset (&sigset); +#ifdef __APPLE__ + DDSRT_WARNING_GNUC_OFF (sign - conversion) +#endif + sigaddset (&sigset, SIGHUP); + sigaddset (&sigset, SIGINT); + sigaddset (&sigset, SIGTERM); +#ifdef __APPLE__ + DDSRT_WARNING_GNUC_ON (sign - conversion) +#endif + sigprocmask (SIG_BLOCK, &sigset, &osigset); + { + ddsrt_threadattr_t tattr; + ddsrt_threadattr_init (&tattr); + ddsrt_thread_create (&sigtid, "sigthread", &tattr, sigthread, &sigset); + } +#endif + + bool termflag = false; + while (!termflag) + { + (void)dds_waitset_wait (waitset, NULL, 0, DDS_INFINITY); + dds_read_guardcondition (termcond, &termflag); + + if (!write_to_mcap (reader, &writer, &channel)) + { + fprintf (stderr, "write_to_mcap failed, exiting..\n"); + break; + } + } + + if (termflag) + { + fprintf (stderr, "signal received, exiting..\n"); + } + +#if _WIN32 + signal_handler (SIGINT); +#elif !DDSRT_WITH_FREERTOS && !__ZEPHYR__ + { + /* get the attention of the signal handler thread */ + void (*osigint) (int); + void (*osigterm) (int); + kill (getpid (), SIGTERM); + ddsrt_thread_join (sigtid, NULL); + osigint = signal (SIGINT, SIG_IGN); + osigterm = signal (SIGTERM, SIG_IGN); + sigprocmask (SIG_SETMASK, &osigset, NULL); + signal (SIGINT, osigint); + signal (SIGINT, osigterm); + } +#endif + + writer.close (); + dds_delete (participant); + return 0; +} diff --git a/examples/recorder_and_replayer/dds_topic_descriptor_serde.c b/examples/recorder_and_replayer/dds_topic_descriptor_serde.c new file mode 100644 index 0000000000..15fb1b5561 --- /dev/null +++ b/examples/recorder_and_replayer/dds_topic_descriptor_serde.c @@ -0,0 +1,381 @@ +/* + * Copyright(c) 2025 ZettaScale Technology and others + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v. 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Eclipse Distribution License + * v. 1.0 which is available at + * http://www.eclipse.org/org/documents/edl-v10.php. + * + * SPDX-License-Identifier: EPL-2.0 OR BSD-3-Clause + */ + +#include "dds_topic_descriptor_serde.h" + +#include +#include +#include + +#include "dds/cdr/dds_cdrstream.h" + +// Helper function to calculate serialized size +size_t dds_topic_descriptor_serialized_size (const dds_topic_descriptor_t *desc) +{ + size_t size = 0; + + // Fixed fields + size += sizeof (desc->m_size); + size += sizeof (desc->m_align); + size += sizeof (desc->m_flagset); + size += sizeof (desc->m_nkeys); + size += sizeof (desc->m_nops); + size += sizeof (desc->restrict_data_representation); + + // m_typename (length + string) + size += sizeof (uint32_t); + if (desc->m_typename) + { + size += strlen (desc->m_typename) + 1; + } + + // m_keys array + size += desc->m_nkeys * sizeof (dds_key_descriptor_t); + for (uint32_t i = 0; i < desc->m_nkeys; i++) + { + size += sizeof (uint32_t); // name length + if (desc->m_keys[i].m_name) + { + size += strlen (desc->m_keys[i].m_name) + 1; + } + } + + // m_ops array + // m_ops array length is not always m_nops, and it seems the internal implement doesn't + // really rely on it, but calculate on the run. + // so we discard the origin m_nops and calculate it instead. + const uint32_t m_ops_len = dds_stream_countops (desc->m_ops, desc->m_nkeys, desc->m_keys); + size += m_ops_len * sizeof (uint32_t); + + // m_meta (length + string) + size += sizeof (uint32_t); + if (desc->m_meta) + { + size += strlen (desc->m_meta) + 1; + } + + // type_information + size += sizeof (uint32_t) + desc->type_information.sz; + + // type_mapping + size += sizeof (uint32_t) + desc->type_mapping.sz; + + return size; +} + +// Serialize dds_topic_descriptor to buffer +size_t dds_topic_descriptor_serialize (const dds_topic_descriptor_t *desc, void *buffer, size_t buffer_size) +{ + if (!desc || !buffer) + { + return 0; + } + + size_t required_size = dds_topic_descriptor_serialized_size (desc); + if (buffer_size < required_size) + { + return 0; + } + + uint8_t *ptr = (uint8_t *)buffer; + + // Serialize fixed fields + memcpy (ptr, &desc->m_size, sizeof (desc->m_size)); + ptr += sizeof (desc->m_size); + + memcpy (ptr, &desc->m_align, sizeof (desc->m_align)); + ptr += sizeof (desc->m_align); + + memcpy (ptr, &desc->m_flagset, sizeof (desc->m_flagset)); + ptr += sizeof (desc->m_flagset); + + memcpy (ptr, &desc->m_nkeys, sizeof (desc->m_nkeys)); + ptr += sizeof (desc->m_nkeys); + + const uint32_t m_ops_len = dds_stream_countops (desc->m_ops, desc->m_nkeys, desc->m_keys); + memcpy (ptr, &m_ops_len, sizeof (m_ops_len)); + ptr += sizeof (m_ops_len); + + memcpy (ptr, &desc->restrict_data_representation, sizeof (desc->restrict_data_representation)); + ptr += sizeof (desc->restrict_data_representation); + + // Serialize m_typename + uint32_t typename_len = desc->m_typename ? (uint32_t)strlen (desc->m_typename) + 1 : 0; + memcpy (ptr, &typename_len, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + if (typename_len > 0) + { + memcpy (ptr, desc->m_typename, typename_len); + ptr += typename_len; + } + + // Serialize m_keys + for (uint32_t i = 0; i < desc->m_nkeys; i++) + { + memcpy (ptr, &desc->m_keys[i].m_offset, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + memcpy (ptr, &desc->m_keys[i].m_idx, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + + uint32_t name_len = desc->m_keys[i].m_name ? (uint32_t)strlen (desc->m_keys[i].m_name) + 1 : 0; + memcpy (ptr, &name_len, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + if (name_len > 0) + { + memcpy (ptr, desc->m_keys[i].m_name, name_len); + ptr += name_len; + } + } + + // Serialize m_ops + if (m_ops_len > 0) + { + memcpy (ptr, desc->m_ops, m_ops_len * sizeof (uint32_t)); + ptr += m_ops_len * sizeof (uint32_t); + } + + // Serialize m_meta + uint32_t meta_len = desc->m_meta ? (uint32_t)strlen (desc->m_meta) + 1 : 0; + memcpy (ptr, &meta_len, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + if (meta_len > 0) + { + memcpy (ptr, desc->m_meta, meta_len); + ptr += meta_len; + } + + // Serialize type_information + memcpy (ptr, &desc->type_information.sz, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + if (desc->type_information.sz > 0) + { + memcpy (ptr, desc->type_information.data, desc->type_information.sz); + ptr += desc->type_information.sz; + } + + // Serialize type_mapping + memcpy (ptr, &desc->type_mapping.sz, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + if (desc->type_mapping.sz > 0) + { + memcpy (ptr, desc->type_mapping.data, desc->type_mapping.sz); + ptr += desc->type_mapping.sz; + } + + return (size_t)(ptr - (uint8_t *)buffer); +} + +// Deserialize dds_topic_descriptor from buffer +dds_topic_descriptor_t *dds_topic_descriptor_deserialize (const void *buffer, size_t buffer_size) +{ + if (!buffer || buffer_size == 0) + { + return NULL; + } + + dds_topic_descriptor_t *desc = (dds_topic_descriptor_t *)calloc (1, sizeof (dds_topic_descriptor_t)); + if (!desc) + { + return NULL; + } + + const uint8_t *ptr = (const uint8_t *)buffer; + const uint8_t *end = ptr + buffer_size; + + // Deserialize fixed fields + if (ptr + sizeof (desc->m_size) > end) + goto error; + memcpy ((void *)&desc->m_size, ptr, sizeof (desc->m_size)); + ptr += sizeof (desc->m_size); + + if (ptr + sizeof (desc->m_align) > end) + goto error; + memcpy ((void *)&desc->m_align, ptr, sizeof (desc->m_align)); + ptr += sizeof (desc->m_align); + + if (ptr + sizeof (desc->m_flagset) > end) + goto error; + memcpy ((void *)&desc->m_flagset, ptr, sizeof (desc->m_flagset)); + ptr += sizeof (desc->m_flagset); + + if (ptr + sizeof (desc->m_nkeys) > end) + goto error; + memcpy ((void *)&desc->m_nkeys, ptr, sizeof (desc->m_nkeys)); + ptr += sizeof (desc->m_nkeys); + + + if (ptr + sizeof (desc->m_nops) > end) + goto error; + memcpy ((void *)&desc->m_nops, ptr, sizeof (desc->m_nops)); + ptr += sizeof (desc->m_nops); + + if (ptr + sizeof (desc->restrict_data_representation) > end) + goto error; + memcpy ((void *)&desc->restrict_data_representation, ptr, sizeof (desc->restrict_data_representation)); + ptr += sizeof (desc->restrict_data_representation); + + // Deserialize m_typename + uint32_t typename_len; + if (ptr + sizeof (uint32_t) > end) + goto error; + memcpy (&typename_len, ptr, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + if (typename_len > 0) + { + if (ptr + typename_len > end) + goto error; + char *typename = (char *)malloc (typename_len); + if (!typename) + goto error; + memcpy (typename, ptr, typename_len); + *(char **)&desc->m_typename = typename; + ptr += typename_len; + } + + // Deserialize m_keys + if (desc->m_nkeys > 0) + { + dds_key_descriptor_t *keys = (dds_key_descriptor_t *)calloc (desc->m_nkeys, sizeof (dds_key_descriptor_t)); + if (!keys) + goto error; + *(dds_key_descriptor_t **)&desc->m_keys = keys; + + for (uint32_t i = 0; i < desc->m_nkeys; i++) + { + if (ptr + sizeof (uint32_t) * 2 > end) + goto error; + memcpy (&keys[i].m_offset, ptr, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + memcpy (&keys[i].m_idx, ptr, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + + uint32_t name_len; + if (ptr + sizeof (uint32_t) > end) + goto error; + memcpy (&name_len, ptr, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + if (name_len > 0) + { + if (ptr + name_len > end) + goto error; + char *name = (char *)malloc (name_len); + if (!name) + goto error; + memcpy (name, ptr, name_len); + keys[i].m_name = name; + ptr += name_len; + } + } + } + + // Deserialize m_ops + if (desc->m_nops > 0) + { + if (ptr + desc->m_nops * sizeof (uint32_t) > end) + goto error; + uint32_t *ops = (uint32_t *)malloc (desc->m_nops * sizeof (uint32_t)); + if (!ops) + goto error; + memcpy (ops, ptr, desc->m_nops * sizeof (uint32_t)); + *(uint32_t **)&desc->m_ops = ops; + ptr += desc->m_nops * sizeof (uint32_t); + } + + // Deserialize m_meta + uint32_t meta_len; + if (ptr + sizeof (uint32_t) > end) + goto error; + memcpy (&meta_len, ptr, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + if (meta_len > 0) + { + if (ptr + meta_len > end) + goto error; + char *meta = (char *)malloc (meta_len); + if (!meta) + goto error; + memcpy (meta, ptr, meta_len); + *(char **)&desc->m_meta = meta; + ptr += meta_len; + } + + // Deserialize type_information + uint32_t type_info_sz; + if (ptr + sizeof (uint32_t) > end) + goto error; + memcpy (&type_info_sz, ptr, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + *(uint32_t *)&desc->type_information.sz = type_info_sz; + if (type_info_sz > 0) + { + if (ptr + type_info_sz > end) + goto error; + unsigned char *type_info_data = (unsigned char *)malloc (type_info_sz); + if (!type_info_data) + goto error; + memcpy (type_info_data, ptr, type_info_sz); + *(const unsigned char **)&desc->type_information.data = type_info_data; + ptr += type_info_sz; + } + + // Deserialize type_mapping + uint32_t type_map_sz; + if (ptr + sizeof (uint32_t) > end) + goto error; + memcpy (&type_map_sz, ptr, sizeof (uint32_t)); + ptr += sizeof (uint32_t); + *(uint32_t *)&desc->type_mapping.sz = type_map_sz; + if (type_map_sz > 0) + { + if (ptr + type_map_sz > end) + goto error; + unsigned char *type_map_data = (unsigned char *)malloc (type_map_sz); + if (!type_map_data) + goto error; + memcpy (type_map_data, ptr, type_map_sz); + *(const unsigned char **)&desc->type_mapping.data = type_map_data; + } + assert (ptr == end); + + return desc; + +error: + dds_topic_descriptor_free (desc); + return NULL; +} + +// Free deserialized dds_topic_descriptor +void dds_topic_descriptor_free (dds_topic_descriptor_t *desc) +{ + if (!desc) + { + return; + } + + free ((void *)desc->m_typename); + + if (desc->m_keys) + { + for (uint32_t i = 0; i < desc->m_nkeys; i++) + { + free ((void *)desc->m_keys[i].m_name); + } + free ((void *)desc->m_keys); + } + + free ((void *)desc->m_ops); + free ((void *)desc->m_meta); + free ((void *)desc->type_information.data); + free ((void *)desc->type_mapping.data); + + free (desc); +} diff --git a/examples/recorder_and_replayer/dds_topic_descriptor_serde.h b/examples/recorder_and_replayer/dds_topic_descriptor_serde.h new file mode 100644 index 0000000000..29e529c2f3 --- /dev/null +++ b/examples/recorder_and_replayer/dds_topic_descriptor_serde.h @@ -0,0 +1,58 @@ +/* + * Copyright(c) 2025 ZettaScale Technology and others + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v. 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Eclipse Distribution License + * v. 1.0 which is available at + * http://www.eclipse.org/org/documents/edl-v10.php. + * + * SPDX-License-Identifier: EPL-2.0 OR BSD-3-Clause + */ + +#ifndef DDS_TOPIC_DESCRIPTOR_SERDE_H +#define DDS_TOPIC_DESCRIPTOR_SERDE_H + +#include + +#include "dds/ddsc/dds_public_impl.h" + +#if defined(__cplusplus) +extern "C" { +#endif + +/** + * @brief Calculate the serialized size of a topic descriptor + * @param[in] desc Topic descriptor to measure + * @return Size in bytes needed to serialize the descriptor + */ +size_t dds_topic_descriptor_serialized_size (const dds_topic_descriptor_t *desc); + +/** + * @brief Serialize a topic descriptor to a buffer + * @param[in] desc Topic descriptor to serialize + * @param[out] buffer Buffer to write serialized data + * @param[in] buffer_size Size of the buffer + * @return Number of bytes written, or 0 on error + */ +size_t dds_topic_descriptor_serialize (const dds_topic_descriptor_t *desc, void *buffer, size_t buffer_size); + +/** + * @brief Deserialize a topic descriptor from a buffer + * @param[in] buffer Buffer containing serialized topic descriptor + * @param[in] buffer_size Size of the buffer + * @return Newly allocated topic descriptor, or NULL on error + */ +dds_topic_descriptor_t *dds_topic_descriptor_deserialize (const void *buffer, size_t buffer_size); + +/** + * @brief Free a deserialized topic descriptor + * @param[in] desc Topic descriptor to free (can be NULL) + */ +void dds_topic_descriptor_free (dds_topic_descriptor_t *desc); + +#if defined(__cplusplus) +} +#endif + +#endif /* DDS_TOPIC_DESCRIPTOR_SERDE_H */