JFJochReceiver: Skip frames if acquisition finished and frames stopped earlier on the first acquisition device

This commit is contained in:
2023-10-22 14:36:53 +02:00
parent bc43921004
commit c1469d1e46
6 changed files with 65 additions and 42 deletions
+28 -18
View File
@@ -6,13 +6,15 @@
#include "../common/JFJochException.h"
AcquisitionCounters::AcquisitionCounters()
: head(max_modules, 0), slowest_head(0), total_packets(0), expected_frames(0), acquisition_finished(false) {}
: curr_frame_number(max_modules, 0), slowest_frame_number(0), fastest_frame_number(0),
total_packets(0), expected_frames(0), acquisition_finished(false) {}
void AcquisitionCounters::Reset(const DiffractionExperiment &experiment, uint16_t data_stream) {
std::unique_lock<std::shared_mutex> ul(m);
acquisition_finished = false;
slowest_head = 0;
slowest_frame_number = 0;
fastest_frame_number = 0;
if ((experiment.GetDetectorMode() == DetectorMode::PedestalG0) ||
(experiment.GetDetectorMode() == DetectorMode::PedestalG1) ||
@@ -25,7 +27,7 @@ void AcquisitionCounters::Reset(const DiffractionExperiment &experiment, uint16_
if (nmodules > max_modules)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "Acquisition counter cannot support that many modules");
for (int i = 0; i < max_modules; i++)
head[i] = 0;
curr_frame_number[i] = 0;
handle_for_frame = std::vector<uint64_t>((expected_frames+1) * nmodules, HandleNotFound);
packets_collected = std::vector<uint16_t>(expected_frames * nmodules);
@@ -44,14 +46,18 @@ void AcquisitionCounters::UpdateCounters(const Completion *c) {
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"UpdateCounters frame number is out of bounds");
else {
if (head.at(c->module_number) < c->frame_number)
head.at(c->module_number) = c->frame_number;
if (curr_frame_number.at(c->module_number) < c->frame_number)
curr_frame_number.at(c->module_number) = c->frame_number;
if (c->frame_number > slowest_head) {
slowest_head = head[0];
if (c->frame_number > slowest_frame_number) {
slowest_frame_number = curr_frame_number[0];
for (int i = 1; i < nmodules; i++)
if (head[i] < slowest_head)
slowest_head = head[i];
if (curr_frame_number[i] < slowest_frame_number)
slowest_frame_number = curr_frame_number[i];
}
if (c->frame_number > fastest_frame_number) {
fastest_frame_number = c->frame_number;
}
packets_collected.at(c->frame_number * nmodules + c->module_number) = c->packet_count;
@@ -96,22 +102,26 @@ uint64_t AcquisitionCounters::GetBufferHandleAndClear(size_t frame, uint16_t mod
return ret_val;
}
uint64_t AcquisitionCounters::GetHead(uint16_t module_number) const {
uint64_t AcquisitionCounters::GetCurrFrameNumber(uint16_t module_number) const {
if (module_number >= max_modules)
throw JFJochException(JFJochExceptionCategory::ArrayOutOfBounds,
"GetHead Wrong module number: " + std::to_string(module_number));
return head[module_number];
"GetCurrFrameNumber Wrong module number: " + std::to_string(module_number));
return curr_frame_number[module_number];
}
uint64_t AcquisitionCounters::GetSlowestHead() const {
return slowest_head;
uint64_t AcquisitionCounters::GetSlowestFrameNumber() const {
return slowest_frame_number;
}
uint64_t AcquisitionCounters::GetFastestFrameNumber() const {
return fastest_frame_number;
}
void AcquisitionCounters::WaitForFrame(size_t curr_frame, uint16_t module_number) const {
uint64_t slowest_head_tmp = (module_number == UINT16_MAX) ? GetSlowestHead() : GetHead(module_number);
uint64_t slowest_head_tmp = (module_number == UINT16_MAX) ? GetSlowestFrameNumber() : GetCurrFrameNumber(module_number);
while (!acquisition_finished && (slowest_head_tmp < curr_frame)) {
std::this_thread::sleep_for(std::chrono::microseconds(100));
slowest_head_tmp = (module_number == UINT16_MAX) ? GetSlowestHead() : GetHead(module_number);
slowest_head_tmp = (module_number == UINT16_MAX) ? GetSlowestFrameNumber() : GetCurrFrameNumber(module_number);
}
}
@@ -119,9 +129,9 @@ int64_t AcquisitionCounters::CalculateDelay(size_t curr_frame, uint16_t module_n
uint64_t slowest_head_tmp;
if (module_number == UINT16_MAX)
slowest_head_tmp = GetSlowestHead();
slowest_head_tmp = GetSlowestFrameNumber();
else
slowest_head_tmp = GetHead(module_number);
slowest_head_tmp = GetCurrFrameNumber(module_number);
return slowest_head_tmp - curr_frame;
}
+6 -4
View File
@@ -29,8 +29,9 @@ class AcquisitionCounters {
uint64_t total_packets;
std::vector<uint64_t> packets_per_module;
uint64_t slowest_head;
std::vector<uint64_t> head;
uint64_t slowest_frame_number;
uint64_t fastest_frame_number;
std::vector<uint64_t> curr_frame_number;
bool acquisition_finished;
uint64_t expected_frames;
uint64_t nmodules = max_modules;
@@ -45,8 +46,9 @@ public:
uint64_t GetBufferHandleAndClear(size_t frame, uint16_t module_number);
uint64_t GetHead(uint16_t module_number) const;
uint64_t GetSlowestHead() const;
uint64_t GetCurrFrameNumber(uint16_t module_number) const;
uint64_t GetSlowestFrameNumber() const;
uint64_t GetFastestFrameNumber() const;
void WaitForFrame(size_t curr_frame, uint16_t module_number = UINT16_MAX) const;
int64_t CalculateDelay(size_t curr_frame, uint16_t module_number = UINT16_MAX) const; // mutex acquired indirectly
+1 -1
View File
@@ -315,7 +315,7 @@ JFJochProtoBuf::FPGAStatus FPGAAcquisitionDevice::GetStatus() const {
ret.set_ethernet_rx_aligned(env.ethernet_aligned);
ret.set_hbm_temp_0_degc(env.hbm_0_temp_C);
ret.set_hbm_temp_1_degc(env.hbm_1_temp_C);
ret.set_slowest_head(counters.GetSlowestHead());
ret.set_slowest_head(counters.GetSlowestFrameNumber());
return ret;
}
+8 -1
View File
@@ -417,6 +417,13 @@ void JFJochReceiver::FrameTransformationThread() {
uint64_t image_number;
while (images_to_go.Get(image_number) != 0) {
try {
// If data acquisition is finished and fastest frame for the first device is behind
size_t frame_number = image_number * experiment.GetSummation();
acquisition_device[0]->Counters().WaitForFrame(frame_number + 2);
if (acquisition_device[0]->Counters().IsAcquisitionFinished() &&
(acquisition_device[0]->Counters().GetFastestFrameNumber() < frame_number))
continue;
DataMessage message{};
message.number = image_number;
message.timestamp_base = 10*1000*1000;
@@ -450,7 +457,7 @@ void JFJochReceiver::FrameTransformationThread() {
message.receiver_aq_dev_delay = max_delay;
} else {
for (int j = 0; j < experiment.GetSummation(); j++) {
size_t frame_number = image_number * experiment.GetSummation() + j;
frame_number = image_number * experiment.GetSummation() + j;
for (int d = 0; d < ndatastreams; d++) {
acquisition_device[d]->Counters().WaitForFrame(frame_number + 2);
+17 -13
View File
@@ -9,10 +9,11 @@ TEST_CASE("AcquisitionCountersTest","[AcquisitionDeviceCounters]") {
AcquisitionCounters counters;
counters.Reset(x, 0);
REQUIRE(counters.GetSlowestHead() == 0);
REQUIRE(counters.GetHead(0) == 0);
REQUIRE(counters.GetHead(1) == 0);
REQUIRE_THROWS(counters.GetHead(32));
REQUIRE(counters.GetSlowestFrameNumber() == 0);
REQUIRE(counters.GetFastestFrameNumber() == 0);
REQUIRE(counters.GetCurrFrameNumber(0) == 0);
REQUIRE(counters.GetCurrFrameNumber(1) == 0);
REQUIRE_THROWS(counters.GetCurrFrameNumber(32));
REQUIRE(counters.CalculateDelay(2) == -2);
REQUIRE(!counters.IsAcquisitionFinished());
@@ -26,9 +27,10 @@ TEST_CASE("AcquisitionCountersTest","[AcquisitionDeviceCounters]") {
c.timestamp = 456;
counters.UpdateCounters(&c);
REQUIRE(counters.GetSlowestHead() == 0);
REQUIRE(counters.GetHead(0) == 0);
REQUIRE(counters.GetHead(1) == 32);
REQUIRE(counters.GetSlowestFrameNumber() == 0);
REQUIRE(counters.GetFastestFrameNumber() == 32);
REQUIRE(counters.GetCurrFrameNumber(0) == 0);
REQUIRE(counters.GetCurrFrameNumber(1) == 32);
REQUIRE(counters.CalculateDelay(31, 1) == 1);
REQUIRE(counters.CalculateDelay(33, 1) == -1);
@@ -42,9 +44,10 @@ TEST_CASE("AcquisitionCountersTest","[AcquisitionDeviceCounters]") {
c.module_number = 0;
counters.UpdateCounters(&c);
REQUIRE(counters.GetSlowestHead() == 15);
REQUIRE(counters.GetHead(0) == 15);
REQUIRE(counters.GetHead(1) == 32);
REQUIRE(counters.GetSlowestFrameNumber() == 15);
REQUIRE(counters.GetFastestFrameNumber() == 32);
REQUIRE(counters.GetCurrFrameNumber(0) == 15);
REQUIRE(counters.GetCurrFrameNumber(1) == 32);
REQUIRE(counters.CalculateDelay(14) == 1);
REQUIRE(counters.CalculateDelay(16) == -1);
@@ -57,9 +60,10 @@ TEST_CASE("AcquisitionCountersTest","[AcquisitionDeviceCounters]") {
counters.SetAcquisitionFinished();
REQUIRE(counters.GetSlowestHead() == 15);
REQUIRE(counters.GetHead(0) == 15);
REQUIRE(counters.GetHead(1) == 32);
REQUIRE(counters.GetFastestFrameNumber() == 32);
REQUIRE(counters.GetSlowestFrameNumber() == 15);
REQUIRE(counters.GetCurrFrameNumber(0) == 15);
REQUIRE(counters.GetCurrFrameNumber(1) == 32);
REQUIRE(counters.CalculateDelay(0) == 15);
REQUIRE(counters.CalculateDelay(50) == 15-50);
}
+5 -5
View File
@@ -99,7 +99,7 @@ TEST_CASE("HLS_C_Simulation_check_raw", "[FPGA][Full]") {
REQUIRE_NOTHROW(test.StartAction(x));
REQUIRE_NOTHROW(test.WaitForActionComplete());
REQUIRE(test.Counters().GetSlowestHead() == 0);
REQUIRE(test.Counters().GetSlowestFrameNumber() == 0);
REQUIRE_NOTHROW(test.OutputStream().read());
REQUIRE(test.OutputStream().size() == 0);
@@ -138,7 +138,7 @@ TEST_CASE("HLS_C_Simulation_check_cancel", "[FPGA][Full]") {
REQUIRE_NOTHROW(test.WaitForActionComplete());
REQUIRE(test.Counters().GetSlowestHead() == 0);
REQUIRE(test.Counters().GetSlowestFrameNumber() == 0);
REQUIRE_NOTHROW(test.OutputStream().read());
REQUIRE(test.OutputStream().size() == 0);
@@ -163,7 +163,7 @@ TEST_CASE("HLS_C_Simulation_check_cancel_conversion", "[FPGA][Full]") {
REQUIRE_NOTHROW(test.WaitForActionComplete());
REQUIRE(test.Counters().GetSlowestHead() == 0);
REQUIRE(test.Counters().GetSlowestFrameNumber() == 0);
REQUIRE_NOTHROW(test.OutputStream().read());
REQUIRE(test.OutputStream().size() == 0);
@@ -614,8 +614,8 @@ TEST_CASE("HLS_C_Simulation_check_2_trigger_convert", "[FPGA][Full]") {
// address properly aligned
REQUIRE((uint64_t) test.GetDeviceOutput(0,0)->pixels % 128 == 0);
REQUIRE(test.Counters().GetSlowestHead() == 0);
REQUIRE(test.Counters().GetHead(0) == 9);
REQUIRE(test.Counters().GetSlowestFrameNumber() == 0);
REQUIRE(test.Counters().GetCurrFrameNumber(0) == 9);
REQUIRE_NOTHROW(test.OutputStream().read());
REQUIRE(test.OutputStream().size() == 0);