codekingpro/portable-devtools
115k
1from __future__ import annotations2 3import inspect4import logging5import os6import tempfile7import time8import weakref9from collections.abc import Callable10from shutil import rmtree11from typing import TYPE_CHECKING, Any, ClassVar12 13from fsspec import filesystem14from fsspec.callbacks import DEFAULT_CALLBACK15from fsspec.compression import compr16from fsspec.core import BaseCache, MMapCache17from fsspec.exceptions import BlocksizeMismatchError18from fsspec.implementations.cache_mapper import create_cache_mapper19from fsspec.implementations.cache_metadata import CacheMetadata20from fsspec.implementations.chained import ChainedFileSystem21from fsspec.implementations.local import LocalFileSystem22from fsspec.spec import AbstractBufferedFile23from fsspec.transaction import Transaction24from fsspec.utils import infer_compression25 26if TYPE_CHECKING:27 from fsspec.implementations.cache_mapper import AbstractCacheMapper28 29logger = logging.getLogger("fsspec.cached")30 31 32class WriteCachedTransaction(Transaction):33 def complete(self, commit=True):34 rpaths = [f.path for f in self.files]35 lpaths = [f.fn for f in self.files]36 if commit:37 self.fs.put(lpaths, rpaths)38 self.files.clear()39 self.fs._intrans = False40 self.fs._transaction = None41 self.fs = None # break cycle42 43 44class CachingFileSystem(ChainedFileSystem):45 """Locally caching filesystem, layer over any other FS46 47 This class implements chunk-wise local storage of remote files, for quick48 access after the initial download. The files are stored in a given49 directory with hashes of URLs for the filenames. If no directory is given,50 a temporary one is used, which should be cleaned up by the OS after the51 process ends. The files themselves are sparse (as implemented in52 :class:`~fsspec.caching.MMapCache`), so only the data which is accessed53 takes up space.54 55 Restrictions:56 57 - the block-size must be the same for each access of a given file, unless58 all blocks of the file have already been read59 - caching can only be applied to file-systems which produce files60 derived from fsspec.spec.AbstractBufferedFile ; LocalFileSystem is also61 allowed, for testing62 """63 64 protocol: ClassVar[str | tuple[str, ...]] = ("blockcache", "cached")65 _strip_tokenize_options = ("fo",)66 67 def __init__(68 self,69 target_protocol=None,70 cache_storage="TMP",71 cache_check=10,72 check_files=False,73 expiry_time=604800,74 target_options=None,75 fs=None,76 same_names: bool | None = None,77 compression=None,78 cache_mapper: AbstractCacheMapper | None = None,79 **kwargs,80 ):81 """82 83 Parameters84 ----------85 target_protocol: str (optional)86 Target filesystem protocol. Provide either this or ``fs``.87 cache_storage: str or list(str)88 Location to store files. If "TMP", this is a temporary directory,89 and will be cleaned up by the OS when this process ends (or later).90 If a list, each location will be tried in the order given, but91 only the last will be considered writable.92 cache_check: int93 Number of seconds between reload of cache metadata94 check_files: bool95 Whether to explicitly see if the UID of the remote file matches96 the stored one before using. Warning: some file systems such as97 HTTP cannot reliably give a unique hash of the contents of some98 path, so be sure to set this option to False.99 expiry_time: int100 The time in seconds after which a local copy is considered useless.101 Set to falsy to prevent expiry. The default is equivalent to one102 week.103 target_options: dict or None104 Passed to the instantiation of the FS, if fs is None.105 fs: filesystem instance106 The target filesystem to run against. Provide this or ``protocol``.107 same_names: bool (optional)108 By default, target URLs are hashed using a ``HashCacheMapper`` so109 that files from different backends with the same basename do not110 conflict. If this argument is ``true``, a ``BasenameCacheMapper``111 is used instead. Other cache mapper options are available by using112 the ``cache_mapper`` keyword argument. Only one of this and113 ``cache_mapper`` should be specified.114 compression: str (optional)115 To decompress on download. Can be 'infer' (guess from the URL name),116 one of the entries in ``fsspec.compression.compr``, or None for no117 decompression.118 cache_mapper: AbstractCacheMapper (optional)119 The object use to map from original filenames to cached filenames.120 Only one of this and ``same_names`` should be specified.121 """122 super().__init__(**kwargs)123 if fs is None and target_protocol is None:124 raise ValueError(125 "Please provide filesystem instance(fs) or target_protocol"126 )127 if not (fs is None) ^ (target_protocol is None):128 raise ValueError(129 "Both filesystems (fs) and target_protocol may not be both given."130 )131 if cache_storage == "TMP":132 tempdir = tempfile.mkdtemp()133 storage = [tempdir]134 weakref.finalize(self, self._remove_tempdir, tempdir)135 else:136 if isinstance(cache_storage, str):137 storage = [cache_storage]138 else:139 storage = cache_storage140 os.makedirs(storage[-1], exist_ok=True)141 self.storage = storage142 self.kwargs = target_options or {}143 self.cache_check = cache_check144 self.check_files = check_files145 self.expiry = expiry_time146 self.compression = compression147 148 # Size of cache in bytes. If None then the size is unknown and will be149 # recalculated the next time cache_size() is called. On writes to the150 # cache this is reset to None.151 self._cache_size = None152 153 if same_names is not None and cache_mapper is not None:154 raise ValueError(155 "Cannot specify both same_names and cache_mapper in "156 "CachingFileSystem.__init__"157 )158 if cache_mapper is not None:159 self._mapper = cache_mapper160 else:161 self._mapper = create_cache_mapper(162 same_names if same_names is not None else False163 )164 165 self.target_protocol = (166 target_protocol167 if isinstance(target_protocol, str)168 else (fs.protocol if isinstance(fs.protocol, str) else fs.protocol[0])169 )170 self._metadata = CacheMetadata(self.storage)171 self.load_cache()172 self.fs = fs if fs is not None else filesystem(target_protocol, **self.kwargs)173 174 def _strip_protocol(path):175 # acts as a method, since each instance has a difference target176 return self.fs._strip_protocol(type(self)._strip_protocol(path))177 178 self._strip_protocol: Callable = _strip_protocol179 180 @staticmethod181 def _remove_tempdir(tempdir):182 try:183 rmtree(tempdir)184 except Exception:185 pass186 187 def _mkcache(self):188 os.makedirs(self.storage[-1], exist_ok=True)189 190 def cache_size(self):191 """Return size of cache in bytes.192 193 If more than one cache directory is in use, only the size of the last194 one (the writable cache directory) is returned.195 """196 if self._cache_size is None:197 cache_dir = self.storage[-1]198 self._cache_size = filesystem("file").du(cache_dir, withdirs=True)199 return self._cache_size200 201 def load_cache(self):202 """Read set of stored blocks from file"""203 self._metadata.load()204 self._mkcache()205 self.last_cache = time.time()206 207 def save_cache(self):208 """Save set of stored blocks from file"""209 self._mkcache()210 self._metadata.save()211 self.last_cache = time.time()212 self._cache_size = None213 214 def _check_cache(self):215 """Reload caches if time elapsed or any disappeared"""216 self._mkcache()217 if not self.cache_check:218 # explicitly told not to bother checking219 return220 timecond = time.time() - self.last_cache > self.cache_check221 existcond = all(os.path.exists(storage) for storage in self.storage)222 if timecond or not existcond:223 self.load_cache()224 225 def _check_file(self, path):226 """Is path in cache and still valid"""227 path = self._strip_protocol(path)228 self._check_cache()229 return self._metadata.check_file(path, self)230 231 def clear_cache(self):232 """Remove all files and metadata from the cache233 234 In the case of multiple cache locations, this clears only the last one,235 which is assumed to be the read/write one.236 """237 rmtree(self.storage[-1])238 self.load_cache()239 self._cache_size = None240 241 def clear_expired_cache(self, expiry_time=None):242 """Remove all expired files and metadata from the cache243 244 In the case of multiple cache locations, this clears only the last one,245 which is assumed to be the read/write one.246 247 Parameters248 ----------249 expiry_time: int250 The time in seconds after which a local copy is considered useless.251 If not defined the default is equivalent to the attribute from the252 file caching instantiation.253 """254 255 if not expiry_time:256 expiry_time = self.expiry257 258 self._check_cache()259 260 expired_files, writable_cache_empty = self._metadata.clear_expired(expiry_time)261 for fn in expired_files:262 if os.path.exists(fn):263 os.remove(fn)264 265 if writable_cache_empty:266 rmtree(self.storage[-1])267 self.load_cache()268 269 self._cache_size = None270 271 def pop_from_cache(self, path):272 """Remove cached version of given file273 274 Deletes local copy of the given (remote) path. If it is found in a cache275 location which is not the last, it is assumed to be read-only, and276 raises PermissionError277 """278 path = self._strip_protocol(path)279 fn = self._metadata.pop_file(path)280 if fn is not None:281 os.remove(fn)282 self._cache_size = None283 284 def _open(285 self,286 path,287 mode="rb",288 block_size=None,289 autocommit=True,290 cache_options=None,291 **kwargs,292 ):293 """Wrap the target _open294 295 If the whole file exists in the cache, just open it locally and296 return that.297 298 Otherwise, open the file on the target FS, and make it have a mmap299 cache pointing to the location which we determine, in our cache.300 The ``blocks`` instance is shared, so as the mmap cache instance301 updates, so does the entry in our ``cached_files`` attribute.302 We monkey-patch this file, so that when it closes, we call303 ``close_and_update`` to save the state of the blocks.304 """305 path = self._strip_protocol(path)306 307 path = self.fs._strip_protocol(path)308 if "r" not in mode:309 return self.fs._open(310 path,311 mode=mode,312 block_size=block_size,313 autocommit=autocommit,314 cache_options=cache_options,315 **kwargs,316 )317 detail = self._check_file(path)318 if detail:319 # file is in cache320 detail, fn = detail321 hash, blocks = detail["fn"], detail["blocks"]322 if blocks is True:323 # stored file is complete324 logger.debug("Opening local copy of %s", path)325 return open(fn, mode)326 # TODO: action where partial file exists in read-only cache327 logger.debug("Opening partially cached copy of %s", path)328 else:329 hash = self._mapper(path)330 fn = os.path.join(self.storage[-1], hash)331 blocks = set()332 detail = {333 "original": path,334 "fn": hash,335 "blocks": blocks,336 "time": time.time(),337 "uid": self.fs.ukey(path),338 }339 self._metadata.update_file(path, detail)340 logger.debug("Creating local sparse file for %s", path)341 342 # explicitly submitting the size to the open call will avoid extra343 # operations when opening. This is particularly relevant344 # for any file that is read over a network, e.g. S3.345 size = detail.get("size")346 347 # call target filesystems open348 self._mkcache()349 f = self.fs._open(350 path,351 mode=mode,352 block_size=block_size,353 autocommit=autocommit,354 cache_options=cache_options,355 cache_type="none",356 size=size,357 **kwargs,358 )359 360 # set size if not already set361 if size is None:362 detail["size"] = f.size363 self._metadata.update_file(path, detail)364 365 if self.compression:366 comp = (367 infer_compression(path)368 if self.compression == "infer"369 else self.compression370 )371 f = compr[comp](f, mode="rb")372 if "blocksize" in detail:373 if detail["blocksize"] != f.blocksize:374 raise BlocksizeMismatchError(375 f"Cached file must be reopened with same block"376 f" size as original (old: {detail['blocksize']},"377 f" new {f.blocksize})"378 )379 else:380 detail["blocksize"] = f.blocksize381 382 def _fetch_ranges(ranges):383 return self.fs.cat_ranges(384 [path] * len(ranges),385 [r[0] for r in ranges],386 [r[1] for r in ranges],387 **kwargs,388 )389 390 multi_fetcher = None if self.compression else _fetch_ranges391 f.cache = MMapCache(392 f.blocksize, f._fetch_range, f.size, fn, blocks, multi_fetcher=multi_fetcher393 )394 close = f.close395 f.close = lambda: self.close_and_update(f, close)396 self.save_cache()397 return f398 399 def _parent(self, path):400 return self.fs._parent(path)401 402 def hash_name(self, path: str, *args: Any) -> str:403 # Kept for backward compatibility with downstream libraries.404 # Ignores extra arguments, previously same_name boolean.405 return self._mapper(path)406 407 def close_and_update(self, f, close):408 """Called when a file is closing, so store the set of blocks"""409 if f.closed:410 return411 path = self._strip_protocol(f.path)412 self._metadata.on_close_cached_file(f, path)413 try:414 logger.debug("going to save")415 self.save_cache()416 logger.debug("saved")417 except OSError:418 logger.debug("Cache saving failed while closing file")419 except NameError:420 logger.debug("Cache save failed due to interpreter shutdown")421 close()422 f.closed = True423 424 def ls(self, path, detail=True):425 return self.fs.ls(path, detail)426 427 def __getattribute__(self, item):428 if item in {429 "load_cache",430 "_get_cached_file_before_open",431 "_open",432 "save_cache",433 "close_and_update",434 "__init__",435 "__getattribute__",436 "__reduce__",437 "_make_local_details",438 "open",439 "cat",440 "cat_file",441 "_cat_file",442 "cat_ranges",443 "_cat_ranges",444 "get",445 "read_block",446 "tail",447 "head",448 "info",449 "ls",450 "exists",451 "isfile",452 "isdir",453 "_check_file",454 "_check_cache",455 "_mkcache",456 "clear_cache",457 "clear_expired_cache",458 "pop_from_cache",459 "local_file",460 "_paths_from_path",461 "get_mapper",462 "open_many",463 "commit_many",464 "hash_name",465 "__hash__",466 "__eq__",467 "to_json",468 "to_dict",469 "cache_size",470 "pipe_file",471 "pipe",472 "start_transaction",473 "end_transaction",474 }:475 # all the methods defined in this class. Note `open` here, since476 # it calls `_open`, but is actually in superclass477 if hasattr(type(self), item):478 return lambda *args, **kw: getattr(type(self), item).__get__(self)(479 *args, **kw480 )481 # method is in the whitelist but not defined on this subclass;482 # fall through to delegate to the wrapped filesystem below483 if item in ["__reduce_ex__"]:484 raise AttributeError485 if item in ["transaction"]:486 # property487 return type(self).transaction.__get__(self)488 if item in {"_cache", "transaction_type", "protocol"}:489 # class attributes490 return getattr(type(self), item)491 if item == "__class__":492 return type(self)493 d = object.__getattribute__(self, "__dict__")494 fs = d.get("fs", None) # fs is not immediately defined495 if item in d:496 return d[item]497 elif fs is not None:498 if item in fs.__dict__:499 # attribute of instance500 return fs.__dict__[item]501 # attributed belonging to the target filesystem502 cls = type(fs)503 m = getattr(cls, item)504 if (inspect.isfunction(m) or inspect.isdatadescriptor(m)) and (505 not hasattr(m, "__self__") or m.__self__ is None506 ):507 # instance method508 return m.__get__(fs, cls)509 return m # class method or attribute510 else:511 # attributes of the superclass, while target is being set up512 return super().__getattribute__(item)513 514 def __eq__(self, other):515 """Test for equality."""516 if self is other:517 return True518 if not isinstance(other, type(self)):519 return False520 return (521 self.storage == other.storage522 and self.kwargs == other.kwargs523 and self.cache_check == other.cache_check524 and self.check_files == other.check_files525 and self.expiry == other.expiry526 and self.compression == other.compression527 and self._mapper == other._mapper528 and self.target_protocol == other.target_protocol529 )530 531 def __hash__(self):532 """Calculate hash."""533 return (534 hash(tuple(self.storage))535 ^ hash(str(self.kwargs))536 ^ hash(self.cache_check)537 ^ hash(self.check_files)538 ^ hash(self.expiry)539 ^ hash(self.compression)540 ^ hash(self._mapper)541 ^ hash(self.target_protocol)542 )543 544 545class WholeFileCacheFileSystem(CachingFileSystem):546 """Caches whole remote files on first access547 548 This class is intended as a layer over any other file system, and549 will make a local copy of each file accessed, so that all subsequent550 reads are local. This is similar to ``CachingFileSystem``, but without551 the block-wise functionality and so can work even when sparse files552 are not allowed. See its docstring for definition of the init553 arguments.554 555 The class still needs access to the remote store for listing files,556 and may refresh cached files.557 """558 559 protocol = "filecache"560 local_file = True561 562 def open_many(self, open_files, **kwargs):563 paths = [of.path for of in open_files]564 if "r" in open_files.mode:565 self._mkcache()566 else:567 return [568 LocalTempFile(569 self.fs,570 path,571 mode=open_files.mode,572 fn=os.path.join(self.storage[-1], self._mapper(path)),573 **kwargs,574 )575 for path in paths576 ]577 578 if self.compression:579 raise NotImplementedError580 details = [self._check_file(sp) for sp in paths]581 downpath = [p for p, d in zip(paths, details) if not d]582 downfn0 = [583 os.path.join(self.storage[-1], self._mapper(p))584 for p, d in zip(paths, details)585 ] # keep these path names for opening later586 downfn = [fn for fn, d in zip(downfn0, details) if not d]587 if downpath:588 # skip if all files are already cached and up to date589 self.fs.get(downpath, downfn)590 591 # update metadata - only happens when downloads are successful592 newdetail = [593 {594 "original": path,595 "fn": self._mapper(path),596 "blocks": True,597 "time": time.time(),598 "uid": self.fs.ukey(path),599 }600 for path in downpath601 ]602 for path, detail in zip(downpath, newdetail):603 self._metadata.update_file(path, detail)604 self.save_cache()605 606 def firstpart(fn):607 # helper to adapt both whole-file and simple-cache608 return fn[1] if isinstance(fn, tuple) else fn609 610 return [611 open(firstpart(fn0) if fn0 else fn1, mode=open_files.mode)612 for fn0, fn1 in zip(details, downfn0)613 ]614 615 def commit_many(self, open_files):616 self.fs.put([f.fn for f in open_files], [f.path for f in open_files])617 [f.close() for f in open_files]618 for f in open_files:619 # in case autocommit is off, and so close did not already delete620 try:621 os.remove(f.name)622 except FileNotFoundError:623 pass624 self._cache_size = None625 626 def _make_local_details(self, path):627 hash = self._mapper(path)628 fn = os.path.join(self.storage[-1], hash)629 detail = {630 "original": path,631 "fn": hash,632 "blocks": True,633 "time": time.time(),634 "uid": self.fs.ukey(path),635 }636 self._metadata.update_file(path, detail)637 logger.debug("Copying %s to local cache", path)638 return fn639 640 def cat(641 self,642 path,643 recursive=False,644 on_error="raise",645 callback=DEFAULT_CALLBACK,646 **kwargs,647 ):648 paths = self.expand_path(649 path, recursive=recursive, maxdepth=kwargs.get("maxdepth")650 )651 getpaths = []652 storepaths = []653 fns = []654 out = {}655 for p in paths.copy():656 try:657 detail = self._check_file(p)658 if not detail:659 fn = self._make_local_details(p)660 getpaths.append(p)661 storepaths.append(fn)662 else:663 detail, fn = detail if isinstance(detail, tuple) else (None, detail)664 fns.append(fn)665 except Exception as e:666 if on_error == "raise":667 raise668 if on_error == "return":669 out[p] = e670 paths.remove(p)671 672 if getpaths:673 self.fs.get(getpaths, storepaths)674 self.save_cache()675 676 callback.set_size(len(paths))677 for p, fn in zip(paths, fns):678 with open(fn, "rb") as f:679 out[p] = f.read()680 callback.relative_update(1)681 if isinstance(path, str) and len(paths) == 1 and recursive is False:682 out = out[paths[0]]683 return out684 685 def _get_cached_file_before_open(self, path, **kwargs):686 fn = self._make_local_details(path)687 # call target filesystems open688 self._mkcache()689 if self.compression:690 with self.fs._open(path, mode="rb", **kwargs) as f, open(fn, "wb") as f2:691 if isinstance(f, AbstractBufferedFile):692 # want no type of caching if just downloading whole thing693 f.cache = BaseCache(0, f.cache.fetcher, f.size)694 comp = (695 infer_compression(path)696 if self.compression == "infer"697 else self.compression698 )699 f = compr[comp](f, mode="rb")700 data = True701 while data:702 block = getattr(f, "blocksize", 5 * 2**20)703 data = f.read(block)704 f2.write(data)705 else:706 self.fs.get_file(path, fn)707 self.save_cache()708 709 def _open(self, path, mode="rb", **kwargs):710 path = self._strip_protocol(path)711 # For read (or append), (try) download from remote712 if "r" in mode or "a" in mode:713 if not self._check_file(path):714 if self.fs.exists(path):715 self._get_cached_file_before_open(path, **kwargs)716 elif "r" in mode:717 raise FileNotFoundError(path)718 719 detail, fn = self._check_file(path)720 _, blocks = detail["fn"], detail["blocks"]721 if blocks is True:722 logger.debug("Opening local copy of %s", path)723 else:724 raise ValueError(725 f"Attempt to open partially cached file {path}"726 f" as a wholly cached file"727 )728 729 # Just reading does not need special file handling730 if "r" in mode and "+" not in mode:731 # In order to support downstream filesystems to be able to732 # infer the compression from the original filename, like733 # the `TarFileSystem`, let's extend the `io.BufferedReader`734 # fileobject protocol by adding a dedicated attribute735 # `original`.736 f = open(fn, mode)737 f.original = detail.get("original")738 return f739 740 hash = self._mapper(path)741 fn = os.path.join(self.storage[-1], hash)742 user_specified_kwargs = {743 k: v744 for k, v in kwargs.items()745 # those kwargs were added by open(), we don't want them746 if k not in ["autocommit", "block_size", "cache_options"]747 }748 return LocalTempFile(self, path, mode=mode, fn=fn, **user_specified_kwargs)749 750 async def _cat_file(self, path, start=None, end=None, **kwargs):751 logger.debug("async cat_file %s", path)752 path = self._strip_protocol(path)753 sha = self._mapper(path)754 fn = self._check_file(path)755 756 if not fn:757 fn = os.path.join(self.storage[-1], sha)758 await self.fs._get_file(path, fn, **kwargs)759 760 with open(fn, "rb") as f: # noqa ASYNC230761 if start:762 f.seek(start)763 size = -1 if end is None else end - f.tell()764 return f.read(size)765 766 async def _cat_ranges(767 self, paths, starts, ends, max_gap=None, on_error="return", **kwargs768 ):769 logger.debug("async cat ranges %s", paths)770 lpaths = []771 rset = set()772 download = []773 rpaths = []774 for p in paths:775 fn = self._check_file(p)776 if fn is None and p not in rset:777 sha = self._mapper(p)778 fn = os.path.join(self.storage[-1], sha)779 download.append(fn)780 rset.add(p)781 rpaths.append(p)782 lpaths.append(fn)783 if download:784 await self.fs._get(rpaths, download, on_error=on_error)785 786 return LocalFileSystem().cat_ranges(787 lpaths, starts, ends, max_gap=max_gap, on_error=on_error, **kwargs788 )789 790 791class SimpleCacheFileSystem(WholeFileCacheFileSystem):792 """Caches whole remote files on first access793 794 This class is intended as a layer over any other file system, and795 will make a local copy of each file accessed, so that all subsequent796 reads are local. This implementation only copies whole files, and797 does not keep any metadata about the download time or file details.798 It is therefore safer to use in multi-threaded/concurrent situations.799 800 This is the only of the caching filesystems that supports write: you will801 be given a real local open file, and upon close and commit, it will be802 uploaded to the target filesystem; the writability or the target URL is803 not checked until that time.804 805 """806 807 protocol = "simplecache"808 local_file = True809 transaction_type = WriteCachedTransaction810 811 def __init__(self, **kwargs):812 kw = kwargs.copy()813 for key in ["cache_check", "expiry_time", "check_files"]:814 kw[key] = False815 super().__init__(**kw)816 for storage in self.storage:817 if not os.path.exists(storage):818 os.makedirs(storage, exist_ok=True)819 820 def _check_file(self, path):821 self._check_cache()822 sha = self._mapper(path)823 for storage in self.storage:824 fn = os.path.join(storage, sha)825 if os.path.exists(fn):826 return fn827 828 def save_cache(self):829 pass830 831 def load_cache(self):832 pass833 834 def pipe_file(self, path, value=None, **kwargs):835 if self._intrans:836 with self.open(path, "wb") as f:837 f.write(value)838 else:839 super().pipe_file(path, value)840 841 def ls(self, path, detail=True, **kwargs):842 path = self._strip_protocol(path)843 details = []844 try:845 details = self.fs.ls(846 path, detail=True, **kwargs847 ).copy() # don't edit original!848 except FileNotFoundError as e:849 ex = e850 else:851 ex = None852 if self._intrans:853 path1 = path.rstrip("/") + "/"854 for f in self.transaction.files:855 if f.path == path:856 details.append(857 {"name": path, "size": f.size or f.tell(), "type": "file"}858 )859 elif f.path.startswith(path1):860 if f.path.count("/") == path1.count("/"):861 details.append(862 {"name": f.path, "size": f.size or f.tell(), "type": "file"}863 )864 else:865 dname = "/".join(f.path.split("/")[: path1.count("/") + 1])866 details.append({"name": dname, "size": 0, "type": "directory"})867 if ex is not None and not details:868 raise ex869 if detail:870 return details871 return sorted(_["name"] for _ in details)872 873 def info(self, path, **kwargs):874 path = self._strip_protocol(path)875 if self._intrans:876 f = [_ for _ in self.transaction.files if _.path == path]877 if f:878 size = os.path.getsize(f[0].fn) if f[0].closed else f[0].tell()879 return {"name": path, "size": size, "type": "file"}880 f = any(_.path.startswith(path + "/") for _ in self.transaction.files)881 if f:882 return {"name": path, "size": 0, "type": "directory"}883 return self.fs.info(path, **kwargs)884 885 def pipe(self, path, value=None, **kwargs):886 if isinstance(path, str):887 self.pipe_file(self._strip_protocol(path), value, **kwargs)888 elif isinstance(path, dict):889 for k, v in path.items():890 self.pipe_file(self._strip_protocol(k), v, **kwargs)891 else:892 raise ValueError("path must be str or dict")893 894 def cat_ranges(895 self, paths, starts, ends, max_gap=None, on_error="return", **kwargs896 ):897 logger.debug("cat ranges %s", paths)898 lpaths = [self._check_file(p) for p in paths]899 rpaths = [p for l, p in zip(lpaths, paths) if l is False]900 lpaths = [l for l, p in zip(lpaths, paths) if l is False]901 self.fs.get(rpaths, lpaths)902 paths = [self._check_file(p) for p in paths]903 return LocalFileSystem().cat_ranges(904 paths, starts, ends, max_gap=max_gap, on_error=on_error, **kwargs905 )906 907 def _get_cached_file_before_open(self, path, **kwargs):908 sha = self._mapper(path)909 fn = os.path.join(self.storage[-1], sha)910 logger.debug("Copying %s to local cache", path)911 912 self._mkcache()913 self._cache_size = None914 915 if self.compression:916 with self.fs._open(path, mode="rb", **kwargs) as f, open(fn, "wb") as f2:917 if isinstance(f, AbstractBufferedFile):918 # want no type of caching if just downloading whole thing919 f.cache = BaseCache(0, f.cache.fetcher, f.size)920 comp = (921 infer_compression(path)922 if self.compression == "infer"923 else self.compression924 )925 f = compr[comp](f, mode="rb")926 data = True927 while data:928 block = getattr(f, "blocksize", 5 * 2**20)929 data = f.read(block)930 f2.write(data)931 else:932 self.fs.get_file(path, fn)933 934 def _open(self, path, mode="rb", **kwargs):935 path = self._strip_protocol(path)936 sha = self._mapper(path)937 938 # For read (or append), (try) download from remote939 if "r" in mode or "a" in mode:940 if not self._check_file(path):941 # append does not require an existing file but read does942 if self.fs.exists(path):943 self._get_cached_file_before_open(path, **kwargs)944 elif "r" in mode:945 raise FileNotFoundError(path)946 947 fn = self._check_file(path)948 # Just reading does not need special file handling949 if "r" in mode and "+" not in mode:950 return open(fn, mode)951 952 fn = os.path.join(self.storage[-1], sha)953 user_specified_kwargs = {954 k: v955 for k, v in kwargs.items()956 if k not in ["autocommit", "block_size", "cache_options"]957 } # those were added by open()958 return LocalTempFile(959 self,960 path,961 mode=mode,962 autocommit=not self._intrans,963 fn=fn,964 **user_specified_kwargs,965 )966 967 968class LocalTempFile:969 """A temporary local file, which will be uploaded on commit"""970 971 def __init__(self, fs, path, fn, mode="wb", autocommit=True, seek=0, **kwargs):972 self.fn = fn973 self.fh = open(fn, mode)974 self.mode = mode975 if seek:976 self.fh.seek(seek)977 self.path = path978 self.size = None979 self.fs = fs980 self.closed = False981 self.autocommit = autocommit982 self.kwargs = kwargs983 984 def __reduce__(self):985 # always open in r+b to allow continuing writing at a location986 return (987 LocalTempFile,988 (self.fs, self.path, self.fn, "r+b", self.autocommit, self.tell()),989 )990 991 def __enter__(self):992 return self.fh993 994 def __exit__(self, exc_type, exc_val, exc_tb):995 self.close()996 997 def close(self):998 # self.size = self.fh.tell()999 if self.closed:1000 return1001 self.fh.close()1002 self.closed = True1003 if self.autocommit:1004 self.commit()1005 1006 def discard(self):1007 self.fh.close()1008 os.remove(self.fn)1009 1010 def commit(self):1011 # calling put() with list arguments avoids path expansion and additional operations1012 # like isdir()1013 self.fs.put([self.fn], [self.path], **self.kwargs)1014 # we do not delete the local copy, it's still in the cache.1015 1016 @property1017 def name(self):1018 return self.fn1019 1020 def __repr__(self) -> str:1021 return f"LocalTempFile: {self.path}"1022 1023 def __getattr__(self, item):1024 return getattr(self.fh, item)1025 