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

635 statements  

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 

13 

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 

19 

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 

25 

26 

27def get_lock(): 

28 """Allocate or return a threading lock. 

29 

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 

36 

37 

38def reset_lock(): 

39 """Reset the global lock. 

40 

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 

45 

46 iothread[0] = None 

47 loop[0] = None 

48 _lock = None 

49 

50 

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() 

61 

62 

63def sync(loop, func, *args, timeout=None, **kwargs): 

64 """ 

65 Make loop run coroutine until it returns. Runs in other thread 

66 

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 

97 

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 

106 

107 

108def sync_wrapper(func, obj=None): 

109 """Given a function, make so can be called in blocking contexts 

110 

111 Leave obj=None if defining within a class. Pass the instance if attaching 

112 as an attribute of the instance. 

113 """ 

114 

115 @functools.wraps(func) 

116 def wrapper(*args, **kwargs): 

117 self = obj or args[0] 

118 return sync(self.loop, func, *args, **kwargs) 

119 

120 return wrapper 

121 

122 

123def async_gen_wrapper(func, obj=None): 

124 """Given a async generator, make so can be called in blocking contexts""" 

125 

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 

135 

136 return wrapper 

137 

138 

139def get_loop(): 

140 """Create or return the default fsspec IO loop 

141 

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] 

155 

156 

157def reset_after_fork(): 

158 global lock 

159 loop[0] = None 

160 iothread[0] = None 

161 lock = None 

162 

163 

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) 

167 

168 

169if TYPE_CHECKING: 

170 import resource 

171 

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) 

181 

182_DEFAULT_BATCH_SIZE = 128 

183_NOFILES_DEFAULT_BATCH_SIZE = 1280 

184 

185 

186def _get_batch_size(nofiles=False): 

187 from fsspec.config import conf 

188 

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 

199 

200 try: 

201 soft_limit, _ = resource.getrlimit(resource.RLIMIT_NOFILE) 

202 except (ImportError, ValueError, ResourceError): 

203 return _DEFAULT_BATCH_SIZE 

204 

205 if soft_limit == resource.RLIM_INFINITY: 

206 return -1 

207 else: 

208 return soft_limit // 8 

209 

210 

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 

218 

219 

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. 

229 

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 """ 

249 

250 if batch_size is None: 

251 batch_size = _get_batch_size(nofiles=nofiles) 

252 

253 if batch_size == -1: 

254 batch_size = len(coros) 

255 elif batch_size <= 0: 

256 raise ValueError 

257 

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) 

267 

268 i = 0 

269 n = len(coros) 

270 results = [None] * n 

271 pending = set() 

272 

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 

277 

278 if not pending: 

279 break 

280 

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 

291 

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 

298 

299 return results 

300 

301 

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] 

324 

325 

326class AsyncFileSystem(AbstractFileSystem): 

327 """Async file operations, default implementations 

328 

329 Passes bulk operations to asyncio.gather for concurrent operation. 

330 

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 """ 

335 

336 # note that methods do not have docstring here; they will be copied 

337 # for _* methods and inferred for overridden methods. 

338 

339 async_impl = True 

340 mirror_sync_methods = True 

341 disable_throttling = False 

342 

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) 

352 

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 

358 

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 

366 

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 ) 

378 

379 async def _cp_file(self, path1, path2, **kwargs): 

380 raise NotImplementedError 

381 

382 async def _mv_file(self, path1, path2): 

383 await self._cp_file(path1, path2) 

384 await self._rm_file(path1) 

385 

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" 

400 

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 

418 

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 ) 

423 

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 ) 

434 

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 ) 

440 

441 for ex in filter(is_exception, result): 

442 if on_error == "ignore" and isinstance(ex, FileNotFoundError): 

443 continue 

444 raise ex 

445 

446 async def _pipe_file(self, path, value, mode="overwrite", **kwargs): 

447 raise NotImplementedError 

448 

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 ) 

458 

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}" 

484 

485 async def _cat_file(self, path, start=None, end=None, **kwargs): 

486 raise NotImplementedError 

487 

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] 

513 

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 

525 

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 

562 

563 async def _put_file(self, lpath, rpath, mode="overwrite", **kwargs): 

564 raise NotImplementedError 

565 

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. 

577 

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. 

581 

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 

605 

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 ) 

610 

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 ) 

622 

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]] 

626 

627 await asyncio.gather(*[self._makedirs(d, exist_ok=True) for d in rdirs]) 

628 batch_size = batch_size or self.batch_size 

629 

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)) 

635 

636 return await _run_coros_in_chunks( 

637 coros, batch_size=batch_size, callback=callback 

638 ) 

639 

640 async def _get_file(self, rpath, lpath, **kwargs): 

641 raise NotImplementedError 

642 

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. 

653 

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. 

658 

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 

686 

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 ) 

692 

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) 

708 

709 [os.makedirs(os.path.dirname(lp), exist_ok=True) for lp in lpaths] 

710 batch_size = kwargs.pop("batch_size", self.batch_size) 

711 

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 ) 

720 

721 async def _isfile(self, path): 

722 try: 

723 return (await self._info(path))["type"] == "file" 

724 except Exception: 

725 return False 

726 

727 async def _isdir(self, path): 

728 try: 

729 return (await self._info(path))["type"] == "directory" 

730 except OSError: 

731 return False 

732 

733 async def _size(self, path): 

734 return (await self._info(path)).get("size", None) 

735 

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 ) 

741 

742 async def _exists(self, path, **kwargs): 

743 try: 

744 await self._info(path, **kwargs) 

745 return True 

746 except FileNotFoundError: 

747 return False 

748 

749 async def _info(self, path, **kwargs): 

750 raise NotImplementedError 

751 

752 async def _ls(self, path, detail=True, **kwargs): 

753 raise NotImplementedError 

754 

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") 

758 

759 path = self._strip_protocol(path) 

760 full_dirs = {} 

761 dirs = {} 

762 files = {} 

763 

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 

777 

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 

792 

793 if not detail: 

794 dirs = list(dirs) 

795 files = list(files) 

796 

797 if topdown: 

798 # Yield before recursion if walking top down 

799 yield path, dirs, files 

800 

801 if maxdepth is not None: 

802 maxdepth -= 1 

803 if maxdepth < 1: 

804 if not topdown: 

805 yield path, dirs, files 

806 return 

807 

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 _ 

817 

818 if not topdown: 

819 # Yield after recursion if walking bottom up 

820 yield path, dirs, files 

821 

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") 

825 

826 import re 

827 

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) 

837 

838 min_idx = min(idx_star, idx_qmark, idx_brace) 

839 

840 detail = kwargs.pop("detail", False) 

841 withdirs = kwargs.pop("withdirs", True) 

842 

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 

868 

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 

876 

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 ) 

884 

885 pattern = glob_translate(path + ("/" if ends_with_sep else "")) 

886 pattern = re.compile(pattern) 

887 

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 } 

897 

898 if detail: 

899 return out 

900 else: 

901 return list(out) 

902 

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 

913 

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) 

918 

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) 

923 

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} 

938 

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") 

944 

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) 

978 

979 async def _mkdir(self, path, create_parents=True, **kwargs): 

980 pass # not necessary to implement, may not have directories 

981 

982 async def _makedirs(self, path, exist_ok=False): 

983 pass # not necessary to implement, may not have directories 

984 

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 

989 

990 

991def mirror_sync_methods(obj): 

992 """Populate sync and async methods for obj 

993 

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. 

997 

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 

1005 

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 ) 

1025 

1026 

1027class FSSpecCoroutineCancel(Exception): 

1028 pass 

1029 

1030 

1031def _dump_running_tasks( 

1032 printout=True, cancel=True, exc=FSSpecCoroutineCancel, with_task=False 

1033): 

1034 import traceback 

1035 

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 

1062 

1063 

1064class AbstractAsyncStreamedFile(AbstractBufferedFile): 

1065 # no read buffering, and always auto-commit 

1066 # TODO: readahead might still be useful here, but needs async version 

1067 

1068 async def read(self, length=-1): 

1069 """ 

1070 Return data from cache, or fetch pieces as necessary 

1071 

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 

1090 

1091 async def write(self, data): 

1092 """ 

1093 Write data to buffer. 

1094 

1095 Buffer only sent on flush() or if buffer is greater than 

1096 or equal to blocksize. 

1097 

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 

1114 

1115 async def close(self): 

1116 """Close file 

1117 

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) 

1129 

1130 if self.fs is not None: 

1131 self.fs.invalidate_cache(self.path) 

1132 self.fs.invalidate_cache(self.fs._parent(self.path)) 

1133 

1134 self.closed = True 

1135 

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 

1143 

1144 if self.mode not in {"wb", "ab"}: 

1145 # no-op to flush on read-mode 

1146 return 

1147 

1148 if not force and self.buffer.tell() < self.blocksize: 

1149 # Defer write on small block 

1150 return 

1151 

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 

1160 

1161 if await self._upload_chunk(final=force) is not False: 

1162 self.offset += self.buffer.seek(0, 2) 

1163 self.buffer = io.BytesIO() 

1164 

1165 async def __aenter__(self): 

1166 return self 

1167 

1168 async def __aexit__(self, exc_type, exc_val, exc_tb): 

1169 await self.close() 

1170 

1171 async def _fetch_range(self, start, end): 

1172 raise NotImplementedError 

1173 

1174 async def _initiate_upload(self): 

1175 pass 

1176 

1177 async def _upload_chunk(self, final=False): 

1178 raise NotImplementedError