Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/fsspec/spec.py: 24%
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 # A negative start or end counts back from the file length. Only
860 # fsspec's own file classes expose `.size`; a plain file (e.g. from
861 # the caching filesystems) does not, so measure with a seek instead.
862 needs_size = (start is not None and start < 0) or (
863 end is not None and end < 0
864 )
865 size = None
866 seeked_for_size = False
867 if needs_size:
868 size = getattr(f, "size", None)
869 if size is None:
870 size = f.seek(0, 2)
871 seeked_for_size = True
872 if start is not None:
873 if start >= 0:
874 f.seek(start)
875 else:
876 f.seek(max(0, size + start))
877 elif seeked_for_size:
878 # measuring the size above moved the cursor to the end; with no
879 # start given, put it back so the read begins at the start. When
880 # .size was used instead, or nothing was measured, the cursor
881 # is untouched, as before, so a plain full read still works
882 # on non-seekable files.
883 f.seek(0)
884 if end is not None:
885 if end < 0:
886 end = size + end
887 return f.read(end - f.tell())
888 return f.read()
890 def pipe_file(self, path, value, mode="overwrite", **kwargs):
891 """Set the bytes of given file"""
892 if mode == "create" and self.exists(path):
893 # non-atomic but simple way; or could use "xb" in open(), which is likely
894 # not as well supported
895 raise FileExistsError
896 with self.open(path, "wb", **kwargs) as f:
897 f.write(value)
899 def pipe(self, path, value=None, **kwargs):
900 """Put value into path
902 (counterpart to ``cat``)
904 Parameters
905 ----------
906 path: string or dict(str, bytes)
907 If a string, a single remote location to put ``value`` bytes; if a dict,
908 a mapping of {path: bytesvalue}.
909 value: bytes, optional
910 If using a single path, these are the bytes to put there. Ignored if
911 ``path`` is a dict
912 """
913 if isinstance(path, str):
914 self.pipe_file(self._strip_protocol(path), value, **kwargs)
915 elif isinstance(path, dict):
916 for k, v in path.items():
917 self.pipe_file(self._strip_protocol(k), v, **kwargs)
918 else:
919 raise ValueError("path must be str or dict")
921 def cat_ranges(
922 self, paths, starts, ends, max_gap=None, on_error="return", **kwargs
923 ):
924 """Get the contents of byte ranges from one or more files
926 Parameters
927 ----------
928 paths: list
929 A list of of filepaths on this filesystems
930 starts, ends: int or list
931 Bytes limits of the read. If using a single int, the same value will be
932 used to read all the specified files.
933 """
934 if max_gap is not None:
935 raise NotImplementedError
936 if not isinstance(paths, list):
937 raise TypeError
938 if not isinstance(starts, list):
939 starts = [starts] * len(paths)
940 if not isinstance(ends, list):
941 ends = [ends] * len(paths)
942 if len(starts) != len(paths) or len(ends) != len(paths):
943 raise ValueError
944 out = []
945 for p, s, e in zip(paths, starts, ends):
946 try:
947 out.append(self.cat_file(p, s, e, **kwargs))
948 except Exception as e:
949 if on_error == "return":
950 out.append(e)
951 else:
952 raise
953 return out
955 def cat(self, path, recursive=False, on_error="raise", **kwargs):
956 """Fetch (potentially multiple) paths' contents
958 Parameters
959 ----------
960 recursive: bool
961 If True, assume the path(s) are directories, and get all the
962 contained files
963 on_error : "raise", "omit", "return"
964 If raise, an underlying exception will be raised (converted to KeyError
965 if the type is in self.missing_exceptions); if omit, keys with exception
966 will simply not be included in the output; if "return", all keys are
967 included in the output, but the value will be bytes or an exception
968 instance.
969 kwargs: passed to cat_file
971 Returns
972 -------
973 dict of {path: contents} if there are multiple paths
974 or the path has been otherwise expanded
975 """
976 paths = self.expand_path(path, recursive=recursive, **kwargs)
977 if (
978 len(paths) > 1
979 or isinstance(path, list)
980 or paths[0] != self._strip_protocol(path)
981 ):
982 out = {}
983 for path in paths:
984 try:
985 out[path] = self.cat_file(path, **kwargs)
986 except Exception as e:
987 if on_error == "raise":
988 raise
989 if on_error == "return":
990 out[path] = e
991 return out
992 else:
993 return self.cat_file(paths[0], **kwargs)
995 def get_file(
996 self, rpath, lpath=None, callback=DEFAULT_CALLBACK, outfile=None, **kwargs
997 ):
998 """Copy single remote file to local"""
999 from .implementations.local import LocalFileSystem
1001 if outfile is None and isfilelike(lpath):
1002 outfile = lpath
1003 elif outfile is None and self.isdir(rpath):
1004 os.makedirs(lpath, exist_ok=True)
1005 return None
1007 if outfile is None:
1008 fs = LocalFileSystem(auto_mkdir=True)
1009 fs.makedirs(fs._parent(lpath), exist_ok=True)
1011 with self.open(rpath, "rb", **kwargs) as f1:
1012 close_outfile = outfile is None
1013 if close_outfile:
1014 outfile = open(lpath, "wb")
1016 try:
1017 callback.set_size(getattr(f1, "size", None))
1018 data = True
1019 while data:
1020 data = f1.read(self.blocksize)
1021 segment_len = outfile.write(data)
1022 if segment_len is None:
1023 segment_len = len(data)
1024 callback.relative_update(segment_len)
1025 finally:
1026 if close_outfile:
1027 outfile.close()
1029 def get(
1030 self,
1031 rpath,
1032 lpath,
1033 recursive=False,
1034 callback=DEFAULT_CALLBACK,
1035 maxdepth=None,
1036 **kwargs,
1037 ):
1038 """Copy file(s) to local.
1040 Copies a specific file or tree of files (if recursive=True). If lpath
1041 ends with a "/", it will be assumed to be a directory, and target files
1042 will go within. Can submit a list of paths, which may be glob-patterns
1043 and will be expanded.
1045 Calls get_file for each source.
1046 """
1047 if isinstance(lpath, list) and isinstance(rpath, list):
1048 # No need to expand paths when both source and destination
1049 # are provided as lists
1050 rpaths = rpath
1051 lpaths = lpath
1052 else:
1053 from .implementations.local import (
1054 LocalFileSystem,
1055 make_path_posix,
1056 trailing_sep,
1057 )
1059 source_is_str = isinstance(rpath, str)
1060 rpaths = self.expand_path(
1061 rpath, recursive=recursive, maxdepth=maxdepth, **kwargs
1062 )
1063 if source_is_str and (not recursive or maxdepth is not None):
1064 # Non-recursive glob does not copy directories
1065 rpaths = [p for p in rpaths if not (trailing_sep(p) or self.isdir(p))]
1066 if not rpaths:
1067 return
1069 if isinstance(lpath, str):
1070 lpath = make_path_posix(lpath)
1072 source_is_file = len(rpaths) == 1
1073 dest_is_dir = isinstance(lpath, str) and (
1074 trailing_sep(lpath) or LocalFileSystem().isdir(lpath)
1075 )
1077 exists = source_is_str and (
1078 (has_magic(rpath) and source_is_file)
1079 or (not has_magic(rpath) and dest_is_dir and not trailing_sep(rpath))
1080 )
1081 lpaths = other_paths(
1082 rpaths,
1083 lpath,
1084 exists=exists,
1085 flatten=not source_is_str,
1086 )
1087 if isinstance(lpath, str):
1088 # The names came from the source listing; ".." in one of them
1089 # would otherwise place the copy above the destination. When
1090 # lpath is a list the caller named every destination itself.
1091 check_contained(lpath, lpaths)
1093 callback.set_size(len(lpaths))
1094 for lpath, rpath in callback.wrap(zip(lpaths, rpaths)):
1095 with callback.branched(rpath, lpath) as child:
1096 self.get_file(rpath, lpath, callback=child, **kwargs)
1098 def put_file(
1099 self, lpath, rpath, callback=DEFAULT_CALLBACK, mode="overwrite", **kwargs
1100 ):
1101 """Copy single file to remote"""
1102 if mode == "create" and self.exists(rpath):
1103 raise FileExistsError
1104 if os.path.isdir(lpath):
1105 self.makedirs(rpath, exist_ok=True)
1106 return None
1108 with open(lpath, "rb") as f1:
1109 size = f1.seek(0, 2)
1110 callback.set_size(size)
1111 f1.seek(0)
1113 self.mkdirs(self._parent(os.fspath(rpath)), exist_ok=True)
1114 with self.open(rpath, "wb", **kwargs) as f2:
1115 while f1.tell() < size:
1116 data = f1.read(self.blocksize)
1117 segment_len = f2.write(data)
1118 if segment_len is None:
1119 segment_len = len(data)
1120 callback.relative_update(segment_len)
1122 def put(
1123 self,
1124 lpath,
1125 rpath,
1126 recursive=False,
1127 callback=DEFAULT_CALLBACK,
1128 maxdepth=None,
1129 **kwargs,
1130 ):
1131 """Copy file(s) from local.
1133 Copies a specific file or tree of files (if recursive=True). If rpath
1134 ends with a "/", it will be assumed to be a directory, and target files
1135 will go within.
1137 Calls put_file for each source.
1138 """
1139 if isinstance(lpath, list) and isinstance(rpath, list):
1140 # No need to expand paths when both source and destination
1141 # are provided as lists
1142 rpaths = rpath
1143 lpaths = lpath
1144 else:
1145 from .implementations.local import (
1146 LocalFileSystem,
1147 make_path_posix,
1148 trailing_sep,
1149 )
1151 source_is_str = isinstance(lpath, str)
1152 if source_is_str:
1153 lpath = make_path_posix(lpath)
1154 fs = LocalFileSystem()
1155 lpaths = fs.expand_path(
1156 lpath, recursive=recursive, maxdepth=maxdepth, **kwargs
1157 )
1158 if source_is_str and (not recursive or maxdepth is not None):
1159 # Non-recursive glob does not copy directories
1160 lpaths = [p for p in lpaths if not (trailing_sep(p) or fs.isdir(p))]
1161 if not lpaths:
1162 return
1164 source_is_file = len(lpaths) == 1
1165 dest_is_dir = isinstance(rpath, str) and (
1166 trailing_sep(rpath) or self.isdir(rpath)
1167 )
1169 rpath = (
1170 self._strip_protocol(rpath)
1171 if isinstance(rpath, str)
1172 else [self._strip_protocol(p) for p in rpath]
1173 )
1174 exists = source_is_str and (
1175 (has_magic(lpath) and source_is_file)
1176 or (not has_magic(lpath) and dest_is_dir and not trailing_sep(lpath))
1177 )
1178 rpaths = other_paths(
1179 lpaths,
1180 rpath,
1181 exists=exists,
1182 flatten=not source_is_str,
1183 )
1185 callback.set_size(len(rpaths))
1186 for lpath, rpath in callback.wrap(zip(lpaths, rpaths)):
1187 with callback.branched(lpath, rpath) as child:
1188 self.put_file(lpath, rpath, callback=child, **kwargs)
1190 def head(self, path, size=1024):
1191 """Get the first ``size`` bytes from file"""
1192 with self.open(path, "rb") as f:
1193 return f.read(size)
1195 def tail(self, path, size=1024):
1196 """Get the last ``size`` bytes from file"""
1197 with self.open(path, "rb") as f:
1198 f.seek(0, 2)
1199 f.seek(max(f.tell() - size, 0))
1200 return f.read()
1202 def cp_file(self, path1, path2, **kwargs):
1203 raise NotImplementedError
1205 def copy(
1206 self, path1, path2, recursive=False, maxdepth=None, on_error=None, **kwargs
1207 ):
1208 """Copy within two locations in the filesystem
1210 on_error : "raise", "ignore"
1211 If raise, any not-found exceptions will be raised; if ignore any
1212 not-found exceptions will cause the path to be skipped; defaults to
1213 raise unless recursive is true, where the default is ignore
1214 """
1215 if on_error is None and recursive:
1216 on_error = "ignore"
1217 elif on_error is None:
1218 on_error = "raise"
1220 if isinstance(path1, list) and isinstance(path2, list):
1221 # No need to expand paths when both source and destination
1222 # are provided as lists
1223 paths1 = path1
1224 paths2 = path2
1225 else:
1226 from .implementations.local import trailing_sep
1228 source_is_str = isinstance(path1, str)
1229 paths1 = self.expand_path(
1230 path1, recursive=recursive, maxdepth=maxdepth, **kwargs
1231 )
1232 if source_is_str and (not recursive or maxdepth is not None):
1233 # Non-recursive glob does not copy directories
1234 paths1 = [p for p in paths1 if not (trailing_sep(p) or self.isdir(p))]
1235 if not paths1:
1236 return
1238 source_is_file = len(paths1) == 1
1239 dest_is_dir = isinstance(path2, str) and (
1240 trailing_sep(path2) or self.isdir(path2)
1241 )
1243 exists = source_is_str and (
1244 (has_magic(path1) and source_is_file)
1245 or (not has_magic(path1) and dest_is_dir and not trailing_sep(path1))
1246 )
1247 paths2 = other_paths(
1248 paths1,
1249 path2,
1250 exists=exists,
1251 flatten=not source_is_str,
1252 )
1254 for p1, p2 in zip(paths1, paths2):
1255 try:
1256 self.cp_file(p1, p2, **kwargs)
1257 except FileNotFoundError:
1258 if on_error == "raise":
1259 raise
1261 def expand_path(
1262 self, path, recursive=False, maxdepth=None, assume_literal=False, **kwargs
1263 ):
1264 """Turn one or more globs or directories into a list of all matching paths
1265 to files or directories.
1267 kwargs are passed to ``glob`` or ``find``, which may in turn call ``ls``
1268 """
1270 if maxdepth is not None and maxdepth < 1:
1271 raise ValueError("maxdepth must be at least 1")
1273 if isinstance(path, (str, os.PathLike)):
1274 out = self.expand_path([path], recursive, maxdepth, **kwargs)
1275 else:
1276 out = set()
1277 path = [self._strip_protocol(p) for p in path]
1278 for p in path:
1279 if not assume_literal and has_magic(p):
1280 bit = set(self.glob(p, maxdepth=maxdepth, **kwargs))
1281 out |= bit
1282 if recursive:
1283 # glob call above expanded one depth so if maxdepth is defined
1284 # then decrement it in expand_path call below. If it is zero
1285 # after decrementing then avoid expand_path call.
1286 if maxdepth is not None and maxdepth <= 1:
1287 continue
1288 out |= set(
1289 self.expand_path(
1290 list(bit),
1291 recursive=recursive,
1292 maxdepth=maxdepth - 1 if maxdepth is not None else None,
1293 assume_literal=True,
1294 **kwargs,
1295 )
1296 )
1297 continue
1298 elif recursive:
1299 rec = set(
1300 self.find(
1301 p, maxdepth=maxdepth, withdirs=True, detail=False, **kwargs
1302 )
1303 )
1304 out |= rec
1305 if p not in out and (recursive is False or self.exists(p)):
1306 # should only check once, for the root
1307 out.add(p)
1308 if not out:
1309 raise FileNotFoundError(path)
1310 return sorted(out)
1312 def mv(self, path1, path2, recursive=False, maxdepth=None, **kwargs):
1313 """Move file(s) from one location to another"""
1314 if path1 == path2:
1315 logger.debug("%s mv: The paths are the same, so no files were moved.", self)
1316 else:
1317 # explicitly raise exception to prevent data corruption
1318 self.copy(
1319 path1, path2, recursive=recursive, maxdepth=maxdepth, on_error="raise"
1320 )
1321 self.rm(path1, recursive=recursive)
1323 def rm_file(self, path):
1324 """Delete a file"""
1325 self._rm(path)
1327 def _rm(self, path):
1328 """Delete one file"""
1329 # this is the old name for the method, prefer rm_file
1330 raise NotImplementedError
1332 def rm(self, path, recursive=False, maxdepth=None):
1333 """Delete files or directories.
1335 Parameters
1336 ----------
1337 path: str or list of str
1338 Files or directories to delete.
1339 recursive: bool
1340 If True, recursively delete directories and their contents.
1341 maxdepth: int or None
1342 Depth to pass to walk for finding files to delete, if recursive.
1343 If None, there will be no limit and infinite recursion may be
1344 possible.
1345 """
1346 path = self.expand_path(path, recursive=recursive, maxdepth=maxdepth)
1347 for p in reversed(path):
1348 self.rm_file(p)
1350 @classmethod
1351 def _parent(cls, path):
1352 path = cls._strip_protocol(path)
1353 if "/" in path:
1354 parent = path.rsplit("/", 1)[0].lstrip(cls.root_marker)
1355 return cls.root_marker + parent
1356 else:
1357 return cls.root_marker
1359 def _open(
1360 self,
1361 path,
1362 mode="rb",
1363 block_size=None,
1364 autocommit=True,
1365 cache_options=None,
1366 **kwargs,
1367 ):
1368 """Return raw bytes-mode file-like from the file-system"""
1369 return AbstractBufferedFile(
1370 self,
1371 path,
1372 mode,
1373 block_size,
1374 autocommit,
1375 cache_options=cache_options,
1376 **kwargs,
1377 )
1379 def open(
1380 self,
1381 path,
1382 mode="rb",
1383 block_size=None,
1384 cache_options=None,
1385 compression=None,
1386 **kwargs,
1387 ):
1388 """
1389 Return a file-like object from the filesystem
1391 The resultant instance must function correctly in a context ``with``
1392 block.
1394 Parameters
1395 ----------
1396 path: str
1397 Target file
1398 mode: str like 'rb', 'w'
1399 See builtin ``open()``
1400 Mode "x" (exclusive write) may be implemented by the backend. Even if
1401 it is, whether it is checked up front or on commit, and whether it is
1402 atomic is implementation-dependent.
1403 block_size: int
1404 Some indication of buffering - this is a value in bytes
1405 cache_options : dict, optional
1406 Extra arguments to pass through to the cache.
1407 compression: string or None
1408 If given, open file using compression codec. Can either be a compression
1409 name (a key in ``fsspec.compression.compr``) or "infer" to guess the
1410 compression from the filename suffix.
1411 encoding, errors, newline: passed on to TextIOWrapper for text mode
1412 """
1413 import io
1415 path = self._strip_protocol(path)
1416 if "b" not in mode:
1417 mode = mode.replace("t", "") + "b"
1419 text_kwargs = {
1420 k: kwargs.pop(k)
1421 for k in ["encoding", "errors", "newline"]
1422 if k in kwargs
1423 }
1424 return io.TextIOWrapper(
1425 self.open(
1426 path,
1427 mode,
1428 block_size=block_size,
1429 cache_options=cache_options,
1430 compression=compression,
1431 **kwargs,
1432 ),
1433 **text_kwargs,
1434 )
1435 else:
1436 ac = kwargs.pop("autocommit", not self._intrans)
1437 f = self._open(
1438 path,
1439 mode=mode,
1440 block_size=block_size,
1441 autocommit=ac,
1442 cache_options=cache_options,
1443 **kwargs,
1444 )
1445 if not ac and "r" not in mode:
1446 self.transaction.files.append(f)
1447 if compression is not None:
1448 from fsspec.compression import compr
1449 from fsspec.core import get_compression
1451 compression = get_compression(path, compression)
1452 compress = compr[compression]
1453 f = compress(f, mode=mode[0])
1454 return f
1456 def touch(self, path, truncate=True, **kwargs):
1457 """Create empty file, or update timestamp
1459 Parameters
1460 ----------
1461 path: str
1462 file location
1463 truncate: bool
1464 If True, always set file size to 0; if False, update timestamp and
1465 leave file unchanged, if backend allows this
1466 """
1467 if truncate or not self.exists(path):
1468 with self.open(path, "wb", **kwargs):
1469 pass
1470 else:
1471 raise NotImplementedError # update timestamp, if possible
1473 def ukey(self, path):
1474 """Hash of file properties, to tell if it has changed"""
1475 return sha256(str(self.info(path)).encode()).hexdigest()
1477 def read_block(self, fn, offset, length, delimiter=None):
1478 """Read a block of bytes from
1480 Starting at ``offset`` of the file, read ``length`` bytes. If
1481 ``delimiter`` is set then we ensure that the read starts and stops at
1482 delimiter boundaries that follow the locations ``offset`` and ``offset
1483 + length``. If ``offset`` is zero then we start at zero. The
1484 bytestring returned WILL include the end delimiter string.
1486 If offset+length is beyond the eof, reads to eof.
1488 Parameters
1489 ----------
1490 fn: string
1491 Path to filename
1492 offset: int
1493 Byte offset to start read
1494 length: int
1495 Number of bytes to read. If None, read to end.
1496 delimiter: bytes (optional)
1497 Ensure reading starts and stops at delimiter bytestring
1499 Examples
1500 --------
1501 >>> fs.read_block('data/file.csv', 0, 13) # doctest: +SKIP
1502 b'Alice, 100\\nBo'
1503 >>> fs.read_block('data/file.csv', 0, 13, delimiter=b'\\n') # doctest: +SKIP
1504 b'Alice, 100\\nBob, 200\\n'
1506 Use ``length=None`` to read to the end of the file.
1507 >>> fs.read_block('data/file.csv', 0, None, delimiter=b'\\n') # doctest: +SKIP
1508 b'Alice, 100\\nBob, 200\\nCharlie, 300'
1510 See Also
1511 --------
1512 :func:`fsspec.utils.read_block`
1513 """
1514 with self.open(fn, "rb") as f:
1515 # Only fsspec's own file classes expose `.size`; a plain file (e.g.
1516 # from the caching filesystems) does not, so fall back to a seek.
1517 size = getattr(f, "size", None)
1518 if size is None:
1519 size = f.seek(0, 2)
1520 f.seek(0)
1521 if length is None:
1522 length = size
1523 if size is not None and offset + length > size:
1524 length = size - offset
1525 return read_block(f, offset, length, delimiter)
1527 def to_json(self, *, include_password: bool = True) -> str:
1528 """
1529 JSON representation of this filesystem instance.
1531 Parameters
1532 ----------
1533 include_password: bool, default True
1534 Whether to include the password (if any) in the output.
1536 Returns
1537 -------
1538 JSON string with keys ``cls`` (the python location of this class),
1539 protocol (text name of this class's protocol, first one in case of
1540 multiple), ``args`` (positional args, usually empty), and all other
1541 keyword arguments as their own keys.
1543 Warnings
1544 --------
1545 Serialized filesystems may contain sensitive information which have been
1546 passed to the constructor, such as passwords and tokens. Make sure you
1547 store and send them in a secure environment!
1548 """
1549 from .json import FilesystemJSONEncoder
1551 return json.dumps(
1552 self,
1553 cls=type(
1554 "_FilesystemJSONEncoder",
1555 (FilesystemJSONEncoder,),
1556 {"include_password": include_password},
1557 ),
1558 )
1560 @staticmethod
1561 def from_json(blob: str) -> AbstractFileSystem:
1562 """
1563 Recreate a filesystem instance from JSON representation.
1565 See ``.to_json()`` for the expected structure of the input.
1567 Parameters
1568 ----------
1569 blob: str
1571 Returns
1572 -------
1573 file system instance, not necessarily of this particular class.
1575 Warnings
1576 --------
1577 This can import arbitrary modules (as determined by the ``cls`` key).
1578 Make sure you haven't installed any modules that may execute malicious code
1579 at import time.
1580 """
1581 from .json import FilesystemJSONDecoder
1583 return json.loads(blob, cls=FilesystemJSONDecoder)
1585 def to_dict(self, *, include_password: bool = True) -> dict[str, Any]:
1586 """
1587 JSON-serializable dictionary representation of this filesystem instance.
1589 Parameters
1590 ----------
1591 include_password: bool, default True
1592 Whether to include the password (if any) in the output.
1594 Returns
1595 -------
1596 Dictionary with keys ``cls`` (the python location of this class),
1597 protocol (text name of this class's protocol, first one in case of
1598 multiple), ``args`` (positional args, usually empty), and all other
1599 keyword arguments as their own keys.
1601 Warnings
1602 --------
1603 Serialized filesystems may contain sensitive information which have been
1604 passed to the constructor, such as passwords and tokens. Make sure you
1605 store and send them in a secure environment!
1606 """
1607 from .json import FilesystemJSONEncoder
1609 json_encoder = FilesystemJSONEncoder()
1611 cls = type(self)
1612 proto = self.protocol
1614 storage_options = dict(self.storage_options)
1615 if not include_password:
1616 storage_options.pop("password", None)
1618 return dict(
1619 cls=f"{cls.__module__}:{cls.__name__}",
1620 protocol=proto[0] if isinstance(proto, (tuple, list)) else proto,
1621 args=json_encoder.make_serializable(self.storage_args),
1622 **json_encoder.make_serializable(storage_options),
1623 )
1625 @staticmethod
1626 def from_dict(dct: dict[str, Any]) -> AbstractFileSystem:
1627 """
1628 Recreate a filesystem instance from dictionary representation.
1630 See ``.to_dict()`` for the expected structure of the input.
1632 Parameters
1633 ----------
1634 dct: Dict[str, Any]
1636 Returns
1637 -------
1638 file system instance, not necessarily of this particular class.
1640 Warnings
1641 --------
1642 This can import arbitrary modules (as determined by the ``cls`` key).
1643 Make sure you haven't installed any modules that may execute malicious code
1644 at import time.
1645 """
1646 from .json import FilesystemJSONDecoder
1648 json_decoder = FilesystemJSONDecoder()
1650 dct = dict(dct) # Defensive copy
1652 cls = FilesystemJSONDecoder.try_resolve_fs_cls(dct)
1653 if cls is None:
1654 raise ValueError("Not a serialized AbstractFileSystem")
1656 dct.pop("cls", None)
1657 dct.pop("protocol", None)
1659 return cls(
1660 *json_decoder.unmake_serializable(dct.pop("args", ())),
1661 **json_decoder.unmake_serializable(dct),
1662 )
1664 def _get_pyarrow_filesystem(self):
1665 """
1666 Make a version of the FS instance which will be acceptable to pyarrow
1667 """
1668 # all instances already also derive from pyarrow
1669 return self
1671 def get_mapper(self, root="", check=False, create=False, missing_exceptions=None):
1672 """Create key/value store based on this file-system
1674 Makes a MutableMapping interface to the FS at the given root path.
1675 See ``fsspec.mapping.FSMap`` for further details.
1676 """
1677 from .mapping import FSMap
1679 return FSMap(
1680 root,
1681 self,
1682 check=check,
1683 create=create,
1684 missing_exceptions=missing_exceptions,
1685 )
1687 @classmethod
1688 def clear_instance_cache(cls):
1689 """
1690 Clear the cache of filesystem instances.
1692 Notes
1693 -----
1694 Unless overridden by setting the ``cachable`` class attribute to False,
1695 the filesystem class stores a reference to newly created instances. This
1696 prevents Python's normal rules around garbage collection from working,
1697 since the instances refcount will not drop to zero until
1698 ``clear_instance_cache`` is called.
1699 """
1700 cls._cache.clear()
1702 def created(self, path):
1703 """Return the created timestamp of a file as a datetime.datetime"""
1704 raise NotImplementedError
1706 def modified(self, path):
1707 """Return the modified timestamp of a file as a datetime.datetime"""
1708 raise NotImplementedError
1710 def tree(
1711 self,
1712 path: str = "/",
1713 recursion_limit: int = 2,
1714 max_display: int = 25,
1715 display_size: bool = False,
1716 prefix: str = "",
1717 is_last: bool = True,
1718 first: bool = True,
1719 indent_size: int = 4,
1720 ) -> str:
1721 """
1722 Return a tree-like structure of the filesystem starting from the given path as a string.
1724 Parameters
1725 ----------
1726 path: Root path to start traversal from
1727 recursion_limit: Maximum depth of directory traversal
1728 max_display: Maximum number of items to display per directory
1729 display_size: Whether to display file sizes
1730 prefix: Current line prefix for visual tree structure
1731 is_last: Whether current item is last in its level
1732 first: Whether this is the first call (displays root path)
1733 indent_size: Number of spaces by indent
1735 Returns
1736 -------
1737 str: A string representing the tree structure.
1739 Example
1740 -------
1741 >>> from fsspec import filesystem
1743 >>> fs = filesystem('ftp', host='test.rebex.net', user='demo', password='password')
1744 >>> tree = fs.tree(display_size=True, recursion_limit=3, indent_size=8, max_display=10)
1745 >>> print(tree)
1746 """
1748 def format_bytes(n: int) -> str:
1749 """Format bytes as text."""
1750 for prefix, k in (
1751 ("P", 2**50),
1752 ("T", 2**40),
1753 ("G", 2**30),
1754 ("M", 2**20),
1755 ("k", 2**10),
1756 ):
1757 if n >= 0.9 * k:
1758 return f"{n / k:.2f} {prefix}b"
1759 return f"{n}B"
1761 result = []
1763 if first:
1764 result.append(path)
1766 if recursion_limit:
1767 indent = " " * indent_size
1768 contents = self.ls(path, detail=True)
1769 contents.sort(
1770 key=lambda x: (x.get("type") != "directory", x.get("name", ""))
1771 )
1773 if max_display is not None and len(contents) > max_display:
1774 displayed_contents = contents[:max_display]
1775 remaining_count = len(contents) - max_display
1776 else:
1777 displayed_contents = contents
1778 remaining_count = 0
1780 for i, item in enumerate(displayed_contents):
1781 is_last_item = (i == len(displayed_contents) - 1) and (
1782 remaining_count == 0
1783 )
1785 branch = (
1786 "└" + ("─" * (indent_size - 2))
1787 if is_last_item
1788 else "├" + ("─" * (indent_size - 2))
1789 )
1790 branch += " "
1791 new_prefix = prefix + (
1792 indent if is_last_item else "│" + " " * (indent_size - 1)
1793 )
1795 name = os.path.basename(item.get("name", ""))
1797 if display_size and item.get("type") == "directory":
1798 sub_contents = self.ls(item.get("name", ""), detail=True)
1799 num_files = sum(
1800 1 for sub_item in sub_contents if sub_item.get("type") == "file"
1801 )
1802 num_folders = sum(
1803 1
1804 for sub_item in sub_contents
1805 if sub_item.get("type") == "directory"
1806 )
1808 if num_files == 0 and num_folders == 0:
1809 size = " (empty folder)"
1810 elif num_files == 0:
1811 size = f" ({num_folders} subfolder{'s' if num_folders > 1 else ''})"
1812 elif num_folders == 0:
1813 size = f" ({num_files} file{'s' if num_files > 1 else ''})"
1814 else:
1815 size = f" ({num_files} file{'s' if num_files > 1 else ''}, {num_folders} subfolder{'s' if num_folders > 1 else ''})"
1816 elif display_size and item.get("type") == "file":
1817 size = f" ({format_bytes(item.get('size', 0))})"
1818 else:
1819 size = ""
1821 result.append(f"{prefix}{branch}{name}{size}")
1823 if item.get("type") == "directory" and recursion_limit > 0:
1824 result.append(
1825 self.tree(
1826 path=item.get("name", ""),
1827 recursion_limit=recursion_limit - 1,
1828 max_display=max_display,
1829 display_size=display_size,
1830 prefix=new_prefix,
1831 is_last=is_last_item,
1832 first=False,
1833 indent_size=indent_size,
1834 )
1835 )
1837 if remaining_count > 0:
1838 more_message = f"{remaining_count} more item(s) not displayed."
1839 result.append(
1840 f"{prefix}{'└' + ('─' * (indent_size - 2))} {more_message}"
1841 )
1843 return "\n".join(_ for _ in result if _)
1845 # ------------------------------------------------------------------------
1846 # Aliases
1848 def read_bytes(self, path, start=None, end=None, **kwargs):
1849 """Alias of `AbstractFileSystem.cat_file`."""
1850 return self.cat_file(path, start=start, end=end, **kwargs)
1852 def write_bytes(self, path, value, **kwargs):
1853 """Alias of `AbstractFileSystem.pipe_file`."""
1854 self.pipe_file(path, value, **kwargs)
1856 def makedir(self, path, create_parents=True, **kwargs):
1857 """Alias of `AbstractFileSystem.mkdir`."""
1858 return self.mkdir(path, create_parents=create_parents, **kwargs)
1860 def mkdirs(self, path, exist_ok=False):
1861 """Alias of `AbstractFileSystem.makedirs`."""
1862 return self.makedirs(path, exist_ok=exist_ok)
1864 def listdir(self, path, detail=True, **kwargs):
1865 """Alias of `AbstractFileSystem.ls`."""
1866 return self.ls(path, detail=detail, **kwargs)
1868 def cp(self, path1, path2, **kwargs):
1869 """Alias of `AbstractFileSystem.copy`."""
1870 return self.copy(path1, path2, **kwargs)
1872 def move(self, path1, path2, **kwargs):
1873 """Alias of `AbstractFileSystem.mv`."""
1874 return self.mv(path1, path2, **kwargs)
1876 def stat(self, path, **kwargs):
1877 """Alias of `AbstractFileSystem.info`."""
1878 return self.info(path, **kwargs)
1880 def disk_usage(self, path, total=True, maxdepth=None, **kwargs):
1881 """Alias of `AbstractFileSystem.du`."""
1882 return self.du(path, total=total, maxdepth=maxdepth, **kwargs)
1884 def rename(self, path1, path2, **kwargs):
1885 """Alias of `AbstractFileSystem.mv`."""
1886 return self.mv(path1, path2, **kwargs)
1888 def delete(self, path, recursive=False, maxdepth=None):
1889 """Alias of `AbstractFileSystem.rm`."""
1890 return self.rm(path, recursive=recursive, maxdepth=maxdepth)
1892 def upload(self, lpath, rpath, recursive=False, **kwargs):
1893 """Alias of `AbstractFileSystem.put`."""
1894 return self.put(lpath, rpath, recursive=recursive, **kwargs)
1896 def download(self, rpath, lpath, recursive=False, **kwargs):
1897 """Alias of `AbstractFileSystem.get`."""
1898 return self.get(rpath, lpath, recursive=recursive, **kwargs)
1900 def sign(self, path, expiration=100, **kwargs):
1901 """Create a signed URL representing the given path
1903 Some implementations allow temporary URLs to be generated, as a
1904 way of delegating credentials.
1906 Parameters
1907 ----------
1908 path : str
1909 The path on the filesystem
1910 expiration : int
1911 Number of seconds to enable the URL for (if supported)
1913 Returns
1914 -------
1915 URL : str
1916 The signed URL
1918 Raises
1919 ------
1920 NotImplementedError : if method is not implemented for a filesystem
1921 """
1922 raise NotImplementedError("Sign is not implemented for this filesystem")
1924 def _isfilestore(self):
1925 # Originally inherited from pyarrow DaskFileSystem. Keeping this
1926 # here for backwards compatibility as long as pyarrow uses its
1927 # legacy fsspec-compatible filesystems and thus accepts fsspec
1928 # filesystems as well
1929 return False
1932class AbstractBufferedFile(io.IOBase):
1933 """Convenient class to derive from to provide buffering
1935 In the case that the backend does not provide a pythonic file-like object
1936 already, this class contains much of the logic to build one. The only
1937 methods that need to be overridden are ``_upload_chunk``,
1938 ``_initiate_upload`` and ``_fetch_range``.
1939 """
1941 DEFAULT_BLOCK_SIZE = 5 * 2**20
1942 _details = None
1944 def __init__(
1945 self,
1946 fs,
1947 path,
1948 mode="rb",
1949 block_size="default",
1950 autocommit=True,
1951 cache_type="readahead",
1952 cache_options=None,
1953 size=None,
1954 **kwargs,
1955 ):
1956 """
1957 Template for files with buffered reading and writing
1959 Parameters
1960 ----------
1961 fs: instance of FileSystem
1962 path: str
1963 location in file-system
1964 mode: str
1965 Normal file modes. Currently only 'wb', 'ab' or 'rb'. Some file
1966 systems may be read-only, and some may not support append.
1967 block_size: int
1968 Buffer size for reading or writing, 'default' for class default
1969 autocommit: bool
1970 Whether to write to final destination; may only impact what
1971 happens when file is being closed.
1972 cache_type: {"readahead", "none", "mmap", "bytes"}, default "readahead"
1973 Caching policy in read mode. See the definitions in ``core``.
1974 cache_options : dict
1975 Additional options passed to the constructor for the cache specified
1976 by `cache_type`.
1977 size: int
1978 If given and in read mode, suppressed having to look up the file size
1979 kwargs:
1980 Gets stored as self.kwargs
1981 """
1982 from .core import caches
1984 self.path = path
1985 self.fs = fs
1986 self.mode = mode
1987 self.blocksize = (
1988 self.DEFAULT_BLOCK_SIZE if block_size in ["default", None] else block_size
1989 )
1990 self.loc = 0
1991 self.autocommit = autocommit
1992 self.end = None
1993 self.start = None
1994 self.closed = False
1996 if cache_options is None:
1997 cache_options = {}
1999 if "trim" in kwargs:
2000 warnings.warn(
2001 "Passing 'trim' to control the cache behavior has been deprecated. "
2002 "Specify it within the 'cache_options' argument instead.",
2003 FutureWarning,
2004 )
2005 cache_options["trim"] = kwargs.pop("trim")
2007 self.kwargs = kwargs
2009 if mode not in {"ab", "rb", "wb", "xb"}:
2010 raise NotImplementedError("File mode not supported")
2011 if mode == "rb":
2012 if size is not None:
2013 self.size = size
2014 else:
2015 self.size = self.details["size"]
2016 self.cache = caches[cache_type](
2017 self.blocksize, self._fetch_range, self.size, **cache_options
2018 )
2019 else:
2020 self.buffer = io.BytesIO()
2021 self.offset = None
2022 self.forced = False
2023 self.location = None
2025 @property
2026 def details(self):
2027 if self._details is None:
2028 self._details = self.fs.info(self.path)
2029 return self._details
2031 @details.setter
2032 def details(self, value):
2033 self._details = value
2034 self.size = value["size"]
2036 @property
2037 def full_name(self):
2038 return _unstrip_protocol(self.path, self.fs)
2040 @property
2041 def closed(self):
2042 # get around this attr being read-only in IOBase
2043 # use getattr here, since this can be called during del
2044 return getattr(self, "_closed", True)
2046 @closed.setter
2047 def closed(self, c):
2048 self._closed = c
2050 def __hash__(self):
2051 if "w" in self.mode:
2052 return id(self)
2053 else:
2054 return int(tokenize(self.details), 16)
2056 def __eq__(self, other):
2057 """Files are equal if they have the same checksum, only in read mode"""
2058 if self is other:
2059 return True
2060 return (
2061 isinstance(other, type(self))
2062 and self.mode == "rb"
2063 and other.mode == "rb"
2064 and hash(self) == hash(other)
2065 )
2067 def commit(self):
2068 """Move from temp to final destination"""
2070 def discard(self):
2071 """Throw away temporary file"""
2073 def info(self):
2074 """File information about this path"""
2075 if self.readable():
2076 return self.details
2077 else:
2078 raise ValueError("Info not available while writing")
2080 def tell(self):
2081 """Current file location"""
2082 return self.loc
2084 def seek(self, loc, whence=0):
2085 """Set current file location
2087 Parameters
2088 ----------
2089 loc: int
2090 byte location
2091 whence: {0, 1, 2}
2092 from start of file, current location or end of file, resp.
2093 """
2094 loc = int(loc)
2095 if not self.mode == "rb":
2096 raise OSError(ESPIPE, "Seek only available in read mode")
2097 if whence == 0:
2098 nloc = loc
2099 elif whence == 1:
2100 nloc = self.loc + loc
2101 elif whence == 2:
2102 nloc = self.size + loc
2103 else:
2104 raise ValueError(f"invalid whence ({whence}, should be 0, 1 or 2)")
2105 if nloc < 0:
2106 raise ValueError("Seek before start of file")
2107 self.loc = nloc
2108 return self.loc
2110 def write(self, data):
2111 """
2112 Write data to buffer.
2114 Buffer only sent on flush() or if buffer is greater than
2115 or equal to blocksize.
2117 Parameters
2118 ----------
2119 data: bytes
2120 Set of bytes to be written.
2121 """
2122 if not self.writable():
2123 raise ValueError("File not in write mode")
2124 if self.closed:
2125 raise ValueError("I/O operation on closed file.")
2126 if self.forced:
2127 raise ValueError("This file has been force-flushed, can only close")
2128 out = self.buffer.write(data)
2129 self.loc += out
2130 if self.buffer.tell() >= self.blocksize:
2131 self.flush()
2132 return out
2134 def flush(self, force=False):
2135 """
2136 Write buffered data to backend store.
2138 Writes the current buffer, if it is larger than the block-size, or if
2139 the file is being closed.
2141 Parameters
2142 ----------
2143 force: bool
2144 When closing, write the last block even if it is smaller than
2145 blocks are allowed to be. Disallows further writing to this file.
2146 """
2148 if self.closed:
2149 raise ValueError("Flush on closed file")
2150 if force and self.forced:
2151 raise ValueError("Force flush cannot be called more than once")
2152 if force:
2153 self.forced = True
2155 if self.readable():
2156 # no-op to flush on read-mode
2157 return
2159 if not force and self.buffer.tell() < self.blocksize:
2160 # Defer write on small block
2161 return
2163 if self.offset is None:
2164 # Initialize a multipart upload
2165 self.offset = 0
2166 try:
2167 self._initiate_upload()
2168 except Exception:
2169 self.closed = True
2170 raise
2172 if self._upload_chunk(final=force) is not False:
2173 self.offset += self.buffer.seek(0, 2)
2174 self.buffer = io.BytesIO()
2176 def _upload_chunk(self, final=False):
2177 """Write one part of a multi-block file upload
2179 Parameters
2180 ==========
2181 final: bool
2182 This is the last block, so should complete file, if
2183 self.autocommit is True.
2184 """
2185 # may not yet have been initialized, may need to call _initialize_upload
2187 def _initiate_upload(self):
2188 """Create remote file/upload"""
2189 pass
2191 def _fetch_range(self, start, end):
2192 """Get the specified set of bytes from remote"""
2193 return self.fs.cat_file(self.path, start=start, end=end)
2195 def read(self, length=-1):
2196 """
2197 Return data from cache, or fetch pieces as necessary
2199 Parameters
2200 ----------
2201 length: int (-1)
2202 Number of bytes to read; if <0, all remaining bytes.
2203 """
2204 length = -1 if length is None else int(length)
2205 if self.mode != "rb":
2206 raise ValueError("File not in read mode")
2207 if length < 0:
2208 length = self.size - self.loc
2209 if self.closed:
2210 raise ValueError("I/O operation on closed file.")
2211 if length == 0:
2212 # don't even bother calling fetch
2213 return b""
2214 out = self.cache._fetch(self.loc, self.loc + length)
2216 logger.debug(
2217 "%s read: %i - %i %s",
2218 self,
2219 self.loc,
2220 self.loc + length,
2221 self.cache._log_stats(),
2222 )
2223 self.loc += len(out)
2224 return out
2226 def readinto(self, b):
2227 """mirrors builtin file's readinto method
2229 https://docs.python.org/3/library/io.html#io.RawIOBase.readinto
2230 """
2231 out = memoryview(b).cast("B")
2232 data = self.read(out.nbytes)
2233 out[: len(data)] = data
2234 return len(data)
2236 def readuntil(self, char=b"\n", blocks=None):
2237 """Return data between current position and first occurrence of char
2239 char is included in the output, except if the end of the tile is
2240 encountered first.
2242 Parameters
2243 ----------
2244 char: bytes
2245 Thing to find
2246 blocks: None or int
2247 How much to read in each go. Defaults to file blocksize - which may
2248 mean a new read on every call.
2249 """
2250 out = []
2251 while True:
2252 start = self.tell()
2253 part = self.read(blocks or self.blocksize)
2254 if len(part) == 0:
2255 break
2256 found = part.find(char)
2257 if found > -1:
2258 out.append(part[: found + len(char)])
2259 self.seek(start + found + len(char))
2260 break
2261 out.append(part)
2262 return b"".join(out)
2264 def readline(self):
2265 """Read until and including the first occurrence of newline character
2267 Note that, because of character encoding, this is not necessarily a
2268 true line ending.
2269 """
2270 return self.readuntil(b"\n")
2272 def __next__(self):
2273 out = self.readline()
2274 if out:
2275 return out
2276 raise StopIteration
2278 def __iter__(self):
2279 return self
2281 def readlines(self):
2282 """Return all data, split by the newline character, including the newline character"""
2283 data = self.read()
2284 lines = data.split(b"\n")
2285 out = [l + b"\n" for l in lines[:-1]]
2286 if data.endswith(b"\n"):
2287 return out
2288 else:
2289 return out + [lines[-1]]
2290 # return list(self) ???
2292 def readinto1(self, b):
2293 return self.readinto(b)
2295 def close(self):
2296 """Close file
2298 Finalizes writes, discards cache
2299 """
2300 if getattr(self, "_unclosable", False):
2301 return
2302 if self.closed:
2303 return
2304 try:
2305 if self.mode == "rb":
2306 cache = getattr(self, "cache", None)
2307 if cache is not None:
2308 close = getattr(cache, "close", None)
2309 if callable(close):
2310 close()
2311 self.cache = None
2312 else:
2313 if not getattr(self, "forced", True):
2314 self.flush(force=True)
2316 if self.fs is not None:
2317 self.fs.invalidate_cache(self.path)
2318 self.fs.invalidate_cache(self.fs._parent(self.path))
2319 finally:
2320 self.closed = True
2322 def readable(self):
2323 """Whether opened for reading"""
2324 return "r" in self.mode and not self.closed
2326 def seekable(self):
2327 """Whether is seekable (only in read mode)"""
2328 return self.readable()
2330 def writable(self):
2331 """Whether opened for writing"""
2332 return self.mode in {"wb", "ab", "xb"} and not self.closed
2334 def __reduce__(self):
2335 if self.mode != "rb":
2336 raise RuntimeError("Pickling a writeable file is not supported")
2338 return reopen, (
2339 self.fs,
2340 self.path,
2341 self.mode,
2342 self.blocksize,
2343 self.loc,
2344 self.size,
2345 self.autocommit,
2346 self.cache.name if self.cache else "none",
2347 self.kwargs,
2348 )
2350 def __del__(self):
2351 if not self.closed:
2352 self.close()
2354 def __str__(self):
2355 return f"<File-like object {type(self.fs).__name__}, {self.path}>"
2357 __repr__ = __str__
2359 def __enter__(self):
2360 return self
2362 def __exit__(self, *args):
2363 self.close()
2366def reopen(fs, path, mode, blocksize, loc, size, autocommit, cache_type, kwargs):
2367 file = fs.open(
2368 path,
2369 mode=mode,
2370 block_size=blocksize,
2371 autocommit=autocommit,
2372 cache_type=cache_type,
2373 size=size,
2374 **kwargs,
2375 )
2376 if loc > 0:
2377 file.seek(loc)
2378 return file