diff --git a/modules/ddscom/src/dds_participant.cpp b/modules/ddscom/src/dds_participant.cpp index e7f93115..cbc70916 100644 --- a/modules/ddscom/src/dds_participant.cpp +++ b/modules/ddscom/src/dds_participant.cpp @@ -518,7 +518,6 @@ namespace dls info.type_information.type_information.complete().typeid_with_size().type_id(), remote_type_object)) { - std::cout << "Cannot get the remote type information" << std::endl; return; } diff --git a/modules/supervisor/include/dls2/supervisor/orchestrator_base.hpp b/modules/supervisor/include/dls2/supervisor/orchestrator_base.hpp index 71c7deb7..fe0d2d38 100644 --- a/modules/supervisor/include/dls2/supervisor/orchestrator_base.hpp +++ b/modules/supervisor/include/dls2/supervisor/orchestrator_base.hpp @@ -55,6 +55,8 @@ namespace dls */ virtual void telemetryMain(const std::vector&) {}; + void activation() override; + // Events std::mutex event_mutex_; logging::EventListener event_listener_; diff --git a/modules/supervisor/include/dls2/supervisor/telemetry_base.hpp b/modules/supervisor/include/dls2/supervisor/telemetry_base.hpp index 4626f131..364de7cf 100644 --- a/modules/supervisor/include/dls2/supervisor/telemetry_base.hpp +++ b/modules/supervisor/include/dls2/supervisor/telemetry_base.hpp @@ -18,6 +18,7 @@ namespace dls public: virtual ~ReaderBindingBase() = default; virtual void copyToField(void* target) = 0; + virtual bool isValid() const = 0; }; template @@ -33,6 +34,10 @@ namespace dls typed->*field_ = reader_->msg; } + bool isValid() const override { + return static_cast(reader_); + } + private: ReaderPtrT reader_; MsgT TargetT::* field_; @@ -43,6 +48,7 @@ namespace dls public: virtual ~WriterBindingBase() = default; virtual void copyFromField(const void* source) = 0; + virtual bool isValid() const = 0; }; template @@ -58,6 +64,10 @@ namespace dls writer_->msg = typed->*field_; } + bool isValid() const override { + return static_cast(writer_); + } + private: WriterPtrT writer_; MsgT SourceT::* field_; diff --git a/modules/supervisor/include/dls2/supervisor/telemetry_base.tpp b/modules/supervisor/include/dls2/supervisor/telemetry_base.tpp index c2e7de8e..97af04ba 100644 --- a/modules/supervisor/include/dls2/supervisor/telemetry_base.tpp +++ b/modules/supervisor/include/dls2/supervisor/telemetry_base.tpp @@ -2,6 +2,8 @@ #include "dls2/supervisor/telemetry_base.hpp" +#include + using namespace dls; inline TelemetryBase::TelemetryBase( @@ -28,14 +30,36 @@ void TelemetryBase::tick(InputT& input, OutputT& output) { { std::lock_guard lock(input.mutex); - for (auto& b : reader_bindings_) { + if (reader_bindings_.size() != readers_.size()) { + std::cerr << "[telemetry tick] reader bindings/readers size mismatch: bindings=" + << reader_bindings_.size() << " readers=" << readers_.size() << std::endl; + } + + for (size_t i = 0; i < reader_bindings_.size(); ++i) { + auto& b = reader_bindings_[i]; + if (!b || !b->isValid() || i >= readers_.size() || !readers_[i]) { + std::cerr << "[telemetry tick] invalid reader binding at index " << i << std::endl; + continue; + } + b->copyToField(&input); - }; + } } { std::lock_guard lock(output.mutex); - for (auto& b : writer_bindings_) { + if (writer_bindings_.size() != writers_.size()) { + std::cerr << "[telemetry tick] writer bindings/writers size mismatch: bindings=" + << writer_bindings_.size() << " writers=" << writers_.size() << std::endl; + } + + for (size_t i = 0; i < writer_bindings_.size(); ++i) { + auto& b = writer_bindings_[i]; + if (!b || !b->isValid() || i >= writers_.size() || !writers_[i]) { + std::cerr << "[telemetry tick] invalid writer binding at index " << i << std::endl; + continue; + } + b->copyFromField(&output); - }; + } } } \ No newline at end of file diff --git a/modules/supervisor/src/orchestrator_base.cpp b/modules/supervisor/src/orchestrator_base.cpp index b4d687b0..90b05d9c 100644 --- a/modules/supervisor/src/orchestrator_base.cpp +++ b/modules/supervisor/src/orchestrator_base.cpp @@ -13,6 +13,9 @@ namespace dls , event_to_publish_(event_to_publish) , telemetry_manager_(telemetry_readers_, telemetry_writers_) , sm_(sm) + {}; + + void OrchestratorBase::activation() { if(!this->telemetry_started_.load()){ try{ @@ -26,14 +29,23 @@ namespace dls throw; } } - }; + dls::App::activation(); + } void OrchestratorBase::telemetryCallback() { while (!should_quit) { - // Reading msgs from Control Station - for(auto& reader : telemetry_readers_){ + for (size_t i = 0; i < telemetry_readers_.size(); ++i) { + auto& reader = telemetry_readers_[i]; + if (!reader) { + std::cerr << "[orchestrator telemetry] null reader at index " << i + << " for " << this->getID() << std::endl; + continue; + } + if(!reader->is_receiving_data()){ + continue; + } reader->read(); } @@ -53,8 +65,14 @@ namespace dls telemetryMain(events_to_publish); - // Sending msgs to Control Station - for(auto& writer : telemetry_writers_){ + for (size_t i = 0; i < telemetry_writers_.size(); ++i) { + auto& writer = telemetry_writers_[i]; + if (!writer) { + std::cerr << "[orchestrator telemetry] null writer at index " << i + << " for " << this->getID() << std::endl; + continue; + } + writer->publish(); } @@ -64,26 +82,38 @@ namespace dls void OrchestratorBase::run(const std::chrono::system_clock::time_point &time) { - read(); - - // Collecting events from DLS2 - static long int idx_read = 0; - const auto events_fifo = event_listener_.readEvents(idx_read); - EventsPriorityQueue events_priority_queue_tmp; + try { - // Update internal events representation - std::lock_guard lock(event_mutex_); + read(); - for(const auto& event : events_fifo){ - events_priority_queue_.push(event); - } + // Collecting events from DLS2 + static long int idx_read = 0; + const auto events_fifo = event_listener_.readEvents(idx_read); + EventsPriorityQueue events_priority_queue_tmp; + { + // Update internal events representation + std::lock_guard lock(event_mutex_); - events_priority_queue_tmp = events_priority_queue_; - } + for(const auto& event : events_fifo){ + events_priority_queue_.push(event); + } - orchestrate(time, events_priority_queue_tmp); + events_priority_queue_tmp = events_priority_queue_; + } - write(); + orchestrate(time, events_priority_queue_tmp); + write(); + } + catch (const std::exception& e) + { + std::cerr << "[orchestrator run] exception in " << this->getID() << ": " << e.what() << std::endl; + throw; + } + catch (...) + { + std::cerr << "[orchestrator run] unknown exception in " << this->getID() << std::endl; + throw; + } } extern "C" PeriodicAppPlugin *create(size_t telemetry_thread_period_ms, size_t event_to_publish, const std::string& ID)