Files
Jungfraujoch/receiver/JFJochReceiverFPGA.cpp
T
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

714 lines
33 KiB
C++

// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include "JFJochReceiverFPGA.h"
#include <thread>
#include "ImageMetadata.h"
#include "../image_analysis/IceScore.h"
JFJochReceiverFPGA::JFJochReceiverFPGA(const DiffractionExperiment &in_experiment,
const PixelMask &in_pixel_mask,
const JFCalibration *in_calibration,
AcquisitionDeviceGroup &in_aq_device,
ImagePusher &in_image_sender,
Logger &in_logger, int64_t in_forward_and_sum_nthreads,
const SpotFindingSettings &in_spot_finding_settings,
JFJochReceiverCurrentStatus &in_current_status,
JFJochReceiverPlots &in_plots,
ImageBuffer &in_send_buf_ctrl,
ZMQPreviewSocket *in_zmq_preview_socket,
ZMQMetadataSocket *in_zmq_metadata_socket,
IndexerThreadPool *indexing_thread_pool)
: JFJochReceiver(in_experiment,
in_send_buf_ctrl,
in_image_sender,
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),
calibration(nullptr),
acquisition_device(in_aq_device),
ndatastreams(experiment.GetDataStreamsNum()),
frame_transformation_nthreads(
in_forward_and_sum_nthreads),
pedestal_nthreads((experiment.GetStorageCellNumber() > 2) ? 1 : 4),
summation_nthreads(2),
frame_transformation_ready((experiment.GetImageNum() > 0) ? frame_transformation_nthreads : 0),
data_acquisition_ready(ndatastreams) {
// The FPGA integration core has a fixed bin-count limit. Reject configurations it cannot
// handle before starting acquisition, unless azimuthal integration is forced onto the CPU.
const auto &azint_settings = experiment.GetAzimuthalIntegrationSettings();
if (!azint_settings.IsForceCPUinFPGAWorkflow() && azint_settings.GetBinCount() > FPGA_INTEGRATION_BIN_COUNT)
throw JFJochException(JFJochExceptionCategory::InputParameterAboveMax,
fmt::format("Azimuthal integration bin count ({}) exceeds the FPGA limit ({}). "
"Enable ForceCPUinFPGAWorkflow to compute it on the CPU.",
azint_settings.GetBinCount(), FPGA_INTEGRATION_BIN_COUNT));
image_buffer.StartMeasurement(experiment);
for (int m = 0; m < experiment.GetModulesNum(); m++)
adu_histogram_module.emplace_back(std::make_unique<ADUHistogram>());
roi_map = experiment.ExportROIMap();
if (experiment.GetDetectorSetup().GetDetectorType() == DetectorType::JUNGFRAU)
calibration = in_calibration;
if (acquisition_device.size() < ndatastreams)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Number of acquisition devices has to match data streams");
if (frame_transformation_nthreads <= 0)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Number of threads must be more than zero");
expected_packets_per_image = 0;
for (int d = 0; d < ndatastreams; d++) {
acquisition_device[d].PrepareAction(experiment);
acquisition_device[d].SetSpotFinderParameters(spot_finding_settings);
expected_packets_per_image += acquisition_device[d].Counters().GetExpectedPacketsPerImage();
logger.Debug("Acquisition device {} prepared", d);
}
if (experiment.IsCPUSummation())
expected_packets_per_image *= experiment.GetSummation();
logger.Info("Data acquisition devices ready");
if ((experiment.GetDetectorMode() == DetectorMode::PedestalG0)
|| (experiment.GetDetectorMode() == DetectorMode::PedestalG1)
|| (experiment.GetDetectorMode() == DetectorMode::PedestalG2)) {
if (experiment.GetImageNum() > 0) {
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Saving and calculating pedestal is not supported for the time being");
}
switch (experiment.GetDetectorMode()) {
case DetectorMode::PedestalG1:
only_2nd_sc_pedestal = !experiment.IsFixedGainG1() && (experiment.GetStorageCellNumber() == 2);
break;
case DetectorMode::PedestalG2:
only_2nd_sc_pedestal = (experiment.GetStorageCellNumber() == 2);
break;
default:
only_2nd_sc_pedestal = false;
break;
}
int64_t pedestal_count = (only_2nd_sc_pedestal)
? experiment.GetModulesNum()
: experiment.GetModulesNum() * experiment.GetStorageCellNumber();
for (int i = 0; i < pedestal_count; i++)
pedestal.emplace_back(std::make_unique<JFPedestalCalc>(experiment));
for (int s = 0; s < experiment.GetStorageCellNumber(); s++) {
bool ignore = only_2nd_sc_pedestal && (s == 0);
for (int d = 0; d < ndatastreams; d++)
for (int m = 0; m < experiment.GetModulesNum(d); m++)
for (int n = 0; n < pedestal_nthreads; n++)
frame_transformation_futures.emplace_back(std::async(std::launch::async,
&JFJochReceiverFPGA::MeasurePedestalThread,
this, d, m, s, n, ignore));
}
logger.Info("Pedestal threads ready ({} threads/module*SC)", pedestal_nthreads);
} else if (experiment.GetImageNum() > 0) {
SendStartMessage();
SendCalibration();
for (int i = 0; i < experiment.GetImageNum(); i++)
images_to_go.Put(i);
// Setup frames summation and forwarding
for (uint32_t i = 0; i < frame_transformation_nthreads; i++) {
auto handle = std::async(std::launch::async, &JFJochReceiverFPGA::FrameTransformationThread,
this, i);
frame_transformation_futures.emplace_back(std::move(handle));
}
logger.Info("Image compression/forwarding threads started ({} threads)", frame_transformation_nthreads);
frame_transformation_ready.wait();
logger.Info("Image compression/forwarding threads ready");
}
for (int d = 0; d < ndatastreams; d++)
data_acquisition_futures.emplace_back(std::async(std::launch::async, &JFJochReceiverFPGA::AcquireThread,
this, d));
data_acquisition_ready.wait();
start_time = std::chrono::system_clock::now();
measurement = std::async(std::launch::async, &JFJochReceiverFPGA::FinalizeMeasurement, this);
logger.Info("Receiving data started");
}
void JFJochReceiverFPGA::SendPedestal(const std::string &prefix, const std::vector<uint8_t> &v, int gain, int sc) {
size_t xpixel = RAW_MODULE_COLS;
size_t ypixel = experiment.GetModulesNum() * RAW_MODULE_LINES;
std::string channel;
if (experiment.GetStorageCellNumber() > 1)
channel = fmt::format("{:s}_g{:d}_sc{:d}", prefix, gain, sc);
else
channel = fmt::format("{:s}_g{:d}", prefix, gain);
CompressedImage image(v.data(), v.size(), xpixel, ypixel, CompressedImageMode::Uint16,
CompressionAlgorithm::BSHUF_LZ4, channel);
image_pusher.SendCalibration(image);
}
void JFJochReceiverFPGA::SendCalibration() {
if ((calibration == nullptr) || !experiment.GetSaveCalibration() || !push_images_to_writer)
return;
JFJochBitShuffleCompressor compressor(CompressionAlgorithm::BSHUF_LZ4);
for (int sc = 0; sc < experiment.GetStorageCellNumber(); sc++) {
for (int gain = 0; gain < 3; gain++) {
if (experiment.IsFixedGainG1() && (gain != 1))
continue;
SendPedestal("pedestal", compressor.Compress(calibration->GetPedestal(gain, sc)), gain, sc);
SendPedestal("pedestal_rms", compressor.Compress(calibration->GetPedestalRMS(gain, sc)), gain, sc);
}
}
}
void JFJochReceiverFPGA::AcquireThread(uint16_t data_stream) {
try {
LoadCalibrationToFPGA(data_stream);
frame_transformation_ready.wait();
logger.Debug("Device thread {} start FPGA action", data_stream);
acquisition_device[data_stream].StartAction(experiment);
} catch (const JFJochException &e) {
Cancel(e);
data_acquisition_ready.count_down();
logger.ErrorException(e);
logger.Warning("Device thread {} done due to an error", data_stream);
return;
}
data_acquisition_ready.count_down();
try {
logger.Debug("Device thread {} wait for FPGA action complete", data_stream);
acquisition_device[data_stream].WaitForActionComplete();
} catch (const JFJochException &e) {
logger.ErrorException(e);
Cancel(e);
logger.ErrorException(e);
logger.Warning("Device thread {} done due to an error", data_stream);
return;
}
logger.Info("Device thread {} done", data_stream);
}
void JFJochReceiverFPGA::MeasurePedestalThread(uint16_t data_stream, uint16_t module_number, uint16_t storage_cell,
uint32_t threadid, bool ignore) {
JFPedestalCalc pedestal_calc(experiment);
uint64_t starting_frame = storage_cell + threadid * experiment.GetStorageCellNumber();
uint64_t frame_stride = experiment.GetStorageCellNumber() * pedestal_nthreads;
uint32_t storage_cell_header = UINT32_MAX;
try {
for (size_t frame = starting_frame; frame < experiment.GetFrameNum(); frame += frame_stride) {
// Frame will be processed only if one already collects frame+2
acquisition_device[data_stream].Counters().WaitForFrame(frame + 2, module_number);
if (acquisition_device[data_stream].Counters().IsFullModuleCollected(frame, module_number) && !ignore) {
auto output = acquisition_device[data_stream].GetDeviceOutput(frame, module_number);
// Partial packets will bring more problems, than benefit
pedestal_calc.AnalyzeImage((uint16_t *) output->pixels);
storage_cell_header = (output->module_statistics.debug >> 8) & 0xF;
}
acquisition_device[data_stream].FrameBufferRelease(frame, module_number);
UpdateMaxDelay(acquisition_device[data_stream].Counters().CalculateDelay(frame, module_number));
current_status.SetProgress(GetProgress());
current_status.SetStatus(GetStatus());
}
uint64_t offset = experiment.GetFirstModuleOfDataStream(data_stream) + module_number;
if (!only_2nd_sc_pedestal)
offset += experiment.GetModulesNum() * storage_cell;
if (!ignore)
*pedestal[offset] += pedestal_calc;
} catch (const JFJochException &e) {
Cancel(e);
}
logger.Debug("Pedestal calculation thread for data stream {} module {} storage cell {} -> header {} done",
data_stream, module_number, storage_cell, storage_cell_header);
}
int64_t JFJochReceiverFPGA::SummationThread(uint16_t data_stream,
int64_t image_number,
uint16_t module_number,
uint32_t threadid,
ModuleSummation &summation) {
ModuleSummation local_summation(experiment);
int64_t starting_frame = image_number * experiment.GetSummation();
for (int64_t i = threadid; i < experiment.GetSummation(); i += summation_nthreads) {
const int64_t frame = starting_frame + i;
// Frame will be processed only if one already collects frame+2
acquisition_device[data_stream].Counters().WaitForFrame(frame + 2, module_number);
if (acquisition_device[data_stream].Counters().IsAnyPacketCollected(frame, module_number)) {
const auto output = acquisition_device[data_stream].GetDeviceOutput(frame, module_number);
local_summation.AddFPGAOutput(*output);
} else
local_summation.AddEmptyOutput();
acquisition_device[data_stream].FrameBufferRelease(frame, module_number);
UpdateMaxDelay(acquisition_device[data_stream].Counters().CalculateDelay(frame, module_number));
current_status.SetProgress(GetProgress());
current_status.SetStatus(GetStatus());
}
if (!summation.empty())
summation.AddFPGAOutput(local_summation.GetOutput(), 4);
return 0;
}
void JFJochReceiverFPGA::FrameTransformationThread(uint32_t threadid) {
std::unique_ptr<MXAnalysisAfterFPGA> analyzer;
try {
analyzer = std::make_unique<MXAnalysisAfterFPGA>(experiment, *az_int_mapping, indexer);
} catch (const JFJochException &e) {
frame_transformation_ready.count_down();
logger.Error("Thread setup error {}", e.what());
Cancel(e);
return;
}
FrameTransformation transformation(experiment);
frame_transformation_ready.count_down();
uint16_t az_int_min_bin = std::floor(az_int_mapping->QToBin(experiment.GetLowQForBkgEstimate_recipA()));
uint16_t az_int_max_bin = std::ceil(az_int_mapping->QToBin(experiment.GetHighQForBkgEstimate_recipA()));
// When forced, azimuthal integration is computed on the CPU from the assembled image
// instead of being read back from the FPGA per module (lifts the FPGA bin-count limit).
const bool force_cpu_azint = experiment.GetAzimuthalIntegrationSettings().IsForceCPUinFPGAWorkflow();
int64_t image_number;
while (images_to_go.Get(image_number) != 0) {
try {
int64_t expected_frame = image_number;
if (experiment.IsCPUSummation())
expected_frame *= experiment.GetSummation();
logger.Debug("Frame transformation thread - trying to get image {}", expected_frame);
// If data acquisition is finished and fastest frame for the first device is behind
acquisition_device[0].Counters().WaitForFrame(expected_frame);
logger.Debug("Frame transformation thread - frame arrived {}", expected_frame);
if (acquisition_device[0].Counters().IsAcquisitionFinished() &&
(acquisition_device[0].Counters().GetFastestFrameNumber() < expected_frame)) {
logger.Debug("Frame transformation thread - skipping image {}", expected_frame);
continue;
}
DataMessage message{};
message.number = image_number;
message.original_number = image_number;
message.user_data = experiment.GetImageAppendix();
message.run_number = experiment.GetRunNumber();
message.run_name = experiment.GetRunName();
ImageMetadata metadata(experiment);
AzimuthalIntegrationProfile az_int_profile_image(*az_int_mapping);
auto local_spot_finding_settings = GetSpotFindingSettings();
const auto preprocessing_start_time = std::chrono::steady_clock::now();
if (experiment.IsCPUSummation()) {
std::vector<std::unique_ptr<ModuleSummation>> summation;
for (int i = 0; i < experiment.GetModulesNum(); i++)
summation.emplace_back(std::make_unique<ModuleSummation>(experiment));
std::vector<std::future<int64_t> > futures;
for (int d = 0; d < ndatastreams; d++) {
for (int m = 0; m < experiment.GetModulesNum(d); m++) {
size_t module_abs_number = experiment.GetFirstModuleOfDataStream(d) + m;
for (int i = 0; i < summation_nthreads; i++) {
futures.emplace_back(
std::async(std::launch::async,
&JFJochReceiverFPGA::SummationThread,
this,
d, image_number, m, i,
std::ref(*summation[module_abs_number]))
);
}
}
}
for (auto &f: futures)
f.get();
for (int d = 0; d < ndatastreams; d++) {
for (int m = 0; m < experiment.GetModulesNum(d); m++) {
size_t i = experiment.GetFirstModuleOfDataStream(d) + m;
if (!summation[i]->empty()) {
adu_histogram_module[i]->Add(summation[i]->GetOutput());
transformation.ProcessModule(&summation[i]->GetOutput(), d);
metadata.Process(&summation[i]->GetOutput());
if (!force_cpu_azint)
az_int_profile_image.Add(summation[i]->GetOutput());
analyzer->ReadFromCPU(&summation[i]->GetOutput(), GetSpotFindingSettings(), i);
} else
transformation.FillNotCollectedModule(m, d);
}
}
} else {
logger.Debug("Frame transformation thread - processing image from FPGA {}", image_number);
for (int d = 0; d < ndatastreams; d++) {
for (int m = 0; m < experiment.GetModulesNum(d); m++) {
acquisition_device[d].Counters().WaitForFrame(image_number + 2, m);
if (acquisition_device[d].Counters().IsAnyPacketCollected(image_number, m)) {
const DeviceOutput *output = acquisition_device[d].GetDeviceOutput(image_number, m);
metadata.Process(output);
size_t module_abs_number = experiment.GetFirstModuleOfDataStream(d) + m;
adu_histogram_module[module_abs_number]->Add(*output);
if (!force_cpu_azint)
az_int_profile_image.Add(*output);
analyzer->ReadFromFPGA(output, local_spot_finding_settings, module_abs_number);
transformation.ProcessModule(output, d);
} else
transformation.FillNotCollectedModule(m, d);
acquisition_device[d].FrameBufferRelease(image_number, m);
}
auto delay = acquisition_device[d].Counters().CalculateDelay(image_number);
UpdateMaxDelay(delay);
if (delay > message.receiver_aq_dev_delay)
message.receiver_aq_dev_delay = delay;
}
}
const auto preprocessing_end_time = std::chrono::steady_clock::now();
message.preprocessing_time_s = std::chrono::duration<float>(preprocessing_end_time - preprocessing_start_time).count();
auto image_start_time = std::chrono::high_resolution_clock::now();
metadata.Export(message, expected_packets_per_image);
if (message.image_collection_efficiency == 0.0f) {
plots.AddEmptyImage(message);
continue;
}
message.image = CompressedImage(transformation.GetImage(),
experiment.GetPixelsNum() * experiment.GetByteDepthImage(),
experiment.GetXPixelsNum(),
experiment.GetYPixelsNum(),
experiment.GetImageMode(),
CompressionAlgorithm::NO_COMPRESSION);
// No-op unless the CPU backend is forced; fills az_int_profile_image from the assembled image
analyzer->RunAzimuthalIntegration(transformation.GetImage(), az_int_profile_image);
analyzer->Process(message, local_spot_finding_settings);
auto status = image_buffer.GetStatus();
message.receiver_buf_available = status.available_slots;
message.receiver_buf_in_preparation = status.preparation_slots;
message.receiver_buf_in_sending = status.sending_slots;
message.az_int_profile = az_int_profile_image.GetResult();
message.az_int_profile_count = az_int_profile_image.GetPixelCount();
if (force_cpu_azint)
message.az_int_profile_std = az_int_profile_image.GetStd();
message.bkg_estimate = az_int_profile_image.GetBkgEstimate(experiment.GetAzimuthalIntegrationSettings());
message.ice_ring_ratio = az_int_profile_image.GetIceRingRatio(
experiment.GetAzimuthalIntegrationSettings(), spot_finding_settings.ice_ring_width_Q_recipA);
// The radial channel of the ice score needs the profile's standard deviation, which the
// FPGA azimuthal integration does not produce - it abstains there and only the spot
// channel contributes. ForceCPUinFPGAWorkflow gives it the standard deviation back.
message.ice_score = IceScore(message.az_int_profile, message.az_int_profile_std,
message.az_int_profile_count,
experiment.GetAzimuthalIntegrationSettings().GetQBinCount(),
experiment.GetAzimuthalIntegrationSettings(), message.spot_d_A_unfiltered,
spot_finding_settings.ice_ring_width_Q_recipA);
scan_result.Add(message);
auto image_end_time = std::chrono::high_resolution_clock::now();
std::chrono::duration<float> duration = image_end_time - image_start_time;
message.processing_time_s = duration.count();
// Store overload/error pixel count
if (message.image_collection_efficiency == 1.0f) {
saturated_pixels.Add(message.saturated_pixel_count);
error_pixels.Add(message.error_pixel_count);
if (message.roi.contains("beam")) {
roi_beam_npixel.Add(message.roi["beam"].pixels);
roi_beam_sum.Add(message.roi["beam"].sum);
}
}
++images_collected;
uncompressed_size += experiment.GetModulesNum() * RAW_MODULE_SIZE * experiment.GetByteDepthImage();
if (!serialmx_filter.ApplyFilter(message))
++images_skipped;
else {
auto loc = image_buffer.GetImageSlot();
if (loc == nullptr) {
// No free buffer locations - continue
writer_queue_full = true;
} else {
auto writer_buffer = (uint8_t *) loc->GetImage();
CBORStream2Serializer serializer(writer_buffer, experiment.GetImageBufferLocationSize());
message.image = CompressedImage(nullptr, 0,
experiment.GetXPixelsNum(),
experiment.GetYPixelsNum(),
experiment.GetImageMode(),
experiment.GetCompressionAlgorithm());
serializer.SerializeImage(message);
const size_t metadata_size = serializer.GetImageAppendOffset();
const size_t slot_size = experiment.GetImageBufferLocationSize();
const auto compression_start_time = std::chrono::steady_clock::now();
try {
// Reserve 32 bytes for close, etc. If the metadata already fills the
// slot the subtraction below would wrap (size_t), so bail out here and
// let the catch drop just this frame.
if (metadata_size + 32 >= slot_size)
throw CompressionBufferTooSmallException("CBOR metadata leaves no room for the compressed image");
size_t image_size = transformation.CompressImage(
writer_buffer + serializer.GetImageAppendOffset(),
slot_size - (metadata_size + 32));
const auto compression_end_time = std::chrono::steady_clock::now();
message.compression_time_s = std::chrono::duration<float>(
compression_end_time - compression_start_time).count();
serializer.AppendImage(image_size);
compressed_size += image_size;
if (image_size > 0)
message.compression_ratio =
static_cast<float>(experiment.GetPixelsNum() * experiment.GetByteDepthImage())
/ static_cast<float>(image_size);
loc->SetImageNumber(image_number);
loc->SetImageSize(serializer.GetBufferSize());
loc->SetIndexed(message.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(message);
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(message.number);
} catch (const CompressionBufferTooSmallException &) {
// The per-image CBOR metadata (spot/reflection lists, ...) leaves
// too little room for the compressed image in the buffer slot. Drop
// just this frame instead of aborting the whole data collection.
logger.Error("Dropping image {} - CBOR metadata too large to store image: "
"metadata {} do not fit in buffer slot {} B",
message.number, metadata_size, slot_size);
loc->release();
}
}
}
logger.Debug("Frame transformation thread - done sending image {} / {}", image_number, message.number);
plots.Add(message, az_int_profile_image);
current_status.SetProgress(GetProgress());
current_status.SetStatus(GetStatus());
} catch (const JFJochException &e) {
logger.ErrorException(e);
Cancel(e);
}
}
logger.Debug("Sum&compression thread done");
}
float JFJochReceiverFPGA::GetEfficiency() const {
uint64_t expected_packets;
if (experiment.GetImageNum() == 0)
expected_packets = expected_packets_per_image * experiment.GetFrameNum();
else
expected_packets = expected_packets_per_image * experiment.GetImageNum();
uint64_t received_packets = 0;
for (int d = 0; d < ndatastreams; d++) {
received_packets += acquisition_device[d].Counters().GetTotalPackets();
}
if ((expected_packets == received_packets) || (expected_packets == 0))
return 1.0;
return received_packets / static_cast<double>(expected_packets);
}
void JFJochReceiverFPGA::Cancel(bool silent) {
JFJochReceiver::Cancel(silent);
for (int d = 0; d < ndatastreams; d++)
acquisition_device[d].Cancel();
}
void JFJochReceiverFPGA::Cancel(const JFJochException &e) {
JFJochReceiver::Cancel(e);
for (int d = 0; d < ndatastreams; d++)
acquisition_device[d].Cancel();
}
float JFJochReceiverFPGA::GetProgress() const {
int64_t frames = experiment.GetImageNum();
if (experiment.IsCPUSummation())
frames *= experiment.GetSummation();
if (frames == 0)
frames = experiment.GetFrameNum();
if ((frames == 0) || (acquisition_device[0].Counters().IsAcquisitionFinished()))
return 1.0;
return static_cast<float>(acquisition_device[0].Counters().GetSlowestFrameNumber()) / static_cast<float>(frames);
}
void JFJochReceiverFPGA::FinalizeMeasurement() {
if (!frame_transformation_futures.empty()) {
for (auto &future: frame_transformation_futures)
future.get();
logger.Info("All processing threads done");
}
current_status.SetProgress(1.0);
current_status.SetStatus(GetStatus());
SendEndMessage();
if (experiment.GetImageNum() > 0) {
for (int d = 0; d < ndatastreams; d++)
acquisition_device[d].Cancel();
}
end_time = std::chrono::system_clock::now();
for (auto &future: data_acquisition_futures)
future.get();
logger.Info("Devices stopped");
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("Receiving data done");
if (push_images_to_writer)
writer_error = image_pusher.Finalize();
current_status.SetProgress({});
current_status.SetStatus(GetStatus());
logger.Info("Writing process finalized");
}
void JFJochReceiverFPGA::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;
for (int i = 0; i < ndatastreams; i++)
acquisition_device[i].SetSpotFinderParameters(spot_finding_settings);
}
void JFJochReceiverFPGA::StopReceiver() {
if (measurement.valid()) {
measurement.get();
logger.Info("Receiver stopped");
}
}
JFJochReceiverFPGA::~JFJochReceiverFPGA() {
try {
if (measurement.valid())
measurement.get();
} catch (const std::exception& e) {
logger.Error("ReceiverFPGA teardown failed: {}", e.what());
} catch (...) {
logger.Error("ReceiverFPGA teardown failed with unknown exception");
}
}
JFJochReceiverOutput JFJochReceiverFPGA::GetFinalStatistics() const {
JFJochReceiverOutput ret = JFJochReceiver::GetFinalStatistics();
for (int d = 0; d < ndatastreams; d++) {
for (int m = 0; m < acquisition_device[d].Counters().GetModuleNumber(); m++) {
if (experiment.IsCPUSummation())
ret.expected_packets.push_back(acquisition_device[d].Counters().GetTotalExpectedPacketsPerModule() * experiment.GetSummation());
else
ret.expected_packets.push_back(acquisition_device[d].Counters().GetTotalExpectedPacketsPerModule());
ret.received_packets.push_back(acquisition_device[d].Counters().GetTotalPackets(m));
}
}
RetrievePedestal(ret.pedestal_result);
return ret;
}
void JFJochReceiverFPGA::RetrievePedestal(std::vector<JFModulePedestal> &output) const {
time_t curr_time = std::chrono::system_clock::to_time_t(start_time);
for (const auto &pc: pedestal) {
JFModulePedestal mp;
if (experiment.GetDetectorMode() == DetectorMode::PedestalG0)
pc->Export(mp, PEDESTAL_G0_WRONG_GAIN_ALLOWED_COUNT);
else
pc->Export(mp);
mp.SetCollectionTime(curr_time);
output.emplace_back(std::move(mp));
}
}
void JFJochReceiverFPGA::LoadCalibrationToFPGA(uint16_t data_stream) {
if (experiment.IsPedestalRun()) {
acquisition_device[data_stream].InitializeEmptyPixelMask(experiment);
return; // No calibration loaded for pedestal
}
if (calibration != nullptr)
acquisition_device[data_stream].InitializeCalibration(experiment, *calibration);
// Initialize pixel_mask
acquisition_device[data_stream].InitializePixelMask(experiment, pixel_mask);
// Initialize roi_map
acquisition_device[data_stream].InitializeROIMap(experiment, roi_map);
// Initialize data processing
acquisition_device[data_stream].InitializeDataProcessing(experiment, *az_int_mapping);
}