Files
Jungfraujoch/writer/FileWriter.cpp
T
jungfrauandClaude Opus 5 c3236eed44 Publish the underload, not the error marker, on the finalized-file socket
The finalized-file notification has carried j["underload"] = error_value since the two meant the
same thing. They stopped meaning the same thing when error_value became the marker the pixels
actually store: UINTx_MAX for an unsigned image, where it used to be GetUnderflow()'s -1, a value
no unsigned pixel can hold and which therefore excluded nothing.

So for an unsigned 16-bit run the key went from -1 to 65535. A facility that forwards it into an
XDS UNDERLOAD or a DIALS trusted range - which is what a key called "underload" is for - would
reject every pixel below 65535, i.e. all of them. Nothing in HDF5 is affected and none of the
writer tests look at this socket, so it fails silently and outside the file.

The start message already carries underload_value, the lowest valid value: 0 for an unsigned image
and INTx_MIN+1 for a signed one. Send that. The signed case moves too, from the marker itself to
one above it, which is what the key has always claimed to be.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_011n8riB6X59oRjkrSHzNPAU
2026-08-23 08:31:35 -04:00

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);
}
}