From 0f6e40378f30245a9877cabaa452f6a53d3d52d5 Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Mon, 14 Sep 2026 11:50:24 -0500 Subject: [PATCH 01/11] first round of changes to help improve the robustness of process startup in integrationtests --- src/integrationtest/async_proc_mgmt.py | 30 +++++++++---------- src/integrationtest/data_classes.py | 31 ++++++++++++-------- src/integrationtest/integrationtest_drunc.py | 4 +-- 3 files changed, 36 insertions(+), 29 deletions(-) diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index 14c9ad4..372ef7d 100644 --- a/src/integrationtest/async_proc_mgmt.py +++ b/src/integrationtest/async_proc_mgmt.py @@ -7,16 +7,13 @@ from integrationtest.data_classes import * from integrationtest.verbosity_helper import * from datetime import datetime, timezone -from typing import Final import functools print = functools.partial(print, flush=True) # always flush print() output -PROCESS_ECHO_STRING: Final[str] = "*** COMMAND HAS COMPLETED ***" - async def read_stream(stream, process_name, app_exe_name, print_proc_name, run_dir, - shared_data: CommandProcessingSharedData, verbosity_level): + shared_data: OutputMonitoringSharedData, verbosity_level): """Asynchronously reads lines from a stream and processes them immediately.""" full_output = "" observed_command_prompt = "" @@ -111,8 +108,8 @@ async def read_stream(stream, process_name, app_exe_name, print_proc_name, run_d return full_output -async def wait_for_console_output_lull(start_time, wait_params: CommandWaitParameters, - shared_data: CommandProcessingSharedData): +async def wait_for_console_output_lull(start_time, wait_params: ConsoleOutputWaitParameters, + shared_data: OutputMonitoringSharedData): now = time.time() while True: async with shared_data.lock: @@ -126,7 +123,7 @@ async def wait_for_console_output_lull(start_time, wait_params: CommandWaitParam now = time.time() -async def send_commands(target_proc_info, proc_name, shared_data: CommandProcessingSharedData, +async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitoringSharedData, cmd_list, wait_params, verbosity_level): target_proc = target_proc_info.process if target_proc.returncode is not None: # Check if process is still running @@ -148,9 +145,9 @@ async def send_commands(target_proc_info, proc_name, shared_data: CommandProcess print(".", end="") # wait for the command(s) to finish, if requested - if not wait_params.wait_for_command_completion: + if wait_params is None: return - if wait_params.style == CommandWaitStyle.TIME_PLUS_EXIT: + if isinstance(wait_params, ProcessExitWaitParameters): # The idea behind this command style is that we want to wait until the process has # exited and we want to support 'exit' timeout values that are not long and arbitrary. # In order to do that, we wait for a lull in the console output before waiting @@ -160,7 +157,7 @@ async def send_commands(target_proc_info, proc_name, shared_data: CommandProcess # waiting for the process to respond to it. But, we tell users that we skipped it. await wait_for_console_output_lull(cmd_start_time, wait_params, shared_data) if "exit" in target_proc_info.supported_commands: - sleep_interval: float = wait_params.timeout_waiting_for_exit / 10 + sleep_interval: float = wait_params.wait_time_after_last_msg / 10 for idx in range(10): if target_proc.returncode is not None: break @@ -172,7 +169,7 @@ async def send_commands(target_proc_info, proc_name, shared_data: CommandProcess if verbosity_level >= IntegtestVerbosityLevels.integtest_debug: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] The {proc_name} process doesn't support the 'exit' command, so waiting for exit was skipped") - elif wait_params.style == CommandWaitStyle.ECHO: + elif isinstance(wait_params, EchoCommandWaitParameters): if "echo" in target_proc_info.supported_commands: shared_data.cmd_cmplt_evt.clear() target_proc.stdin.write((f"echo '{PROCESS_ECHO_STRING}'\n").encode()) @@ -186,7 +183,7 @@ async def send_commands(target_proc_info, proc_name, shared_data: CommandProcess now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] The {proc_name} process doesn't support the 'echo' command, using TIME wait instead'") await wait_for_console_output_lull(cmd_start_time, wait_params, shared_data) - else: # treat everything else as wait_params.style == CommandWaitStyle.TIME: + else: await wait_for_console_output_lull(cmd_start_time, wait_params, shared_data) @@ -196,7 +193,7 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r tasks = {} command_completion_event = asyncio.Event() proc_results = {} - shared_data: CommandProcessingSharedData = CommandProcessingSharedData() + shared_data: OutputMonitoringSharedData = OutputMonitoringSharedData() # 1. Start all subprocesses for session_app in daq_session_ingredients.applications: @@ -237,7 +234,7 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r # determine the supported commands for each app (using the 'help' command) help_cmd = ["help"] - help_cmd_wait_params = CommandWaitParameters(timeout_waiting_for_first_msg=2) + help_cmd_wait_params = ConsoleOutputWaitParameters(timeout_waiting_for_first_msg=2) await wait_for_console_output_lull(time.time(), help_cmd_wait_params, shared_data) for proc_name, proc_info in processes.items(): async with shared_data.lock: @@ -302,8 +299,11 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r # so, we create a new list from the iterator.) reformatted_cmd_list = list(reversed(working_cmd_list)) + wait_params: ConsoleOutputWaitParameters = None + if cmd_set.wait_for_command_completion: + wait_params = cmd_set.wait_params await send_commands(proc_info, target, shared_data, reformatted_cmd_list, - cmd_set.wait_params, verbosity_level) + wait_params, verbosity_level) else: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") diff --git a/src/integrationtest/data_classes.py b/src/integrationtest/data_classes.py index b5da940..db71b67 100644 --- a/src/integrationtest/data_classes.py +++ b/src/integrationtest/data_classes.py @@ -1,6 +1,9 @@ from dataclasses import dataclass, field from enum import Enum import asyncio +from typing import Final + +PROCESS_ECHO_STRING: Final[str] = "*** COMMAND HAS COMPLETED ***" @dataclass class DROMap_config: @@ -129,19 +132,22 @@ class CreateConfigResult: trmon_data_dirs: list[str] -class CommandWaitStyle(Enum): - ECHO = "echo" - TIME = "time" - TIME_PLUS_EXIT = "time_plus_exit" - NONE = "none" - @dataclass -class CommandWaitParameters: - wait_for_command_completion: bool = True - style: CommandWaitStyle = CommandWaitStyle.TIME +class ConsoleOutputWaitParameters: timeout_waiting_for_first_msg: int = 2 # seconds wait_time_after_last_msg: int = 2 # seconds - timeout_waiting_for_exit: int = 5 # seconds + +@dataclass +class SearchPhraseWaitParameters(ConsoleOutputWaitParameters): + search_phrase: str = None + +@dataclass +class EchoCommandWaitParameters(ConsoleOutputWaitParameters): + search_phrase: str = PROCESS_ECHO_STRING + +@dataclass +class ProcessExitWaitParameters(ConsoleOutputWaitParameters): + process: asyncio.subprocess.Process = None @dataclass class DAQControlApplication: @@ -153,7 +159,8 @@ class DAQControlApplication: class DAQCommandSet: target: str command_list: list[str] - wait_params: CommandWaitParameters = field(default_factory=lambda: CommandWaitParameters()) + wait_params: ConsoleOutputWaitParameters + wait_for_command_completion: bool = True @dataclass class DAQSessionIngredients: @@ -166,7 +173,7 @@ class RunningProcessInfo: supported_commands: list[str] = field(default_factory=list) @dataclass -class CommandProcessingSharedData: +class OutputMonitoringSharedData: lock: asyncio.Lock = field(default_factory=asyncio.Lock, repr=False) cmd_cmplt_evt: asyncio.Event = field(default_factory=asyncio.Event, repr=False) last_msg_time: int = 0 diff --git a/src/integrationtest/integrationtest_drunc.py b/src/integrationtest/integrationtest_drunc.py index a6e7bdd..9683119 100644 --- a/src/integrationtest/integrationtest_drunc.py +++ b/src/integrationtest/integrationtest_drunc.py @@ -598,7 +598,7 @@ class RunResult: if verbosity_level < IntegtestVerbosityLevels.full_output: sys.stdout = original_stdout - exit_cmd = DAQCommandSet("drunc", [ "exit" ], CommandWaitParameters(style=CommandWaitStyle.TIME_PLUS_EXIT)) + exit_cmd = DAQCommandSet("drunc", [ "exit" ], ProcessExitWaitParameters()) if user_supplied_apps: dsi = copy.deepcopy(run_control_commands) for app in dsi.applications: @@ -647,7 +647,7 @@ class RunResult: dsapp = DAQControlApplication("drunc", popen_command_list) - requested_cmds = DAQCommandSet("drunc", run_control_commands, CommandWaitParameters(style=CommandWaitStyle.ECHO)) + requested_cmds = DAQCommandSet("drunc", run_control_commands, EchoCommandWaitParameters()) app_list = [ dsapp ] cmd_set_list = [ requested_cmds, exit_cmd ] From 993caf199e9d0963be764d7b4700f944ae0a5ff7 Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Mon, 14 Sep 2026 15:02:39 -0500 Subject: [PATCH 02/11] second round of changes to help improve the robustness of process startup in integrationtests --- src/integrationtest/async_proc_mgmt.py | 65 +++++++++++++------- src/integrationtest/data_classes.py | 27 +++++--- src/integrationtest/integrationtest_drunc.py | 3 +- 3 files changed, 62 insertions(+), 33 deletions(-) diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index 372ef7d..39a54e3 100644 --- a/src/integrationtest/async_proc_mgmt.py +++ b/src/integrationtest/async_proc_mgmt.py @@ -28,11 +28,12 @@ async def read_stream(stream, process_name, app_exe_name, print_proc_name, run_d async with shared_data.lock: shared_data.last_msg_time = time.time() - # if the special end-of-command string has been echo-ed by the process, - # send the relevant signal to any waiting task by setting the completion event - if PROCESS_ECHO_STRING in decoded_line: - shared_data.cmd_cmplt_evt.set() - continue + # if we find a requested phrase in the output, set the relevant flag + async with shared_data.lock: + if shared_data.phrase_searching_in_progress: + clean_line = re.sub(r"\x1b\[[0-9;]*m", "", decoded_line) + if shared_data.search_phrase in clean_line: + shared_data.phrase_has_been_found = True # process the output of the "help" command, if requested async with shared_data.lock: @@ -108,8 +109,17 @@ async def read_stream(stream, process_name, app_exe_name, print_proc_name, run_d return full_output -async def wait_for_console_output_lull(start_time, wait_params: ConsoleOutputWaitParameters, +async def wait_for_requested_condition(start_time, wait_params: ConditionalWaitParameters, shared_data: OutputMonitoringSharedData): + + if (isinstance(wait_params, EchoCommandWaitParameters) or \ + isinstance(wait_params, KeyPhraseWaitParameters)) and \ + wait_params.search_phrase is not None: + async with shared_data.lock: + shared_data.phrase_has_been_found = False + shared_data.search_phrase = wait_params.search_phrase + shared_data.phrase_searching_in_progress = True + now = time.time() while True: async with shared_data.lock: @@ -119,9 +129,21 @@ async def wait_for_console_output_lull(start_time, wait_params: ConsoleOutputWai else: if now - shared_data.last_msg_time >= wait_params.wait_time_after_last_msg: break + if shared_data.phrase_has_been_found: + break + if isinstance(wait_params, ProcessExitWaitParameters): + if wait_params.process.returncode is not None: + break await asyncio.sleep(0.25) now = time.time() + if (isinstance(wait_params, EchoCommandWaitParameters) or \ + isinstance(wait_params, KeyPhraseWaitParameters)) and \ + wait_params.search_phrase is not None: + async with shared_data.lock: + shared_data.phrase_searching_in_progress = False + shared_data.phrase_has_been_found = False + async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitoringSharedData, cmd_list, wait_params, verbosity_level): @@ -155,13 +177,9 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori # finishing of the console output. # Of course, if the app doesn't support the "exit" command, there is no sense in # waiting for the process to respond to it. But, we tell users that we skipped it. - await wait_for_console_output_lull(cmd_start_time, wait_params, shared_data) if "exit" in target_proc_info.supported_commands: - sleep_interval: float = wait_params.wait_time_after_last_msg / 10 - for idx in range(10): - if target_proc.returncode is not None: - break - await asyncio.sleep(sleep_interval) + wait_params.process = target_proc + await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) if target_proc.returncode is None: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] WARNING: timeout waiting for {proc_name} to exit in response to {cmd_list}") @@ -171,27 +189,28 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori print(f"[integtest_proc_mgmt {now_string}] The {proc_name} process doesn't support the 'exit' command, so waiting for exit was skipped") elif isinstance(wait_params, EchoCommandWaitParameters): if "echo" in target_proc_info.supported_commands: - shared_data.cmd_cmplt_evt.clear() - target_proc.stdin.write((f"echo '{PROCESS_ECHO_STRING}'\n").encode()) + target_proc.stdin.write((f"echo '{wait_params.search_phrase}'\n").encode()) await target_proc.stdin.drain() if verbosity_level >= IntegtestVerbosityLevels.integtest_debug: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") - print(f"[integtest_proc_mgmt {now_string}] Sent command to {proc_name}: echo '{PROCESS_ECHO_STRING}'") - await shared_data.cmd_cmplt_evt.wait() - shared_data.cmd_cmplt_evt.clear() + print(f"[integtest_proc_mgmt {now_string}] Sent command to {proc_name}: echo '{wait_params.search_phrase}'") + await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) else: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] The {proc_name} process doesn't support the 'echo' command, using TIME wait instead'") - await wait_for_console_output_lull(cmd_start_time, wait_params, shared_data) + wait_params.search_phrase = None + await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) else: - await wait_for_console_output_lull(cmd_start_time, wait_params, shared_data) + await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) + +# add handling of KeyPhrase +# add background task? async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, run_dir, verbosity_level): processes = {} tasks = {} - command_completion_event = asyncio.Event() proc_results = {} shared_data: OutputMonitoringSharedData = OutputMonitoringSharedData() @@ -220,7 +239,7 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r run_dir, shared_data, verbosity_level )) - time.sleep(session_app.wait_time_after_start) + await wait_for_requested_condition(time.time(), session_app.startup_wait_params, shared_data) if verbosity_level >= IntegtestVerbosityLevels.integtest_debug: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") @@ -235,7 +254,7 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r # determine the supported commands for each app (using the 'help' command) help_cmd = ["help"] help_cmd_wait_params = ConsoleOutputWaitParameters(timeout_waiting_for_first_msg=2) - await wait_for_console_output_lull(time.time(), help_cmd_wait_params, shared_data) + #await wait_for_requested_condition(time.time(), help_cmd_wait_params, shared_data) for proc_name, proc_info in processes.items(): async with shared_data.lock: shared_data.results_of_parsing_help_output = [] @@ -299,7 +318,7 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r # so, we create a new list from the iterator.) reformatted_cmd_list = list(reversed(working_cmd_list)) - wait_params: ConsoleOutputWaitParameters = None + wait_params: ConditionalWaitParameters = None if cmd_set.wait_for_command_completion: wait_params = cmd_set.wait_params await send_commands(proc_info, target, shared_data, reformatted_cmd_list, diff --git a/src/integrationtest/data_classes.py b/src/integrationtest/data_classes.py index db71b67..0930381 100644 --- a/src/integrationtest/data_classes.py +++ b/src/integrationtest/data_classes.py @@ -1,9 +1,6 @@ from dataclasses import dataclass, field from enum import Enum import asyncio -from typing import Final - -PROCESS_ECHO_STRING: Final[str] = "*** COMMAND HAS COMPLETED ***" @dataclass class DROMap_config: @@ -133,33 +130,43 @@ class CreateConfigResult: @dataclass -class ConsoleOutputWaitParameters: +class ConditionalWaitParameters: + pass + +@dataclass +class ConsoleOutputWaitParameters(ConditionalWaitParameters): timeout_waiting_for_first_msg: int = 2 # seconds wait_time_after_last_msg: int = 2 # seconds @dataclass -class SearchPhraseWaitParameters(ConsoleOutputWaitParameters): +class KeyPhraseWaitParameters(ConsoleOutputWaitParameters): + timeout_waiting_for_first_msg: int = 30 # seconds + wait_time_after_last_msg: int = 30 # seconds search_phrase: str = None @dataclass class EchoCommandWaitParameters(ConsoleOutputWaitParameters): - search_phrase: str = PROCESS_ECHO_STRING + timeout_waiting_for_first_msg: int = 20 # seconds + wait_time_after_last_msg: int = 20 # seconds + search_phrase: str = "*** COMMAND HAS COMPLETED ***" @dataclass class ProcessExitWaitParameters(ConsoleOutputWaitParameters): + timeout_waiting_for_first_msg: int = 10 # seconds + wait_time_after_last_msg: int = 10 # seconds process: asyncio.subprocess.Process = None @dataclass class DAQControlApplication: alias: str startup_strings: list[str] - wait_time_after_start: int = 2 # seconds + startup_wait_params: ConditionalWaitParameters @dataclass class DAQCommandSet: target: str command_list: list[str] - wait_params: ConsoleOutputWaitParameters + wait_params: ConditionalWaitParameters wait_for_command_completion: bool = True @dataclass @@ -175,8 +182,10 @@ class RunningProcessInfo: @dataclass class OutputMonitoringSharedData: lock: asyncio.Lock = field(default_factory=asyncio.Lock, repr=False) - cmd_cmplt_evt: asyncio.Event = field(default_factory=asyncio.Event, repr=False) last_msg_time: int = 0 number_of_lines_printed_to_the_console: int = 0 + search_phrase: str = "nullnullnull" + phrase_searching_in_progress: bool = False + phrase_has_been_found: bool = False parsing_of_help_output_in_progress: bool = False results_of_parsing_help_output: list[str] = field(default_factory=list) diff --git a/src/integrationtest/integrationtest_drunc.py b/src/integrationtest/integrationtest_drunc.py index 9683119..8944cab 100644 --- a/src/integrationtest/integrationtest_drunc.py +++ b/src/integrationtest/integrationtest_drunc.py @@ -645,7 +645,8 @@ class RunResult: + [str(create_config_files.integtest_params.config_session_name)] \ + [str(create_config_files.integtest_params.daq_session_name)] - dsapp = DAQControlApplication("drunc", popen_command_list) + dsapp = DAQControlApplication("drunc", popen_command_list, + KeyPhraseWaitParameters(search_phrase="unified_shell ready")) requested_cmds = DAQCommandSet("drunc", run_control_commands, EchoCommandWaitParameters()) From e19bb41c93dae80dbf57add9b0dda50260058571 Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Mon, 14 Sep 2026 20:07:51 -0500 Subject: [PATCH 03/11] some additional changes to help improve the robustness of process startup in integrationtests --- src/integrationtest/async_proc_mgmt.py | 12 ++++++------ src/integrationtest/data_classes.py | 4 ++-- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index 39a54e3..bd7da05 100644 --- a/src/integrationtest/async_proc_mgmt.py +++ b/src/integrationtest/async_proc_mgmt.py @@ -33,7 +33,7 @@ async def read_stream(stream, process_name, app_exe_name, print_proc_name, run_d if shared_data.phrase_searching_in_progress: clean_line = re.sub(r"\x1b\[[0-9;]*m", "", decoded_line) if shared_data.search_phrase in clean_line: - shared_data.phrase_has_been_found = True + shared_data.search_phrase_has_been_found = True # process the output of the "help" command, if requested async with shared_data.lock: @@ -116,7 +116,7 @@ async def wait_for_requested_condition(start_time, wait_params: ConditionalWaitP isinstance(wait_params, KeyPhraseWaitParameters)) and \ wait_params.search_phrase is not None: async with shared_data.lock: - shared_data.phrase_has_been_found = False + shared_data.search_phrase_has_been_found = False shared_data.search_phrase = wait_params.search_phrase shared_data.phrase_searching_in_progress = True @@ -129,7 +129,7 @@ async def wait_for_requested_condition(start_time, wait_params: ConditionalWaitP else: if now - shared_data.last_msg_time >= wait_params.wait_time_after_last_msg: break - if shared_data.phrase_has_been_found: + if shared_data.search_phrase_has_been_found: break if isinstance(wait_params, ProcessExitWaitParameters): if wait_params.process.returncode is not None: @@ -142,7 +142,7 @@ async def wait_for_requested_condition(start_time, wait_params: ConditionalWaitP wait_params.search_phrase is not None: async with shared_data.lock: shared_data.phrase_searching_in_progress = False - shared_data.phrase_has_been_found = False + shared_data.search_phrase_has_been_found = False async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitoringSharedData, @@ -201,10 +201,10 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori wait_params.search_phrase = None await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) else: + # KeyPhrase gets handled automatically here, along with ConsoleOutput await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) -# add handling of KeyPhrase -# add background task? +# add background task to avoid race condition? async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, run_dir, diff --git a/src/integrationtest/data_classes.py b/src/integrationtest/data_classes.py index 0930381..9ed8464 100644 --- a/src/integrationtest/data_classes.py +++ b/src/integrationtest/data_classes.py @@ -166,7 +166,7 @@ class DAQControlApplication: class DAQCommandSet: target: str command_list: list[str] - wait_params: ConditionalWaitParameters + wait_params: ConditionalWaitParameters = None wait_for_command_completion: bool = True @dataclass @@ -186,6 +186,6 @@ class OutputMonitoringSharedData: number_of_lines_printed_to_the_console: int = 0 search_phrase: str = "nullnullnull" phrase_searching_in_progress: bool = False - phrase_has_been_found: bool = False + search_phrase_has_been_found: bool = False parsing_of_help_output_in_progress: bool = False results_of_parsing_help_output: list[str] = field(default_factory=list) From 945c0256c8e95964be1797d27ae7c38a588e6b65 Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Tue, 15 Sep 2026 14:13:07 -0500 Subject: [PATCH 04/11] minor updates to the changes to help improve the robustness of process startup in integrationtests --- src/integrationtest/async_proc_mgmt.py | 27 +++++++++++++++++++------- src/integrationtest/data_classes.py | 22 +++++++++------------ 2 files changed, 29 insertions(+), 20 deletions(-) diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index bd7da05..578e9be 100644 --- a/src/integrationtest/async_proc_mgmt.py +++ b/src/integrationtest/async_proc_mgmt.py @@ -30,7 +30,8 @@ async def read_stream(stream, process_name, app_exe_name, print_proc_name, run_d # if we find a requested phrase in the output, set the relevant flag async with shared_data.lock: - if shared_data.phrase_searching_in_progress: + if shared_data.phrase_searching_in_progress and \ + shared_data.search_phrase is not None: clean_line = re.sub(r"\x1b\[[0-9;]*m", "", decoded_line) if shared_data.search_phrase in clean_line: shared_data.search_phrase_has_been_found = True @@ -109,7 +110,13 @@ async def read_stream(stream, process_name, app_exe_name, print_proc_name, run_d return full_output -async def wait_for_requested_condition(start_time, wait_params: ConditionalWaitParameters, +# The purpose of this function is to wait until one of the requested conditions has been +# satified. The conditions are specified in "wait parameter" objects. The baseline wait +# parameter class provides time values that are used to watch for quiet times in the console +# output from the control process(es). Wait parameter classes that build on the +# ConsoleOutputWaitParameters class add conditions that allow the waiting to end earlier +# than the wait times specified in the base class. +async def wait_for_requested_condition(start_time, wait_params: ConsoleOutputWaitParameters, shared_data: OutputMonitoringSharedData): if (isinstance(wait_params, EchoCommandWaitParameters) or \ @@ -131,7 +138,8 @@ async def wait_for_requested_condition(start_time, wait_params: ConditionalWaitP break if shared_data.search_phrase_has_been_found: break - if isinstance(wait_params, ProcessExitWaitParameters): + if isinstance(wait_params, ProcessExitWaitParameters) and \ + wait_params.process is not None: if wait_params.process.returncode is not None: break await asyncio.sleep(0.25) @@ -197,8 +205,8 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) else: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") - print(f"[integtest_proc_mgmt {now_string}] The {proc_name} process doesn't support the 'echo' command, using TIME wait instead'") - wait_params.search_phrase = None + print(f"[integtest_proc_mgmt {now_string}] WARNING: The {proc_name} process doesn't support the 'echo' command, using TIME wait instead'") + wait_params = KeyPhraseWaitParameters() await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) else: # KeyPhrase gets handled automatically here, along with ConsoleOutput @@ -239,7 +247,12 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r run_dir, shared_data, verbosity_level )) - await wait_for_requested_condition(time.time(), session_app.startup_wait_params, shared_data) + # Wait for the process to start up. If the user has not specified wait parameters, + # default-construct ones that make use of the console output. + wait_params = session_app.startup_wait_params + if wait_params is None: + wait_params = ConsoleOutputWaitParams() + await wait_for_requested_condition(time.time(), wait_params, shared_data) if verbosity_level >= IntegtestVerbosityLevels.integtest_debug: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") @@ -318,7 +331,7 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r # so, we create a new list from the iterator.) reformatted_cmd_list = list(reversed(working_cmd_list)) - wait_params: ConditionalWaitParameters = None + wait_params: ConsoleOutputWaitParameters = None if cmd_set.wait_for_command_completion: wait_params = cmd_set.wait_params await send_commands(proc_info, target, shared_data, reformatted_cmd_list, diff --git a/src/integrationtest/data_classes.py b/src/integrationtest/data_classes.py index 9ed8464..d0301f9 100644 --- a/src/integrationtest/data_classes.py +++ b/src/integrationtest/data_classes.py @@ -130,43 +130,39 @@ class CreateConfigResult: @dataclass -class ConditionalWaitParameters: - pass - -@dataclass -class ConsoleOutputWaitParameters(ConditionalWaitParameters): +class ConsoleOutputWaitParameters: timeout_waiting_for_first_msg: int = 2 # seconds wait_time_after_last_msg: int = 2 # seconds @dataclass class KeyPhraseWaitParameters(ConsoleOutputWaitParameters): - timeout_waiting_for_first_msg: int = 30 # seconds - wait_time_after_last_msg: int = 30 # seconds + timeout_waiting_for_first_msg: int = 60 # seconds + wait_time_after_last_msg: int = 60 # seconds search_phrase: str = None @dataclass class EchoCommandWaitParameters(ConsoleOutputWaitParameters): - timeout_waiting_for_first_msg: int = 20 # seconds - wait_time_after_last_msg: int = 20 # seconds + timeout_waiting_for_first_msg: int = 999999 # seconds + wait_time_after_last_msg: int = 999999 # seconds search_phrase: str = "*** COMMAND HAS COMPLETED ***" @dataclass class ProcessExitWaitParameters(ConsoleOutputWaitParameters): - timeout_waiting_for_first_msg: int = 10 # seconds - wait_time_after_last_msg: int = 10 # seconds + timeout_waiting_for_first_msg: int = 30 # seconds + wait_time_after_last_msg: int = 30 # seconds process: asyncio.subprocess.Process = None @dataclass class DAQControlApplication: alias: str startup_strings: list[str] - startup_wait_params: ConditionalWaitParameters + startup_wait_params: ConsoleOutputWaitParameters = None @dataclass class DAQCommandSet: target: str command_list: list[str] - wait_params: ConditionalWaitParameters = None + wait_params: ConsoleOutputWaitParameters = None wait_for_command_completion: bool = True @dataclass From 86717bda589efcdce17ba80c9cbaeb60c7e8beec Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Thu, 17 Sep 2026 17:37:45 +0200 Subject: [PATCH 05/11] minor updates --- src/integrationtest/async_proc_mgmt.py | 2 +- src/integrationtest/data_classes.py | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index 578e9be..4b753db 100644 --- a/src/integrationtest/async_proc_mgmt.py +++ b/src/integrationtest/async_proc_mgmt.py @@ -251,7 +251,7 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r # default-construct ones that make use of the console output. wait_params = session_app.startup_wait_params if wait_params is None: - wait_params = ConsoleOutputWaitParams() + wait_params = ConsoleOutputWaitParameters() await wait_for_requested_condition(time.time(), wait_params, shared_data) if verbosity_level >= IntegtestVerbosityLevels.integtest_debug: diff --git a/src/integrationtest/data_classes.py b/src/integrationtest/data_classes.py index d0301f9..ac41d9c 100644 --- a/src/integrationtest/data_classes.py +++ b/src/integrationtest/data_classes.py @@ -136,8 +136,8 @@ class ConsoleOutputWaitParameters: @dataclass class KeyPhraseWaitParameters(ConsoleOutputWaitParameters): - timeout_waiting_for_first_msg: int = 60 # seconds - wait_time_after_last_msg: int = 60 # seconds + timeout_waiting_for_first_msg: int = 30 # seconds + wait_time_after_last_msg: int = 30 # seconds search_phrase: str = None @dataclass From 5d89a7ed183d7a8ad9d6295a8f95fe525c7fdb97 Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Thu, 17 Sep 2026 21:09:44 +0200 Subject: [PATCH 06/11] Run some of the wait conditions in the background to avoid race conditions --- src/integrationtest/async_proc_mgmt.py | 38 +++++++++++++++++++------- 1 file changed, 28 insertions(+), 10 deletions(-) diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index 4b753db..fe138eb 100644 --- a/src/integrationtest/async_proc_mgmt.py +++ b/src/integrationtest/async_proc_mgmt.py @@ -155,14 +155,22 @@ async def wait_for_requested_condition(start_time, wait_params: ConsoleOutputWai async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitoringSharedData, cmd_list, wait_params, verbosity_level): + # Check if the process is still running; return if not target_proc = target_proc_info.process - if target_proc.returncode is not None: # Check if process is still running + if target_proc.returncode is not None: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] Error: {proc_name} has already exited, unable to send \"{cmd_list}\".") return + cmd_start_time = time.time() + + # for KeyPhrase wait conditions, start watching the console output *before* we send + # the requested commands (to avoid a race condition) + key_phrase_bg_task = None + if wait_params is not None and isinstance(wait_params, KeyPhraseWaitParameters) and \ + wait_params.search_phrase is not None: + key_phrase_bg_task = asyncio.create_task(wait_for_requested_condition(cmd_start_time, wait_params, shared_data)) # send the requested commands - cmd_start_time = time.time() for cmd in cmd_list: target_proc.stdin.write((cmd + "\n").encode()) await target_proc.stdin.drain() @@ -170,6 +178,8 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] Sent command to {proc_name}: {cmd}") else: + # in order to indicate to the user that the program is not stalled, + # we print out dots if nothing else has been printed async with shared_data.lock: if shared_data.number_of_lines_printed_to_the_console == 0: print(".", end="") @@ -178,11 +188,8 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori if wait_params is None: return if isinstance(wait_params, ProcessExitWaitParameters): - # The idea behind this command style is that we want to wait until the process has - # exited and we want to support 'exit' timeout values that are not long and arbitrary. - # In order to do that, we wait for a lull in the console output before waiting - # for the process exit. So, the exit timeout can hopefully be relative to the - # finishing of the console output. + # The idea behind this wait style is that we want to wait until the process has + # exited, and if it fails to exit, we want to time out after a reasonable time. # Of course, if the app doesn't support the "exit" command, there is no sense in # waiting for the process to respond to it. But, we tell users that we skipped it. if "exit" in target_proc_info.supported_commands: @@ -194,20 +201,31 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori else: if verbosity_level >= IntegtestVerbosityLevels.integtest_debug: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") - print(f"[integtest_proc_mgmt {now_string}] The {proc_name} process doesn't support the 'exit' command, so waiting for exit was skipped") + print(f"[integtest_proc_mgmt {now_string}] The {proc_name} process doesn't support the 'exit' command, so we won't wait for a response") elif isinstance(wait_params, EchoCommandWaitParameters): if "echo" in target_proc_info.supported_commands: + # start a background task to watch for the echo command output + bg_task = asyncio.create_task(wait_for_requested_condition(cmd_start_time, wait_params, shared_data)) + # send the echo command to the process target_proc.stdin.write((f"echo '{wait_params.search_phrase}'\n").encode()) await target_proc.stdin.drain() if verbosity_level >= IntegtestVerbosityLevels.integtest_debug: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] Sent command to {proc_name}: echo '{wait_params.search_phrase}'") - await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) + # wait until the echo command output shows up in the console output + await bg_task else: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") - print(f"[integtest_proc_mgmt {now_string}] WARNING: The {proc_name} process doesn't support the 'echo' command, using TIME wait instead'") + print(f"[integtest_proc_mgmt {now_string}] WARNING: The {proc_name} process doesn't support the 'echo' command, using time-based wait instead'") + # we use a KeyPhrase wait parameter set here (without setting a search phrase) + # because its default timeout values are longer than a couple of seconds but not too + # long. In any case, this choice is a poor substitute for a test string to be echo-ed + # by the application, and there is a non-trivial chance that the timeout values are + # not well-matched to the console output that is produced by the process. wait_params = KeyPhraseWaitParameters() await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) + elif isinstance(wait_params, KeyPhraseWaitParameters) and key_phrase_bg_task is not None: + await key_phrase_bg_task else: # KeyPhrase gets handled automatically here, along with ConsoleOutput await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) From 80c98326eb5186881e7ef2c6898a4b60fa342d3a Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Mon, 21 Sep 2026 11:05:03 -0500 Subject: [PATCH 07/11] first update to the InformationAboutSpecialVariables.md documentation to include the updated command wait styles. --- docs/InformationAboutSpecialVariables.md | 83 +++++++++++++++--------- 1 file changed, 53 insertions(+), 30 deletions(-) diff --git a/docs/InformationAboutSpecialVariables.md b/docs/InformationAboutSpecialVariables.md index d44d8e5..03b6220 100644 --- a/docs/InformationAboutSpecialVariables.md +++ b/docs/InformationAboutSpecialVariables.md @@ -62,38 +62,51 @@ class DAQSessionIngredients: class DAQControlApplication: alias: str # a short-hand name for the process that is started startup_strings: list[str] # the elements of the command string that should be used to start the application - wait_time_after_start: int = 2 # seconds to sleep after spawning the process + startup_wait_params: ConsoleOutputWaitParameters = None @dataclass class DAQCommandSet: target: str # the name of the process that should receive the commands - command_list: list[str] $ the list of commands, e.g. ["boot", "conf"] - wait_params: CommandWaitParameters = field(default_factory=lambda: CommandWaitParameters()) + command_list: list[str] # the list of commands, e.g. ["boot", "conf"] + wait_params: ConsoleOutputWaitParameters = None + wait_for_command_completion: bool = True @dataclass -class CommandWaitParameters: # please see the comments below for information about this class, etc. - wait_for_command_completion: bool = True - style: CommandWaitStyle = CommandWaitStyle.TIME +class ConsoleOutputWaitParameters: timeout_waiting_for_first_msg: int = 2 # seconds wait_time_after_last_msg: int = 2 # seconds - timeout_waiting_for_exit: int = 5 # seconds -class CommandWaitStyle(Enum): - ECHO = "echo" - TIME = "time" - TIME_PLUS_EXIT = "time_plus_exit" - NONE = "none" +@dataclass +class KeyPhraseWaitParameters(ConsoleOutputWaitParameters): + timeout_waiting_for_first_msg: int = 30 # seconds + wait_time_after_last_msg: int = 30 # seconds + search_phrase: str = None + +@dataclass +class EchoCommandWaitParameters(ConsoleOutputWaitParameters): + timeout_waiting_for_first_msg: int = 999999 # seconds + wait_time_after_last_msg: int = 999999 # seconds + search_phrase: str = "*** COMMAND HAS COMPLETED ***" + +@dataclass +class ProcessExitWaitParameters(ConsoleOutputWaitParameters): + timeout_waiting_for_first_msg: int = 30 # seconds + wait_time_after_last_msg: int = 30 # seconds + process: asyncio.subprocess.Process = None ``` -* Here is some additional information about `CommandWaitParameters`: +* Here is some additional information about `ConsoleOutputParameters`: * the commands that are specified in a `DAQCommandSet` are sent individually to the target process without any delay between them. So, we typically send all of the commands in the set in a fraction of a second, while the target process could take tens of seconds to execute all of them. - * when there is only one control process in an integtest, this rapid-fire approach may be all that we need, because a single process handles the throttling of the commands, running them one after another. However, when there are multiple control processes in an integtest, we may want to send a set of commands to Process1, wait for those to finish, and only then send a set of commands to Process2. This demonstrates a need to allow an `integrationtest` developer to specify whether they want the integrationtest infrastructure to wait for each command set to finish before moving on to the next set of commands, and if so, what style of waiting they would like be used. This is the motivation for the `CommandWaitParameters` class. + * when there is only one control process in an integtest, this rapid-fire approach may be all that we need, because a single process handles the throttling of the commands, running them one after another. However, when there are multiple control processes in an integtest, we may want to send a set of commands to Process1, wait for those to finish, and only then send a set of commands to Process2. This demonstrates a need to allow an `integrationtest` developer to specify whether they want the integrationtest infrastructure to wait for each command set to finish before moving on to the next set of commands, and if so, what style of waiting they would like be used. This is the motivation for the `ConsoleOutputWaitParameters` class and its child classes. * of course, there are also situations in which we want to wait for all of the requested commands to finish running even when there is only one control process in the integtest. For example, we will likely want to allow a single process to finish executing all of the requested commands before the `integrationtest` infrastructure starts shutting down that process. - * the currently-supported wait styles are ECHO, TIME, and TIME_PLUS_EXIT. - * the ECHO wait style makes use of the `echo` command that is available in some of our control applications to clearly identify when a set of commands has finished. So, if a user specifies a command set that contains commands `['boot', 'conf']` and has a wait style of ECHO, the `integrationtest` infrastructure appends an `echo` command with a special string to the set, i.e. `['boot', 'conf', 'echo ""']`. When the `integrationtest` infrastructure sees the special string in the output of the target process, it knows that the command set has finished. - * this wait style is the most robust since we know that all of the commands before the `echo` command have been run when the `echo` results are seen in the process output. However, some applications don't provide `echo` functionality. - * the TIME wait style simply waits for configured amounts of time for console output to start and then stop. The idea here is to use the console output as an indicator of activity, and when the console output stops, presume that activity related to the requested command(s) has stopped. - * the TIME_PLUS_EXIT wait style is intended to be used with "exit" commands. The idea here is to wait for console output to stop and then wait for the process to exit (within a configurable timeout). + * the currently-supported wait styles are console-output, echo-command, key-phrase, and process-exit. + * the console-output wait style simply waits for configured amounts of time for console output to start and then stop. The idea here is to use the console output as an indicator of activity, and when the console output stops, we presume that activity related to the requested command(s) has stopped. + * the echo-command wait style makes use of the `echo` command that is available in some of our control applications to clearly identify when a set of commands has finished. So, if a user specifies a command set that contains commands `['boot', 'conf']` and has a wait style of echo-command, the `integrationtest` infrastructure appends an `echo` command with a special string to the set, i.e. `['boot', 'conf', 'echo ""']`. When the `integrationtest` infrastructure sees the special string in the output of the target process, it knows that the command set has finished. + * this wait style is robust since we know that all of the commands before the `echo` command have been run when the `echo` results are seen in the process output. However, some applications don't provide `echo` functionality. + * this wait style inherits from the console-output wait style, so, in principle, it will time out if the special echo string is not seen in the console output. However, the default values for the console output timeouts are set very long so that we don't accidentlly time out too soon (for example, if an integtest includes a 300-second data-taking run). Of course, integtest developers can choose smaller timeout values for special situations. + * the key-phrase wait style looks for a specific phrase in the console output, and the `integrationtest` infrastructure stops waiting when it sees that phrase. + * the process-exit wait style is intended to be used with "exit" commands. The idea here is to wait for console output to stop and then wait for the process to exit (within a configurable timeout). + * if, for some reason, the process does not exit in response to the 'exit' command, this wait style fails back on the console-output wait style. In this way, we avoid having the `integrationtest` infrastructure wait forever for a process to exit. * There are several strings that are dynamically determined by the `integrationtest` infrastructure that we may want to include in the `startup_strings` field in our `DAQControlApplication` declarations. To take this into account, placeholder strings have been defined. These placeholder strings can be used in `DAQControlApplication` declarations and the `integrationtest` infrastructure will substitute the appropriate value at runtime. The placeholders that are currently available are the following: * `` - the process manager type that should be used in the DAQ session * recall that the `integrationtest` infrastructure has support for user-specified (dynamic) process manager types. If we don't want to make use of that functionality, we can hard-code the process manager type in our `DAQControlApplication.startup_strings`. Of course, that reduces flexibility, but there may be cases where it would make sense. @@ -110,10 +123,12 @@ Here is a snippet of code from the `basic_multapp_test.py` that shows how the `D ```python # The commands to run in dunerc and the process manager shell dunerc_commands_1 = ( - "boot conf start --run-number 101 wait 1 enable-triggers wait ".split() - + [str(run_duration)] + ["disable-triggers"] + "boot conf start --run-number 101 wait 1".split() ) dunerc_commands_2 = ( + "enable-triggers wait".split() + [str(run_duration)] + ["disable-triggers"] +) +dunerc_commands_3 = ( "drain-dataflow stop-trigger-sources stop wait 2 scrap terminate".split() ) pmshell_command = ["ps"] @@ -123,23 +138,31 @@ pm_port = find_free_port(50020, 52000) # The command lines that should be used to start the applications procmsg_startup_commands = ["drunc-process-manager", "", str(pm_port)] -pmapp = DAQControlApplication("pm", procmsg_startup_commands) +pmapp = idc.DAQControlApplication("pm", procmsg_startup_commands, + idc.KeyPhraseWaitParameters(search_phrase="communicating through", + timeout_waiting_for_first_msg=5, + wait_time_after_last_msg=5)) pmshell_startup_commands = ["drunc-process-manager-shell", f"grpc://localhost:{pm_port}"] -pmshellapp = DAQControlApplication("pmshell", pmshell_startup_commands) +pmshellapp = idc.DAQControlApplication("pmshell", pmshell_startup_commands, + idc.KeyPhraseWaitParameters(search_phrase="Ready")) -drunc_startup_commands = ["drunc-unified-shell", f"grpc://localhost:{pm_port}", "", "", ""] -druncapp = DAQControlApplication("drunc", drunc_startup_commands) +drunc_startup_commands = ["drunc-unified-shell", f"grpc://localhost:{pm_port}", + "", "", ""] +druncapp = idc.DAQControlApplication("drunc", drunc_startup_commands, + idc.KeyPhraseWaitParameters(search_phrase="unified_shell ready")) # Packaging up the commands into DAQCommandSets -cmd_set_1 = DAQCommandSet("drunc", dunerc_commands_1, CommandWaitParameters(style=CommandWaitStyle.ECHO)) -cmd_set_2 = DAQCommandSet("pmshell", pmshell_command, CommandWaitParameters(style=CommandWaitStyle.TIME)) -cmd_set_3 = DAQCommandSet("drunc", dunerc_commands_2, CommandWaitParameters(style=CommandWaitStyle.ECHO)) +cmd_set_1 = idc.DAQCommandSet("drunc", dunerc_commands_1, idc.EchoCommandWaitParameters()) +cmd_set_2 = idc.DAQCommandSet("pmshell", pmshell_command, wait_for_command_completion=False) +cmd_set_3 = idc.DAQCommandSet("drunc", dunerc_commands_2, idc.EchoCommandWaitParameters()) +cmd_set_4 = idc.DAQCommandSet("pmshell", pmshell_command, idc.KeyPhraseWaitParameters(search_phrase="mlt")) +cmd_set_5 = idc.DAQCommandSet("drunc", dunerc_commands_3, idc.EchoCommandWaitParameters()) # Putting everything together into a DAQSessionIngredients object app_list = [ pmapp, pmshellapp, druncapp ] -cmd_set_list = [ cmd_set_1, cmd_set_2, cmd_set_3 ] -dsi = DAQSessionIngredients(app_list, cmd_set_list) +cmd_set_list = [ cmd_set_1, cmd_set_2, cmd_set_3, cmd_set_4, cmd_set_5 ] +dsi = idc.DAQSessionIngredients(app_list, cmd_set_list) # Declare the special variable that tells the integrationtest infrastructure what we want to run daq_session_ingredients = {"MultiRCAppSession": dsi} From b0b6b25023fbc2c93e2da35bdf9827da96205b30 Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Mon, 21 Sep 2026 11:22:58 -0500 Subject: [PATCH 08/11] changes to InformationAboutSpecialVariables.md --- docs/InformationAboutSpecialVariables.md | 21 +++++++++++---------- 1 file changed, 11 insertions(+), 10 deletions(-) diff --git a/docs/InformationAboutSpecialVariables.md b/docs/InformationAboutSpecialVariables.md index 03b6220..f81a924 100644 --- a/docs/InformationAboutSpecialVariables.md +++ b/docs/InformationAboutSpecialVariables.md @@ -1,6 +1,6 @@ # Special variables that are used by the integrationtest infrastructure -18-Aug-2026, Kurt Biery +21-Sep-2026, Kurt Biery ## Introduction @@ -97,16 +97,17 @@ class ProcessExitWaitParameters(ConsoleOutputWaitParameters): * Here is some additional information about `ConsoleOutputParameters`: * the commands that are specified in a `DAQCommandSet` are sent individually to the target process without any delay between them. So, we typically send all of the commands in the set in a fraction of a second, while the target process could take tens of seconds to execute all of them. - * when there is only one control process in an integtest, this rapid-fire approach may be all that we need, because a single process handles the throttling of the commands, running them one after another. However, when there are multiple control processes in an integtest, we may want to send a set of commands to Process1, wait for those to finish, and only then send a set of commands to Process2. This demonstrates a need to allow an `integrationtest` developer to specify whether they want the integrationtest infrastructure to wait for each command set to finish before moving on to the next set of commands, and if so, what style of waiting they would like be used. This is the motivation for the `ConsoleOutputWaitParameters` class and its child classes. + * when there is only one control process in an integtest, this rapid-fire approach may be all that we need, because a single process handles the throttling of the commands, running them one after another. However, when there are multiple control processes in an integtest, we may want to send a set of commands to Process1, wait for those to finish, and only then send a set of commands to Process2. This demonstrates a need to allow an `integrationtest` developer to specify whether they want the `integrationtest` infrastructure to wait for each command set to finish before moving on to the next set of commands, and if so, what style of waiting they would like be used. This is the motivation for the `ConsoleOutputWaitParameters` class and its child classes. * of course, there are also situations in which we want to wait for all of the requested commands to finish running even when there is only one control process in the integtest. For example, we will likely want to allow a single process to finish executing all of the requested commands before the `integrationtest` infrastructure starts shutting down that process. - * the currently-supported wait styles are console-output, echo-command, key-phrase, and process-exit. - * the console-output wait style simply waits for configured amounts of time for console output to start and then stop. The idea here is to use the console output as an indicator of activity, and when the console output stops, we presume that activity related to the requested command(s) has stopped. - * the echo-command wait style makes use of the `echo` command that is available in some of our control applications to clearly identify when a set of commands has finished. So, if a user specifies a command set that contains commands `['boot', 'conf']` and has a wait style of echo-command, the `integrationtest` infrastructure appends an `echo` command with a special string to the set, i.e. `['boot', 'conf', 'echo ""']`. When the `integrationtest` infrastructure sees the special string in the output of the target process, it knows that the command set has finished. - * this wait style is robust since we know that all of the commands before the `echo` command have been run when the `echo` results are seen in the process output. However, some applications don't provide `echo` functionality. - * this wait style inherits from the console-output wait style, so, in principle, it will time out if the special echo string is not seen in the console output. However, the default values for the console output timeouts are set very long so that we don't accidentlly time out too soon (for example, if an integtest includes a 300-second data-taking run). Of course, integtest developers can choose smaller timeout values for special situations. - * the key-phrase wait style looks for a specific phrase in the console output, and the `integrationtest` infrastructure stops waiting when it sees that phrase. - * the process-exit wait style is intended to be used with "exit" commands. The idea here is to wait for console output to stop and then wait for the process to exit (within a configurable timeout). - * if, for some reason, the process does not exit in response to the 'exit' command, this wait style fails back on the console-output wait style. In this way, we avoid having the `integrationtest` infrastructure wait forever for a process to exit. + * the currently-supported wait styles are _console-output_, _echo-command_, _key-phrase_, and _process-exit_. + * the _console-output_ wait style simply waits for configured amounts of time for console output to start and then stop. The idea here is to use the console output as an indicator of activity, and when the console output stops, we presume that activity related to the requested command(s) has stopped. + * the _echo-command_ wait style makes use of the `echo` command that is available in some of our control applications to clearly identify when a set of commands has finished. So, if a user specifies a command set that contains commands `['boot', 'conf']` and has a wait style of _echo-command_, the `integrationtest` infrastructure appends an `echo` command with a special string to the set, i.e. `['boot', 'conf', 'echo ""']`. When the `integrationtest` infrastructure sees the special string in the output of the target process, it knows that the command set has finished. + * this wait style is quite robust since we know that all of the commands before the `echo` command have been run when the `echo` results are seen in the process output. However, some applications don't provide `echo` functionality. In those cases, the `integrationtest` infrastructure switches to a _console-output_ wait style with timeout values taken from the _key-phrase_ defaults. + * this wait style inherits from the _console-output_ wait style, so, in principle, it will time out if the special echo string is not seen in the console output. However, the default values for the console output timeouts are set very long so that we don't accidentlly time out too soon (for example, if an integtest includes a 300-second data-taking run). Of course, integtest developers can choose smaller timeout values for special situations. + * the _key-phrase_ wait style looks for a specific phrase in the console output, and the `integrationtest` infrastructure stops waiting when it sees that phrase. + * if the phrase is not found before the _console-output_ timeout values in the `KeyPhraseWaitParameters` instance, then the infrastructure will stop waiting and print out a warning message. The default values for _console-output_ timeout values in `KeyPhraseWaitParameters` instances is + * the _process-exit_ wait style is intended to be used with "exit" commands. The idea here is to wait for console output to stop and then wait for the process to exit (within a configurable timeout). + * if, for some reason, the process does not exit in response to the 'exit' command, this wait style fails back on the _console-output_ wait style. In this way, we avoid having the `integrationtest` infrastructure wait forever for a process to exit. * There are several strings that are dynamically determined by the `integrationtest` infrastructure that we may want to include in the `startup_strings` field in our `DAQControlApplication` declarations. To take this into account, placeholder strings have been defined. These placeholder strings can be used in `DAQControlApplication` declarations and the `integrationtest` infrastructure will substitute the appropriate value at runtime. The placeholders that are currently available are the following: * `` - the process manager type that should be used in the DAQ session * recall that the `integrationtest` infrastructure has support for user-specified (dynamic) process manager types. If we don't want to make use of that functionality, we can hard-code the process manager type in our `DAQControlApplication.startup_strings`. Of course, that reduces flexibility, but there may be cases where it would make sense. From 14c9fcedfe3966017bdc2fcdf25305a986e35231 Mon Sep 17 00:00:00 2001 From: bieryAtFnal <36311946+bieryAtFnal@users.noreply.github.com> Date: Mon, 21 Sep 2026 12:01:45 -0500 Subject: [PATCH 09/11] Updated InformationAboutSpecialVariables.md Updated documentation to reflect changes in ConsoleOutputWaitParameters and its child classes. Clarified wait styles and their behaviors. --- docs/InformationAboutSpecialVariables.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/InformationAboutSpecialVariables.md b/docs/InformationAboutSpecialVariables.md index f81a924..e21d1bf 100644 --- a/docs/InformationAboutSpecialVariables.md +++ b/docs/InformationAboutSpecialVariables.md @@ -95,7 +95,7 @@ class ProcessExitWaitParameters(ConsoleOutputWaitParameters): process: asyncio.subprocess.Process = None ``` -* Here is some additional information about `ConsoleOutputParameters`: +* Here is some additional information about `ConsoleOutputWaitParameters` and its child classes: * the commands that are specified in a `DAQCommandSet` are sent individually to the target process without any delay between them. So, we typically send all of the commands in the set in a fraction of a second, while the target process could take tens of seconds to execute all of them. * when there is only one control process in an integtest, this rapid-fire approach may be all that we need, because a single process handles the throttling of the commands, running them one after another. However, when there are multiple control processes in an integtest, we may want to send a set of commands to Process1, wait for those to finish, and only then send a set of commands to Process2. This demonstrates a need to allow an `integrationtest` developer to specify whether they want the `integrationtest` infrastructure to wait for each command set to finish before moving on to the next set of commands, and if so, what style of waiting they would like be used. This is the motivation for the `ConsoleOutputWaitParameters` class and its child classes. * of course, there are also situations in which we want to wait for all of the requested commands to finish running even when there is only one control process in the integtest. For example, we will likely want to allow a single process to finish executing all of the requested commands before the `integrationtest` infrastructure starts shutting down that process. @@ -103,9 +103,9 @@ class ProcessExitWaitParameters(ConsoleOutputWaitParameters): * the _console-output_ wait style simply waits for configured amounts of time for console output to start and then stop. The idea here is to use the console output as an indicator of activity, and when the console output stops, we presume that activity related to the requested command(s) has stopped. * the _echo-command_ wait style makes use of the `echo` command that is available in some of our control applications to clearly identify when a set of commands has finished. So, if a user specifies a command set that contains commands `['boot', 'conf']` and has a wait style of _echo-command_, the `integrationtest` infrastructure appends an `echo` command with a special string to the set, i.e. `['boot', 'conf', 'echo ""']`. When the `integrationtest` infrastructure sees the special string in the output of the target process, it knows that the command set has finished. * this wait style is quite robust since we know that all of the commands before the `echo` command have been run when the `echo` results are seen in the process output. However, some applications don't provide `echo` functionality. In those cases, the `integrationtest` infrastructure switches to a _console-output_ wait style with timeout values taken from the _key-phrase_ defaults. - * this wait style inherits from the _console-output_ wait style, so, in principle, it will time out if the special echo string is not seen in the console output. However, the default values for the console output timeouts are set very long so that we don't accidentlly time out too soon (for example, if an integtest includes a 300-second data-taking run). Of course, integtest developers can choose smaller timeout values for special situations. + * this wait style inherits from the _console-output_ wait style, so, in principle, it will time out if the special echo string is not seen in the console output. However, the default values for the console output timeouts are set very long so that we don't accidentally time out too soon (for example, if an integtest includes a 300-second data-taking run). Of course, integtest developers can choose smaller timeout values for special situations. * the _key-phrase_ wait style looks for a specific phrase in the console output, and the `integrationtest` infrastructure stops waiting when it sees that phrase. - * if the phrase is not found before the _console-output_ timeout values in the `KeyPhraseWaitParameters` instance, then the infrastructure will stop waiting and print out a warning message. The default values for _console-output_ timeout values in `KeyPhraseWaitParameters` instances is + * if the phrase is not found before the console output times out based on the timeout values in the `KeyPhraseWaitParameters` instance, then the infrastructure will stop waiting and print out a warning message. * the _process-exit_ wait style is intended to be used with "exit" commands. The idea here is to wait for console output to stop and then wait for the process to exit (within a configurable timeout). * if, for some reason, the process does not exit in response to the 'exit' command, this wait style fails back on the _console-output_ wait style. In this way, we avoid having the `integrationtest` infrastructure wait forever for a process to exit. * There are several strings that are dynamically determined by the `integrationtest` infrastructure that we may want to include in the `startup_strings` field in our `DAQControlApplication` declarations. To take this into account, placeholder strings have been defined. These placeholder strings can be used in `DAQControlApplication` declarations and the `integrationtest` infrastructure will substitute the appropriate value at runtime. The placeholders that are currently available are the following: From 465ab8c933b4339fc1849a17c285d88c854f9cd5 Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Mon, 21 Sep 2026 12:29:37 -0500 Subject: [PATCH 10/11] Added warning messages when commands time out in async_proc_mgmt.py. --- src/integrationtest/async_proc_mgmt.py | 18 +++++++++++++++--- 1 file changed, 15 insertions(+), 3 deletions(-) diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index fe138eb..6182d6f 100644 --- a/src/integrationtest/async_proc_mgmt.py +++ b/src/integrationtest/async_proc_mgmt.py @@ -116,8 +116,11 @@ async def read_stream(stream, process_name, app_exe_name, print_proc_name, run_d # output from the control process(es). Wait parameter classes that build on the # ConsoleOutputWaitParameters class add conditions that allow the waiting to end earlier # than the wait times specified in the base class. +# Return codes are 0 for console output timeout, 1 for finding a search phrase in the console +# output, and 2 for when the process has exited. async def wait_for_requested_condition(start_time, wait_params: ConsoleOutputWaitParameters, shared_data: OutputMonitoringSharedData): + retcode = 0 if (isinstance(wait_params, EchoCommandWaitParameters) or \ isinstance(wait_params, KeyPhraseWaitParameters)) and \ @@ -137,10 +140,12 @@ async def wait_for_requested_condition(start_time, wait_params: ConsoleOutputWai if now - shared_data.last_msg_time >= wait_params.wait_time_after_last_msg: break if shared_data.search_phrase_has_been_found: + retcode = 1 break if isinstance(wait_params, ProcessExitWaitParameters) and \ wait_params.process is not None: if wait_params.process.returncode is not None: + retcode = 2 break await asyncio.sleep(0.25) now = time.time() @@ -152,6 +157,8 @@ async def wait_for_requested_condition(start_time, wait_params: ConsoleOutputWai shared_data.phrase_searching_in_progress = False shared_data.search_phrase_has_been_found = False + return retcode + async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitoringSharedData, cmd_list, wait_params, verbosity_level): @@ -213,7 +220,10 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] Sent command to {proc_name}: echo '{wait_params.search_phrase}'") # wait until the echo command output shows up in the console output - await bg_task + retcode = await bg_task + if retcode != 1: + now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") + print(f"[integtest_proc_mgmt {now_string}] WARNING: timeout waiting for {proc_name} to echo '{wait_params.search_phrase}' after executing {cmd_list}") else: now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") print(f"[integtest_proc_mgmt {now_string}] WARNING: The {proc_name} process doesn't support the 'echo' command, using time-based wait instead'") @@ -225,9 +235,11 @@ async def send_commands(target_proc_info, proc_name, shared_data: OutputMonitori wait_params = KeyPhraseWaitParameters() await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) elif isinstance(wait_params, KeyPhraseWaitParameters) and key_phrase_bg_task is not None: - await key_phrase_bg_task + retcode = await key_phrase_bg_task + if retcode != 1: + now_string = datetime.now(timezone.utc).strftime("%H:%M:%SZ") + print(f"[integtest_proc_mgmt {now_string}] WARNING: timeout waiting for {proc_name} to print out '{wait_params.search_phrase}' as part of executing {cmd_list}") else: - # KeyPhrase gets handled automatically here, along with ConsoleOutput await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) # add background task to avoid race condition? From 9912c47ab30647eca37de6739923d78f5416bc0f Mon Sep 17 00:00:00 2001 From: Kurt Biery Date: Tue, 22 Sep 2026 14:12:42 -0500 Subject: [PATCH 11/11] minor tweaks --- docs/InformationAboutSpecialVariables.md | 8 ++++---- src/integrationtest/async_proc_mgmt.py | 1 - 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/docs/InformationAboutSpecialVariables.md b/docs/InformationAboutSpecialVariables.md index e21d1bf..3dfaca7 100644 --- a/docs/InformationAboutSpecialVariables.md +++ b/docs/InformationAboutSpecialVariables.md @@ -101,13 +101,13 @@ class ProcessExitWaitParameters(ConsoleOutputWaitParameters): * of course, there are also situations in which we want to wait for all of the requested commands to finish running even when there is only one control process in the integtest. For example, we will likely want to allow a single process to finish executing all of the requested commands before the `integrationtest` infrastructure starts shutting down that process. * the currently-supported wait styles are _console-output_, _echo-command_, _key-phrase_, and _process-exit_. * the _console-output_ wait style simply waits for configured amounts of time for console output to start and then stop. The idea here is to use the console output as an indicator of activity, and when the console output stops, we presume that activity related to the requested command(s) has stopped. - * the _echo-command_ wait style makes use of the `echo` command that is available in some of our control applications to clearly identify when a set of commands has finished. So, if a user specifies a command set that contains commands `['boot', 'conf']` and has a wait style of _echo-command_, the `integrationtest` infrastructure appends an `echo` command with a special string to the set, i.e. `['boot', 'conf', 'echo ""']`. When the `integrationtest` infrastructure sees the special string in the output of the target process, it knows that the command set has finished. - * this wait style is quite robust since we know that all of the commands before the `echo` command have been run when the `echo` results are seen in the process output. However, some applications don't provide `echo` functionality. In those cases, the `integrationtest` infrastructure switches to a _console-output_ wait style with timeout values taken from the _key-phrase_ defaults. + * the _echo-command_ wait style makes use of the `echo` command that is available in some of our control applications to clearly identify when a set of commands has finished. So, if a user specifies a command set that contains commands `['boot', 'conf']` and has a wait style of _echo-command_, the `integrationtest` infrastructure appends an `echo` command with a special string to the set, i.e. `['boot', 'conf', 'echo ""']`. When the `integrationtest` infrastructure sees that special string in the output of the target process, it knows that the command set has finished. + * this wait style is quite robust since we know that all of the commands before the `echo` command have been run when the `echo` results are seen in the process output. However, some applications don't provide `echo` functionality. In the unlikely even that this wait style is requested from an application type that doesn't support it, the `integrationtest` infrastructure will switch to a _console-output_ wait style with timeout values taken from the _key-phrase_ defaults. * this wait style inherits from the _console-output_ wait style, so, in principle, it will time out if the special echo string is not seen in the console output. However, the default values for the console output timeouts are set very long so that we don't accidentally time out too soon (for example, if an integtest includes a 300-second data-taking run). Of course, integtest developers can choose smaller timeout values for special situations. * the _key-phrase_ wait style looks for a specific phrase in the console output, and the `integrationtest` infrastructure stops waiting when it sees that phrase. * if the phrase is not found before the console output times out based on the timeout values in the `KeyPhraseWaitParameters` instance, then the infrastructure will stop waiting and print out a warning message. - * the _process-exit_ wait style is intended to be used with "exit" commands. The idea here is to wait for console output to stop and then wait for the process to exit (within a configurable timeout). - * if, for some reason, the process does not exit in response to the 'exit' command, this wait style fails back on the _console-output_ wait style. In this way, we avoid having the `integrationtest` infrastructure wait forever for a process to exit. + * the _process-exit_ wait style is intended to be used with "exit" commands. The idea here is to wait for console output to stop and the process to exit (within a configurable timeout). + * if, for some reason, the process does not exit in response to the 'exit' command, the timeout values in the ProcessExitWaitParameters object are used to stop waiting in a reasonable amount of time. * There are several strings that are dynamically determined by the `integrationtest` infrastructure that we may want to include in the `startup_strings` field in our `DAQControlApplication` declarations. To take this into account, placeholder strings have been defined. These placeholder strings can be used in `DAQControlApplication` declarations and the `integrationtest` infrastructure will substitute the appropriate value at runtime. The placeholders that are currently available are the following: * `` - the process manager type that should be used in the DAQ session * recall that the `integrationtest` infrastructure has support for user-specified (dynamic) process manager types. If we don't want to make use of that functionality, we can hard-code the process manager type in our `DAQControlApplication.startup_strings`. Of course, that reduces flexibility, but there may be cases where it would make sense. diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index 6182d6f..dab6972 100644 --- a/src/integrationtest/async_proc_mgmt.py +++ b/src/integrationtest/async_proc_mgmt.py @@ -297,7 +297,6 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r # determine the supported commands for each app (using the 'help' command) help_cmd = ["help"] help_cmd_wait_params = ConsoleOutputWaitParameters(timeout_waiting_for_first_msg=2) - #await wait_for_requested_condition(time.time(), help_cmd_wait_params, shared_data) for proc_name, proc_info in processes.items(): async with shared_data.lock: shared_data.results_of_parsing_help_output = []