Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/fsspec/spec.py: 25%
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
1from __future__ import annotations
3import io
4import json
5import logging
6import os
7import threading
8import warnings
9import weakref
10from errno import ESPIPE
11from glob import has_magic
12from hashlib import sha256
13from typing import Any
15from .callbacks import DEFAULT_CALLBACK
16from .config import apply_config, conf
17from .dircache import DirCache
18from .transaction import Transaction
19from .utils import (
20 _unstrip_protocol,
21 check_contained,
22 glob_translate,
23 isfilelike,
24 other_paths,
25 read_block,
26 stringify_path,
27 tokenize,
28)
30logger = logging.getLogger("fsspec")
33def make_instance(cls, args, kwargs):
34 return cls(*args, **kwargs)
37FORK_AVAILABLE = hasattr(os, "register_at_fork")
40if FORK_AVAILABLE:
41 _registered_classes = weakref.WeakSet()
43 def _reset_instances_lock():
44 for cls in _registered_classes:
45 cls._instantiation_lock = threading.RLock()
46 cls._cache.clear()
47 cls._pid = os.getpid()
49 os.register_at_fork(after_in_child=_reset_instances_lock)
52class _Cached(type):
53 """
54 Metaclass for caching file system instances.
56 Notes
57 -----
58 Instances are cached according to
60 * The values of the class attributes listed in `_extra_tokenize_attributes`
61 * The arguments passed to ``__init__``.
63 This creates an additional reference to the filesystem, which prevents the
64 filesystem from being garbage collected when all *user* references go away.
65 A call to the :meth:`AbstractFileSystem.clear_instance_cache` must *also*
66 be made for a filesystem instance to be garbage collected.
67 """
69 def __init__(cls, *args, **kwargs):
70 super().__init__(*args, **kwargs)
72 # Note: we intentionally create a reference here, to avoid garbage
73 # collecting instances when all other references are gone. To really
74 # delete a FileSystem, the cache must be cleared.
75 if conf.get("weakref_instance_cache"): # pragma: no cover
76 # debug option for analysing fork/spawn conditions
77 cls._cache = weakref.WeakValueDictionary()
78 else:
79 cls._cache = {}
80 cls._pid = os.getpid()
81 cls._instantiation_lock = threading.RLock()
83 if FORK_AVAILABLE:
84 _registered_classes.add(cls)
86 def _check_instance_cache(cls, token):
87 inst = cls._cache.get(token)
88 if inst is not None:
89 cls._latest = token
90 return inst
92 def __call__(cls, *args, **kwargs):
93 kwargs = apply_config(cls, kwargs)
94 extra_tokens = tuple(
95 getattr(cls, attr, None) for attr in cls._extra_tokenize_attributes
96 )
97 strip_tokenize_options = {
98 k: kwargs.pop(k) for k in cls._strip_tokenize_options if k in kwargs
99 }
100 pid = os.getpid()
102 if getattr(cls, "async_impl", False) and not kwargs.get("asynchronous", False):
103 token = tokenize(cls, pid, *args, *extra_tokens, **kwargs)
104 else:
105 token = tokenize(
106 cls, pid, threading.get_ident(), *args, *extra_tokens, **kwargs
107 )
108 skip = kwargs.pop("skip_instance_cache", False)
110 if pid != cls._pid:
111 with cls._instantiation_lock:
112 if pid != cls._pid:
113 cls._cache.clear()
114 cls._pid = pid
116 if not skip and cls.cachable:
117 inst = cls._check_instance_cache(token)
118 if inst is not None:
119 return inst
121 with cls._instantiation_lock:
122 # protect against the race condition that a new instance was created
123 # and inserted into the cache since the initial check just above
124 inst = cls._check_instance_cache(token)
125 if inst is not None:
126 return inst
128 obj = super().__call__(*args, **kwargs, **strip_tokenize_options)
129 # Setting _fs_token here causes some static linters to complain.
130 obj._fs_token_ = token
131 obj.storage_args = args
132 obj.storage_options = kwargs
133 if obj.async_impl and obj.mirror_sync_methods:
134 from .asyn import mirror_sync_methods
136 mirror_sync_methods(obj)
138 if cls.cachable and not skip:
139 with cls._instantiation_lock:
140 # another thread may have created the instance while we were calling
141 # super().__call__(), so we check again.
142 inst = cls._check_instance_cache(token)
143 if inst is not None:
144 return inst
146 cls._latest = token
147 cls._cache[token] = obj
148 return obj
151class AbstractFileSystem(metaclass=_Cached):
152 """
153 An abstract super-class for pythonic file-systems
155 Implementations are expected to be compatible with or, better, subclass
156 from here.
157 """
159 cachable = True # this class can be cached, instances reused
160 _cached = False
161 blocksize = 2**22
162 sep = "/"
163 # Implementations may select their protocol per instance (for example, a
164 # single adapter class backed by different storage implementations).
165 protocol: str | tuple[str, ...] = "abstract"
166 _latest = None
167 async_impl = False
168 mirror_sync_methods = False
169 root_marker = "" # For some FSs, may require leading '/' or other character
170 transaction_type = Transaction
172 #: Extra *class attributes* that should be considered when hashing.
173 _extra_tokenize_attributes = ()
174 #: *storage options* that should not be considered when hashing.
175 _strip_tokenize_options = ()
177 # Set by _Cached metaclass
178 storage_args: tuple[Any, ...]
179 storage_options: dict[str, Any]
181 def __init__(self, *args, **storage_options):
182 """Create and configure file-system instance
184 Instances may be cachable, so if similar enough arguments are seen
185 a new instance is not required. The token attribute exists to allow
186 implementations to cache instances if they wish.
188 A reasonable default should be provided if there are no arguments.
190 Subclasses should call this method.
192 Parameters
193 ----------
194 use_listings_cache, listings_expiry_time, max_paths:
195 passed to ``DirCache``, if the implementation supports
196 directory listing caching. Pass use_listings_cache=False
197 to disable such caching.
198 skip_instance_cache: bool
199 If this is a cachable implementation, pass True here to force
200 creating a new instance even if a matching instance exists, and prevent
201 storing this instance.
202 asynchronous: bool
203 loop: asyncio-compatible IOLoop or None
204 """
205 if self._cached:
206 # reusing instance, don't change
207 return
208 self._cached = True
209 self._intrans = False
210 self._transaction = None
211 self._invalidated_caches_in_transaction = []
212 self.dircache = DirCache(**storage_options)
214 if storage_options.pop("add_docs", None):
215 warnings.warn("add_docs is no longer supported.", FutureWarning)
217 if storage_options.pop("add_aliases", None):
218 warnings.warn("add_aliases has been removed.", FutureWarning)
219 # This is set in _Cached
220 self._fs_token_ = None
222 @property
223 def fsid(self):
224 """Persistent filesystem id that can be used to compare filesystems
225 across sessions.
226 """
227 raise NotImplementedError
229 @property
230 def _fs_token(self):
231 return self._fs_token_
233 def __dask_tokenize__(self):
234 return self._fs_token
236 def __hash__(self):
237 return int(self._fs_token, 16)
239 def __eq__(self, other):
240 return isinstance(other, type(self)) and self._fs_token == other._fs_token
242 def __reduce__(self):
243 return make_instance, (type(self), self.storage_args, self.storage_options)
245 @classmethod
246 def _strip_protocol(cls, path):
247 """Turn path from fully-qualified to file-system-specific
249 May require FS-specific handling, e.g., for relative paths or links.
250 """
251 if isinstance(path, list):
252 return [cls._strip_protocol(p) for p in path]
253 path = stringify_path(path)
254 protos = (cls.protocol,) if isinstance(cls.protocol, str) else cls.protocol
255 for protocol in protos:
256 if path.startswith(protocol + "://"):
257 path = path[len(protocol) + 3 :]
258 elif path.startswith(protocol + "::"):
259 path = path[len(protocol) + 2 :]
260 path = path.rstrip("/")
261 # use of root_marker to make minimum required path, e.g., "/"
262 return path or cls.root_marker
264 def unstrip_protocol(self, name: str) -> str:
265 """Format FS-specific path to generic, including protocol"""
266 protos = (self.protocol,) if isinstance(self.protocol, str) else self.protocol
267 for protocol in protos:
268 if name.startswith(f"{protocol}://"):
269 return name
270 return f"{protos[0]}://{name}"
272 @staticmethod
273 def _get_kwargs_from_urls(path):
274 """If kwargs can be encoded in the paths, extract them here
276 This should happen before instantiation of the class; incoming paths
277 then should be amended to strip the options in methods.
279 Examples may look like an sftp path "sftp://user@host:/my/path", where
280 the user and host should become kwargs and later get stripped.
281 """
282 # by default, nothing happens
283 return {}
285 @classmethod
286 def current(cls):
287 """Return the most recently instantiated FileSystem
289 If no instance has been created, then create one with defaults
290 """
291 inst = cls._cache.get(cls._latest)
292 if inst is not None:
293 return inst
294 return cls()
296 @property
297 def transaction(self):
298 """A context within which files are committed together upon exit
300 Requires the file class to implement `.commit()` and `.discard()`
301 for the normal and exception cases.
302 """
303 if self._transaction is None:
304 self._transaction = self.transaction_type(self)
305 return self._transaction
307 def start_transaction(self):
308 """Begin write transaction for deferring files, non-context version"""
309 self._intrans = True
310 self._transaction = self.transaction_type(self)
311 return self.transaction
313 def end_transaction(self):
314 """Finish write transaction, non-context version"""
315 self.transaction.complete()
316 self._transaction = None
317 # The invalid cache must be cleared after the transaction is completed.
318 for path in self._invalidated_caches_in_transaction:
319 self.invalidate_cache(path)
320 self._invalidated_caches_in_transaction.clear()
322 def invalidate_cache(self, path=None):
323 """
324 Discard any cached directory information
326 Parameters
327 ----------
328 path: string or None
329 If None, clear all listings cached else listings at or under given
330 path.
331 """
332 # Not necessary to implement invalidation mechanism, may have no cache.
333 # But if have, you should call this method of parent class from your
334 # subclass to ensure expiring caches after transacations correctly.
335 # See the implementation of FTPFileSystem in ftp.py
336 if self._intrans:
337 self._invalidated_caches_in_transaction.append(path)
339 def mkdir(self, path, create_parents=True, **kwargs):
340 """
341 Create directory entry at path
343 For systems that don't have true directories, may create an for
344 this instance only and not touch the real filesystem
346 Parameters
347 ----------
348 path: str
349 location
350 create_parents: bool
351 if True, this is equivalent to ``makedirs``
352 kwargs:
353 may be permissions, etc.
354 """
355 pass # not necessary to implement, may not have directories
357 def makedirs(self, path, exist_ok=False):
358 """Recursively make directories
360 Creates directory at path and any intervening required directories.
361 Raises exception if, for instance, the path already exists but is a
362 file.
364 Parameters
365 ----------
366 path: str
367 leaf directory name
368 exist_ok: bool (False)
369 If False, will error if the target already exists
370 """
371 pass # not necessary to implement, may not have directories
373 def rmdir(self, path):
374 """Remove a directory, if empty"""
375 pass # not necessary to implement, may not have directories
377 def ls(self, path, detail=True, **kwargs):
378 """List objects at path.
380 This should include subdirectories and files at that location. The
381 difference between a file and a directory must be clear when details
382 are requested.
384 The specific keys, or perhaps a FileInfo class, or similar, is TBD,
385 but must be consistent across implementations.
386 Must include:
388 - full path to the entry (without protocol)
389 - size of the entry, in bytes. If the value cannot be determined, will
390 be ``None``.
391 - type of entry, "file", "directory" or other
393 Additional information
394 may be present, appropriate to the file-system, e.g., generation,
395 checksum, etc.
397 May use refresh=True|False to allow use of self._ls_from_cache to
398 check for a saved listing and avoid calling the backend. This would be
399 common where listing may be expensive.
401 Parameters
402 ----------
403 path: str
404 detail: bool
405 if True, gives a list of dictionaries, where each is the same as
406 the result of ``info(path)``. If False, gives a list of paths
407 (str).
408 kwargs: may have additional backend-specific options, such as version
409 information
411 Returns
412 -------
413 List of strings if detail is False, or list of directory information
414 dicts if detail is True.
415 """
416 raise NotImplementedError
418 def _ls_from_cache(self, path):
419 """Check cache for listing
421 Returns listing, if found (may be empty list for a directly that exists
422 but contains nothing), None if not in cache.
423 """
424 parent = self._parent(path)
425 try:
426 return self.dircache[path.rstrip("/")]
427 except KeyError:
428 pass
429 try:
430 files = [
431 f
432 for f in self.dircache[parent]
433 if f["name"] == path
434 or (f["name"] == path.rstrip("/") and f["type"] == "directory")
435 ]
436 if len(files) == 0:
437 # parent dir was listed but did not contain this file
438 raise FileNotFoundError(path)
439 return files
440 except KeyError:
441 pass
443 def walk(self, path, maxdepth=None, topdown=True, on_error="omit", **kwargs):
444 """Return all files under the given path.
446 List all files, recursing into subdirectories; output is iterator-style,
447 like ``os.walk()``. For a simple list of files, ``find()`` is available.
449 When topdown is True, the caller can modify the dirnames list in-place (perhaps
450 using del or slice assignment), and walk() will
451 only recurse into the subdirectories whose names remain in dirnames;
452 this can be used to prune the search, impose a specific order of visiting,
453 or even to inform walk() about directories the caller creates or renames before
454 it resumes walk() again.
455 Modifying dirnames when topdown is False has no effect. (see os.walk)
457 Note that the "files" outputted will include anything that is not
458 a directory, such as links.
460 Parameters
461 ----------
462 path: str
463 Root to recurse into
464 maxdepth: int
465 Maximum recursion depth. None means limitless, but not recommended
466 on link-based file-systems.
467 topdown: bool (True)
468 Whether to walk the directory tree from the top downwards or from
469 the bottom upwards.
470 on_error: "omit", "raise", a callable
471 if omit (default), path with exception will simply be empty;
472 If raise, an underlying exception will be raised;
473 if callable, it will be called with a single OSError instance as argument
474 kwargs: passed to ``ls``
475 """
476 if maxdepth is not None and maxdepth < 1:
477 raise ValueError("maxdepth must be at least 1")
479 path = self._strip_protocol(path)
480 full_dirs = {}
481 dirs = {}
482 files = {}
484 detail = kwargs.pop("detail", False)
485 try:
486 listing = self.ls(path, detail=True, **kwargs)
487 except (FileNotFoundError, OSError) as e:
488 if on_error == "raise":
489 raise
490 if callable(on_error):
491 on_error(e)
492 return
494 for info in listing:
495 # each info name must be at least [path]/part , but here
496 # we check also for names like [path]/part/
497 pathname = info["name"].rstrip("/")
498 name = pathname.rsplit("/", 1)[-1]
499 if info["type"] == "directory" and pathname != path:
500 # do not include "self" path
501 full_dirs[name] = pathname
502 dirs[name] = info
503 elif pathname == path:
504 # file-like with same name as give path
505 files[""] = info
506 else:
507 files[name] = info
509 if not detail:
510 dirs = list(dirs)
511 files = list(files)
513 if topdown:
514 # Yield before recursion if walking top down
515 yield path, dirs, files
517 if maxdepth is not None:
518 maxdepth -= 1
519 if maxdepth < 1:
520 if not topdown:
521 yield path, dirs, files
522 return
524 for d in dirs:
525 yield from self.walk(
526 full_dirs[d],
527 maxdepth=maxdepth,
528 detail=detail,
529 topdown=topdown,
530 **kwargs,
531 )
533 if not topdown:
534 # Yield after recursion if walking bottom up
535 yield path, dirs, files
537 def find(self, path, maxdepth=None, withdirs=False, detail=False, **kwargs):
538 """List all files below path.
540 Like posix ``find`` command without conditions
542 Parameters
543 ----------
544 path : str
545 maxdepth: int or None
546 If not None, the maximum number of levels to descend
547 withdirs: bool
548 Whether to include directory paths in the output. This is True
549 when used by glob, but users usually only want files.
550 kwargs are passed to ``ls``.
551 """
552 # TODO: allow equivalent of -name parameter
553 path = self._strip_protocol(path)
554 out = {}
556 # Add the root directory if withdirs is requested
557 # This is needed for posix glob compliance
558 if withdirs and path != "" and self.isdir(path):
559 out[path] = self.info(path)
561 for _, dirs, files in self.walk(path, maxdepth, detail=True, **kwargs):
562 if withdirs:
563 files.update(dirs)
564 out.update({info["name"]: info for name, info in files.items()})
565 if not out and self.isfile(path):
566 # walk works on directories, but find should also return [path]
567 # when path happens to be a file
568 out[path] = {}
569 names = sorted(out)
570 if not detail:
571 return names
572 else:
573 return {name: out[name] for name in names}
575 def du(self, path, total=True, maxdepth=None, withdirs=False, **kwargs):
576 """Space used by files and optionally directories within a path
578 Directory size does not include the size of its contents.
580 Parameters
581 ----------
582 path: str
583 total: bool
584 Whether to sum all the file sizes
585 maxdepth: int or None
586 Maximum number of directory levels to descend, None for unlimited.
587 withdirs: bool
588 Whether to include directory paths in the output.
589 kwargs: passed to ``find``
591 Returns
592 -------
593 Dict of {path: size} if total=False, or int otherwise, where numbers
594 refer to bytes used.
595 """
596 sizes = {}
597 if withdirs and self.isdir(path):
598 # Include top-level directory in output
599 info = self.info(path)
600 sizes[info["name"]] = info["size"]
601 for f in self.find(path, maxdepth=maxdepth, withdirs=withdirs, **kwargs):
602 info = self.info(f)
603 sizes[info["name"]] = info["size"]
604 if total:
605 return sum(sizes.values())
606 else:
607 return sizes
609 def glob(self, path, maxdepth=None, **kwargs):
610 """Find files by glob-matching.
612 Pattern matching capabilities for finding files that match the given pattern.
614 Parameters
615 ----------
616 path: str
617 The glob pattern to match against
618 maxdepth: int or None
619 Maximum depth for ``'**'`` patterns. Applied on the first ``'**'`` found.
620 Must be at least 1 if provided.
621 kwargs:
622 Additional arguments passed to ``find`` (e.g., detail=True)
624 Returns
625 -------
626 List of matched paths, or dict of paths and their info if detail=True
628 Notes
629 -----
630 Supported patterns:
631 - '*': Matches any sequence of characters within a single directory level
632 - ``'**'``: Matches any number of directory levels (must be an entire path component)
633 - '?': Matches exactly one character
634 - '[abc]': Matches any character in the set
635 - '[a-z]': Matches any character in the range
636 - '[!abc]': Matches any character NOT in the set
638 Special behaviors:
639 - If the path ends with '/', only folders are returned
640 - Consecutive '*' characters are compressed into a single '*'
641 - Empty set '[]' or negated empty negated set '[!]' never match anything
642 - Special characters in character classes are escaped properly
644 Limitations:
645 - ``'**'`` must be a complete path component (e.g., ``'a/**/b'``, not ``'a**b'``)
646 - No brace expansion ('{a,b}.txt')
647 - No extended glob patterns ('+(pattern)', '!(pattern)')
648 """
649 if maxdepth is not None and maxdepth < 1:
650 raise ValueError("maxdepth must be at least 1")
652 import re
654 seps = (os.path.sep, os.path.altsep) if os.path.altsep else (os.path.sep,)
655 ends_with_sep = path.endswith(seps) # _strip_protocol strips trailing slash
656 path = self._strip_protocol(path)
657 append_slash_to_dirname = ends_with_sep or path.endswith(
658 tuple(sep + "**" for sep in seps)
659 )
660 idx_star = path.find("*") if path.find("*") >= 0 else len(path)
661 idx_qmark = path.find("?") if path.find("?") >= 0 else len(path)
662 idx_brace = path.find("[") if path.find("[") >= 0 else len(path)
664 min_idx = min(idx_star, idx_qmark, idx_brace)
666 detail = kwargs.pop("detail", False)
667 withdirs = kwargs.pop("withdirs", True)
669 if not has_magic(path):
670 if self.exists(path, **kwargs):
671 if not detail:
672 return [path]
673 else:
674 return {path: self.info(path, **kwargs)}
675 else:
676 if not detail:
677 return [] # glob of non-existent returns empty
678 else:
679 return {}
680 elif "/" in path[:min_idx]:
681 min_idx = path[:min_idx].rindex("/")
682 root = path[: min_idx + 1]
683 depth = path[min_idx + 1 :].count("/") + 1
684 else:
685 root = ""
686 depth = path[min_idx + 1 :].count("/") + 1
688 if "**" in path:
689 if maxdepth is not None:
690 idx_double_stars = path.find("**")
691 depth_double_stars = path[idx_double_stars:].count("/") + 1
692 depth = depth - depth_double_stars + maxdepth
693 else:
694 depth = None
696 allpaths = self.find(
697 root, maxdepth=depth, withdirs=withdirs, detail=True, **kwargs
698 )
700 pattern = glob_translate(path + ("/" if ends_with_sep else ""))
701 pattern = re.compile(pattern)
703 out = {
704 p: info
705 for p, info in sorted(allpaths.items())
706 if pattern.match(
707 p + "/"
708 if append_slash_to_dirname and info["type"] == "directory"
709 else p
710 )
711 }
713 if detail:
714 return out
715 else:
716 return list(out)
718 def exists(self, path, **kwargs):
719 """Is there a file at the given path"""
720 try:
721 self.info(path, **kwargs)
722 return True
723 except: # noqa: E722
724 # any exception allowed bar FileNotFoundError?
725 return False
727 def lexists(self, path, **kwargs):
728 """If there is a file at the given path (including
729 broken links)"""
730 return self.exists(path)
732 def info(self, path, **kwargs):
733 """Give details of entry at path
735 Returns a single dictionary, with exactly the same information as ``ls``
736 would with ``detail=True``.
738 The default implementation calls ls and could be overridden by a
739 shortcut. kwargs are passed on to ```ls()``.
741 Some file systems might not be able to measure the file's size, in
742 which case, the returned dict will include ``'size': None``.
744 Returns
745 -------
746 dict with keys: name (full path in the FS), size (in bytes), type (file,
747 directory, or something else) and other FS-specific keys.
748 """
749 path = self._strip_protocol(path)
750 out = self.ls(self._parent(path), detail=True, **kwargs)
751 out = [o for o in out if o["name"].rstrip("/") == path]
752 if out:
753 return out[0]
754 out = self.ls(path, detail=True, **kwargs)
755 path = path.rstrip("/")
756 out1 = [o for o in out if o["name"].rstrip("/") == path]
757 if len(out1) == 1:
758 if "size" not in out1[0]:
759 out1[0]["size"] = None
760 return out1[0]
761 elif len(out1) > 1 or out:
762 return {"name": path, "size": 0, "type": "directory"}
763 else:
764 raise FileNotFoundError(path)
766 def checksum(self, path):
767 """Unique value for current version of file
769 If the checksum is the same from one moment to another, the contents
770 are guaranteed to be the same. If the checksum changes, the contents
771 *might* have changed.
773 This should normally be overridden; default will probably capture
774 creation/modification timestamp (which would be good) or maybe
775 access timestamp (which would be bad)
776 """
777 return int(tokenize(self.info(path)), 16)
779 def size(self, path):
780 """Size in bytes of file"""
781 return self.info(path).get("size", None)
783 def sizes(self, paths):
784 """Size in bytes of each file in a list of paths"""
785 return [self.size(p) for p in paths]
787 def isdir(self, path):
788 """Is this entry directory-like?"""
789 try:
790 return self.info(path)["type"] == "directory"
791 except OSError:
792 return False
794 def isfile(self, path):
795 """Is this entry file-like?"""
796 try:
797 return self.info(path)["type"] == "file"
798 except Exception:
799 return False
801 def read_text(self, path, encoding=None, errors=None, newline=None, **kwargs):
802 """Get the contents of the file as a string.
804 Parameters
805 ----------
806 path: str
807 URL of file on this filesystems
808 encoding, errors, newline: same as `open`.
809 """
810 with self.open(
811 path,
812 mode="r",
813 encoding=encoding,
814 errors=errors,
815 newline=newline,
816 **kwargs,
817 ) as f:
818 return f.read()
820 def write_text(
821 self, path, value, encoding=None, errors=None, newline=None, **kwargs
822 ):
823 """Write the text to the given file.
825 An existing file will be overwritten.
827 Parameters
828 ----------
829 path: str
830 URL of file on this filesystems
831 value: str
832 Text to write.
833 encoding, errors, newline: same as `open`.
834 """
835 with self.open(
836 path,
837 mode="w",
838 encoding=encoding,
839 errors=errors,
840 newline=newline,
841 **kwargs,
842 ) as f:
843 return f.write(value)
845 def cat_file(self, path, start=None, end=None, **kwargs):
846 """Get the content of a file
848 Parameters
849 ----------
850 path: URL of file on this filesystems
851 start, end: int
852 Bytes limits of the read. If negative, backwards from end,
853 like usual python slices. Either can be None for start or
854 end of file, respectively
855 kwargs: passed to ``open()``.
856 """
857 # explicitly set buffering off?
858 with self.open(path, "rb", **kwargs) as f:
859 if start is not None:
860 if start >= 0:
861 f.seek(start)
862 else:
863 f.seek(max(0, f.size + start))
864 if end is not None:
865 if end < 0:
866 end = f.size + end
867 return f.read(end - f.tell())
868 return f.read()
870 def pipe_file(self, path, value, mode="overwrite", **kwargs):
871 """Set the bytes of given file"""
872 if mode == "create" and self.exists(path):
873 # non-atomic but simple way; or could use "xb" in open(), which is likely
874 # not as well supported
875 raise FileExistsError
876 with self.open(path, "wb", **kwargs) as f:
877 f.write(value)
879 def pipe(self, path, value=None, **kwargs):
880 """Put value into path
882 (counterpart to ``cat``)
884 Parameters
885 ----------
886 path: string or dict(str, bytes)
887 If a string, a single remote location to put ``value`` bytes; if a dict,
888 a mapping of {path: bytesvalue}.
889 value: bytes, optional
890 If using a single path, these are the bytes to put there. Ignored if
891 ``path`` is a dict
892 """
893 if isinstance(path, str):
894 self.pipe_file(self._strip_protocol(path), value, **kwargs)
895 elif isinstance(path, dict):
896 for k, v in path.items():
897 self.pipe_file(self._strip_protocol(k), v, **kwargs)
898 else:
899 raise ValueError("path must be str or dict")
901 def cat_ranges(
902 self, paths, starts, ends, max_gap=None, on_error="return", **kwargs
903 ):
904 """Get the contents of byte ranges from one or more files
906 Parameters
907 ----------
908 paths: list
909 A list of of filepaths on this filesystems
910 starts, ends: int or list
911 Bytes limits of the read. If using a single int, the same value will be
912 used to read all the specified files.
913 """
914 if max_gap is not None:
915 raise NotImplementedError
916 if not isinstance(paths, list):
917 raise TypeError
918 if not isinstance(starts, list):
919 starts = [starts] * len(paths)
920 if not isinstance(ends, list):
921 ends = [ends] * len(paths)
922 if len(starts) != len(paths) or len(ends) != len(paths):
923 raise ValueError
924 out = []
925 for p, s, e in zip(paths, starts, ends):
926 try:
927 out.append(self.cat_file(p, s, e, **kwargs))
928 except Exception as e:
929 if on_error == "return":
930 out.append(e)
931 else:
932 raise
933 return out
935 def cat(self, path, recursive=False, on_error="raise", **kwargs):
936 """Fetch (potentially multiple) paths' contents
938 Parameters
939 ----------
940 recursive: bool
941 If True, assume the path(s) are directories, and get all the
942 contained files
943 on_error : "raise", "omit", "return"
944 If raise, an underlying exception will be raised (converted to KeyError
945 if the type is in self.missing_exceptions); if omit, keys with exception
946 will simply not be included in the output; if "return", all keys are
947 included in the output, but the value will be bytes or an exception
948 instance.
949 kwargs: passed to cat_file
951 Returns
952 -------
953 dict of {path: contents} if there are multiple paths
954 or the path has been otherwise expanded
955 """
956 paths = self.expand_path(path, recursive=recursive, **kwargs)
957 if (
958 len(paths) > 1
959 or isinstance(path, list)
960 or paths[0] != self._strip_protocol(path)
961 ):
962 out = {}
963 for path in paths:
964 try:
965 out[path] = self.cat_file(path, **kwargs)
966 except Exception as e:
967 if on_error == "raise":
968 raise
969 if on_error == "return":
970 out[path] = e
971 return out
972 else:
973 return self.cat_file(paths[0], **kwargs)
975 def get_file(
976 self, rpath, lpath=None, callback=DEFAULT_CALLBACK, outfile=None, **kwargs
977 ):
978 """Copy single remote file to local"""
979 from .implementations.local import LocalFileSystem
981 if outfile is None and isfilelike(lpath):
982 outfile = lpath
983 elif outfile is None and self.isdir(rpath):
984 os.makedirs(lpath, exist_ok=True)
985 return None
987 if outfile is None:
988 fs = LocalFileSystem(auto_mkdir=True)
989 fs.makedirs(fs._parent(lpath), exist_ok=True)
991 with self.open(rpath, "rb", **kwargs) as f1:
992 close_outfile = outfile is None
993 if close_outfile:
994 outfile = open(lpath, "wb")
996 try:
997 callback.set_size(getattr(f1, "size", None))
998 data = True
999 while data:
1000 data = f1.read(self.blocksize)
1001 segment_len = outfile.write(data)
1002 if segment_len is None:
1003 segment_len = len(data)
1004 callback.relative_update(segment_len)
1005 finally:
1006 if close_outfile:
1007 outfile.close()
1009 def get(
1010 self,
1011 rpath,
1012 lpath,
1013 recursive=False,
1014 callback=DEFAULT_CALLBACK,
1015 maxdepth=None,
1016 **kwargs,
1017 ):
1018 """Copy file(s) to local.
1020 Copies a specific file or tree of files (if recursive=True). If lpath
1021 ends with a "/", it will be assumed to be a directory, and target files
1022 will go within. Can submit a list of paths, which may be glob-patterns
1023 and will be expanded.
1025 Calls get_file for each source.
1026 """
1027 if isinstance(lpath, list) and isinstance(rpath, list):
1028 # No need to expand paths when both source and destination
1029 # are provided as lists
1030 rpaths = rpath
1031 lpaths = lpath
1032 else:
1033 from .implementations.local import (
1034 LocalFileSystem,
1035 make_path_posix,
1036 trailing_sep,
1037 )
1039 source_is_str = isinstance(rpath, str)
1040 rpaths = self.expand_path(
1041 rpath, recursive=recursive, maxdepth=maxdepth, **kwargs
1042 )
1043 if source_is_str and (not recursive or maxdepth is not None):
1044 # Non-recursive glob does not copy directories
1045 rpaths = [p for p in rpaths if not (trailing_sep(p) or self.isdir(p))]
1046 if not rpaths:
1047 return
1049 if isinstance(lpath, str):
1050 lpath = make_path_posix(lpath)
1052 source_is_file = len(rpaths) == 1
1053 dest_is_dir = isinstance(lpath, str) and (
1054 trailing_sep(lpath) or LocalFileSystem().isdir(lpath)
1055 )
1057 exists = source_is_str and (
1058 (has_magic(rpath) and source_is_file)
1059 or (not has_magic(rpath) and dest_is_dir and not trailing_sep(rpath))
1060 )
1061 lpaths = other_paths(
1062 rpaths,
1063 lpath,
1064 exists=exists,
1065 flatten=not source_is_str,
1066 )
1067 if isinstance(lpath, str):
1068 # The names came from the source listing; ".." in one of them
1069 # would otherwise place the copy above the destination. When
1070 # lpath is a list the caller named every destination itself.
1071 check_contained(lpath, lpaths)
1073 callback.set_size(len(lpaths))
1074 for lpath, rpath in callback.wrap(zip(lpaths, rpaths)):
1075 with callback.branched(rpath, lpath) as child:
1076 self.get_file(rpath, lpath, callback=child, **kwargs)
1078 def put_file(
1079 self, lpath, rpath, callback=DEFAULT_CALLBACK, mode="overwrite", **kwargs
1080 ):
1081 """Copy single file to remote"""
1082 if mode == "create" and self.exists(rpath):
1083 raise FileExistsError
1084 if os.path.isdir(lpath):
1085 self.makedirs(rpath, exist_ok=True)
1086 return None
1088 with open(lpath, "rb") as f1:
1089 size = f1.seek(0, 2)
1090 callback.set_size(size)
1091 f1.seek(0)
1093 self.mkdirs(self._parent(os.fspath(rpath)), exist_ok=True)
1094 with self.open(rpath, "wb", **kwargs) as f2:
1095 while f1.tell() < size:
1096 data = f1.read(self.blocksize)
1097 segment_len = f2.write(data)
1098 if segment_len is None:
1099 segment_len = len(data)
1100 callback.relative_update(segment_len)
1102 def put(
1103 self,
1104 lpath,
1105 rpath,
1106 recursive=False,
1107 callback=DEFAULT_CALLBACK,
1108 maxdepth=None,
1109 **kwargs,
1110 ):
1111 """Copy file(s) from local.
1113 Copies a specific file or tree of files (if recursive=True). If rpath
1114 ends with a "/", it will be assumed to be a directory, and target files
1115 will go within.
1117 Calls put_file for each source.
1118 """
1119 if isinstance(lpath, list) and isinstance(rpath, list):
1120 # No need to expand paths when both source and destination
1121 # are provided as lists
1122 rpaths = rpath
1123 lpaths = lpath
1124 else:
1125 from .implementations.local import (
1126 LocalFileSystem,
1127 make_path_posix,
1128 trailing_sep,
1129 )
1131 source_is_str = isinstance(lpath, str)
1132 if source_is_str:
1133 lpath = make_path_posix(lpath)
1134 fs = LocalFileSystem()
1135 lpaths = fs.expand_path(
1136 lpath, recursive=recursive, maxdepth=maxdepth, **kwargs
1137 )
1138 if source_is_str and (not recursive or maxdepth is not None):
1139 # Non-recursive glob does not copy directories
1140 lpaths = [p for p in lpaths if not (trailing_sep(p) or fs.isdir(p))]
1141 if not lpaths:
1142 return
1144 source_is_file = len(lpaths) == 1
1145 dest_is_dir = isinstance(rpath, str) and (
1146 trailing_sep(rpath) or self.isdir(rpath)
1147 )
1149 rpath = (
1150 self._strip_protocol(rpath)
1151 if isinstance(rpath, str)
1152 else [self._strip_protocol(p) for p in rpath]
1153 )
1154 exists = source_is_str and (
1155 (has_magic(lpath) and source_is_file)
1156 or (not has_magic(lpath) and dest_is_dir and not trailing_sep(lpath))
1157 )
1158 rpaths = other_paths(
1159 lpaths,
1160 rpath,
1161 exists=exists,
1162 flatten=not source_is_str,
1163 )
1165 callback.set_size(len(rpaths))
1166 for lpath, rpath in callback.wrap(zip(lpaths, rpaths)):
1167 with callback.branched(lpath, rpath) as child:
1168 self.put_file(lpath, rpath, callback=child, **kwargs)
1170 def head(self, path, size=1024):
1171 """Get the first ``size`` bytes from file"""
1172 with self.open(path, "rb") as f:
1173 return f.read(size)
1175 def tail(self, path, size=1024):
1176 """Get the last ``size`` bytes from file"""
1177 with self.open(path, "rb") as f:
1178 f.seek(max(-size, -f.size), 2)
1179 return f.read()
1181 def cp_file(self, path1, path2, **kwargs):
1182 raise NotImplementedError
1184 def copy(
1185 self, path1, path2, recursive=False, maxdepth=None, on_error=None, **kwargs
1186 ):
1187 """Copy within two locations in the filesystem
1189 on_error : "raise", "ignore"
1190 If raise, any not-found exceptions will be raised; if ignore any
1191 not-found exceptions will cause the path to be skipped; defaults to
1192 raise unless recursive is true, where the default is ignore
1193 """
1194 if on_error is None and recursive:
1195 on_error = "ignore"
1196 elif on_error is None:
1197 on_error = "raise"
1199 if isinstance(path1, list) and isinstance(path2, list):
1200 # No need to expand paths when both source and destination
1201 # are provided as lists
1202 paths1 = path1
1203 paths2 = path2
1204 else:
1205 from .implementations.local import trailing_sep
1207 source_is_str = isinstance(path1, str)
1208 paths1 = self.expand_path(
1209 path1, recursive=recursive, maxdepth=maxdepth, **kwargs
1210 )
1211 if source_is_str and (not recursive or maxdepth is not None):
1212 # Non-recursive glob does not copy directories
1213 paths1 = [p for p in paths1 if not (trailing_sep(p) or self.isdir(p))]
1214 if not paths1:
1215 return
1217 source_is_file = len(paths1) == 1
1218 dest_is_dir = isinstance(path2, str) and (
1219 trailing_sep(path2) or self.isdir(path2)
1220 )
1222 exists = source_is_str and (
1223 (has_magic(path1) and source_is_file)
1224 or (not has_magic(path1) and dest_is_dir and not trailing_sep(path1))
1225 )
1226 paths2 = other_paths(
1227 paths1,
1228 path2,
1229 exists=exists,
1230 flatten=not source_is_str,
1231 )
1233 for p1, p2 in zip(paths1, paths2):
1234 try:
1235 self.cp_file(p1, p2, **kwargs)
1236 except FileNotFoundError:
1237 if on_error == "raise":
1238 raise
1240 def expand_path(
1241 self, path, recursive=False, maxdepth=None, assume_literal=False, **kwargs
1242 ):
1243 """Turn one or more globs or directories into a list of all matching paths
1244 to files or directories.
1246 kwargs are passed to ``glob`` or ``find``, which may in turn call ``ls``
1247 """
1249 if maxdepth is not None and maxdepth < 1:
1250 raise ValueError("maxdepth must be at least 1")
1252 if isinstance(path, (str, os.PathLike)):
1253 out = self.expand_path([path], recursive, maxdepth, **kwargs)
1254 else:
1255 out = set()
1256 path = [self._strip_protocol(p) for p in path]
1257 for p in path:
1258 if not assume_literal and has_magic(p):
1259 bit = set(self.glob(p, maxdepth=maxdepth, **kwargs))
1260 out |= bit
1261 if recursive:
1262 # glob call above expanded one depth so if maxdepth is defined
1263 # then decrement it in expand_path call below. If it is zero
1264 # after decrementing then avoid expand_path call.
1265 if maxdepth is not None and maxdepth <= 1:
1266 continue
1267 out |= set(
1268 self.expand_path(
1269 list(bit),
1270 recursive=recursive,
1271 maxdepth=maxdepth - 1 if maxdepth is not None else None,
1272 assume_literal=True,
1273 **kwargs,
1274 )
1275 )
1276 continue
1277 elif recursive:
1278 rec = set(
1279 self.find(
1280 p, maxdepth=maxdepth, withdirs=True, detail=False, **kwargs
1281 )
1282 )
1283 out |= rec
1284 if p not in out and (recursive is False or self.exists(p)):
1285 # should only check once, for the root
1286 out.add(p)
1287 if not out:
1288 raise FileNotFoundError(path)
1289 return sorted(out)
1291 def mv(self, path1, path2, recursive=False, maxdepth=None, **kwargs):
1292 """Move file(s) from one location to another"""
1293 if path1 == path2:
1294 logger.debug("%s mv: The paths are the same, so no files were moved.", self)
1295 else:
1296 # explicitly raise exception to prevent data corruption
1297 self.copy(
1298 path1, path2, recursive=recursive, maxdepth=maxdepth, on_error="raise"
1299 )
1300 self.rm(path1, recursive=recursive)
1302 def rm_file(self, path):
1303 """Delete a file"""
1304 self._rm(path)
1306 def _rm(self, path):
1307 """Delete one file"""
1308 # this is the old name for the method, prefer rm_file
1309 raise NotImplementedError
1311 def rm(self, path, recursive=False, maxdepth=None):
1312 """Delete files or directories.
1314 Parameters
1315 ----------
1316 path: str or list of str
1317 Files or directories to delete.
1318 recursive: bool
1319 If True, recursively delete directories and their contents.
1320 maxdepth: int or None
1321 Depth to pass to walk for finding files to delete, if recursive.
1322 If None, there will be no limit and infinite recursion may be
1323 possible.
1324 """
1325 path = self.expand_path(path, recursive=recursive, maxdepth=maxdepth)
1326 for p in reversed(path):
1327 self.rm_file(p)
1329 @classmethod
1330 def _parent(cls, path):
1331 path = cls._strip_protocol(path)
1332 if "/" in path:
1333 parent = path.rsplit("/", 1)[0].lstrip(cls.root_marker)
1334 return cls.root_marker + parent
1335 else:
1336 return cls.root_marker
1338 def _open(
1339 self,
1340 path,
1341 mode="rb",
1342 block_size=None,
1343 autocommit=True,
1344 cache_options=None,
1345 **kwargs,
1346 ):
1347 """Return raw bytes-mode file-like from the file-system"""
1348 return AbstractBufferedFile(
1349 self,
1350 path,
1351 mode,
1352 block_size,
1353 autocommit,
1354 cache_options=cache_options,
1355 **kwargs,
1356 )
1358 def open(
1359 self,
1360 path,
1361 mode="rb",
1362 block_size=None,
1363 cache_options=None,
1364 compression=None,
1365 **kwargs,
1366 ):
1367 """
1368 Return a file-like object from the filesystem
1370 The resultant instance must function correctly in a context ``with``
1371 block.
1373 Parameters
1374 ----------
1375 path: str
1376 Target file
1377 mode: str like 'rb', 'w'
1378 See builtin ``open()``
1379 Mode "x" (exclusive write) may be implemented by the backend. Even if
1380 it is, whether it is checked up front or on commit, and whether it is
1381 atomic is implementation-dependent.
1382 block_size: int
1383 Some indication of buffering - this is a value in bytes
1384 cache_options : dict, optional
1385 Extra arguments to pass through to the cache.
1386 compression: string or None
1387 If given, open file using compression codec. Can either be a compression
1388 name (a key in ``fsspec.compression.compr``) or "infer" to guess the
1389 compression from the filename suffix.
1390 encoding, errors, newline: passed on to TextIOWrapper for text mode
1391 """
1392 import io
1394 path = self._strip_protocol(path)
1395 if "b" not in mode:
1396 mode = mode.replace("t", "") + "b"
1398 text_kwargs = {
1399 k: kwargs.pop(k)
1400 for k in ["encoding", "errors", "newline"]
1401 if k in kwargs
1402 }
1403 return io.TextIOWrapper(
1404 self.open(
1405 path,
1406 mode,
1407 block_size=block_size,
1408 cache_options=cache_options,
1409 compression=compression,
1410 **kwargs,
1411 ),
1412 **text_kwargs,
1413 )
1414 else:
1415 ac = kwargs.pop("autocommit", not self._intrans)
1416 f = self._open(
1417 path,
1418 mode=mode,
1419 block_size=block_size,
1420 autocommit=ac,
1421 cache_options=cache_options,
1422 **kwargs,
1423 )
1424 if not ac and "r" not in mode:
1425 self.transaction.files.append(f)
1426 if compression is not None:
1427 from fsspec.compression import compr
1428 from fsspec.core import get_compression
1430 compression = get_compression(path, compression)
1431 compress = compr[compression]
1432 f = compress(f, mode=mode[0])
1433 return f
1435 def touch(self, path, truncate=True, **kwargs):
1436 """Create empty file, or update timestamp
1438 Parameters
1439 ----------
1440 path: str
1441 file location
1442 truncate: bool
1443 If True, always set file size to 0; if False, update timestamp and
1444 leave file unchanged, if backend allows this
1445 """
1446 if truncate or not self.exists(path):
1447 with self.open(path, "wb", **kwargs):
1448 pass
1449 else:
1450 raise NotImplementedError # update timestamp, if possible
1452 def ukey(self, path):
1453 """Hash of file properties, to tell if it has changed"""
1454 return sha256(str(self.info(path)).encode()).hexdigest()
1456 def read_block(self, fn, offset, length, delimiter=None):
1457 """Read a block of bytes from
1459 Starting at ``offset`` of the file, read ``length`` bytes. If
1460 ``delimiter`` is set then we ensure that the read starts and stops at
1461 delimiter boundaries that follow the locations ``offset`` and ``offset
1462 + length``. If ``offset`` is zero then we start at zero. The
1463 bytestring returned WILL include the end delimiter string.
1465 If offset+length is beyond the eof, reads to eof.
1467 Parameters
1468 ----------
1469 fn: string
1470 Path to filename
1471 offset: int
1472 Byte offset to start read
1473 length: int
1474 Number of bytes to read. If None, read to end.
1475 delimiter: bytes (optional)
1476 Ensure reading starts and stops at delimiter bytestring
1478 Examples
1479 --------
1480 >>> fs.read_block('data/file.csv', 0, 13) # doctest: +SKIP
1481 b'Alice, 100\\nBo'
1482 >>> fs.read_block('data/file.csv', 0, 13, delimiter=b'\\n') # doctest: +SKIP
1483 b'Alice, 100\\nBob, 200\\n'
1485 Use ``length=None`` to read to the end of the file.
1486 >>> fs.read_block('data/file.csv', 0, None, delimiter=b'\\n') # doctest: +SKIP
1487 b'Alice, 100\\nBob, 200\\nCharlie, 300'
1489 See Also
1490 --------
1491 :func:`fsspec.utils.read_block`
1492 """
1493 with self.open(fn, "rb") as f:
1494 size = f.size
1495 if length is None:
1496 length = size
1497 if size is not None and offset + length > size:
1498 length = size - offset
1499 return read_block(f, offset, length, delimiter)
1501 def to_json(self, *, include_password: bool = True) -> str:
1502 """
1503 JSON representation of this filesystem instance.
1505 Parameters
1506 ----------
1507 include_password: bool, default True
1508 Whether to include the password (if any) in the output.
1510 Returns
1511 -------
1512 JSON string with keys ``cls`` (the python location of this class),
1513 protocol (text name of this class's protocol, first one in case of
1514 multiple), ``args`` (positional args, usually empty), and all other
1515 keyword arguments as their own keys.
1517 Warnings
1518 --------
1519 Serialized filesystems may contain sensitive information which have been
1520 passed to the constructor, such as passwords and tokens. Make sure you
1521 store and send them in a secure environment!
1522 """
1523 from .json import FilesystemJSONEncoder
1525 return json.dumps(
1526 self,
1527 cls=type(
1528 "_FilesystemJSONEncoder",
1529 (FilesystemJSONEncoder,),
1530 {"include_password": include_password},
1531 ),
1532 )
1534 @staticmethod
1535 def from_json(blob: str) -> AbstractFileSystem:
1536 """
1537 Recreate a filesystem instance from JSON representation.
1539 See ``.to_json()`` for the expected structure of the input.
1541 Parameters
1542 ----------
1543 blob: str
1545 Returns
1546 -------
1547 file system instance, not necessarily of this particular class.
1549 Warnings
1550 --------
1551 This can import arbitrary modules (as determined by the ``cls`` key).
1552 Make sure you haven't installed any modules that may execute malicious code
1553 at import time.
1554 """
1555 from .json import FilesystemJSONDecoder
1557 return json.loads(blob, cls=FilesystemJSONDecoder)
1559 def to_dict(self, *, include_password: bool = True) -> dict[str, Any]:
1560 """
1561 JSON-serializable dictionary representation of this filesystem instance.
1563 Parameters
1564 ----------
1565 include_password: bool, default True
1566 Whether to include the password (if any) in the output.
1568 Returns
1569 -------
1570 Dictionary with keys ``cls`` (the python location of this class),
1571 protocol (text name of this class's protocol, first one in case of
1572 multiple), ``args`` (positional args, usually empty), and all other
1573 keyword arguments as their own keys.
1575 Warnings
1576 --------
1577 Serialized filesystems may contain sensitive information which have been
1578 passed to the constructor, such as passwords and tokens. Make sure you
1579 store and send them in a secure environment!
1580 """
1581 from .json import FilesystemJSONEncoder
1583 json_encoder = FilesystemJSONEncoder()
1585 cls = type(self)
1586 proto = self.protocol
1588 storage_options = dict(self.storage_options)
1589 if not include_password:
1590 storage_options.pop("password", None)
1592 return dict(
1593 cls=f"{cls.__module__}:{cls.__name__}",
1594 protocol=proto[0] if isinstance(proto, (tuple, list)) else proto,
1595 args=json_encoder.make_serializable(self.storage_args),
1596 **json_encoder.make_serializable(storage_options),
1597 )
1599 @staticmethod
1600 def from_dict(dct: dict[str, Any]) -> AbstractFileSystem:
1601 """
1602 Recreate a filesystem instance from dictionary representation.
1604 See ``.to_dict()`` for the expected structure of the input.
1606 Parameters
1607 ----------
1608 dct: Dict[str, Any]
1610 Returns
1611 -------
1612 file system instance, not necessarily of this particular class.
1614 Warnings
1615 --------
1616 This can import arbitrary modules (as determined by the ``cls`` key).
1617 Make sure you haven't installed any modules that may execute malicious code
1618 at import time.
1619 """
1620 from .json import FilesystemJSONDecoder
1622 json_decoder = FilesystemJSONDecoder()
1624 dct = dict(dct) # Defensive copy
1626 cls = FilesystemJSONDecoder.try_resolve_fs_cls(dct)
1627 if cls is None:
1628 raise ValueError("Not a serialized AbstractFileSystem")
1630 dct.pop("cls", None)
1631 dct.pop("protocol", None)
1633 return cls(
1634 *json_decoder.unmake_serializable(dct.pop("args", ())),
1635 **json_decoder.unmake_serializable(dct),
1636 )
1638 def _get_pyarrow_filesystem(self):
1639 """
1640 Make a version of the FS instance which will be acceptable to pyarrow
1641 """
1642 # all instances already also derive from pyarrow
1643 return self
1645 def get_mapper(self, root="", check=False, create=False, missing_exceptions=None):
1646 """Create key/value store based on this file-system
1648 Makes a MutableMapping interface to the FS at the given root path.
1649 See ``fsspec.mapping.FSMap`` for further details.
1650 """
1651 from .mapping import FSMap
1653 return FSMap(
1654 root,
1655 self,
1656 check=check,
1657 create=create,
1658 missing_exceptions=missing_exceptions,
1659 )
1661 @classmethod
1662 def clear_instance_cache(cls):
1663 """
1664 Clear the cache of filesystem instances.
1666 Notes
1667 -----
1668 Unless overridden by setting the ``cachable`` class attribute to False,
1669 the filesystem class stores a reference to newly created instances. This
1670 prevents Python's normal rules around garbage collection from working,
1671 since the instances refcount will not drop to zero until
1672 ``clear_instance_cache`` is called.
1673 """
1674 cls._cache.clear()
1676 def created(self, path):
1677 """Return the created timestamp of a file as a datetime.datetime"""
1678 raise NotImplementedError
1680 def modified(self, path):
1681 """Return the modified timestamp of a file as a datetime.datetime"""
1682 raise NotImplementedError
1684 def tree(
1685 self,
1686 path: str = "/",
1687 recursion_limit: int = 2,
1688 max_display: int = 25,
1689 display_size: bool = False,
1690 prefix: str = "",
1691 is_last: bool = True,
1692 first: bool = True,
1693 indent_size: int = 4,
1694 ) -> str:
1695 """
1696 Return a tree-like structure of the filesystem starting from the given path as a string.
1698 Parameters
1699 ----------
1700 path: Root path to start traversal from
1701 recursion_limit: Maximum depth of directory traversal
1702 max_display: Maximum number of items to display per directory
1703 display_size: Whether to display file sizes
1704 prefix: Current line prefix for visual tree structure
1705 is_last: Whether current item is last in its level
1706 first: Whether this is the first call (displays root path)
1707 indent_size: Number of spaces by indent
1709 Returns
1710 -------
1711 str: A string representing the tree structure.
1713 Example
1714 -------
1715 >>> from fsspec import filesystem
1717 >>> fs = filesystem('ftp', host='test.rebex.net', user='demo', password='password')
1718 >>> tree = fs.tree(display_size=True, recursion_limit=3, indent_size=8, max_display=10)
1719 >>> print(tree)
1720 """
1722 def format_bytes(n: int) -> str:
1723 """Format bytes as text."""
1724 for prefix, k in (
1725 ("P", 2**50),
1726 ("T", 2**40),
1727 ("G", 2**30),
1728 ("M", 2**20),
1729 ("k", 2**10),
1730 ):
1731 if n >= 0.9 * k:
1732 return f"{n / k:.2f} {prefix}b"
1733 return f"{n}B"
1735 result = []
1737 if first:
1738 result.append(path)
1740 if recursion_limit:
1741 indent = " " * indent_size
1742 contents = self.ls(path, detail=True)
1743 contents.sort(
1744 key=lambda x: (x.get("type") != "directory", x.get("name", ""))
1745 )
1747 if max_display is not None and len(contents) > max_display:
1748 displayed_contents = contents[:max_display]
1749 remaining_count = len(contents) - max_display
1750 else:
1751 displayed_contents = contents
1752 remaining_count = 0
1754 for i, item in enumerate(displayed_contents):
1755 is_last_item = (i == len(displayed_contents) - 1) and (
1756 remaining_count == 0
1757 )
1759 branch = (
1760 "└" + ("─" * (indent_size - 2))
1761 if is_last_item
1762 else "├" + ("─" * (indent_size - 2))
1763 )
1764 branch += " "
1765 new_prefix = prefix + (
1766 indent if is_last_item else "│" + " " * (indent_size - 1)
1767 )
1769 name = os.path.basename(item.get("name", ""))
1771 if display_size and item.get("type") == "directory":
1772 sub_contents = self.ls(item.get("name", ""), detail=True)
1773 num_files = sum(
1774 1 for sub_item in sub_contents if sub_item.get("type") == "file"
1775 )
1776 num_folders = sum(
1777 1
1778 for sub_item in sub_contents
1779 if sub_item.get("type") == "directory"
1780 )
1782 if num_files == 0 and num_folders == 0:
1783 size = " (empty folder)"
1784 elif num_files == 0:
1785 size = f" ({num_folders} subfolder{'s' if num_folders > 1 else ''})"
1786 elif num_folders == 0:
1787 size = f" ({num_files} file{'s' if num_files > 1 else ''})"
1788 else:
1789 size = f" ({num_files} file{'s' if num_files > 1 else ''}, {num_folders} subfolder{'s' if num_folders > 1 else ''})"
1790 elif display_size and item.get("type") == "file":
1791 size = f" ({format_bytes(item.get('size', 0))})"
1792 else:
1793 size = ""
1795 result.append(f"{prefix}{branch}{name}{size}")
1797 if item.get("type") == "directory" and recursion_limit > 0:
1798 result.append(
1799 self.tree(
1800 path=item.get("name", ""),
1801 recursion_limit=recursion_limit - 1,
1802 max_display=max_display,
1803 display_size=display_size,
1804 prefix=new_prefix,
1805 is_last=is_last_item,
1806 first=False,
1807 indent_size=indent_size,
1808 )
1809 )
1811 if remaining_count > 0:
1812 more_message = f"{remaining_count} more item(s) not displayed."
1813 result.append(
1814 f"{prefix}{'└' + ('─' * (indent_size - 2))} {more_message}"
1815 )
1817 return "\n".join(_ for _ in result if _)
1819 # ------------------------------------------------------------------------
1820 # Aliases
1822 def read_bytes(self, path, start=None, end=None, **kwargs):
1823 """Alias of `AbstractFileSystem.cat_file`."""
1824 return self.cat_file(path, start=start, end=end, **kwargs)
1826 def write_bytes(self, path, value, **kwargs):
1827 """Alias of `AbstractFileSystem.pipe_file`."""
1828 self.pipe_file(path, value, **kwargs)
1830 def makedir(self, path, create_parents=True, **kwargs):
1831 """Alias of `AbstractFileSystem.mkdir`."""
1832 return self.mkdir(path, create_parents=create_parents, **kwargs)
1834 def mkdirs(self, path, exist_ok=False):
1835 """Alias of `AbstractFileSystem.makedirs`."""
1836 return self.makedirs(path, exist_ok=exist_ok)
1838 def listdir(self, path, detail=True, **kwargs):
1839 """Alias of `AbstractFileSystem.ls`."""
1840 return self.ls(path, detail=detail, **kwargs)
1842 def cp(self, path1, path2, **kwargs):
1843 """Alias of `AbstractFileSystem.copy`."""
1844 return self.copy(path1, path2, **kwargs)
1846 def move(self, path1, path2, **kwargs):
1847 """Alias of `AbstractFileSystem.mv`."""
1848 return self.mv(path1, path2, **kwargs)
1850 def stat(self, path, **kwargs):
1851 """Alias of `AbstractFileSystem.info`."""
1852 return self.info(path, **kwargs)
1854 def disk_usage(self, path, total=True, maxdepth=None, **kwargs):
1855 """Alias of `AbstractFileSystem.du`."""
1856 return self.du(path, total=total, maxdepth=maxdepth, **kwargs)
1858 def rename(self, path1, path2, **kwargs):
1859 """Alias of `AbstractFileSystem.mv`."""
1860 return self.mv(path1, path2, **kwargs)
1862 def delete(self, path, recursive=False, maxdepth=None):
1863 """Alias of `AbstractFileSystem.rm`."""
1864 return self.rm(path, recursive=recursive, maxdepth=maxdepth)
1866 def upload(self, lpath, rpath, recursive=False, **kwargs):
1867 """Alias of `AbstractFileSystem.put`."""
1868 return self.put(lpath, rpath, recursive=recursive, **kwargs)
1870 def download(self, rpath, lpath, recursive=False, **kwargs):
1871 """Alias of `AbstractFileSystem.get`."""
1872 return self.get(rpath, lpath, recursive=recursive, **kwargs)
1874 def sign(self, path, expiration=100, **kwargs):
1875 """Create a signed URL representing the given path
1877 Some implementations allow temporary URLs to be generated, as a
1878 way of delegating credentials.
1880 Parameters
1881 ----------
1882 path : str
1883 The path on the filesystem
1884 expiration : int
1885 Number of seconds to enable the URL for (if supported)
1887 Returns
1888 -------
1889 URL : str
1890 The signed URL
1892 Raises
1893 ------
1894 NotImplementedError : if method is not implemented for a filesystem
1895 """
1896 raise NotImplementedError("Sign is not implemented for this filesystem")
1898 def _isfilestore(self):
1899 # Originally inherited from pyarrow DaskFileSystem. Keeping this
1900 # here for backwards compatibility as long as pyarrow uses its
1901 # legacy fsspec-compatible filesystems and thus accepts fsspec
1902 # filesystems as well
1903 return False
1906class AbstractBufferedFile(io.IOBase):
1907 """Convenient class to derive from to provide buffering
1909 In the case that the backend does not provide a pythonic file-like object
1910 already, this class contains much of the logic to build one. The only
1911 methods that need to be overridden are ``_upload_chunk``,
1912 ``_initiate_upload`` and ``_fetch_range``.
1913 """
1915 DEFAULT_BLOCK_SIZE = 5 * 2**20
1916 _details = None
1918 def __init__(
1919 self,
1920 fs,
1921 path,
1922 mode="rb",
1923 block_size="default",
1924 autocommit=True,
1925 cache_type="readahead",
1926 cache_options=None,
1927 size=None,
1928 **kwargs,
1929 ):
1930 """
1931 Template for files with buffered reading and writing
1933 Parameters
1934 ----------
1935 fs: instance of FileSystem
1936 path: str
1937 location in file-system
1938 mode: str
1939 Normal file modes. Currently only 'wb', 'ab' or 'rb'. Some file
1940 systems may be read-only, and some may not support append.
1941 block_size: int
1942 Buffer size for reading or writing, 'default' for class default
1943 autocommit: bool
1944 Whether to write to final destination; may only impact what
1945 happens when file is being closed.
1946 cache_type: {"readahead", "none", "mmap", "bytes"}, default "readahead"
1947 Caching policy in read mode. See the definitions in ``core``.
1948 cache_options : dict
1949 Additional options passed to the constructor for the cache specified
1950 by `cache_type`.
1951 size: int
1952 If given and in read mode, suppressed having to look up the file size
1953 kwargs:
1954 Gets stored as self.kwargs
1955 """
1956 from .core import caches
1958 self.path = path
1959 self.fs = fs
1960 self.mode = mode
1961 self.blocksize = (
1962 self.DEFAULT_BLOCK_SIZE if block_size in ["default", None] else block_size
1963 )
1964 self.loc = 0
1965 self.autocommit = autocommit
1966 self.end = None
1967 self.start = None
1968 self.closed = False
1970 if cache_options is None:
1971 cache_options = {}
1973 if "trim" in kwargs:
1974 warnings.warn(
1975 "Passing 'trim' to control the cache behavior has been deprecated. "
1976 "Specify it within the 'cache_options' argument instead.",
1977 FutureWarning,
1978 )
1979 cache_options["trim"] = kwargs.pop("trim")
1981 self.kwargs = kwargs
1983 if mode not in {"ab", "rb", "wb", "xb"}:
1984 raise NotImplementedError("File mode not supported")
1985 if mode == "rb":
1986 if size is not None:
1987 self.size = size
1988 else:
1989 self.size = self.details["size"]
1990 self.cache = caches[cache_type](
1991 self.blocksize, self._fetch_range, self.size, **cache_options
1992 )
1993 else:
1994 self.buffer = io.BytesIO()
1995 self.offset = None
1996 self.forced = False
1997 self.location = None
1999 @property
2000 def details(self):
2001 if self._details is None:
2002 self._details = self.fs.info(self.path)
2003 return self._details
2005 @details.setter
2006 def details(self, value):
2007 self._details = value
2008 self.size = value["size"]
2010 @property
2011 def full_name(self):
2012 return _unstrip_protocol(self.path, self.fs)
2014 @property
2015 def closed(self):
2016 # get around this attr being read-only in IOBase
2017 # use getattr here, since this can be called during del
2018 return getattr(self, "_closed", True)
2020 @closed.setter
2021 def closed(self, c):
2022 self._closed = c
2024 def __hash__(self):
2025 if "w" in self.mode:
2026 return id(self)
2027 else:
2028 return int(tokenize(self.details), 16)
2030 def __eq__(self, other):
2031 """Files are equal if they have the same checksum, only in read mode"""
2032 if self is other:
2033 return True
2034 return (
2035 isinstance(other, type(self))
2036 and self.mode == "rb"
2037 and other.mode == "rb"
2038 and hash(self) == hash(other)
2039 )
2041 def commit(self):
2042 """Move from temp to final destination"""
2044 def discard(self):
2045 """Throw away temporary file"""
2047 def info(self):
2048 """File information about this path"""
2049 if self.readable():
2050 return self.details
2051 else:
2052 raise ValueError("Info not available while writing")
2054 def tell(self):
2055 """Current file location"""
2056 return self.loc
2058 def seek(self, loc, whence=0):
2059 """Set current file location
2061 Parameters
2062 ----------
2063 loc: int
2064 byte location
2065 whence: {0, 1, 2}
2066 from start of file, current location or end of file, resp.
2067 """
2068 loc = int(loc)
2069 if not self.mode == "rb":
2070 raise OSError(ESPIPE, "Seek only available in read mode")
2071 if whence == 0:
2072 nloc = loc
2073 elif whence == 1:
2074 nloc = self.loc + loc
2075 elif whence == 2:
2076 nloc = self.size + loc
2077 else:
2078 raise ValueError(f"invalid whence ({whence}, should be 0, 1 or 2)")
2079 if nloc < 0:
2080 raise ValueError("Seek before start of file")
2081 self.loc = nloc
2082 return self.loc
2084 def write(self, data):
2085 """
2086 Write data to buffer.
2088 Buffer only sent on flush() or if buffer is greater than
2089 or equal to blocksize.
2091 Parameters
2092 ----------
2093 data: bytes
2094 Set of bytes to be written.
2095 """
2096 if not self.writable():
2097 raise ValueError("File not in write mode")
2098 if self.closed:
2099 raise ValueError("I/O operation on closed file.")
2100 if self.forced:
2101 raise ValueError("This file has been force-flushed, can only close")
2102 out = self.buffer.write(data)
2103 self.loc += out
2104 if self.buffer.tell() >= self.blocksize:
2105 self.flush()
2106 return out
2108 def flush(self, force=False):
2109 """
2110 Write buffered data to backend store.
2112 Writes the current buffer, if it is larger than the block-size, or if
2113 the file is being closed.
2115 Parameters
2116 ----------
2117 force: bool
2118 When closing, write the last block even if it is smaller than
2119 blocks are allowed to be. Disallows further writing to this file.
2120 """
2122 if self.closed:
2123 raise ValueError("Flush on closed file")
2124 if force and self.forced:
2125 raise ValueError("Force flush cannot be called more than once")
2126 if force:
2127 self.forced = True
2129 if self.readable():
2130 # no-op to flush on read-mode
2131 return
2133 if not force and self.buffer.tell() < self.blocksize:
2134 # Defer write on small block
2135 return
2137 if self.offset is None:
2138 # Initialize a multipart upload
2139 self.offset = 0
2140 try:
2141 self._initiate_upload()
2142 except Exception:
2143 self.closed = True
2144 raise
2146 if self._upload_chunk(final=force) is not False:
2147 self.offset += self.buffer.seek(0, 2)
2148 self.buffer = io.BytesIO()
2150 def _upload_chunk(self, final=False):
2151 """Write one part of a multi-block file upload
2153 Parameters
2154 ==========
2155 final: bool
2156 This is the last block, so should complete file, if
2157 self.autocommit is True.
2158 """
2159 # may not yet have been initialized, may need to call _initialize_upload
2161 def _initiate_upload(self):
2162 """Create remote file/upload"""
2163 pass
2165 def _fetch_range(self, start, end):
2166 """Get the specified set of bytes from remote"""
2167 return self.fs.cat_file(self.path, start=start, end=end)
2169 def read(self, length=-1):
2170 """
2171 Return data from cache, or fetch pieces as necessary
2173 Parameters
2174 ----------
2175 length: int (-1)
2176 Number of bytes to read; if <0, all remaining bytes.
2177 """
2178 length = -1 if length is None else int(length)
2179 if self.mode != "rb":
2180 raise ValueError("File not in read mode")
2181 if length < 0:
2182 length = self.size - self.loc
2183 if self.closed:
2184 raise ValueError("I/O operation on closed file.")
2185 if length == 0:
2186 # don't even bother calling fetch
2187 return b""
2188 out = self.cache._fetch(self.loc, self.loc + length)
2190 logger.debug(
2191 "%s read: %i - %i %s",
2192 self,
2193 self.loc,
2194 self.loc + length,
2195 self.cache._log_stats(),
2196 )
2197 self.loc += len(out)
2198 return out
2200 def readinto(self, b):
2201 """mirrors builtin file's readinto method
2203 https://docs.python.org/3/library/io.html#io.RawIOBase.readinto
2204 """
2205 out = memoryview(b).cast("B")
2206 data = self.read(out.nbytes)
2207 out[: len(data)] = data
2208 return len(data)
2210 def readuntil(self, char=b"\n", blocks=None):
2211 """Return data between current position and first occurrence of char
2213 char is included in the output, except if the end of the tile is
2214 encountered first.
2216 Parameters
2217 ----------
2218 char: bytes
2219 Thing to find
2220 blocks: None or int
2221 How much to read in each go. Defaults to file blocksize - which may
2222 mean a new read on every call.
2223 """
2224 out = []
2225 while True:
2226 start = self.tell()
2227 part = self.read(blocks or self.blocksize)
2228 if len(part) == 0:
2229 break
2230 found = part.find(char)
2231 if found > -1:
2232 out.append(part[: found + len(char)])
2233 self.seek(start + found + len(char))
2234 break
2235 out.append(part)
2236 return b"".join(out)
2238 def readline(self):
2239 """Read until and including the first occurrence of newline character
2241 Note that, because of character encoding, this is not necessarily a
2242 true line ending.
2243 """
2244 return self.readuntil(b"\n")
2246 def __next__(self):
2247 out = self.readline()
2248 if out:
2249 return out
2250 raise StopIteration
2252 def __iter__(self):
2253 return self
2255 def readlines(self):
2256 """Return all data, split by the newline character, including the newline character"""
2257 data = self.read()
2258 lines = data.split(b"\n")
2259 out = [l + b"\n" for l in lines[:-1]]
2260 if data.endswith(b"\n"):
2261 return out
2262 else:
2263 return out + [lines[-1]]
2264 # return list(self) ???
2266 def readinto1(self, b):
2267 return self.readinto(b)
2269 def close(self):
2270 """Close file
2272 Finalizes writes, discards cache
2273 """
2274 if getattr(self, "_unclosable", False):
2275 return
2276 if self.closed:
2277 return
2278 try:
2279 if self.mode == "rb":
2280 cache = getattr(self, "cache", None)
2281 if cache is not None:
2282 close = getattr(cache, "close", None)
2283 if callable(close):
2284 close()
2285 self.cache = None
2286 else:
2287 if not getattr(self, "forced", True):
2288 self.flush(force=True)
2290 if self.fs is not None:
2291 self.fs.invalidate_cache(self.path)
2292 self.fs.invalidate_cache(self.fs._parent(self.path))
2293 finally:
2294 self.closed = True
2296 def readable(self):
2297 """Whether opened for reading"""
2298 return "r" in self.mode and not self.closed
2300 def seekable(self):
2301 """Whether is seekable (only in read mode)"""
2302 return self.readable()
2304 def writable(self):
2305 """Whether opened for writing"""
2306 return self.mode in {"wb", "ab", "xb"} and not self.closed
2308 def __reduce__(self):
2309 if self.mode != "rb":
2310 raise RuntimeError("Pickling a writeable file is not supported")
2312 return reopen, (
2313 self.fs,
2314 self.path,
2315 self.mode,
2316 self.blocksize,
2317 self.loc,
2318 self.size,
2319 self.autocommit,
2320 self.cache.name if self.cache else "none",
2321 self.kwargs,
2322 )
2324 def __del__(self):
2325 if not self.closed:
2326 self.close()
2328 def __str__(self):
2329 return f"<File-like object {type(self.fs).__name__}, {self.path}>"
2331 __repr__ = __str__
2333 def __enter__(self):
2334 return self
2336 def __exit__(self, *args):
2337 self.close()
2340def reopen(fs, path, mode, blocksize, loc, size, autocommit, cache_type, kwargs):
2341 file = fs.open(
2342 path,
2343 mode=mode,
2344 block_size=blocksize,
2345 autocommit=autocommit,
2346 cache_type=cache_type,
2347 size=size,
2348 **kwargs,
2349 )
2350 if loc > 0:
2351 file.seek(loc)
2352 return file