From ff2cfb9c016f566f4ca04fe2209996c347e14595 Mon Sep 17 00:00:00 2001 From: Michael Orlov Date: Sat, 15 Apr 2023 13:38:26 -0700 Subject: [PATCH 1/4] Expose Writer::close() as public API and add Recorder::stop() API Signed-off-by: Michael Orlov Co-authored-by: aditya --- rosbag2_cpp/include/rosbag2_cpp/writer.hpp | 5 +++++ rosbag2_cpp/src/rosbag2_cpp/writer.cpp | 5 +++++ .../include/rosbag2_transport/recorder.hpp | 5 +++++ .../src/rosbag2_transport/recorder.cpp | 17 +++++++++++++++-- 4 files changed, 30 insertions(+), 2 deletions(-) diff --git a/rosbag2_cpp/include/rosbag2_cpp/writer.hpp b/rosbag2_cpp/include/rosbag2_cpp/writer.hpp index c1cb03c54f..4a738ee226 100644 --- a/rosbag2_cpp/include/rosbag2_cpp/writer.hpp +++ b/rosbag2_cpp/include/rosbag2_cpp/writer.hpp @@ -88,6 +88,11 @@ class ROSBAG2_CPP_PUBLIC Writer const rosbag2_storage::StorageOptions & storage_options, const ConverterOptions & converter_options = ConverterOptions()); + /** + * \brief Close the current bag file and write metadata.yaml file + */ + void close(); + /** * Create a new topic in the underlying storage. Needs to be called for every topic used within * a message which is passed to write(...). diff --git a/rosbag2_cpp/src/rosbag2_cpp/writer.cpp b/rosbag2_cpp/src/rosbag2_cpp/writer.cpp index ea00baf5aa..355977d03d 100644 --- a/rosbag2_cpp/src/rosbag2_cpp/writer.cpp +++ b/rosbag2_cpp/src/rosbag2_cpp/writer.cpp @@ -65,6 +65,11 @@ void Writer::open( writer_impl_->open(storage_options, converter_options); } +void Writer::close() +{ + writer_impl_->close(); +} + void Writer::create_topic(const rosbag2_storage::TopicMetadata & topic_with_type) { std::lock_guard writer_lock(writer_mutex_); diff --git a/rosbag2_transport/include/rosbag2_transport/recorder.hpp b/rosbag2_transport/include/rosbag2_transport/recorder.hpp index 7f043a0793..bf6c4c0a93 100644 --- a/rosbag2_transport/include/rosbag2_transport/recorder.hpp +++ b/rosbag2_transport/include/rosbag2_transport/recorder.hpp @@ -85,6 +85,11 @@ class Recorder : public rclcpp::Node ROSBAG2_TRANSPORT_PUBLIC void record(); + /// @brief Stopping recording and closing writer. + /// The record() can be called again after stop(). + ROSBAG2_TRANSPORT_PUBLIC + void stop(); + ROSBAG2_TRANSPORT_PUBLIC const std::unordered_set & topics_using_fallback_qos() const; diff --git a/rosbag2_transport/src/rosbag2_transport/recorder.cpp b/rosbag2_transport/src/rosbag2_transport/recorder.cpp index df6b5f6aeb..bc54f84cc2 100644 --- a/rosbag2_transport/src/rosbag2_transport/recorder.cpp +++ b/rosbag2_transport/src/rosbag2_transport/recorder.cpp @@ -57,6 +57,10 @@ class RecorderImpl void record(); + /// @brief Stopping recording and closing writer. + /// The record() can be called again after stop(). + void stop(); + const rosbag2_cpp::Writer & get_writer_handle(); /// Pause the recording. @@ -195,6 +199,11 @@ RecorderImpl::~RecorderImpl() } } +void RecorderImpl::stop() +{ + writer_->close(); +} + void RecorderImpl::record() { topic_qos_profile_overrides_ = record_options_.topic_qos_profile_overrides; @@ -615,12 +624,16 @@ Recorder::Recorder( Recorder::~Recorder() {} -void -Recorder::record() +void Recorder::record() { pimpl_->record(); } +void Recorder::stop() +{ + pimpl_->stop(); +} + const std::unordered_set & Recorder::topics_using_fallback_qos() const { From 56e34873c55b7b8efee89f6a3657d2f6ab36b537 Mon Sep 17 00:00:00 2001 From: Michael Orlov Date: Sat, 15 Apr 2023 13:40:02 -0700 Subject: [PATCH 2/4] Move routines from recorder's destructor to RecorderImpl::stop() Signed-off-by: Michael Orlov --- .../src/rosbag2_transport/recorder.cpp | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/rosbag2_transport/src/rosbag2_transport/recorder.cpp b/rosbag2_transport/src/rosbag2_transport/recorder.cpp index bc54f84cc2..c6d5ebef5e 100644 --- a/rosbag2_transport/src/rosbag2_transport/recorder.cpp +++ b/rosbag2_transport/src/rosbag2_transport/recorder.cpp @@ -182,12 +182,19 @@ RecorderImpl::RecorderImpl( RecorderImpl::~RecorderImpl() { keyboard_handler_->delete_key_press_callback(toggle_paused_key_callback_handle_); + stop(); +} + + +void RecorderImpl::stop() +{ stop_discovery_ = true; if (discovery_future_.valid()) { discovery_future_.wait(); } - + paused_ = true; subscriptions_.clear(); + writer_->close(); // Call writer->close() to finalize current bag file and write metadata { std::lock_guard lock(event_publisher_thread_mutex_); @@ -199,13 +206,9 @@ RecorderImpl::~RecorderImpl() } } -void RecorderImpl::stop() -{ - writer_->close(); -} - void RecorderImpl::record() { + paused_ = record_options_.start_paused; topic_qos_profile_overrides_ = record_options_.topic_qos_profile_overrides; if (record_options_.rmw_serialization_format.empty()) { throw std::runtime_error("No serialization format specified!"); From 75fabe6bb15dec0c0c583ab95cfc70dc9e355f95 Mon Sep 17 00:00:00 2001 From: Michael Orlov Date: Wed, 19 Apr 2023 20:37:05 -0700 Subject: [PATCH 3/4] Add can_record_again_after_stop unit test in rosbag2_transport Signed-off-by: Michael Orlov --- .../mock_sequential_writer.hpp | 12 ++++- .../test/rosbag2_transport/test_record.cpp | 54 ++++++++++++++++++- 2 files changed, 64 insertions(+), 2 deletions(-) diff --git a/rosbag2_transport/test/rosbag2_transport/mock_sequential_writer.hpp b/rosbag2_transport/test/rosbag2_transport/mock_sequential_writer.hpp index 209019c181..fba48d4987 100644 --- a/rosbag2_transport/test/rosbag2_transport/mock_sequential_writer.hpp +++ b/rosbag2_transport/test/rosbag2_transport/mock_sequential_writer.hpp @@ -33,9 +33,13 @@ class MockSequentialWriter : public rosbag2_cpp::writer_interfaces::BaseWriterIn snapshot_mode_ = storage_options.snapshot_mode; (void) storage_options; (void) converter_options; + writer_close_called_ = false; } - void close() override {} + void close() override + { + writer_close_called_ = true; + } void create_topic(const rosbag2_storage::TopicMetadata & topic_with_type) override { @@ -131,6 +135,11 @@ class MockSequentialWriter : public rosbag2_cpp::writer_interfaces::BaseWriterIn return max_messages_per_file_; } + bool closed_was_called() const + { + return writer_close_called_; + } + private: std::unordered_map< std::string, @@ -144,6 +153,7 @@ class MockSequentialWriter : public rosbag2_cpp::writer_interfaces::BaseWriterIn rosbag2_cpp::bag_events::EventCallbackManager callback_manager_; size_t file_number_ = 0; size_t max_messages_per_file_ = 0; + bool writer_close_called_{false}; }; #endif // ROSBAG2_TRANSPORT__MOCK_SEQUENTIAL_WRITER_HPP_ diff --git a/rosbag2_transport/test/rosbag2_transport/test_record.cpp b/rosbag2_transport/test/rosbag2_transport/test_record.cpp index 93e01c3049..baff0b1583 100644 --- a/rosbag2_transport/test/rosbag2_transport/test_record.cpp +++ b/rosbag2_transport/test/rosbag2_transport/test_record.cpp @@ -46,7 +46,7 @@ TEST_F(RecordIntegrationTestFixture, published_messages_from_multiple_topics_are pub_manager.setup_publisher(string_topic, string_message, 2); rosbag2_transport::RecordOptions record_options = - {false, false, {string_topic, array_topic}, "rmw_format", 100ms}; + {false, false, {string_topic, array_topic}, "rmw_format", 50ms}; auto recorder = std::make_shared( std::move(writer_), storage_options_, record_options); recorder->record(); @@ -88,6 +88,58 @@ TEST_F(RecordIntegrationTestFixture, published_messages_from_multiple_topics_are EXPECT_THAT(array_messages[0]->float32_values, Eq(array_message->float32_values)); } +TEST_F(RecordIntegrationTestFixture, can_record_again_after_stop) +{ + auto string_message = get_messages_strings()[1]; + std::string string_topic = "/string_topic"; + + rosbag2_test_common::PublicationManager pub_manager; + pub_manager.setup_publisher(string_topic, string_message, 2); + + rosbag2_transport::RecordOptions record_options = + {false, false, {string_topic}, "rmw_format", 50ms}; + auto recorder = std::make_shared( + std::move(writer_), storage_options_, record_options); + recorder->record(); + + auto & writer = recorder->get_writer_handle(); + auto & mock_writer = dynamic_cast(writer.get_implementation_handle()); + + start_async_spin(recorder); + ASSERT_TRUE(pub_manager.wait_for_matched(string_topic.c_str())); + + pub_manager.run_publishers(); + + EXPECT_FALSE(mock_writer.closed_was_called()); + recorder->stop(); + EXPECT_TRUE(mock_writer.closed_was_called()); + + // Record one more time after stop() + recorder->record(); + + ASSERT_TRUE(pub_manager.wait_for_matched(string_topic.c_str())); + pub_manager.run_publishers(); + + size_t expected_messages = 4; // 4 because was running recorder-record() and publishers twice + auto ret = rosbag2_test_common::wait_until_shutdown( + std::chrono::seconds(5), + [&mock_writer, &expected_messages]() { + return mock_writer.get_messages().size() >= expected_messages; + }); + auto recorded_messages = mock_writer.get_messages(); + EXPECT_TRUE(ret) << "failed to capture expected messages in time"; + EXPECT_THAT(recorded_messages, SizeIs(expected_messages)); + + auto recorded_topics = mock_writer.get_topics(); + ASSERT_THAT(recorded_topics, SizeIs(1)) << "size=" << recorded_topics.size(); + EXPECT_THAT(recorded_topics.at(string_topic).first.serialization_format, Eq("rmw_format")); + ASSERT_THAT(recorded_messages, SizeIs(expected_messages)); + auto string_messages = filter_messages( + recorded_messages, string_topic); + ASSERT_THAT(string_messages, SizeIs(4)); + EXPECT_THAT(string_messages[0]->string_value, Eq(string_message->string_value)); +} + TEST_F(RecordIntegrationTestFixture, qos_is_stored_in_metadata) { auto string_message = get_messages_strings()[1]; From a0713f3a3b45dce6a84448bb3c8084f36c1641a0 Mon Sep 17 00:00:00 2001 From: Michael Orlov Date: Tue, 2 May 2023 14:43:38 -0700 Subject: [PATCH 4/4] Update Doxygen comments for stop() and pause() API in recorder.hpp Signed-off-by: Michael Orlov --- rosbag2_transport/include/rosbag2_transport/recorder.hpp | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/rosbag2_transport/include/rosbag2_transport/recorder.hpp b/rosbag2_transport/include/rosbag2_transport/recorder.hpp index bf6c4c0a93..0b09f3e6bf 100644 --- a/rosbag2_transport/include/rosbag2_transport/recorder.hpp +++ b/rosbag2_transport/include/rosbag2_transport/recorder.hpp @@ -85,8 +85,9 @@ class Recorder : public rclcpp::Node ROSBAG2_TRANSPORT_PUBLIC void record(); - /// @brief Stopping recording and closing writer. - /// The record() can be called again after stop(). + /// @brief Stopping recording. + /// @details The stop() is opposite to the record() operation. It will stop recording, dump + /// all buffers to the disk and close writer. The record() can be called again after stop(). ROSBAG2_TRANSPORT_PUBLIC void stop(); @@ -101,7 +102,8 @@ class Recorder : public rclcpp::Node ROSBAG2_TRANSPORT_PUBLIC const rosbag2_cpp::Writer & get_writer_handle(); - /// Pause the recording. + /// @brief Pause the recording. + /// @details Will keep writer open and skip messages upon arrival on subscriptions. ROSBAG2_TRANSPORT_PUBLIC void pause();