qq_lib.run

Execution utilities for running qq jobs inside the batch environment.

This module defines the Runner class, which prepares the execution environment, launches the user's job script, updates qq's state tracking, and performs cleanup on success, failure, or interruption. It handles both shared and scratch working directories, loop-job archiving, resubmission, communication with the batch system, and SIGTERM-safe shutdown.

 1# Released under MIT License.
 2# Copyright (c) 2025-2026 Ladislav Bartos and Robert Vacha Lab
 3
 4"""
 5Execution utilities for running qq jobs inside the batch environment.
 6
 7This module defines the `Runner` class, which prepares the execution
 8environment, launches the user's job script, updates qq's state tracking,
 9and performs cleanup on success, failure, or interruption. It handles both
10shared and scratch working directories, loop-job archiving, resubmission,
11communication with the batch system, and SIGTERM-safe shutdown.
12"""
13
14from .runner import Runner
15
16__all__ = [
17    "Runner",
18]
class Runner:
 41class Runner:
 42    """
 43    Manages the setup, execution, and cleanup of scripts within the qq batch environment.
 44
 45    The Runner class is responsible for:
 46      - Preparing a working directory (shared or scratch space)
 47      - Executing a provided job script
 48      - Updating the job info file with run state, success, or failure
 49      - Cleaning up resources when execution is finished
 50    """
 51
 52    def __init__(self, info_file: Path, host: str):
 53        """
 54        Initialize a new Runner instance.
 55
 56        Args:
 57            info_file (Path): Path to the qq info file that contains job metadata.
 58            host (str): The hostname of the input machine from which the job was submitted.
 59
 60        Raises:
 61            QQRunFatalError: If loading the QQ info file fails fatally during initialization.
 62        """
 63        # install a signal handler
 64        signal.signal(signal.SIGTERM, self._handle_sigterm)
 65
 66        # process running the wrapped script
 67        self._process: subprocess.Popen[str] | None = None
 68
 69        self._info_file = Path(info_file)
 70        logger.debug(f"Info file: '{self._info_file}'.")
 71
 72        self._input_machine = host
 73        logger.debug(f"Input machine: '{self._input_machine}'.")
 74
 75        # load the info file or raise a fatal qq error if this fails
 76        try:
 77            # get the batch system from the environment variable (or guess it)
 78            self._batch_system = BatchInterface.from_env_var_or_guess()
 79            logger.debug(f"Batch system: {str(self._batch_system)}.")
 80
 81            # get the id of the job from the batch system
 82            if not (job_id := self._batch_system.get_job_id()):
 83                raise QQError("Job has no associated job id")
 84
 85            # load the info file
 86            self._informer: Informer = Retryer(
 87                Informer.from_file,
 88                self._info_file,
 89                host=self._input_machine,
 90                max_tries=CFG.runner.retry_tries,
 91                wait_seconds=CFG.runner.retry_wait,
 92            ).run()
 93
 94            # check that the id of this job matches the job id in the info file
 95            if not self._informer.matches_job(job_id):
 96                raise QQJobMismatchError(
 97                    "Info file does not correspond to the current job"
 98                )
 99
100            # check that the batch system in info file matches the one loaded from the environment variable
101            if self._batch_system != self._informer.batch_system:
102                raise QQError(
103                    f"Batch system mismatch - env var: '{str(self._batch_system)}', info file: '{self._informer.batch_system}'"
104                )
105
106        except Exception as e:
107            raise QQRunFatalError(
108                f"Unable to load valid qq info file '{self._info_file}' on '{self._input_machine}': {e}"
109            ) from e
110
111        logger.info(
112            f"[qq-{str(self._batch_system)} v{qq_lib.__version__}] Initializing "
113            f"job '{self._informer.info.job_id}' on host '{socket.getfqdn()}'."
114        )
115
116        # get input directory
117        self._input_dir = Path(self._informer.info.input_dir)
118        logger.debug(f"Input directory: {self._input_dir}.")
119
120        # should the scratch directory be used?
121        self._use_scratch = self._informer.uses_scratch()
122        logger.debug(f"Use scratch: {self._use_scratch}.")
123
124        # initialize archiver, if this is a loop job
125        if loop_info := self._informer.info.loop_info:
126            self._archiver = Archiver(
127                archive=loop_info.archive,
128                archive_format=loop_info.archive_format,
129                input_machine=self._informer.info.input_machine,
130                input_dir=self._informer.info.input_dir,
131                batch_system=self._batch_system,
132                included_files=self._informer.info.included_files,
133                excluded_files=self._informer.info.excluded_files,
134                ignored_files=self._informer.info.ignored_files,
135            )
136            self._should_resubmit = True
137        else:
138            self._archiver = None
139
140        if self._informer.info.job_type == JobType.CONTINUOUS:
141            self._should_resubmit = True
142
143    def prepare(self) -> None:
144        """
145        Prepare the script for execution, setting up the archive
146        and archiving run time files (if this is a loop job) and
147        preparing working directory.
148
149        Raises:
150            QQError: If working directory setup fails.
151        """
152        if self._archiver:
153            assert self._informer.info.loop_info is not None
154            # prepare the directory for archiving
155            self._archiver.make_archive_dir()
156
157            # archive runtime files from the previous cycle
158            # this has to be done before the working directory is prepared,
159            # otherwise the runtime files would get copied to the working directory
160            logger.debug(
161                f"Archiving run time files from cycle {self._informer.info.loop_info.current - 1}."
162            )
163            self._archiver.archive_runtime_files(
164                # we need to escape the '+' character
165                construct_loop_job_name(
166                    self._informer.info.script_name,
167                    self._informer.info.loop_info.current - 1,
168                ).replace("+", "\\+"),
169                self._informer.info.loop_info.current - 1,
170            )
171
172        if self._use_scratch:
173            self._set_up_scratch_dir()
174        else:
175            self._set_up_shared_dir()
176
177        if self._archiver:
178            assert self._informer.info.loop_info is not None
179            # fetch files for the current cycle of the loop job from the archive
180            self._archiver.from_archive(
181                self._work_dir, self._informer.info.loop_info.current
182            )
183
184    def execute(self) -> int:
185        """
186        Execute the job script in the working directory.
187
188        Returns:
189            int: The exit code from the executed script.
190
191        Raises:
192            QQError: If execution fails or info file cannot be updated.
193        """
194        # update the qqinfo file
195        self._update_info_running()
196
197        # get the actual name of the script to execute
198        script = logical_resolve(Path(self._informer.info.script_name))
199
200        # get paths to output files
201        stdout_log = self._informer.info.stdout_file
202        stderr_log = self._informer.info.stderr_file
203
204        logger.info(f"Executing script '{script}'.")
205
206        # get intepreter if configured, otherwise use the default
207        interpreter = self._informer.info.interpreter or Interpreter()
208
209        # get the command to execute
210        command_list = [*interpreter.to_command_list(), str(script)]
211        logger.debug(f"Command executed using subprocess.Popen: {command_list}")
212
213        try:
214            with Path(stdout_log).open("w") as out, Path(stderr_log).open("w") as err:
215                self._process = subprocess.Popen(
216                    command_list,
217                    stdout=out,
218                    stderr=err,
219                    text=True,
220                )
221
222                # wait for the process to finish in a non-blocking manner
223                while self._process.poll() is None:
224                    sleep(CFG.runner.subprocess_checks_wait_time)
225
226        except Exception as e:
227            raise QQError(f"Failed to execute script '{script}': {e}") from e
228
229        # if the script returns an exit code corresponding to CFG.exit_codes.qq_run_no_resubmit,
230        # do not submit the next cycle of the job but return 0
231        if (
232            self._informer.info.job_type in [JobType.LOOP, JobType.CONTINUOUS]
233            and self._process.returncode == CFG.exit_codes.qq_run_no_resubmit
234        ):
235            logger.debug(
236                f"Detected an exit code of '{self._process.returncode}'. Replacing with '0' and will not submit the next cycle of the job."
237            )
238            self._process.returncode = 0
239            self._should_resubmit = False
240
241        return self._process.returncode
242
243    def finalize(self) -> None:
244        """
245        Finalize the execution of the job script.
246
247        Handles post-processing of the job based on the script's exit code and the
248        configured transfer and archive modes. The specific actions taken depend on
249        the job's transfer mode, archive mode, and whether scratch directory is being used.
250
251        Specifically, this method:
252
253        1. Archives files from the working directory if archiving is enabled and the
254            archive mode allows it for the given exit code (loop jobs only).
255        2. Transfers or handles files based on whether scratch directory is used:
256            - If using scratch and transfer mode allows: Syncs the entire working
257                directory back to the input directory (excluding explicitly included files)
258                and removes the working directory from scratch.
259            - If using scratch and transfer mode disallows: Copies only runtime files
260                to the input directory and preserves the working directory.
261            - If not using scratch: No file operations are performed.
262        3. Updates the qq info file to "finished" (exit code 0) or "failed" (non-zero
263            exit code).
264        4. Resubmits the job if it is a loop or continuous job and was completed successfully.
265
266        Raises:
267            QQError: If copying, deletion, or archiving of files fails or if the resubmission fails.
268        """
269        logger.info("Finalizing the execution.")
270        assert self._process is not None
271
272        # archive files
273        if self._archiver and self._informer.should_archive_files(
274            self._process.returncode
275        ):
276            logger.debug(
277                f"Script exit code is '{self._process.returncode}'. Archiving files."
278            )
279            self._archive_files_from_work_dir()
280
281        # transfer files back to the input (submission) directory
282        if self._use_scratch:
283            if self._informer.should_transfer_files(self._process.returncode):
284                logger.debug(
285                    f"Script exit code is '{self._process.returncode}'. Transferring files from working directory."
286                )
287
288                Retryer(
289                    self._batch_system.sync_with_exclusions,
290                    self._work_dir,
291                    self._input_dir,
292                    socket.getfqdn(),
293                    self._informer.info.input_machine,
294                    # exclude files that were specifically included via the `--include` option
295                    # and files that were specifically chosen to be ignored via `--ignore` option
296                    self._get_excluded_from_input_dir(),
297                    max_tries=CFG.runner.retry_tries,
298                    wait_seconds=CFG.runner.retry_wait,
299                ).run()
300
301                # remove the working directory from scratch
302                self._delete_work_dir()
303            else:
304                # copy only the runtime files to input directory
305                # and keep the working directory
306                self._copy_runtime_files_to_input_dir(retry=True)
307
308        if self._process.returncode == 0:
309            # update the qqinfo file
310            self._update_info_finished()
311
312            # if this is a loop/continuous job
313            if self._informer.info.job_type in [JobType.LOOP, JobType.CONTINUOUS]:
314                self._resubmit()
315        else:
316            # update the qqinfo file
317            self._update_info_failed(self._process.returncode)
318
319        logger.info(f"Job completed with an exit code of {self._process.returncode}.")
320
321    def log_failure_and_exit(self, exception: BaseException) -> NoReturn:
322        """
323        Record a failure state into the qq info file and exit the program.
324
325        Args:
326            exception (BaseException): The exception to log.
327
328        Raises:
329            SystemExit: Always exits with the exit code associated with the given exception.
330        """
331        exit_code = getattr(exception, "exit_code", CFG.exit_codes.unexpected_error)
332        try:
333            self._update_info_failed(exit_code)
334            logger.error(terminate(str(exception)))
335            sys.exit(exit_code)
336        except Exception as e:
337            # unable to log the current state into the info file
338            log_fatal_error_and_exit(e)  # exits here
339
340    def _set_up_shared_dir(self) -> None:
341        """
342        Configure the input directory as the working directory.
343        """
344        # set qq working directory to the input dir
345        self._work_dir = self._input_dir
346
347        # move to the working directory
348        Retryer(
349            os.chdir,
350            self._work_dir,
351            max_tries=CFG.runner.retry_tries,
352            wait_seconds=CFG.runner.retry_wait,
353        ).run()
354
355    def _set_up_scratch_dir(self) -> None:
356        """
357        Configure a scratch directory as the working directory.
358
359        Copies all files from the job directory to the working directory
360        (excluding the qq info file).
361
362        Raises:
363            QQError: If scratch directory cannot be determined.
364        """
365        # get path to the working directory (created by the batch system)
366        self._work_dir: Path = Retryer(
367            self._batch_system.create_work_dir_on_scratch,
368            self._informer.info.job_id,
369            max_tries=CFG.runner.retry_tries,
370            wait_seconds=CFG.runner.retry_wait,
371        ).run()
372
373        logger.info(f"Setting up working directory in '{self._work_dir}'.")
374
375        # move to the working directory
376        Retryer(
377            os.chdir,
378            self._work_dir,
379            max_tries=CFG.runner.retry_tries,
380            wait_seconds=CFG.runner.retry_wait,
381        ).run()
382
383        # files excluded from copying to the working directory
384        excluded = self._get_excluded_from_work_dir()
385        logger.debug(
386            f"Files excluded from being copied to the working directory: {excluded}."
387        )
388
389        # copy files from the input directory to the working directory
390        Retryer(
391            self._batch_system.sync_with_exclusions,
392            self._input_dir,
393            self._work_dir,
394            self._informer.info.input_machine,
395            socket.getfqdn(),
396            excluded,
397            max_tries=CFG.runner.retry_tries,
398            wait_seconds=CFG.runner.retry_wait,
399        ).run()
400
401        # copy explicitly included files to the working directory
402        # this will copy files that were specified with the --include option, even if they are also in the list of excluded files
403        logger.debug(
404            f"Files explicitly requested to be copied to the working directory: {self._informer.info.included_files}."
405        )
406        Retryer(
407            self._copy_files,
408            self._informer.info.included_files,
409            max_tries=CFG.runner.retry_tries,
410            wait_seconds=CFG.runner.retry_wait,
411        ).run()
412
413    def _delete_work_dir(self) -> None:
414        """
415        Delete the entire working directory.
416
417        Used only after successful execution in scratch space.
418        """
419        logger.debug(f"Removing working directory '{self._work_dir}'.")
420        Retryer(
421            shutil.rmtree,
422            self._work_dir,
423            max_tries=CFG.runner.retry_tries,
424            wait_seconds=CFG.runner.retry_wait,
425        ).run()
426
427    def _update_info_running(self) -> None:
428        """
429        Update the qq info file to mark the job as running.
430
431        Raises:
432            QQRunCommunicationError: If the job was killed without informing Runner.
433            QQError: If the info file cannot be updated.
434        """
435        logger.debug(f"Updating '{self._info_file}' at job start.")
436        self._reload_info_and_ensure_valid()
437
438        try:
439            nodes = Retryer(
440                self._get_nodes,
441                max_tries=CFG.runner.retry_tries,
442                wait_seconds=CFG.runner.retry_wait,
443            ).run()
444
445            self._informer.set_running(
446                datetime.now(),
447                socket.getfqdn(),
448                nodes,
449                self._work_dir,
450            )
451
452            Retryer(
453                self._informer.to_file,
454                self._info_file,
455                host=self._input_machine,
456                max_tries=CFG.runner.retry_tries,
457                wait_seconds=CFG.runner.retry_wait,
458            ).run()
459        except Exception as e:
460            raise QQError(
461                f"Could not update qqinfo file '{self._info_file}' at JOB START: {e}"
462            ) from e
463
464    def _get_nodes(self) -> list[str]:
465        """
466        Get a list of nodes used to execute this job. The nodes are obtained by
467        querying the batch system.
468
469        If the batch server is not available and only one node was requested, uses
470        `socket.getfqdn()` instead and prints warning.
471
472        Returns:
473            list[str]: Names of nodes used to execute the job.
474
475        Raises:
476            QQError: If the batch system is unable to provide information about the nodes after retries
477                and more than one node is used.
478        """
479        nodes = self._informer.get_nodes()
480        if not nodes:
481            # if the batch server is not reachable but the requested number of nodes is one,
482            # we assume that only one node is actually being used and Runner thus runs on this node
483            # we can then get the node name from socket
484            # this avoids issues with occasional inaccessibility of the batch server in
485            # the unstable Metacentrum environment
486            if self._informer.info.resources.nnodes == 1:
487                node = socket.getfqdn()
488                logger.warning(
489                    f"Could not get the list of used nodes from the batch server. Assuming the only node is the current node '{node}'."
490                )
491                return [node]
492
493            raise QQError("Could not get the list of used nodes from the batch server")
494
495        return nodes
496
497    def _update_info_finished(self) -> None:
498        """
499        Update the qq info file to mark the job as successfully finished.
500
501        Logs errors as warnings if updating fails.
502
503        Raises:
504            QQRunCommunicationError: If the job was killed without informing Runner.
505        """
506        logger.debug(f"Updating '{self._info_file}' at job completion.")
507        self._reload_info_and_ensure_valid()
508
509        try:
510            self._informer.set_finished(datetime.now())
511            Retryer(
512                self._informer.to_file,
513                self._info_file,
514                host=self._input_machine,
515                max_tries=CFG.runner.retry_tries,
516                wait_seconds=CFG.runner.retry_wait,
517            ).run()
518        except Exception as e:
519            logger.warning(
520                f"Could not update qqinfo file '{self._info_file}' at JOB COMPLETION: {e}."
521            )
522
523    def _update_info_failed(self, return_code: int) -> None:
524        """
525        Update the qq info file to mark the job as failed.
526
527        Args:
528            return_code (int): Exit code from the failed job.
529
530        Logs errors as warnings if updating fails.
531
532        Raises:
533            QQRunCommunicationError: If the job was killed without informing Runner.
534        """
535        logger.debug(f"Updating '{self._info_file}' at job failure.")
536        self._reload_info_and_ensure_valid()
537
538        try:
539            self._informer.set_failed(datetime.now(), return_code)
540            Retryer(
541                self._informer.to_file,
542                self._info_file,
543                host=self._input_machine,
544                max_tries=CFG.runner.retry_tries,
545                wait_seconds=CFG.runner.retry_wait,
546            ).run()
547        except Exception as e:
548            logger.warning(
549                f"Could not update qqinfo file '{self._info_file}' at JOB FAILURE: {e}."
550            )
551
552    def _update_info_killed(self) -> None:
553        """
554        Update the qq info file to mark the job as killed.
555
556        Used during SIGTERM cleanup.
557
558        Logs errors as warnings if updating fails.
559
560        No retrying since there is no time for that.
561        """
562        logger.debug(f"Updating '{self._info_file}' at job kill.")
563        self._reload_info_and_ensure_valid(retry=False)
564
565        try:
566            self._informer.set_killed(datetime.now())
567            # no retrying here since we cannot afford multiple attempts here
568            self._informer.to_file(self._info_file, host=self._input_machine)
569        except Exception as e:
570            logger.warning(
571                f"Could not update qqinfo file '{self._info_file}' at JOB KILL: {e}."
572            )
573
574    def _copy_runtime_files_to_input_dir(self, retry: bool = True) -> None:
575        """
576        Copy .out and .err runtime files from the working directory to the input directory.
577
578        Args:
579            retry (bool): Retry the copying if it fails.
580
581        Raises:
582            QQError: If the files could not be copied after retrying.
583        """
584        files_to_copy = [
585            logical_resolve(Path(self._informer.info.stdout_file)),
586            logical_resolve(Path(self._informer.info.stderr_file)),
587        ]
588
589        logger.debug(f"Copying runtime files '{files_to_copy}' to input directory.")
590
591        if retry:
592            Retryer(
593                self._batch_system.sync_selected,
594                self._work_dir,
595                self._input_dir,
596                socket.getfqdn(),
597                self._informer.info.input_machine,
598                include_files=files_to_copy,
599                max_tries=CFG.runner.retry_tries,
600                wait_seconds=CFG.runner.retry_wait,
601            ).run()
602        else:
603            self._batch_system.sync_selected(
604                self._work_dir,
605                self._input_dir,
606                socket.getfqdn(),
607                self._informer.info.input_machine,
608                files_to_copy,
609            )
610
611    def _reload_info(self, retry: bool = True) -> None:
612        """
613        Reload the qq job info file for this job.
614
615        Args:
616            retry (bool): Retry the loading operation if it fails.
617
618        Raises:
619            QQError: If the qq info file cannot be reach or read after retrying.
620        """
621        if retry:
622            self._informer = Retryer(
623                Informer.from_file,
624                self._info_file,
625                host=self._input_machine,
626                max_tries=CFG.runner.retry_tries,
627                wait_seconds=CFG.runner.retry_wait,
628            ).run()
629        else:
630            self._informer = Informer.from_file(self._info_file, self._input_machine)
631
632    def _ensure_matches_job(self, job_id: str) -> None:
633        """
634        Ensure that the provided job_id matches the job id in the wrapped informer.
635
636        Raises:
637            QQJobMismatchError: If the info file corresponds to a different job.
638        """
639        if not self._informer.matches_job(job_id):
640            raise QQJobMismatchError(
641                f"Info file '{self._info_file}' does not correspond to job '{job_id}'"
642            )
643
644    def _ensure_not_killed(self) -> None:
645        """
646        Ensure that the job has not been killed.
647
648        Raises:
649            QQRunCommunicationError: If the job state is `KILLED`.
650        """
651        if self._informer.info.job_state == NaiveState.KILLED:
652            raise QQRunCommunicationError(
653                "Job has been killed without informing qq run. Aborting the job"
654            )
655
656    def _reload_info_and_ensure_valid(self, retry: bool = False) -> None:
657        """
658        Reload the qq job info file and check that it corresponds to the current job
659        by comparing job ids.
660
661        Then check the job's state and ensure it is not killed.
662
663        Args:
664            retry (bool): Retry the loading operation if it fails.
665
666        Raises:
667            QQJobMismatchError: If the info file corresponds to a different job.
668            QQRunCommunicationError: If the job state is `KILLED`.
669            QQError: If the qq info file cannot be reached or read.
670        """
671        job_id = self._informer.info.job_id
672        self._reload_info(retry)
673        self._ensure_matches_job(job_id)
674        self._ensure_not_killed()
675
676    def _resubmit(self) -> None:
677        """
678        Resubmit the current job if either of the following is true:
679            a) it is a loop job and additional cycles remain,
680            b) it is a continuous job that should be resubmitted.
681
682        Raises:
683            QQError: If the job cannot be resubmitted.
684        """
685        if not self._should_resubmit:
686            logger.info(
687                f"The script finished with an exit code of '{CFG.exit_codes.qq_run_no_resubmit}' indicating that the next cycle of the job should not be submitted. Not resubmitting."
688            )
689            return
690
691        if self._informer.info.job_type == JobType.LOOP:
692            if not (loop_info := self._informer.info.loop_info):
693                raise QQError(
694                    "Loop info is undefined while resubmiting a loop job. This is a bug, please report it"
695                )
696                return
697
698            if loop_info.current >= loop_info.end:
699                logger.info(
700                    "This was the final cycle of the loop job. Not resubmitting."
701                )
702                return
703
704        logger.info("Resubmitting the job.")
705        resubmitter = Resubmitter.from_informer(self._informer)
706        job_id = resubmitter.resubmit()
707
708        logger.info(f"Job resubmitted successfully as '{job_id}'.")
709
710    def _archive_files_from_work_dir(self) -> None:
711        """
712        Archive files from the working directory.
713
714        If no file exists for the next loop cycle, creates an empty init file to ensure the loop job continues normally.
715        """
716        if not self._archiver:
717            raise QQError(
718                "Archiver is undefined while archiving files. This is a bug, please report it"
719            )
720
721        if not (loop_info := self._informer.info.loop_info):
722            raise QQError(
723                "Loop info is undefined while archiving files. This is a bug, please report it"
724            )
725
726        # get the files to archive corresponding to the next loop job cycle
727        files_matching_pattern = self._archiver.get_files_matching_pattern(
728            self._work_dir,
729            None,
730            loop_info.archive_format,
731            loop_info.current + 1,
732            False,
733        )
734        exclude = self._get_excluded_from_input_dir()
735        logger.debug(f"Files excluded from archiving: {exclude}.")
736        files = [f for f in files_matching_pattern if f not in exclude]
737
738        if not files:
739            # if there are no files matching the next loop job cycle
740            # (which are not excluded via `--include` or `--ignore` options)
741            # create an empty .init file
742            # so that the loop job continues normally
743            logger.debug(
744                f"Creating .init file for loop job cycle {loop_info.current + 1}."
745            )
746            self._archiver.create_init_file(loop_info.current + 1)
747
748        # archive all files matching the archive format
749        # note that if the .init file created in the previous block of code
750        # is in a list of ignored files, it will not be included in the archive
751        # thus, its creation is pointless
752        self._archiver.to_archive(self._work_dir)
753
754    def _copy_files(self, files: list[Path]):
755        """
756        Copy files and directories using the provided absolute paths to the working directory.
757        """
758        for file in files:
759            # we rsync each file or directory individually because each file can be provided in a different directory
760            # this may be very slow if there is a large amount of files/directories to include
761            self._batch_system.sync_selected(
762                file.parent,
763                self._work_dir,
764                self._informer.info.input_machine,
765                socket.getfqdn(),
766                [file],
767            )
768
769    def _get_excluded_from_work_dir(self) -> list[Path]:
770        """
771        Return paths that must not be copied to the working directory.
772
773        Collects the files excluded and ignored by the user, the qq info file,
774        the qq output file, and the archive if the job is a loop job.
775        Duplicates are removed, preserving the order of first occurrence.
776
777        Returns:
778            list[Path]: Paths that should not be copied to the working directory.
779        """
780        info = self._informer.info
781
782        qq_out = (info.input_dir / info.job_name).with_suffix(CFG.suffixes.qq_out)
783
784        excluded = [
785            *info.excluded_files,
786            *info.ignored_files,
787            self._info_file,
788            qq_out,
789        ]
790
791        if self._archiver:
792            excluded.append(self._archiver.archive)
793
794        return list(dict.fromkeys(excluded))
795
796    def _get_excluded_from_input_dir(self) -> list[Path]:
797        """
798        Return paths that must not be copied to the input directory.
799
800        Collects explicitly included files and ignored files, and the
801        archive if the job is a loop job. Duplicates are removed,
802        preserving the order of first occurrence.
803
804        Returns:
805            list[Path]: Paths that should not be copied to the input directory.
806        """
807        info = self._informer.info
808
809        excluded = [*info.included_files, *info.ignored_files]
810        if self._archiver:
811            excluded.append(self._archiver.archive)
812
813        return list(dict.fromkeys(relocate_by_name(excluded, self._work_dir)))
814
815    def _cleanup(self) -> None:
816        """
817        Clean up after execution is interrupted or killed.
818
819        - Copies .out and .err file to the input directory.
820        - Marks job as killed in the info file.
821        - Terminates the subprocess.
822        """
823        # update the qq info file
824        self._update_info_killed()
825
826        # send SIGTERM to the running process, if there is any
827        # this may potentially not even be called -- the subprocess might be already terminated
828        if self._process and self._process.poll() is None:
829            logger.info("Cleaning up: terminating subprocess.")
830            self._process.terminate()
831
832            # wait for the subprocess to exit, then SIGKILL it
833            sleep(CFG.runner.sigterm_to_sigkill)
834            if self._process and self._process.poll() is None:
835                self._process.kill()
836
837        # copy runtime files to input dir without retrying
838        if self._use_scratch:
839            self._copy_runtime_files_to_input_dir(retry=False)
840
841    def _handle_sigterm(self, _signum: int, _frame: FrameType | None) -> NoReturn:
842        """
843        Signal handler for SIGTERM.
844
845        Performs cleanup, logs termination, and exits.
846        """
847        logger.info("Received SIGTERM, initiating shutdown.")
848        self._cleanup()
849        logger.error("Execution was terminated by SIGTERM.")
850        # this may get ignored by the batch system
851        # so you should not rely on this specific exit code
852        sys.exit(143)

Manages the setup, execution, and cleanup of scripts within the qq batch environment.

The Runner class is responsible for:
  • Preparing a working directory (shared or scratch space)
  • Executing a provided job script
  • Updating the job info file with run state, success, or failure
  • Cleaning up resources when execution is finished
Runner(info_file: pathlib._local.Path, host: str)
 52    def __init__(self, info_file: Path, host: str):
 53        """
 54        Initialize a new Runner instance.
 55
 56        Args:
 57            info_file (Path): Path to the qq info file that contains job metadata.
 58            host (str): The hostname of the input machine from which the job was submitted.
 59
 60        Raises:
 61            QQRunFatalError: If loading the QQ info file fails fatally during initialization.
 62        """
 63        # install a signal handler
 64        signal.signal(signal.SIGTERM, self._handle_sigterm)
 65
 66        # process running the wrapped script
 67        self._process: subprocess.Popen[str] | None = None
 68
 69        self._info_file = Path(info_file)
 70        logger.debug(f"Info file: '{self._info_file}'.")
 71
 72        self._input_machine = host
 73        logger.debug(f"Input machine: '{self._input_machine}'.")
 74
 75        # load the info file or raise a fatal qq error if this fails
 76        try:
 77            # get the batch system from the environment variable (or guess it)
 78            self._batch_system = BatchInterface.from_env_var_or_guess()
 79            logger.debug(f"Batch system: {str(self._batch_system)}.")
 80
 81            # get the id of the job from the batch system
 82            if not (job_id := self._batch_system.get_job_id()):
 83                raise QQError("Job has no associated job id")
 84
 85            # load the info file
 86            self._informer: Informer = Retryer(
 87                Informer.from_file,
 88                self._info_file,
 89                host=self._input_machine,
 90                max_tries=CFG.runner.retry_tries,
 91                wait_seconds=CFG.runner.retry_wait,
 92            ).run()
 93
 94            # check that the id of this job matches the job id in the info file
 95            if not self._informer.matches_job(job_id):
 96                raise QQJobMismatchError(
 97                    "Info file does not correspond to the current job"
 98                )
 99
100            # check that the batch system in info file matches the one loaded from the environment variable
101            if self._batch_system != self._informer.batch_system:
102                raise QQError(
103                    f"Batch system mismatch - env var: '{str(self._batch_system)}', info file: '{self._informer.batch_system}'"
104                )
105
106        except Exception as e:
107            raise QQRunFatalError(
108                f"Unable to load valid qq info file '{self._info_file}' on '{self._input_machine}': {e}"
109            ) from e
110
111        logger.info(
112            f"[qq-{str(self._batch_system)} v{qq_lib.__version__}] Initializing "
113            f"job '{self._informer.info.job_id}' on host '{socket.getfqdn()}'."
114        )
115
116        # get input directory
117        self._input_dir = Path(self._informer.info.input_dir)
118        logger.debug(f"Input directory: {self._input_dir}.")
119
120        # should the scratch directory be used?
121        self._use_scratch = self._informer.uses_scratch()
122        logger.debug(f"Use scratch: {self._use_scratch}.")
123
124        # initialize archiver, if this is a loop job
125        if loop_info := self._informer.info.loop_info:
126            self._archiver = Archiver(
127                archive=loop_info.archive,
128                archive_format=loop_info.archive_format,
129                input_machine=self._informer.info.input_machine,
130                input_dir=self._informer.info.input_dir,
131                batch_system=self._batch_system,
132                included_files=self._informer.info.included_files,
133                excluded_files=self._informer.info.excluded_files,
134                ignored_files=self._informer.info.ignored_files,
135            )
136            self._should_resubmit = True
137        else:
138            self._archiver = None
139
140        if self._informer.info.job_type == JobType.CONTINUOUS:
141            self._should_resubmit = True

Initialize a new Runner instance.

Arguments:
  • info_file (Path): Path to the qq info file that contains job metadata.
  • host (str): The hostname of the input machine from which the job was submitted.
Raises:
  • QQRunFatalError: If loading the QQ info file fails fatally during initialization.
def prepare(self) -> None:
143    def prepare(self) -> None:
144        """
145        Prepare the script for execution, setting up the archive
146        and archiving run time files (if this is a loop job) and
147        preparing working directory.
148
149        Raises:
150            QQError: If working directory setup fails.
151        """
152        if self._archiver:
153            assert self._informer.info.loop_info is not None
154            # prepare the directory for archiving
155            self._archiver.make_archive_dir()
156
157            # archive runtime files from the previous cycle
158            # this has to be done before the working directory is prepared,
159            # otherwise the runtime files would get copied to the working directory
160            logger.debug(
161                f"Archiving run time files from cycle {self._informer.info.loop_info.current - 1}."
162            )
163            self._archiver.archive_runtime_files(
164                # we need to escape the '+' character
165                construct_loop_job_name(
166                    self._informer.info.script_name,
167                    self._informer.info.loop_info.current - 1,
168                ).replace("+", "\\+"),
169                self._informer.info.loop_info.current - 1,
170            )
171
172        if self._use_scratch:
173            self._set_up_scratch_dir()
174        else:
175            self._set_up_shared_dir()
176
177        if self._archiver:
178            assert self._informer.info.loop_info is not None
179            # fetch files for the current cycle of the loop job from the archive
180            self._archiver.from_archive(
181                self._work_dir, self._informer.info.loop_info.current
182            )

Prepare the script for execution, setting up the archive and archiving run time files (if this is a loop job) and preparing working directory.

Raises:
  • QQError: If working directory setup fails.
def execute(self) -> int:
184    def execute(self) -> int:
185        """
186        Execute the job script in the working directory.
187
188        Returns:
189            int: The exit code from the executed script.
190
191        Raises:
192            QQError: If execution fails or info file cannot be updated.
193        """
194        # update the qqinfo file
195        self._update_info_running()
196
197        # get the actual name of the script to execute
198        script = logical_resolve(Path(self._informer.info.script_name))
199
200        # get paths to output files
201        stdout_log = self._informer.info.stdout_file
202        stderr_log = self._informer.info.stderr_file
203
204        logger.info(f"Executing script '{script}'.")
205
206        # get intepreter if configured, otherwise use the default
207        interpreter = self._informer.info.interpreter or Interpreter()
208
209        # get the command to execute
210        command_list = [*interpreter.to_command_list(), str(script)]
211        logger.debug(f"Command executed using subprocess.Popen: {command_list}")
212
213        try:
214            with Path(stdout_log).open("w") as out, Path(stderr_log).open("w") as err:
215                self._process = subprocess.Popen(
216                    command_list,
217                    stdout=out,
218                    stderr=err,
219                    text=True,
220                )
221
222                # wait for the process to finish in a non-blocking manner
223                while self._process.poll() is None:
224                    sleep(CFG.runner.subprocess_checks_wait_time)
225
226        except Exception as e:
227            raise QQError(f"Failed to execute script '{script}': {e}") from e
228
229        # if the script returns an exit code corresponding to CFG.exit_codes.qq_run_no_resubmit,
230        # do not submit the next cycle of the job but return 0
231        if (
232            self._informer.info.job_type in [JobType.LOOP, JobType.CONTINUOUS]
233            and self._process.returncode == CFG.exit_codes.qq_run_no_resubmit
234        ):
235            logger.debug(
236                f"Detected an exit code of '{self._process.returncode}'. Replacing with '0' and will not submit the next cycle of the job."
237            )
238            self._process.returncode = 0
239            self._should_resubmit = False
240
241        return self._process.returncode

Execute the job script in the working directory.

Returns:

int: The exit code from the executed script.

Raises:
  • QQError: If execution fails or info file cannot be updated.
def finalize(self) -> None:
243    def finalize(self) -> None:
244        """
245        Finalize the execution of the job script.
246
247        Handles post-processing of the job based on the script's exit code and the
248        configured transfer and archive modes. The specific actions taken depend on
249        the job's transfer mode, archive mode, and whether scratch directory is being used.
250
251        Specifically, this method:
252
253        1. Archives files from the working directory if archiving is enabled and the
254            archive mode allows it for the given exit code (loop jobs only).
255        2. Transfers or handles files based on whether scratch directory is used:
256            - If using scratch and transfer mode allows: Syncs the entire working
257                directory back to the input directory (excluding explicitly included files)
258                and removes the working directory from scratch.
259            - If using scratch and transfer mode disallows: Copies only runtime files
260                to the input directory and preserves the working directory.
261            - If not using scratch: No file operations are performed.
262        3. Updates the qq info file to "finished" (exit code 0) or "failed" (non-zero
263            exit code).
264        4. Resubmits the job if it is a loop or continuous job and was completed successfully.
265
266        Raises:
267            QQError: If copying, deletion, or archiving of files fails or if the resubmission fails.
268        """
269        logger.info("Finalizing the execution.")
270        assert self._process is not None
271
272        # archive files
273        if self._archiver and self._informer.should_archive_files(
274            self._process.returncode
275        ):
276            logger.debug(
277                f"Script exit code is '{self._process.returncode}'. Archiving files."
278            )
279            self._archive_files_from_work_dir()
280
281        # transfer files back to the input (submission) directory
282        if self._use_scratch:
283            if self._informer.should_transfer_files(self._process.returncode):
284                logger.debug(
285                    f"Script exit code is '{self._process.returncode}'. Transferring files from working directory."
286                )
287
288                Retryer(
289                    self._batch_system.sync_with_exclusions,
290                    self._work_dir,
291                    self._input_dir,
292                    socket.getfqdn(),
293                    self._informer.info.input_machine,
294                    # exclude files that were specifically included via the `--include` option
295                    # and files that were specifically chosen to be ignored via `--ignore` option
296                    self._get_excluded_from_input_dir(),
297                    max_tries=CFG.runner.retry_tries,
298                    wait_seconds=CFG.runner.retry_wait,
299                ).run()
300
301                # remove the working directory from scratch
302                self._delete_work_dir()
303            else:
304                # copy only the runtime files to input directory
305                # and keep the working directory
306                self._copy_runtime_files_to_input_dir(retry=True)
307
308        if self._process.returncode == 0:
309            # update the qqinfo file
310            self._update_info_finished()
311
312            # if this is a loop/continuous job
313            if self._informer.info.job_type in [JobType.LOOP, JobType.CONTINUOUS]:
314                self._resubmit()
315        else:
316            # update the qqinfo file
317            self._update_info_failed(self._process.returncode)
318
319        logger.info(f"Job completed with an exit code of {self._process.returncode}.")

Finalize the execution of the job script.

Handles post-processing of the job based on the script's exit code and the configured transfer and archive modes. The specific actions taken depend on the job's transfer mode, archive mode, and whether scratch directory is being used.

Specifically, this method:

  1. Archives files from the working directory if archiving is enabled and the archive mode allows it for the given exit code (loop jobs only).
  2. Transfers or handles files based on whether scratch directory is used:
    • If using scratch and transfer mode allows: Syncs the entire working directory back to the input directory (excluding explicitly included files) and removes the working directory from scratch.
    • If using scratch and transfer mode disallows: Copies only runtime files to the input directory and preserves the working directory.
    • If not using scratch: No file operations are performed.
  3. Updates the qq info file to "finished" (exit code 0) or "failed" (non-zero exit code).
  4. Resubmits the job if it is a loop or continuous job and was completed successfully.
Raises:
  • QQError: If copying, deletion, or archiving of files fails or if the resubmission fails.
def log_failure_and_exit(self, exception: BaseException) -> NoReturn:
321    def log_failure_and_exit(self, exception: BaseException) -> NoReturn:
322        """
323        Record a failure state into the qq info file and exit the program.
324
325        Args:
326            exception (BaseException): The exception to log.
327
328        Raises:
329            SystemExit: Always exits with the exit code associated with the given exception.
330        """
331        exit_code = getattr(exception, "exit_code", CFG.exit_codes.unexpected_error)
332        try:
333            self._update_info_failed(exit_code)
334            logger.error(terminate(str(exception)))
335            sys.exit(exit_code)
336        except Exception as e:
337            # unable to log the current state into the info file
338            log_fatal_error_and_exit(e)  # exits here

Record a failure state into the qq info file and exit the program.

Arguments:
  • exception (BaseException): The exception to log.
Raises:
  • SystemExit: Always exits with the exit code associated with the given exception.