From f7704173bc0fb9a0e6fea120509ebe0388a0cd2b Mon Sep 17 00:00:00 2001 From: leonarski_f Date: Mon, 2 Mar 2026 14:45:41 +0100 Subject: [PATCH] ZMQStream2Pusher: Add error message for writer died during data collection --- image_pusher/TCPStreamPusher.cpp | 4 ++++ image_pusher/TCPStreamPusher.h | 2 +- image_pusher/ZMQStream2Pusher.cpp | 4 ++++ image_pusher/ZMQStream2Pusher.h | 1 + 4 files changed, 10 insertions(+), 1 deletion(-) diff --git a/image_pusher/TCPStreamPusher.cpp b/image_pusher/TCPStreamPusher.cpp index c013d3cc..a369164a 100644 --- a/image_pusher/TCPStreamPusher.cpp +++ b/image_pusher/TCPStreamPusher.cpp @@ -27,6 +27,7 @@ void TCPStreamPusher::StartDataCollection(StartMessage &message) { images_per_file = message.images_per_file; run_number = message.run_number; run_name = message.run_name; + transmission_error = false; for (size_t i = 0; i < socket.size(); i++) { if (!socket[i]->AcceptConnection(std::chrono::seconds(5))) @@ -93,11 +94,14 @@ bool TCPStreamPusher::EndDataCollection(const EndMessage &message) { else if (!s->Send(serialization_buffer.data(), serializer.GetBufferSize(), TCPFrameType::END)) ret = false; } + transmission_error = !ret; return ret; } std::string TCPStreamPusher::Finalize() { std::string ret; + if (transmission_error) + ret += "Timeout sending images (e.g., writer disabled during data collection);"; if (writer_notification_socket) { for (size_t i = 0; i < socket.size(); i++) { auto n = writer_notification_socket->Receive(run_number, run_name); diff --git a/image_pusher/TCPStreamPusher.h b/image_pusher/TCPStreamPusher.h index 170490f5..3d574c22 100644 --- a/image_pusher/TCPStreamPusher.h +++ b/image_pusher/TCPStreamPusher.h @@ -16,7 +16,7 @@ class TCPStreamPusher : public ImagePusher { int64_t images_per_file = 1; uint64_t run_number = 0; std::string run_name; - + std::atomic transmission_error = false; public: explicit TCPStreamPusher(const std::vector& addr, std::optional send_buffer_size = {}, diff --git a/image_pusher/ZMQStream2Pusher.cpp b/image_pusher/ZMQStream2Pusher.cpp index 94ab1b00..77817691 100644 --- a/image_pusher/ZMQStream2Pusher.cpp +++ b/image_pusher/ZMQStream2Pusher.cpp @@ -43,6 +43,7 @@ void ZMQStream2Pusher::StartDataCollection(StartMessage& message) { images_per_file = message.images_per_file; run_number = message.run_number; run_name = message.run_name; + transmission_error = false; for (int i = 0; i < socket.size(); i++) { message.socket_number = i; @@ -77,6 +78,7 @@ bool ZMQStream2Pusher::EndDataCollection(const EndMessage& message) { if (!s->Send(serialization_buffer.data(), serializer.GetBufferSize())) ret = false; } + transmission_error = !ret; return ret; } @@ -89,6 +91,8 @@ std::vector ZMQStream2Pusher::GetAddress() { std::string ZMQStream2Pusher::Finalize() { std::string ret; + if (transmission_error) + ret += "Timeout sending images (e.g., writer disabled during data collection);"; if (writer_notification_socket) { for (int i = 0; i < socket.size(); i++) { auto n = writer_notification_socket->Receive(run_number, run_name); diff --git a/image_pusher/ZMQStream2Pusher.h b/image_pusher/ZMQStream2Pusher.h index d466898a..5007c5d4 100644 --- a/image_pusher/ZMQStream2Pusher.h +++ b/image_pusher/ZMQStream2Pusher.h @@ -21,6 +21,7 @@ class ZMQStream2Pusher : public ImagePusher { int64_t images_per_file = 1; uint64_t run_number = 0; std::string run_name; + std::atomic transmission_error = false; public: explicit ZMQStream2Pusher(const std::vector& addr, std::optional send_buffer_high_watermark = {},