diff --git a/docs/InformationAboutSpecialVariables.md b/docs/InformationAboutSpecialVariables.md index d44d8e5..3dfaca7 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 @@ -62,38 +62,52 @@ 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 `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 `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 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 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. @@ -110,10 +124,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 +139,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} diff --git a/src/integrationtest/async_proc_mgmt.py b/src/integrationtest/async_proc_mgmt.py index 14c9ad4..dab6972 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 = "" @@ -31,11 +28,13 @@ 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 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 # process the output of the "help" command, if requested async with shared_data.lock: @@ -111,8 +110,26 @@ 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): +# 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. +# 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 \ + wait_params.search_phrase is not None: + async with shared_data.lock: + shared_data.search_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: @@ -122,20 +139,45 @@ async def wait_for_console_output_lull(start_time, wait_params: CommandWaitParam else: 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() + 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.search_phrase_has_been_found = False + + return retcode + -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): + # 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() @@ -143,60 +185,72 @@ 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}] 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="") # 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: - # 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. + if isinstance(wait_params, ProcessExitWaitParameters): + # 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. - 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 - 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}") 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") - elif wait_params.style == CommandWaitStyle.ECHO: + 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: - shared_data.cmd_cmplt_evt.clear() - target_proc.stdin.write((f"echo '{PROCESS_ECHO_STRING}'\n").encode()) + # 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 '{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}'") + # wait until the echo command output shows up in the console output + 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}] 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: - await wait_for_console_output_lull(cmd_start_time, wait_params, shared_data) + 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: + 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: + await wait_for_requested_condition(cmd_start_time, wait_params, shared_data) + +# add background task to avoid race condition? async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, run_dir, verbosity_level): processes = {} 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: @@ -223,7 +277,12 @@ async def intg_process_manager(daq_session_ingredients: DAQSessionIngredients, r run_dir, shared_data, verbosity_level )) - time.sleep(session_app.wait_time_after_start) + # 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 = ConsoleOutputWaitParameters() + 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") @@ -237,8 +296,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) - await wait_for_console_output_lull(time.time(), help_cmd_wait_params, shared_data) + help_cmd_wait_params = ConsoleOutputWaitParameters(timeout_waiting_for_first_msg=2) for proc_name, proc_info in processes.items(): async with shared_data.lock: shared_data.results_of_parsing_help_output = [] @@ -302,8 +360,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..ac41d9c 100644 --- a/src/integrationtest/data_classes.py +++ b/src/integrationtest/data_classes.py @@ -129,31 +129,41 @@ 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 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 @dataclass class DAQControlApplication: alias: str startup_strings: list[str] - wait_time_after_start: int = 2 # seconds + startup_wait_params: ConsoleOutputWaitParameters = None @dataclass class DAQCommandSet: target: str command_list: list[str] - wait_params: CommandWaitParameters = field(default_factory=lambda: CommandWaitParameters()) + wait_params: ConsoleOutputWaitParameters = None + wait_for_command_completion: bool = True @dataclass class DAQSessionIngredients: @@ -166,10 +176,12 @@ 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 number_of_lines_printed_to_the_console: int = 0 + search_phrase: str = "nullnullnull" + phrase_searching_in_progress: 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) diff --git a/src/integrationtest/integrationtest_drunc.py b/src/integrationtest/integrationtest_drunc.py index a6e7bdd..8944cab 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: @@ -645,9 +645,10 @@ 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, CommandWaitParameters(style=CommandWaitStyle.ECHO)) + requested_cmds = DAQCommandSet("drunc", run_control_commands, EchoCommandWaitParameters()) app_list = [ dsapp ] cmd_set_list = [ requested_cmds, exit_cmd ]