qq_lib.submit
Utilities for submitting qq jobs.
This module integrates three main components - Parser, Submitter, and
SubmitterFactory - that collectively interpret submission settings, construct
job metadata, and hand off execution to the batch system.
Parser extracts qq directives declared inside the script (via # qq ...
lines) and normalizes them into structured submission parameters such as
resources, dependencies, file include/exclude rules and loop-job fields.
Submitter validates the script, prevents accidental duplicate submissions,
constructs the qq info file, sets up environment variables needed by qq run,
and finally invokes the batch system's submission mechanism.
SubmitterFactory coordinates command-line arguments with script-embedded
directives, merges and resolves resources, determines the batch system and
queue, constructs loop-job settings, and ultimately produces a fully configured
Submitter. It ensures a consistent and unified interpretation of submission
parameters from all available sources.
1# Released under MIT License. 2# Copyright (c) 2025-2026 Ladislav Bartos and Robert Vacha Lab 3 4""" 5Utilities for submitting qq jobs. 6 7This module integrates three main components - `Parser`, `Submitter`, and 8`SubmitterFactory` - that collectively interpret submission settings, construct 9job metadata, and hand off execution to the batch system. 10 11`Parser` extracts qq directives declared inside the script (via `# qq ...` 12lines) and normalizes them into structured submission parameters such as 13resources, dependencies, file include/exclude rules and loop-job fields. 14 15`Submitter` validates the script, prevents accidental duplicate submissions, 16constructs the qq info file, sets up environment variables needed by `qq run`, 17and finally invokes the batch system's submission mechanism. 18 19`SubmitterFactory` coordinates command-line arguments with script-embedded 20directives, merges and resolves resources, determines the batch system and 21queue, constructs loop-job settings, and ultimately produces a fully configured 22`Submitter`. It ensures a consistent and unified interpretation of submission 23parameters from all available sources. 24""" 25 26from .factory import SubmitterFactory 27from .parser import Parser 28from .submitter import Submitter 29 30__all__ = ["SubmitterFactory", "Parser", "Submitter"]
27class SubmitterFactory: 28 """ 29 Factory class to construct a Submitter instance based on parameters from 30 the command-line and from the script itself. 31 """ 32 33 def __init__(self, script: Path, **kwargs): 34 """ 35 Initialize the factory with the script, command-line parameters, and additional options. 36 37 Args: 38 script (Path): Path to the script to submit. 39 **kwargs: Keyword arguments from the command line. 40 """ 41 from qq_lib.submit.cli import submit 42 43 self._parser = Parser(script, submit.params) 44 self._script = script 45 self._input_dir = script.parent 46 self._kwargs = kwargs 47 48 def make_submitter(self) -> Submitter: 49 """ 50 Construct and return a Submitter instance. 51 52 Returns: 53 Submitter: A fully initialized submitter object ready to submit a job. 54 55 Raises: 56 QQError: If required information, such as the submission queue, is missing. 57 """ 58 self._parser.parse() 59 60 BatchSystem = self._get_batch_system() 61 queue = self._get_queue() 62 63 if (job_type := self._get_job_type()) == JobType.LOOP: 64 loop_info = self._get_loop_info() 65 else: 66 # tell the user that any loop job-specific options will be ignored 67 # because the job is not a loop job 68 self._print_warning_if_loop_info_defined(job_type) 69 loop_info = None 70 71 if job_type not in (JobType.LOOP, JobType.CONTINUOUS): 72 # tell the user that 'resubmit_from' will be ignored if the 73 # job type is not 'loop' or 'continuous' 74 self._print_warning_if_resubmit_from_defined(job_type) 75 76 server = self._get_server() 77 78 return Submitter( 79 batch_system=BatchSystem, 80 queue=queue, 81 account=self._get_account(), 82 script=self._script, 83 job_type=job_type, 84 resources=self._get_resources(BatchSystem, queue, server), 85 loop_info=loop_info, 86 exclude=self._get_exclude(), 87 include=self._get_include(), 88 ignore=self._get_ignore(), 89 depend=self._get_depend(), 90 transfer_mode=self._get_transfer_mode(), 91 server=server, 92 interpreter=self._get_interpreter(), 93 resubmit_from=self._get_resubmit_from(BatchSystem) 94 if job_type in {JobType.LOOP, JobType.CONTINUOUS} 95 else None, 96 ) 97 98 def _get_batch_system(self) -> type[BatchInterface]: 99 """ 100 Determine which batch system to use for the job submission. 101 102 Priority: 103 1. Command-line specification 104 2. Batch system specified in the script 105 3. Environment variable 106 4. Guessed batch system 107 108 Returns: 109 type[BatchInterface]: The selected batch system class. 110 """ 111 if batch_system := self._kwargs.get("batch_system"): 112 return BatchInterface.from_str(batch_system) 113 return self._parser.get_batch_system() or BatchInterface.from_env_var_or_guess() 114 115 def _get_job_type(self) -> JobType: 116 """ 117 Determine the type of job to submit. 118 119 Priority: 120 1. Command-line specification 121 2. Job type specified in the script 122 3. Default to `JobType.STANDARD` 123 124 Returns: 125 JobType: The determined job type. 126 """ 127 if job_type := self._kwargs.get("job_type"): 128 return JobType.from_str(job_type) 129 return self._parser.get_job_type() or JobType.STANDARD 130 131 def _get_queue(self) -> str: 132 """ 133 Determine the submission queue to use. 134 135 Priority: 136 1. Command-line specification 137 2. Queue specified in the script 138 139 Returns: 140 str: Name of the submission queue. 141 142 Raises: 143 QQError: If no queue is specified either in kwargs or in the script. 144 """ 145 if not (queue := self._kwargs.get("queue") or self._parser.get_queue()): 146 raise QQError("Submission queue not specified") 147 return queue 148 149 def _get_resources( 150 self, BatchSystem: type[BatchInterface], queue: str, server: str | None 151 ) -> Resources: 152 """ 153 Get the resource requirements for the job by merging the requirements specified on the command 154 line with requirements specified inside the submitted script. 155 156 The resources are then further modified to conform to the provided `BatchSystem` and submission `queue`. 157 158 Args: 159 BatchSystem (type[BatchInterface]): The batch system class to use. 160 queue (str): The submission queue. 161 server (str | None): The submission server. `None` = the current main server. 162 163 Returns: 164 Resources: A merged Resources object containing the final resource requirements. 165 """ 166 field_names = {f.name for f in fields(Resources)} 167 command_line_resources = Resources( 168 **{k: v for k, v in self._kwargs.items() if k in field_names} 169 ) 170 171 return BatchSystem.transform_resources( 172 queue, 173 server, 174 Resources.merge_resources( 175 command_line_resources, self._parser.get_resources() 176 ), 177 ) 178 179 def _get_loop_info(self) -> LoopInfo: 180 """ 181 Construct LoopInfo holding information about the loop job. 182 183 Returns: 184 LoopInfo: An object containing loop job parameters. 185 186 Raises: 187 QQError: If required loop job parameters are missing or invalid. 188 """ 189 return LoopInfo( 190 self._kwargs.get("loop_start") or self._parser.get_loop_start() or 1, 191 self._kwargs.get("loop_end") or self._parser.get_loop_end(), 192 self._input_dir 193 / ( 194 self._kwargs.get("archive") 195 or self._parser.get_archive() 196 or CFG.loop_jobs.archive_dir 197 ), 198 self._kwargs.get("archive_format") 199 or self._parser.get_archive_format() 200 or CFG.loop_jobs.archive_format, 201 input_dir=self._input_dir, 202 archive_mode=TransferMode.multi_from_str( 203 self._kwargs.get("archive_mode") or "" 204 ) 205 or self._parser.get_archive_mode(), 206 ) 207 208 def _print_warning_if_loop_info_defined(self, job_type: JobType) -> None: 209 """ 210 Print warning(s) if any of the loop-job specific options 211 are defined either on the command line or in the script itself. 212 This should only be used if the job is not a loop job. 213 """ 214 for option, parser_func in zip( 215 ["loop_start", "loop_end", "archive", "archive_format", "archive_mode"], 216 [ 217 self._parser.get_loop_start, 218 self._parser.get_loop_end, 219 self._parser.get_archive, 220 self._parser.get_archive_format, 221 self._parser.get_archive_mode, 222 ], 223 ): 224 if self._kwargs.get(option) or parser_func(): 225 logger.warning( 226 f"Option '{option}' is specified but job type is '{str(job_type)}', not 'loop' - '{option}' will be ignored." 227 ) 228 229 def _print_warning_if_resubmit_from_defined(self, job_type: JobType) -> None: 230 """ 231 Print a warning if the 'resubmit_from' option is defined but the job type is not 'loop' or 'continuous'. 232 """ 233 if self._kwargs.get("resubmit_from") or self._parser.get_resubmit_from(): 234 logger.warning( 235 f"Option 'resubmit_from' is specified but job type is '{str(job_type)}', not 'loop' or 'continuous' - 'resubmit_from' will be ignored." 236 ) 237 238 def _get_exclude(self) -> list[str]: 239 """ 240 Determine the files to exclude from being copied to the job's working directory. 241 242 Priority: 243 1. Excluded files specified on the command line. 244 2. Excluded files specified inside the submitted script. 245 246 The lists are NOT merged. 247 248 Returns: 249 list[str]: List of files or glob patterns to exclude. 250 """ 251 return ( 252 split_string_list(self._kwargs.get("exclude")) or self._parser.get_exclude() 253 ) 254 255 def _get_include(self) -> list[str]: 256 """ 257 Determine the files to explicitly copy to the job's working directory. 258 259 Priority: 260 1. Included files specified on the command line. 261 2. Included files specified inside the submitted script. 262 263 The lists are NOT merged. 264 265 Returns: 266 list[str]: List of files or glob patterns to include. 267 """ 268 return ( 269 split_string_list(self._kwargs.get("include")) or self._parser.get_include() 270 ) 271 272 def _get_ignore(self) -> list[str]: 273 """ 274 Determine the files that transfer operations should ignore completely. 275 276 Priority: 277 1. Ignored files specified on the command line. 278 2. Ignored files specified inside the submitted script. 279 280 The lists are NOT merged. 281 282 Returns: 283 list[str]: List of files or glob patterns to ignore. 284 """ 285 return ( 286 split_string_list(self._kwargs.get("ignore")) or self._parser.get_ignore() 287 ) 288 289 def _get_depend(self) -> list[Depend]: 290 """ 291 Determine the list of dependencies. 292 293 Priority: 294 1. Dependencies specified on the command line. 295 2. Dependencies specified inside the submitted script. 296 297 The lists are NOT merged. 298 299 Returns: 300 list[Depend]: List of job dependencies. 301 """ 302 return ( 303 Depend.multi_from_str(self._kwargs.get("depend") or "") 304 or self._parser.get_depend() 305 ) 306 307 def _get_account(self) -> str | None: 308 """ 309 Determine the account name to use for the job. 310 311 Returns: 312 str | None: The account name or None if not defined. 313 """ 314 return self._kwargs.get("account") or self._parser.get_account() 315 316 def _get_transfer_mode(self) -> list[TransferMode]: 317 """ 318 Determine the mode specifying when files should be 319 transferred from the working directory to the input directory. 320 321 Priority: 322 1. Transfer modes specified on the command line. 323 2. Transfer modes specified inside the submitted script. 324 325 The lists are NOT merged. 326 327 Returns: 328 list[TransferMode]: List of transfer modes. 329 """ 330 return ( 331 TransferMode.multi_from_str(self._kwargs.get("transfer_mode") or "") 332 or self._parser.get_transfer_mode() 333 ) 334 335 def _get_server(self) -> str | None: 336 """ 337 Determine the batch server to submit the job to. 338 339 Priority: 340 1. Command-line specification 341 2. Batch server specified in the script 342 3. None - the current batch server 343 344 Returns: 345 str | None: The full name of the batch server to use, 346 or `None` for the current batch server. 347 """ 348 if raw := (self._kwargs.get("server") or self._parser.get_server()): 349 return translate_server(raw) 350 351 return None 352 353 def _get_interpreter(self) -> Interpreter | None: 354 """ 355 Determine the interpreter to use for running the script. 356 357 Priority: 358 1. Command-line specification 359 2. Interpreter specified in the script 360 3. None - the default intepreter 361 362 Returns: 363 Interpreter | None: The interpreter to use for running the script 364 or `None` to use the default intepreter. 365 """ 366 if (raw := self._kwargs.get("interpreter")) is not None: 367 return Interpreter.from_str(raw) 368 369 return self._parser.get_interpreter() 370 371 def _get_resubmit_from(self, BatchSystem: AnyBatchClass) -> list[ResubmitHost]: 372 """ 373 Determine the list of resubmission hosts to be used to resubmit loop/continuous job. 374 375 Priority: 376 1. Resubmission hosts specified on the command line. 377 2. Resubmission hosts specified inside the submitted script. 378 3. Resubmission hosts specified in the configuration file. 379 4. Default resubmission hosts provided by the batch system. 380 381 The lists are NOT merged. 382 383 Args: 384 BatchSystem (AnyBatchClass): The batch system used for job submission. 385 386 Returns: 387 list[ResubmitHost]: List of resubmission hosts. 388 """ 389 return ( 390 ResubmitHost.multi_from_str(self._kwargs.get("resubmit_from") or "") 391 or self._parser.get_resubmit_from() 392 or ResubmitHost.multi_from_str(CFG.resubmitter.default_resubmit_hosts or "") 393 or BatchSystem.get_default_resubmit_hosts() 394 )
Factory class to construct a Submitter instance based on parameters from the command-line and from the script itself.
33 def __init__(self, script: Path, **kwargs): 34 """ 35 Initialize the factory with the script, command-line parameters, and additional options. 36 37 Args: 38 script (Path): Path to the script to submit. 39 **kwargs: Keyword arguments from the command line. 40 """ 41 from qq_lib.submit.cli import submit 42 43 self._parser = Parser(script, submit.params) 44 self._script = script 45 self._input_dir = script.parent 46 self._kwargs = kwargs
Initialize the factory with the script, command-line parameters, and additional options.
Arguments:
- script (Path): Path to the script to submit.
- **kwargs: Keyword arguments from the command line.
48 def make_submitter(self) -> Submitter: 49 """ 50 Construct and return a Submitter instance. 51 52 Returns: 53 Submitter: A fully initialized submitter object ready to submit a job. 54 55 Raises: 56 QQError: If required information, such as the submission queue, is missing. 57 """ 58 self._parser.parse() 59 60 BatchSystem = self._get_batch_system() 61 queue = self._get_queue() 62 63 if (job_type := self._get_job_type()) == JobType.LOOP: 64 loop_info = self._get_loop_info() 65 else: 66 # tell the user that any loop job-specific options will be ignored 67 # because the job is not a loop job 68 self._print_warning_if_loop_info_defined(job_type) 69 loop_info = None 70 71 if job_type not in (JobType.LOOP, JobType.CONTINUOUS): 72 # tell the user that 'resubmit_from' will be ignored if the 73 # job type is not 'loop' or 'continuous' 74 self._print_warning_if_resubmit_from_defined(job_type) 75 76 server = self._get_server() 77 78 return Submitter( 79 batch_system=BatchSystem, 80 queue=queue, 81 account=self._get_account(), 82 script=self._script, 83 job_type=job_type, 84 resources=self._get_resources(BatchSystem, queue, server), 85 loop_info=loop_info, 86 exclude=self._get_exclude(), 87 include=self._get_include(), 88 ignore=self._get_ignore(), 89 depend=self._get_depend(), 90 transfer_mode=self._get_transfer_mode(), 91 server=server, 92 interpreter=self._get_interpreter(), 93 resubmit_from=self._get_resubmit_from(BatchSystem) 94 if job_type in {JobType.LOOP, JobType.CONTINUOUS} 95 else None, 96 )
Construct and return a Submitter instance.
Returns:
Submitter: A fully initialized submitter object ready to submit a job.
Raises:
- QQError: If required information, such as the submission queue, is missing.
26class Parser: 27 """ 28 Parser for qq job submission options (qq directives) specified in a script. 29 """ 30 31 def __init__(self, script: Path, params: list[Parameter]): 32 """ 33 Initialize the parser. 34 35 Args: 36 script (Path): Path to the qq job script to parse. 37 params (list[Parameter]): List of click Parameter objects defining 38 valid options. Only `GroupedOption` names are considered. 39 """ 40 self._script = script 41 self._known_options = { 42 p.name 43 for p in params 44 if isinstance(p, GroupedOption) and p.name is not None 45 } 46 logger.debug( 47 f"Known options for Parser: {self._known_options} ({len(self._known_options)} options)." 48 ) 49 50 self._options: dict[str, object] = {} 51 52 def parse(self) -> None: 53 """ 54 Extract and parse `qq` options from the script. 55 56 The method processes the script line by line, skipping the first line (shebang). 57 It continues reading until it encounters a line that is not a `qq` directive, 58 is non-empty, and is not a comment. Empty or commented lines are ignored. 59 60 Each valid `qq` line is parsed into key-value pairs, normalized to `snake_case`, 61 and stored in `self._options`. 62 63 Raises: 64 QQError: If the script cannot be read, or if an option line is malformed or 65 contains an unknown option. 66 """ 67 if not self._script.is_file(): 68 raise QQError(f"Could not open '{self._script}' as a file") 69 70 with self._script.open() as f: 71 # skip the first line (shebang) 72 next(f, None) 73 74 for line in f: 75 stripped = line.strip() 76 if stripped == "": 77 logger.debug("Parser: skipping empty line.") 78 continue # skip empty lines 79 80 # check whether this is a qq command 81 if not re.match(r"#\s*qq", stripped, re.IGNORECASE): 82 if stripped.startswith("#"): 83 logger.debug(f"Parser: skipping commented line '{line}'.") 84 continue # skip commented lines 85 logger.debug(f"Parser: ending parsing at line '{line}'.") 86 break # stop parsing at other lines 87 88 # remove the leading '# qq' and split by whitespace or '=' 89 parts = Parser._strip_and_split(line) 90 if len(parts) < 2: 91 raise QQError( 92 f"Invalid qq submit option line in '{str(self._script)}': {line}" 93 ) 94 95 key, value = parts[-2], parts[-1] 96 snake_case_key = to_snake_case(key) 97 98 # handle workdir and worksize where two forms of the keyword are allowed 99 snake_case_key = snake_case_key.replace("workdir", "work_dir").replace( 100 "worksize", "work_size" 101 ) 102 103 # is this a known option? 104 if snake_case_key in self._known_options: 105 try: 106 self._options[snake_case_key] = int(value) 107 except ValueError: 108 self._options[snake_case_key] = value 109 else: 110 raise QQError( 111 f"Unknown qq submit option '{key}' in '{str(self._script)}': {line.strip()}.\nKnown options are '{' '.join(self._known_options)}'" 112 ) 113 114 logger.debug(f"Parsed options from '{self._script}': {self._options}.") 115 116 def get_batch_system(self) -> type[BatchInterface] | None: 117 """ 118 Return the batch system class specified in the script. 119 120 Returns: 121 type[BatchInterface] | None: The batch system class if specified, otherwise None. 122 """ 123 if (batch_system := self._options.get("batch_system")) is not None: 124 return BatchInterface.from_str(str(batch_system)) 125 126 return None 127 128 def get_queue(self) -> str | None: 129 """ 130 Return the queue specified for the job. 131 132 Returns: 133 str | None: Queue name, or None if not set. 134 """ 135 if (queue := self._options.get("queue")) is not None: 136 return str(queue) 137 return None 138 139 def get_job_type(self) -> JobType | None: 140 """ 141 Return the job type specified in the script. 142 143 Returns: 144 JobType | None: Enum value representing the job type, or None if not set. 145 """ 146 if (job_type := self._options.get("job_type")) is not None: 147 return JobType.from_str(str(job_type)) 148 149 return None 150 151 def get_resources(self) -> Resources: 152 """ 153 Return the job resource specifications parsed from the script. 154 155 Returns: 156 Resources: Resource requirements for the job. 157 """ 158 field_names = {f.name for f in fields(Resources)} 159 # only select fields that are part of Resources 160 return Resources(**{k: v for k, v in self._options.items() if k in field_names}) # ty: ignore[invalid-argument-type] 161 162 def get_exclude(self) -> list[str]: 163 """ 164 Determine the files to exclude from being copied to the job's working directory. 165 166 Returns: 167 list[str]: List of excluded files or glob patterns. 168 Returns an empty list if none specified. 169 """ 170 if (exclude := self._options.get("exclude")) is not None: 171 return split_string_list(str(exclude)) 172 173 return [] 174 175 def get_include(self) -> list[str]: 176 """ 177 Determine the files to explicitly copy to the job's working directory. 178 179 Returns: 180 list[Path]: List of included files or glob patterns. Returns an empty list if none specified. 181 """ 182 if (include := self._options.get("include")) is not None: 183 return split_string_list(str(include)) 184 185 return [] 186 187 def get_ignore(self) -> list[str]: 188 """ 189 Determine the files to completely ignore during transfer operations. 190 191 Returns: 192 list[Path]: List of ignored files or glob patterns. Returns an empty list if none specified. 193 """ 194 if (ignore := self._options.get("ignore")) is not None: 195 return split_string_list(str(ignore)) 196 197 return [] 198 199 def get_loop_start(self) -> int | None: 200 """ 201 Return the starting cycle number for loop jobs. 202 203 Returns: 204 int | None: Start cycle, or None if not specified. 205 """ 206 if isinstance(loop_start := self._options.get("loop_start"), int): 207 return loop_start 208 return None 209 210 def get_loop_end(self) -> int | None: 211 """ 212 Return the ending cycle number for loop jobs. 213 214 Returns: 215 int | None: End cycle, or None if not specified. 216 """ 217 if isinstance(loop_end := self._options.get("loop_end"), int): 218 return loop_end 219 return None 220 221 def get_archive(self) -> Path | None: 222 """ 223 Return the archive directory path specified in the script. 224 225 Returns: 226 Path | None: Archive directory path, or None if not set. 227 """ 228 if (archive := self._options.get("archive")) is not None: 229 return Path(str(archive)) 230 231 return None 232 233 def get_archive_format(self) -> str | None: 234 """ 235 Return the file naming format used for archived files. 236 237 Returns: 238 str | None: Archive filename format string, or None if not set. 239 """ 240 if (archive_format := self._options.get("archive_format")) is not None: 241 return str(archive_format) 242 return None 243 244 def get_archive_mode(self) -> list[TransferMode]: 245 """ 246 Get the mode specifying when the files should be archived. 247 248 Returns: 249 list[TransferMode]: List of transfer modes. 250 """ 251 if (raw := self._options.get("archive_mode")) is not None: 252 return TransferMode.multi_from_str(str(raw)) 253 254 return [] 255 256 def get_depend(self) -> list[Depend]: 257 """ 258 Return the list of job dependencies. 259 260 Returns: 261 list[Depend]: List of job dependencies. 262 """ 263 if (raw := self._options.get("depend")) is not None: 264 return Depend.multi_from_str(str(raw)) 265 266 return [] 267 268 def get_account(self) -> str | None: 269 """ 270 Get the account name to use for the job. 271 272 Returns: 273 str | None: The account name or None if not defined. 274 """ 275 if (account := self._options.get("account")) is not None: 276 return str(account) 277 278 return None 279 280 def get_transfer_mode(self) -> list[TransferMode]: 281 """ 282 Get the mode specifying when the files should be transferred 283 from the working directory to the input directory. 284 285 Returns: 286 list[TransferMode]: List of transfer modes. 287 """ 288 if (raw := self._options.get("transfer_mode")) is not None: 289 return TransferMode.multi_from_str(str(raw)) 290 291 return [] 292 293 def get_server(self) -> str | None: 294 """ 295 Get the batch server to which the job should be submitted. 296 297 Note that this function returns the raw name of the server 298 as provided by the user. It should be then translated using 299 the `translate_server` function. 300 301 Returns: 302 str | None: The name or shortcut of the batch server or `None` if not specified. 303 """ 304 if (server := self._options.get("server")) is not None: 305 return str(server) 306 307 return None 308 309 def get_interpreter(self) -> Interpreter | None: 310 """ 311 Get the interpreter that should be used to run the script. 312 313 Returns: 314 Interpreter | None: The interpreter or `None` if not specified. 315 """ 316 if (raw := self._options.get("interpreter")) is not None: 317 return Interpreter.from_str(str(raw)) 318 319 return None 320 321 def get_resubmit_from(self) -> list[ResubmitHost]: 322 """ 323 Return the list of resubmission hosts. 324 325 Returns: 326 list[ResubmitHost]: List of job dependencies. 327 """ 328 if (raw := self._options.get("resubmit_from")) is not None: 329 return ResubmitHost.multi_from_str(str(raw)) 330 331 return [] 332 333 @staticmethod 334 def _strip_and_split(string: str) -> list[str]: 335 """ 336 Remove the leading `# qq` directive from a line, extract content before the next `#` 337 (if any), and split the remaining content. 338 339 Args: 340 string (str): Input line to process. 341 342 Returns: 343 list[str]: A list with one or two elements depending on whether a split occurred. 344 """ 345 match = re.search( 346 r"^#\s*qq\s*(.*?)\s*(?:#|$)", string.strip(), flags=re.IGNORECASE 347 ) 348 content = match.group(1).strip() if match else string.strip() 349 350 # split by whitespace or '=' 351 return re.split(r"[=\s]+", content, maxsplit=1)
Parser for qq job submission options (qq directives) specified in a script.
31 def __init__(self, script: Path, params: list[Parameter]): 32 """ 33 Initialize the parser. 34 35 Args: 36 script (Path): Path to the qq job script to parse. 37 params (list[Parameter]): List of click Parameter objects defining 38 valid options. Only `GroupedOption` names are considered. 39 """ 40 self._script = script 41 self._known_options = { 42 p.name 43 for p in params 44 if isinstance(p, GroupedOption) and p.name is not None 45 } 46 logger.debug( 47 f"Known options for Parser: {self._known_options} ({len(self._known_options)} options)." 48 ) 49 50 self._options: dict[str, object] = {}
Initialize the parser.
Arguments:
- script (Path): Path to the qq job script to parse.
- params (list[Parameter]): List of click Parameter objects defining
valid options. Only
GroupedOptionnames are considered.
52 def parse(self) -> None: 53 """ 54 Extract and parse `qq` options from the script. 55 56 The method processes the script line by line, skipping the first line (shebang). 57 It continues reading until it encounters a line that is not a `qq` directive, 58 is non-empty, and is not a comment. Empty or commented lines are ignored. 59 60 Each valid `qq` line is parsed into key-value pairs, normalized to `snake_case`, 61 and stored in `self._options`. 62 63 Raises: 64 QQError: If the script cannot be read, or if an option line is malformed or 65 contains an unknown option. 66 """ 67 if not self._script.is_file(): 68 raise QQError(f"Could not open '{self._script}' as a file") 69 70 with self._script.open() as f: 71 # skip the first line (shebang) 72 next(f, None) 73 74 for line in f: 75 stripped = line.strip() 76 if stripped == "": 77 logger.debug("Parser: skipping empty line.") 78 continue # skip empty lines 79 80 # check whether this is a qq command 81 if not re.match(r"#\s*qq", stripped, re.IGNORECASE): 82 if stripped.startswith("#"): 83 logger.debug(f"Parser: skipping commented line '{line}'.") 84 continue # skip commented lines 85 logger.debug(f"Parser: ending parsing at line '{line}'.") 86 break # stop parsing at other lines 87 88 # remove the leading '# qq' and split by whitespace or '=' 89 parts = Parser._strip_and_split(line) 90 if len(parts) < 2: 91 raise QQError( 92 f"Invalid qq submit option line in '{str(self._script)}': {line}" 93 ) 94 95 key, value = parts[-2], parts[-1] 96 snake_case_key = to_snake_case(key) 97 98 # handle workdir and worksize where two forms of the keyword are allowed 99 snake_case_key = snake_case_key.replace("workdir", "work_dir").replace( 100 "worksize", "work_size" 101 ) 102 103 # is this a known option? 104 if snake_case_key in self._known_options: 105 try: 106 self._options[snake_case_key] = int(value) 107 except ValueError: 108 self._options[snake_case_key] = value 109 else: 110 raise QQError( 111 f"Unknown qq submit option '{key}' in '{str(self._script)}': {line.strip()}.\nKnown options are '{' '.join(self._known_options)}'" 112 ) 113 114 logger.debug(f"Parsed options from '{self._script}': {self._options}.")
Extract and parse qq options from the script.
The method processes the script line by line, skipping the first line (shebang).
It continues reading until it encounters a line that is not a qq directive,
is non-empty, and is not a comment. Empty or commented lines are ignored.
Each valid qq line is parsed into key-value pairs, normalized to snake_case,
and stored in self._options.
Raises:
- QQError: If the script cannot be read, or if an option line is malformed or contains an unknown option.
116 def get_batch_system(self) -> type[BatchInterface] | None: 117 """ 118 Return the batch system class specified in the script. 119 120 Returns: 121 type[BatchInterface] | None: The batch system class if specified, otherwise None. 122 """ 123 if (batch_system := self._options.get("batch_system")) is not None: 124 return BatchInterface.from_str(str(batch_system)) 125 126 return None
Return the batch system class specified in the script.
Returns:
type[BatchInterface] | None: The batch system class if specified, otherwise None.
128 def get_queue(self) -> str | None: 129 """ 130 Return the queue specified for the job. 131 132 Returns: 133 str | None: Queue name, or None if not set. 134 """ 135 if (queue := self._options.get("queue")) is not None: 136 return str(queue) 137 return None
Return the queue specified for the job.
Returns:
str | None: Queue name, or None if not set.
139 def get_job_type(self) -> JobType | None: 140 """ 141 Return the job type specified in the script. 142 143 Returns: 144 JobType | None: Enum value representing the job type, or None if not set. 145 """ 146 if (job_type := self._options.get("job_type")) is not None: 147 return JobType.from_str(str(job_type)) 148 149 return None
Return the job type specified in the script.
Returns:
JobType | None: Enum value representing the job type, or None if not set.
151 def get_resources(self) -> Resources: 152 """ 153 Return the job resource specifications parsed from the script. 154 155 Returns: 156 Resources: Resource requirements for the job. 157 """ 158 field_names = {f.name for f in fields(Resources)} 159 # only select fields that are part of Resources 160 return Resources(**{k: v for k, v in self._options.items() if k in field_names}) # ty: ignore[invalid-argument-type]
Return the job resource specifications parsed from the script.
Returns:
Resources: Resource requirements for the job.
162 def get_exclude(self) -> list[str]: 163 """ 164 Determine the files to exclude from being copied to the job's working directory. 165 166 Returns: 167 list[str]: List of excluded files or glob patterns. 168 Returns an empty list if none specified. 169 """ 170 if (exclude := self._options.get("exclude")) is not None: 171 return split_string_list(str(exclude)) 172 173 return []
Determine the files to exclude from being copied to the job's working directory.
Returns:
list[str]: List of excluded files or glob patterns. Returns an empty list if none specified.
175 def get_include(self) -> list[str]: 176 """ 177 Determine the files to explicitly copy to the job's working directory. 178 179 Returns: 180 list[Path]: List of included files or glob patterns. Returns an empty list if none specified. 181 """ 182 if (include := self._options.get("include")) is not None: 183 return split_string_list(str(include)) 184 185 return []
Determine the files to explicitly copy to the job's working directory.
Returns:
list[Path]: List of included files or glob patterns. Returns an empty list if none specified.
187 def get_ignore(self) -> list[str]: 188 """ 189 Determine the files to completely ignore during transfer operations. 190 191 Returns: 192 list[Path]: List of ignored files or glob patterns. Returns an empty list if none specified. 193 """ 194 if (ignore := self._options.get("ignore")) is not None: 195 return split_string_list(str(ignore)) 196 197 return []
Determine the files to completely ignore during transfer operations.
Returns:
list[Path]: List of ignored files or glob patterns. Returns an empty list if none specified.
199 def get_loop_start(self) -> int | None: 200 """ 201 Return the starting cycle number for loop jobs. 202 203 Returns: 204 int | None: Start cycle, or None if not specified. 205 """ 206 if isinstance(loop_start := self._options.get("loop_start"), int): 207 return loop_start 208 return None
Return the starting cycle number for loop jobs.
Returns:
int | None: Start cycle, or None if not specified.
210 def get_loop_end(self) -> int | None: 211 """ 212 Return the ending cycle number for loop jobs. 213 214 Returns: 215 int | None: End cycle, or None if not specified. 216 """ 217 if isinstance(loop_end := self._options.get("loop_end"), int): 218 return loop_end 219 return None
Return the ending cycle number for loop jobs.
Returns:
int | None: End cycle, or None if not specified.
221 def get_archive(self) -> Path | None: 222 """ 223 Return the archive directory path specified in the script. 224 225 Returns: 226 Path | None: Archive directory path, or None if not set. 227 """ 228 if (archive := self._options.get("archive")) is not None: 229 return Path(str(archive)) 230 231 return None
Return the archive directory path specified in the script.
Returns:
Path | None: Archive directory path, or None if not set.
233 def get_archive_format(self) -> str | None: 234 """ 235 Return the file naming format used for archived files. 236 237 Returns: 238 str | None: Archive filename format string, or None if not set. 239 """ 240 if (archive_format := self._options.get("archive_format")) is not None: 241 return str(archive_format) 242 return None
Return the file naming format used for archived files.
Returns:
str | None: Archive filename format string, or None if not set.
244 def get_archive_mode(self) -> list[TransferMode]: 245 """ 246 Get the mode specifying when the files should be archived. 247 248 Returns: 249 list[TransferMode]: List of transfer modes. 250 """ 251 if (raw := self._options.get("archive_mode")) is not None: 252 return TransferMode.multi_from_str(str(raw)) 253 254 return []
Get the mode specifying when the files should be archived.
Returns:
list[TransferMode]: List of transfer modes.
256 def get_depend(self) -> list[Depend]: 257 """ 258 Return the list of job dependencies. 259 260 Returns: 261 list[Depend]: List of job dependencies. 262 """ 263 if (raw := self._options.get("depend")) is not None: 264 return Depend.multi_from_str(str(raw)) 265 266 return []
Return the list of job dependencies.
Returns:
list[Depend]: List of job dependencies.
268 def get_account(self) -> str | None: 269 """ 270 Get the account name to use for the job. 271 272 Returns: 273 str | None: The account name or None if not defined. 274 """ 275 if (account := self._options.get("account")) is not None: 276 return str(account) 277 278 return None
Get the account name to use for the job.
Returns:
str | None: The account name or None if not defined.
280 def get_transfer_mode(self) -> list[TransferMode]: 281 """ 282 Get the mode specifying when the files should be transferred 283 from the working directory to the input directory. 284 285 Returns: 286 list[TransferMode]: List of transfer modes. 287 """ 288 if (raw := self._options.get("transfer_mode")) is not None: 289 return TransferMode.multi_from_str(str(raw)) 290 291 return []
Get the mode specifying when the files should be transferred from the working directory to the input directory.
Returns:
list[TransferMode]: List of transfer modes.
293 def get_server(self) -> str | None: 294 """ 295 Get the batch server to which the job should be submitted. 296 297 Note that this function returns the raw name of the server 298 as provided by the user. It should be then translated using 299 the `translate_server` function. 300 301 Returns: 302 str | None: The name or shortcut of the batch server or `None` if not specified. 303 """ 304 if (server := self._options.get("server")) is not None: 305 return str(server) 306 307 return None
Get the batch server to which the job should be submitted.
Note that this function returns the raw name of the server
as provided by the user. It should be then translated using
the translate_server function.
Returns:
str | None: The name or shortcut of the batch server or
Noneif not specified.
309 def get_interpreter(self) -> Interpreter | None: 310 """ 311 Get the interpreter that should be used to run the script. 312 313 Returns: 314 Interpreter | None: The interpreter or `None` if not specified. 315 """ 316 if (raw := self._options.get("interpreter")) is not None: 317 return Interpreter.from_str(str(raw)) 318 319 return None
Get the interpreter that should be used to run the script.
Returns:
Interpreter | None: The interpreter or
Noneif not specified.
321 def get_resubmit_from(self) -> list[ResubmitHost]: 322 """ 323 Return the list of resubmission hosts. 324 325 Returns: 326 list[ResubmitHost]: List of job dependencies. 327 """ 328 if (raw := self._options.get("resubmit_from")) is not None: 329 return ResubmitHost.multi_from_str(str(raw)) 330 331 return []
Return the list of resubmission hosts.
Returns:
list[ResubmitHost]: List of job dependencies.
39class Submitter: 40 """ 41 Class to submit jobs to a batch system. 42 43 Responsibilities: 44 - Validate that the script exists and has a proper shebang. 45 - Guard against multiple submissions from the same directory. 46 - Set environment variables required for `qq run`. 47 - Create a qq info file for tracking job state and metadata. 48 49 Note that Submitter ignores qq directives in the submitted script. 50 To handle them, you have to build a Submitter using the SubmitterFactory. 51 """ 52 53 def __init__( 54 self, 55 batch_system: AnyBatchClass, 56 queue: str, 57 account: str | None, 58 script: Path, 59 job_type: JobType, 60 resources: Resources, 61 loop_info: LoopInfo | None = None, 62 exclude: list[str] | None = None, 63 include: list[str] | None = None, 64 ignore: list[str] | None = None, 65 depend: list[Depend] | None = None, 66 transfer_mode: list[TransferMode] | None = None, 67 server: str | None = None, 68 interpreter: Interpreter | None = None, 69 resubmit_from: list[ResubmitHost] | None = None, 70 ): 71 """ 72 Initialize a Submitter instance. 73 74 Args: 75 batch_system (AnyBatchClass): The batch system class implementing 76 the BatchInterface used for job submission. 77 queue (str): The name of the batch system queue to which the job will be submitted. 78 account (str | None): The name of the account to use for the job. 79 script (Path): Path to the job script to submit. 80 job_type (JobType): Type of the job to submit (e.g. standard, loop). 81 resources (Resources): Job resource requirements (e.g., CPUs, memory, walltime). 82 loop_info (LoopInfo | None): Optional information for loop jobs. Pass None if not applicable. 83 exclude (list[str] | None): Optional list of files or glob patterns which should not be copied to the working directory. 84 Paths are provided relative to the input directory or absolute. 85 include (list[str] | None): Optional list of files or glob patterns which should be copied to the working directory 86 even though they are not part of the job's input directory. 87 Paths are provided either absolute or relative to the input directory. 88 ignore (list[str] | None): Optional list of files or glob patterns which should be ignored completely. 89 These files will not be copied to the working directory and if they are created in the working directory, they are also not copied back. 90 Paths are provided either absolute or relative to the input directory. 91 depend (list[Depend] | None): Optional list of job dependencies. 92 transfer_mode (list[TransferMode] | None): Mode specifying when files whould be transferred from the 93 working directory to the input directory. Defaults to [`Success()`]. 94 server (str | None): Optional name of the server to which the job should be submitted. 95 If `None`, the default batch server, as configured by the batch system is used. 96 intepreter (Interpreter | None): Optional interpreter specification to use to execute the script. 97 If not specified, the config default is used. 98 resubmit_from (list[ResubmitHost] | None): List of hosts from which a loop/continuous job should be resubmitted. 99 Must only be specified for loop/continuous jobs! 100 101 Raises: 102 QQError: If the script does not exist or has an invalid shebang line. 103 """ 104 105 self._batch_system = batch_system 106 self._job_type = job_type 107 self._queue = queue 108 self._server = server 109 self._account = account 110 self._loop_info = loop_info 111 self._script = script 112 self._input_dir = logical_resolve(script).parent 113 self._script_name = script.name 114 self._job_name = self._construct_job_name() 115 self._info_file = construct_info_file_path(self._input_dir, self._job_name) 116 self._resources = resources 117 self._exclude = expand_paths(exclude or [], self._input_dir) 118 self._include = expand_paths(include or [], self._input_dir) 119 self._ignore = expand_paths(ignore or [], self._input_dir) 120 self._depend = depend or [] 121 self._transfer_mode = transfer_mode or TransferMode.multi_from_str( 122 CFG.transfer_files_options.default_transfer_mode 123 ) 124 self._interpreter = interpreter 125 self._resubmit_from = resubmit_from or [] 126 127 # script must exist 128 if not self._script.is_file(): 129 raise QQError(f"Script '{script}' does not exist or is not a file") 130 131 # script must have a valid qq shebang 132 if not self._has_valid_shebang(self._script): 133 raise QQError( 134 f"Script '{self._script}' has an invalid shebang. The first line of the script should be '#!/usr/bin/env -S {CFG.binary_name} run'" 135 ) 136 137 def submit(self, remote: str | None = None) -> str: 138 """ 139 Submit the script to the batch system. 140 141 Sets required environment variables, calls the batch system's 142 job submission mechanism, and creates an info file with job metadata. 143 144 This method is thread-safe, if the submission is done from the current machine. 145 146 Args: 147 remote (str | None): Name of the machine from which the job should be submitted. 148 If `None`, the current machine is used. 149 150 Returns: 151 str: The job ID of the submitted job. 152 153 Raises: 154 QQError: If job submission fails. 155 """ 156 job_id = self._batch_system.job_submit( 157 self._resources, 158 self._queue, 159 self._script, 160 self._job_name, 161 self._depend, 162 self._create_env_vars_dict(), 163 self._account, 164 self._server, 165 remote_host=remote, 166 ) 167 168 # create job qq info file 169 # we create the info file from the current machine no matter 170 # whether we are submiting from the current machine or from the remote machine 171 # the input directory should be available on both concerned machines, 172 # so this should be okay 173 Info( 174 batch_system=self._batch_system, 175 qq_version=qq_lib.__version__, 176 username=getpass.getuser(), 177 job_id=job_id, 178 job_name=self._job_name, 179 script_name=self._script_name, 180 queue=self._queue, 181 job_type=self._job_type, 182 input_machine=socket.getfqdn(remote or ""), 183 input_dir=self._input_dir, 184 job_state=NaiveState.QUEUED, 185 submission_time=datetime.now(), 186 stdout_file=str(Path(self._job_name).with_suffix(CFG.suffixes.stdout)), 187 stderr_file=str(Path(self._job_name).with_suffix(CFG.suffixes.stderr)), 188 resources=self._resources, 189 loop_info=self._loop_info, 190 excluded_files=self._exclude, 191 included_files=self._include, 192 ignored_files=self._ignore, 193 depend=self._depend, 194 account=self._account, 195 transfer_mode=self._transfer_mode, 196 server=self._server, 197 interpreter=self._interpreter, 198 resubmit_from=self._resubmit_from, 199 ).to_file(self._info_file) 200 201 return job_id 202 203 def continues_loop(self) -> bool: 204 """ 205 Determine whether the submitted job is a continuation of a loop/continuous job. 206 207 Returns: 208 bool: True if the job is a valid continuation of a previous loop/continuous job, 209 False otherwise. 210 """ 211 try: 212 # there should only be one info file for both loop jobs (runtime files are archived) 213 # and continuous jobs (runtime files overwrite each other) 214 info_file = get_info_file(self._input_dir) 215 informer = Informer.from_file(info_file) 216 217 if self._loop_job_continues_loop( 218 informer 219 ) or self._continuous_job_continues_loop(informer): 220 logger.debug("Valid loop job with a correct cycle or a continuous job.") 221 return True 222 logger.debug( 223 "Detected info file does not correspond to a resubmittable job." 224 ) 225 return False 226 except QQError as e: 227 logger.debug(f"Could not read an info file: {e}") 228 return False 229 230 def _loop_job_continues_loop(self, previous: Informer) -> bool: 231 """ 232 Determine whether the submitted job is a continuation of a loop job. 233 234 Args: 235 previous (Informer): Informer associated with the previous job. 236 237 Returns: 238 bool: True if the job is a valid continuation of a previous loop job, False otherwise. 239 """ 240 return ( 241 # both the previous job and the current job must be loop jobs 242 previous.info.loop_info is not None 243 and self._loop_info is not None 244 # previous job must be successfully finished 245 and previous.info.job_state == NaiveState.FINISHED 246 # the cycle of the current job is one more than the cycle of the previous job 247 and previous.info.loop_info.current == self._loop_info.current - 1 248 ) 249 250 def _continuous_job_continues_loop(self, previous: Informer) -> bool: 251 """ 252 Determine whether the submitted job is a continuation of a continuous job. 253 254 Args: 255 previous (Informer): Informer associated with the previous job. 256 257 Returns: 258 bool: True if the job is a valid continuation of a previous continuous job, False otherwise. 259 """ 260 return ( 261 # both the previous and the current job must be continuous jobs 262 previous.info.job_type == JobType.CONTINUOUS 263 and self._job_type == JobType.CONTINUOUS 264 # previous job must be successfully finished 265 and previous.info.job_state == NaiveState.FINISHED 266 ) 267 268 def get_input_dir(self) -> Path: 269 """ 270 Get path to the job's input directory. 271 272 Returns: 273 Path: Path to the job's input directory. 274 """ 275 return self._input_dir 276 277 def get_batch_system(self) -> AnyBatchClass: 278 """Get the batch system used for submiting.""" 279 return self._batch_system 280 281 def get_job_name(self) -> str: 282 """Get the name of the job.""" 283 return self._job_name 284 285 def get_queue(self) -> str: 286 """Get the submission queue.""" 287 return self._queue 288 289 def get_account(self) -> str | None: 290 """Get the user's account.""" 291 return self._account 292 293 def get_script(self) -> Path: 294 """Get absolute (logical) path to the submitted script.""" 295 return self._script 296 297 def get_job_type(self) -> JobType: 298 """Get type of the job.""" 299 return self._job_type 300 301 def get_resources(self) -> Resources: 302 """Get resources requested for the job.""" 303 return self._resources 304 305 def get_loop_info(self) -> LoopInfo | None: 306 """Get loop job information.""" 307 return self._loop_info 308 309 def get_exclude(self) -> list[Path]: 310 """Get a list of excluded files.""" 311 return self._exclude 312 313 def get_include(self) -> list[Path]: 314 """Get a list of included files.""" 315 return self._include 316 317 def get_ignore(self) -> list[Path]: 318 """Get a list of ignored files.""" 319 return self._ignore 320 321 def get_depend(self) -> list[Depend]: 322 """Get the list of dependencies.""" 323 return self._depend 324 325 def get_transfer_mode(self) -> list[TransferMode]: 326 """Get the list of transfer modes.""" 327 return self._transfer_mode 328 329 def get_server(self) -> str | None: 330 """Get the submission server.""" 331 return self._server 332 333 def get_interpreter(self) -> Interpreter | None: 334 """Get the interpreter to use for running the script.""" 335 return self._interpreter 336 337 def get_resubmit_from(self) -> list[ResubmitHost] | None: 338 """Get the list of hosts to resubmit the job from.""" 339 return self._resubmit_from 340 341 def _create_env_vars_dict(self) -> dict[str, str]: 342 """ 343 Create a dictionary of environment variables provided to qq runtime. 344 345 Returns 346 dict[str, str]: Dictionary of environment variables and their values. 347 """ 348 env_vars = {} 349 350 # propagate qq debug environment 351 if os.environ.get(CFG.env_vars.debug_mode): 352 env_vars[CFG.env_vars.debug_mode] = "true" 353 354 # indicates that the job is running in a qq environment 355 env_vars[CFG.env_vars.guard] = "true" 356 357 # contains a path to the qq info file 358 env_vars[CFG.env_vars.info_file] = str(self._info_file) 359 360 # contains the name of the input host 361 env_vars[CFG.env_vars.input_machine] = socket.getfqdn() 362 363 # contains the name of the used batch system 364 env_vars[CFG.env_vars.batch_system] = str(self._batch_system) 365 366 # contains the path to the input directory 367 env_vars[CFG.env_vars.input_dir] = str(self._input_dir) 368 369 # environment variables for resources 370 nnodes = self._resources.nnodes or 1 371 if ncpus := self._resources.ncpus: 372 env_vars[CFG.env_vars.ncpus] = str(ncpus) 373 elif ncpus_per_node := self._resources.ncpus_per_node: 374 env_vars[CFG.env_vars.ncpus] = str(ncpus_per_node * nnodes) 375 else: 376 env_vars[CFG.env_vars.ncpus] = "1" 377 378 if ngpus := self._resources.ngpus: 379 env_vars[CFG.env_vars.ngpus] = str(ngpus) 380 elif ngpus_per_node := self._resources.ngpus_per_node: 381 env_vars[CFG.env_vars.ngpus] = str(ngpus_per_node * nnodes) 382 else: 383 env_vars[CFG.env_vars.ngpus] = "0" 384 385 env_vars[CFG.env_vars.nnodes] = str(nnodes) 386 env_vars[CFG.env_vars.walltime] = str( 387 hhmmss_to_duration(self._resources.walltime or "00:00:00").total_seconds() 388 / 3600 389 ) 390 391 # loop job-specific environment variables 392 if self._loop_info: 393 env_vars[CFG.env_vars.loop_current] = str(self._loop_info.current) 394 env_vars[CFG.env_vars.loop_next] = str(self._loop_info.current + 1) 395 env_vars[CFG.env_vars.loop_start] = str(self._loop_info.start) 396 env_vars[CFG.env_vars.loop_end] = str(self._loop_info.end) 397 env_vars[CFG.env_vars.archive_format] = self._loop_info.archive_format 398 env_vars[CFG.env_vars.archive_current] = self._make_pattern( 399 self._loop_info.archive_format, self._loop_info.current 400 ) 401 env_vars[CFG.env_vars.archive_next] = self._make_pattern( 402 self._loop_info.archive_format, self._loop_info.current + 1 403 ) 404 405 # loop job- or continuous job-specific environment variables 406 if self._job_type in [JobType.LOOP, JobType.CONTINUOUS]: 407 env_vars[CFG.env_vars.no_resubmit] = str(CFG.exit_codes.qq_run_no_resubmit) 408 409 return env_vars 410 411 @staticmethod 412 def _make_pattern(archive_format: str, cycle: int) -> str: 413 """ 414 Create a pattern for archived files in the specified cycle. 415 416 If the archive_format is not a printf pattern, returns an empty string. 417 418 Args: 419 archive_format (str): The provided archive format. 420 cycle (int): Cycle number to use. 421 422 Returns: 423 str: The pattern or an empty string if the archive format is not a printf pattern. 424 """ 425 if is_printf_pattern(archive_format): 426 return archive_format % cycle 427 428 return "" 429 430 def _has_valid_shebang(self, script: Path) -> bool: 431 """ 432 Verify that the script has a valid shebang for qq run. 433 434 Args: 435 script (Path): Path to the script file. 436 437 Returns: 438 bool: True if the first line starts with '#!' and ends with 'qq run'. 439 """ 440 with Path.open(script) as file: 441 first_line = file.readline() 442 return first_line.startswith("#!") and first_line.strip().endswith( 443 f"{CFG.binary_name} run" 444 ) 445 446 def _construct_job_name(self) -> str: 447 """ 448 Construct the job name for submission. 449 450 Returns: 451 str: The constructed job name. 452 """ 453 # for standard jobs, use script name 454 if not self._loop_info: 455 return self._script_name 456 457 # for loop jobs, use script_name with cycle number 458 return construct_loop_job_name(self._script_name, self._loop_info.current)
Class to submit jobs to a batch system.
Responsibilities:
- Validate that the script exists and has a proper shebang.
- Guard against multiple submissions from the same directory.
- Set environment variables required for
qq run.- Create a qq info file for tracking job state and metadata.
Note that Submitter ignores qq directives in the submitted script. To handle them, you have to build a Submitter using the SubmitterFactory.
53 def __init__( 54 self, 55 batch_system: AnyBatchClass, 56 queue: str, 57 account: str | None, 58 script: Path, 59 job_type: JobType, 60 resources: Resources, 61 loop_info: LoopInfo | None = None, 62 exclude: list[str] | None = None, 63 include: list[str] | None = None, 64 ignore: list[str] | None = None, 65 depend: list[Depend] | None = None, 66 transfer_mode: list[TransferMode] | None = None, 67 server: str | None = None, 68 interpreter: Interpreter | None = None, 69 resubmit_from: list[ResubmitHost] | None = None, 70 ): 71 """ 72 Initialize a Submitter instance. 73 74 Args: 75 batch_system (AnyBatchClass): The batch system class implementing 76 the BatchInterface used for job submission. 77 queue (str): The name of the batch system queue to which the job will be submitted. 78 account (str | None): The name of the account to use for the job. 79 script (Path): Path to the job script to submit. 80 job_type (JobType): Type of the job to submit (e.g. standard, loop). 81 resources (Resources): Job resource requirements (e.g., CPUs, memory, walltime). 82 loop_info (LoopInfo | None): Optional information for loop jobs. Pass None if not applicable. 83 exclude (list[str] | None): Optional list of files or glob patterns which should not be copied to the working directory. 84 Paths are provided relative to the input directory or absolute. 85 include (list[str] | None): Optional list of files or glob patterns which should be copied to the working directory 86 even though they are not part of the job's input directory. 87 Paths are provided either absolute or relative to the input directory. 88 ignore (list[str] | None): Optional list of files or glob patterns which should be ignored completely. 89 These files will not be copied to the working directory and if they are created in the working directory, they are also not copied back. 90 Paths are provided either absolute or relative to the input directory. 91 depend (list[Depend] | None): Optional list of job dependencies. 92 transfer_mode (list[TransferMode] | None): Mode specifying when files whould be transferred from the 93 working directory to the input directory. Defaults to [`Success()`]. 94 server (str | None): Optional name of the server to which the job should be submitted. 95 If `None`, the default batch server, as configured by the batch system is used. 96 intepreter (Interpreter | None): Optional interpreter specification to use to execute the script. 97 If not specified, the config default is used. 98 resubmit_from (list[ResubmitHost] | None): List of hosts from which a loop/continuous job should be resubmitted. 99 Must only be specified for loop/continuous jobs! 100 101 Raises: 102 QQError: If the script does not exist or has an invalid shebang line. 103 """ 104 105 self._batch_system = batch_system 106 self._job_type = job_type 107 self._queue = queue 108 self._server = server 109 self._account = account 110 self._loop_info = loop_info 111 self._script = script 112 self._input_dir = logical_resolve(script).parent 113 self._script_name = script.name 114 self._job_name = self._construct_job_name() 115 self._info_file = construct_info_file_path(self._input_dir, self._job_name) 116 self._resources = resources 117 self._exclude = expand_paths(exclude or [], self._input_dir) 118 self._include = expand_paths(include or [], self._input_dir) 119 self._ignore = expand_paths(ignore or [], self._input_dir) 120 self._depend = depend or [] 121 self._transfer_mode = transfer_mode or TransferMode.multi_from_str( 122 CFG.transfer_files_options.default_transfer_mode 123 ) 124 self._interpreter = interpreter 125 self._resubmit_from = resubmit_from or [] 126 127 # script must exist 128 if not self._script.is_file(): 129 raise QQError(f"Script '{script}' does not exist or is not a file") 130 131 # script must have a valid qq shebang 132 if not self._has_valid_shebang(self._script): 133 raise QQError( 134 f"Script '{self._script}' has an invalid shebang. The first line of the script should be '#!/usr/bin/env -S {CFG.binary_name} run'" 135 )
Initialize a Submitter instance.
Arguments:
- batch_system (AnyBatchClass): The batch system class implementing the BatchInterface used for job submission.
- queue (str): The name of the batch system queue to which the job will be submitted.
- account (str | None): The name of the account to use for the job.
- script (Path): Path to the job script to submit.
- job_type (JobType): Type of the job to submit (e.g. standard, loop).
- resources (Resources): Job resource requirements (e.g., CPUs, memory, walltime).
- loop_info (LoopInfo | None): Optional information for loop jobs. Pass None if not applicable.
- exclude (list[str] | None): Optional list of files or glob patterns which should not be copied to the working directory. Paths are provided relative to the input directory or absolute.
- include (list[str] | None): Optional list of files or glob patterns which should be copied to the working directory even though they are not part of the job's input directory. Paths are provided either absolute or relative to the input directory.
- ignore (list[str] | None): Optional list of files or glob patterns which should be ignored completely. These files will not be copied to the working directory and if they are created in the working directory, they are also not copied back. Paths are provided either absolute or relative to the input directory.
- depend (list[Depend] | None): Optional list of job dependencies.
- transfer_mode (list[TransferMode] | None): Mode specifying when files whould be transferred from the
working directory to the input directory. Defaults to [
Success()]. - server (str | None): Optional name of the server to which the job should be submitted.
If
None, the default batch server, as configured by the batch system is used. - intepreter (Interpreter | None): Optional interpreter specification to use to execute the script. If not specified, the config default is used.
- resubmit_from (list[ResubmitHost] | None): List of hosts from which a loop/continuous job should be resubmitted. Must only be specified for loop/continuous jobs!
Raises:
- QQError: If the script does not exist or has an invalid shebang line.
137 def submit(self, remote: str | None = None) -> str: 138 """ 139 Submit the script to the batch system. 140 141 Sets required environment variables, calls the batch system's 142 job submission mechanism, and creates an info file with job metadata. 143 144 This method is thread-safe, if the submission is done from the current machine. 145 146 Args: 147 remote (str | None): Name of the machine from which the job should be submitted. 148 If `None`, the current machine is used. 149 150 Returns: 151 str: The job ID of the submitted job. 152 153 Raises: 154 QQError: If job submission fails. 155 """ 156 job_id = self._batch_system.job_submit( 157 self._resources, 158 self._queue, 159 self._script, 160 self._job_name, 161 self._depend, 162 self._create_env_vars_dict(), 163 self._account, 164 self._server, 165 remote_host=remote, 166 ) 167 168 # create job qq info file 169 # we create the info file from the current machine no matter 170 # whether we are submiting from the current machine or from the remote machine 171 # the input directory should be available on both concerned machines, 172 # so this should be okay 173 Info( 174 batch_system=self._batch_system, 175 qq_version=qq_lib.__version__, 176 username=getpass.getuser(), 177 job_id=job_id, 178 job_name=self._job_name, 179 script_name=self._script_name, 180 queue=self._queue, 181 job_type=self._job_type, 182 input_machine=socket.getfqdn(remote or ""), 183 input_dir=self._input_dir, 184 job_state=NaiveState.QUEUED, 185 submission_time=datetime.now(), 186 stdout_file=str(Path(self._job_name).with_suffix(CFG.suffixes.stdout)), 187 stderr_file=str(Path(self._job_name).with_suffix(CFG.suffixes.stderr)), 188 resources=self._resources, 189 loop_info=self._loop_info, 190 excluded_files=self._exclude, 191 included_files=self._include, 192 ignored_files=self._ignore, 193 depend=self._depend, 194 account=self._account, 195 transfer_mode=self._transfer_mode, 196 server=self._server, 197 interpreter=self._interpreter, 198 resubmit_from=self._resubmit_from, 199 ).to_file(self._info_file) 200 201 return job_id
Submit the script to the batch system.
Sets required environment variables, calls the batch system's job submission mechanism, and creates an info file with job metadata.
This method is thread-safe, if the submission is done from the current machine.
Arguments:
- remote (str | None): Name of the machine from which the job should be submitted.
If
None, the current machine is used.
Returns:
str: The job ID of the submitted job.
Raises:
- QQError: If job submission fails.
203 def continues_loop(self) -> bool: 204 """ 205 Determine whether the submitted job is a continuation of a loop/continuous job. 206 207 Returns: 208 bool: True if the job is a valid continuation of a previous loop/continuous job, 209 False otherwise. 210 """ 211 try: 212 # there should only be one info file for both loop jobs (runtime files are archived) 213 # and continuous jobs (runtime files overwrite each other) 214 info_file = get_info_file(self._input_dir) 215 informer = Informer.from_file(info_file) 216 217 if self._loop_job_continues_loop( 218 informer 219 ) or self._continuous_job_continues_loop(informer): 220 logger.debug("Valid loop job with a correct cycle or a continuous job.") 221 return True 222 logger.debug( 223 "Detected info file does not correspond to a resubmittable job." 224 ) 225 return False 226 except QQError as e: 227 logger.debug(f"Could not read an info file: {e}") 228 return False
Determine whether the submitted job is a continuation of a loop/continuous job.
Returns:
bool: True if the job is a valid continuation of a previous loop/continuous job, False otherwise.
268 def get_input_dir(self) -> Path: 269 """ 270 Get path to the job's input directory. 271 272 Returns: 273 Path: Path to the job's input directory. 274 """ 275 return self._input_dir
Get path to the job's input directory.
Returns:
Path: Path to the job's input directory.
277 def get_batch_system(self) -> AnyBatchClass: 278 """Get the batch system used for submiting.""" 279 return self._batch_system
Get the batch system used for submiting.
293 def get_script(self) -> Path: 294 """Get absolute (logical) path to the submitted script.""" 295 return self._script
Get absolute (logical) path to the submitted script.
301 def get_resources(self) -> Resources: 302 """Get resources requested for the job.""" 303 return self._resources
Get resources requested for the job.
305 def get_loop_info(self) -> LoopInfo | None: 306 """Get loop job information.""" 307 return self._loop_info
Get loop job information.
309 def get_exclude(self) -> list[Path]: 310 """Get a list of excluded files.""" 311 return self._exclude
Get a list of excluded files.
313 def get_include(self) -> list[Path]: 314 """Get a list of included files.""" 315 return self._include
Get a list of included files.
317 def get_ignore(self) -> list[Path]: 318 """Get a list of ignored files.""" 319 return self._ignore
Get a list of ignored files.
321 def get_depend(self) -> list[Depend]: 322 """Get the list of dependencies.""" 323 return self._depend
Get the list of dependencies.
325 def get_transfer_mode(self) -> list[TransferMode]: 326 """Get the list of transfer modes.""" 327 return self._transfer_mode
Get the list of transfer modes.
329 def get_server(self) -> str | None: 330 """Get the submission server.""" 331 return self._server
Get the submission server.
333 def get_interpreter(self) -> Interpreter | None: 334 """Get the interpreter to use for running the script.""" 335 return self._interpreter
Get the interpreter to use for running the script.