From 1e6b5c48e7eab62f54d0f709624e41cfcfd31a8b Mon Sep 17 00:00:00 2001 From: Nathan Brei Date: Mon, 24 Aug 2026 12:04:51 -0400 Subject: [PATCH 1/3] Move JEvent::Clear() logic into JEventPool::Ingest() This addresses issue #515. --- src/libraries/JANA/JEvent.cc | 8 ++++++-- src/libraries/JANA/JEvent.h | 2 +- src/libraries/JANA/JEventSource.cc | 9 +++++++-- src/libraries/JANA/Topology/JArrow.cc | 1 - src/libraries/JANA/Topology/JEventPool.cc | 2 ++ 5 files changed, 16 insertions(+), 6 deletions(-) diff --git a/src/libraries/JANA/JEvent.cc b/src/libraries/JANA/JEvent.cc index d8086cf83..3807a19cb 100644 --- a/src/libraries/JANA/JEvent.cc +++ b/src/libraries/JANA/JEvent.cc @@ -148,10 +148,14 @@ int JEvent::GetChildCount() { return mReferenceCount; } -void JEvent::Clear(bool processed_successfully) { - if (processed_successfully && mEventSource != nullptr) { +void JEvent::Clear() { + // Calling JEvent::FinishEvent() without a preceding Emit()/GetEvent() is not allowed. + // (This is not academic -- GlueX relies on this behavior) + // We test whether the event was in fact emitted by checking if mEventSource==nullptr. + if (mEventSource != nullptr) { mEventSource->DoFinishEvent(*this); mIsWarmedUp = true; + mEventSource = nullptr; } mFactorySet.Clear(); mInspector.Reset(); diff --git a/src/libraries/JANA/JEvent.h b/src/libraries/JANA/JEvent.h index 2ab8e7241..3d4371078 100644 --- a/src/libraries/JANA/JEvent.h +++ b/src/libraries/JANA/JEvent.h @@ -97,7 +97,7 @@ class JEvent : public std::enable_shared_from_this { void SetParentNumber(JEventLevel level, uint64_t number); // Lifecycle - void Clear(bool processed_successfully=true); + void Clear(); void Finish(); JFactory* GetFactory(const std::string& object_name, const std::string& tag) const; diff --git a/src/libraries/JANA/JEventSource.cc b/src/libraries/JANA/JEventSource.cc index f3790fd35..8a900e805 100644 --- a/src/libraries/JANA/JEventSource.cc +++ b/src/libraries/JANA/JEventSource.cc @@ -95,7 +95,6 @@ JEventSource::Result JEventSource::DoNext(std::shared_ptr event) { // We configure the event event->SetEventNumber(m_events_emitted); // Default event number to event count - event->SetJEventSource(this); event->SetSequential(false); event->GetJCallGraphRecorder()->Reset(); @@ -159,6 +158,10 @@ JEventSource::Result JEventSource::DoNext(std::shared_ptr event) { if (result == Result::Success) { m_events_emitted += 1; + // We only set the JEventSource once Emit/GetEvent succeeded because it also controls + // m_events_processed and EventSource::FinishEvent + event->SetJEventSource(this); + // We end up here if we read an entry in our file or retrieved a message from our socket, // and believe we could obtain another one immediately if we wanted to for (auto* output : GetOutputs()) { @@ -213,6 +216,7 @@ std::pair JEventSource::Skip(JEvent& event, size_t while (events_to_skip > 0 && result == Result::Success) { try { + event.SetJEventSource(this); auto previous_origin = event.GetJCallGraphRecorder()->SetInsertDataOrigin( JCallGraphRecorder::ORIGIN_FROM_SOURCE); // (see note at top of JCallGraphRecorder.h) if (m_callback_style == CallbackStyle::LegacyMode) { GetEvent(event.shared_from_this()); @@ -226,7 +230,8 @@ std::pair JEventSource::Skip(JEvent& event, size_t if (m_enable_finish_event) { CallWithJExceptionWrapper("JEventSource::FinishEvent", [&](){ FinishEvent(event); }); } - event.Clear(false); + event.SetJEventSource(nullptr); // This tells event::Clear() _not_ to call FinishEvent() + event.Clear(); events_to_skip -= 1; } catch (RETURN_STATUS rs) { diff --git a/src/libraries/JANA/Topology/JArrow.cc b/src/libraries/JANA/Topology/JArrow.cc index 864c7b18f..88533e307 100644 --- a/src/libraries/JANA/Topology/JArrow.cc +++ b/src/libraries/JANA/Topology/JArrow.cc @@ -49,7 +49,6 @@ void JArrow::Push(OutputData& outputs, size_t output_count, size_t location_id) port.GetQueue()->Push(event, location_id); } else if (port.GetPool() != nullptr) { - event->Clear(!port.GetSkipFinishEvent()); port.GetPool()->Ingest(event, location_id); } else { diff --git a/src/libraries/JANA/Topology/JEventPool.cc b/src/libraries/JANA/Topology/JEventPool.cc index 9389960c9..5be964661 100644 --- a/src/libraries/JANA/Topology/JEventPool.cc +++ b/src/libraries/JANA/Topology/JEventPool.cc @@ -87,6 +87,7 @@ void JEventPool::Ingest(JEvent* event, size_t location) { if (event->GetChildCount() == 0) { // There's no way for additional children to appear because Ingest takes the "original" parent //LOG << "JEventPool::Ingest: " << toString(m_level) << " event is pushed"; + event->Clear(); Push(event, location); } else { @@ -102,6 +103,7 @@ void JEventPool::NotifyThatAllChildrenFinished(JEvent* event, size_t location) { //LOG << "JEventPool::Notify called for level " << toString(m_level); size_t was_present = m_pending.erase(event); if (was_present == 1) { + event->Clear(); Push(event, location); //LOG << "JEventPool at level " << toString(m_level) << " has pushed a parent event"; } From 94e57b5f67f90c182445786e20351798286a6114 Mon Sep 17 00:00:00 2001 From: Nathan Brei Date: Mon, 24 Aug 2026 12:17:51 -0400 Subject: [PATCH 2/3] Add test case for hierarchical ownership --- src/programs/integration_tests/CMakeLists.txt | 1 + .../HierarchicalOwnership.cc | 70 +++++++++++++++++++ 2 files changed, 71 insertions(+) create mode 100644 src/programs/integration_tests/HierarchicalOwnership.cc diff --git a/src/programs/integration_tests/CMakeLists.txt b/src/programs/integration_tests/CMakeLists.txt index acba6b103..d5c3fc2d3 100644 --- a/src/programs/integration_tests/CMakeLists.txt +++ b/src/programs/integration_tests/CMakeLists.txt @@ -3,6 +3,7 @@ set(JANA2_INTEGRATION_TEST_SOURCES SimpleOffloading.cc BatchedArrow.cc + HierarchicalOwnership.cc ) add_jana_test(jana-integration-tests SOURCES ${JANA2_INTEGRATION_TEST_SOURCES}) diff --git a/src/programs/integration_tests/HierarchicalOwnership.cc b/src/programs/integration_tests/HierarchicalOwnership.cc new file mode 100644 index 000000000..a8122cebf --- /dev/null +++ b/src/programs/integration_tests/HierarchicalOwnership.cc @@ -0,0 +1,70 @@ +#include +#include +#include +#include + +#include +#include +#include +#include + +namespace jana::integration_tests::hierarchical_ownership { + +static constexpr uint64_t CANARY = 0xABCDEF0123456789ULL; + +struct BlockData { // inserted into the parent (timeslice) by the source + std::vector samples; + uint64_t canary = CANARY; +}; + +struct FrameRef { // child payload: pointer INTO the parent's BlockData, + const BlockData* block; // valid per the parent-lifetime guarantee + int frame_index; +}; + +struct TimesliceSource : JEventSource { + TimesliceSource() { + SetLevel(JEventLevel::Timeslice); + SetCallbackStyle(CallbackStyle::ExpertMode); + } + Result Emit(JEvent& event) override { + auto* block = new BlockData; + block->samples.assign(100000, event.GetEventNumber()); + event.Insert(block); + return Result::Success; + } +}; + +struct FrameUnfolder : JEventUnfolder { + FrameUnfolder() { + SetParentLevel(JEventLevel::Timeslice); + SetChildLevel(JEventLevel::PhysicsEvent); + } + Result Unfold(const JEvent& parent, JEvent& child, int child_idx) override { + child.Insert(new FrameRef{parent.GetSingle(), child_idx}); + // Three children per timeslice; the last one takes the NextChildNextParent + // branch, which pushes the parent to its pool in the same firing. + return child_idx == 2 ? Result::NextChildNextParent : Result::NextChildKeepParent; + } +}; + +struct FrameProcessor : JEventProcessor { + FrameProcessor() { SetCallbackStyle(CallbackStyle::ExpertMode); } + void ProcessSequential(const JEvent& event) override { + const auto* ref = event.GetSingle(); + REQUIRE(ref->block->canary == CANARY); + } +}; + +TEST_CASE("ParentEventOutlivesChildren") { + + JApplication app; + app.SetParameterValue("jana:nevents", 10); + app.SetParameterValue("nthreads", 1); + app.Add(new TimesliceSource); + app.Add(new FrameUnfolder); + app.Add(new FrameProcessor); + app.Run(); +} + +} // namespace From 75dac60d60a933a3790676271ac4496f62955d86 Mon Sep 17 00:00:00 2001 From: Nathan Brei Date: Mon, 24 Aug 2026 12:46:10 -0400 Subject: [PATCH 3/3] Quick fix --- src/programs/unit_tests/Components/PodioTests.cc | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/src/programs/unit_tests/Components/PodioTests.cc b/src/programs/unit_tests/Components/PodioTests.cc index 91a2439e9..caca8b123 100644 --- a/src/programs/unit_tests/Components/PodioTests.cc +++ b/src/programs/unit_tests/Components/PodioTests.cc @@ -187,7 +187,7 @@ TEST_CASE("PodioClearData_Test") { REQUIRE(clusters->at(0).Hits_size() == 1); REQUIRE(clusters->at(0).Hits().at(0).cellID() == 22); - event->Clear(true); + event->Clear(); ExampleHitCollection hits2; auto h2 = hits2.create(); @@ -218,7 +218,7 @@ TEST_CASE("PodioClearData_Test") { REQUIRE(clusters->at(0).Hits_size() == 1); REQUIRE(clusters->at(0).Hits().at(0).cellID() == 22); - event->Clear(true); + event->Clear(); ExampleHitCollection hits2; auto h2 = hits2.create(); @@ -297,7 +297,7 @@ TEST_CASE("PodioMultifactoryClearData_Test") { REQUIRE(other_clusters.size() == 1); - event->Clear(true); + event->Clear(); ExampleHitCollection hits2; auto h2 = hits2.create(); @@ -331,7 +331,7 @@ TEST_CASE("PodioMultifactoryClearData_Test") { REQUIRE(clusters->at(0).Hits_size() == 1); REQUIRE(clusters->at(0).Hits().at(0).cellID() == 22); - event->Clear(true); + event->Clear(); ExampleHitCollection hits2; auto h2 = hits2.create(); @@ -409,7 +409,7 @@ TEST_CASE("PodioTests_ExceptionInFactoryInit") { auto frame = event.Get().at(0); REQUIRE(frame->get("ExampleHit").size() == 0); - event.Clear(true); + event.Clear(); try { event.GetCollectionBase("ExampleHit"); @@ -440,7 +440,7 @@ TEST_CASE("PodioTests_ExceptionInFactoryInit") { auto frame = event.Get().at(0); REQUIRE(frame->get("ExampleHit").size() == 0); - event.Clear(true); + event.Clear(); try { event.GetCollectionBase("ExampleHit"); @@ -478,7 +478,7 @@ TEST_CASE("PodioTests_ExceptionInFactoryInit") { REQUIRE(frame->get("ExampleHit").size() == 1); std::cout << "Clearing event" << std::endl; - event.Clear(true); + event.Clear(); event.SetEventNumber(23); std::cout << "Calling failing factory" << std::endl;