Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/git_utils/repo_cloner.py: 53%
97 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 11:54 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 11:54 +0000
1import asyncio
2import os
3import shutil
4import uuid
5from functools import partial
6from pathlib import Path
7from typing import Generator, Optional
9from git import GitCommandError, GitError, Repo
10from opal_common.config import opal_common_config
11from opal_common.git_utils.env import provide_git_ssh_environment
12from opal_common.git_utils.exceptions import GitFailed
13from opal_common.http_utils import redact_url, redact_url_in_text
14from opal_common.logger import logger
15from opal_common.utils import get_filepaths_with_glob
16from tenacity import RetryError, retry, stop, wait
19class CloneResult:
20 """Wraps a git.Repo instance but knows if the repo was initialized with a
21 url and cloned from a remote repo, or was initialed from a local `.git`
22 repo."""
24 def __init__(self, repo: Repo):
25 self._repo = repo
27 @property
28 def repo(self) -> Repo:
29 """The wrapped repo instance."""
30 return self._repo
33class RepoClonePathFinder:
34 """
35 We are cloning the policy repo into a unique random subdirectory of a base path.
36 Args:
37 base_clone_path (str): parent directory for the repoistory clone
38 clone_subdirectory_prefix (str): the prefix for the randomized repository dir, or the dir name itself when `use_fixes_path=true`
39 use_fixed_path (bool): if set, random suffix won't be added to `clone_subdirectory_prefix` (if the path already exists, it would be reused)
41 This class knows how to such clones, so we can discard previous ones, but also so
42 that siblings workers (who are not the master who decided where to clone) can also
43 find the current clone by globing the base dir.
44 """
46 def __init__(
47 self, base_clone_path: str, clone_subdirectory_prefix: str, use_fixed_path: bool
48 ):
49 if not base_clone_path:
50 raise ValueError("base_clone_path cannot be empty!")
52 if not clone_subdirectory_prefix:
53 raise ValueError("clone_subdirectory_prefix cannot be empty!")
55 self._base_clone_path = os.path.expanduser(base_clone_path)
56 self._clone_subdirectory_prefix = clone_subdirectory_prefix
57 self._use_fixed_path = use_fixed_path
59 def _get_randomized_clone_subdirectories(self) -> Generator[str, None, None]:
60 """A generator yielding all the randomized subdirectories of the base
61 clone path that are matching the clone pattern.
63 Yields:
64 the next subdirectory matching the pattern
65 """
66 folders_with_pattern = get_filepaths_with_glob(
67 self._base_clone_path, f"{self._clone_subdirectory_prefix}-*"
68 )
69 for folder in folders_with_pattern: 69 ↛ 70line 69 didn't jump to line 70 because the loop on line 69 never started
70 yield folder
72 def _get_single_existing_random_clone_path(self) -> Optional[str]:
73 """Searches for the single randomly-suffixed clone subdirectory in
74 existence.
76 If found no such subdirectory or if found more than one (multiple matching subdirectories) - will return None.
77 otherwise: will return the single and only clone.
78 """
79 subdirectories = list(self._get_randomized_clone_subdirectories())
80 if len(subdirectories) != 1: 80 ↛ 82line 80 didn't jump to line 82 because the condition on line 80 was always true
81 return None
82 return subdirectories[0]
84 def _generate_randomized_clone_path(self) -> str:
85 folder_name = f"{self._clone_subdirectory_prefix}-{uuid.uuid4().hex}"
86 full_local_repo_path = os.path.join(self._base_clone_path, folder_name)
87 return full_local_repo_path
89 def _get_fixed_clone_path(self) -> str:
90 return os.path.join(self._base_clone_path, self._clone_subdirectory_prefix)
92 def get_clone_path(self) -> Optional[str]:
93 """Get the clone path (fixed or randomized) if it exists."""
94 if self._use_fixed_path:
95 fixed_path = self._get_fixed_clone_path()
96 if os.path.exists(fixed_path):
97 return fixed_path
98 else:
99 return None
100 else:
101 return self._get_single_existing_random_clone_path()
103 def create_new_clone_path(self) -> str:
104 """
105 If using a fixed path - simply creates it.
106 If using a randomized suffix -
107 takes the base path from server config and create new folder with unique name for the local clone.
108 The folder name is looks like /<base-path>/<folder-prefix>-<uuid>
109 If such folders already exist they would be removed.
110 """
111 if self._use_fixed_path:
112 # When using fixed path - just use old path without cleanup
113 full_local_repo_path = self._get_fixed_clone_path()
114 else:
115 # Remove old randomized subdirectories
116 for folder in self._get_randomized_clone_subdirectories():
117 logger.warning(
118 "Found previous policy repo clone: {folder_name}, removing it to avoid conflicts.",
119 folder_name=folder,
120 )
121 shutil.rmtree(folder)
122 full_local_repo_path = self._generate_randomized_clone_path()
124 os.makedirs(full_local_repo_path, exist_ok=True)
125 return full_local_repo_path
128class RepoCloner:
129 """Simple wrapper for git.Repo() to simplify other classes that need to
130 deal with the case where a repo must be cloned from url *only if* the repo
131 does not already exists locally, and otherwise initialize the repo instance
132 from the repo already existing on the filesystem."""
134 # wait indefinitely until successful
135 DEFAULT_RETRY_CONFIG = {
136 "wait": wait.wait_random_exponential(multiplier=0.5, max=30),
137 }
139 def __init__(
140 self,
141 repo_url: str,
142 clone_path: str,
143 branch_name: str = "master",
144 retry_config=None,
145 ssh_key: Optional[str] = None,
146 ssh_key_file_path: Optional[str] = None,
147 clone_timeout: int = 0,
148 ):
149 """Inits the repo cloner.
151 Args:
152 repo_url (str): the url to the remote repo we want to clone
153 clone_path (str): the target local path in our file system we want the
154 repo to be cloned to
155 retry_config (dict): Tenacity.retry config (@see https://tenacity.readthedocs.io/en/latest/api.html#retry-main-api)
156 ssh_key (str, optional): private ssh key used to gain access to the cloned repo
157 ssh_key_file_path (str, optional): local path to save the private ssh key contents
158 """
159 if repo_url is None:
160 raise ValueError("must provide repo url!")
162 self.url = repo_url
163 self.path = os.path.expanduser(clone_path)
164 self.branch_name = branch_name
165 self._ssh_key = ssh_key
166 self._ssh_key_file_path = (
167 ssh_key_file_path or opal_common_config.GIT_SSH_KEY_FILE
168 )
169 self._retry_config = (
170 retry_config if retry_config is not None else self.DEFAULT_RETRY_CONFIG
171 )
172 if clone_timeout > 0:
173 self._retry_config.update({"stop": stop.stop_after_delay(clone_timeout)})
175 async def clone(self) -> CloneResult:
176 """Initializes a git.Repo and returns the clone result. it either:
178 - does not found a cloned repo locally and clones from remote url
179 - finds a cloned repo locally and does not clone from remote.
180 """
181 logger.info(
182 "Cloning repo from '{url}' to '{to_path}' (branch: '{branch}')",
183 url=redact_url(self.url),
184 to_path=self.path,
185 branch=self.branch_name,
186 )
187 loop = asyncio.get_running_loop()
188 return await loop.run_in_executor(None, self._attempt_clone_from_url)
190 def _attempt_clone_from_url(self) -> CloneResult:
191 """Clones the repo from url or throws GitFailed."""
192 env = provide_git_ssh_environment(self.url, self._ssh_key)
193 _clone_func = partial(self._clone, env=env)
194 _clone_with_retries = retry(**self._retry_config)(_clone_func)
195 try:
196 repo: Repo = _clone_with_retries()
197 except (GitError, GitCommandError) as e:
198 raise GitFailed(e)
199 except RetryError as e:
200 logger.exception(
201 "cannot clone policy repo: {error}",
202 error=redact_url_in_text(str(e), self.url),
203 )
204 raise GitFailed(e)
205 else:
206 logger.info("Clone succeeded", repo_path=self.path)
207 return CloneResult(repo)
209 def _clone(self, env) -> Repo:
210 try:
211 return Repo.clone_from(
212 url=self.url, to_path=self.path, branch=self.branch_name, env=env
213 )
214 except (GitError, GitCommandError) as e:
215 logger.error(
216 "cannot clone policy repo: {error}",
217 error=redact_url_in_text(str(e), self.url),
218 )
219 raise