Coverage for python/lsst/resources/gs.py: 15%
220 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-23 02:10 -0700
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-23 02:10 -0700
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.
12"""Accessing Google Cloud Storage resources."""
14from __future__ import annotations
16__all__ = ("GSResourcePath",)
18import contextlib
19import datetime
20import logging
21import re
22from collections.abc import Generator, Iterator
23from typing import TYPE_CHECKING
25from ._resourceHandles._baseResourceHandle import ResourceHandleProtocol
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
44 # Must also fake the exception classes.
45 class ClientError(Exception):
46 """Generic client error."""
48 pass
50 class NotFound(ClientError): # noqa: N818
51 """Resource not found error."""
53 pass
55 class TooManyRequests(ClientError): # noqa: N818
56 """Too many requests error."""
58 pass
60 class InternalServerError(ClientError):
61 """Internal server error."""
63 pass
65 class BadGateway(ClientError): # noqa: N818
66 """Bad gateway error."""
68 pass
70 class ServiceUnavailable(ClientError): # noqa: N818
71 """Service unavailable error."""
73 pass
76from lsst.utils.timer import time_this
78from ._resourcePath import ResourceInfo, ResourcePath
80if TYPE_CHECKING:
81 from .utils import TransactionProtocol
83log = logging.getLogger(__name__)
86_RETRIEVABLE_TYPES = (
87 TooManyRequests, # 429
88 InternalServerError, # 500
89 BadGateway, # 502
90 ServiceUnavailable, # 503
91)
94def is_retryable(exc: Exception) -> bool:
95 """Report if the given exception is a condition that can be retried.
97 Parameters
98 ----------
99 exc : `Exception`
100 Exception to check.
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)
111_RETRY_POLICY = retry.Retry(predicate=is_retryable) if retry is not None else None
114_client = None
115"""Cached client connection."""
118def _coerce_gcs_datetime(value: datetime.datetime | str | None) -> datetime.datetime | None:
119 """Convert GCS timestamp values to timezone-aware UTC datetimes.
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)
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
145class GSResourcePath(ResourcePath):
146 """Access Google Cloud Storage resources."""
148 _bucket: storage.Bucket | None = None
149 _blob: storage.Blob | None = None
150 _client: storage.Client | None = None
152 @property
153 def client(self) -> storage.Client:
154 return _get_client()
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
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
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)
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
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 )
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 )
215 try:
216 self.blob.reload(retry=_RETRY_POLICY)
217 except NotFound:
218 raise FileNotFoundError(f"Resource {self} does not exist") from None
220 size = self.blob.size
221 if size is None:
222 raise FileNotFoundError(f"Resource {self} does not exist")
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
230 try:
231 updated = _coerce_gcs_datetime(self.blob.updated)
232 except ValueError:
233 updated = _coerce_gcs_datetime(self.blob._properties.get("updated"))
235 return ResourceInfo(
236 uri=str(self),
237 is_file=True,
238 size=size,
239 last_modified=updated,
240 checksums=checksums,
241 )
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
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
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)
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}!")
273 if not self.dirLike:
274 raise NotADirectoryError(f"Can not create a 'directory' for a file-like URI {self}")
276 # GCS does not have directory objects, so mkdir is a no-op once the
277 # bucket exists.
278 return
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
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}")
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 )
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
327 if not overwrite and self.exists():
328 raise FileExistsError(f"Destination path '{self}' already exists.")
330 if transfer == "auto":
331 transfer = self.transferDefault
333 timer_msg = "Transfer from %s to %s"
334 timer_args = (src, self)
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)
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()
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")
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
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
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}")
405 if isinstance(file_filter, str):
406 file_filter = re.compile(file_filter)
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
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)
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
437 filenames.extend(found_files)
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
445 for dir in sorted(dirnames):
446 new_uri = self.join(dir)
447 yield from new_uri.walk(file_filter)