Coverage for python/lsst/resources/gs.py: 15%

220 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-19 09:04 +0000

1# This file is part of lsst-resources. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

5# (https://www.lsst.org). 

6# See the COPYRIGHT file at the top-level directory of this distribution 

7# for details of code ownership. 

8# 

9# Use of this source code is governed by a 3-clause BSD-style 

10# license that can be found in the LICENSE file. 

11 

12"""Accessing Google Cloud Storage resources.""" 

13 

14from __future__ import annotations 

15 

16__all__ = ("GSResourcePath",) 

17 

18import contextlib 

19import datetime 

20import logging 

21import re 

22from collections.abc import Generator, Iterator 

23from typing import TYPE_CHECKING 

24 

25from ._resourceHandles._baseResourceHandle import ResourceHandleProtocol 

26 

27try: 

28 import google.api_core.retry as retry 

29 import google.cloud.storage as storage 

30 from google.cloud.exceptions import ( 

31 BadGateway, 

32 InternalServerError, 

33 NotFound, 

34 ServiceUnavailable, 

35 TooManyRequests, 

36 ) 

37except ImportError: 

38 # Hidden from type checkers so that the names above keep the types they 

39 # have when google-cloud-storage is installed. 

40 if not TYPE_CHECKING: 

41 storage = None 

42 retry = None 

43 

44 # Must also fake the exception classes. 

45 class ClientError(Exception): 

46 """Generic client error.""" 

47 

48 pass 

49 

50 class NotFound(ClientError): # noqa: N818 

51 """Resource not found error.""" 

52 

53 pass 

54 

55 class TooManyRequests(ClientError): # noqa: N818 

56 """Too many requests error.""" 

57 

58 pass 

59 

60 class InternalServerError(ClientError): 

61 """Internal server error.""" 

62 

63 pass 

64 

65 class BadGateway(ClientError): # noqa: N818 

66 """Bad gateway error.""" 

67 

68 pass 

69 

70 class ServiceUnavailable(ClientError): # noqa: N818 

71 """Service unavailable error.""" 

72 

73 pass 

74 

75 

76from lsst.utils.timer import time_this 

77 

78from ._resourcePath import ResourceInfo, ResourcePath 

79 

80if TYPE_CHECKING: 

81 from .utils import TransactionProtocol 

82 

83log = logging.getLogger(__name__) 

84 

85 

86_RETRIEVABLE_TYPES = ( 

87 TooManyRequests, # 429 

88 InternalServerError, # 500 

89 BadGateway, # 502 

90 ServiceUnavailable, # 503 

91) 

92 

93 

94def is_retryable(exc: Exception) -> bool: 

95 """Report if the given exception is a condition that can be retried. 

96 

97 Parameters 

98 ---------- 

99 exc : `Exception` 

100 Exception to check. 

101 

102 Returns 

103 ------- 

104 `bool` 

105 Returns `True` if the given exception is a condition that can be 

106 retried. 

107 """ 

108 return isinstance(exc, _RETRIEVABLE_TYPES) 

109 

110 

111_RETRY_POLICY = retry.Retry(predicate=is_retryable) if retry is not None else None 

112 

113 

114_client = None 

115"""Cached client connection.""" 

116 

117 

118def _coerce_gcs_datetime(value: datetime.datetime | str | None) -> datetime.datetime | None: 

119 """Convert GCS timestamp values to timezone-aware UTC datetimes. 

120 

121 Some emulators return RFC3339 timestamps with an explicit UTC offset 

122 instead of a trailing ``Z``, which the google-cloud-storage property 

123 accessors do not always accept. 

124 """ 

125 if value is None: 

126 return None 

127 if isinstance(value, datetime.datetime): 

128 if value.tzinfo is None: 

129 return value.replace(tzinfo=datetime.UTC) 

130 return value.astimezone(datetime.UTC) 

131 if value.endswith("Z"): 

132 value = value[:-1] + "+00:00" 

133 return datetime.datetime.fromisoformat(value).astimezone(datetime.UTC) 

134 

135 

136def _get_client() -> storage.Client: 

137 global _client 

138 if storage is None: 

139 raise ImportError("google-cloud-storage package not installed. Unable to communicate with GCS.") 

140 if _client is None: 

141 _client = storage.Client() 

142 return _client 

143 

144 

145class GSResourcePath(ResourcePath): 

146 """Access Google Cloud Storage resources.""" 

147 

148 _bucket: storage.Bucket | None = None 

149 _blob: storage.Blob | None = None 

150 _client: storage.Client | None = None 

151 

152 @property 

153 def client(self) -> storage.Client: 

154 return _get_client() 

155 

156 @property 

157 def bucket(self) -> storage.Bucket: 

158 if self._bucket is None: 

159 self._bucket = self.client.bucket(self.netloc) 

160 return self._bucket 

161 

162 @property 

163 def blob(self) -> storage.Blob: 

164 if self._blob is None: 

165 self._blob = self.bucket.blob(self.relativeToPathRoot) 

166 return self._blob 

167 

168 def exists(self) -> bool: 

169 if self.is_root: 

170 return self.bucket.exists(retry=_RETRY_POLICY) 

171 if self.dirLike: 

172 # GCS does not have concrete directory objects; treat any 

173 # directory-like path within an existing bucket as existing. 

174 return self.bucket.exists(retry=_RETRY_POLICY) 

175 return self.blob.exists(retry=_RETRY_POLICY) 

176 

177 def size(self) -> int: 

178 if self.dirLike: 

179 return 0 

180 # The first time this is called we need to sync from the remote. 

181 # Force the blob to be recalculated. 

182 try: 

183 self.blob.reload(retry=_RETRY_POLICY) 

184 except NotFound: 

185 raise FileNotFoundError(f"Resource {self} does not exist") from None 

186 size = self.blob.size 

187 if size is None: 

188 raise FileNotFoundError(f"Resource {self} does not exist") 

189 return size 

190 

191 def get_info(self) -> ResourceInfo: 

192 """Return lightweight metadata about this GCS resource.""" 

193 if self.is_root: 

194 if not self.bucket.exists(retry=_RETRY_POLICY): 

195 raise FileNotFoundError(f"Resource {self} does not exist") 

196 return ResourceInfo( 

197 uri=str(self), 

198 is_file=False, 

199 size=0, 

200 last_modified=None, 

201 checksums={}, 

202 ) 

203 

204 if self.dirLike: 

205 if not self.exists(): 

206 raise FileNotFoundError(f"Resource {self} does not exist") 

207 return ResourceInfo( 

208 uri=str(self), 

209 is_file=False, 

210 size=0, 

211 last_modified=None, 

212 checksums={}, 

213 ) 

214 

215 try: 

216 self.blob.reload(retry=_RETRY_POLICY) 

217 except NotFound: 

218 raise FileNotFoundError(f"Resource {self} does not exist") from None 

219 

220 size = self.blob.size 

221 if size is None: 

222 raise FileNotFoundError(f"Resource {self} does not exist") 

223 

224 checksums = {} 

225 if self.blob.md5_hash: 

226 checksums["md5"] = self.blob.md5_hash 

227 if self.blob.crc32c: 

228 checksums["crc32c"] = self.blob.crc32c 

229 

230 try: 

231 updated = _coerce_gcs_datetime(self.blob.updated) 

232 except ValueError: 

233 updated = _coerce_gcs_datetime(self.blob._properties.get("updated")) 

234 

235 return ResourceInfo( 

236 uri=str(self), 

237 is_file=True, 

238 size=size, 

239 last_modified=updated, 

240 checksums=checksums, 

241 ) 

242 

243 def remove(self) -> None: 

244 try: 

245 self.blob.delete(retry=_RETRY_POLICY) 

246 except NotFound as e: 

247 raise FileNotFoundError(f"No such resource: {self}") from e 

248 

249 def read(self, size: int = -1) -> bytes: 

250 if size < 0: 

251 start = None 

252 end = None 

253 else: 

254 start = 0 

255 end = size - 1 

256 try: 

257 with time_this(log, msg="Read from %s", args=(self,)): 

258 body = self.blob.download_as_bytes(start=start, end=end, retry=_RETRY_POLICY) 

259 except NotFound as e: 

260 raise FileNotFoundError(f"No such resource: {self}") from e 

261 return body 

262 

263 def write(self, data: bytes, overwrite: bool = True) -> None: 

264 if not overwrite and self.exists(): 

265 raise FileExistsError(f"Remote resource {self} exists and overwrite has been disabled") 

266 with time_this(log, msg="Write to %s", args=(self,)): 

267 self.blob.upload_from_string(data, retry=_RETRY_POLICY) 

268 

269 def mkdir(self) -> None: 

270 if not self.bucket.exists(retry=_RETRY_POLICY): 

271 raise ValueError(f"Bucket {self.netloc} does not exist for {self}!") 

272 

273 if not self.dirLike: 

274 raise NotADirectoryError(f"Can not create a 'directory' for a file-like URI {self}") 

275 

276 # GCS does not have directory objects, so mkdir is a no-op once the 

277 # bucket exists. 

278 return 

279 

280 @contextlib.contextmanager 

281 def _as_local( 

282 self, multithreaded: bool = True, tmpdir: ResourcePath | None = None 

283 ) -> Generator[ResourcePath]: 

284 with ( 

285 ResourcePath.temporary_uri(prefix=tmpdir, suffix=self.getExtension(), delete=True) as tmp_uri, 

286 time_this(log, msg="Downloading %s to local file", args=(self,)), 

287 ): 

288 try: 

289 with tmp_uri.open("wb") as tmpFile: 

290 self.blob.download_to_file(tmpFile, retry=_RETRY_POLICY) 

291 yield tmp_uri 

292 except NotFound as e: 

293 raise FileNotFoundError(f"No such resource: {self}") from e 

294 

295 def transfer_from( 

296 self, 

297 src: ResourcePath, 

298 transfer: str = "copy", 

299 overwrite: bool = False, 

300 transaction: TransactionProtocol | None = None, 

301 multithreaded: bool = True, 

302 ) -> None: 

303 if transfer not in self.transferModes: 

304 raise ValueError(f"Transfer mode '{transfer}' not supported by URI scheme {self.scheme}") 

305 

306 # Existence checks cost time so do not call this unless we know 

307 # that debugging is enabled. 

308 if log.isEnabledFor(logging.DEBUG): 

309 log.debug( 

310 "Transferring %s [exists: %s] -> %s [exists: %s] (transfer=%s)", 

311 src, 

312 src.exists(), 

313 self, 

314 self.exists(), 

315 transfer, 

316 ) 

317 

318 # Short circuit if the URIs are identical immediately. 

319 if self == src: 

320 log.debug( 

321 "Target and destination URIs are identical: %s, returning immediately." 

322 " No further action required.", 

323 self, 

324 ) 

325 return 

326 

327 if not overwrite and self.exists(): 

328 raise FileExistsError(f"Destination path '{self}' already exists.") 

329 

330 if transfer == "auto": 

331 transfer = self.transferDefault 

332 

333 timer_msg = "Transfer from %s to %s" 

334 timer_args = (src, self) 

335 

336 if isinstance(src, type(self)): 

337 # Looks like a GS remote uri so we can use direct copy 

338 with time_this(log, msg=timer_msg, args=timer_args): 

339 rewrite_token = None 

340 while True: 

341 try: 

342 rewrite_token, bytes_copied, total_bytes = self.blob.rewrite( 

343 src.blob, token=rewrite_token, retry=_RETRY_POLICY 

344 ) 

345 except NotFound as e: 

346 raise FileNotFoundError("No such resource to transfer: {self}") from e 

347 log.debug("Copied %d bytes out of %d (%s to %s)", bytes_copied, total_bytes, src, self) 

348 if rewrite_token is None: 

349 # Copy has completed 

350 break 

351 else: 

352 # Use local file and upload it 

353 with ( 

354 src.as_local(multithreaded=multithreaded) as local_uri, 

355 time_this(log, msg=timer_msg, args=timer_args), 

356 ): 

357 self.blob.upload_from_filename(local_uri.ospath, retry=_RETRY_POLICY) 

358 

359 # This was an explicit move requested from a remote resource 

360 # try to remove that resource 

361 if transfer == "move": 

362 # Transactions do not work here 

363 src.remove() 

364 

365 @contextlib.contextmanager 

366 def open( 

367 self, 

368 mode: str = "r", 

369 *, 

370 encoding: str | None = None, 

371 prefer_file_temporary: bool = False, 

372 ) -> Generator[ResourceHandleProtocol]: 

373 # Docstring inherited 

374 if self.isdir() or self.is_root: 

375 raise IsADirectoryError(f"Can not 'open' a directory URI: {self}") 

376 if "x" in mode: 

377 if self.exists(): 

378 raise FileExistsError(f"File at {self} already exists.") 

379 mode = mode.replace("x", "w") 

380 

381 # Clear the blob before calling open if we are in write mode. 

382 # This ensures that everything is resynced. 

383 if "w" in mode: 

384 self._blob = None 

385 

386 # The GCS API does not support append or read/write modes so for 

387 # those we use the base class implementation. 

388 # There seems to be a bug in the Google open() API where it does not 

389 # properly write a BOM at the start of the file in UTF-16 encoding 

390 # which leads to python not being able to read the contents back. 

391 if "+" in mode or "a" in mode or ("w" in mode and encoding == "utf-16"): 

392 with super().open(mode, encoding=encoding, prefer_file_temporary=prefer_file_temporary) as buffer: 

393 yield buffer 

394 else: 

395 with self.blob.open(mode, encoding=encoding, retry=_RETRY_POLICY) as buffer: 

396 yield buffer 

397 

398 def walk( 

399 self, file_filter: str | re.Pattern | None = None 

400 ) -> Iterator[list | tuple[ResourcePath, list[str], list[str]]]: 

401 # We pretend that GCS uses directories and files and not simply keys. 

402 if not (self.isdir() or self.is_root): 

403 raise ValueError(f"Can not walk a non-directory URI: {self}") 

404 

405 if isinstance(file_filter, str): 

406 file_filter = re.compile(file_filter) 

407 

408 # Limit each query to a single "directory" to match os.walk 

409 # We could download all keys at once with no delimiter and work 

410 # it out locally but this could potentially lead to large memory 

411 # usage for millions of keys. It will also make the initial call 

412 # to this method potentially very slow. If making this method look 

413 # like os.walk was not required, we could query all keys with 

414 # pagination and return them in groups of 1000, but that would 

415 # be a different interface since we can't guarantee we would get 

416 # them all grouped properly across the 1000 limit boundary. 

417 prefix = self.relativeToPathRoot if not self.is_root else "" 

418 prefix_len = len(prefix) 

419 dirnames: set[str] = set() 

420 filenames = [] 

421 files_there = False 

422 

423 blobs = self.client.list_blobs(self.bucket, prefix=prefix, delimiter="/", retry=_RETRY_POLICY) 

424 for page in blobs.pages: 

425 # "Sub-directories" turn up as prefixes in each page. 

426 dirnames.update(dir[prefix_len:] for dir in page.prefixes) 

427 

428 # Files are reported for this "directory" only. 

429 # The prefix itself can be included as a file because we write 

430 # a zero-length file for mkdir(). These must be filtered out. 

431 found_files = [f.name[prefix_len:] for f in page if f.name != prefix] 

432 if file_filter is not None: 

433 found_files = [f for f in found_files if file_filter.search(f)] 

434 if found_files: 

435 files_there = True 

436 

437 filenames.extend(found_files) 

438 

439 if not dirnames and not files_there: 

440 # Nothing found so match os.walk and return immediately. 

441 return 

442 else: 

443 yield self, sorted(dirnames), filenames 

444 

445 for dir in sorted(dirnames): 

446 new_uri = self.join(dir) 

447 yield from new_uri.walk(file_filter)