Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
a314394
Add DataflowStatusModule, update token creation in DataWriterModule
eflumerf Jun 29, 2026
9d6ef93
Merge remote-tracking branch 'origin/develop' into eflumerf/DFOProtoc…
eflumerf Jun 29, 2026
942782d
Start working on the DFO Module
eflumerf Jul 6, 2026
aed8509
Merge remote-tracking branch 'origin/develop' into eflumerf/DFOProtoc…
eflumerf Jul 6, 2026
cb2aa32
Update DFOModule to use DataflowStatus-based protocol. Implement basic
eflumerf Jul 8, 2026
377ae27
Configuration updates for schema completeness
eflumerf Jul 9, 2026
4bb1e0f
Add TRBModule_test. Adapt TRBModule for TRBCompletion reporting. Track
eflumerf Jul 13, 2026
50f9857
Implement status watchdog thread in DFOModule, which watches for
eflumerf Jul 14, 2026
f9f9aeb
Merge remote-tracking branch 'origin/develop' into eflumerf/DFOProtoc…
eflumerf Aug 3, 2026
69fb72f
Fix issues with status locking scope. Format files
eflumerf Aug 3, 2026
9740d6c
Add initial version of dfo_protocol_test integration test
eflumerf Aug 3, 2026
953b9da
Fill in decision_destination in TriggerInhibit messages. Send a Trigg…
eflumerf Aug 4, 2026
c9fb938
Receive TriggerInhibit sent by start
eflumerf Aug 4, 2026
215976e
Clear assigned trigger decisions list after determining remnants
eflumerf Aug 18, 2026
0342628
Clarify usage of TriggerId to mean only trigger_number/run_number.
eflumerf Aug 19, 2026
d1984c9
Since DataflowStatusModule tracks sequences, DataWriterModule no longer
eflumerf Aug 20, 2026
6848b17
Clarify log messages, recognize that sequenced triggers can both buil…
eflumerf Aug 20, 2026
e7a2435
Make log messages more consistent
eflumerf Aug 20, 2026
d70f537
Don't notify while processing
eflumerf Aug 20, 2026
7f2ad96
Improve logging around is_busy
eflumerf Aug 20, 2026
712fdad
Add missing includes
eflumerf Aug 28, 2026
5adf9d5
More linting fixes
eflumerf Aug 28, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 20 additions & 11 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -22,20 +22,21 @@ find_package(Boost COMPONENTS iostreams unit_test_framework REQUIRED)
daq_protobuf_codegen( opmon/*.proto )

##############################################################################
daq_add_library( TriggerInhibitAgent.cpp TriggerRecordBuilderData.cpp TPBundleHandler.cpp
LINK_LIBRARIES
daq_add_library( TriggerInhibitAgent.cpp TPBundleHandler.cpp
LINK_LIBRARIES
opmonlib::opmonlib ers::ers HighFive appfwk::appfwk logging::logging stdc++fs dfmessages::dfmessages utilities::utilities trigger::trigger detdataformats::detdataformats trgdataformats::trgdataformats)


daq_add_plugin( HDF5DataStore duneDataStore LINK_LIBRARIES dfmodules hdf5libs::hdf5libs stdc++fs)

daq_add_plugin( FragmentAggregatorModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( DataWriterModule duneDAQModule LINK_LIBRARIES dfmodules hdf5libs::hdf5libs iomanager::iomanager )
daq_add_plugin( DFOModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( TRBModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( TRMonRequestorModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( FakeDataProdModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager)
daq_add_plugin( TPStreamWriterModule duneDAQModule LINK_LIBRARIES dfmodules hdf5libs::hdf5libs trigger::trigger Boost::iostreams )
daq_add_plugin( FragmentAggregatorModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( DataWriterModule duneDAQModule LINK_LIBRARIES dfmodules hdf5libs::hdf5libs iomanager::iomanager )
daq_add_plugin( DataflowStatusModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( DFOModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( TRBModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( TRMonRequestorModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( FakeDataProdModule duneDAQModule LINK_LIBRARIES dfmodules iomanager::iomanager )
daq_add_plugin( TPStreamWriterModule duneDAQModule LINK_LIBRARIES dfmodules hdf5libs::hdf5libs trigger::trigger Boost::iostreams )

##############################################################################
daq_add_unit_test( HDF5FileUtils_test LINK_LIBRARIES dfmodules )
Expand All @@ -46,9 +47,17 @@ add_dependencies( HDF5Write_test dfmodules_HDF5DataStore_duneDataStore )
daq_add_unit_test( DFOModule_test LINK_LIBRARIES dfmodules )
add_dependencies( DFOModule_test dfmodules_DFOModule_duneDAQModule)

daq_add_unit_test( TriggerRecordBuilderData_test LINK_LIBRARIES dfmodules)
daq_add_unit_test( DataflowStatusModule_test LINK_LIBRARIES dfmodules )
add_dependencies( DataflowStatusModule_test dfmodules_DataflowStatusModule_duneDAQModule)

daq_add_unit_test( DFOProtocol_test LINK_LIBRARIES dfmodules )
add_dependencies( DFOProtocol_test dfmodules_DFOModule_duneDAQModule dfmodules_DataflowStatusModule_duneDAQModule)

daq_add_unit_test( DataStoreFactory_test LINK_LIBRARIES dfmodules)

daq_add_unit_test( TRBModule_test LINK_LIBRARIES dfmodules )
add_dependencies( TRBModule_test dfmodules_TRBModule_duneDAQModule)

##############################################################################

daq_install()
247 changes: 247 additions & 0 deletions integtest/dfo_protocol_test.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,247 @@
"""
Integration Test for DFO Protocol

This test module validates DAQ system behavior while using multiple DFO applications.
It also verifies that the protocol correctly handles DFO and DF application crashes.
"""

import copy
import conffwk
import os
import pathlib
import pytest
import random
import string

import integrationtest.data_classes as data_classes
import integrationtest.data_file_checks as data_file_checks
import integrationtest.log_file_checks as log_file_checks
import integrationtest.resource_validation as resource_validation
import integrationtest.utility_functions as utility_functions
from integrationtest.get_pytest_tmpdir import get_pytest_tmpdir
from integrationtest.verbosity_helper import IntegtestVerbosityLevels

import functools

print = functools.partial(print, flush=True) # always flush print() output

pytest_plugins = "integrationtest.integrationtest_drunc"

# Run setup
run_duration = 30 # seconds
check_for_logfile_errors = True

# Default values for validation parameters
number_of_dataflow_apps = 3
number_of_data_producers = 4
number_of_readout_apps = 2
trigger_rate = 4.0
expected_number_of_data_files = number_of_dataflow_apps
check_for_logfile_errors = True
expected_event_count = run_duration * trigger_rate / number_of_dataflow_apps
expected_event_count_tolerance = expected_event_count / 10
ta_prescale = 1000

wibeth_frag_params = {
"fragment_type_description": "WIBEth",
"fragment_type": "WIBEth",
"expected_fragment_count": number_of_readout_apps * number_of_data_producers,
"min_size_bytes": 7272,
"max_size_bytes": 28872,
}
# sizes: 128 is for one TC with zero TAs inside it (72+56)
# 208 is for one TC with one TA inside it (72+56+80)
# 264 is for two TCs with one TA in one of them (72+56+80+56)
triggercandidate_frag_params = {
"fragment_type_description": "Trigger Candidate",
"fragment_type": "Trigger_Candidate",
"expected_fragment_count": 1,
"min_size_bytes": 128,
"max_size_bytes": 264,
"debug_mask": 0x0,
"frag_sizes_by_TC_type": {"kPrescale": {"min_size_bytes": 208, "max_size_bytes": 264},
"kRandom": {"min_size_bytes": 128, "max_size_bytes": 264},
"default": {"min_size_bytes": 128, "max_size_bytes": 264} }
}
# sizes: 72 is for an empty TP fragment
# 168 is for a fragment with four TPs in it (72+24+24+24+24)
triggerprimitive_frag_params = {
"fragment_type_description": "Trigger Primitive",
"fragment_type": "Trigger_Primitive",
"expected_fragment_count": 3 * number_of_readout_apps,
"min_size_bytes": 72,
"max_size_bytes": 168,
}
# 03-Jul-2025, KAB: changing the default max size from 72 to 100 to handle cases in which there
# was a Random or Prescale trigger along with a coincidental HSI event within the readout window.
hsi_frag_params = {
"fragment_type_description": "HSI",
"fragment_type": "Hardware_Signal",
"expected_fragment_count": 1,
"min_size_bytes": 72,
"max_size_bytes": 100,
"frag_sizes_by_TC_type": {"kTiming": {"min_size_bytes": 100, "max_size_bytes": 100},
"default": {"min_size_bytes": 72, "max_size_bytes": 100} }
}
ignored_logfile_problems = {
"-controller": [
],
"local-connection-server": [
"errorlog: -",
],
# 04-Mar-2026, KAB: added the absl::InitializeLog warning message to the ignored list for
# all DAQ processes, given that we currently don't have a way suppress it at its source.
r".*": [
r"WARNING: All log messages before absl::InitializeLog\(\) is called are written to STDERR"
]
}

# Determine if this computer has enough resources for these tests
resource_validator = resource_validation.ResourceValidator()
resource_validator.cpu_count_needs(
15, 30
) # 3 for each data source (incl TPG) plus 3 more for everything else
resource_validator.free_memory_needs(
9, 14
) # 30% more than what we observe being used ('free -h')
actual_output_path = get_pytest_tmpdir()
resource_validator.free_disk_space_needs(
actual_output_path, 1
) # more than what we observe

### Config setup
common_config_obj = data_classes.integtest_params_for_predefined_dunedaq_config()
common_config_obj.op_env = "test"
common_config_obj.predefined_config_db ="config/daqsystemtest/example-configs.data.xml"

common_config_obj.config_substitutions.append(
data_classes.attribute_substitution(
obj_class="TCDataProcessor", # 12-Nov-2025, KAB: turned off the merging of
obj_id="def-tc-processor", # overlapping TCs so that we get more consistent
updates={ # numbers of TriggerRecords in the output files.
"merge_overlapping_tcs": False
},)
)

# Get default config
multidfo_local_conf = copy.deepcopy(common_config_obj)
multidfo_local_conf.config_session_name = "local-multidfo-config"

# Prep configs
stopped_dfo_conf = copy.deepcopy(multidfo_local_conf)
killed_dfo_conf = copy.deepcopy(multidfo_local_conf)
stopped_df_conf = copy.deepcopy(multidfo_local_conf)
killed_df_conf = copy.deepcopy(multidfo_local_conf)

stopped_dfo_conf.system_signal_configs = [
data_classes.system_signal_config(
application_label="dfo-02", signal=data_classes.PosixSignal.SIGSTOP, delay_s=25
),
data_classes.system_signal_config(
application_label="dfo-02", signal=data_classes.PosixSignal.SIGCONT, delay_s=30
),
]


killed_dfo_conf.system_signal_configs = [
data_classes.system_signal_config(
application_label="dfo-02", signal=data_classes.PosixSignal.SIGKILL, delay_s=25
),
]

stopped_df_conf.system_signal_configs = [
data_classes.system_signal_config(
application_label="df-02", signal=data_classes.PosixSignal.SIGSTOP, delay_s=25
),
data_classes.system_signal_config(
application_label="df-02", signal=data_classes.PosixSignal.SIGCONT, delay_s=30
),
]

killed_df_conf.system_signal_configs = [
data_classes.system_signal_config(
application_label="df-02", signal=data_classes.PosixSignal.SIGKILL, delay_s=25
),
]

# Finally store configs in map
confgen_arguments = {
"default": multidfo_local_conf,
"stopped-dfo": stopped_dfo_conf,
"killed-dfo": killed_dfo_conf,
"stopped-df": stopped_df_conf,
"killed-df": killed_df_conf,
}

# The commands to run in dunerc, as a list
dunerc_command_list = (
"boot wait 2 conf start --run-number 101 wait 1 enable-triggers wait ".split()
+ [str(run_duration)]
+ "disable-triggers wait 2 drain-dataflow wait 2 stop-trigger-sources stop scrap terminate".split()
)


### Tests

def test_dunerc_success(run_dunerc, caplog):
# checks for run control success, problems during pytest setup, etc.
utility_functions.basic_checks(run_dunerc, caplog, print_test_name=True)


def test_log_files(run_dunerc):
if check_for_logfile_errors:
# Check that there are no warnings or errors in the log files
assert log_file_checks.logs_are_error_free(
run_dunerc.log_files, True, True, ignored_logfile_problems,
verbosity_helper=run_dunerc.verbosity_helper
)


def test_data_files(run_dunerc):
local_expected_event_count = expected_event_count
local_event_count_tolerance = expected_event_count_tolerance
low_number_of_files = expected_number_of_data_files
high_number_of_files = expected_number_of_data_files
fragment_check_list = [triggercandidate_frag_params, hsi_frag_params, wibeth_frag_params]

local_expected_event_count += (
(6250.0 / ta_prescale)
* number_of_data_producers
* number_of_readout_apps
* run_duration
/ (100 * number_of_dataflow_apps)
)
local_event_count_tolerance += (
(250.0 / ta_prescale)
* number_of_data_producers
* number_of_readout_apps
* run_duration
/ (100 * number_of_dataflow_apps)
)
fragment_check_list.append(triggerprimitive_frag_params)
nontrig_fragment_check_list = [hsi_frag_params, wibeth_frag_params]

# Run some tests on the output data file
assert (
len(run_dunerc.data_files) == high_number_of_files
or len(run_dunerc.data_files) == low_number_of_files
)

all_ok = True
for idx in range(len(run_dunerc.data_files)):
data_file = data_file_checks.DataFile(run_dunerc.data_files[idx], run_dunerc.verbosity_helper)
all_ok &= data_file_checks.sanity_check(data_file)
all_ok &= data_file_checks.check_file_attributes(data_file)
all_ok &= data_file_checks.check_event_count(
data_file, local_expected_event_count, local_event_count_tolerance
)
for jdx in range(len(fragment_check_list)):
all_ok &= data_file_checks.check_fragment_count(
data_file, fragment_check_list[jdx]
)
all_ok &= data_file_checks.check_fragment_sizes(
data_file, fragment_check_list[jdx]
)
for kdx in range(len(nontrig_fragment_check_list)):
all_ok &= data_file_checks.check_fragment_error_flags( data_file, nontrig_fragment_check_list[kdx])
assert all_ok
Loading
Loading