Files
Jungfraujoch/tests/ZMQImagePusherTest.cpp
leonarski_fandClaude Opus 5 db5068b2f8 tests: a puller whose queue filled must recover, not go silent
Covers the half of the deadlock fix that decides whether a collection poisons the ones after
it. outside_fifo fills whenever the detector keeps streaming past the point where the run
stopped consuming - the analysis threads have gone, and Suspend() only arrives once Stop() has
finished with the writer - and the CBOR thread then parks in PutBlocking on the full queue.
Suspend() cannot release it, since the flag is tested before the put, so ResumeAndClear at the
next start has to; nothing else ever would, because a Get on an empty queue notifies no
producer. Without that notify the puller is dead for good: cbor_fifo fills behind it, the
puller thread blocks in turn, and every later collection receives nothing at all - no images
and no start message - while the puller itself looks healthy.

The test fails with the c_full notify in ThreadSafeFIFO::Clear removed.

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

423 lines
16 KiB
C++

// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include <random>
#include <catch2/catch_all.hpp>
#include "../image_puller/ZMQImagePuller.h"
#include "../image_pusher/ZMQStream2Pusher.h"
#include "../common/DiffractionSpot.h"
void test_puller(ZMQImagePuller *puller,
const DiffractionExperiment& x,
const std::vector<uint16_t> &image1,
int64_t nwriter,
int64_t writer_id,
std::vector<size_t> &diff_split,
std::vector<size_t> &diff_size,
std::vector<size_t> &diff_content,
std::vector<size_t> &nimages) {
auto timeout = std::chrono::minutes(3);
auto img = puller->PollImage(timeout);
if (!img || !img->cbor || !img->cbor->start_message) {
diff_content[writer_id]++;
return;
}
if ((!img->cbor->start_message->write_master_file) || (img->cbor->start_message->write_master_file.value() != (writer_id == 0)))
diff_content[writer_id]++;
img = puller->PollImage(timeout);
while (img && img->cbor && !img->cbor->end_message) {
if (img->cbor->data_message) {
if ((nwriter > 1) && ((img->cbor->data_message->number / 16) % nwriter != writer_id))
diff_split[writer_id]++;
if (img->cbor->data_message->image.GetCompressedSize() != x.GetPixelsNum() * sizeof(uint16_t))
diff_size[writer_id]++;
else if (memcmp(img->cbor->data_message->image.GetCompressed(), image1.data() + img->cbor->data_message->number * x.GetPixelsNum(),
x.GetPixelsNum() * sizeof(uint16_t)) != 0)
diff_content[writer_id]++;
if (img->cbor->data_message->image.GetWidth() != RAW_MODULE_COLS)
diff_content[writer_id]++;
if (img->cbor->data_message->image.GetHeight() != RAW_MODULE_LINES)
diff_content[writer_id]++;
if (img->cbor->data_message->image.GetByteDepth() != 2)
diff_content[writer_id]++;
if (img->cbor->data_message->image.GetCompressionAlgorithm() != CompressionAlgorithm::NO_COMPRESSION)
diff_content[writer_id]++;
nimages[writer_id]++;
}
img = puller->PollImage(timeout);
}
}
TEST_CASE("ZMQImageCommTest_1Writer","[ZeroMQ]") {
const size_t nframes = 256;
Logger logger(Catch::getResultCapture().getCurrentTestName());
DiffractionExperiment x(DetJF(1));
x.Raw();
x.PedestalG0Frames(0).NumTriggers(1).UseInternalPacketGenerator(false).IncidentEnergy_keV(12.4)
.ImagesPerTrigger(nframes).Compression(CompressionAlgorithm::NO_COMPRESSION);
std::vector<DiffractionSpot> empty_spot_vector;
std::vector<float> empty_rad_int_profile;
REQUIRE(x.GetImageNum() == nframes);
std::mt19937 g1(1387);
std::uniform_int_distribution<uint16_t> dist;
std::vector<uint16_t> image1(x.GetPixelsNum()*nframes);
for (auto &i: image1) i = dist(g1);
// Puller needs to be declared first, but both objects need to exist till communication finished
// TODO: ImageSender should not allow if there are still completions to be done
ZMQStream2Pusher pusher({"ipc://*"});
std::vector<size_t> diff_size(1), diff_content(1), diff_split(1), nimages(1);
auto pusher_addr = pusher.GetAddress();
ZMQImagePuller puller(pusher_addr[0]);
std::thread sender_thread = std::thread([&] {
std::vector<uint8_t> serialization_buffer(16*1024*1024);
CBORStream2Serializer serializer(serialization_buffer.data(), serialization_buffer.size());
StartMessage message {
.images_per_file = 16,
.write_master_file = true
};
EndMessage end_message{};
pusher.StartDataCollection(message);
for (int i = 0; i < nframes; i++) {
DataMessage data_message;
data_message.number = i;
data_message.image = CompressedImage(image1.data() + i * x.GetPixelsNum(),
x.GetPixelsNum() * sizeof(uint16_t),
x.GetXPixelsNum(),
x.GetYPixelsNum(),
x.GetImageMode(),
x.GetCompressionAlgorithm());
serializer.SerializeImage(data_message);
pusher.SendImage(serialization_buffer.data(), serializer.GetBufferSize(), i);
}
pusher.EndDataCollection(end_message);
});
std::thread puller_thread(test_puller, &puller, std::cref(x), std::cref(image1), 1, 0,
std::ref(diff_split), std::ref(diff_size), std::ref(diff_content),
std::ref(nimages));
sender_thread.join();
puller_thread.join();
puller.Disconnect();
REQUIRE(nimages[0] == nframes);
REQUIRE(diff_size[0] == 0);
REQUIRE(diff_content[0] == 0);
}
TEST_CASE("ZMQImageCommTest_2Writers","[ZeroMQ]") {
const size_t nframes = 256;
Logger logger(Catch::getResultCapture().getCurrentTestName());
DiffractionExperiment x(DetJF(1));
x.Raw();
x.PedestalG0Frames(0).NumTriggers(1).UseInternalPacketGenerator(false).IncidentEnergy_keV(12.4)
.ImagesPerTrigger(nframes).Compression(CompressionAlgorithm::NO_COMPRESSION);
REQUIRE(x.GetImageNum() == nframes);
std::mt19937 g1(1387);
std::uniform_int_distribution<uint16_t> dist;
std::vector<uint16_t> image1(x.GetPixelsNum()*nframes);
for (auto &i: image1) i = dist(g1);
std::vector<DiffractionSpot> empty_spot_vector;
std::vector<float> empty_rad_int_profile;
std::vector<std::string> zmq_addr;
int64_t npullers = 2;
for (int i = 0; i < npullers; i++)
zmq_addr.push_back("ipc://*");
ZMQStream2Pusher pusher(zmq_addr);
// Puller needs to be declared first, but both objects need to exist till communication finished
// TODO: ImageSender should not allow if there are still completions to be done
std::vector<std::unique_ptr<ZMQImagePuller> > puller;
auto pusher_addr = pusher.GetAddress();
REQUIRE(pusher_addr.size() == 2);
for (int i = 0; i < npullers; i++) {
puller.push_back(std::make_unique<ZMQImagePuller>(pusher_addr[i]));
}
std::vector<size_t> diff_size(npullers), diff_content(npullers), diff_split(npullers), nimages(npullers);
std::thread sender_thread = std::thread([&] {
std::vector<uint8_t> serialization_buffer(16*1024*1024);
CBORStream2Serializer serializer(serialization_buffer.data(), serialization_buffer.size());
StartMessage message {
.images_per_file = 16,
.write_master_file = true
};
EndMessage end_message{};
pusher.StartDataCollection(message);
for (int i = 0; i < nframes; i++) {
DataMessage data_message;
data_message.number = i;
data_message.image = CompressedImage(image1.data() + i * x.GetPixelsNum(),
x.GetPixelsNum() * sizeof(uint16_t),
x.GetXPixelsNum(),
x.GetYPixelsNum(),
x.GetImageMode(),
x.GetCompressionAlgorithm());
serializer.SerializeImage(data_message);
pusher.SendImage(serialization_buffer.data(), serializer.GetBufferSize(), i);
}
pusher.EndDataCollection(end_message);
});
std::vector<std::thread> puller_threads;
for (int i = 0; i < npullers; i++)
puller_threads.emplace_back(test_puller, puller[i].get(), std::cref(x),
std::cref(image1), npullers, i,
std::ref(diff_split), std::ref(diff_size), std::ref(diff_content), std::ref(nimages));
for (int i = 0; i < npullers; i++)
puller_threads[i].join();
sender_thread.join();
REQUIRE_NOTHROW(puller[0]->Disconnect());
REQUIRE_NOTHROW(puller[1]->Disconnect());
REQUIRE(nimages[0] == nframes / 2);
REQUIRE(nimages[1] == nframes / 2);
REQUIRE(diff_size[0] == 0);
REQUIRE(diff_content[0] == 0);
REQUIRE(diff_size[1] == 0);
REQUIRE(diff_content[1] == 0);
REQUIRE(diff_split[0] == 0);
REQUIRE(diff_split[1] == 0);
}
TEST_CASE("ZMQImageCommTest_4Writers","[ZeroMQ]") {
const size_t nframes = 255;
Logger logger(Catch::getResultCapture().getCurrentTestName());
DiffractionExperiment x(DetJF(1));
x.Raw();
x.PedestalG0Frames(0).NumTriggers(1).UseInternalPacketGenerator(false).IncidentEnergy_keV(12.4)
.ImagesPerTrigger(nframes).Compression(CompressionAlgorithm::NO_COMPRESSION);
REQUIRE(x.GetImageNum() == nframes);
std::mt19937 g1(1387);
std::uniform_int_distribution<uint16_t> dist;
std::vector<uint16_t> image1(x.GetPixelsNum()*nframes);
for (auto &i: image1) i = dist(g1);
std::vector<DiffractionSpot> empty_spot_vector;
std::vector<float> empty_rad_int_profile;
std::vector<std::string> zmq_addr;
int64_t npullers = 4;
for (int i = 0; i < npullers; i++)
zmq_addr.push_back("ipc://*");
ZMQStream2Pusher pusher(zmq_addr);
auto pusher_addr = pusher.GetAddress();
REQUIRE(pusher_addr.size() == npullers);
// Puller needs to be declared first, but both objects need to exist till communication finished
// TODO: ImageSender should not allow if there are still completions to be done
std::vector<std::unique_ptr<ZMQImagePuller> > puller;
for (int i = 0; i < npullers; i++) {
puller.push_back(std::make_unique<ZMQImagePuller>(pusher_addr[i]));
}
std::vector<size_t> diff_size(npullers), diff_content(npullers), diff_split(npullers), nimages(npullers);
std::thread sender_thread = std::thread([&] {
std::vector<uint8_t> serialization_buffer(16*1024*1024);
CBORStream2Serializer serializer(serialization_buffer.data(), serialization_buffer.size());
StartMessage message {
.images_per_file = 16,
.write_master_file = true
};
EndMessage end_message{};
pusher.StartDataCollection(message);
for (int i = 0; i < nframes; i++) {
DataMessage data_message;
data_message.number = i;
data_message.image = CompressedImage(image1.data() + i * x.GetPixelsNum(),
x.GetPixelsNum() * sizeof(uint16_t),
x.GetXPixelsNum(),
x.GetYPixelsNum(),
x.GetImageMode(),
x.GetCompressionAlgorithm());
serializer.SerializeImage(data_message);
pusher.SendImage(serialization_buffer.data(), serializer.GetBufferSize(), i);
}
pusher.EndDataCollection(end_message);
});
std::vector<std::thread> puller_threads;
for (int i = 0; i < npullers; i++)
puller_threads.emplace_back(test_puller, puller[i].get(), std::cref(x),
std::cref(image1), npullers, i,
std::ref(diff_split), std::ref(diff_size), std::ref(diff_content), std::ref(nimages));
for (int i = 0; i < npullers; i++)
puller_threads[i].join();
sender_thread.join();
REQUIRE_NOTHROW(puller[0]->Disconnect());
REQUIRE_NOTHROW(puller[1]->Disconnect());
REQUIRE_NOTHROW(puller[2]->Disconnect());
REQUIRE_NOTHROW(puller[3]->Disconnect());
REQUIRE(nimages[0] == 64);
REQUIRE(nimages[1] == 64);
REQUIRE(nimages[2] == 64);
REQUIRE(nimages[3] == 63);
for (int i = 0; i < npullers; i++) {
REQUIRE(diff_size[i] == 0);
REQUIRE(diff_content[i] == 0);
REQUIRE(diff_split[i] == 0);
}
}
TEST_CASE("ZMQImageCommTest_NoWriter","[ZeroMQ]") {
Logger logger(Catch::getResultCapture().getCurrentTestName());
ZMQStream2Pusher pusher({"ipc://*"});
StartMessage msg{};
REQUIRE_THROWS(pusher.StartDataCollection(msg));
std::vector<uint8_t> test(512*1024, 11);
CompressedImage image(test, 1024, 512);
DataMessage data_message{
.number = 1,
.image = image
};
std::vector<uint8_t> v(16*1024*1024);
CBORStream2Serializer serializer(v.data(), v.size());
serializer.SerializeImage(data_message);
REQUIRE(!pusher.SendImage(v.data(), serializer.GetBufferSize(), 1));
EndMessage end_message{};
REQUIRE(!pusher.EndDataCollection(end_message));
}
// A puller whose queue nobody drains must still shut down. Before this was fixed, the CBOR thread
// parked in PutBlocking on the full outside_fifo, the puller thread backed up behind it in
// cbor_fifo, and Disconnect() hung for ever on the joins - which froze the whole broker, because
// the receiver service destroys the previous run's puller from Start(), holding its state mutex,
// which in a calibration sequence is itself held under the state machine's mutex.
TEST_CASE("ZMQImagePuller_DisconnectWithFullQueue","[ZeroMQ]") {
Logger logger(Catch::getResultCapture().getCurrentTestName());
ZMQSocket sender(ZMQSocketType::Push);
sender.SendWaterMark(4 * ImagePuller::DefaultQueueSize);
sender.Bind("ipc://*");
auto puller = std::make_unique<ZMQImagePuller>(sender.GetEndpointName());
std::vector<uint8_t> pixels(1024, 11);
DataMessage data_message{
.number = 1,
.image = CompressedImage(pixels, 32, 32)
};
std::vector<uint8_t> v(64 * 1024);
CBORStream2Serializer serializer(v.data(), v.size());
serializer.SerializeImage(data_message);
// Twice the queue size, so the CBOR thread is certainly blocked on a full queue by the end
for (int i = 0; i < 2 * ImagePuller::DefaultQueueSize; i++)
sender.Send(v.data(), serializer.GetBufferSize());
while (puller->GetCurrentFifoUtilization() < ImagePuller::DefaultQueueSize)
std::this_thread::sleep_for(std::chrono::milliseconds(10));
auto shutdown = std::async(std::launch::async, [&] { puller.reset(); });
REQUIRE(shutdown.wait_for(std::chrono::seconds(30)) == std::future_status::ready);
}
// A run that ends while the detector is still streaming leaves outside_fifo full and the CBOR
// thread blocked in PutBlocking on it. Suspend() cannot release it - the flag is only tested
// before the put - so the next run's ResumeAndClear has to. Nothing else ever would: a Get on an
// empty queue notifies no producer, so a queue emptied under a waiting producer leaves it asleep
// for good, and every collection after that one silently receives nothing at all - no images and
// no start message - while the puller looks perfectly healthy.
TEST_CASE("ZMQImagePuller_RecoversFromFullQueue","[ZeroMQ]") {
Logger logger(Catch::getResultCapture().getCurrentTestName());
ZMQSocket sender(ZMQSocketType::Push);
sender.SendWaterMark(8 * ImagePuller::DefaultQueueSize);
sender.Bind("ipc://*");
ZMQImagePuller puller(sender.GetEndpointName());
std::vector<uint8_t> pixels(1024, 11);
auto send_image = [&](int64_t number) {
DataMessage data_message{
.number = number,
.image = CompressedImage(pixels, 32, 32)
};
std::vector<uint8_t> v(64 * 1024);
CBORStream2Serializer serializer(v.data(), v.size());
serializer.SerializeImage(data_message);
sender.Send(v.data(), serializer.GetBufferSize());
};
// Images still arriving after the last run stopped consuming them
for (int i = 0; i < 2 * ImagePuller::DefaultQueueSize; i++)
send_image(1);
while (puller.GetCurrentFifoUtilization() < ImagePuller::DefaultQueueSize)
std::this_thread::sleep_for(std::chrono::milliseconds(10));
puller.Suspend(); // what JFJochServices::Stop does when the run ends
puller.ResumeAndClear(); // ... and what JFJochServices::Start does for the run after it
// The run after it must still be able to receive
send_image(4242);
bool received = false;
auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(30);
while (!received && (std::chrono::steady_clock::now() < deadline)) {
auto msg = puller.PollImage(std::chrono::milliseconds(10));
if (msg && msg->cbor && msg->cbor->data_message && (msg->cbor->data_message->number == 4242))
received = true;
}
REQUIRE(received);
puller.Disconnect();
}