Files
Jungfraujoch/receiver/JFJochReceiverLite.cpp
T
leonarski_f 538f3504d3
Build Packages / build:windows:nocuda (push) Successful in 20m4s
Build Packages / Unit tests (push) Skipped
Build Packages / build:viewer-tgz:cpu (push) Successful in 16m5s
Build Packages / build:viewer-tgz:cuda (push) Successful in 17m26s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 27m46s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 20m17s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 26m13s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 23m17s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 28m11s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 19m30s
Build Packages / build:rpm (rocky8) (push) Successful in 24m34s
Build Packages / build:rpm (rocky9) (push) Successful in 21m30s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 23m33s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 20m18s
Build Packages / DIALS test (push) Successful in 18m23s
Build Packages / XDS test (durin plugin) (push) Successful in 11m30s
Build Packages / XDS test (JFJoch plugin) (push) Successful in 10m16s
Build Packages / XDS test (neggia plugin) (push) Successful in 8m2s
Build Packages / Generate python client (push) Successful in 49s
Build Packages / Build documentation (push) Successful in 1m21s
Build Packages / Create release (push) Skipped
Build Packages / build:windows:cuda (push) Successful in 29m45s
v1.0.0.rc-161 (#71)
This is an UNSTABLE release. It includes many experimental features, as well as many AI generated fixes. We recommend using rc.152 for production use.

* **rugnux: significantly better quality of results, and faster.** A large rework of integration, scaling, merging, geometry refinement and space-group determination, together with measurements the program previously made no attempt at - the direct beam before indexing, the beam stop, the goniometer rotation scale, and the stretches of a sweep the crystal did not deliver. A rotation dataset typically gains observations at better <I/sigma> and R_meas, and every `mx` and `scale` run writes a `<prefix>_report.txt` results report modelled on XDS's `CORRECT.LP`. Many defaults moved with it: spot detection is self-calibrating, beam-stop detection and rotation geometry post-refinement are on, resolution limits default to as far as the detector reaches, and ice-ring handling engages only where the crystal is measured to have ice.
* **jfjoch_viewer:** the beam-stop shadow, the detector calibration and the beam-centre measurement are reachable from "Analyze dataset"; the settings panel reports how the sample moved and how polarized the beam was; image rendering and interaction are faster.
* **Performance:** bitshuffle+LZ4 images are decoded on the GPU rather than on the host, with the bitshuffle inverse fused into preprocessing so the decompressed frame is never held in device memory.
* **Broker, writer, packaging and build:** image-slot lifetime and locking fixes, per-image datasets sized by the images actually written, the Debian/Ubuntu broker package renamed to `jfjoch`, and `image_analysis` compiling under MSVC again.

**Breaking change to the rugnux command line:**
* `--azint-only` and `--scale` are **removed**, replaced by `--mode azint` and `--mode scale`; the full pipeline is `--mode mx` and remains the default. A script passing the old flags now fails with the list of valid modes rather than silently running the wrong one.
* `-t`/`--stride` is **refused on rotation data**: skipping frames cuts every reflection's rocking curve, so the combined fulls and their partiality would be measured over frames the sweep never recorded. Select a contiguous range with `-s`/`-e` instead. `--mode azint` and `--force-still` still take a stride.

**Breaking changes to OpenAPI** - regenerate the client (`jfjoch-client` 1.0.0-rc.161, `frontend/src/client`) or read the affected fields as optional:
* `image_scale_b` is removed from the `plot_type` enum, so a client requesting that plot now gets an error rather than a curve.
* `azim_int_settings.high_q_recipA`, `spot_finding_settings.high_resolution_limit` and `spot_finding_settings.low_resolution_limit` are no longer `required`. All three mean "no limit at that end" when unset and are omitted from the response instead of carrying a placeholder value, which raises in a client generated from an rc.160-or-earlier spec. A value of 0 is still accepted and means the same thing.

**Breaking changes to the stored formats** - a consumer reading these fields must treat them as optional:
* The per-image image-scale B factor is no longer computed, so `/entry/MX/imageScaleBFactor` is absent from newly written HDF5 files and the corresponding key is absent from the CBOR DataMessage and END blocks. Files written by rc.160 and earlier still contain it and still open; nothing in the pipeline reads it any more.
* `_reflns.jfjoch_diffrn_ISa` now carries the whole-range `1/sqrt(a*b)` that XDS's ISa denotes, and the error-model `a` and `b` are reported in XDS's convention; the strong-reflection asymptote moves to `_reflns.jfjoch_diffrn_ISa_asymptotic`. **A file written by an earlier version carries the asymptote under the plain `ISa` name.**

Reviewed-on: #71
Co-authored-by: Filip Leonarski <filip.leonarski@psi.ch>
2026-08-13 17:03:10 +02:00

432 lines
18 KiB
C++

// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include "JFJochReceiverLite.h"
#include "../image_analysis/indexing/IndexerFactory.h"
#include "../common/CUDAWrapper.h"
using namespace std::chrono_literals;
namespace {
// Serialize a processed image into its buffer slot. The slot holds the compressed
// image plus the per-image CBOR metadata (spot list, reflection list, ...); if an
// unusually rich frame makes that metadata too large to fit alongside the image,
// drop just this frame (log + recycle the slot) instead of letting the serializer
// throw, which would abort the whole data collection. Returns true when the image
// was serialized and the slot is ready for the caller to commit.
bool SerializeImageOrDrop(CBORStream2Serializer &serializer, DataMessage &msg,
ZeroCopyReturnValue &loc, size_t slot_size, Logger &logger) {
try {
serializer.SerializeImage(msg);
return true;
} catch (const JFJochException &) {
// Confirm this is a slot-overflow (not some other CBOR error) by re-serializing
// the metadata alone with an empty image placeholder: that always fits when the
// image is what overflowed, and gives the exact metadata size for the log. If it
// throws too, it is not a simple overflow and the exception propagates.
const auto image = msg.image;
msg.image = CompressedImage(nullptr, 0, image.GetWidth(), image.GetHeight(),
image.GetMode(), image.GetCompressionAlgorithm(),
image.GetChannel());
CBORStream2Serializer meas(static_cast<uint8_t *>(loc.GetImage()), slot_size);
meas.SerializeImage(msg);
const size_t metadata_size = meas.GetImageAppendOffset();
msg.image = image;
logger.Error("Dropping image {} - CBOR metadata too large to store image: "
"metadata {} B + image {} B do not fit in buffer slot {} B",
msg.number, metadata_size, image.GetCompressedSize(), slot_size);
loc.release();
return false;
}
}
}
int64_t JFJochReceiverLite::NumberOfDataAnalysisThreads(int64_t requested_thread_number,
const DiffractionExperiment& in_experiment) {
int64_t number_of_images = in_experiment.GetImageNum();
auto image_time = in_experiment.GetImageTime();
if (requested_thread_number <= 0)
return 1;
// For very small datasets no need to go multithreaded
if (number_of_images < 4)
return 1;
if (number_of_images < 16)
return std::min<int64_t>(2, requested_thread_number);
// For 100 Hz, there is no reason to go for a very large number of threads
if (number_of_images < 256 || (image_time >= 10ms))
return std::min<int64_t>(16, requested_thread_number);
return requested_thread_number;
}
JFJochReceiverLite::JFJochReceiverLite(const DiffractionExperiment &in_experiment,
const PixelMask &in_pixel_mask,
ImagePuller &in_image_puller,
ImagePusher &in_image_pusher,
Logger &in_logger,
int64_t forward_and_sum_nthreads,
const SpotFindingSettings &in_spot_finding_settings,
JFJochReceiverCurrentStatus &in_current_status,
JFJochReceiverPlots &in_plots,
ImageBuffer &in_image_buffer,
ZMQPreviewSocket *in_zmq_preview_socket,
ZMQMetadataSocket *in_zmq_metadata_socket,
IndexerThreadPool *indexing_thread_pool)
: JFJochReceiver(in_experiment,
in_image_buffer,
in_image_pusher,
in_current_status,
in_plots,
in_spot_finding_settings,
in_logger,
in_pixel_mask,
in_zmq_preview_socket,
in_zmq_metadata_socket,
indexing_thread_pool),
image_puller(in_image_puller),
data_analysis_nthreads(NumberOfDataAnalysisThreads(forward_and_sum_nthreads, in_experiment)),
data_analysis_started(data_analysis_nthreads),
measurement_started(1),
dark_mask_analysis(in_experiment.GetDarkMaskSettings(), in_experiment.GetPixelsNum()) {
logger.Info("Starting {} data analysis threads", data_analysis_nthreads);
if (experiment.GetDetectorMode() == DetectorMode::DarkMask) {
for (int i = 0; i < data_analysis_nthreads; i++)
data_analysis_futures.emplace_back(
std::async(std::launch::async, &JFJochReceiverLite::MaskThread, this, i)
);
} else {
// Start frame transformation threads
for (int i = 0; i < data_analysis_nthreads; i++)
data_analysis_futures.emplace_back(
std::async(std::launch::async, &JFJochReceiverLite::DataAnalysisThread, this, i)
);
}
measurement = std::async(std::launch::async, &JFJochReceiverLite::MeasurementThread, this);
data_analysis_started.wait();
logger.Info("Data analysis threads ready");
}
JFJochReceiverLite::~JFJochReceiverLite() {
cancelled = true;
try {
if (measurement.valid())
measurement.get();
} catch (const std::exception& e) {
logger.Error("ReceiverLite teardown failed: {}", e.what());
} catch (...) {
logger.Error("ReceiverLite teardown failed with unknown exception");
}
}
void JFJochReceiverLite::MeasurementThread() {
try {
// Wait for start message to arrive
auto msg = image_puller.PollImage();
while (!cancelled && (!msg.has_value() || !msg->cbor || !msg->cbor->start_message))
msg = image_puller.PollImage();
if (cancelled) {
current_status.SetProgress({});
current_status.SetStatus(GetStatus());
measurement_started.count_down();
return; // just quit the function - no measurement, no results
}
Configure(msg->cbor->start_message.value());
image_buffer.StartMeasurement(experiment); // Only at this point we know bit-depth of the images
start_time = std::chrono::system_clock::now();
// Send new start message out
SendStartMessage();
measurement_started.count_down();
} catch (const JFJochException &e) {
logger.ErrorException(e);
Cancel(e);
measurement_started.count_down();
throw;
}
logger.Info("Receiving started");
// Analysis is running
// ...
// Till it is done.
// Combine frame transformation threads
for (auto &f: data_analysis_futures)
f.get();
logger.Info("Data analysis threads finished");
current_status.SetProgress(1.0);
current_status.SetStatus(GetStatus());
// Send end message out
SendEndMessage();
if (!image_buffer.CheckIfBufferReturned(std::chrono::seconds(10))) {
logger.Error("Send commands not finalized in 10 seconds");
throw JFJochException(JFJochExceptionCategory::ZeroMQ, "Send commands not finalized in 10 seconds");
}
logger.Info("All images sent through ZeroMQ");
if (push_images_to_writer)
writer_error = image_pusher.Finalize();
logger.Info("Writing process finalized");
end_time = std::chrono::system_clock::now();
current_status.SetProgress({});
current_status.SetStatus(GetStatus());
}
void JFJochReceiverLite::Configure(const StartMessage &msg) {
if ((experiment.GetXPixelsNum() != msg.image_size_x)
|| (experiment.GetYPixelsNum() != msg.image_size_y))
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Mismatch in detector size");
if (fabs(experiment.GetPixelSize_mm() - msg.pixel_size_x * 1e3) > 1e-7)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Mismatch in pixel size");
experiment.Detector().Description(msg.detector_description);
experiment.Detector().SerialNumber(msg.detector_serial_number);
experiment.Detector().SensorMaterial(msg.sensor_material);
experiment.Detector().SensorThickness_um(msg.sensor_thickness * 1e6);
experiment.Detector().SaturationLimit(msg.saturation_value);
experiment.Detector().BitDepthImage(msg.bit_depth_image);
if (msg.bit_depth_readout)
experiment.Detector().BitDepthReadout(msg.bit_depth_readout.value());
}
void JFJochReceiverLite::MaskThread(uint32_t id) {
std::vector<uint8_t> buffer;
data_analysis_started.count_down();
measurement_started.wait();
while (!cancelled && !end_message_received) {
try {
auto msg = image_puller.PollImage(std::chrono::milliseconds(5));
if (!msg.has_value() || !msg->cbor) {
// Message not received or not parsed
continue;
} else if (msg->cbor->end_message) {
end_message_received = true;
logger.Debug("Thread {} end message received in JFJochReceiverLite", id);
} else if (msg->cbor->data_message) {
++images_collected;
DataMessage data_msg = msg->cbor->data_message.value();
dark_mask_analysis.AnalyzeImage(data_msg, buffer);
compressed_size += data_msg.image.GetCompressedSize();
uncompressed_size += data_msg.image.GetUncompressedSize();
UpdateMaxImageReceived(data_msg.number);
auto loc = image_buffer.GetImageSlot();
if (loc == nullptr)
writer_queue_full = true;
else {
auto writer_buffer = (uint8_t *) loc->GetImage();
CBORStream2Serializer serializer(writer_buffer, experiment.GetImageBufferLocationSize());
if (SerializeImageOrDrop(serializer, data_msg, *loc,
experiment.GetImageBufferLocationSize(), logger)) {
loc->SetImageNumber(data_msg.number);
loc->SetImageSize(serializer.GetBufferSize());
loc->SetIndexed(false);
loc->release();
}
}
}
current_status.SetProgress(GetProgress());
current_status.SetStatus(GetStatus());
} catch (const JFJochException &e) {
logger.ErrorException(e);
Cancel(e);
}
}
}
void JFJochReceiverLite::DataAnalysisThread(uint32_t id) {
std::unique_ptr<MXAnalysisWithoutFPGA> analysis;
logger.Debug("Thread {} started", id);
try {
pin_gpu();
} catch (const JFJochException &e) {
logger.Warning("Error pinning GPU {}", e.what());
}
data_analysis_started.count_down();
measurement_started.wait();
try {
analysis = std::make_unique<MXAnalysisWithoutFPGA>(experiment, *az_int_mapping, pixel_mask, indexer,
/*enable_fused_adaptive_gpu=*/true);
} catch (const JFJochException &e) {
Cancel(e);
return;
}
while (!cancelled && !end_message_received) {
try {
auto msg = image_puller.PollImage(std::chrono::milliseconds(5));
if (!msg.has_value() || !msg->cbor) {
// Message not received or not parsed
continue;
} else if (msg->cbor->end_message) {
// Message is end message
end_message_received = true;
logger.Debug("Thread {} End message received in JFJochReceiverLite", id);
} else if (msg->cbor->calibration) {
// Calibration messages are just forwarded
if (push_images_to_writer)
image_pusher.SendCalibration(msg->cbor->calibration.value());
} else if (msg->cbor->data_message) {
auto start = std::chrono::high_resolution_clock::now();
DataMessage data_msg = msg->cbor->data_message.value();
++images_collected;
auto compressed_size_img = data_msg.image.GetCompressedSize();
auto uncompressed_size_img = data_msg.image.GetUncompressedSize();
compressed_size += compressed_size_img;
uncompressed_size += uncompressed_size_img;
UpdateMaxImageReceived(data_msg.number);
auto image_start_time = std::chrono::high_resolution_clock::now();
AzimuthalIntegrationProfile profile(*az_int_mapping);
analysis->Analyze(data_msg, profile, GetSpotFindingSettings());
auto image_end_time = std::chrono::high_resolution_clock::now();
std::chrono::duration<float> image_duration = image_end_time - image_start_time;
data_msg.processing_time_s = image_duration.count();
data_msg.original_number = data_msg.number;
data_msg.user_data = experiment.GetImageAppendix();
data_msg.run_number = experiment.GetRunNumber();
data_msg.run_name = experiment.GetRunName();
data_msg.receiver_buf_available = image_buffer.GetAvailSlots();
data_msg.receiver_aq_dev_delay = image_puller.GetCurrentFifoUtilization();
// Compression ratio is uncompressed/compressed (e.g. 7x), matching the FPGA path
// (JFJochReceiverFPGA) and the aggregate (JFJochReceiver::GetFinalStatistics) - not the
// reciprocal.
if (compressed_size_img > 0)
data_msg.compression_ratio = static_cast<float>(uncompressed_size_img)
/ static_cast<float>(compressed_size_img);
saturated_pixels.Add(data_msg.saturated_pixel_count);
error_pixels.Add(data_msg.error_pixel_count);
if (data_msg.roi.contains("beam")) {
roi_beam_npixel.Add(data_msg.roi["beam"].pixels);
roi_beam_sum.Add(data_msg.roi["beam"].sum);
}
plots.Add(data_msg, profile);
scan_result.Add(data_msg);
if (!serialmx_filter.ApplyFilter(data_msg))
++images_skipped;
else {
auto loc = image_buffer.GetImageSlot();
if (loc == nullptr)
writer_queue_full = true;
else {
auto writer_buffer = (uint8_t *) loc->GetImage();
CBORStream2Serializer serializer(writer_buffer, experiment.GetImageBufferLocationSize());
if (SerializeImageOrDrop(serializer, data_msg, *loc,
experiment.GetImageBufferLocationSize(), logger)) {
loc->SetImageNumber(data_msg.number);
loc->SetImageSize(serializer.GetBufferSize());
loc->SetIndexed(data_msg.indexing_result.value_or(false));
loc->ReadyToSend();
if (zmq_preview_socket != nullptr)
zmq_preview_socket->SendImage(writer_buffer, serializer.GetBufferSize());
if (zmq_metadata_socket != nullptr)
zmq_metadata_socket->AddDataMessage(data_msg);
if (push_images_to_writer) {
// Only count images the pusher accepted; TCP can drop on a
// broken/absent connection or an expired enqueue deadline.
if (image_pusher.SendImage(*loc))
++images_sent;
} else
loc->release();
UpdateMaxImageSent(data_msg.number);
}
}
}
auto end = std::chrono::high_resolution_clock::now();
std::chrono::duration<float> duration_total = end - start;
logger.Debug("Thread {} Image {:6d} processing {:8.04f} s analysis {:8.04f} s",
id, data_msg.number, duration_total.count(), image_duration.count());
}
UpdateMaxDelay(image_puller.GetCurrentFifoUtilization());
current_status.SetProgress(GetProgress());
current_status.SetStatus(GetStatus());
} catch (const JFJochException &e) {
logger.ErrorException(e);
Cancel(e);
}
}
logger.Debug("Thread {} finished", id);
}
void JFJochReceiverLite::StopReceiver() {
if (measurement.valid()) {
measurement.get();
logger.Info("Receiver stopped");
}
}
float JFJochReceiverLite::GetEfficiency() const {
if (experiment.GetFrameNum() == 0)
return 0;
return static_cast<float>(images_collected) / static_cast<float>(experiment.GetFrameNum());
}
float JFJochReceiverLite::GetProgress() const {
if (experiment.GetFrameNum() == 0)
return 0.0;
return static_cast<float>(max_image_number_received) / static_cast<float>(experiment.GetFrameNum());
}
void JFJochReceiverLite::SetSpotFindingSettings(const SpotFindingSettings &in_spot_finding_settings) {
std::unique_lock ul(spot_finding_settings_mutex);
DiffractionExperiment::CheckDataProcessingSettings(in_spot_finding_settings);
spot_finding_settings = in_spot_finding_settings;
}
JFJochReceiverOutput JFJochReceiverLite::GetFinalStatistics() const {
JFJochReceiverOutput ret = JFJochReceiver::GetFinalStatistics();
if (experiment.GetDetectorMode() == DetectorMode::DarkMask)
ret.dark_mask_result = dark_mask_analysis.GetMask();
return ret;
}