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]
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
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.
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.
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.
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:
- Archives files from the working directory if archiving is enabled and the archive mode allows it for the given exit code (loop jobs only).
- 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.
- Updates the qq info file to "finished" (exit code 0) or "failed" (non-zero exit code).
- 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.
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.