Files
Jungfraujoch/tests/ZMQImagePusherTest.cpp
T
leonarski_f 9aae0c2ba7
Build Packages / Create release (push) Successful in 21s
Build Packages / build:rugnux-tgz (x86_64) (push) Successful in 9m40s
Build Packages / build:rugnux:aarch64 (cross) (push) Successful in 9m49s
Build Packages / build:viewer-tgz:cpu (push) Successful in 11m37s
Build Packages / build:viewer-tgz:cuda (push) Successful in 12m40s
Build Packages / build:windows:nocuda (push) Successful in 17m44s
Build Packages / build:windows:cuda (push) Successful in 20m13s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 14m41s
Build Packages / HDF5 consumer tests (DIALS, XDS) (push) Successful in 25m59s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 15m5s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 14m35s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 15m53s
Build Packages / build:rugnux:windows (push) Successful in 11m29s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 18m51s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 18m43s
Build Packages / Generate python client (push) Successful in 51s
Build Packages / build:rpm (rocky8) (push) Successful in 18m51s
Build Packages / Build documentation (push) Successful in 1m21s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 18m38s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 18m24s
Build Packages / build:rpm (rocky9) (push) Successful in 19m19s
Build Packages / Unit tests (push) Successful in 1h37m15s
v1.0.0-rc.169 (#79)
* Building Jungfraujoch no longer needs zlib or Eigen installed on the machine, and the dependencies the build fetches are pinned and updated to current releases.
* rugnux: improvements in indexing, lattice selection and geometry post-refinement, which index crystals that previously returned no lattice and keep the better of the two geometries a run measures.
* rugnux: improvements in beam-centre measurement, beam-stop detection and space-group determination.
* rugnux: the unit cell reported with a determined space group now obeys that group - a cell whose symmetry was confirmed from the intensities is re-refined under it, and a cell the group cannot describe is reported with a warning rather than as it stands.
* rugnux drops the stretches of a rotation sweep whose removal measurably improves the merged intensities and reports what became of every frame, and decides the resolution cut on the crystal's own diffraction rather than on its ice rings.
* The rugnux results report is machine-readable - every line that is not `KEY= value` data starts with `#` - and states the build it was written by, its authorship and its terms of use (`REPORT_VERSION= 8`).
* `jfjoch_viewer`: improvements in the file manager (CBF frames beside HDF5 datasets, a remembered root), the dataset plots, the inspector and the image statistics, plus a settable font size, a view of the rugnux results report, usable performance over a remote display (`ssh -X`) and a reset of all settings to defaults; the reciprocal-space window is removed.
* Broker fixes around DECTRIS collections and dark-mask calibration: re-initialising after a run that never started no longer freezes the broker, a cancelled calibration is abandoned instead of reported as done, and a collection whose start message never arrives ends by itself.

Reviewed-on: #79
Co-authored-by: Filip Leonarski <filip.leonarski@psi.ch>
2026-09-15 17:09:31 +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();
}