Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/fsspec/asyn.py: 26%
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
1import asyncio
2import asyncio.events
3import functools
4import inspect
5import io
6import numbers
7import os
8import re
9import threading
10from collections.abc import Iterable
11from glob import has_magic
12from typing import TYPE_CHECKING
14from .callbacks import DEFAULT_CALLBACK
15from .exceptions import FSTimeoutError
16from .implementations.local import LocalFileSystem, make_path_posix, trailing_sep
17from .spec import AbstractBufferedFile, AbstractFileSystem
18from .utils import check_contained, glob_translate, is_exception, other_paths
20private = re.compile("_[^_]")
21iothread = [None] # dedicated fsspec IO thread
22loop = [None] # global event loop for any non-async instance
23_lock = None # global lock placeholder
24get_running_loop = asyncio.get_running_loop
27def get_lock():
28 """Allocate or return a threading lock.
30 The lock is allocated on first use to allow setting one lock per forked process.
31 """
32 global _lock
33 if not _lock:
34 _lock = threading.Lock()
35 return _lock
38def reset_lock():
39 """Reset the global lock.
41 This should be called only on the init of a forked process to reset the lock to
42 None, enabling the new forked process to get a new lock.
43 """
44 global _lock
46 iothread[0] = None
47 loop[0] = None
48 _lock = None
51async def _runner(event, coro, result, timeout=None):
52 timeout = timeout if timeout else None # convert 0 or 0.0 to None
53 if timeout is not None:
54 coro = asyncio.wait_for(coro, timeout=timeout)
55 try:
56 result[0] = await coro
57 except Exception as ex:
58 result[0] = ex
59 finally:
60 event.set()
63def sync(loop, func, *args, timeout=None, **kwargs):
64 """
65 Make loop run coroutine until it returns. Runs in other thread
67 Examples
68 --------
69 >>> fsspec.asyn.sync(fsspec.asyn.get_loop(), func, *args,
70 timeout=timeout, **kwargs)
71 """
72 timeout = timeout if timeout else None # convert 0 or 0.0 to None
73 # NB: if the loop is not running *yet*, it is OK to submit work
74 # and we will wait for it
75 if loop is None or loop.is_closed():
76 raise RuntimeError("Loop is not running")
77 try:
78 loop0 = asyncio.events.get_running_loop()
79 if loop0 is loop:
80 raise NotImplementedError("Calling sync() from within a running loop")
81 except NotImplementedError:
82 raise
83 except RuntimeError:
84 pass
85 coro = func(*args, **kwargs)
86 result = [None]
87 event = threading.Event()
88 asyncio.run_coroutine_threadsafe(_runner(event, coro, result, timeout), loop)
89 while True:
90 # this loops allows thread to get interrupted
91 if event.wait(1):
92 break
93 if timeout is not None:
94 timeout -= 1
95 if timeout < 0:
96 raise FSTimeoutError
98 return_result = result[0]
99 if isinstance(return_result, asyncio.TimeoutError):
100 # suppress asyncio.TimeoutError, raise FSTimeoutError
101 raise FSTimeoutError from return_result
102 elif isinstance(return_result, BaseException):
103 raise return_result
104 else:
105 return return_result
108def sync_wrapper(func, obj=None):
109 """Given a function, make so can be called in blocking contexts
111 Leave obj=None if defining within a class. Pass the instance if attaching
112 as an attribute of the instance.
113 """
115 @functools.wraps(func)
116 def wrapper(*args, **kwargs):
117 self = obj or args[0]
118 return sync(self.loop, func, *args, **kwargs)
120 return wrapper
123def async_gen_wrapper(func, obj=None):
124 """Given a async generator, make so can be called in blocking contexts"""
126 @functools.wraps(func)
127 def wrapper(*args, **kwargs):
128 self = obj or args[0]
129 gen = func(*args, **kwargs)
130 while True:
131 try:
132 yield sync(self.loop, gen.__anext__)
133 except StopAsyncIteration:
134 break
136 return wrapper
139def get_loop():
140 """Create or return the default fsspec IO loop
142 The loop will be running on a separate thread.
143 """
144 if loop[0] is None:
145 with get_lock():
146 # repeat the check just in case the loop got filled between the
147 # previous two calls from another thread
148 if loop[0] is None:
149 loop[0] = asyncio.new_event_loop()
150 th = threading.Thread(target=loop[0].run_forever, name="fsspecIO")
151 th.daemon = True
152 th.start()
153 iothread[0] = th
154 return loop[0]
157def reset_after_fork():
158 global lock
159 loop[0] = None
160 iothread[0] = None
161 lock = None
164if hasattr(os, "register_at_fork"):
165 # should be posix; this will do nothing for spawn or forkserver subprocesses
166 os.register_at_fork(after_in_child=reset_after_fork)
169if TYPE_CHECKING:
170 import resource
172 ResourceError = resource.error
173else:
174 try:
175 import resource
176 except ImportError:
177 resource = None
178 ResourceError = OSError
179 else:
180 ResourceError = getattr(resource, "error", OSError)
182_DEFAULT_BATCH_SIZE = 128
183_NOFILES_DEFAULT_BATCH_SIZE = 1280
186def _get_batch_size(nofiles=False):
187 from fsspec.config import conf
189 if nofiles:
190 if "nofiles_gather_batch_size" in conf:
191 return conf["nofiles_gather_batch_size"]
192 else:
193 if "gather_batch_size" in conf:
194 return conf["gather_batch_size"]
195 if nofiles:
196 return _NOFILES_DEFAULT_BATCH_SIZE
197 if resource is None:
198 return _DEFAULT_BATCH_SIZE
200 try:
201 soft_limit, _ = resource.getrlimit(resource.RLIMIT_NOFILE)
202 except (ImportError, ValueError, ResourceError):
203 return _DEFAULT_BATCH_SIZE
205 if soft_limit == resource.RLIM_INFINITY:
206 return -1
207 else:
208 return soft_limit // 8
211def running_async() -> bool:
212 """Being executed by an event loop?"""
213 try:
214 asyncio.get_running_loop()
215 return True
216 except RuntimeError:
217 return False
220async def _run_coros_in_chunks(
221 coros,
222 batch_size=None,
223 callback=DEFAULT_CALLBACK,
224 timeout=None,
225 return_exceptions=False,
226 nofiles=False,
227):
228 """Run the given coroutines in chunks.
230 Parameters
231 ----------
232 coros: list of coroutines to run
233 batch_size: int or None
234 Number of coroutines to submit/wait on simultaneously.
235 If -1, then it will not be any throttling. If
236 None, it will be inferred from _get_batch_size()
237 callback: fsspec.callbacks.Callback instance
238 Gets a relative_update when each coroutine completes
239 timeout: number or None
240 If given, each coroutine times out after this time. Note that, since
241 there are multiple batches, the total run time of this function will in
242 general be longer
243 return_exceptions: bool
244 Same meaning as in asyncio.gather
245 nofiles: bool
246 If inferring the batch_size, does this operation involve local files?
247 If yes, you normally expect smaller batches.
248 """
250 if batch_size is None:
251 batch_size = _get_batch_size(nofiles=nofiles)
253 if batch_size == -1:
254 batch_size = len(coros)
255 elif batch_size <= 0:
256 raise ValueError
258 async def _run_coro(coro, i):
259 try:
260 return await asyncio.wait_for(coro, timeout=timeout), i
261 except Exception as e:
262 if not return_exceptions:
263 raise
264 return e, i
265 finally:
266 callback.relative_update(1)
268 i = 0
269 n = len(coros)
270 results = [None] * n
271 pending = set()
273 while pending or i < n:
274 while len(pending) < batch_size and i < n:
275 pending.add(asyncio.ensure_future(_run_coro(coros[i], i)))
276 i += 1
278 if not pending:
279 break
281 done, pending = await asyncio.wait(pending, return_when=asyncio.FIRST_COMPLETED)
282 first_exc = None
283 while done:
284 task = done.pop()
285 try:
286 result, k = await task
287 results[k] = result
288 except Exception as exc:
289 if first_exc is None:
290 first_exc = exc
292 if first_exc is not None:
293 for task in pending:
294 task.cancel()
295 if pending:
296 await asyncio.gather(*pending, return_exceptions=True)
297 raise first_exc
299 return results
302# these methods should be implemented as async by any async-able backend
303async_methods = [
304 "_ls",
305 "_cat_file",
306 "_get_file",
307 "_put_file",
308 "_rm_file",
309 "_cp_file",
310 "_pipe_file",
311 "_expand_path",
312 "_info",
313 "_isfile",
314 "_isdir",
315 "_exists",
316 "_walk",
317 "_glob",
318 "_find",
319 "_du",
320 "_size",
321 "_mkdir",
322 "_makedirs",
323]
326class AsyncFileSystem(AbstractFileSystem):
327 """Async file operations, default implementations
329 Passes bulk operations to asyncio.gather for concurrent operation.
331 Implementations that have concurrent batch operations and/or async methods
332 should inherit from this class instead of AbstractFileSystem. Docstrings are
333 copied from the un-underscored method in AbstractFileSystem, if not given.
334 """
336 # note that methods do not have docstring here; they will be copied
337 # for _* methods and inferred for overridden methods.
339 async_impl = True
340 mirror_sync_methods = True
341 disable_throttling = False
343 def __init__(self, *args, asynchronous=False, loop=None, batch_size=None, **kwargs):
344 self.asynchronous = asynchronous
345 self._pid = os.getpid()
346 if not asynchronous:
347 self._loop = loop or get_loop()
348 else:
349 self._loop = None
350 self.batch_size = batch_size
351 super().__init__(*args, **kwargs)
353 @property
354 def loop(self):
355 if self._pid != os.getpid():
356 raise RuntimeError("This class is not fork-safe")
357 return self._loop
359 async def _rm_file(self, path, **kwargs):
360 if (
361 inspect.iscoroutinefunction(self._rm)
362 and type(self)._rm is not AsyncFileSystem._rm
363 ):
364 return await self._rm(path, recursive=False, batch_size=1, **kwargs)
365 raise NotImplementedError
367 async def _rm(
368 self, path, recursive=False, batch_size=None, maxdepth=None, **kwargs
369 ):
370 # TODO: implement on_error
371 batch_size = batch_size or self.batch_size
372 path = await self._expand_path(path, recursive=recursive, maxdepth=maxdepth)
373 return await _run_coros_in_chunks(
374 [self._rm_file(p, **kwargs) for p in reversed(path)],
375 batch_size=batch_size,
376 nofiles=True,
377 )
379 async def _cp_file(self, path1, path2, **kwargs):
380 raise NotImplementedError
382 async def _mv_file(self, path1, path2):
383 await self._cp_file(path1, path2)
384 await self._rm_file(path1)
386 async def _copy(
387 self,
388 path1,
389 path2,
390 recursive=False,
391 on_error=None,
392 maxdepth=None,
393 batch_size=None,
394 **kwargs,
395 ):
396 if on_error is None and recursive:
397 on_error = "ignore"
398 elif on_error is None:
399 on_error = "raise"
401 if isinstance(path1, list) and isinstance(path2, list):
402 # No need to expand paths when both source and destination
403 # are provided as lists
404 paths1 = path1
405 paths2 = path2
406 else:
407 source_is_str = isinstance(path1, str)
408 paths1 = await self._expand_path(
409 path1, maxdepth=maxdepth, recursive=recursive
410 )
411 if source_is_str and (not recursive or maxdepth is not None):
412 # Non-recursive glob does not copy directories
413 paths1 = [
414 p for p in paths1 if not (trailing_sep(p) or await self._isdir(p))
415 ]
416 if not paths1:
417 return
419 source_is_file = len(paths1) == 1
420 dest_is_dir = isinstance(path2, str) and (
421 trailing_sep(path2) or await self._isdir(path2)
422 )
424 exists = source_is_str and (
425 (has_magic(path1) and source_is_file)
426 or (not has_magic(path1) and dest_is_dir and not trailing_sep(path1))
427 )
428 paths2 = other_paths(
429 paths1,
430 path2,
431 exists=exists,
432 flatten=not source_is_str,
433 )
435 batch_size = batch_size or self.batch_size
436 coros = [self._cp_file(p1, p2, **kwargs) for p1, p2 in zip(paths1, paths2)]
437 result = await _run_coros_in_chunks(
438 coros, batch_size=batch_size, return_exceptions=True, nofiles=True
439 )
441 for ex in filter(is_exception, result):
442 if on_error == "ignore" and isinstance(ex, FileNotFoundError):
443 continue
444 raise ex
446 async def _pipe_file(self, path, value, mode="overwrite", **kwargs):
447 raise NotImplementedError
449 async def _pipe(self, path, value=None, batch_size=None, **kwargs):
450 if isinstance(path, str):
451 path = {path: value}
452 batch_size = batch_size or self.batch_size
453 return await _run_coros_in_chunks(
454 [self._pipe_file(k, v, **kwargs) for k, v in path.items()],
455 batch_size=batch_size,
456 nofiles=True,
457 )
459 async def _process_limits(self, url, start, end):
460 """Helper for "Range"-based _cat_file"""
461 size = None
462 suff = False
463 if start is not None and start < 0:
464 # if start is negative and end None, end is the "suffix length"
465 if end is None:
466 end = -start
467 start = ""
468 suff = True
469 else:
470 size = size or (await self._info(url))["size"]
471 start = size + start
472 elif start is None:
473 start = 0
474 if not suff:
475 if end is not None and end < 0:
476 if start is not None:
477 size = size or (await self._info(url))["size"]
478 end = size + end
479 elif end is None:
480 end = ""
481 if isinstance(end, numbers.Integral):
482 end -= 1 # bytes range is inclusive
483 return f"bytes={start}-{end}"
485 async def _cat_file(self, path, start=None, end=None, **kwargs):
486 raise NotImplementedError
488 async def _cat(
489 self, path, recursive=False, on_error="raise", batch_size=None, **kwargs
490 ):
491 paths = await self._expand_path(path, recursive=recursive)
492 coros = [self._cat_file(path, **kwargs) for path in paths]
493 batch_size = batch_size or self.batch_size
494 out = await _run_coros_in_chunks(
495 coros, batch_size=batch_size, nofiles=True, return_exceptions=True
496 )
497 if on_error == "raise":
498 ex = next(filter(is_exception, out), False)
499 if ex:
500 raise ex
501 if (
502 len(paths) > 1
503 or isinstance(path, list)
504 or paths[0] != self._strip_protocol(path)
505 ):
506 return {
507 k: v
508 for k, v in zip(paths, out)
509 if on_error != "omit" or not is_exception(v)
510 }
511 else:
512 return out[0]
514 async def _cat_ranges(
515 self,
516 paths,
517 starts,
518 ends,
519 max_gap=None,
520 batch_size=None,
521 on_error="return",
522 **kwargs,
523 ):
524 """Get the contents of byte ranges from one or more files
526 Parameters
527 ----------
528 paths: list
529 A list of of filepaths on this filesystems
530 starts, ends: int or list
531 Bytes limits of the read. If using a single int, the same value will be
532 used to read all the specified files.
533 on_error: "return" or "raise"
534 If "return" (default), any per-range exception is placed in the output
535 list at the corresponding position. Otherwise the first such exception
536 is raised. Matches ``AbstractFileSystem.cat_ranges``.
537 """
538 if max_gap is not None:
539 # use utils.merge_offset_ranges
540 raise NotImplementedError
541 if not isinstance(paths, list):
542 raise TypeError
543 if not isinstance(starts, Iterable):
544 starts = [starts] * len(paths)
545 if not isinstance(ends, Iterable):
546 ends = [ends] * len(paths)
547 if len(starts) != len(paths) or len(ends) != len(paths):
548 raise ValueError
549 coros = [
550 self._cat_file(p, start=s, end=e, **kwargs)
551 for p, s, e in zip(paths, starts, ends)
552 ]
553 batch_size = batch_size or self.batch_size
554 out = await _run_coros_in_chunks(
555 coros, batch_size=batch_size, nofiles=True, return_exceptions=True
556 )
557 if on_error != "return":
558 ex = next(filter(is_exception, out), None)
559 if ex is not None:
560 raise ex
561 return out
563 async def _put_file(self, lpath, rpath, mode="overwrite", **kwargs):
564 raise NotImplementedError
566 async def _put(
567 self,
568 lpath,
569 rpath,
570 recursive=False,
571 callback=DEFAULT_CALLBACK,
572 batch_size=None,
573 maxdepth=None,
574 **kwargs,
575 ):
576 """Copy file(s) from local.
578 Copies a specific file or tree of files (if recursive=True). If rpath
579 ends with a "/", it will be assumed to be a directory, and target files
580 will go within.
582 The put_file method will be called concurrently on a batch of files. The
583 batch_size option can configure the amount of futures that can be executed
584 at the same time. If it is -1, then all the files will be uploaded concurrently.
585 The default can be set for this instance by passing "batch_size" in the
586 constructor, or for all instances by setting the "gather_batch_size" key
587 in ``fsspec.config.conf``, falling back to 1/8th of the system limit .
588 """
589 if isinstance(lpath, list) and isinstance(rpath, list):
590 # No need to expand paths when both source and destination
591 # are provided as lists
592 rpaths = rpath
593 lpaths = lpath
594 else:
595 source_is_str = isinstance(lpath, str)
596 if source_is_str:
597 lpath = make_path_posix(lpath)
598 fs = LocalFileSystem()
599 lpaths = fs.expand_path(lpath, recursive=recursive, maxdepth=maxdepth)
600 if source_is_str and (not recursive or maxdepth is not None):
601 # Non-recursive glob does not copy directories
602 lpaths = [p for p in lpaths if not (trailing_sep(p) or fs.isdir(p))]
603 if not lpaths:
604 return
606 source_is_file = len(lpaths) == 1
607 dest_is_dir = isinstance(rpath, str) and (
608 trailing_sep(rpath) or await self._isdir(rpath)
609 )
611 rpath = self._strip_protocol(rpath)
612 exists = source_is_str and (
613 (has_magic(lpath) and source_is_file)
614 or (not has_magic(lpath) and dest_is_dir and not trailing_sep(lpath))
615 )
616 rpaths = other_paths(
617 lpaths,
618 rpath,
619 exists=exists,
620 flatten=not source_is_str,
621 )
623 is_dir = {l: os.path.isdir(l) for l in lpaths}
624 rdirs = [r for l, r in zip(lpaths, rpaths) if is_dir[l]]
625 file_pairs = [(l, r) for l, r in zip(lpaths, rpaths) if not is_dir[l]]
627 await asyncio.gather(*[self._makedirs(d, exist_ok=True) for d in rdirs])
628 batch_size = batch_size or self.batch_size
630 coros = []
631 callback.set_size(len(file_pairs))
632 for lfile, rfile in file_pairs:
633 put_file = callback.branch_coro(self._put_file)
634 coros.append(put_file(lfile, rfile, **kwargs))
636 return await _run_coros_in_chunks(
637 coros, batch_size=batch_size, callback=callback
638 )
640 async def _get_file(self, rpath, lpath, **kwargs):
641 raise NotImplementedError
643 async def _get(
644 self,
645 rpath,
646 lpath,
647 recursive=False,
648 callback=DEFAULT_CALLBACK,
649 maxdepth=None,
650 **kwargs,
651 ):
652 """Copy file(s) to local.
654 Copies a specific file or tree of files (if recursive=True). If lpath
655 ends with a "/", it will be assumed to be a directory, and target files
656 will go within. Can submit a list of paths, which may be glob-patterns
657 and will be expanded.
659 The get_file method will be called concurrently on a batch of files. The
660 batch_size option can configure the amount of futures that can be executed
661 at the same time. If it is -1, then all the files will be uploaded concurrently.
662 The default can be set for this instance by passing "batch_size" in the
663 constructor, or for all instances by setting the "gather_batch_size" key
664 in ``fsspec.config.conf``, falling back to 1/8th of the system limit .
665 """
666 if isinstance(lpath, list) and isinstance(rpath, list):
667 # No need to expand paths when both source and destination
668 # are provided as lists
669 rpaths = rpath
670 lpaths = lpath
671 else:
672 source_is_str = isinstance(rpath, str)
673 # First check for rpath trailing slash as _strip_protocol removes it.
674 source_not_trailing_sep = source_is_str and not trailing_sep(rpath)
675 rpath = self._strip_protocol(rpath)
676 rpaths = await self._expand_path(
677 rpath, recursive=recursive, maxdepth=maxdepth
678 )
679 if source_is_str and (not recursive or maxdepth is not None):
680 # Non-recursive glob does not copy directories
681 rpaths = [
682 p for p in rpaths if not (trailing_sep(p) or await self._isdir(p))
683 ]
684 if not rpaths:
685 return
687 lpath = make_path_posix(lpath)
688 source_is_file = len(rpaths) == 1
689 dest_is_dir = isinstance(lpath, str) and (
690 trailing_sep(lpath) or LocalFileSystem().isdir(lpath)
691 )
693 exists = source_is_str and (
694 (has_magic(rpath) and source_is_file)
695 or (not has_magic(rpath) and dest_is_dir and source_not_trailing_sep)
696 )
697 lpaths = other_paths(
698 rpaths,
699 lpath,
700 exists=exists,
701 flatten=not source_is_str,
702 )
703 if isinstance(lpath, str):
704 # The names came from the source listing; ".." in one of them
705 # would otherwise place the copy above the destination. When
706 # lpath is a list the caller named every destination itself.
707 check_contained(lpath, lpaths)
709 [os.makedirs(os.path.dirname(lp), exist_ok=True) for lp in lpaths]
710 batch_size = kwargs.pop("batch_size", self.batch_size)
712 coros = []
713 callback.set_size(len(lpaths))
714 for lpath, rpath in zip(lpaths, rpaths):
715 get_file = callback.branch_coro(self._get_file)
716 coros.append(get_file(rpath, lpath, **kwargs))
717 return await _run_coros_in_chunks(
718 coros, batch_size=batch_size, callback=callback
719 )
721 async def _isfile(self, path):
722 try:
723 return (await self._info(path))["type"] == "file"
724 except Exception:
725 return False
727 async def _isdir(self, path):
728 try:
729 return (await self._info(path))["type"] == "directory"
730 except OSError:
731 return False
733 async def _size(self, path):
734 return (await self._info(path)).get("size", None)
736 async def _sizes(self, paths, batch_size=None):
737 batch_size = batch_size or self.batch_size
738 return await _run_coros_in_chunks(
739 [self._size(p) for p in paths], batch_size=batch_size
740 )
742 async def _exists(self, path, **kwargs):
743 try:
744 await self._info(path, **kwargs)
745 return True
746 except FileNotFoundError:
747 return False
749 async def _info(self, path, **kwargs):
750 raise NotImplementedError
752 async def _ls(self, path, detail=True, **kwargs):
753 raise NotImplementedError
755 async def _walk(self, path, maxdepth=None, topdown=True, on_error="omit", **kwargs):
756 if maxdepth is not None and maxdepth < 1:
757 raise ValueError("maxdepth must be at least 1")
759 path = self._strip_protocol(path)
760 full_dirs = {}
761 dirs = {}
762 files = {}
764 detail = kwargs.pop("detail", False)
765 try:
766 listing = await self._ls(path, detail=True, **kwargs)
767 except (FileNotFoundError, OSError) as e:
768 if on_error == "raise":
769 raise
770 elif callable(on_error):
771 on_error(e)
772 if detail:
773 yield path, {}, {}
774 else:
775 yield path, [], []
776 return
778 for info in listing:
779 # each info name must be at least [path]/part , but here
780 # we check also for names like [path]/part/
781 pathname = info["name"].rstrip("/")
782 name = pathname.rsplit("/", 1)[-1]
783 if info["type"] == "directory" and pathname != path:
784 # do not include "self" path
785 full_dirs[name] = pathname
786 dirs[name] = info
787 elif pathname == path:
788 # file-like with same name as give path
789 files[""] = info
790 else:
791 files[name] = info
793 if not detail:
794 dirs = list(dirs)
795 files = list(files)
797 if topdown:
798 # Yield before recursion if walking top down
799 yield path, dirs, files
801 if maxdepth is not None:
802 maxdepth -= 1
803 if maxdepth < 1:
804 if not topdown:
805 yield path, dirs, files
806 return
808 for d in dirs:
809 async for _ in self._walk(
810 full_dirs[d],
811 maxdepth=maxdepth,
812 detail=detail,
813 topdown=topdown,
814 **kwargs,
815 ):
816 yield _
818 if not topdown:
819 # Yield after recursion if walking bottom up
820 yield path, dirs, files
822 async def _glob(self, path, maxdepth=None, **kwargs):
823 if maxdepth is not None and maxdepth < 1:
824 raise ValueError("maxdepth must be at least 1")
826 import re
828 seps = (os.path.sep, os.path.altsep) if os.path.altsep else (os.path.sep,)
829 ends_with_sep = path.endswith(seps) # _strip_protocol strips trailing slash
830 path = self._strip_protocol(path)
831 append_slash_to_dirname = ends_with_sep or path.endswith(
832 tuple(sep + "**" for sep in seps)
833 )
834 idx_star = path.find("*") if path.find("*") >= 0 else len(path)
835 idx_qmark = path.find("?") if path.find("?") >= 0 else len(path)
836 idx_brace = path.find("[") if path.find("[") >= 0 else len(path)
838 min_idx = min(idx_star, idx_qmark, idx_brace)
840 detail = kwargs.pop("detail", False)
841 withdirs = kwargs.pop("withdirs", True)
843 if not has_magic(path):
844 if await self._exists(path, **kwargs):
845 if not detail:
846 return [path]
847 else:
848 return {path: await self._info(path, **kwargs)}
849 else:
850 if not detail:
851 return [] # glob of non-existent returns empty
852 else:
853 return {}
854 elif "/" in path[:min_idx]:
855 first_wildcard_idx = min_idx
856 min_idx = path[:min_idx].rindex("/")
857 root = path[
858 : min_idx + 1
859 ] # everything up to the last / before the first wildcard
860 prefix = path[
861 min_idx + 1 : first_wildcard_idx
862 ] # stem between last "/" and first wildcard
863 depth = path[min_idx + 1 :].count("/") + 1
864 else:
865 root = ""
866 prefix = path[:min_idx] # stem up to the first wildcard
867 depth = path[min_idx + 1 :].count("/") + 1
869 if "**" in path:
870 if maxdepth is not None:
871 idx_double_stars = path.find("**")
872 depth_double_stars = path[idx_double_stars:].count("/") + 1
873 depth = depth - depth_double_stars + maxdepth
874 else:
875 depth = None
877 # Pass the filename stem as prefix= so backends that support it such as
878 # gcsfs, s3fs and adlfs can filter server-side up to the first wildcard.
879 if prefix:
880 kwargs["prefix"] = prefix
881 allpaths = await self._find(
882 root, maxdepth=depth, withdirs=withdirs, detail=True, **kwargs
883 )
885 pattern = glob_translate(path + ("/" if ends_with_sep else ""))
886 pattern = re.compile(pattern)
888 out = {
889 p: info
890 for p, info in sorted(allpaths.items())
891 if pattern.match(
892 p + "/"
893 if append_slash_to_dirname and info["type"] == "directory"
894 else p
895 )
896 }
898 if detail:
899 return out
900 else:
901 return list(out)
903 async def _du(self, path, total=True, maxdepth=None, **kwargs):
904 sizes = {}
905 # async for?
906 for f in await self._find(path, maxdepth=maxdepth, **kwargs):
907 info = await self._info(f)
908 sizes[info["name"]] = info["size"]
909 if total:
910 return sum(sizes.values())
911 else:
912 return sizes
914 async def _find(self, path, maxdepth=None, withdirs=False, **kwargs):
915 path = self._strip_protocol(path)
916 out = {}
917 detail = kwargs.pop("detail", False)
919 # Add the root directory if withdirs is requested
920 # This is needed for posix glob compliance
921 if withdirs and path != "" and await self._isdir(path):
922 out[path] = await self._info(path)
924 # async for?
925 async for _, dirs, files in self._walk(path, maxdepth, detail=True, **kwargs):
926 if withdirs:
927 files.update(dirs)
928 out.update({info["name"]: info for name, info in files.items()})
929 if not out and (await self._isfile(path)):
930 # walk works on directories, but find should also return [path]
931 # when path happens to be a file
932 out[path] = {}
933 names = sorted(out)
934 if not detail:
935 return names
936 else:
937 return {name: out[name] for name in names}
939 async def _expand_path(
940 self, path, recursive=False, maxdepth=None, assume_literal=False
941 ):
942 if maxdepth is not None and maxdepth < 1:
943 raise ValueError("maxdepth must be at least 1")
945 if isinstance(path, str):
946 out = await self._expand_path([path], recursive, maxdepth)
947 else:
948 out = set()
949 path = [self._strip_protocol(p) for p in path]
950 for p in path: # can gather here
951 if not assume_literal and has_magic(p):
952 bit = set(await self._glob(p, maxdepth=maxdepth))
953 out |= bit
954 if recursive:
955 # glob call above expanded one depth so if maxdepth is defined
956 # then decrement it in expand_path call below. If it is zero
957 # after decrementing then avoid expand_path call.
958 if maxdepth is not None and maxdepth <= 1:
959 continue
960 out |= set(
961 await self._expand_path(
962 list(bit),
963 recursive=recursive,
964 maxdepth=maxdepth - 1 if maxdepth is not None else None,
965 assume_literal=True,
966 )
967 )
968 continue
969 elif recursive:
970 rec = set(await self._find(p, maxdepth=maxdepth, withdirs=True))
971 out |= rec
972 if p not in out and (recursive is False or (await self._exists(p))):
973 # should only check once, for the root
974 out.add(p)
975 if not out:
976 raise FileNotFoundError(path)
977 return sorted(out)
979 async def _mkdir(self, path, create_parents=True, **kwargs):
980 pass # not necessary to implement, may not have directories
982 async def _makedirs(self, path, exist_ok=False):
983 pass # not necessary to implement, may not have directories
985 async def open_async(self, path, mode="rb", **kwargs):
986 if "b" not in mode or kwargs.get("compression"):
987 raise ValueError
988 raise NotImplementedError
991def mirror_sync_methods(obj):
992 """Populate sync and async methods for obj
994 For each method will create a sync version if the name refers to an async method
995 (coroutine) and there is no override in the child class; will create an async
996 method for the corresponding sync method if there is no implementation.
998 Uses the methods specified in
999 - async_methods: the set that an implementation is expected to provide
1000 - default_async_methods: that can be derived from their sync version in
1001 AbstractFileSystem
1002 - AsyncFileSystem: async-specific default coroutines
1003 """
1004 from fsspec import AbstractFileSystem
1006 for method in set(async_methods + dir(AsyncFileSystem)):
1007 if not method.startswith("_"):
1008 continue
1009 smethod = method[1:]
1010 if private.match(method):
1011 isco = inspect.iscoroutinefunction(getattr(obj, method, None))
1012 unsync = getattr(getattr(obj, smethod, False), "__func__", None)
1013 is_default = unsync is getattr(AbstractFileSystem, smethod, "")
1014 if isco and is_default:
1015 mth = sync_wrapper(getattr(obj, method), obj=obj)
1016 elif inspect.isasyncgenfunction(getattr(obj, method, None)) and is_default:
1017 mth = async_gen_wrapper(getattr(obj, method), obj=obj)
1018 else:
1019 continue
1020 setattr(obj, smethod, mth)
1021 if not mth.__doc__:
1022 mth.__doc__ = getattr(
1023 getattr(AbstractFileSystem, smethod, None), "__doc__", ""
1024 )
1027class FSSpecCoroutineCancel(Exception):
1028 pass
1031def _dump_running_tasks(
1032 printout=True, cancel=True, exc=FSSpecCoroutineCancel, with_task=False
1033):
1034 import traceback
1036 tasks = [t for t in asyncio.tasks.all_tasks(loop[0]) if not t.done()]
1037 if printout:
1038 [task.print_stack() for task in tasks]
1039 out = [
1040 {
1041 "locals": task._coro.cr_frame.f_locals,
1042 "file": task._coro.cr_frame.f_code.co_filename,
1043 "firstline": task._coro.cr_frame.f_code.co_firstlineno,
1044 "linelo": task._coro.cr_frame.f_lineno,
1045 "stack": traceback.format_stack(task._coro.cr_frame),
1046 "task": task if with_task else None,
1047 }
1048 for task in tasks
1049 ]
1050 if cancel:
1051 for t in tasks:
1052 cbs = t._callbacks
1053 t.cancel()
1054 asyncio.futures.Future.set_exception(t, exc)
1055 asyncio.futures.Future.cancel(t)
1056 [cb[0](t) for cb in cbs] # cancels any dependent concurrent.futures
1057 try:
1058 t._coro.throw(exc) # exits coro, unless explicitly handled
1059 except exc:
1060 pass
1061 return out
1064class AbstractAsyncStreamedFile(AbstractBufferedFile):
1065 # no read buffering, and always auto-commit
1066 # TODO: readahead might still be useful here, but needs async version
1068 async def read(self, length=-1):
1069 """
1070 Return data from cache, or fetch pieces as necessary
1072 Parameters
1073 ----------
1074 length: int (-1)
1075 Number of bytes to read; if <0, all remaining bytes.
1076 """
1077 length = -1 if length is None else int(length)
1078 if self.mode != "rb":
1079 raise ValueError("File not in read mode")
1080 if length < 0:
1081 length = self.size - self.loc
1082 if self.closed:
1083 raise ValueError("I/O operation on closed file.")
1084 if length == 0:
1085 # don't even bother calling fetch
1086 return b""
1087 out = await self._fetch_range(self.loc, self.loc + length)
1088 self.loc += len(out)
1089 return out
1091 async def write(self, data):
1092 """
1093 Write data to buffer.
1095 Buffer only sent on flush() or if buffer is greater than
1096 or equal to blocksize.
1098 Parameters
1099 ----------
1100 data: bytes
1101 Set of bytes to be written.
1102 """
1103 if self.mode not in {"wb", "ab"}:
1104 raise ValueError("File not in write mode")
1105 if self.closed:
1106 raise ValueError("I/O operation on closed file.")
1107 if self.forced:
1108 raise ValueError("This file has been force-flushed, can only close")
1109 out = self.buffer.write(data)
1110 self.loc += out
1111 if self.buffer.tell() >= self.blocksize:
1112 await self.flush()
1113 return out
1115 async def close(self):
1116 """Close file
1118 Finalizes writes, discards cache
1119 """
1120 if getattr(self, "_unclosable", False):
1121 return
1122 if self.closed:
1123 return
1124 if self.mode == "rb":
1125 self.cache = None
1126 else:
1127 if not self.forced:
1128 await self.flush(force=True)
1130 if self.fs is not None:
1131 self.fs.invalidate_cache(self.path)
1132 self.fs.invalidate_cache(self.fs._parent(self.path))
1134 self.closed = True
1136 async def flush(self, force=False):
1137 if self.closed:
1138 raise ValueError("Flush on closed file")
1139 if force and self.forced:
1140 raise ValueError("Force flush cannot be called more than once")
1141 if force:
1142 self.forced = True
1144 if self.mode not in {"wb", "ab"}:
1145 # no-op to flush on read-mode
1146 return
1148 if not force and self.buffer.tell() < self.blocksize:
1149 # Defer write on small block
1150 return
1152 if self.offset is None:
1153 # Initialize a multipart upload
1154 self.offset = 0
1155 try:
1156 await self._initiate_upload()
1157 except:
1158 self.closed = True
1159 raise
1161 if await self._upload_chunk(final=force) is not False:
1162 self.offset += self.buffer.seek(0, 2)
1163 self.buffer = io.BytesIO()
1165 async def __aenter__(self):
1166 return self
1168 async def __aexit__(self, exc_type, exc_val, exc_tb):
1169 await self.close()
1171 async def _fetch_range(self, start, end):
1172 raise NotImplementedError
1174 async def _initiate_upload(self):
1175 pass
1177 async def _upload_chunk(self, final=False):
1178 raise NotImplementedError