// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute // SPDX-License-Identifier: GPL-3.0-only #include "ZMQImagePuller.h" #include "../frame_serialize/CBORStream2Serializer.h" ZMQImagePuller::ZMQImagePuller(const std::string &in_addr, const std::string &repub_address, const std::optional &rcv_watermark, const std::optional &repub_watermark) : socket (ZMQSocketType::Pull), addr(in_addr) { puller_thread = std::thread(&ZMQImagePuller::PullerThread, this); cbor_thread = std::thread(&ZMQImagePuller::CBORThread, this); socket.ReceiveWaterMark(rcv_watermark.value_or(default_receive_watermark)); socket.ReceiveTimeout(ReceiveTimeout); socket.Connect(addr); if (!repub_address.empty()) { repub_socket = std::make_unique(ZMQSocketType::Push); repub_socket->SendWaterMark(repub_watermark.value_or(default_repub_watermark)); repub_socket->SendTimeout(RepubTimeout); repub_socket->Bind(repub_address); repub_thread = std::thread(&ZMQImagePuller::RepubThread, this); } } ZMQImagePuller::~ZMQImagePuller() { Disconnect(); } void ZMQImagePuller::Disconnect() { disconnect = 1; if (puller_thread.joinable()) puller_thread.join(); if (cbor_thread.joinable()) cbor_thread.join(); if (repub_thread.joinable()) repub_thread.join(); if (!addr.empty()) socket.Disconnect(addr); addr = ""; } void ZMQImagePuller::PullerThread() { while (true) { ImagePullerOutput ret; ret.msg = std::make_shared(); bool received = false; while (!received) { if (disconnect) { cbor_fifo.PutBlocking(ImagePullerOutput{}); return; } try { received = socket.Receive(*ret.msg, false); if (!received) std::this_thread::sleep_for(std::chrono::milliseconds(1)); } catch (const JFJochException &e) { logger.ErrorException(e); } } cbor_fifo.PutBlocking(ret); } } void ZMQImagePuller::CBORThread() { auto ret = cbor_fifo.GetBlocking(); while (ret.msg) { try { ret.cbor = CBORStream2Deserialize(ret.msg->data(), ret.msg->size()); if (ret.cbor->msg_type == CBORImageType::END) logger.Info("Received END"); outside_fifo.PutBlocking(ret); if (repub_socket) { if ((ret.cbor->msg_type == CBORImageType::START) || (ret.cbor->msg_type == CBORImageType::END)) repub_fifo.PutBlocking(ret); else repub_fifo.Put(ret); } } catch (const JFJochException &e) { logger.ErrorException(e); } ret = cbor_fifo.GetBlocking(); } if (repub_socket) repub_fifo.PutBlocking(ret); outside_fifo.PutBlocking(ret); } void ZMQImagePuller::RepubThread() { auto ret = repub_fifo.GetBlocking(); bool repub_active = false; while (ret.msg) { try { if (ret.cbor->msg_type == CBORImageType::START) { // Start message needs to be cleaned when running republish StartMessage msg = ret.cbor->start_message.value(); msg.writer_notification_zmq_addr = ""; std::vector serialization_buffer(256*1024*1024); CBORStream2Serializer serializer(serialization_buffer.data(), serialization_buffer.size()); serializer.SerializeSequenceStart(msg); repub_active = repub_socket->Send(serialization_buffer.data(), serializer.GetBufferSize(), true); if (repub_active) logger.Info("Republish active"); } else { if (repub_active) repub_socket->Send(ret.msg->data(), ret.msg->size(), true); } } catch (const JFJochException &e) { logger.ErrorException(e); } ret = repub_fifo.GetBlocking(); } if (repub_active) logger.Info("Republish finished"); }