Build Packages / build:rpm (rocky9_sls9) (push) Successful in 18m57s
Build Packages / Unit tests (push) Skipped
Build Packages / build:windows:nocuda (push) Successful in 16m55s
Build Packages / build:windows:cuda (push) Successful in 18m48s
Build Packages / build:viewer-tgz:cpu (push) Successful in 13m10s
Build Packages / build:viewer-tgz:cuda (push) Successful in 14m45s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 22m23s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 20m12s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 23m7s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 20m43s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 23m9s
Build Packages / XDS test (durin plugin) (push) Successful in 12m26s
Build Packages / build:rpm (rocky9) (push) Successful in 24m58s
Build Packages / Generate python client (push) Successful in 50s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 23m20s
Build Packages / Create release (push) Skipped
Build Packages / XDS test (JFJoch plugin) (push) Successful in 12m37s
Build Packages / build:rpm (rocky8) (push) Successful in 27m58s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 25m38s
Build Packages / Build documentation (push) Successful in 59s
Build Packages / DIALS test (push) Successful in 23m16s
Build Packages / XDS test (neggia plugin) (push) Successful in 6m38s
**Files written by Jungfraujoch now import correctly in DIALS, XDS and pyFAI.** A tilted detector, a grid scan, a still recorded at a goniometer position, and saturated or unreadable pixels were each described in a way that a third-party program acted on wrongly. If you process Jungfraujoch data outside Jungfraujoch, prefer this release to any earlier one. * HDF5: the detector tilt (`rot1`/`rot2`/`rot3`) is exported correctly in the NXmx transformation chain; untilted geometries are unaffected. * HDF5: a still recorded at a goniometer position is no longer read back as a single image, and a grid scan records a stationary spindle so a program that requires a rotation axis can open it. * HDF5: the sample transformation chain is written in mounting order, with a Smargon head position told apart from the spindle, one entry per image, `module_offset` as a float unit vector, and `offset_units` on every offset. * HDF5: saturated, underloaded and unreadable pixels are described so a downstream program masks them - `saturation_value`, `underload_value`, `error_value` and `bit_depth_readout` are written correctly, and a data file missing next to a VDS master reads as the error marker rather than as zero counts. * HDF5: the rotation axis is read back under whatever name it carries, and `mirror_y` records whether the assembled image is mirrored in Y relative to the detector's raw readout. * A grid scan and a goniometer axis can both be set; they are no longer alternatives. * `images_per_file` is chosen from the acquisition when it is not given: a rotation sweep of at most 20000 images goes into a single data file, a grid scan splits on whole fast-axis rows, and stills and serial keep 1000. * The writer refuses a stream whose start message declares a different pixel format than its images carry, and a DECTRIS detector sending signed images is no longer declared unsigned. * The image stream can carry the sample transformation chain (`transformations`, in the END message); a producer that does not send it gets the same chain built by the writer. * rugnux: fixing the space group with `-S` no longer prevents the lattice from being found - a lattice indexed in a different setting is reindexed into that group's own setting, and a run whose crystal does not have that group's lattice stops and names the cell it indexed as, rather than reporting statistics that cannot describe it. * rugnux: the per-image resolution estimate now predicts the resolution the merged data reach rather than the highest-resolution spot found, and is reported as `SPOT_RESOLUTION_ESTIMATE`. * rugnux: two runs of the same command on the same images produce the same merged intensities; the azimuthal profile written alongside them is not yet reproducible in the same way. * rugnux: the offline lattice refinement is bounded by iterations rather than by a wall clock, so a loaded machine can no longer refine to a different lattice; a live acquisition keeps its real-time bound. * rugnux: the detector-frame modulation correction is fitted on a grid spanning the detector, so whether it is applied no longer depends on how far integration reached. * rugnux: the geometry pre-pass no longer writes `<prefix>_01.mtz`, `_01.cif`, `_01.hkl` and `_01_image.dat`; the refined second pass writes those files under `<prefix>`, and that is the result to use. * rugnux: `_process.h5` describes the pixel format of the images it links to, and is written on a thread of its own. * rugnux: the detector geometry is also logged in XDS's convention (`ORGX`/`ORGY`, detector axis vectors, rotation axis), so it can be compared with an XDS refinement. * rugnux: an image integrated in pyFAI through the `.poni` file written by `--mode calibration` comes out with the correct azimuth, and the file declares pyFAI's `orientation`, which needs pyFAI 2024.01 or newer. Radial integration is unchanged. * rugnux: a rotation run is substantially faster throughout - beam-stop detection, first-pass indexing, geometry refinement, integration, scaling and merging - and observations outside the scaling resolution range are dropped as they are ingested. The refined geometry, the space group chosen and the merged statistics are unchanged. * Faster spot finding and indexing, on the broker as well as in rugnux; the spots found and the lattices indexed are unchanged. * A run reserves substantially less GPU memory: nothing is allocated for buffers that are never read, and a worker builds only the engines it uses. * rugnux: with `-N` left at its default the per-image loop of `--mode mx` uses at most 16 workers per GPU, rather than one per hardware thread; an explicit `-N` is obeyed as given. * CUDA 12 builds now contain device code for Volta, so the RHEL 8 packages and the portable Linux `.tgz` run on a V100; the CUDA 13 artefacts (RHEL 9, Ubuntu, Windows) remain Turing and newer. * The build resolves a single Eigen for the whole project, and refuses to configure if Ceres picks up a different one; a build that mixed two Eigen versions was undefined behaviour and crashed at -O2. * Documentation: a security page, and the supported GPU generations and minimum NVIDIA driver version of every released artefact. **Breaking change to OpenAPI** - regenerate the client (`jfjoch-client` 1.0.0-rc.162, `frontend/src/client`): * `dataset_settings.images_per_file` is no longer `default: 1000` and no longer accepts `0`; it is optional, and its minimum is 1. A client sending `0` (previously "one file for the whole run") is now rejected - omit the field instead, which for a rotation sweep gives the same single file. * `file_writer_format` now defaults to `NXmxVDS`, matching the server's own default and the layout recommended for DIALS, XDS and CrystFEL. A generated client that fills in schema defaults and does not set the format explicitly will write VDS masters where it previously wrote legacy ones; set `NXmxLegacy` explicitly to keep them. --------- Co-authored-by: jungfrau <jungfrau@mx-aare-test.psi.ch> Reviewed-on: #72 Co-authored-by: Filip Leonarski <filip.leonarski@psi.ch>
268 lines
9.4 KiB
C++
268 lines
9.4 KiB
C++
// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
|
|
// SPDX-License-Identifier: GPL-3.0-only
|
|
|
|
#include "FileWriter.h"
|
|
#include <filesystem>
|
|
#include <nlohmann/json.hpp>
|
|
#include "MakeDirectory.h"
|
|
#include "../common/CheckPath.h"
|
|
#include "../common/Logger.h"
|
|
#include "../common/JFJochException.h"
|
|
#include "JFJochDecompress.h"
|
|
|
|
FileWriter::FileWriter(const StartMessage &request, bool check_overwrite_at_start, bool trusted_path)
|
|
: start_message(request) {
|
|
if (start_message.file_format)
|
|
format = start_message.file_format.value();
|
|
|
|
if (start_message.images_per_file <= 0)
|
|
start_message.images_per_file = default_images_per_file;
|
|
|
|
// trusted_path skips the multi-user guard so an offline caller can write to an absolute path;
|
|
// MakeDirectory still runs (it handles absolute prefixes). See the constructor comment in the header.
|
|
if (!trusted_path)
|
|
CheckPath(start_message.file_prefix);
|
|
MakeDirectory(start_message.file_prefix);
|
|
if (check_overwrite_at_start)
|
|
CheckOutputFilesAvailable();
|
|
if (start_message.write_master_file && start_message.write_master_file.value()) {
|
|
switch (format) {
|
|
case FileWriterFormat::NXmxLegacy:
|
|
case FileWriterFormat::NXmxVDS:
|
|
case FileWriterFormat::NXmxIntegrated:
|
|
CreateHDF5MasterFile(request);
|
|
break;
|
|
default:
|
|
// Do nothing
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
|
|
void FileWriter::Write(const DataMessage &msg) {
|
|
switch (format) {
|
|
case FileWriterFormat::DataOnly:
|
|
case FileWriterFormat::NXmxLegacy:
|
|
case FileWriterFormat::NXmxVDS:
|
|
case FileWriterFormat::NXmxIntegrated:
|
|
WriteHDF5(msg);
|
|
break;
|
|
case FileWriterFormat::NoFile:
|
|
// Do nothing
|
|
break;
|
|
}
|
|
}
|
|
|
|
void FileWriter::WriteHDF5(const DataMessage& msg) {
|
|
std::lock_guard<std::mutex> lock(hdf5_mutex);
|
|
if (msg.image.GetCompressedSize() == 0)
|
|
return;
|
|
|
|
if (msg.number < 0)
|
|
throw JFJochException(JFJochExceptionCategory::ArrayOutOfBounds, "No support for negative images");
|
|
|
|
if (format == FileWriterFormat::NXmxIntegrated && master_file) {
|
|
if (files.empty() )
|
|
files.resize(1);
|
|
if (!files[0]) {
|
|
files[0] = std::make_unique<HDF5DataFile>(start_message, 0, HDF5Metadata::MasterFileName(start_message));
|
|
files[0]->CreateFile(msg, master_file->GetFile());
|
|
}
|
|
files[0]->Write(msg, msg.number);
|
|
} else {
|
|
const uint64_t file_number = (start_message.images_per_file == 0) ? 0 : msg.number / start_message.images_per_file;
|
|
const uint64_t image_number = (start_message.images_per_file == 0) ? msg.number : msg.number % start_message.images_per_file;
|
|
|
|
if (closed_files.contains(file_number))
|
|
return;
|
|
|
|
if (files.size() <= file_number)
|
|
files.resize(file_number + 1);
|
|
|
|
if (!files[file_number])
|
|
files[file_number] = std::make_unique<HDF5DataFile>(start_message, file_number,
|
|
HDF5Metadata::DataFileName(start_message, file_number));
|
|
files[file_number]->Write(msg, image_number);
|
|
|
|
if (files[file_number]->GetNumImages() == start_message.images_per_file) {
|
|
CloseFile(file_number);
|
|
} else {
|
|
CloseOldFiles(static_cast<uint64_t>(msg.number));
|
|
}
|
|
}
|
|
}
|
|
|
|
void FileWriter::CloseFile(uint64_t file_number) {
|
|
if (file_number >= files.size())
|
|
return;
|
|
if (!files[file_number])
|
|
return;
|
|
if (closed_files.contains(file_number))
|
|
return;
|
|
|
|
auto file_stats = files[file_number]->Close();
|
|
files[file_number].reset();
|
|
closed_files.insert(file_number);
|
|
|
|
AddStats(file_stats);
|
|
}
|
|
|
|
void FileWriter::CloseOldFiles(uint64_t current_image_number) {
|
|
if (start_message.images_per_file == 0)
|
|
return;
|
|
|
|
for (uint64_t f = 0; f < files.size(); ++f) {
|
|
if (!files[f] || closed_files.contains(f))
|
|
continue;
|
|
|
|
const uint64_t file_end_image = (f + 1) * start_message.images_per_file - 1;
|
|
if (current_image_number > file_end_image + close_file_lag_images) {
|
|
CloseFile(f);
|
|
}
|
|
}
|
|
}
|
|
|
|
std::vector<HDF5DataFileStatistics> FileWriter::Finalize() {
|
|
std::lock_guard<std::mutex> lock(hdf5_mutex);
|
|
|
|
std::exception_ptr first_exception;
|
|
|
|
for (uint64_t f = 0; f < files.size(); ++f) {
|
|
if (files[f] && !closed_files.contains(f)) {
|
|
try {
|
|
CloseFile(f);
|
|
} catch (...) {
|
|
if (!first_exception)
|
|
first_exception = std::current_exception();
|
|
}
|
|
}
|
|
}
|
|
|
|
if (master_file) {
|
|
try {
|
|
master_file.reset();
|
|
} catch (...) {
|
|
if (!first_exception)
|
|
first_exception = std::current_exception();
|
|
}
|
|
}
|
|
|
|
if (first_exception)
|
|
std::rethrow_exception(first_exception);
|
|
|
|
return stats;
|
|
}
|
|
|
|
void FileWriter::AddStats(const std::optional<HDF5DataFileStatistics>& s) {
|
|
if (!s)
|
|
return;
|
|
|
|
stats.push_back(*s);
|
|
if (finalized_file_socket) {
|
|
nlohmann::json j;
|
|
j["filename"] = s->filename;
|
|
j["nimages"] = s->total_images;
|
|
j["file_number"] = s->file_number;
|
|
|
|
j["detector_distance_m"] = start_message.detector_distance;
|
|
j["beam_x_pxl"] = start_message.beam_center_x;
|
|
j["beam_y_pxl"] = start_message.beam_center_y;
|
|
j["pixel_size_m"] = start_message.pixel_size_x;
|
|
j["detector_width_pxl"] = start_message.image_size_x;
|
|
j["detector_height_pxl"] = start_message.image_size_y;
|
|
j["incident_energy_eV"] = start_message.incident_energy;
|
|
j["saturation"] = start_message.saturation_value;
|
|
j["sample_name"] = start_message.sample_name;
|
|
j["run_number"] = start_message.run_number;
|
|
j["run_name"] = start_message.run_name;
|
|
|
|
if (!start_message.experiment_group.empty())
|
|
j["experiment_group"] = start_message.experiment_group;
|
|
|
|
if (start_message.unit_cell) {
|
|
j["unit_cell"]["a"] = start_message.unit_cell->a;
|
|
j["unit_cell"]["b"] = start_message.unit_cell->b;
|
|
j["unit_cell"]["c"] = start_message.unit_cell->c;
|
|
j["unit_cell"]["alpha"] = start_message.unit_cell->alpha;
|
|
j["unit_cell"]["beta"] = start_message.unit_cell->beta;
|
|
j["unit_cell"]["gamma"] = start_message.unit_cell->gamma;
|
|
}
|
|
if (start_message.space_group_number)
|
|
j["space_group_number"] = start_message.space_group_number.value();
|
|
// The lowest VALID value, not the error marker. These were the same field until the marker
|
|
// became the value the pixels actually carry: for an unsigned image that is UINTx_MAX, and a
|
|
// consumer feeding this key into an XDS UNDERLOAD or a DIALS trusted range would then reject
|
|
// every pixel below 65535.
|
|
if (start_message.underload_value)
|
|
j["underload"] = start_message.underload_value.value();
|
|
|
|
j["user_data"] = start_message.user_data;
|
|
finalized_file_socket->Send(j.dump());
|
|
}
|
|
}
|
|
|
|
void FileWriter::SetupFinalizedFileSocket(const std::string &addr) {
|
|
finalized_file_socket = std::make_unique<ZMQSocket>(ZMQSocketType::Pub);
|
|
finalized_file_socket->Bind(addr);
|
|
}
|
|
|
|
std::optional<std::string> FileWriter::GetZMQAddr() {
|
|
if (finalized_file_socket) {
|
|
return finalized_file_socket->GetEndpointName();
|
|
} else
|
|
return {};
|
|
}
|
|
|
|
void FileWriter::CreateHDF5MasterFile(const StartMessage &msg) {
|
|
std::lock_guard<std::mutex> lock(hdf5_mutex);
|
|
master_file = std::make_unique<NXmx>(msg);
|
|
}
|
|
|
|
void FileWriter::CheckOutputFilesAvailable() const {
|
|
if (start_message.overwrite.value_or(false))
|
|
return;
|
|
|
|
// Only the master file is checked, and only by the single writer that owns it
|
|
// (write_master_file - index 0 in a multi-writer TCP/ZMQ setup). Data files are
|
|
// staggered across writers by file number, so enumerating them here would make
|
|
// every writer stat files it never writes and race sibling writers that are
|
|
// already creating them; those conflicts are caught per-writer at finalize.
|
|
const bool nxmx = format == FileWriterFormat::NXmxLegacy
|
|
|| format == FileWriterFormat::NXmxVDS
|
|
|| format == FileWriterFormat::NXmxIntegrated;
|
|
|
|
if (nxmx && start_message.write_master_file.value_or(false)) {
|
|
const std::string name = HDF5Metadata::MasterFileName(start_message);
|
|
if (std::filesystem::exists(name))
|
|
throw JFJochException(JFJochExceptionCategory::FileWriteError,
|
|
"Output file already exists and overwrite is off: " + name);
|
|
}
|
|
}
|
|
|
|
void FileWriter::WriteHDF5(const CompressedImage &msg) {
|
|
if (master_file) {
|
|
std::lock_guard<std::mutex> lock(hdf5_mutex);
|
|
try {
|
|
master_file->WriteCalibration(msg);
|
|
} catch (const JFJochException &e) {
|
|
spdlog::error("Calibration {} not written {}", msg.GetChannel(), e.what());
|
|
}
|
|
}
|
|
}
|
|
|
|
void FileWriter::WriteHDF5(const EndMessage &msg) {
|
|
if (master_file) {
|
|
std::lock_guard<std::mutex> lock(hdf5_mutex);
|
|
|
|
if (format == FileWriterFormat::NXmxIntegrated) {
|
|
try {
|
|
CloseFile(0);
|
|
} catch (...) {
|
|
throw;
|
|
}
|
|
}
|
|
|
|
master_file->Finalize(msg);
|
|
}
|
|
}
|