Files
Jungfraujoch/tests/StreamWriterTest.cpp
T
leonarski_fandClaude Opus 5.5 f970972163 Writer and broker: the master file appears only after every data file it links to
The final master marks a finished collection for anyone watching the directory (rsync, autoprocessing,
a reader following the collection), so it must not appear before a data file it links to.

- FileWriter: at END every data file is closed (renamed into place) before the master is finalized
  and renamed. A data file that fails to close does not cost the master; the failure is reported
  after it. For NXmxIntegrated, file 0 is the master's own image dataset, so there a failure still
  stops the master.
- With several writers, the master was written by the writer of the first socket as soon as it got
  END, while the others could still be closing their data files. END now goes to every other writer
  first:
  - TCP: the connections are taken in reverse order, and each END waits for its acknowledgement,
    which a writer sends only after closing its files;
  - ZeroMQ, with a writer notification socket: END to every other socket, then their notifications
    (the same ones Finalize collected before, now collected here), then END to the first socket.
    Without a notification socket there is nothing to wait on and the order is all it changes.

Tests: TCPImageCommTest_MasterWriterGetsEndLast (three mock writers, the master's END after the
other two acknowledged), StreamWriterTest_ZMQ_MasterAfterAllDataFiles (two real writers, the second
started only after the collection was sent; the master's ctime is after every data file's). Both
fail without this change.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi
2026-10-08 14:55:56 +02:00

249 lines
10 KiB
C++

// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include <catch2/catch_all.hpp>
#include <filesystem>
#include <future>
#include <thread>
#include <sys/stat.h>
#include "../writer/StreamWriter.h"
#include "../image_pusher/ZMQStream2Pusher.h"
#include "../receiver/JFJochReceiverService.h"
#include "../image_pusher/ZMQWriterNotificationPuller.h"
TEST_CASE("StreamWriterTest_ZMQ", "[StreamWriter]") {
RegisterHDF5Filter();
Logger logger("StreamWriterTest_ZMQ");
DiffractionExperiment x(DetJF(2));
x.FilePrefix("subdir/StreamWriterTest").NumTriggers(1).ImagesPerTrigger(5)
.UseInternalPacketGenerator(true).Raw().PedestalG0Frames(0).OverwriteExistingFiles(true);
PixelMask pixel_mask(x);
JFModuleGainCalibration gain;
AcquisitionDeviceGroup aq_devices;
for (int i = 0; i < x.GetDataStreamsNum(); i++)
aq_devices.AddHLSDevice(64);
ZMQStream2Pusher pusher({"ipc://*"});
JFJochReceiverService fpga_receiver_service(aq_devices, logger, pusher);
std::unique_ptr<StreamWriter> writer;
REQUIRE(x.GetImageNum() == 5);
auto pusher_addr = pusher.GetAddress();
REQUIRE(pusher_addr.size() == 1);
ZMQImagePuller puller(pusher_addr[0]);
REQUIRE_NOTHROW(writer = std::make_unique<StreamWriter>(logger, puller));
CHECK(writer->GetStatistics().state == StreamWriterState::Idle);
REQUIRE_NOTHROW(fpga_receiver_service.Start(x, pixel_mask, nullptr));
REQUIRE_NOTHROW(writer->Run());
REQUIRE_NOTHROW(fpga_receiver_service.Stop());
REQUIRE(fpga_receiver_service.GetStatus()->images_collected == 5);
REQUIRE(fpga_receiver_service.GetStatus()->images_sent == 5);
CHECK(writer->GetStatistics().state == StreamWriterState::Finalized);
CHECK(writer->GetStatistics().processed_images == 5);
CHECK(writer->GetStatistics().file_prefix == x.GetFilePrefix());
// HDF5 file can be opened
std::unique_ptr<HDF5ReadOnlyFile> file;
REQUIRE_NOTHROW(file = std::make_unique<HDF5ReadOnlyFile>("subdir/StreamWriterTest_data_000001.h5"));
std::unique_ptr<HDF5DataSet> dataset;
REQUIRE_NOTHROW(dataset = std::make_unique<HDF5DataSet>(*file, "/entry/data/data"));
std::unique_ptr<HDF5DataSpace> dataspace;
REQUIRE_NOTHROW(dataspace = std::make_unique<HDF5DataSpace>(*dataset));
REQUIRE(dataspace->GetNumOfDimensions() == 3);
REQUIRE(dataspace->GetDimensions()[0] == 5);
REQUIRE(dataspace->GetDimensions()[1] == RAW_MODULE_COLS);
REQUIRE(dataspace->GetDimensions()[2] == 2*RAW_MODULE_LINES);
REQUIRE(std::filesystem::remove("subdir/StreamWriterTest_master.h5"));
REQUIRE(std::filesystem::remove("subdir/StreamWriterTest_data_000001.h5"));
REQUIRE(std::filesystem::remove("subdir"));
}
TEST_CASE("StreamWriterTest_ZMQ_Update", "[StreamWriter]") {
RegisterHDF5Filter();
Logger logger("StreamWriterTest_ZMQ_Update");
DatasetSettings d;
d.FilePrefix("subdir/StreamWriterTest2").NumTriggers(1).ImagesPerTrigger(5).RunName("run1").RunNumber(256);
DiffractionExperiment x(DetJF(2));
x.UseInternalPacketGenerator(true).Raw().PedestalG0Frames(0)
.ImportDatasetSettings(d).OverwriteExistingFiles(true);
PixelMask pixel_mask(x);
JFModuleGainCalibration gain;
AcquisitionDeviceGroup aq_devices;
for (int i = 0; i < x.GetDataStreamsNum(); i++)
aq_devices.AddHLSDevice(64);
ZMQStream2Pusher pusher({"ipc://*"});
pusher.WriterNotificationSocket("ipc://*");
JFJochReceiverService fpga_receiver_service(aq_devices, logger, pusher);
std::unique_ptr<StreamWriter> writer;
REQUIRE(x.GetImageNum() == 5);
auto pusher_addr = pusher.GetAddress();
REQUIRE(pusher_addr.size() == 1);
ZMQImagePuller puller(pusher_addr[0]);
REQUIRE_NOTHROW(writer = std::make_unique<StreamWriter>(logger, puller));
CHECK(writer->GetStatistics().state == StreamWriterState::Idle);
REQUIRE_NOTHROW(fpga_receiver_service.Start(x, pixel_mask, nullptr));
REQUIRE_NOTHROW(writer->Run());
REQUIRE_NOTHROW(fpga_receiver_service.Stop());
REQUIRE(fpga_receiver_service.GetStatus()->images_collected == 5);
REQUIRE(fpga_receiver_service.GetStatus()->images_sent == 5);
CHECK(writer->GetStatistics().state == StreamWriterState::Finalized);
CHECK(writer->GetStatistics().processed_images == 5);
CHECK(writer->GetStatistics().file_prefix == x.GetFilePrefix());
// HDF5 file can be opened
std::unique_ptr<HDF5ReadOnlyFile> file;
REQUIRE_NOTHROW(file = std::make_unique<HDF5ReadOnlyFile>("subdir/StreamWriterTest2_data_000001.h5"));
std::unique_ptr<HDF5DataSet> dataset;
REQUIRE_NOTHROW(dataset = std::make_unique<HDF5DataSet>(*file, "/entry/data/data"));
std::unique_ptr<HDF5DataSpace> dataspace;
REQUIRE_NOTHROW(dataspace = std::make_unique<HDF5DataSpace>(*dataset));
REQUIRE(dataspace->GetNumOfDimensions() == 3);
REQUIRE(dataspace->GetDimensions()[0] == 5);
REQUIRE(dataspace->GetDimensions()[1] == RAW_MODULE_COLS);
REQUIRE(dataspace->GetDimensions()[2] == 2*RAW_MODULE_LINES);
REQUIRE(std::filesystem::remove("subdir/StreamWriterTest2_master.h5"));
REQUIRE(std::filesystem::remove("subdir/StreamWriterTest2_data_000001.h5"));
REQUIRE(std::filesystem::remove("subdir"));
}
TEST_CASE("StreamWriterTest_ZMQ_Update_NoNotification", "[StreamWriter]") {
// This tests simulates what happens if writer notification about writing end is missing
// Expected end result: receiver ends with an exception
RegisterHDF5Filter();
Logger logger("StreamWriterTest_ZMQ_Update_NoNotification");
DatasetSettings d;
d.FilePrefix("subdir/StreamWriterTest3").NumTriggers(1).ImagesPerTrigger(5).RunName("run1").RunNumber(256);
DiffractionExperiment x(DetJF(2));
x.UseInternalPacketGenerator(true).Raw().PedestalG0Frames(0)
.ImportDatasetSettings(d).OverwriteExistingFiles(true);
PixelMask pixel_mask(x);
JFModuleGainCalibration gain;
AcquisitionDeviceGroup aq_devices;
for (int i = 0; i < x.GetDataStreamsNum(); i++)
aq_devices.AddHLSDevice(64);
ZMQStream2Pusher pusher({"ipc://*"});
pusher.WriterNotificationSocket("ipc://*");
JFJochReceiverService fpga_receiver_service(aq_devices, logger, pusher);
std::unique_ptr<StreamWriter> writer;
REQUIRE(x.GetImageNum() == 5);
auto pusher_addr = pusher.GetAddress();
REQUIRE(pusher_addr.size() == 1);
ZMQImagePuller puller(pusher_addr[0]);
REQUIRE_NOTHROW(writer = std::make_unique<StreamWriter>(logger, puller));
writer->DebugSkipWriteNotification(true);
CHECK(writer->GetStatistics().state == StreamWriterState::Idle);
REQUIRE_NOTHROW(fpga_receiver_service.Start(x, pixel_mask, nullptr));
REQUIRE_NOTHROW(writer->Run());
JFJochReceiverOutput r;
REQUIRE_NOTHROW(r = fpga_receiver_service.Stop());
REQUIRE(!r.writer_err.empty());
REQUIRE(fpga_receiver_service.GetStatus()->images_collected == 5);
REQUIRE(fpga_receiver_service.GetStatus()->images_sent == 5);
CHECK(writer->GetStatistics().state == StreamWriterState::Finalized);
CHECK(writer->GetStatistics().processed_images == 5);
CHECK(writer->GetStatistics().file_prefix == x.GetFilePrefix());
// HDF5 file can be opened
std::unique_ptr<HDF5ReadOnlyFile> file;
REQUIRE_NOTHROW(file = std::make_unique<HDF5ReadOnlyFile>("subdir/StreamWriterTest3_data_000001.h5"));
std::unique_ptr<HDF5DataSet> dataset;
REQUIRE_NOTHROW(dataset = std::make_unique<HDF5DataSet>(*file, "/entry/data/data"));
std::unique_ptr<HDF5DataSpace> dataspace;
REQUIRE_NOTHROW(dataspace = std::make_unique<HDF5DataSpace>(*dataset));
REQUIRE(dataspace->GetNumOfDimensions() == 3);
REQUIRE(dataspace->GetDimensions()[0] == 5);
REQUIRE(dataspace->GetDimensions()[1] == RAW_MODULE_COLS);
REQUIRE(dataspace->GetDimensions()[2] == 2*RAW_MODULE_LINES);
REQUIRE(std::filesystem::remove("subdir/StreamWriterTest3_master.h5"));
REQUIRE(std::filesystem::remove("subdir/StreamWriterTest3_data_000001.h5"));
REQUIRE(std::filesystem::remove("subdir"));
}
// Two writers, the second one writing every other data file: the master, written by the first, must
// appear only after the second has closed its data files - here the second one starts only after the
// whole collection has been sent, so it is still writing when END goes out.
TEST_CASE("StreamWriterTest_ZMQ_MasterAfterAllDataFiles", "[StreamWriter]") {
RegisterHDF5Filter();
Logger logger("StreamWriterTest_ZMQ_MasterAfterAllDataFiles");
DatasetSettings d;
d.FilePrefix("subdir4/StreamWriterTest4").NumTriggers(1).ImagesPerTrigger(4).ImagesPerFile(2)
.RunName("run1").RunNumber(257);
DiffractionExperiment x(DetJF(2));
x.UseInternalPacketGenerator(true).Raw().PedestalG0Frames(0)
.ImportDatasetSettings(d).OverwriteExistingFiles(true);
PixelMask pixel_mask(x);
AcquisitionDeviceGroup aq_devices;
for (int i = 0; i < x.GetDataStreamsNum(); i++)
aq_devices.AddHLSDevice(64);
ZMQStream2Pusher pusher({"ipc://*", "ipc://*"});
pusher.WriterNotificationSocket("ipc://*");
JFJochReceiverService fpga_receiver_service(aq_devices, logger, pusher);
auto pusher_addr = pusher.GetAddress();
REQUIRE(pusher_addr.size() == 2);
ZMQImagePuller puller0(pusher_addr[0]), puller1(pusher_addr[1]);
StreamWriter writer0(logger, puller0), writer1(logger, puller1);
REQUIRE_NOTHROW(fpga_receiver_service.Start(x, pixel_mask, nullptr));
auto w0 = std::async(std::launch::async, [&] { return writer0.Run(); });
for (int i = 0; i < 1200; i++) {
const auto status = fpga_receiver_service.GetStatus();
if (status && status->images_sent == 4)
break;
std::this_thread::sleep_for(std::chrono::milliseconds(100));
}
auto stop = std::async(std::launch::async, [&] { return fpga_receiver_service.Stop(); });
std::this_thread::sleep_for(std::chrono::seconds(2));
auto w1 = std::async(std::launch::async, [&] { return writer1.Run(); });
JFJochReceiverOutput r;
REQUIRE_NOTHROW(r = stop.get());
w0.get();
w1.get();
CHECK(r.writer_err.empty());
auto changed = [](const std::string &path) {
struct stat st{};
REQUIRE(stat(path.c_str(), &st) == 0);
return std::make_pair(st.st_ctim.tv_sec, st.st_ctim.tv_nsec);
};
const auto master = changed("subdir4/StreamWriterTest4_master.h5");
for (int f = 1; f <= 2; f++)
CHECK(changed("subdir4/StreamWriterTest4_data_00000" + std::to_string(f) + ".h5") <= master);
std::filesystem::remove_all("subdir4");
}