diff --git a/broker/JFJochBrokerHttp.cpp b/broker/JFJochBrokerHttp.cpp index 25316a96..6117fe33 100644 --- a/broker/JFJochBrokerHttp.cpp +++ b/broker/JFJochBrokerHttp.cpp @@ -403,9 +403,10 @@ void JFJochBrokerHttp::wait_until_running_post(const std::optional &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 &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"); diff --git a/broker/JFJochStateMachine.cpp b/broker/JFJochStateMachine.cpp index d04867f1..5189335c 100644 --- a/broker/JFJochStateMachine.cpp +++ b/broker/JFJochStateMachine.cpp @@ -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)); } diff --git a/broker/JFJochStateMachine.h b/broker/JFJochStateMachine.h index 40f49e5e..17ad3e58 100644 --- a/broker/JFJochStateMachine.h +++ b/broker/JFJochStateMachine.h @@ -100,7 +100,10 @@ class JFJochStateMachine { PixelMask pixel_mask; int64_t current_detector_setup; // Lock only on change std::optional 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; diff --git a/docs/CHANGELOG.md b/docs/CHANGELOG.md index dc236b5e..241a1adb 100644 --- a/docs/CHANGELOG.md +++ b/docs/CHANGELOG.md @@ -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 diff --git a/receiver/JFJochReceiverService.cpp b/receiver/JFJochReceiverService.cpp index e2ec17c2..1f810c25 100644 --- a/receiver/JFJochReceiverService.cpp +++ b/receiver/JFJochReceiverService.cpp @@ -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; } diff --git a/tests/JFJochStateMachineTest.cpp b/tests/JFJochStateMachineTest.cpp index c27d2c1f..884a39e3 100644 --- a/tests/JFJochStateMachineTest.cpp +++ b/tests/JFJochStateMachineTest.cpp @@ -3,6 +3,8 @@ #include #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(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); +}