Files
leonarski_fandClaude Opus 5 e381d2fd50
Build Packages / build:viewer-tgz:cpu (push) Successful in 12m14s
Build Packages / build:rugnux:aarch64 (cross) (push) Successful in 8m24s
Build Packages / build:rpm (rocky8_nocuda) (push) Failing after 5m36s
Build Packages / build:rugnux-tgz (x86_64) (push) Successful in 12m55s
Build Packages / build:viewer-tgz:cuda (push) Successful in 14m11s
Build Packages / build:rpm (rocky9_nocuda) (push) Failing after 5m6s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Failing after 4m11s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Failing after 3m58s
Build Packages / build:rpm (rocky8) (push) Failing after 3m22s
Build Packages / build:rpm (rocky9_sls9) (push) Failing after 4m12s
Build Packages / build:rpm (rocky9) (push) Failing after 3m20s
Build Packages / build:rpm (rocky8_sls9) (push) Failing after 4m23s
Build Packages / build:rpm (ubuntu2204) (push) Failing after 3m36s
Build Packages / build:windows:nocuda (push) Successful in 17m9s
Build Packages / build:rpm (ubuntu2404) (push) Failing after 4m8s
Build Packages / Generate python client (push) Successful in 35s
Build Packages / Build documentation (push) Successful in 52s
Build Packages / Create release (push) Skipped
Build Packages / build:windows:cuda (push) Successful in 19m32s
Build Packages / XDS test (durin plugin) (push) Successful in 7m22s
Build Packages / XDS test (JFJoch plugin) (push) Successful in 7m28s
Build Packages / XDS test (neggia plugin) (push) Successful in 7m15s
Build Packages / build:rugnux:windows (push) Successful in 12m10s
Build Packages / DIALS test (push) Successful in 17m32s
Build Packages / Unit tests (push) Successful in 1h16m29s
grid scan: review fixes - the ice channel sees its own spots, and a needle is not mirrored
Two independent reviews of the merged grid-scan work. The findings that changed
behaviour:

The ice score's spot channel was fed a list the spot budget had already stripped.
FilterSpotsByCount orders ice-band spots LAST when indexing is not to use them, so on
a frame with more spots than the budget the ice spots are the first discarded - and
the channel that exists for "ice arrives as discrete spots and leaves the radial
profile flat" then read zero on exactly the frames it was written for. Probed at 3000
spots with 1200 on the hexagonal radii and a budget of 1000: 1.000 before the cap,
0.000 after. IceScore now takes d-spacings and is handed the list from before the cap.

The viewer scaled the crystal box by the SIGNED grid step, where every other consumer
takes the magnitude. On a negative step that mirrors the box - +30 deg drawn as -30 -
and hands QRectF a negative width.

rugnux --mode raster never put its settings on the experiment, so the indexing switch
was read at its default while a deprecated per-run flag did the actual work; and
RugnuxCommandLine emitted no --mode for Grid, so a raster job copied to a cluster ran
the default mx - indexing, integrating and merging every cell of the raster.

res_A is NaN where nothing in a blob measured a resolution, and nlohmann writes NaN as
null, which the schema and the generated clients both reject. It is now left unset.

The broker's configuration example named a key that does not exist (calibration, not
calibration_settings); nlohmann ignores unknown keys, so a user copying it got a
silently ignored block. The changelog had lost the rc.166 heading and 21 rc.167
entries to a bad edit of mine, and three entries had been filed under rc.166.

Also: a warning where mode Grid meets a dataset with no grid scan, which was silent
and indistinguishable from finding nothing; the viewer combo still named the retired
ice_ring_score; and the claim that growth "cannot invent a crystal" was too strong -
it cannot start a patch, but the cell count is read over the grown patch, so it does
decide which patches pass.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EFEJG6WBQv8th4UJFNe53N
2026-09-08 09:04:52 +02:00

261 lines
12 KiB
C++

// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include "JFJochReceiver.h"
#include "../common/CUDAWrapper.h"
#include "../common/time_utc.h"
#include "JFJochCompressor.h"
#include "../image_analysis/grid_scan_analysis/AnalyzeGridScan.h"
JFJochReceiver::JFJochReceiver(const DiffractionExperiment &in_experiment,
ImageBuffer &in_image_buffer,
ImagePusher &in_image_pusher,
JFJochReceiverCurrentStatus &in_current_status,
JFJochReceiverPlots &in_plots,
const SpotFindingSettings &spot_finding_settings,
Logger &logger,
const PixelMask &in_pixel_mask,
ZMQPreviewSocket *in_zmq_preview_socket,
ZMQMetadataSocket *in_zmq_metadata_socket,
IndexerThreadPool *indexing_thread_pool)
: logger(logger),
experiment(in_experiment),
spot_finding_settings(spot_finding_settings),
image_buffer(in_image_buffer),
image_pusher(in_image_pusher),
zmq_preview_socket(in_zmq_preview_socket),
zmq_metadata_socket(in_zmq_metadata_socket),
current_status(in_current_status),
plots(in_plots),
scan_result(in_experiment),
serialmx_filter(in_experiment),
pixel_mask(in_pixel_mask),
// retain_outcomes=false for now: online scaling/merge at the end is not enabled yet, so the
// whole-run integration_outcome is unused here. Flip back to true when online scaling lands.
indexer(experiment, indexing_thread_pool, false, /*real_time=*/true) {
logger.Info("Initializing receiver");
// Ensure there is nothing running for now
if (!image_buffer.Finalize(std::chrono::seconds(1)))
throw JFJochException(JFJochExceptionCategory::WrongDAQState,
"There are unfinished preview/sending jobs in the buffer");
logger.Info("Image buffer from previous run finalized");
current_status.SetProgress(0);
current_status.SetEfficiency({});
current_status.SetStatus(JFJochReceiverStatus{}); // GetStatus() is virtual function and cannot be called yet!
auto start_time_point = std::chrono::steady_clock::now();
az_int_mapping = std::make_unique<AzimuthalIntegrationMapping>(experiment, pixel_mask);
auto end_time_point = std::chrono::steady_clock::now();
auto duration = std::chrono::duration<float>(end_time_point - start_time_point);
logger.Info("Azimuthal integration mapping done in {:.5f} s with {} threads", duration.count(), az_int_mapping->GetNThreads());
plots.Setup(experiment, *az_int_mapping);
push_images_to_writer = (experiment.GetImageNum() > 0) && (!experiment.GetFilePrefix().empty());
}
JFJochReceiver::~JFJochReceiver() = default;
SpotFindingSettings JFJochReceiver::GetSpotFindingSettings() {
std::unique_lock ul(spot_finding_settings_mutex);
return spot_finding_settings;
}
void JFJochReceiver::UpdateMaxImageSent(int64_t image_number) {
std::unique_lock ul(max_image_number_sent_mutex);
if (image_number + 1 > max_image_number_sent)
max_image_number_sent = image_number + 1;
}
void JFJochReceiver::UpdateMaxImageReceived(int64_t image_number) {
std::unique_lock ul(max_image_number_received_mutex);
if (image_number + 1 > max_image_number_received)
max_image_number_received = image_number + 1;
}
void JFJochReceiver::UpdateMaxDelay(uint64_t delay) {
std::unique_lock ul(max_delay_mutex);
if (!max_delay || (delay > max_delay))
max_delay = delay;
}
JFJochReceiverStatus JFJochReceiver::GetStatus() const {
JFJochReceiverStatus ret;
ret.indexing_rate = plots.GetIndexingRate();
ret.bkg_estimate = plots.GetBkgEstimate();
if ((experiment.GetImageNum() > 0) && (compressed_size > 0)) {
ret.compressed_ratio = static_cast<double>(uncompressed_size) / static_cast<double>(compressed_size);
}
ret.saturated_pixels = saturated_pixels.Read();
ret.error_pixels = error_pixels.Read();
ret.roi_beam_npixel = roi_beam_npixel.Read();
ret.roi_beam_sum = roi_beam_sum.Read();
ret.compressed_size = compressed_size;
ret.max_receive_delay = max_delay;
ret.max_image_number_sent = max_image_number_sent;
ret.images_collected = images_collected;
ret.images_sent = images_sent;
ret.images_skipped = images_skipped;
ret.images_written = image_pusher.GetImagesWritten();
ret.cancelled = cancelled;
ret.efficiency = GetEfficiency();
return ret;
}
void JFJochReceiver::SendStartMessage() {
StartMessage message{};
experiment.FillMessage(message);
message.arm_date = time_UTC(std::chrono::system_clock::now());
message.az_int_q_bin_count = az_int_mapping->GetQBinCount();
message.az_int_bin_to_q = az_int_mapping->GetBinToQ();
message.az_int_bin_to_two_theta = az_int_mapping->GetBinToTwoTheta();
message.az_int_phi_bin_count = az_int_mapping->GetAzimuthalBinCount();
if (az_int_mapping->GetAzimuthalBinCount() > 1) {
message.az_int_bin_to_phi = az_int_mapping->GetBinToPhi();
message.az_int_map = az_int_mapping->GetPixelToBin();
}
message.writer_notification_zmq_addr = image_pusher.GetWriterNotificationSocketAddress();
message.rois = experiment.ROI().ExportMetadata();
if (!experiment.ROI().empty())
message.roi_map = experiment.ExportROIMap();
message.max_spot_count = experiment.GetMaxSpotCount();
std::vector<uint32_t> nexus_mask;
message.pixel_mask["default"] = pixel_mask.GetMask(experiment);
SaveStartMessageToImageBuffer(message);
if (push_images_to_writer)
image_pusher.StartDataCollection(message);
if (zmq_preview_socket != nullptr)
zmq_preview_socket->StartDataCollection(message);
if (zmq_metadata_socket != nullptr)
zmq_metadata_socket->StartDataCollection(message);
}
void JFJochReceiver::SaveStartMessageToImageBuffer(const StartMessage &msg) {
std::vector<uint8_t> buffer(MESSAGE_SIZE_FOR_START_END);
CBORStream2Serializer serializer(buffer.data(), buffer.size());
serializer.SerializeSequenceStart(msg);
buffer.resize(serializer.GetBufferSize());
image_buffer.SaveStartMessage(buffer);
}
void JFJochReceiver::SendEndMessage() {
EndMessage message{};
message.max_image_number = max_image_number_sent;
message.images_collected_count = images_collected;
message.images_sent_to_write_count = images_sent;
message.max_receiver_delay = max_delay;
message.efficiency = GetEfficiency();
message.end_date = time_UTC(std::chrono::system_clock::now());
message.run_number = experiment.GetRunNumber();
message.run_name = experiment.GetRunName();
message.bkg_estimate = plots.GetBkgEstimate();
message.spindle_blind_fraction = plots.GetSpindleBlindFraction();
message.ice_ring_ratio_mean = plots.GetIceRingRatio();
message.protein_score = plots.GetProteinScore();
message.ice_score = plots.GetIceScore();
message.indexing_rate = plots.GetIndexingRate();
message.az_int_result["dataset"] = plots.GetAzIntProfile();
const auto rotation_indexer_ret = indexer.FinalizeRotationIndexing();
if (rotation_indexer_ret.has_value()) {
message.rotation_lattice = rotation_indexer_ret->lattice;
message.rotation_lattice_type = LatticeMessage{
.centering = rotation_indexer_ret->search_result.centering,
.niggli_class = rotation_indexer_ret->search_result.niggli_class,
.crystal_system = rotation_indexer_ret->search_result.system
};
message.rotation_extra_lattices = rotation_indexer_ret->extra_lattices;
rotation_indexing_lattice = rotation_indexer_ret->lattice;
rotation_indexing_lattice_type = message.rotation_lattice_type;
}
message.unit_cell = indexer.GetConsensusUnitCell();
for (int i = 0; i < adu_histogram_module.size(); i++)
message.adu_histogram["module" + std::to_string(i)] = adu_histogram_module[i]->GetHistogram();
scan_result.FillEndMessage(message);
// A grid scan's crystals are found once, from the completed map. The decision is about the whole
// raster and cannot be taken while it runs: a patch that looks like the best crystal in row three
// is often the shoulder of a better one two rows further down, and a scan exists precisely to see
// the whole loop before choosing. Done here so one list reaches both the written file (through
// this end message) and the API (through GetFinalStatistics).
if (experiment.GetAnalysisMode() == AnalysisMode::Grid && !experiment.GetGridScan().has_value())
logger.Warning("Analysis mode is Grid but the dataset has no grid scan: images were scored, "
"but there is no map to find crystals in and the crystal list is empty");
if (const auto grid = experiment.GetGridScan();
grid.has_value() && experiment.GetAnalysisMode() == AnalysisMode::Grid) {
grid_scan_result = AnalyzeGridScan(scan_result.GetResult(), grid.value(),
experiment.GetBeamSizeX_um().value_or(0.0f),
experiment.GetBeamSizeY_um().value_or(0.0f),
experiment.GetGridScanAnalysisSettings());
message.grid_crystals = grid_scan_result->crystals;
logger.Info("Grid scan: {} crystal(s) found in {} grid points",
grid_scan_result->crystals.size(), message.max_image_number);
}
if (push_images_to_writer) {
if (!image_pusher.EndDataCollection(message))
logger.Error("End message not sent via ZeroMQ (time-out)");
logger.Info("Disconnected from writers");
}
if (zmq_metadata_socket != nullptr)
zmq_metadata_socket->EndDataCollection(message);
if (zmq_preview_socket != nullptr)
zmq_preview_socket->EndDataCollection(message);
}
JFJochReceiverOutput JFJochReceiver::GetFinalStatistics() const {
JFJochReceiverOutput ret;
ret.efficiency = GetEfficiency();
ret.start_time_ms = std::chrono::duration_cast<std::chrono::milliseconds>(start_time.time_since_epoch()).count();
ret.end_time_ms = std::chrono::duration_cast<std::chrono::milliseconds>(end_time.time_since_epoch()).count();
ret.writer_queue_full_warning = writer_queue_full;
ret.status = GetStatus();
ret.writer_err = writer_error;
ret.scan_result = scan_result.GetResult();
ret.scan_result.grid = grid_scan_result;
ret.scan_result.rotation_lattice = rotation_indexing_lattice;
if (rotation_indexing_lattice_type) {
ret.scan_result.rotation_crystal_system = rotation_indexing_lattice_type->crystal_system;
ret.scan_result.rotation_centering = rotation_indexing_lattice_type->centering;
}
ret.images_written = images_written;
ret.processing_time = plots.GetMeanProcessingTime();
return ret;
}
void JFJochReceiver::Cancel(bool silent) {
if (!silent) {
// Remote abort: This tells FPGAs to stop but doesn't do anything to CPU code
logger.Warning("Cancelling on request");
cancelled = true;
}
}
void JFJochReceiver::Cancel(const JFJochException &e) {
logger.Error("Cancelling data collection due to exception");
logger.ErrorException(e);
// Error abort: This tells FPGAs to stop and also prevents deadlock in CPU code by setting abort to 1
cancelled = true;
}