Files
Jungfraujoch/receiver/JFJochReceiverLite.cpp
T
leonarski_f 9aae0c2ba7
Build Packages / Create release (push) Successful in 21s
Build Packages / build:rugnux-tgz (x86_64) (push) Successful in 9m40s
Build Packages / build:rugnux:aarch64 (cross) (push) Successful in 9m49s
Build Packages / build:viewer-tgz:cpu (push) Successful in 11m37s
Build Packages / build:viewer-tgz:cuda (push) Successful in 12m40s
Build Packages / build:windows:nocuda (push) Successful in 17m44s
Build Packages / build:windows:cuda (push) Successful in 20m13s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 14m41s
Build Packages / HDF5 consumer tests (DIALS, XDS) (push) Successful in 25m59s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 15m5s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 14m35s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 15m53s
Build Packages / build:rugnux:windows (push) Successful in 11m29s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 18m51s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 18m43s
Build Packages / Generate python client (push) Successful in 51s
Build Packages / build:rpm (rocky8) (push) Successful in 18m51s
Build Packages / Build documentation (push) Successful in 1m21s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 18m38s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 18m24s
Build Packages / build:rpm (rocky9) (push) Successful in 19m19s
Build Packages / Unit tests (push) Successful in 1h37m15s
v1.0.0-rc.169 (#79)
* Building Jungfraujoch no longer needs zlib or Eigen installed on the machine, and the dependencies the build fetches are pinned and updated to current releases.
* rugnux: improvements in indexing, lattice selection and geometry post-refinement, which index crystals that previously returned no lattice and keep the better of the two geometries a run measures.
* rugnux: improvements in beam-centre measurement, beam-stop detection and space-group determination.
* rugnux: the unit cell reported with a determined space group now obeys that group - a cell whose symmetry was confirmed from the intensities is re-refined under it, and a cell the group cannot describe is reported with a warning rather than as it stands.
* rugnux drops the stretches of a rotation sweep whose removal measurably improves the merged intensities and reports what became of every frame, and decides the resolution cut on the crystal's own diffraction rather than on its ice rings.
* The rugnux results report is machine-readable - every line that is not `KEY= value` data starts with `#` - and states the build it was written by, its authorship and its terms of use (`REPORT_VERSION= 8`).
* `jfjoch_viewer`: improvements in the file manager (CBF frames beside HDF5 datasets, a remembered root), the dataset plots, the inspector and the image statistics, plus a settable font size, a view of the rugnux results report, usable performance over a remote display (`ssh -X`) and a reset of all settings to defaults; the reciprocal-space window is removed.
* Broker fixes around DECTRIS collections and dark-mask calibration: re-initialising after a run that never started no longer freezes the broker, a cancelled calibration is abandoned instead of reported as done, and a collection whose start message never arrives ends by itself.

Reviewed-on: #79
Co-authored-by: Filip Leonarski <filip.leonarski@psi.ch>
2026-09-15 17:09:31 +02:00

452 lines
19 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));
// 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;
}