Broker: report a failed asynchronous start instead of a timeout
An acquisition started with async_start returns from /start before the measurement thread has run, so a failure there had no caller to raise it to. MeasurementThread stored it in start_exception, but only the synchronous branch of Start() ever read it. The ordinary-failure path - a writer refusing to overwrite an existing file, say - sets the state to Idle rather than Error, and Idle is what wait_until_running_post maps to 504 "timeout, need to restart". So the one case the async workflow exists to report came back as a timeout with an empty body, and the writer's message was dropped. /wait_till_done was worse: Idle is its success case, so it answered 200 for a run that never started. Only a critical detector fault, which leaves the state at Error, was reported at all. The wait functions now rethrow start_exception, so both endpoints produce the same 500 and the same message as a synchronous start. rethrow_exception does not consume the exception_ptr, so repeated calls all report the same failure. It is cleared by every entry point that begins new work - Start (before ImportDatasetSettings, which can throw), Initialize, Pedestal, LoadDetectorSettings, SetDarkMaskSettings - so a pending failure is never attributed to the operation after it. The synchronous branch no longer clears it, so a wait call made after a failed /start reports the failure rather than an apparent timeout. wait_until_running_post and wait_till_done_post drop their timeout == 0 special case, which called GetStatus() directly and so bypassed the rethrow: ?timeout=0 returned 200 from /wait_till_done where ?timeout=1 returned the error. wait_for evaluates the predicate before expiring, so the state reported is unchanged. JFJochReceiverService::Start clears the receiver status when the start fails. JFJochReceiver's constructor sets progress to 0 and a zeroed status before the throw, and nothing cleared it, so /status reported Idle with progress 0 - a stalled acquisition, to the frontend - and /statistics reported an all-zero run that never happened, until the next start overwrote it. Its catch widens to std::exception, so a non-JFJochException gets the cleanup and the logging it skipped before. Nothing else is left dirty by a refused start: the receiver keeps its null receiver and Idle state, PrepareAction only resets counters (the FPGA is armed in StartAction, in a thread launched after the start message goes out), and SendStartMessage precedes every std::async in the receiver constructor. The test starts four more times afterwards, one of them a complete acquisition with no re-initialisation in between. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Uwv9ScHtDH6g8tYgfSuApo
This commit is contained in:
@@ -403,9 +403,10 @@ void JFJochBrokerHttp::wait_until_running_post(const std::optional<int32_t> &tim
|
||||
response.status = 400;
|
||||
response.set_content("timeout must be in range 0..3600", "text/plain");
|
||||
return;
|
||||
} else if (timeout.value() == 0)
|
||||
status = state_machine.GetStatus();
|
||||
else
|
||||
} else
|
||||
// A zero timeout still goes through the wait function rather than GetStatus(): the predicate
|
||||
// is evaluated before the wait expires, so the answer is the same, but a start that failed
|
||||
// asynchronously is reported instead of being read back as a plain state.
|
||||
status = state_machine.WaitTillNotBusy(std::chrono::seconds(timeout.value()));
|
||||
|
||||
logger.Info("Wait until running");
|
||||
@@ -435,9 +436,7 @@ void JFJochBrokerHttp::wait_till_done_post(const std::optional<int32_t> &timeout
|
||||
response.status = 400;
|
||||
response.set_content("timeout must be in range 0..3600", "text/plain");
|
||||
return;
|
||||
} else if (timeout.value() == 0)
|
||||
status = state_machine.GetStatus();
|
||||
else
|
||||
} else
|
||||
status = state_machine.WaitTillMeasurementDone(std::chrono::seconds(timeout.value()));
|
||||
|
||||
logger.Info("Wait till done");
|
||||
|
||||
@@ -289,6 +289,7 @@ void JFJochStateMachine::Initialize() {
|
||||
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "Detector information not provided");
|
||||
|
||||
ResetError(); // Clear error, we don't care what was it
|
||||
start_exception = nullptr; // Re-initialising discards a pending start failure
|
||||
|
||||
logger.Info("Initialize");
|
||||
SetState(JFJochState::Busy, "Configuring indexing threads", BrokerStatus::MessageSeverity::Info);
|
||||
@@ -311,6 +312,7 @@ void JFJochStateMachine::Pedestal() {
|
||||
if (state != JFJochState::Idle)
|
||||
throw WrongDAQStateException("Must be idle to take pedestal");
|
||||
|
||||
start_exception = nullptr; // A new operation supersedes a pending start failure
|
||||
SetState(JFJochState::Busy, "Updating calibration", BrokerStatus::MessageSeverity::Info);
|
||||
|
||||
measurement = std::async(std::launch::async, &JFJochStateMachine::CalibrateDetector, this, std::move(ul));
|
||||
@@ -352,6 +354,10 @@ void JFJochStateMachine::Start(const DatasetSettings &settings, bool async) {
|
||||
if (measurement.valid())
|
||||
measurement.get(); // In case measurement was running - clear thread
|
||||
|
||||
// Clear before ImportDatasetSettings, which can throw: a rejected /start must not leave the
|
||||
// previous run's failure behind for the next wait call to report.
|
||||
start_exception = nullptr;
|
||||
|
||||
experiment.ImportDatasetSettings(settings);
|
||||
|
||||
cancel_sequence = false;
|
||||
@@ -362,24 +368,28 @@ void JFJochStateMachine::Start(const DatasetSettings &settings, bool async) {
|
||||
|
||||
experiment.IncrementRunNumber();
|
||||
|
||||
start_exception = nullptr;
|
||||
SetState(JFJochState::Busy, "Preparing measurement", BrokerStatus::MessageSeverity::Info);
|
||||
measurement = std::async(std::launch::async, &JFJochStateMachine::MeasurementThread, this);
|
||||
if (!async) {
|
||||
c.wait(ul, [&]() { return state != JFJochState::Busy; });
|
||||
// A synchronous start propagates the failure to the caller. The state has already been set
|
||||
// by MeasurementThread (Idle for an ordinary failure, Error for a critical detector fault).
|
||||
if (start_exception) {
|
||||
auto e = start_exception;
|
||||
start_exception = nullptr;
|
||||
std::rethrow_exception(e);
|
||||
}
|
||||
// start_exception is left in place - the next Start() or Initialize() clears it - so that a
|
||||
// wait call made afterwards reports the same failure instead of an apparent timeout.
|
||||
if (start_exception)
|
||||
std::rethrow_exception(start_exception);
|
||||
}
|
||||
}
|
||||
|
||||
BrokerStatus JFJochStateMachine::WaitTillNotBusy(std::chrono::milliseconds timeout) {
|
||||
std::unique_lock ul(m);
|
||||
c.wait_for(ul, timeout, [&]() { return state != JFJochState::Busy; });
|
||||
// An asynchronous start reports its failure here, since /start itself returned before the
|
||||
// measurement thread ran. Without this the state is plain Idle and the caller cannot tell a
|
||||
// failed start from a timeout. rethrow_exception does not consume the exception_ptr, so
|
||||
// repeated calls all report the same failure.
|
||||
if (start_exception)
|
||||
std::rethrow_exception(start_exception);
|
||||
return GetStatus();
|
||||
}
|
||||
|
||||
@@ -606,6 +616,7 @@ void JFJochStateMachine::LoadDetectorSettings(const DetectorSettings &settings)
|
||||
break;
|
||||
case JFJochState::Idle:
|
||||
if (ImportDetectorSettings(settings)) {
|
||||
start_exception = nullptr; // A new operation supersedes a pending start failure
|
||||
SetState(JFJochState::Busy, "Loading settings", BrokerStatus::MessageSeverity::Info);
|
||||
measurement = std::async(std::launch::async, &JFJochStateMachine::CalibrateDetector, this, std::move(ul));
|
||||
} else {
|
||||
@@ -801,6 +812,11 @@ BrokerStatus JFJochStateMachine::WaitTillMeasurementDone() {
|
||||
|
||||
c.wait(ul, [&] { return !IsRunning(); });
|
||||
|
||||
// A start that failed asynchronously never reached Measuring, so the state is Idle and would
|
||||
// otherwise be reported as a successfully finished collection.
|
||||
if (start_exception)
|
||||
std::rethrow_exception(start_exception);
|
||||
|
||||
return GetStatus();
|
||||
}
|
||||
|
||||
@@ -809,6 +825,9 @@ BrokerStatus JFJochStateMachine::WaitTillMeasurementDone(std::chrono::millisecon
|
||||
|
||||
c.wait_for(ul, timeout, [&] { return !IsRunning(); });
|
||||
|
||||
if (start_exception)
|
||||
std::rethrow_exception(start_exception);
|
||||
|
||||
return GetStatus();
|
||||
}
|
||||
|
||||
@@ -1158,6 +1177,7 @@ void JFJochStateMachine::SetDarkMaskSettings(const DarkMaskSettings &settings) {
|
||||
}
|
||||
if ((experiment.GetDetectorType() == DetectorType::DECTRIS) && (state == JFJochState::Idle)) {
|
||||
// Need to redo the calibration
|
||||
start_exception = nullptr; // A new operation supersedes a pending start failure
|
||||
SetState(JFJochState::Busy, "Loading settings", BrokerStatus::MessageSeverity::Info);
|
||||
measurement = std::async(std::launch::async, &JFJochStateMachine::CalibrateDetector, this, std::move(ul));
|
||||
}
|
||||
|
||||
@@ -100,7 +100,10 @@ class JFJochStateMachine {
|
||||
PixelMask pixel_mask;
|
||||
int64_t current_detector_setup; // Lock only on change
|
||||
std::optional<ScanResult> scan_result;
|
||||
// Set by MeasurementThread when a synchronous Start fails, so Start() can rethrow to the caller
|
||||
// Set by MeasurementThread when a Start fails. A synchronous Start() rethrows it directly; an
|
||||
// asynchronous one has already returned, so the wait functions rethrow it instead. Every entry
|
||||
// point that begins new work clears it, so it is reported to every caller asking in between but
|
||||
// never attributed to the operation after it.
|
||||
std::exception_ptr start_exception;
|
||||
|
||||
mutable std::mutex calibration_statistics_mutex;
|
||||
|
||||
@@ -16,6 +16,7 @@ This is an UNSTABLE release. It includes many experimental features, as well as
|
||||
* Jungfraujoch needs six fewer shared libraries on the machine - libopenblas and libmetis, and libgfortran, libquadmath, libgomp and libz behind them - because the Ceres LAPACK, METIS and SuiteSparse back-ends are no longer built. Nothing in the code ever selected them, and results are unchanged.
|
||||
* The PCIe driver DKMS package builds for the kernel it is being installed for instead of the running one, so a module built while a kernel update is being applied loads after the reboot.
|
||||
* The PCIe driver builds on RHEL 9.5 and later, and on their CentOS Stream, Rocky and AlmaLinux equivalents, where the `vm_flags` kernel interface was backported into the 5.14 kernel.
|
||||
* A data collection started with `async_start` that fails to start - a writer refusing to overwrite an existing file, for instance - is reported as an error by `/wait_until_running` and `/wait_till_done` instead of as a timeout and a successful collection respectively. The error message is the one the writer gave.
|
||||
* The results report's `REPORT_VERSION` is 3, two sections having been added. Existing key names and table columns are unchanged.
|
||||
|
||||
### 1.0.0-rc.163
|
||||
|
||||
@@ -103,7 +103,12 @@ void JFJochReceiverService::Start(const DiffractionExperiment &experiment,
|
||||
}
|
||||
measurement = std::async(std::launch::async, &JFJochReceiverService::FinalizeMeasurement, this);
|
||||
state = ReceiverState::Running;
|
||||
} catch (const JFJochException &e) {
|
||||
} catch (const std::exception &e) {
|
||||
// The receiver never started, so drop the status its base constructor had already reset -
|
||||
// otherwise /status and /statistics keep reporting a zero-progress run that never happened,
|
||||
// until the next start overwrites it.
|
||||
receiver_status.Clear();
|
||||
receiver_status.SetProgress({});
|
||||
logger.ErrorException(e);
|
||||
throw;
|
||||
}
|
||||
|
||||
@@ -3,6 +3,8 @@
|
||||
|
||||
#include <catch2/catch_all.hpp>
|
||||
#include "../broker/JFJochStateMachine.h"
|
||||
#include "../acquisition_device/HLSSimulatedDevice.h"
|
||||
#include "../receiver/JFJochReceiverService.h"
|
||||
|
||||
using namespace std::literals::chrono_literals;
|
||||
|
||||
@@ -199,3 +201,117 @@ TEST_CASE("JFJochStateMachine_LoadDetectorSettings_Error") {
|
||||
REQUIRE_THROWS(state_machine.LoadDetectorSettings(settings));
|
||||
REQUIRE(state_machine.GetStatus().state == JFJochState::Idle);
|
||||
}
|
||||
|
||||
namespace {
|
||||
// Stands in for a writer that refuses to start the run - e.g. the output file exists and cannot be
|
||||
// overwritten. The refusal reaches the broker as an ordinary exception thrown while the start
|
||||
// message is being sent, which is what separates it from a critical detector fault.
|
||||
class RefusingImagePusher : public ImagePusher {
|
||||
public:
|
||||
bool refuse = true;
|
||||
void StartDataCollection(StartMessage &) override {
|
||||
if (refuse)
|
||||
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "writer_refused_1234");
|
||||
}
|
||||
bool EndDataCollection(const EndMessage &) override { return true; }
|
||||
bool SendImage(const uint8_t *, size_t, int64_t) override { return true; }
|
||||
bool SendCalibration(const CompressedImage &) override { return true; }
|
||||
std::string PrintSetup() const override { return "RefusingImagePusher"; }
|
||||
ImagePusherType GetType() const override { return ImagePusherType::Test; }
|
||||
};
|
||||
}
|
||||
|
||||
// An asynchronous start returns before the measurement thread runs, so a failure there has to be
|
||||
// reported to whoever asks next. Without this the state is a plain Idle, indistinguishable from a
|
||||
// timeout, and the writer's message is lost.
|
||||
TEST_CASE("JFJochStateMachine_AsyncStartFailure") {
|
||||
Logger logger("JFJochStateMachine_AsyncStartFailure");
|
||||
|
||||
DiffractionExperiment experiment(DetJF(2));
|
||||
experiment.Conversion().PedestalG0Frames(0).NumTriggers(1).UseInternalPacketGenerator(true)
|
||||
.ImagesPerTrigger(4).IncidentEnergy_keV(12.4);
|
||||
|
||||
AcquisitionDeviceGroup aq_devices;
|
||||
for (int i = 0; i < experiment.GetDataStreamsNum(); i++)
|
||||
aq_devices.Add(std::make_unique<HLSSimulatedDevice>(i, 64));
|
||||
|
||||
RefusingImagePusher pusher;
|
||||
JFJochReceiverService receiver_service(aq_devices, logger, pusher);
|
||||
|
||||
JFJochServices services(logger);
|
||||
services.Receiver(&receiver_service);
|
||||
|
||||
JFJochStateMachine state_machine(experiment, services, logger);
|
||||
state_machine.AddDetectorSetup(DetJF(2));
|
||||
state_machine.DebugOnly_SetState(JFJochState::Idle);
|
||||
|
||||
DatasetSettings setup;
|
||||
setup.ImagesPerTrigger(4).NumTriggers(1);
|
||||
|
||||
// The asynchronous start itself succeeds - it only launches the measurement thread.
|
||||
REQUIRE_NOTHROW(state_machine.Start(setup, true));
|
||||
|
||||
// Both wait functions report the failure, and repeatedly: the exception_ptr is not consumed.
|
||||
REQUIRE_THROWS_WITH(state_machine.WaitTillNotBusy(std::chrono::seconds(10)),
|
||||
Catch::Matchers::ContainsSubstring("writer_refused_1234"));
|
||||
REQUIRE_THROWS_WITH(state_machine.WaitTillNotBusy(std::chrono::seconds(10)),
|
||||
Catch::Matchers::ContainsSubstring("writer_refused_1234"));
|
||||
REQUIRE_THROWS_WITH(state_machine.WaitTillMeasurementDone(std::chrono::seconds(10)),
|
||||
Catch::Matchers::ContainsSubstring("writer_refused_1234"));
|
||||
|
||||
// An ordinary failure leaves the detector usable, so the run can be retried from Idle. The
|
||||
// status keeps the writer's message, and nothing is left of the run that never happened - a
|
||||
// progress figure in particular would show up in the frontend as a stalled acquisition.
|
||||
auto status = state_machine.GetStatus();
|
||||
REQUIRE(status.state == JFJochState::Idle);
|
||||
REQUIRE(status.message_severity == BrokerStatus::MessageSeverity::Error);
|
||||
REQUIRE_THAT(status.message.value_or(""), Catch::Matchers::ContainsSubstring("writer_refused_1234"));
|
||||
REQUIRE_FALSE(status.progress.has_value());
|
||||
|
||||
auto statistics = state_machine.GetMeasurementStatistics();
|
||||
REQUIRE(statistics.has_value());
|
||||
REQUIRE(statistics->images_collected == 0);
|
||||
REQUIRE(statistics->images_sent == 0);
|
||||
REQUIRE_FALSE(statistics->collection_efficiency.has_value());
|
||||
|
||||
// A synchronous start reports the same failure directly, and leaves it visible to the wait
|
||||
// functions afterwards.
|
||||
REQUIRE_THROWS_WITH(state_machine.Start(setup),
|
||||
Catch::Matchers::ContainsSubstring("writer_refused_1234"));
|
||||
REQUIRE_THROWS_WITH(state_machine.WaitTillNotBusy(std::chrono::seconds(10)),
|
||||
Catch::Matchers::ContainsSubstring("writer_refused_1234"));
|
||||
|
||||
// A start rejected before the measurement thread is launched must not leave the previous
|
||||
// failure behind for the next wait call to report.
|
||||
DatasetSettings rejected_setup;
|
||||
rejected_setup.ImagesPerTrigger(4).NumTriggers(1)
|
||||
.ImageTime(state_machine.Experiment().GetFrameTime() + std::chrono::nanoseconds(1));
|
||||
REQUIRE_THROWS(state_machine.Start(rejected_setup, true));
|
||||
REQUIRE_NOTHROW(state_machine.WaitTillNotBusy(std::chrono::seconds(10)));
|
||||
|
||||
// Any other operation started from Idle supersedes the pending failure - it must not be
|
||||
// attributed to the pedestal, which is allowed from Idle and has nothing to do with it.
|
||||
REQUIRE_THROWS_WITH(state_machine.Start(setup),
|
||||
Catch::Matchers::ContainsSubstring("writer_refused_1234"));
|
||||
REQUIRE_NOTHROW(state_machine.Pedestal());
|
||||
REQUIRE_NOTHROW(state_machine.WaitTillMeasurementDone());
|
||||
|
||||
// Nothing is left half-started: once the writer accepts the run, the next start goes through
|
||||
// and completes without any intervening re-initialisation.
|
||||
pusher.refuse = false;
|
||||
REQUIRE_NOTHROW(state_machine.Start(setup));
|
||||
REQUIRE(state_machine.GetStatus().state == JFJochState::Measuring);
|
||||
REQUIRE_NOTHROW(state_machine.WaitTillMeasurementDone());
|
||||
REQUIRE(state_machine.GetStatus().state == JFJochState::Idle);
|
||||
REQUIRE(state_machine.GetMeasurementStatistics()->images_sent == 4);
|
||||
|
||||
// Re-initialising discards a pending failure too, so a detector brought back up does not report
|
||||
// the failure of the run before it.
|
||||
pusher.refuse = true;
|
||||
REQUIRE_THROWS_WITH(state_machine.Start(setup),
|
||||
Catch::Matchers::ContainsSubstring("writer_refused_1234"));
|
||||
services.Receiver(nullptr);
|
||||
REQUIRE_NOTHROW(state_machine.Initialize());
|
||||
REQUIRE_NOTHROW(state_machine.WaitTillMeasurementDone());
|
||||
REQUIRE(state_machine.GetStatus().state == JFJochState::Idle);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user