// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute // SPDX-License-Identifier: GPL-3.0-only #include #include #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 &image1, int64_t nwriter, int64_t writer_id, std::vector &diff_split, std::vector &diff_size, std::vector &diff_content, std::vector &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 empty_spot_vector; std::vector empty_rad_int_profile; REQUIRE(x.GetImageNum() == nframes); std::mt19937 g1(1387); std::uniform_int_distribution dist; std::vector 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 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 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 dist; std::vector image1(x.GetPixelsNum()*nframes); for (auto &i: image1) i = dist(g1); std::vector empty_spot_vector; std::vector empty_rad_int_profile; std::vector 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 > puller; auto pusher_addr = pusher.GetAddress(); REQUIRE(pusher_addr.size() == 2); for (int i = 0; i < npullers; i++) { puller.push_back(std::make_unique(pusher_addr[i])); } std::vector diff_size(npullers), diff_content(npullers), diff_split(npullers), nimages(npullers); std::thread sender_thread = std::thread([&] { std::vector 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 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 dist; std::vector image1(x.GetPixelsNum()*nframes); for (auto &i: image1) i = dist(g1); std::vector empty_spot_vector; std::vector empty_rad_int_profile; std::vector 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 > puller; for (int i = 0; i < npullers; i++) { puller.push_back(std::make_unique(pusher_addr[i])); } std::vector diff_size(npullers), diff_content(npullers), diff_split(npullers), nimages(npullers); std::thread sender_thread = std::thread([&] { std::vector 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 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 test(512*1024, 11); CompressedImage image(test, 1024, 512); DataMessage data_message{ .number = 1, .image = image }; std::vector 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(sender.GetEndpointName()); std::vector pixels(1024, 11); DataMessage data_message{ .number = 1, .image = CompressedImage(pixels, 32, 32) }; std::vector 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 pixels(1024, 11); auto send_image = [&](int64_t number) { DataMessage data_message{ .number = number, .image = CompressedImage(pixels, 32, 32) }; std::vector 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(); }