diff --git a/.gitignore b/.gitignore index 826e2b4..eb85c24 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,6 @@ # This file is used to ignore files which are generated # ---------------------------------------------------------------------------- - +.vscode *~ *.autosave *.a diff --git a/data_tamer_cpp/CMakeLists.txt b/data_tamer_cpp/CMakeLists.txt index 9632704..af57a6d 100644 --- a/data_tamer_cpp/CMakeLists.txt +++ b/data_tamer_cpp/CMakeLists.txt @@ -5,6 +5,7 @@ project(data_tamer_cpp VERSION 0.9.4) set(CMAKE_CXX_STANDARD_REQUIRED ON) set(CMAKE_EXPORT_COMPILE_COMMANDS ON) set(CMAKE_CXX_STANDARD 17) +set(CMAKE_WINDOWS_EXPORT_ALL_SYMBOLS ON) if(${CMAKE_PROJECT_NAME} STREQUAL ${PROJECT_NAME}) option(DATA_TAMER_BUILD_TESTS "Build tests" ON) diff --git a/data_tamer_cpp/include/data_tamer/channel.hpp b/data_tamer_cpp/include/data_tamer/channel.hpp index c0fe2c9..f8e3510 100644 --- a/data_tamer_cpp/include/data_tamer/channel.hpp +++ b/data_tamer_cpp/include/data_tamer/channel.hpp @@ -202,6 +202,16 @@ class LogChannel : public std::enable_shared_from_this */ void addDataSink(std::shared_ptr sink); + /** + * @brief removeDataSink remove a sink, i.e. a class collecting our snapshots. + */ + void removeDataSink(std::shared_ptr sink); + + /** + * @brief getNumberOfSink returns the number of registered sinks. + */ + size_t getNumberOfSinks() const; + /** * @brief takeSnapshot copies the current value of all your registered values * and send an instance of Snapshot to all your Sinks. diff --git a/data_tamer_cpp/include/data_tamer/contrib/SerializeMe.hpp b/data_tamer_cpp/include/data_tamer/contrib/SerializeMe.hpp index 413332e..98a0575 100644 --- a/data_tamer_cpp/include/data_tamer/contrib/SerializeMe.hpp +++ b/data_tamer_cpp/include/data_tamer/contrib/SerializeMe.hpp @@ -477,9 +477,10 @@ inline void SerializeIntoBuffer(SpanBytes& buffer, T const& value) throw std::runtime_error("SerializeIntoBuffer: buffer overflow"); } #if SERIALIZE_LITTLEENDIAN == 0 - *(reinterpret_cast(buffer.data())) = EndianSwap(value); + T swapped = EndianSwap(value); + std::memcpy(buffer.data(), &swapped, S); #else - *(reinterpret_cast(buffer.data())) = value; + std::memcpy(buffer.data(), &value, S); #endif buffer.trimFront(S); // NOLINT } diff --git a/data_tamer_cpp/include/data_tamer/data_sink.hpp b/data_tamer_cpp/include/data_tamer/data_sink.hpp index 8474262..a921c4d 100644 --- a/data_tamer_cpp/include/data_tamer/data_sink.hpp +++ b/data_tamer_cpp/include/data_tamer/data_sink.hpp @@ -96,6 +96,12 @@ class DataSinkBase void stopThread(); + void stopAcceptingSnapshots(); + + void processQueuedSnapshots(); + + void startAcceptingSnapshots(); + private: struct Pimpl; std::unique_ptr _p; diff --git a/data_tamer_cpp/include/data_tamer/details/locked_reference.hpp b/data_tamer_cpp/include/data_tamer/details/locked_reference.hpp index db3afb7..e5b2157 100644 --- a/data_tamer_cpp/include/data_tamer/details/locked_reference.hpp +++ b/data_tamer_cpp/include/data_tamer/details/locked_reference.hpp @@ -41,6 +41,7 @@ class LockedRef { std::swap(ref_, other.ref_); std::swap(mutex_, other.mutex_); + return *this; } operator bool() const { return ref_ != nullptr; } diff --git a/data_tamer_cpp/include/data_tamer/sinks/mcap_sink.hpp b/data_tamer_cpp/include/data_tamer/sinks/mcap_sink.hpp index 059a9fc..dfde566 100644 --- a/data_tamer_cpp/include/data_tamer/sinks/mcap_sink.hpp +++ b/data_tamer_cpp/include/data_tamer/sinks/mcap_sink.hpp @@ -49,6 +49,10 @@ class MCAPSink : public DataSinkBase /// Stop recording and save the file void stopRecording(); + /// Stop taking snapshots, finish the existing queue, then `stopRecording` + /// will block for at least 250 us to ensure the queue is empty + void finishQueueAndStop(); + /** * @brief restartRecording saves the current file (unless we did it already, * calling stopRecording) and start recording into a new one. diff --git a/data_tamer_cpp/include/data_tamer_parser/data_tamer_parser.hpp b/data_tamer_cpp/include/data_tamer_parser/data_tamer_parser.hpp index 2c680cf..68eed2e 100644 --- a/data_tamer_cpp/include/data_tamer_parser/data_tamer_parser.hpp +++ b/data_tamer_cpp/include/data_tamer_parser/data_tamer_parser.hpp @@ -209,7 +209,7 @@ inline bool GetBit(BufferSpan mask, size_t index) return hash; } -bool TypeField::operator==(const TypeField& other) const +inline bool TypeField::operator==(const TypeField& other) const { return is_vector == other.is_vector && type == other.type && array_size == other.array_size && field_name == other.field_name && diff --git a/data_tamer_cpp/src/channel.cpp b/data_tamer_cpp/src/channel.cpp index 4acf480..c088b60 100644 --- a/data_tamer_cpp/src/channel.cpp +++ b/data_tamer_cpp/src/channel.cpp @@ -31,6 +31,7 @@ struct LogChannel::Pimpl Schema schema; bool logging_started = false; + mutable Mutex sinks_mutex; std::unordered_set> sinks; }; @@ -159,9 +160,29 @@ void LogChannel::unregister(const RegistrationID& id) void LogChannel::addDataSink(std::shared_ptr sink) { + std::lock_guard const lock_sinks(_p->sinks_mutex); + + // if we haven't already started logging, then takeSnapshot() handles adding the channel + // otherwise it must be done here so the sink knows about the existing schema + if (_p->logging_started) + { + sink->addChannel(_p->channel_name, _p->schema); + } _p->sinks.insert(sink); } +void LogChannel::removeDataSink(std::shared_ptr sink) +{ + std::lock_guard const lock(_p->sinks_mutex); + _p->sinks.erase(sink); +} + +size_t LogChannel::getNumberOfSinks() const +{ + std::lock_guard const lock(_p->sinks_mutex); + return _p->sinks.size(); +} + Schema LogChannel::getSchema() const { std::lock_guard const lock(_p->mutex); @@ -182,12 +203,16 @@ void LogChannel::addCustomType(const std::string& custom_type_name, bool LogChannel::takeSnapshot(std::chrono::nanoseconds timestamp) { { - std::lock_guard const lock(_p->mutex); - + std::lock_guard const lock_sinks(_p->sinks_mutex); if (_p->sinks.empty()) { return false; } + } + + { + std::lock_guard const lock(_p->mutex); + // update the _p->snapshot.active_mask if necessary if (_p->mask_dirty) { @@ -214,11 +239,15 @@ bool LogChannel::takeSnapshot(std::chrono::nanoseconds timestamp) } _p->snapshot.payload.resize(payload_size); - // call sink->addChannel (usually done once) + // set up the channel if we haven't begun logging if (!_p->logging_started) { - _p->logging_started = true; _p->snapshot.schema_hash = _p->schema.hash; + + std::lock_guard const lock_sinks(_p->sinks_mutex); + // start logging inside the sinks_mutex so that addDataSink does not have an + // incorrect value due to a race condition + _p->logging_started = true; for (auto const& sink : _p->sinks) { sink->addChannel(_p->channel_name, _p->schema); @@ -242,9 +271,12 @@ bool LogChannel::takeSnapshot(std::chrono::nanoseconds timestamp) } bool all_pushed = true; - for (auto& sink : _p->sinks) { - all_pushed &= sink->pushSnapshot(_p->snapshot); + std::lock_guard const lock_sinks(_p->sinks_mutex); + for (auto& sink : _p->sinks) + { + all_pushed &= sink->pushSnapshot(_p->snapshot); + } } return all_pushed; } diff --git a/data_tamer_cpp/src/data_sink.cpp b/data_tamer_cpp/src/data_sink.cpp index 601345f..df5fd0c 100644 --- a/data_tamer_cpp/src/data_sink.cpp +++ b/data_tamer_cpp/src/data_sink.cpp @@ -29,6 +29,7 @@ struct DataSinkBase::Pimpl std::thread thread; std::atomic_bool run = true; + std::atomic_bool accept_snapshots = true; moodycamel::ConcurrentQueue queue; }; @@ -41,7 +42,33 @@ DataSinkBase::~DataSinkBase() bool DataSinkBase::pushSnapshot(const Snapshot& snapshot) { - return _p->queue.enqueue(snapshot); + if(_p->accept_snapshots) + { + return _p->queue.enqueue(snapshot); + } + else + { + return false; + } +} + +void DataSinkBase::stopAcceptingSnapshots() +{ + _p->accept_snapshots = false; +} + +void DataSinkBase::startAcceptingSnapshots() +{ + _p->accept_snapshots = true; +} + +void DataSinkBase::processQueuedSnapshots() +{ + Snapshot snapshot_copy; + while(_p->queue.try_dequeue(snapshot_copy)) + { + this->storeSnapshot(snapshot_copy); + } } void DataSinkBase::stopThread() diff --git a/data_tamer_cpp/src/sinks/mcap_sink.cpp b/data_tamer_cpp/src/sinks/mcap_sink.cpp index f3d40a9..caa61c1 100644 --- a/data_tamer_cpp/src/sinks/mcap_sink.cpp +++ b/data_tamer_cpp/src/sinks/mcap_sink.cpp @@ -3,6 +3,8 @@ #include #include +#include +#include #ifndef USING_ROS2 #define MCAP_IMPLEMENTATION @@ -138,6 +140,22 @@ void MCAPSink::stopRecording() writer_.reset(); } +void MCAPSink::finishQueueAndStop() +{ + // stop accepting new snapshots + stopAcceptingSnapshots(); + + // finish any that are queued + processQueuedSnapshots(); + + // sleep and process any that were missed by previous processing + std::this_thread::sleep_for(std::chrono::microseconds(250)); + processQueuedSnapshots(); + + // now stop the recording as normal + stopRecording(); +} + void MCAPSink::restartRecording(const std::string& filepath, bool do_compression) { std::scoped_lock lk(mutex_); @@ -150,6 +168,9 @@ void MCAPSink::restartRecording(const std::string& filepath, bool do_compression { addChannel(name, schema); } + + // start accepting snapshots again in case they were stopped + startAcceptingSnapshots(); } } // namespace DataTamer diff --git a/data_tamer_cpp/tests/CMakeLists.txt b/data_tamer_cpp/tests/CMakeLists.txt index b56949c..b5197d5 100644 --- a/data_tamer_cpp/tests/CMakeLists.txt +++ b/data_tamer_cpp/tests/CMakeLists.txt @@ -4,7 +4,8 @@ include(GoogleTest) add_executable(datatamer_test dt_tests.cpp custom_types_tests.cpp - parser_tests.cpp) + parser_tests.cpp + add_remove_sink_tests.cpp) gtest_discover_tests(datatamer_test DISCOVERY_MODE PRE_TEST) target_include_directories(datatamer_test diff --git a/data_tamer_cpp/tests/add_remove_sink_tests.cpp b/data_tamer_cpp/tests/add_remove_sink_tests.cpp new file mode 100644 index 0000000..f33450e --- /dev/null +++ b/data_tamer_cpp/tests/add_remove_sink_tests.cpp @@ -0,0 +1,71 @@ +#include "data_tamer/channel.hpp" +#include "data_tamer/sinks/dummy_sink.hpp" + +#include +#include +#include + +using namespace DataTamer; + +void take_snapshots(std::shared_ptr channel, int count) +{ + for(int i = 0; i < count; i++) + { + channel->takeSnapshot(); + std::this_thread::sleep_for(std::chrono::microseconds(50)); + } + std::this_thread::sleep_for(std::chrono::milliseconds(1)); +} + +TEST(DataTamerSinkRegistry, AddSinkIncreasesCountAndRef) +{ + auto channel = LogChannel::create("chan"); + auto sink = std::make_shared(); + channel->addDataSink(sink); + + std::vector dummyData = { 10, 11, 12 }; + channel->registerValue("valsA", &dummyData); + + ASSERT_EQ(channel->getNumberOfSinks(), 1); +} + +TEST(DataTamerSinkRegistry, SnapshotsAreRecordedWhileSinkPresent) +{ + auto channel = LogChannel::create("chan"); + auto sink = std::make_shared(); + channel->addDataSink(sink); + + std::vector dummyData = { 10, 11, 12 }; + channel->registerValue("valsA", &dummyData); + + const int snapshot_count = 10; + take_snapshots(channel, snapshot_count); + + const auto hash = channel->getSchema().hash; + ASSERT_EQ(sink->snapshots_count[hash], snapshot_count); +} + +TEST(DataTamerSinkRegistry, RemoveSinkStopsRecording) +{ + auto channel = LogChannel::create("chan"); + auto sink = std::make_shared(); + channel->addDataSink(sink); + + std::vector dummyData = { 10, 11, 12 }; + channel->registerValue("valsA", &dummyData); + + const int snapshot_count = 10; + take_snapshots(channel, snapshot_count); + + const auto hash = channel->getSchema().hash; + ASSERT_EQ(sink->snapshots_count[hash], snapshot_count); + + channel->removeDataSink(sink); + + ASSERT_EQ(channel->getNumberOfSinks(), 0); + + // Taking more snapshots, should not be recorded in the sink (i.e does not increase snapshots_count) + take_snapshots(channel, snapshot_count); + + ASSERT_EQ(sink->snapshots_count[hash], snapshot_count); +} \ No newline at end of file diff --git a/data_tamer_cpp/tests/dt_tests.cpp b/data_tamer_cpp/tests/dt_tests.cpp index 8366173..d472e57 100644 --- a/data_tamer_cpp/tests/dt_tests.cpp +++ b/data_tamer_cpp/tests/dt_tests.cpp @@ -1,10 +1,12 @@ #include "data_tamer/data_tamer.hpp" #include "data_tamer/sinks/dummy_sink.hpp" +#include "data_tamer/sinks/mcap_sink.hpp" #include "../examples/geometry_types.hpp" #include +#include #include #include #include @@ -270,3 +272,33 @@ TEST(DataTamerBasic, VectorWithChangingSize) ASSERT_EQ(sink->latest_snapshot.payload.size(), vect.size() * sizeof(float) + sizeof(uint32_t)); } + +TEST(DataTamerBasic, FinishQueue) +{ + auto channel = LogChannel::create("chan"); + auto const temp_path = + std::filesystem::temp_directory_path() / std::filesystem::path("data_tamer_test." + "mcap"); + auto sink = std::make_shared(temp_path.string()); + channel->addDataSink(sink); + + double const value = 1.; + channel->registerValue("value", &value); + + EXPECT_TRUE(channel->takeSnapshot()); + + sink->finishQueueAndStop(); + + // now we shouldn't be able to take more snapshots + EXPECT_FALSE(channel->takeSnapshot()); + + // restart the recording + sink->restartRecording(temp_path); + + EXPECT_TRUE(channel->takeSnapshot()); + + sink->stopRecording(); + + // since we just stopped recording but not snapshots, we'll still be able to take a snapshot (but it won't be written to disk) + EXPECT_TRUE(channel->takeSnapshot()); +}