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"]
class SubmitterFactory:
 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.

SubmitterFactory(script: pathlib._local.Path, **kwargs)
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.
def make_submitter(self) -> Submitter:
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.
class Parser:
 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.

Parser(script: pathlib._local.Path, params: list[click.core.Parameter])
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 GroupedOption names are considered.
def parse(self) -> None:
 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.
def get_batch_system(self) -> type[qq_lib.batch.interface.BatchInterface] | None:
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.

def get_queue(self) -> str | 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.

def get_job_type(self) -> qq_lib.properties.job_type.JobType | None:
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.

def get_resources(self) -> qq_lib.properties.resources.Resources:
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.

def get_exclude(self) -> list[str]:
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.

def get_include(self) -> list[str]:
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.

def get_ignore(self) -> list[str]:
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.

def get_loop_start(self) -> int | None:
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.

def get_loop_end(self) -> int | None:
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.

def get_archive(self) -> pathlib._local.Path | None:
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.

def get_archive_format(self) -> str | None:
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.

def get_archive_mode(self) -> list[qq_lib.properties.transfer_mode.TransferMode]:
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.

def get_depend(self) -> list[qq_lib.properties.depend.Depend]:
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.

def get_account(self) -> str | None:
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.

def get_transfer_mode(self) -> list[qq_lib.properties.transfer_mode.TransferMode]:
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.

def get_server(self) -> str | None:
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 None if not specified.

def get_interpreter(self) -> qq_lib.properties.interpreter.Interpreter | None:
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 None if not specified.

def get_resubmit_from(self) -> list[qq_lib.properties.resubmit_host.ResubmitHost]:
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.

class Submitter:
 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.

Submitter( batch_system: AnyBatchClass, queue: str, account: str | None, script: pathlib._local.Path, job_type: qq_lib.properties.job_type.JobType, resources: qq_lib.properties.resources.Resources, loop_info: qq_lib.properties.loop.LoopInfo | None = None, exclude: list[str] | None = None, include: list[str] | None = None, ignore: list[str] | None = None, depend: list[qq_lib.properties.depend.Depend] | None = None, transfer_mode: list[qq_lib.properties.transfer_mode.TransferMode] | None = None, server: str | None = None, interpreter: qq_lib.properties.interpreter.Interpreter | None = None, resubmit_from: list[qq_lib.properties.resubmit_host.ResubmitHost] | None = None)
 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.
def submit(self, remote: str | None = None) -> str:
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.
def continues_loop(self) -> bool:
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.

def get_input_dir(self) -> pathlib._local.Path:
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.

def get_batch_system(self) -> AnyBatchClass:
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.

def get_job_name(self) -> str:
281    def get_job_name(self) -> str:
282        """Get the name of the job."""
283        return self._job_name

Get the name of the job.

def get_queue(self) -> str:
285    def get_queue(self) -> str:
286        """Get the submission queue."""
287        return self._queue

Get the submission queue.

def get_account(self) -> str | None:
289    def get_account(self) -> str | None:
290        """Get the user's account."""
291        return self._account

Get the user's account.

def get_script(self) -> pathlib._local.Path:
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.

def get_job_type(self) -> qq_lib.properties.job_type.JobType:
297    def get_job_type(self) -> JobType:
298        """Get type of the job."""
299        return self._job_type

Get type of the job.

def get_resources(self) -> qq_lib.properties.resources.Resources:
301    def get_resources(self) -> Resources:
302        """Get resources requested for the job."""
303        return self._resources

Get resources requested for the job.

def get_loop_info(self) -> qq_lib.properties.loop.LoopInfo | None:
305    def get_loop_info(self) -> LoopInfo | None:
306        """Get loop job information."""
307        return self._loop_info

Get loop job information.

def get_exclude(self) -> list[pathlib._local.Path]:
309    def get_exclude(self) -> list[Path]:
310        """Get a list of excluded files."""
311        return self._exclude

Get a list of excluded files.

def get_include(self) -> list[pathlib._local.Path]:
313    def get_include(self) -> list[Path]:
314        """Get a list of included files."""
315        return self._include

Get a list of included files.

def get_ignore(self) -> list[pathlib._local.Path]:
317    def get_ignore(self) -> list[Path]:
318        """Get a list of ignored files."""
319        return self._ignore

Get a list of ignored files.

def get_depend(self) -> list[qq_lib.properties.depend.Depend]:
321    def get_depend(self) -> list[Depend]:
322        """Get the list of dependencies."""
323        return self._depend

Get the list of dependencies.

def get_transfer_mode(self) -> list[qq_lib.properties.transfer_mode.TransferMode]:
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.

def get_server(self) -> str | None:
329    def get_server(self) -> str | None:
330        """Get the submission server."""
331        return self._server

Get the submission server.

def get_interpreter(self) -> qq_lib.properties.interpreter.Interpreter | None:
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.

def get_resubmit_from(self) -> list[qq_lib.properties.resubmit_host.ResubmitHost] | None:
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

Get the list of hosts to resubmit the job from.