v1.0.0-rc.36
This commit is contained in:
121
image_puller/ZMQImagePuller.cpp
Normal file
121
image_puller/ZMQImagePuller.cpp
Normal file
@@ -0,0 +1,121 @@
|
||||
// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
|
||||
// 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<int32_t> &rcv_watermark,
|
||||
const std::optional<int32_t> &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<ZMQSocket>(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<ZMQMessage>();
|
||||
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());
|
||||
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<uint8_t> 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");
|
||||
}
|
||||
Reference in New Issue
Block a user