Build Packages / build:rugnux:aarch64 (cross) (push) Successful in 8m55s
Build Packages / build:windows:nocuda (push) Successful in 17m1s
Build Packages / build:rugnux-tgz (x86_64) (push) Successful in 19m6s
Build Packages / build:windows:cuda (push) Successful in 19m11s
Build Packages / build:viewer-tgz:cpu (push) Successful in 22m19s
Build Packages / build:viewer-tgz:cuda (push) Successful in 23m19s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 23m35s
Build Packages / build:rugnux:windows (push) Successful in 10m50s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 28m32s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 28m50s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 20m7s
Build Packages / build:windows:nocuda (pull_request) Successful in 16m9s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 25m52s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 22m57s
Build Packages / build:rpm (rocky9) (push) Successful in 23m30s
Build Packages / build:windows:cuda (pull_request) Successful in 21m21s
Build Packages / build:rpm (rocky8) (push) Successful in 27m47s
Build Packages / build:rugnux:windows (pull_request) Successful in 15m53s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 22m47s
Build Packages / Generate python client (push) Successful in 39s
Build Packages / Create release (push) Skipped
Build Packages / Build documentation (push) Successful in 1m56s
Build Packages / XDS test (JFJoch plugin) (push) Successful in 10m23s
Build Packages / XDS test (durin plugin) (push) Successful in 10m34s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 27m0s
Build Packages / DIALS test (push) Successful in 26m37s
Build Packages / XDS test (neggia plugin) (push) Successful in 9m15s
Build Packages / build:rugnux:aarch64 (cross) (pull_request) Successful in 9m8s
Build Packages / build:viewer-tgz:cpu (pull_request) Successful in 15m35s
Build Packages / build:rugnux-tgz (x86_64) (pull_request) Successful in 16m50s
Build Packages / build:viewer-tgz:cuda (pull_request) Successful in 18m44s
Build Packages / build:rpm (rocky9_nocuda) (pull_request) Successful in 19m23s
Build Packages / build:rpm (rocky8_nocuda) (pull_request) Successful in 21m53s
Build Packages / build:rpm (ubuntu2204_nocuda) (pull_request) Successful in 21m5s
Build Packages / build:rpm (ubuntu2404_nocuda) (pull_request) Successful in 18m54s
Build Packages / build:rpm (rocky9_sls9) (pull_request) Successful in 21m30s
Build Packages / build:rpm (rocky8_sls9) (pull_request) Successful in 24m55s
Build Packages / build:rpm (rocky9) (pull_request) Successful in 21m46s
Build Packages / build:rpm (rocky8) (pull_request) Successful in 24m50s
Build Packages / build:rpm (ubuntu2404) (pull_request) Successful in 19m20s
Build Packages / Generate python client (pull_request) Successful in 37s
Build Packages / XDS test (durin plugin) (pull_request) Successful in 11m6s
Build Packages / Create release (pull_request) Skipped
Build Packages / Build documentation (pull_request) Successful in 1m24s
Build Packages / build:rpm (ubuntu2204) (pull_request) Successful in 23m48s
Build Packages / XDS test (JFJoch plugin) (pull_request) Successful in 11m1s
Build Packages / XDS test (neggia plugin) (pull_request) Successful in 9m50s
Build Packages / DIALS test (pull_request) Successful in 20m23s
Build Packages / Unit tests (push) Failing after 2h48m34s
Build Packages / Unit tests (pull_request) Failing after 2h33m16s
The pre-flight only checked the master file name, so a data file already sitting in the output's way was refused only when the run had been written in full and the first colliding data file failed to rename into place at the end of it. Checking every candidate name individually would cost one round trip per file on a network filesystem, so a single directory listing is taken once and compared against the run's expected file names in memory.
312 lines
12 KiB
C++
312 lines
12 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(start_message, format);
|
|
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 StartMessage &msg, FileWriterFormat format) {
|
|
if (msg.overwrite.value_or(false))
|
|
return;
|
|
|
|
// Checked only by the single writer that owns the master file (write_master_file - index 0 in
|
|
// a multi-writer TCP/ZMQ setup); every writer of a run shares the same output directory, so
|
|
// one of them answers for all of it.
|
|
if (!msg.write_master_file.value_or(false))
|
|
return;
|
|
|
|
const bool nxmx = format == FileWriterFormat::NXmxLegacy
|
|
|| format == FileWriterFormat::NXmxVDS
|
|
|| format == FileWriterFormat::NXmxIntegrated;
|
|
if (!nxmx)
|
|
return;
|
|
|
|
const std::string master_name = HDF5Metadata::MasterFileName(msg);
|
|
if (std::filesystem::exists(master_name))
|
|
throw JFJochException(JFJochExceptionCategory::FileWriteError,
|
|
"Output file already exists and overwrite is off: " + master_name);
|
|
|
|
if (format == FileWriterFormat::NXmxIntegrated)
|
|
return; // Data lives in the master file, already checked above.
|
|
|
|
// A data file in the way is refused too - this is the collision that used to slip past the
|
|
// pre-flight and only surface when the first colliding data file was renamed into place at
|
|
// the end of a fully-written run. Stat-ing every candidate name individually would cost one
|
|
// round trip per file: thousands of lookups for a long run with a small images_per_file, which
|
|
// on a network filesystem overruns the pre-flight's ACK budget. A single directory listing
|
|
// costs one round trip regardless of how many files the run will write, so collect what is
|
|
// actually there once and compare names in memory.
|
|
std::filesystem::path prefix_path(msg.file_prefix);
|
|
std::filesystem::path dir = prefix_path.has_parent_path() ? prefix_path.parent_path()
|
|
: std::filesystem::path(".");
|
|
|
|
std::unordered_set<std::string> existing;
|
|
for (const auto &entry : std::filesystem::directory_iterator(dir))
|
|
if (entry.is_regular_file())
|
|
existing.insert(entry.path().filename().string());
|
|
|
|
const int64_t images_per_file = (msg.images_per_file <= 0)
|
|
? static_cast<int64_t>(default_images_per_file) : msg.images_per_file;
|
|
const int64_t num_files = (static_cast<int64_t>(msg.number_of_images) + images_per_file - 1)
|
|
/ images_per_file;
|
|
|
|
for (int64_t file_number = 0; file_number < num_files; file_number++) {
|
|
const std::string name = HDF5Metadata::DataFileName(msg, file_number);
|
|
if (existing.contains(std::filesystem::path(name).filename().string()))
|
|
throw JFJochException(JFJochExceptionCategory::FileWriteError,
|
|
"Output file already exists and overwrite is off: " + name);
|
|
}
|
|
}
|
|
|
|
void FileWriter::Preflight(const StartMessage &request, bool trusted_path) {
|
|
if (!trusted_path)
|
|
CheckPath(request.file_prefix);
|
|
MakeDirectory(request.file_prefix);
|
|
|
|
FileWriterFormat format = FileWriterFormat::NXmxLegacy;
|
|
if (request.file_format)
|
|
format = request.file_format.value();
|
|
|
|
CheckOutputFilesAvailable(request, format);
|
|
}
|
|
|
|
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);
|
|
}
|
|
}
|