Files
leonarski_fandClaude Opus 5 c40e0eb11d broker: fix deadlock when re-initialising after a DECTRIS run that never started
A re-initialisation left the previous run's ZMQImagePuller connected: JFJochServices::On
replaces its own shared_ptr, but JFJochReceiverService keeps a second reference to the puller
of the last run. Two PULL sockets on one PUSH peer share messages round-robin, so the orphaned
puller silently took half of the DECTRIS stream - the start message included - and the receiver
then waited for a start message that had already been discarded.

Once such a run was cancelled, nothing drained the orphaned puller's outside_fifo any more. Its
CBOR thread parked in PutBlocking on the full queue (suspend is only tested before the put), the
puller thread backed up behind it in cbor_fifo, and neither could reach the disconnect flag
again. The next JFJochReceiverService::Start dropped the last reference to that puller, so
~ZMQImagePuller joined two threads that could never exit - while holding state_mutex, inside a
calibration sequence that itself holds the state machine's mutex. The whole control plane froze
with no way to cancel; only a restart got out of it.

* ThreadSafeFIFO gains Stop(), which releases every waiter and makes further blocking operations
  return at once. Clear() now notifies c_full as well: clearing a full queue used to leave the
  blocked producer asleep, since the next Get on an empty queue notifies no one.
* ZMQImagePuller::Disconnect and TCPImagePuller::Disconnect stop their queues before joining, so
  a puller whose consumer is gone can always shut down.
* JFJochServices::On and ::Off disconnect the previous puller explicitly rather than relying on
  the shared_ptr going away, so two readers never share the detector stream.

ZMQImagePuller_DisconnectWithFullQueue covers the shutdown; it hangs on the previous code.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012bew392LTGP2fkhfRJsMcB
2026-09-11 18:53:39 +02:00

147 lines
5.2 KiB
C++

// 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) {
auto start_time = std::chrono::steady_clock::now();
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);
}
auto end_time = std::chrono::steady_clock::now();
auto duration = std::chrono::duration<float>(end_time - start_time);
logger.Info("ZMQImagePuller connected to {} in {:.3f} s", addr, duration.count());
}
ZMQImagePuller::~ZMQImagePuller() {
ZMQImagePuller::Disconnect();
}
void ZMQImagePuller::Disconnect() {
disconnect = 1;
// The threads below hand messages on through bounded blocking queues, and by the time a puller
// is disconnected nothing drains outside_fifo any more - the receiver that did is gone. A full
// queue would then park CBORThread in PutBlocking for ever, back up cbor_fifo and park
// PullerThread too, and the joins below would never return. Stop() releases them.
cbor_fifo.Stop();
repub_fifo.Stop();
outside_fifo.Stop();
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.zmq_msg = std::make_shared<ZMQMessage>();
bool received = false;
while (!received) {
if (disconnect) {
cbor_fifo.PutBlocking(ImagePullerOutput{});
return;
}
try {
received = socket.Receive(*ret.zmq_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.zmq_msg) {
try {
ret.cbor = CBORStream2Deserialize(ret.zmq_msg->data(), ret.zmq_msg->size());
// Even if we suspend consuming of the messages by the receiver,
// it is still reasonable to have them republished
// so republish functionality is not affected by Suspend()
if (!suspend)
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.zmq_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.zmq_msg->data(), ret.zmq_msg->size(), true);
}
} catch (const JFJochException &e) {
logger.ErrorException(e);
}
ret = repub_fifo.GetBlocking();
}
if (repub_active)
logger.Info("Republish finished");
}
void ZMQImagePuller::Suspend() {
suspend = true;
}
void ZMQImagePuller::ResumeAndClear() {
outside_fifo.Clear();
outside_fifo.ClearMaxUtilization();
suspend = false;
}