// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute // SPDX-License-Identifier: GPL-3.0-only #include "FileWriter.h" #include #include #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 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(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(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(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 FileWriter::Finalize() { std::lock_guard 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& 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(ZMQSocketType::Pub); finalized_file_socket->Bind(addr); } std::optional FileWriter::GetZMQAddr() { if (finalized_file_socket) { return finalized_file_socket->GetEndpointName(); } else return {}; } void FileWriter::CreateHDF5MasterFile(const StartMessage &msg) { std::lock_guard lock(hdf5_mutex); master_file = std::make_unique(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 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(default_images_per_file) : msg.images_per_file; const int64_t num_files = (static_cast(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 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 lock(hdf5_mutex); if (format == FileWriterFormat::NXmxIntegrated) { try { CloseFile(0); } catch (...) { throw; } } master_file->Finalize(msg); } }