Files
Jungfraujoch/reader/JFJochStartMessageReader.cpp
T
leonarski_fandClaude Opus 5.5 21f30549a7 rugnux: follow a collection from the broker while it is written
rugnux <path>/<prefix>_master.h5 http://<broker>:<port>/ - when the master file is not on disk yet -
processes the broker's current collection as it is written:

- the start message comes from GET /live/start.cbor; the given path has to end in the collection's
  master file name, and the data files are read next to it;
- JFJochStartMessageReader builds the dataset from the start message (detector, geometry,
  goniometer, mask) and reads the images straight from the data files; a read of an image whose file
  the broker has not reported waits for it;
- BrokerFeed reads GET /live/events in a thread of its own and hands the reader each closed file. On
  end, the run finishes with the images the collection delivered (a read past them returns no image).
  A broken stream, another collection, or a refused follower makes every waiting read throw
  LiveCollectionLost, which ends the run like a fatal resource error;
- the token is taken from JUNGFRAUJOCH_HTTP_TOKEN (the variable the viewer uses) and sent only in
  the Authorization header; without it, and with no master file on disk, rugnux stops at once;
- rotation data and --mode mx only; plain http only (TLS is not linked into rugnux).

With the master file on disk, the URL is ignored and the run is offline as before. The streamed run
does not read the master file and is not compared with it: its result can differ from an offline run
of the same files in the last bits of the metadata.

Tests ([Live]): LiveCollection, the server-sent-event parser, the FileWriter callback, the TCP
pusher's closed-file reports with two real writers, and JFJochStartMessageReader against the
master the writer produced.

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

197 lines
8.2 KiB
C++

// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include "JFJochStartMessageReader.h"
#include <filesystem>
#include "spdlog/fmt/fmt.h"
#include "../common/JFJochException.h"
namespace {
// The data file name the writer gives file `file_number` (HDF5Metadata::DataFileName), without
// the directory: the data files sit next to the master file.
std::string DataFileName(const StartMessage &start, uint64_t file_number) {
const std::string prefix = std::filesystem::path(start.file_prefix).filename().string();
return fmt::format("{:s}_data_{:06d}.h5", prefix, file_number + 1);
}
// The dataset as the start message describes it - the same quantities an HDF5 master states, taken
// from the message the writer wrote it from.
std::shared_ptr<JFJochReaderDataset> Dataset(const StartMessage &start, const DiffractionExperiment &defaults) {
auto dataset = std::make_shared<JFJochReaderDataset>();
dataset->arm_date = start.arm_date;
dataset->jfjoch_release = start.jfjoch_release;
dataset->error_value = start.error_value;
dataset->file_detect_ice_rings = start.detect_ice_rings;
auto &x = dataset->experiment;
x = defaults;
// The reader hands every image out as signed 32-bit, as JFJochHDF5Reader does.
x.BitDepthImage(32);
x.PixelSigned(true);
x.FilePrefix(std::filesystem::path(start.file_prefix).filename().string());
x.BeamX_pxl(start.beam_center_x);
x.BeamY_pxl(start.beam_center_y);
x.DetectorDistance_mm(start.detector_distance * 1000.0f);
x.PoniRot1_rad(start.poni_rot1.value_or(0.0f));
x.PoniRot2_rad(start.poni_rot2.value_or(0.0f));
x.PoniRot3_rad(start.poni_rot3.value_or(0.0f));
x.IncidentEnergy_keV(WVL_1A_IN_KEV / start.incident_wavelength);
if (start.incident_wavelength_spread)
x.BandwidthFWHM(start.incident_wavelength_spread.value() / start.incident_wavelength);
x.DetectIceRings(start.detect_ice_rings.value_or(false));
DetectorSetup detector = DetDECTRIS(start.image_size_x, start.image_size_y, start.detector_description, {});
detector.PixelSize_um(start.pixel_size_x * 1e6f);
detector.MirrorY(start.mirror_y);
detector.ImageOrientation(DetectorOrientation(start.detector_orientation_mirror_y,
start.detector_orientation_quarter_turns));
detector.SensorThickness_um(start.sensor_thickness * 1e6f);
detector.SensorMaterial(start.sensor_material);
detector.SaturationLimit(SaturationLimitFromValue(start.saturation_value));
detector.BitDepthImage(32);
detector.MinFrameTime(std::chrono::microseconds(0));
detector.MinCountTime(std::chrono::microseconds(0));
detector.ReadOutTime(std::chrono::nanoseconds(0));
x.Detector(detector);
x.FrameTime(std::chrono::duration_cast<std::chrono::nanoseconds>(std::chrono::duration<float>(start.frame_time)),
std::chrono::duration_cast<std::chrono::nanoseconds>(std::chrono::duration<float>(start.count_time)));
x.Goniometer(start.goniometer);
x.GridScan(start.grid_scan);
x.Smargon(start.smargon_position);
x.NumTriggers(1);
x.ImagesPerTrigger(static_cast<int64_t>(start.number_of_images));
InstrumentMetadata instrument;
instrument.SourceName(start.source_name);
instrument.SourceType(start.source_type);
instrument.InstrumentName(start.instrument_name);
x.ImportInstrumentMetadata(instrument);
x.SampleName(start.sample_name);
x.SampleTemperature_K(start.sample_temperature_K);
x.RingCurrent_mA(start.ring_current_mA);
x.TotalFlux(start.total_flux);
x.AttenuatorTransmission(start.attenuator_transmission);
if (start.beam_size_x && start.beam_size_y) {
x.BeamSizeX_um(start.beam_size_x.value() * 1e6f);
x.BeamSizeY_um(start.beam_size_y.value() * 1e6f);
}
if (!start.pixel_mask.empty())
dataset->pixel_mask = std::make_shared<const PixelMask>(start.pixel_mask.begin()->second);
else
dataset->pixel_mask = std::make_shared<const PixelMask>(
std::vector<uint32_t>(start.image_size_x * start.image_size_y));
return dataset;
}
}
void JFJochStartMessageReader::Open(const StartMessage &start, const std::string &master_path) {
if (start.images_per_file <= 0)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"The start message gives no images per file");
number_of_images_ = start.number_of_images;
images_per_file_ = static_cast<uint64_t>(start.images_per_file);
const uint64_t files = (number_of_images_ + images_per_file_ - 1) / images_per_file_;
const std::filesystem::path directory = std::filesystem::path(master_path).parent_path();
HDF5ImageLocator::Layout layout;
layout.format = FileWriterFormat::NXmxLegacy;
layout.images_per_file = images_per_file_;
for (uint64_t f = 0; f < files; f++)
layout.legacy_files.push_back({(directory / DataFileName(start, f)).string(), "/entry/data/data"});
{
std::unique_lock ul(hdf5_mutex);
image_source_.Configure(std::move(layout));
}
{
std::unique_lock ul(m_);
file_images_.assign(files, -1);
ended_ = false;
lost_.clear();
}
SetStartMessage(Dataset(start, default_experiment));
}
void JFJochStartMessageReader::FileClosed(uint64_t file_number, uint64_t image_count) {
{
std::unique_lock ul(m_);
if (file_number < file_images_.size())
file_images_[file_number] = static_cast<int64_t>(image_count);
}
cv_.notify_all();
}
void JFJochStartMessageReader::Ended() {
{
std::unique_lock ul(m_);
ended_ = true;
}
cv_.notify_all();
}
void JFJochStartMessageReader::Lost(const std::string &why) {
{
std::unique_lock ul(m_);
if (!ended_)
lost_ = why;
}
cv_.notify_all();
}
uint64_t JFJochStartMessageReader::GetNumberOfImages() const {
return number_of_images_;
}
void JFJochStartMessageReader::Close() {
std::unique_lock ul(hdf5_mutex);
image_source_.Clear();
SetStartMessage({});
}
bool JFJochStartMessageReader::WaitForImage(int64_t image_number) {
if (image_number < 0 || static_cast<uint64_t>(image_number) >= number_of_images_)
throw JFJochException(JFJochExceptionCategory::HDF5, "Image out of bounds");
const uint64_t f = static_cast<uint64_t>(image_number) / images_per_file_;
std::unique_lock ul(m_);
cv_.wait(ul, [&] { return file_images_[f] >= 0 || ended_ || !lost_.empty(); });
if (file_images_[f] < 0) {
if (!lost_.empty())
throw LiveCollectionLost(lost_);
return false; // the collection ended without this file
}
return static_cast<int64_t>(static_cast<uint64_t>(image_number) % images_per_file_) < file_images_[f];
}
bool JFJochStartMessageReader::ReadRawImage(int64_t image_number, JFJochReaderRawImage &ret) {
if (!WaitForImage(image_number))
return false;
// As JFJochHDF5Reader: only the chunk lookup under the HDF5 lock, the read after it.
std::optional<HDF5ImageSource::DirectChunk> chunk;
{
std::unique_lock ul(hdf5_mutex);
const auto loc = image_source_.Resolve(image_number);
chunk = image_source_.PrepareDirectRead(loc);
if (!chunk) {
ret.image = image_source_.ReadImageAt(ret.image_buffer, loc);
return true;
}
}
ret.image = HDF5ImageSource::ReadDirect(ret.image_buffer, *chunk);
return true;
}
bool JFJochStartMessageReader::LoadImage_i(std::shared_ptr<JFJochReaderDataset> &dataset, DataMessage &message,
std::vector<uint8_t> &buffer, int64_t image_number, bool) {
if (!dataset || !WaitForImage(image_number))
return false;
std::unique_lock ul(hdf5_mutex);
message.image = image_source_.ReadImageAt(buffer, image_source_.Resolve(image_number));
message.number = image_number;
return true;
}