qq_lib.archive

Utilities for archiving and retrieving job-related files.

This module provides the Archiver class, which coordinates the movement of files between working directory and the job archive.

 1# Released under MIT License.
 2# Copyright (c) 2025-2026 Ladislav Bartos and Robert Vacha Lab
 3
 4"""
 5Utilities for archiving and retrieving job-related files.
 6
 7This module provides the `Archiver` class, which coordinates the movement
 8of files between working directory and the job archive.
 9"""
10
11from .archiver import Archiver
12
13__all__ = [
14    "Archiver",
15]
class Archiver:
 21class Archiver:
 22    """
 23    Manages archiving and retrieval of job-related files.
 24    """
 25
 26    def __init__(
 27        self,
 28        archive: Path,
 29        archive_format: str,
 30        input_machine: str,
 31        input_dir: Path,
 32        batch_system: AnyBatchClass,
 33        included_files: list[Path],
 34        excluded_files: list[Path],
 35        ignored_files: list[Path],
 36    ):
 37        """
 38        Initialize the Archiver.
 39
 40        Args:
 41            archive (Path): Absolute path to the job's archive directory.
 42            archive_format (str): Printf-style or regex pattern describing archived filenames.
 43            input_machine (str): The hostname from which the job was submitted.
 44            input_dir (Path): The directory from which the job was submitted.
 45            batch_system (AnyBatchClass): The batch system which manages the job.
 46            included_files (list[Path]): List that were explicitly included
 47                in the working directory and should not be archived.
 48            excluded_files (list[Path]): List of files that were explicitly excluded
 49                from the working directory and should not be fetched from archive.
 50            ignored_files (list[Path]): List of files that are ignored and should be neither
 51                archived nor fetched from archive.
 52        """
 53        self._batch_system = batch_system
 54        self._archive = archive
 55        self._archive_format = archive_format
 56        self._input_machine = input_machine
 57        self._input_dir = input_dir
 58        self._included_files = included_files
 59        self._excluded_files = excluded_files
 60        self._ignored_files = ignored_files
 61
 62    @property
 63    def archive(self) -> Path:
 64        """
 65        Returns the absolute path to the archive.
 66        """
 67        return self._archive
 68
 69    def make_archive_dir(self) -> None:
 70        """
 71        Create the archive directory in the job's input directory if it does not already exist.
 72        """
 73        logger.debug(
 74            f"Attempting to create an archive '{self._archive}' on '{self._input_machine}'."
 75        )
 76        self._batch_system.make_remote_dir(self._input_machine, self._archive)
 77
 78    def from_archive(self, dir: Path, cycle: int | None = None) -> None:
 79        """
 80        Fetch files from the archive to job's working directory.
 81
 82        This method retrieves files from the archive that match the
 83        configured archive pattern. If a cycle number is provided, only
 84        files corresponding to that cycle (for printf-style patterns) are
 85        fetched. If no cycle is provided, all files matching the pattern
 86        in the archive are fetched.
 87
 88        Files that were explicitly excluded via the `exclude` submission option are not fetched.
 89
 90        Args:
 91            dir (Path): The directory where files will be copied to.
 92            cycle (int | None): The cycle number to filter files for.
 93                Only relevant for printf-style patterns. If `None`, all
 94                matching files are fetched. Defaults to `None`.
 95
 96        Raises:
 97            QQError: If file transfer fails.
 98        """
 99        if not (
100            files := self.get_files_matching_pattern(
101                self._archive, self._input_machine, self._archive_format, cycle, False
102            )
103        ):
104            logger.debug("Nothing to fetch from archive.")
105            return
106
107        # files that were explicitly excluded via the `exclude` submission option are not fetched
108        # as are not fetched files that are explicitly ignored via the `ignore` submission option
109        exclude = self._get_excluded_from_copying_from_archive()
110        logger.debug(
111            f"Files that are excluded or ignored from being copied from the archive: {exclude}."
112        )
113
114        files = [file for file in files if file not in exclude]
115        logger.debug(f"Files to fetch from archive: {files}.")
116
117        Retryer(
118            self._batch_system.sync_selected,
119            self._archive,
120            dir,
121            self._input_machine,
122            socket.getfqdn(),
123            files,
124            max_tries=CFG.archiver.retry_tries,
125            wait_seconds=CFG.archiver.retry_wait,
126        ).run()
127
128    def to_archive(self, dir: Path) -> None:
129        """
130        Archive all files matching the archive format in the specified directory.
131
132        Copies all files matching the archive pattern from directory
133        `dir` to the archive directory. After successfully transferring
134        the files, they are removed from the working directory.
135
136        Files that were explicitly included via the `include` submission option are not archived.
137
138        Args:
139            work_dir (Path): The directory containing files to archive.
140
141        Raises:
142            QQError: If file transfer or removal fails.
143        """
144        if not (
145            files := self.get_files_matching_pattern(
146                dir, None, self._archive_format, None, False
147            )
148        ):
149            logger.debug("Nothing to archive.")
150            return
151
152        # files that were explicitly included via the `include` submission option are not archived
153        # as well as files that are explicitly ignored via the `ignore` submission option
154        exclude = self._get_excluded_from_copying_to_archive(dir)
155        logger.debug(
156            f"Files that are excluded or ignored from being copied to the archive: {exclude}."
157        )
158        files = [file for file in files if file not in exclude]
159
160        logger.debug(f"Files to archive: {files}.")
161
162        Retryer(
163            self._batch_system.sync_selected,
164            dir,
165            self._archive,
166            socket.getfqdn(),
167            self._input_machine,
168            files,
169            max_tries=CFG.archiver.retry_tries,
170            wait_seconds=CFG.archiver.retry_wait,
171        ).run()
172
173        # remove the archived files
174        Retryer(
175            self._remove_files,
176            files,
177            max_tries=CFG.archiver.retry_tries,
178            wait_seconds=CFG.archiver.retry_wait,
179        ).run()
180
181    def archive_runtime_files(self, job_name: str, cycle: int) -> None:
182        """
183        Archive qq runtime files from a specific job located in the input directory.
184
185        The archived files are moved from the input directory to the archive directory.
186
187        Ensure that `job_name` does not contain special regex characters, or that any such
188        characters are properly escaped.
189
190        This function will archive all files whose names match `job_name`, regardless
191        of whether they have any qq-specific suffixes.
192
193        Args:
194            job_name (str): The name of the job.
195            cycle (int): Cycle number for which the files should be archived.
196
197        Raises:
198            QQError: If moving the runtime files fails.
199        """
200        if not (
201            files := self.get_files_matching_pattern(
202                self._input_dir,
203                self._input_machine,
204                # only use the stem of the job name, the extension will not be matched
205                job_name.split(".", maxsplit=1)[0],
206                # we do not need to use the cycle number here since the job_name should already be expanded
207                cycle=None,
208                include_qq_files=True,
209            )
210        ):
211            logger.debug("No qq runtime files to archive.")
212            return
213
214        # the files are renamed to conform the the archive format
215        moved_files = [
216            self._archive / f"{self._archive_format % cycle}{f.suffix}" for f in files
217        ]
218
219        logger.debug(f"qq runtime files to archive: {files}.")
220        logger.debug(f"qq runtime files after moving: {moved_files}.")
221
222        Retryer(
223            self._batch_system.move_remote_files,
224            self._input_machine,
225            files,
226            moved_files,
227            max_tries=CFG.archiver.retry_tries,
228            wait_seconds=CFG.archiver.retry_wait,
229        ).run()
230
231    def get_files_matching_pattern(
232        self,
233        directory: Path,
234        host: str | None,
235        pattern: str,
236        cycle: int | None = None,
237        include_qq_files: bool = False,
238    ) -> list[Path]:
239        """
240        Determine which files in a directory match a given pattern.
241
242        Args:
243            directory (Path): Directory to search for files.
244            host (str | None): Hostname if the directory is remote,
245                or None if it is available from the current machine.
246            pattern (str): A printf-style or regex pattern to match file stems.
247            cycle (int | None): Optional cycle number for printf-style patterns.
248                If provided, only files corresponding to that loop are returned.
249                If `None`, all matching files are returned. Defaults to `None`.
250            include_qq_files (bool): Whether to include qq runtime files. Defaults to False.
251
252        Returns:
253            list[Path]: A list of absolute (logical) paths to matching files.
254        """
255        if cycle and is_printf_pattern(pattern):
256            try:
257                # try inserting the loop number into the printf pattern
258                regex = re.compile(f"{pattern % cycle}")
259            except Exception:
260                logger.debug(
261                    f"Ignoring loop number since the provided pattern ('{pattern}') does not support it."
262                )
263                regex = Archiver._prepare_regex_pattern(pattern)
264        else:
265            logger.debug(
266                f"Loop number not specified or the provided pattern ('{pattern}') does not support it."
267            )
268            regex = Archiver._prepare_regex_pattern(pattern)
269
270        logger.debug(f"Regex for matching: {regex}.")
271
272        # the directory must exist
273        if host and host != socket.getfqdn():
274            # remote directory
275            available_files: list[Path] = Retryer(
276                self._batch_system.list_remote_dir,
277                host,
278                directory,
279                max_tries=CFG.archiver.retry_tries,
280                wait_seconds=CFG.archiver.retry_wait,
281            ).run()
282        else:
283            # local directory
284            available_files = list(directory.iterdir())
285
286        logger.debug(f"All available files: {available_files}.")
287        if include_qq_files:
288            # the stem of the file must contain the regex pattern
289            return [logical_resolve(f) for f in available_files if regex.search(f.stem)]
290        return [
291            logical_resolve(f)
292            for f in available_files
293            if regex.search(f.stem) and f.suffix not in CFG.suffixes.all_suffixes
294        ]
295
296    def _get_excluded_from_copying_to_archive(self, dir: Path) -> list[Path]:
297        """
298        Return paths that must not be copied to the archive.
299
300        Collects the files ignored and explicitly included by the user,
301        and the archive directory itself. Duplicates are removed,
302        preserving the order of first occurrence.
303
304        Args:
305            dir (Path): The directory from which we are archiving the files.
306
307        Returns:
308            list[Path]: Paths that should not be copied to the archive.
309        """
310
311        return list(
312            dict.fromkeys(
313                relocate_by_name(
314                    [
315                        *self._ignored_files,
316                        *self._included_files,
317                        self._archive,
318                    ],
319                    dir,
320                )
321            )
322        )
323
324    def _get_excluded_from_copying_from_archive(self) -> list[Path]:
325        """
326        Return paths that must not be copied from the archive.
327
328        Collects files ignored and explicitly excluded by the user.
329        Duplicates are removed, preserving the order of first occurrence.
330
331        Returns:
332            list[Path]: Paths that should not be copied from the archive.
333        """
334
335        return list(
336            dict.fromkeys(
337                [
338                    *self._ignored_files,
339                    *self._excluded_files,
340                ]
341            )
342        )
343
344    def create_init_file(self, cycle: int) -> None:
345        """
346        Create an empty init file for the given cycle.
347        Used as a fallback when no valid archive file is produced, ensuring
348        the next iteration of the loop job can proceed normally.
349        Args:
350            cycle (int): The index of the next cycle of the loop job.
351        """
352        Path(f"{self._archive_format % cycle}.init").touch()
353
354    @staticmethod
355    def _prepare_regex_pattern(pattern: str) -> re.Pattern[str]:
356        """
357        Convert a printf-style pattern or regex string into a compiled regex.
358
359        Args:
360            pattern (str): The pattern to convert.
361
362        Returns:
363            re.Pattern[str]: Compiled regex pattern that can be used for matching.
364        """
365        if is_printf_pattern(pattern):
366            pattern = printf_to_regex(pattern)
367
368        return re.compile(pattern)
369
370    @staticmethod
371    def _remove_files(files: Iterable[Path]) -> None:
372        """
373        Remove a list of files or directories from the filesystem.
374
375        Args:
376            files (Iterable[Path]): Files or directories to delete.
377
378        Raises:
379            OSError: If file or directory removal fails for any file.
380        """
381        for file in files:
382            if file.is_symlink() or not file.is_dir():
383                file.unlink()
384            else:
385                shutil.rmtree(file)

Manages archiving and retrieval of job-related files.

Archiver( archive: pathlib._local.Path, archive_format: str, input_machine: str, input_dir: pathlib._local.Path, batch_system: AnyBatchClass, included_files: list[pathlib._local.Path], excluded_files: list[pathlib._local.Path], ignored_files: list[pathlib._local.Path])
26    def __init__(
27        self,
28        archive: Path,
29        archive_format: str,
30        input_machine: str,
31        input_dir: Path,
32        batch_system: AnyBatchClass,
33        included_files: list[Path],
34        excluded_files: list[Path],
35        ignored_files: list[Path],
36    ):
37        """
38        Initialize the Archiver.
39
40        Args:
41            archive (Path): Absolute path to the job's archive directory.
42            archive_format (str): Printf-style or regex pattern describing archived filenames.
43            input_machine (str): The hostname from which the job was submitted.
44            input_dir (Path): The directory from which the job was submitted.
45            batch_system (AnyBatchClass): The batch system which manages the job.
46            included_files (list[Path]): List that were explicitly included
47                in the working directory and should not be archived.
48            excluded_files (list[Path]): List of files that were explicitly excluded
49                from the working directory and should not be fetched from archive.
50            ignored_files (list[Path]): List of files that are ignored and should be neither
51                archived nor fetched from archive.
52        """
53        self._batch_system = batch_system
54        self._archive = archive
55        self._archive_format = archive_format
56        self._input_machine = input_machine
57        self._input_dir = input_dir
58        self._included_files = included_files
59        self._excluded_files = excluded_files
60        self._ignored_files = ignored_files

Initialize the Archiver.

Arguments:
  • archive (Path): Absolute path to the job's archive directory.
  • archive_format (str): Printf-style or regex pattern describing archived filenames.
  • input_machine (str): The hostname from which the job was submitted.
  • input_dir (Path): The directory from which the job was submitted.
  • batch_system (AnyBatchClass): The batch system which manages the job.
  • included_files (list[Path]): List that were explicitly included in the working directory and should not be archived.
  • excluded_files (list[Path]): List of files that were explicitly excluded from the working directory and should not be fetched from archive.
  • ignored_files (list[Path]): List of files that are ignored and should be neither archived nor fetched from archive.
archive: pathlib._local.Path
62    @property
63    def archive(self) -> Path:
64        """
65        Returns the absolute path to the archive.
66        """
67        return self._archive

Returns the absolute path to the archive.

def make_archive_dir(self) -> None:
69    def make_archive_dir(self) -> None:
70        """
71        Create the archive directory in the job's input directory if it does not already exist.
72        """
73        logger.debug(
74            f"Attempting to create an archive '{self._archive}' on '{self._input_machine}'."
75        )
76        self._batch_system.make_remote_dir(self._input_machine, self._archive)

Create the archive directory in the job's input directory if it does not already exist.

def from_archive(self, dir: pathlib._local.Path, cycle: int | None = None) -> None:
 78    def from_archive(self, dir: Path, cycle: int | None = None) -> None:
 79        """
 80        Fetch files from the archive to job's working directory.
 81
 82        This method retrieves files from the archive that match the
 83        configured archive pattern. If a cycle number is provided, only
 84        files corresponding to that cycle (for printf-style patterns) are
 85        fetched. If no cycle is provided, all files matching the pattern
 86        in the archive are fetched.
 87
 88        Files that were explicitly excluded via the `exclude` submission option are not fetched.
 89
 90        Args:
 91            dir (Path): The directory where files will be copied to.
 92            cycle (int | None): The cycle number to filter files for.
 93                Only relevant for printf-style patterns. If `None`, all
 94                matching files are fetched. Defaults to `None`.
 95
 96        Raises:
 97            QQError: If file transfer fails.
 98        """
 99        if not (
100            files := self.get_files_matching_pattern(
101                self._archive, self._input_machine, self._archive_format, cycle, False
102            )
103        ):
104            logger.debug("Nothing to fetch from archive.")
105            return
106
107        # files that were explicitly excluded via the `exclude` submission option are not fetched
108        # as are not fetched files that are explicitly ignored via the `ignore` submission option
109        exclude = self._get_excluded_from_copying_from_archive()
110        logger.debug(
111            f"Files that are excluded or ignored from being copied from the archive: {exclude}."
112        )
113
114        files = [file for file in files if file not in exclude]
115        logger.debug(f"Files to fetch from archive: {files}.")
116
117        Retryer(
118            self._batch_system.sync_selected,
119            self._archive,
120            dir,
121            self._input_machine,
122            socket.getfqdn(),
123            files,
124            max_tries=CFG.archiver.retry_tries,
125            wait_seconds=CFG.archiver.retry_wait,
126        ).run()

Fetch files from the archive to job's working directory.

This method retrieves files from the archive that match the configured archive pattern. If a cycle number is provided, only files corresponding to that cycle (for printf-style patterns) are fetched. If no cycle is provided, all files matching the pattern in the archive are fetched.

Files that were explicitly excluded via the exclude submission option are not fetched.

Arguments:
  • dir (Path): The directory where files will be copied to.
  • cycle (int | None): The cycle number to filter files for. Only relevant for printf-style patterns. If None, all matching files are fetched. Defaults to None.
Raises:
  • QQError: If file transfer fails.
def to_archive(self, dir: pathlib._local.Path) -> None:
128    def to_archive(self, dir: Path) -> None:
129        """
130        Archive all files matching the archive format in the specified directory.
131
132        Copies all files matching the archive pattern from directory
133        `dir` to the archive directory. After successfully transferring
134        the files, they are removed from the working directory.
135
136        Files that were explicitly included via the `include` submission option are not archived.
137
138        Args:
139            work_dir (Path): The directory containing files to archive.
140
141        Raises:
142            QQError: If file transfer or removal fails.
143        """
144        if not (
145            files := self.get_files_matching_pattern(
146                dir, None, self._archive_format, None, False
147            )
148        ):
149            logger.debug("Nothing to archive.")
150            return
151
152        # files that were explicitly included via the `include` submission option are not archived
153        # as well as files that are explicitly ignored via the `ignore` submission option
154        exclude = self._get_excluded_from_copying_to_archive(dir)
155        logger.debug(
156            f"Files that are excluded or ignored from being copied to the archive: {exclude}."
157        )
158        files = [file for file in files if file not in exclude]
159
160        logger.debug(f"Files to archive: {files}.")
161
162        Retryer(
163            self._batch_system.sync_selected,
164            dir,
165            self._archive,
166            socket.getfqdn(),
167            self._input_machine,
168            files,
169            max_tries=CFG.archiver.retry_tries,
170            wait_seconds=CFG.archiver.retry_wait,
171        ).run()
172
173        # remove the archived files
174        Retryer(
175            self._remove_files,
176            files,
177            max_tries=CFG.archiver.retry_tries,
178            wait_seconds=CFG.archiver.retry_wait,
179        ).run()

Archive all files matching the archive format in the specified directory.

Copies all files matching the archive pattern from directory dir to the archive directory. After successfully transferring the files, they are removed from the working directory.

Files that were explicitly included via the include submission option are not archived.

Arguments:
  • work_dir (Path): The directory containing files to archive.
Raises:
  • QQError: If file transfer or removal fails.
def archive_runtime_files(self, job_name: str, cycle: int) -> None:
181    def archive_runtime_files(self, job_name: str, cycle: int) -> None:
182        """
183        Archive qq runtime files from a specific job located in the input directory.
184
185        The archived files are moved from the input directory to the archive directory.
186
187        Ensure that `job_name` does not contain special regex characters, or that any such
188        characters are properly escaped.
189
190        This function will archive all files whose names match `job_name`, regardless
191        of whether they have any qq-specific suffixes.
192
193        Args:
194            job_name (str): The name of the job.
195            cycle (int): Cycle number for which the files should be archived.
196
197        Raises:
198            QQError: If moving the runtime files fails.
199        """
200        if not (
201            files := self.get_files_matching_pattern(
202                self._input_dir,
203                self._input_machine,
204                # only use the stem of the job name, the extension will not be matched
205                job_name.split(".", maxsplit=1)[0],
206                # we do not need to use the cycle number here since the job_name should already be expanded
207                cycle=None,
208                include_qq_files=True,
209            )
210        ):
211            logger.debug("No qq runtime files to archive.")
212            return
213
214        # the files are renamed to conform the the archive format
215        moved_files = [
216            self._archive / f"{self._archive_format % cycle}{f.suffix}" for f in files
217        ]
218
219        logger.debug(f"qq runtime files to archive: {files}.")
220        logger.debug(f"qq runtime files after moving: {moved_files}.")
221
222        Retryer(
223            self._batch_system.move_remote_files,
224            self._input_machine,
225            files,
226            moved_files,
227            max_tries=CFG.archiver.retry_tries,
228            wait_seconds=CFG.archiver.retry_wait,
229        ).run()

Archive qq runtime files from a specific job located in the input directory.

The archived files are moved from the input directory to the archive directory.

Ensure that job_name does not contain special regex characters, or that any such characters are properly escaped.

This function will archive all files whose names match job_name, regardless of whether they have any qq-specific suffixes.

Arguments:
  • job_name (str): The name of the job.
  • cycle (int): Cycle number for which the files should be archived.
Raises:
  • QQError: If moving the runtime files fails.
def get_files_matching_pattern( self, directory: pathlib._local.Path, host: str | None, pattern: str, cycle: int | None = None, include_qq_files: bool = False) -> list[pathlib._local.Path]:
231    def get_files_matching_pattern(
232        self,
233        directory: Path,
234        host: str | None,
235        pattern: str,
236        cycle: int | None = None,
237        include_qq_files: bool = False,
238    ) -> list[Path]:
239        """
240        Determine which files in a directory match a given pattern.
241
242        Args:
243            directory (Path): Directory to search for files.
244            host (str | None): Hostname if the directory is remote,
245                or None if it is available from the current machine.
246            pattern (str): A printf-style or regex pattern to match file stems.
247            cycle (int | None): Optional cycle number for printf-style patterns.
248                If provided, only files corresponding to that loop are returned.
249                If `None`, all matching files are returned. Defaults to `None`.
250            include_qq_files (bool): Whether to include qq runtime files. Defaults to False.
251
252        Returns:
253            list[Path]: A list of absolute (logical) paths to matching files.
254        """
255        if cycle and is_printf_pattern(pattern):
256            try:
257                # try inserting the loop number into the printf pattern
258                regex = re.compile(f"{pattern % cycle}")
259            except Exception:
260                logger.debug(
261                    f"Ignoring loop number since the provided pattern ('{pattern}') does not support it."
262                )
263                regex = Archiver._prepare_regex_pattern(pattern)
264        else:
265            logger.debug(
266                f"Loop number not specified or the provided pattern ('{pattern}') does not support it."
267            )
268            regex = Archiver._prepare_regex_pattern(pattern)
269
270        logger.debug(f"Regex for matching: {regex}.")
271
272        # the directory must exist
273        if host and host != socket.getfqdn():
274            # remote directory
275            available_files: list[Path] = Retryer(
276                self._batch_system.list_remote_dir,
277                host,
278                directory,
279                max_tries=CFG.archiver.retry_tries,
280                wait_seconds=CFG.archiver.retry_wait,
281            ).run()
282        else:
283            # local directory
284            available_files = list(directory.iterdir())
285
286        logger.debug(f"All available files: {available_files}.")
287        if include_qq_files:
288            # the stem of the file must contain the regex pattern
289            return [logical_resolve(f) for f in available_files if regex.search(f.stem)]
290        return [
291            logical_resolve(f)
292            for f in available_files
293            if regex.search(f.stem) and f.suffix not in CFG.suffixes.all_suffixes
294        ]

Determine which files in a directory match a given pattern.

Arguments:
  • directory (Path): Directory to search for files.
  • host (str | None): Hostname if the directory is remote, or None if it is available from the current machine.
  • pattern (str): A printf-style or regex pattern to match file stems.
  • cycle (int | None): Optional cycle number for printf-style patterns. If provided, only files corresponding to that loop are returned. If None, all matching files are returned. Defaults to None.
  • include_qq_files (bool): Whether to include qq runtime files. Defaults to False.
Returns:

list[Path]: A list of absolute (logical) paths to matching files.

def create_init_file(self, cycle: int) -> None:
344    def create_init_file(self, cycle: int) -> None:
345        """
346        Create an empty init file for the given cycle.
347        Used as a fallback when no valid archive file is produced, ensuring
348        the next iteration of the loop job can proceed normally.
349        Args:
350            cycle (int): The index of the next cycle of the loop job.
351        """
352        Path(f"{self._archive_format % cycle}.init").touch()

Create an empty init file for the given cycle. Used as a fallback when no valid archive file is produced, ensuring the next iteration of the loop job can proceed normally.

Arguments:
  • cycle (int): The index of the next cycle of the loop job.