Files
Jungfraujoch/receiver/JFJochReceiverLite.cpp
T
leonarski_f 84228bf8be
Build Packages / Create release (push) Successful in 24s
Build Packages / build:viewer:macos-arm64:nocuda (push) Successful in 3m29s
Build Packages / build:rugnux:macos-arm64:nocuda (push) Successful in 2m43s
Build Packages / build:rugnux:linux-aarch64:cuda (push) Successful in 8m27s
Build Packages / build:rugnux:linux-x86_64:cuda (push) Successful in 9m53s
Build Packages / build:viewer:linux-x86_64:nocuda (push) Successful in 9m58s
Build Packages / build:viewer:linux-x86_64:cuda (push) Successful in 11m22s
Build Packages / build:jfjoch:rocky8:nocuda (push) Successful in 13m39s
Build Packages / build:viewer:windows-x86_64:nocuda (push) Successful in 18m37s
Build Packages / build:jfjoch:rocky9:nocuda (push) Successful in 16m32s
Build Packages / build:viewer:windows-x86_64:cuda (push) Successful in 24m11s
Build Packages / HDF5 consumer tests (DIALS, XDS) (push) Successful in 25m30s
Build Packages / build:jfjoch:ubuntu2404:nocuda (push) Successful in 19m3s
Build Packages / build:jfjoch:ubuntu2204:nocuda (push) Successful in 20m23s
Build Packages / build:jfjoch:rocky8:cuda-sls9 (push) Successful in 19m41s
Build Packages / Generate python client (push) Successful in 50s
Build Packages / Build documentation (push) Successful in 1m16s
Build Packages / build:jfjoch:rocky9:cuda-sls9 (push) Successful in 21m0s
Build Packages / build:jfjoch:rocky8:cuda (push) Successful in 18m38s
Build Packages / build:rugnux:windows-x86_64:cuda (push) Successful in 14m33s
Build Packages / build:jfjoch:rocky9:cuda (push) Successful in 17m55s
Build Packages / build:jfjoch:ubuntu2204:cuda (push) Successful in 20m50s
Build Packages / build:jfjoch:ubuntu2404:cuda (push) Successful in 18m38s
Build Packages / Unit tests (push) Successful in 1h46m14s
v1.0.0-rc.173 (#83)
* jfjoch_broker: Optional per-dataset authentication - statistics, images and plots can require a bearer token, which jfjoch_viewer supports.
* jfjoch_viewer: Dark mode and a theme-matched colour scheme, a magnifier panel, and simpler contrast and background controls.
* Rugnux: Multiple performance improvements on GPU and CPU (CPU-only processing up to 40% faster, faster image decoding on ARM), with unchanged results.
* Rugnux: `--model` rigid-body refinement runs on the GPU, and the model-validation check is faster and more reliable.
* Rugnux: Improved scaling and merging - error model, outlier rejection, absorption correction and French-Wilson amplitudes now agree more closely with XDS and ctruncate.
* Rugnux: Improved integration - radial background on powder and ice rings, crowded rotation data keep their reflections, and CPU-only builds integrate large unit cells as GPU builds do.
* Rugnux: More robust detector geometry - measured beam centre, X-ray bandwidth and goniometer rate, and geometry refinement accepted only on significant evidence.
* Rugnux: Merged files are written in the standard setting, or in the setting of a reference MTZ, structure-factor mmCIF or model, with its free-R flags.
* Rugnux: Richer report - ice and powder rings, further lattices, superstructure candidates and mosaicity, with warnings worded as prompts to check.
* Rugnux: Clear error messages when a data set needs more GPU or host memory than is available.

Reviewed-on: #83
Co-authored-by: Filip Leonarski <filip.leonarski@psi.ch>
2026-09-29 15:57:32 +02:00

462 lines
20 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. Once the detector has reported it is done, nothing
// more is coming: allow a start message that is still in flight to land, then give up.
// Without this the poll runs for ever - a detector that never streamed leaves Stop()
// waiting on this thread with no way out other than a cancel from the operator.
std::optional<std::chrono::steady_clock::time_point> give_up_at;
auto msg = image_puller.PollImage();
while (!cancelled && (!msg.has_value() || !msg->cbor || !msg->cbor->start_message)) {
if (detector_finished && !give_up_at)
give_up_at = std::chrono::steady_clock::now() + StartMessageGrace;
if (give_up_at && (std::chrono::steady_clock::now() > *give_up_at)) {
logger.Error("Detector finished, but no start message arrived - abandoning measurement");
cancelled = true;
break;
}
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(SaturationLimitFromValue(msg.saturation_value));
// The detector decides the corrections applied to the pixels (DECTRIS enables them by default);
// the outgoing start message is rebuilt from the experiment, so without copying them NXmx would
// record them as not applied.
experiment.Detector().CountRateCorrectionApplied(msg.countrate_correction_enabled);
experiment.Detector().CountRateCorrectionLookupTable(msg.countrate_correction_lookup_table);
experiment.Detector().FlatfieldApplied(msg.flatfield_enabled);
experiment.Detector().VirtualPixelInterpolationApplied(msg.virtual_pixel_interpolation_enabled);
// The mask written with the data is ours (the detector's plus user masking); it is never
// uploaded to the detector, so the pixels arrive without it applied.
experiment.ApplyPixelMask(false);
// Images are forwarded byte-for-byte, so the stream's own image_dtype - not anything configured
// locally - decides both the width and the sign the outgoing metadata must declare. Taking only
// the width used to leave a signed stream declared unsigned in NXmx.
experiment.Detector().BitDepthImage(msg.bit_depth_image);
experiment.PixelSigned(msg.pixel_signed);
}
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::Cancel(bool silent) {
JFJochReceiver::Cancel(silent);
// A silent cancel is not a request to abort - it says the detector has finished - so it must
// not set cancelled. It does bound the wait for the start message: see MeasurementThread.
detector_finished = true;
}
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;
}