Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/fsspec/spec.py: 25%

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

927 statements  

1from __future__ import annotations 

2 

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 

14 

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) 

29 

30logger = logging.getLogger("fsspec") 

31 

32 

33def make_instance(cls, args, kwargs): 

34 return cls(*args, **kwargs) 

35 

36 

37FORK_AVAILABLE = hasattr(os, "register_at_fork") 

38 

39 

40if FORK_AVAILABLE: 

41 _registered_classes = weakref.WeakSet() 

42 

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

48 

49 os.register_at_fork(after_in_child=_reset_instances_lock) 

50 

51 

52class _Cached(type): 

53 """ 

54 Metaclass for caching file system instances. 

55 

56 Notes 

57 ----- 

58 Instances are cached according to 

59 

60 * The values of the class attributes listed in `_extra_tokenize_attributes` 

61 * The arguments passed to ``__init__``. 

62 

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

68 

69 def __init__(cls, *args, **kwargs): 

70 super().__init__(*args, **kwargs) 

71 

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

82 

83 if FORK_AVAILABLE: 

84 _registered_classes.add(cls) 

85 

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 

91 

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

101 

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) 

109 

110 if pid != cls._pid: 

111 with cls._instantiation_lock: 

112 if pid != cls._pid: 

113 cls._cache.clear() 

114 cls._pid = pid 

115 

116 if not skip and cls.cachable: 

117 inst = cls._check_instance_cache(token) 

118 if inst is not None: 

119 return inst 

120 

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 

127 

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 

135 

136 mirror_sync_methods(obj) 

137 

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 

145 

146 cls._latest = token 

147 cls._cache[token] = obj 

148 return obj 

149 

150 

151class AbstractFileSystem(metaclass=_Cached): 

152 """ 

153 An abstract super-class for pythonic file-systems 

154 

155 Implementations are expected to be compatible with or, better, subclass 

156 from here. 

157 """ 

158 

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 

171 

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

176 

177 # Set by _Cached metaclass 

178 storage_args: tuple[Any, ...] 

179 storage_options: dict[str, Any] 

180 

181 def __init__(self, *args, **storage_options): 

182 """Create and configure file-system instance 

183 

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. 

187 

188 A reasonable default should be provided if there are no arguments. 

189 

190 Subclasses should call this method. 

191 

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) 

213 

214 if storage_options.pop("add_docs", None): 

215 warnings.warn("add_docs is no longer supported.", FutureWarning) 

216 

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 

221 

222 @property 

223 def fsid(self): 

224 """Persistent filesystem id that can be used to compare filesystems 

225 across sessions. 

226 """ 

227 raise NotImplementedError 

228 

229 @property 

230 def _fs_token(self): 

231 return self._fs_token_ 

232 

233 def __dask_tokenize__(self): 

234 return self._fs_token 

235 

236 def __hash__(self): 

237 return int(self._fs_token, 16) 

238 

239 def __eq__(self, other): 

240 return isinstance(other, type(self)) and self._fs_token == other._fs_token 

241 

242 def __reduce__(self): 

243 return make_instance, (type(self), self.storage_args, self.storage_options) 

244 

245 @classmethod 

246 def _strip_protocol(cls, path): 

247 """Turn path from fully-qualified to file-system-specific 

248 

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 

263 

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

271 

272 @staticmethod 

273 def _get_kwargs_from_urls(path): 

274 """If kwargs can be encoded in the paths, extract them here 

275 

276 This should happen before instantiation of the class; incoming paths 

277 then should be amended to strip the options in methods. 

278 

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

284 

285 @classmethod 

286 def current(cls): 

287 """Return the most recently instantiated FileSystem 

288 

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

295 

296 @property 

297 def transaction(self): 

298 """A context within which files are committed together upon exit 

299 

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 

306 

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 

312 

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

321 

322 def invalidate_cache(self, path=None): 

323 """ 

324 Discard any cached directory information 

325 

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) 

338 

339 def mkdir(self, path, create_parents=True, **kwargs): 

340 """ 

341 Create directory entry at path 

342 

343 For systems that don't have true directories, may create an for 

344 this instance only and not touch the real filesystem 

345 

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 

356 

357 def makedirs(self, path, exist_ok=False): 

358 """Recursively make directories 

359 

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. 

363 

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 

372 

373 def rmdir(self, path): 

374 """Remove a directory, if empty""" 

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

376 

377 def ls(self, path, detail=True, **kwargs): 

378 """List objects at path. 

379 

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. 

383 

384 The specific keys, or perhaps a FileInfo class, or similar, is TBD, 

385 but must be consistent across implementations. 

386 Must include: 

387 

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 

392 

393 Additional information 

394 may be present, appropriate to the file-system, e.g., generation, 

395 checksum, etc. 

396 

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. 

400 

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 

410 

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 

417 

418 def _ls_from_cache(self, path): 

419 """Check cache for listing 

420 

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 

442 

443 def walk(self, path, maxdepth=None, topdown=True, on_error="omit", **kwargs): 

444 """Return all files under the given path. 

445 

446 List all files, recursing into subdirectories; output is iterator-style, 

447 like ``os.walk()``. For a simple list of files, ``find()`` is available. 

448 

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) 

456 

457 Note that the "files" outputted will include anything that is not 

458 a directory, such as links. 

459 

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

478 

479 path = self._strip_protocol(path) 

480 full_dirs = {} 

481 dirs = {} 

482 files = {} 

483 

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 

493 

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 

508 

509 if not detail: 

510 dirs = list(dirs) 

511 files = list(files) 

512 

513 if topdown: 

514 # Yield before recursion if walking top down 

515 yield path, dirs, files 

516 

517 if maxdepth is not None: 

518 maxdepth -= 1 

519 if maxdepth < 1: 

520 if not topdown: 

521 yield path, dirs, files 

522 return 

523 

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 ) 

532 

533 if not topdown: 

534 # Yield after recursion if walking bottom up 

535 yield path, dirs, files 

536 

537 def find(self, path, maxdepth=None, withdirs=False, detail=False, **kwargs): 

538 """List all files below path. 

539 

540 Like posix ``find`` command without conditions 

541 

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 = {} 

555 

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) 

560 

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} 

574 

575 def du(self, path, total=True, maxdepth=None, withdirs=False, **kwargs): 

576 """Space used by files and optionally directories within a path 

577 

578 Directory size does not include the size of its contents. 

579 

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

590 

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 

608 

609 def glob(self, path, maxdepth=None, **kwargs): 

610 """Find files by glob-matching. 

611 

612 Pattern matching capabilities for finding files that match the given pattern. 

613 

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) 

623 

624 Returns 

625 ------- 

626 List of matched paths, or dict of paths and their info if detail=True 

627 

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 

637 

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 

643 

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

651 

652 import re 

653 

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) 

663 

664 min_idx = min(idx_star, idx_qmark, idx_brace) 

665 

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

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

668 

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 

687 

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 

695 

696 allpaths = self.find( 

697 root, maxdepth=depth, withdirs=withdirs, detail=True, **kwargs 

698 ) 

699 

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

701 pattern = re.compile(pattern) 

702 

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 } 

712 

713 if detail: 

714 return out 

715 else: 

716 return list(out) 

717 

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 

726 

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) 

731 

732 def info(self, path, **kwargs): 

733 """Give details of entry at path 

734 

735 Returns a single dictionary, with exactly the same information as ``ls`` 

736 would with ``detail=True``. 

737 

738 The default implementation calls ls and could be overridden by a 

739 shortcut. kwargs are passed on to ```ls()``. 

740 

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``. 

743 

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) 

765 

766 def checksum(self, path): 

767 """Unique value for current version of file 

768 

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. 

772 

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) 

778 

779 def size(self, path): 

780 """Size in bytes of file""" 

781 return self.info(path).get("size", None) 

782 

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] 

786 

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 

793 

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 

800 

801 def read_text(self, path, encoding=None, errors=None, newline=None, **kwargs): 

802 """Get the contents of the file as a string. 

803 

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

819 

820 def write_text( 

821 self, path, value, encoding=None, errors=None, newline=None, **kwargs 

822 ): 

823 """Write the text to the given file. 

824 

825 An existing file will be overwritten. 

826 

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) 

844 

845 def cat_file(self, path, start=None, end=None, **kwargs): 

846 """Get the content of a file 

847 

848 Parameters 

849 ---------- 

850 path: URL of file on this filesystems 

851 start, end: int 

852 Bytes limits of the read. If negative, backwards from end, 

853 like usual python slices. Either can be None for start or 

854 end of file, respectively 

855 kwargs: passed to ``open()``. 

856 """ 

857 # explicitly set buffering off? 

858 with self.open(path, "rb", **kwargs) as f: 

859 if start is not None: 

860 if start >= 0: 

861 f.seek(start) 

862 else: 

863 f.seek(max(0, f.size + start)) 

864 if end is not None: 

865 if end < 0: 

866 end = f.size + end 

867 return f.read(end - f.tell()) 

868 return f.read() 

869 

870 def pipe_file(self, path, value, mode="overwrite", **kwargs): 

871 """Set the bytes of given file""" 

872 if mode == "create" and self.exists(path): 

873 # non-atomic but simple way; or could use "xb" in open(), which is likely 

874 # not as well supported 

875 raise FileExistsError 

876 with self.open(path, "wb", **kwargs) as f: 

877 f.write(value) 

878 

879 def pipe(self, path, value=None, **kwargs): 

880 """Put value into path 

881 

882 (counterpart to ``cat``) 

883 

884 Parameters 

885 ---------- 

886 path: string or dict(str, bytes) 

887 If a string, a single remote location to put ``value`` bytes; if a dict, 

888 a mapping of {path: bytesvalue}. 

889 value: bytes, optional 

890 If using a single path, these are the bytes to put there. Ignored if 

891 ``path`` is a dict 

892 """ 

893 if isinstance(path, str): 

894 self.pipe_file(self._strip_protocol(path), value, **kwargs) 

895 elif isinstance(path, dict): 

896 for k, v in path.items(): 

897 self.pipe_file(self._strip_protocol(k), v, **kwargs) 

898 else: 

899 raise ValueError("path must be str or dict") 

900 

901 def cat_ranges( 

902 self, paths, starts, ends, max_gap=None, on_error="return", **kwargs 

903 ): 

904 """Get the contents of byte ranges from one or more files 

905 

906 Parameters 

907 ---------- 

908 paths: list 

909 A list of of filepaths on this filesystems 

910 starts, ends: int or list 

911 Bytes limits of the read. If using a single int, the same value will be 

912 used to read all the specified files. 

913 """ 

914 if max_gap is not None: 

915 raise NotImplementedError 

916 if not isinstance(paths, list): 

917 raise TypeError 

918 if not isinstance(starts, list): 

919 starts = [starts] * len(paths) 

920 if not isinstance(ends, list): 

921 ends = [ends] * len(paths) 

922 if len(starts) != len(paths) or len(ends) != len(paths): 

923 raise ValueError 

924 out = [] 

925 for p, s, e in zip(paths, starts, ends): 

926 try: 

927 out.append(self.cat_file(p, s, e, **kwargs)) 

928 except Exception as e: 

929 if on_error == "return": 

930 out.append(e) 

931 else: 

932 raise 

933 return out 

934 

935 def cat(self, path, recursive=False, on_error="raise", **kwargs): 

936 """Fetch (potentially multiple) paths' contents 

937 

938 Parameters 

939 ---------- 

940 recursive: bool 

941 If True, assume the path(s) are directories, and get all the 

942 contained files 

943 on_error : "raise", "omit", "return" 

944 If raise, an underlying exception will be raised (converted to KeyError 

945 if the type is in self.missing_exceptions); if omit, keys with exception 

946 will simply not be included in the output; if "return", all keys are 

947 included in the output, but the value will be bytes or an exception 

948 instance. 

949 kwargs: passed to cat_file 

950 

951 Returns 

952 ------- 

953 dict of {path: contents} if there are multiple paths 

954 or the path has been otherwise expanded 

955 """ 

956 paths = self.expand_path(path, recursive=recursive, **kwargs) 

957 if ( 

958 len(paths) > 1 

959 or isinstance(path, list) 

960 or paths[0] != self._strip_protocol(path) 

961 ): 

962 out = {} 

963 for path in paths: 

964 try: 

965 out[path] = self.cat_file(path, **kwargs) 

966 except Exception as e: 

967 if on_error == "raise": 

968 raise 

969 if on_error == "return": 

970 out[path] = e 

971 return out 

972 else: 

973 return self.cat_file(paths[0], **kwargs) 

974 

975 def get_file( 

976 self, rpath, lpath=None, callback=DEFAULT_CALLBACK, outfile=None, **kwargs 

977 ): 

978 """Copy single remote file to local""" 

979 from .implementations.local import LocalFileSystem 

980 

981 if outfile is None and isfilelike(lpath): 

982 outfile = lpath 

983 elif outfile is None and self.isdir(rpath): 

984 os.makedirs(lpath, exist_ok=True) 

985 return None 

986 

987 if outfile is None: 

988 fs = LocalFileSystem(auto_mkdir=True) 

989 fs.makedirs(fs._parent(lpath), exist_ok=True) 

990 

991 with self.open(rpath, "rb", **kwargs) as f1: 

992 close_outfile = outfile is None 

993 if close_outfile: 

994 outfile = open(lpath, "wb") 

995 

996 try: 

997 callback.set_size(getattr(f1, "size", None)) 

998 data = True 

999 while data: 

1000 data = f1.read(self.blocksize) 

1001 segment_len = outfile.write(data) 

1002 if segment_len is None: 

1003 segment_len = len(data) 

1004 callback.relative_update(segment_len) 

1005 finally: 

1006 if close_outfile: 

1007 outfile.close() 

1008 

1009 def get( 

1010 self, 

1011 rpath, 

1012 lpath, 

1013 recursive=False, 

1014 callback=DEFAULT_CALLBACK, 

1015 maxdepth=None, 

1016 **kwargs, 

1017 ): 

1018 """Copy file(s) to local. 

1019 

1020 Copies a specific file or tree of files (if recursive=True). If lpath 

1021 ends with a "/", it will be assumed to be a directory, and target files 

1022 will go within. Can submit a list of paths, which may be glob-patterns 

1023 and will be expanded. 

1024 

1025 Calls get_file for each source. 

1026 """ 

1027 if isinstance(lpath, list) and isinstance(rpath, list): 

1028 # No need to expand paths when both source and destination 

1029 # are provided as lists 

1030 rpaths = rpath 

1031 lpaths = lpath 

1032 else: 

1033 from .implementations.local import ( 

1034 LocalFileSystem, 

1035 make_path_posix, 

1036 trailing_sep, 

1037 ) 

1038 

1039 source_is_str = isinstance(rpath, str) 

1040 rpaths = self.expand_path( 

1041 rpath, recursive=recursive, maxdepth=maxdepth, **kwargs 

1042 ) 

1043 if source_is_str and (not recursive or maxdepth is not None): 

1044 # Non-recursive glob does not copy directories 

1045 rpaths = [p for p in rpaths if not (trailing_sep(p) or self.isdir(p))] 

1046 if not rpaths: 

1047 return 

1048 

1049 if isinstance(lpath, str): 

1050 lpath = make_path_posix(lpath) 

1051 

1052 source_is_file = len(rpaths) == 1 

1053 dest_is_dir = isinstance(lpath, str) and ( 

1054 trailing_sep(lpath) or LocalFileSystem().isdir(lpath) 

1055 ) 

1056 

1057 exists = source_is_str and ( 

1058 (has_magic(rpath) and source_is_file) 

1059 or (not has_magic(rpath) and dest_is_dir and not trailing_sep(rpath)) 

1060 ) 

1061 lpaths = other_paths( 

1062 rpaths, 

1063 lpath, 

1064 exists=exists, 

1065 flatten=not source_is_str, 

1066 ) 

1067 if isinstance(lpath, str): 

1068 # The names came from the source listing; ".." in one of them 

1069 # would otherwise place the copy above the destination. When 

1070 # lpath is a list the caller named every destination itself. 

1071 check_contained(lpath, lpaths) 

1072 

1073 callback.set_size(len(lpaths)) 

1074 for lpath, rpath in callback.wrap(zip(lpaths, rpaths)): 

1075 with callback.branched(rpath, lpath) as child: 

1076 self.get_file(rpath, lpath, callback=child, **kwargs) 

1077 

1078 def put_file( 

1079 self, lpath, rpath, callback=DEFAULT_CALLBACK, mode="overwrite", **kwargs 

1080 ): 

1081 """Copy single file to remote""" 

1082 if mode == "create" and self.exists(rpath): 

1083 raise FileExistsError 

1084 if os.path.isdir(lpath): 

1085 self.makedirs(rpath, exist_ok=True) 

1086 return None 

1087 

1088 with open(lpath, "rb") as f1: 

1089 size = f1.seek(0, 2) 

1090 callback.set_size(size) 

1091 f1.seek(0) 

1092 

1093 self.mkdirs(self._parent(os.fspath(rpath)), exist_ok=True) 

1094 with self.open(rpath, "wb", **kwargs) as f2: 

1095 while f1.tell() < size: 

1096 data = f1.read(self.blocksize) 

1097 segment_len = f2.write(data) 

1098 if segment_len is None: 

1099 segment_len = len(data) 

1100 callback.relative_update(segment_len) 

1101 

1102 def put( 

1103 self, 

1104 lpath, 

1105 rpath, 

1106 recursive=False, 

1107 callback=DEFAULT_CALLBACK, 

1108 maxdepth=None, 

1109 **kwargs, 

1110 ): 

1111 """Copy file(s) from local. 

1112 

1113 Copies a specific file or tree of files (if recursive=True). If rpath 

1114 ends with a "/", it will be assumed to be a directory, and target files 

1115 will go within. 

1116 

1117 Calls put_file for each source. 

1118 """ 

1119 if isinstance(lpath, list) and isinstance(rpath, list): 

1120 # No need to expand paths when both source and destination 

1121 # are provided as lists 

1122 rpaths = rpath 

1123 lpaths = lpath 

1124 else: 

1125 from .implementations.local import ( 

1126 LocalFileSystem, 

1127 make_path_posix, 

1128 trailing_sep, 

1129 ) 

1130 

1131 source_is_str = isinstance(lpath, str) 

1132 if source_is_str: 

1133 lpath = make_path_posix(lpath) 

1134 fs = LocalFileSystem() 

1135 lpaths = fs.expand_path( 

1136 lpath, recursive=recursive, maxdepth=maxdepth, **kwargs 

1137 ) 

1138 if source_is_str and (not recursive or maxdepth is not None): 

1139 # Non-recursive glob does not copy directories 

1140 lpaths = [p for p in lpaths if not (trailing_sep(p) or fs.isdir(p))] 

1141 if not lpaths: 

1142 return 

1143 

1144 source_is_file = len(lpaths) == 1 

1145 dest_is_dir = isinstance(rpath, str) and ( 

1146 trailing_sep(rpath) or self.isdir(rpath) 

1147 ) 

1148 

1149 rpath = ( 

1150 self._strip_protocol(rpath) 

1151 if isinstance(rpath, str) 

1152 else [self._strip_protocol(p) for p in rpath] 

1153 ) 

1154 exists = source_is_str and ( 

1155 (has_magic(lpath) and source_is_file) 

1156 or (not has_magic(lpath) and dest_is_dir and not trailing_sep(lpath)) 

1157 ) 

1158 rpaths = other_paths( 

1159 lpaths, 

1160 rpath, 

1161 exists=exists, 

1162 flatten=not source_is_str, 

1163 ) 

1164 

1165 callback.set_size(len(rpaths)) 

1166 for lpath, rpath in callback.wrap(zip(lpaths, rpaths)): 

1167 with callback.branched(lpath, rpath) as child: 

1168 self.put_file(lpath, rpath, callback=child, **kwargs) 

1169 

1170 def head(self, path, size=1024): 

1171 """Get the first ``size`` bytes from file""" 

1172 with self.open(path, "rb") as f: 

1173 return f.read(size) 

1174 

1175 def tail(self, path, size=1024): 

1176 """Get the last ``size`` bytes from file""" 

1177 with self.open(path, "rb") as f: 

1178 f.seek(max(-size, -f.size), 2) 

1179 return f.read() 

1180 

1181 def cp_file(self, path1, path2, **kwargs): 

1182 raise NotImplementedError 

1183 

1184 def copy( 

1185 self, path1, path2, recursive=False, maxdepth=None, on_error=None, **kwargs 

1186 ): 

1187 """Copy within two locations in the filesystem 

1188 

1189 on_error : "raise", "ignore" 

1190 If raise, any not-found exceptions will be raised; if ignore any 

1191 not-found exceptions will cause the path to be skipped; defaults to 

1192 raise unless recursive is true, where the default is ignore 

1193 """ 

1194 if on_error is None and recursive: 

1195 on_error = "ignore" 

1196 elif on_error is None: 

1197 on_error = "raise" 

1198 

1199 if isinstance(path1, list) and isinstance(path2, list): 

1200 # No need to expand paths when both source and destination 

1201 # are provided as lists 

1202 paths1 = path1 

1203 paths2 = path2 

1204 else: 

1205 from .implementations.local import trailing_sep 

1206 

1207 source_is_str = isinstance(path1, str) 

1208 paths1 = self.expand_path( 

1209 path1, recursive=recursive, maxdepth=maxdepth, **kwargs 

1210 ) 

1211 if source_is_str and (not recursive or maxdepth is not None): 

1212 # Non-recursive glob does not copy directories 

1213 paths1 = [p for p in paths1 if not (trailing_sep(p) or self.isdir(p))] 

1214 if not paths1: 

1215 return 

1216 

1217 source_is_file = len(paths1) == 1 

1218 dest_is_dir = isinstance(path2, str) and ( 

1219 trailing_sep(path2) or self.isdir(path2) 

1220 ) 

1221 

1222 exists = source_is_str and ( 

1223 (has_magic(path1) and source_is_file) 

1224 or (not has_magic(path1) and dest_is_dir and not trailing_sep(path1)) 

1225 ) 

1226 paths2 = other_paths( 

1227 paths1, 

1228 path2, 

1229 exists=exists, 

1230 flatten=not source_is_str, 

1231 ) 

1232 

1233 for p1, p2 in zip(paths1, paths2): 

1234 try: 

1235 self.cp_file(p1, p2, **kwargs) 

1236 except FileNotFoundError: 

1237 if on_error == "raise": 

1238 raise 

1239 

1240 def expand_path( 

1241 self, path, recursive=False, maxdepth=None, assume_literal=False, **kwargs 

1242 ): 

1243 """Turn one or more globs or directories into a list of all matching paths 

1244 to files or directories. 

1245 

1246 kwargs are passed to ``glob`` or ``find``, which may in turn call ``ls`` 

1247 """ 

1248 

1249 if maxdepth is not None and maxdepth < 1: 

1250 raise ValueError("maxdepth must be at least 1") 

1251 

1252 if isinstance(path, (str, os.PathLike)): 

1253 out = self.expand_path([path], recursive, maxdepth, **kwargs) 

1254 else: 

1255 out = set() 

1256 path = [self._strip_protocol(p) for p in path] 

1257 for p in path: 

1258 if not assume_literal and has_magic(p): 

1259 bit = set(self.glob(p, maxdepth=maxdepth, **kwargs)) 

1260 out |= bit 

1261 if recursive: 

1262 # glob call above expanded one depth so if maxdepth is defined 

1263 # then decrement it in expand_path call below. If it is zero 

1264 # after decrementing then avoid expand_path call. 

1265 if maxdepth is not None and maxdepth <= 1: 

1266 continue 

1267 out |= set( 

1268 self.expand_path( 

1269 list(bit), 

1270 recursive=recursive, 

1271 maxdepth=maxdepth - 1 if maxdepth is not None else None, 

1272 assume_literal=True, 

1273 **kwargs, 

1274 ) 

1275 ) 

1276 continue 

1277 elif recursive: 

1278 rec = set( 

1279 self.find( 

1280 p, maxdepth=maxdepth, withdirs=True, detail=False, **kwargs 

1281 ) 

1282 ) 

1283 out |= rec 

1284 if p not in out and (recursive is False or self.exists(p)): 

1285 # should only check once, for the root 

1286 out.add(p) 

1287 if not out: 

1288 raise FileNotFoundError(path) 

1289 return sorted(out) 

1290 

1291 def mv(self, path1, path2, recursive=False, maxdepth=None, **kwargs): 

1292 """Move file(s) from one location to another""" 

1293 if path1 == path2: 

1294 logger.debug("%s mv: The paths are the same, so no files were moved.", self) 

1295 else: 

1296 # explicitly raise exception to prevent data corruption 

1297 self.copy( 

1298 path1, path2, recursive=recursive, maxdepth=maxdepth, on_error="raise" 

1299 ) 

1300 self.rm(path1, recursive=recursive) 

1301 

1302 def rm_file(self, path): 

1303 """Delete a file""" 

1304 self._rm(path) 

1305 

1306 def _rm(self, path): 

1307 """Delete one file""" 

1308 # this is the old name for the method, prefer rm_file 

1309 raise NotImplementedError 

1310 

1311 def rm(self, path, recursive=False, maxdepth=None): 

1312 """Delete files or directories. 

1313 

1314 Parameters 

1315 ---------- 

1316 path: str or list of str 

1317 Files or directories to delete. 

1318 recursive: bool 

1319 If True, recursively delete directories and their contents. 

1320 maxdepth: int or None 

1321 Depth to pass to walk for finding files to delete, if recursive. 

1322 If None, there will be no limit and infinite recursion may be 

1323 possible. 

1324 """ 

1325 path = self.expand_path(path, recursive=recursive, maxdepth=maxdepth) 

1326 for p in reversed(path): 

1327 self.rm_file(p) 

1328 

1329 @classmethod 

1330 def _parent(cls, path): 

1331 path = cls._strip_protocol(path) 

1332 if "/" in path: 

1333 parent = path.rsplit("/", 1)[0].lstrip(cls.root_marker) 

1334 return cls.root_marker + parent 

1335 else: 

1336 return cls.root_marker 

1337 

1338 def _open( 

1339 self, 

1340 path, 

1341 mode="rb", 

1342 block_size=None, 

1343 autocommit=True, 

1344 cache_options=None, 

1345 **kwargs, 

1346 ): 

1347 """Return raw bytes-mode file-like from the file-system""" 

1348 return AbstractBufferedFile( 

1349 self, 

1350 path, 

1351 mode, 

1352 block_size, 

1353 autocommit, 

1354 cache_options=cache_options, 

1355 **kwargs, 

1356 ) 

1357 

1358 def open( 

1359 self, 

1360 path, 

1361 mode="rb", 

1362 block_size=None, 

1363 cache_options=None, 

1364 compression=None, 

1365 **kwargs, 

1366 ): 

1367 """ 

1368 Return a file-like object from the filesystem 

1369 

1370 The resultant instance must function correctly in a context ``with`` 

1371 block. 

1372 

1373 Parameters 

1374 ---------- 

1375 path: str 

1376 Target file 

1377 mode: str like 'rb', 'w' 

1378 See builtin ``open()`` 

1379 Mode "x" (exclusive write) may be implemented by the backend. Even if 

1380 it is, whether it is checked up front or on commit, and whether it is 

1381 atomic is implementation-dependent. 

1382 block_size: int 

1383 Some indication of buffering - this is a value in bytes 

1384 cache_options : dict, optional 

1385 Extra arguments to pass through to the cache. 

1386 compression: string or None 

1387 If given, open file using compression codec. Can either be a compression 

1388 name (a key in ``fsspec.compression.compr``) or "infer" to guess the 

1389 compression from the filename suffix. 

1390 encoding, errors, newline: passed on to TextIOWrapper for text mode 

1391 """ 

1392 import io 

1393 

1394 path = self._strip_protocol(path) 

1395 if "b" not in mode: 

1396 mode = mode.replace("t", "") + "b" 

1397 

1398 text_kwargs = { 

1399 k: kwargs.pop(k) 

1400 for k in ["encoding", "errors", "newline"] 

1401 if k in kwargs 

1402 } 

1403 return io.TextIOWrapper( 

1404 self.open( 

1405 path, 

1406 mode, 

1407 block_size=block_size, 

1408 cache_options=cache_options, 

1409 compression=compression, 

1410 **kwargs, 

1411 ), 

1412 **text_kwargs, 

1413 ) 

1414 else: 

1415 ac = kwargs.pop("autocommit", not self._intrans) 

1416 f = self._open( 

1417 path, 

1418 mode=mode, 

1419 block_size=block_size, 

1420 autocommit=ac, 

1421 cache_options=cache_options, 

1422 **kwargs, 

1423 ) 

1424 if not ac and "r" not in mode: 

1425 self.transaction.files.append(f) 

1426 if compression is not None: 

1427 from fsspec.compression import compr 

1428 from fsspec.core import get_compression 

1429 

1430 compression = get_compression(path, compression) 

1431 compress = compr[compression] 

1432 f = compress(f, mode=mode[0]) 

1433 return f 

1434 

1435 def touch(self, path, truncate=True, **kwargs): 

1436 """Create empty file, or update timestamp 

1437 

1438 Parameters 

1439 ---------- 

1440 path: str 

1441 file location 

1442 truncate: bool 

1443 If True, always set file size to 0; if False, update timestamp and 

1444 leave file unchanged, if backend allows this 

1445 """ 

1446 if truncate or not self.exists(path): 

1447 with self.open(path, "wb", **kwargs): 

1448 pass 

1449 else: 

1450 raise NotImplementedError # update timestamp, if possible 

1451 

1452 def ukey(self, path): 

1453 """Hash of file properties, to tell if it has changed""" 

1454 return sha256(str(self.info(path)).encode()).hexdigest() 

1455 

1456 def read_block(self, fn, offset, length, delimiter=None): 

1457 """Read a block of bytes from 

1458 

1459 Starting at ``offset`` of the file, read ``length`` bytes. If 

1460 ``delimiter`` is set then we ensure that the read starts and stops at 

1461 delimiter boundaries that follow the locations ``offset`` and ``offset 

1462 + length``. If ``offset`` is zero then we start at zero. The 

1463 bytestring returned WILL include the end delimiter string. 

1464 

1465 If offset+length is beyond the eof, reads to eof. 

1466 

1467 Parameters 

1468 ---------- 

1469 fn: string 

1470 Path to filename 

1471 offset: int 

1472 Byte offset to start read 

1473 length: int 

1474 Number of bytes to read. If None, read to end. 

1475 delimiter: bytes (optional) 

1476 Ensure reading starts and stops at delimiter bytestring 

1477 

1478 Examples 

1479 -------- 

1480 >>> fs.read_block('data/file.csv', 0, 13) # doctest: +SKIP 

1481 b'Alice, 100\\nBo' 

1482 >>> fs.read_block('data/file.csv', 0, 13, delimiter=b'\\n') # doctest: +SKIP 

1483 b'Alice, 100\\nBob, 200\\n' 

1484 

1485 Use ``length=None`` to read to the end of the file. 

1486 >>> fs.read_block('data/file.csv', 0, None, delimiter=b'\\n') # doctest: +SKIP 

1487 b'Alice, 100\\nBob, 200\\nCharlie, 300' 

1488 

1489 See Also 

1490 -------- 

1491 :func:`fsspec.utils.read_block` 

1492 """ 

1493 with self.open(fn, "rb") as f: 

1494 size = f.size 

1495 if length is None: 

1496 length = size 

1497 if size is not None and offset + length > size: 

1498 length = size - offset 

1499 return read_block(f, offset, length, delimiter) 

1500 

1501 def to_json(self, *, include_password: bool = True) -> str: 

1502 """ 

1503 JSON representation of this filesystem instance. 

1504 

1505 Parameters 

1506 ---------- 

1507 include_password: bool, default True 

1508 Whether to include the password (if any) in the output. 

1509 

1510 Returns 

1511 ------- 

1512 JSON string with keys ``cls`` (the python location of this class), 

1513 protocol (text name of this class's protocol, first one in case of 

1514 multiple), ``args`` (positional args, usually empty), and all other 

1515 keyword arguments as their own keys. 

1516 

1517 Warnings 

1518 -------- 

1519 Serialized filesystems may contain sensitive information which have been 

1520 passed to the constructor, such as passwords and tokens. Make sure you 

1521 store and send them in a secure environment! 

1522 """ 

1523 from .json import FilesystemJSONEncoder 

1524 

1525 return json.dumps( 

1526 self, 

1527 cls=type( 

1528 "_FilesystemJSONEncoder", 

1529 (FilesystemJSONEncoder,), 

1530 {"include_password": include_password}, 

1531 ), 

1532 ) 

1533 

1534 @staticmethod 

1535 def from_json(blob: str) -> AbstractFileSystem: 

1536 """ 

1537 Recreate a filesystem instance from JSON representation. 

1538 

1539 See ``.to_json()`` for the expected structure of the input. 

1540 

1541 Parameters 

1542 ---------- 

1543 blob: str 

1544 

1545 Returns 

1546 ------- 

1547 file system instance, not necessarily of this particular class. 

1548 

1549 Warnings 

1550 -------- 

1551 This can import arbitrary modules (as determined by the ``cls`` key). 

1552 Make sure you haven't installed any modules that may execute malicious code 

1553 at import time. 

1554 """ 

1555 from .json import FilesystemJSONDecoder 

1556 

1557 return json.loads(blob, cls=FilesystemJSONDecoder) 

1558 

1559 def to_dict(self, *, include_password: bool = True) -> dict[str, Any]: 

1560 """ 

1561 JSON-serializable dictionary representation of this filesystem instance. 

1562 

1563 Parameters 

1564 ---------- 

1565 include_password: bool, default True 

1566 Whether to include the password (if any) in the output. 

1567 

1568 Returns 

1569 ------- 

1570 Dictionary with keys ``cls`` (the python location of this class), 

1571 protocol (text name of this class's protocol, first one in case of 

1572 multiple), ``args`` (positional args, usually empty), and all other 

1573 keyword arguments as their own keys. 

1574 

1575 Warnings 

1576 -------- 

1577 Serialized filesystems may contain sensitive information which have been 

1578 passed to the constructor, such as passwords and tokens. Make sure you 

1579 store and send them in a secure environment! 

1580 """ 

1581 from .json import FilesystemJSONEncoder 

1582 

1583 json_encoder = FilesystemJSONEncoder() 

1584 

1585 cls = type(self) 

1586 proto = self.protocol 

1587 

1588 storage_options = dict(self.storage_options) 

1589 if not include_password: 

1590 storage_options.pop("password", None) 

1591 

1592 return dict( 

1593 cls=f"{cls.__module__}:{cls.__name__}", 

1594 protocol=proto[0] if isinstance(proto, (tuple, list)) else proto, 

1595 args=json_encoder.make_serializable(self.storage_args), 

1596 **json_encoder.make_serializable(storage_options), 

1597 ) 

1598 

1599 @staticmethod 

1600 def from_dict(dct: dict[str, Any]) -> AbstractFileSystem: 

1601 """ 

1602 Recreate a filesystem instance from dictionary representation. 

1603 

1604 See ``.to_dict()`` for the expected structure of the input. 

1605 

1606 Parameters 

1607 ---------- 

1608 dct: Dict[str, Any] 

1609 

1610 Returns 

1611 ------- 

1612 file system instance, not necessarily of this particular class. 

1613 

1614 Warnings 

1615 -------- 

1616 This can import arbitrary modules (as determined by the ``cls`` key). 

1617 Make sure you haven't installed any modules that may execute malicious code 

1618 at import time. 

1619 """ 

1620 from .json import FilesystemJSONDecoder 

1621 

1622 json_decoder = FilesystemJSONDecoder() 

1623 

1624 dct = dict(dct) # Defensive copy 

1625 

1626 cls = FilesystemJSONDecoder.try_resolve_fs_cls(dct) 

1627 if cls is None: 

1628 raise ValueError("Not a serialized AbstractFileSystem") 

1629 

1630 dct.pop("cls", None) 

1631 dct.pop("protocol", None) 

1632 

1633 return cls( 

1634 *json_decoder.unmake_serializable(dct.pop("args", ())), 

1635 **json_decoder.unmake_serializable(dct), 

1636 ) 

1637 

1638 def _get_pyarrow_filesystem(self): 

1639 """ 

1640 Make a version of the FS instance which will be acceptable to pyarrow 

1641 """ 

1642 # all instances already also derive from pyarrow 

1643 return self 

1644 

1645 def get_mapper(self, root="", check=False, create=False, missing_exceptions=None): 

1646 """Create key/value store based on this file-system 

1647 

1648 Makes a MutableMapping interface to the FS at the given root path. 

1649 See ``fsspec.mapping.FSMap`` for further details. 

1650 """ 

1651 from .mapping import FSMap 

1652 

1653 return FSMap( 

1654 root, 

1655 self, 

1656 check=check, 

1657 create=create, 

1658 missing_exceptions=missing_exceptions, 

1659 ) 

1660 

1661 @classmethod 

1662 def clear_instance_cache(cls): 

1663 """ 

1664 Clear the cache of filesystem instances. 

1665 

1666 Notes 

1667 ----- 

1668 Unless overridden by setting the ``cachable`` class attribute to False, 

1669 the filesystem class stores a reference to newly created instances. This 

1670 prevents Python's normal rules around garbage collection from working, 

1671 since the instances refcount will not drop to zero until 

1672 ``clear_instance_cache`` is called. 

1673 """ 

1674 cls._cache.clear() 

1675 

1676 def created(self, path): 

1677 """Return the created timestamp of a file as a datetime.datetime""" 

1678 raise NotImplementedError 

1679 

1680 def modified(self, path): 

1681 """Return the modified timestamp of a file as a datetime.datetime""" 

1682 raise NotImplementedError 

1683 

1684 def tree( 

1685 self, 

1686 path: str = "/", 

1687 recursion_limit: int = 2, 

1688 max_display: int = 25, 

1689 display_size: bool = False, 

1690 prefix: str = "", 

1691 is_last: bool = True, 

1692 first: bool = True, 

1693 indent_size: int = 4, 

1694 ) -> str: 

1695 """ 

1696 Return a tree-like structure of the filesystem starting from the given path as a string. 

1697 

1698 Parameters 

1699 ---------- 

1700 path: Root path to start traversal from 

1701 recursion_limit: Maximum depth of directory traversal 

1702 max_display: Maximum number of items to display per directory 

1703 display_size: Whether to display file sizes 

1704 prefix: Current line prefix for visual tree structure 

1705 is_last: Whether current item is last in its level 

1706 first: Whether this is the first call (displays root path) 

1707 indent_size: Number of spaces by indent 

1708 

1709 Returns 

1710 ------- 

1711 str: A string representing the tree structure. 

1712 

1713 Example 

1714 ------- 

1715 >>> from fsspec import filesystem 

1716 

1717 >>> fs = filesystem('ftp', host='test.rebex.net', user='demo', password='password') 

1718 >>> tree = fs.tree(display_size=True, recursion_limit=3, indent_size=8, max_display=10) 

1719 >>> print(tree) 

1720 """ 

1721 

1722 def format_bytes(n: int) -> str: 

1723 """Format bytes as text.""" 

1724 for prefix, k in ( 

1725 ("P", 2**50), 

1726 ("T", 2**40), 

1727 ("G", 2**30), 

1728 ("M", 2**20), 

1729 ("k", 2**10), 

1730 ): 

1731 if n >= 0.9 * k: 

1732 return f"{n / k:.2f} {prefix}b" 

1733 return f"{n}B" 

1734 

1735 result = [] 

1736 

1737 if first: 

1738 result.append(path) 

1739 

1740 if recursion_limit: 

1741 indent = " " * indent_size 

1742 contents = self.ls(path, detail=True) 

1743 contents.sort( 

1744 key=lambda x: (x.get("type") != "directory", x.get("name", "")) 

1745 ) 

1746 

1747 if max_display is not None and len(contents) > max_display: 

1748 displayed_contents = contents[:max_display] 

1749 remaining_count = len(contents) - max_display 

1750 else: 

1751 displayed_contents = contents 

1752 remaining_count = 0 

1753 

1754 for i, item in enumerate(displayed_contents): 

1755 is_last_item = (i == len(displayed_contents) - 1) and ( 

1756 remaining_count == 0 

1757 ) 

1758 

1759 branch = ( 

1760 "└" + ("─" * (indent_size - 2)) 

1761 if is_last_item 

1762 else "├" + ("─" * (indent_size - 2)) 

1763 ) 

1764 branch += " " 

1765 new_prefix = prefix + ( 

1766 indent if is_last_item else "│" + " " * (indent_size - 1) 

1767 ) 

1768 

1769 name = os.path.basename(item.get("name", "")) 

1770 

1771 if display_size and item.get("type") == "directory": 

1772 sub_contents = self.ls(item.get("name", ""), detail=True) 

1773 num_files = sum( 

1774 1 for sub_item in sub_contents if sub_item.get("type") == "file" 

1775 ) 

1776 num_folders = sum( 

1777 1 

1778 for sub_item in sub_contents 

1779 if sub_item.get("type") == "directory" 

1780 ) 

1781 

1782 if num_files == 0 and num_folders == 0: 

1783 size = " (empty folder)" 

1784 elif num_files == 0: 

1785 size = f" ({num_folders} subfolder{'s' if num_folders > 1 else ''})" 

1786 elif num_folders == 0: 

1787 size = f" ({num_files} file{'s' if num_files > 1 else ''})" 

1788 else: 

1789 size = f" ({num_files} file{'s' if num_files > 1 else ''}, {num_folders} subfolder{'s' if num_folders > 1 else ''})" 

1790 elif display_size and item.get("type") == "file": 

1791 size = f" ({format_bytes(item.get('size', 0))})" 

1792 else: 

1793 size = "" 

1794 

1795 result.append(f"{prefix}{branch}{name}{size}") 

1796 

1797 if item.get("type") == "directory" and recursion_limit > 0: 

1798 result.append( 

1799 self.tree( 

1800 path=item.get("name", ""), 

1801 recursion_limit=recursion_limit - 1, 

1802 max_display=max_display, 

1803 display_size=display_size, 

1804 prefix=new_prefix, 

1805 is_last=is_last_item, 

1806 first=False, 

1807 indent_size=indent_size, 

1808 ) 

1809 ) 

1810 

1811 if remaining_count > 0: 

1812 more_message = f"{remaining_count} more item(s) not displayed." 

1813 result.append( 

1814 f"{prefix}{'└' + ('─' * (indent_size - 2))} {more_message}" 

1815 ) 

1816 

1817 return "\n".join(_ for _ in result if _) 

1818 

1819 # ------------------------------------------------------------------------ 

1820 # Aliases 

1821 

1822 def read_bytes(self, path, start=None, end=None, **kwargs): 

1823 """Alias of `AbstractFileSystem.cat_file`.""" 

1824 return self.cat_file(path, start=start, end=end, **kwargs) 

1825 

1826 def write_bytes(self, path, value, **kwargs): 

1827 """Alias of `AbstractFileSystem.pipe_file`.""" 

1828 self.pipe_file(path, value, **kwargs) 

1829 

1830 def makedir(self, path, create_parents=True, **kwargs): 

1831 """Alias of `AbstractFileSystem.mkdir`.""" 

1832 return self.mkdir(path, create_parents=create_parents, **kwargs) 

1833 

1834 def mkdirs(self, path, exist_ok=False): 

1835 """Alias of `AbstractFileSystem.makedirs`.""" 

1836 return self.makedirs(path, exist_ok=exist_ok) 

1837 

1838 def listdir(self, path, detail=True, **kwargs): 

1839 """Alias of `AbstractFileSystem.ls`.""" 

1840 return self.ls(path, detail=detail, **kwargs) 

1841 

1842 def cp(self, path1, path2, **kwargs): 

1843 """Alias of `AbstractFileSystem.copy`.""" 

1844 return self.copy(path1, path2, **kwargs) 

1845 

1846 def move(self, path1, path2, **kwargs): 

1847 """Alias of `AbstractFileSystem.mv`.""" 

1848 return self.mv(path1, path2, **kwargs) 

1849 

1850 def stat(self, path, **kwargs): 

1851 """Alias of `AbstractFileSystem.info`.""" 

1852 return self.info(path, **kwargs) 

1853 

1854 def disk_usage(self, path, total=True, maxdepth=None, **kwargs): 

1855 """Alias of `AbstractFileSystem.du`.""" 

1856 return self.du(path, total=total, maxdepth=maxdepth, **kwargs) 

1857 

1858 def rename(self, path1, path2, **kwargs): 

1859 """Alias of `AbstractFileSystem.mv`.""" 

1860 return self.mv(path1, path2, **kwargs) 

1861 

1862 def delete(self, path, recursive=False, maxdepth=None): 

1863 """Alias of `AbstractFileSystem.rm`.""" 

1864 return self.rm(path, recursive=recursive, maxdepth=maxdepth) 

1865 

1866 def upload(self, lpath, rpath, recursive=False, **kwargs): 

1867 """Alias of `AbstractFileSystem.put`.""" 

1868 return self.put(lpath, rpath, recursive=recursive, **kwargs) 

1869 

1870 def download(self, rpath, lpath, recursive=False, **kwargs): 

1871 """Alias of `AbstractFileSystem.get`.""" 

1872 return self.get(rpath, lpath, recursive=recursive, **kwargs) 

1873 

1874 def sign(self, path, expiration=100, **kwargs): 

1875 """Create a signed URL representing the given path 

1876 

1877 Some implementations allow temporary URLs to be generated, as a 

1878 way of delegating credentials. 

1879 

1880 Parameters 

1881 ---------- 

1882 path : str 

1883 The path on the filesystem 

1884 expiration : int 

1885 Number of seconds to enable the URL for (if supported) 

1886 

1887 Returns 

1888 ------- 

1889 URL : str 

1890 The signed URL 

1891 

1892 Raises 

1893 ------ 

1894 NotImplementedError : if method is not implemented for a filesystem 

1895 """ 

1896 raise NotImplementedError("Sign is not implemented for this filesystem") 

1897 

1898 def _isfilestore(self): 

1899 # Originally inherited from pyarrow DaskFileSystem. Keeping this 

1900 # here for backwards compatibility as long as pyarrow uses its 

1901 # legacy fsspec-compatible filesystems and thus accepts fsspec 

1902 # filesystems as well 

1903 return False 

1904 

1905 

1906class AbstractBufferedFile(io.IOBase): 

1907 """Convenient class to derive from to provide buffering 

1908 

1909 In the case that the backend does not provide a pythonic file-like object 

1910 already, this class contains much of the logic to build one. The only 

1911 methods that need to be overridden are ``_upload_chunk``, 

1912 ``_initiate_upload`` and ``_fetch_range``. 

1913 """ 

1914 

1915 DEFAULT_BLOCK_SIZE = 5 * 2**20 

1916 _details = None 

1917 

1918 def __init__( 

1919 self, 

1920 fs, 

1921 path, 

1922 mode="rb", 

1923 block_size="default", 

1924 autocommit=True, 

1925 cache_type="readahead", 

1926 cache_options=None, 

1927 size=None, 

1928 **kwargs, 

1929 ): 

1930 """ 

1931 Template for files with buffered reading and writing 

1932 

1933 Parameters 

1934 ---------- 

1935 fs: instance of FileSystem 

1936 path: str 

1937 location in file-system 

1938 mode: str 

1939 Normal file modes. Currently only 'wb', 'ab' or 'rb'. Some file 

1940 systems may be read-only, and some may not support append. 

1941 block_size: int 

1942 Buffer size for reading or writing, 'default' for class default 

1943 autocommit: bool 

1944 Whether to write to final destination; may only impact what 

1945 happens when file is being closed. 

1946 cache_type: {"readahead", "none", "mmap", "bytes"}, default "readahead" 

1947 Caching policy in read mode. See the definitions in ``core``. 

1948 cache_options : dict 

1949 Additional options passed to the constructor for the cache specified 

1950 by `cache_type`. 

1951 size: int 

1952 If given and in read mode, suppressed having to look up the file size 

1953 kwargs: 

1954 Gets stored as self.kwargs 

1955 """ 

1956 from .core import caches 

1957 

1958 self.path = path 

1959 self.fs = fs 

1960 self.mode = mode 

1961 self.blocksize = ( 

1962 self.DEFAULT_BLOCK_SIZE if block_size in ["default", None] else block_size 

1963 ) 

1964 self.loc = 0 

1965 self.autocommit = autocommit 

1966 self.end = None 

1967 self.start = None 

1968 self.closed = False 

1969 

1970 if cache_options is None: 

1971 cache_options = {} 

1972 

1973 if "trim" in kwargs: 

1974 warnings.warn( 

1975 "Passing 'trim' to control the cache behavior has been deprecated. " 

1976 "Specify it within the 'cache_options' argument instead.", 

1977 FutureWarning, 

1978 ) 

1979 cache_options["trim"] = kwargs.pop("trim") 

1980 

1981 self.kwargs = kwargs 

1982 

1983 if mode not in {"ab", "rb", "wb", "xb"}: 

1984 raise NotImplementedError("File mode not supported") 

1985 if mode == "rb": 

1986 if size is not None: 

1987 self.size = size 

1988 else: 

1989 self.size = self.details["size"] 

1990 self.cache = caches[cache_type]( 

1991 self.blocksize, self._fetch_range, self.size, **cache_options 

1992 ) 

1993 else: 

1994 self.buffer = io.BytesIO() 

1995 self.offset = None 

1996 self.forced = False 

1997 self.location = None 

1998 

1999 @property 

2000 def details(self): 

2001 if self._details is None: 

2002 self._details = self.fs.info(self.path) 

2003 return self._details 

2004 

2005 @details.setter 

2006 def details(self, value): 

2007 self._details = value 

2008 self.size = value["size"] 

2009 

2010 @property 

2011 def full_name(self): 

2012 return _unstrip_protocol(self.path, self.fs) 

2013 

2014 @property 

2015 def closed(self): 

2016 # get around this attr being read-only in IOBase 

2017 # use getattr here, since this can be called during del 

2018 return getattr(self, "_closed", True) 

2019 

2020 @closed.setter 

2021 def closed(self, c): 

2022 self._closed = c 

2023 

2024 def __hash__(self): 

2025 if "w" in self.mode: 

2026 return id(self) 

2027 else: 

2028 return int(tokenize(self.details), 16) 

2029 

2030 def __eq__(self, other): 

2031 """Files are equal if they have the same checksum, only in read mode""" 

2032 if self is other: 

2033 return True 

2034 return ( 

2035 isinstance(other, type(self)) 

2036 and self.mode == "rb" 

2037 and other.mode == "rb" 

2038 and hash(self) == hash(other) 

2039 ) 

2040 

2041 def commit(self): 

2042 """Move from temp to final destination""" 

2043 

2044 def discard(self): 

2045 """Throw away temporary file""" 

2046 

2047 def info(self): 

2048 """File information about this path""" 

2049 if self.readable(): 

2050 return self.details 

2051 else: 

2052 raise ValueError("Info not available while writing") 

2053 

2054 def tell(self): 

2055 """Current file location""" 

2056 return self.loc 

2057 

2058 def seek(self, loc, whence=0): 

2059 """Set current file location 

2060 

2061 Parameters 

2062 ---------- 

2063 loc: int 

2064 byte location 

2065 whence: {0, 1, 2} 

2066 from start of file, current location or end of file, resp. 

2067 """ 

2068 loc = int(loc) 

2069 if not self.mode == "rb": 

2070 raise OSError(ESPIPE, "Seek only available in read mode") 

2071 if whence == 0: 

2072 nloc = loc 

2073 elif whence == 1: 

2074 nloc = self.loc + loc 

2075 elif whence == 2: 

2076 nloc = self.size + loc 

2077 else: 

2078 raise ValueError(f"invalid whence ({whence}, should be 0, 1 or 2)") 

2079 if nloc < 0: 

2080 raise ValueError("Seek before start of file") 

2081 self.loc = nloc 

2082 return self.loc 

2083 

2084 def write(self, data): 

2085 """ 

2086 Write data to buffer. 

2087 

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

2089 or equal to blocksize. 

2090 

2091 Parameters 

2092 ---------- 

2093 data: bytes 

2094 Set of bytes to be written. 

2095 """ 

2096 if not self.writable(): 

2097 raise ValueError("File not in write mode") 

2098 if self.closed: 

2099 raise ValueError("I/O operation on closed file.") 

2100 if self.forced: 

2101 raise ValueError("This file has been force-flushed, can only close") 

2102 out = self.buffer.write(data) 

2103 self.loc += out 

2104 if self.buffer.tell() >= self.blocksize: 

2105 self.flush() 

2106 return out 

2107 

2108 def flush(self, force=False): 

2109 """ 

2110 Write buffered data to backend store. 

2111 

2112 Writes the current buffer, if it is larger than the block-size, or if 

2113 the file is being closed. 

2114 

2115 Parameters 

2116 ---------- 

2117 force: bool 

2118 When closing, write the last block even if it is smaller than 

2119 blocks are allowed to be. Disallows further writing to this file. 

2120 """ 

2121 

2122 if self.closed: 

2123 raise ValueError("Flush on closed file") 

2124 if force and self.forced: 

2125 raise ValueError("Force flush cannot be called more than once") 

2126 if force: 

2127 self.forced = True 

2128 

2129 if self.readable(): 

2130 # no-op to flush on read-mode 

2131 return 

2132 

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

2134 # Defer write on small block 

2135 return 

2136 

2137 if self.offset is None: 

2138 # Initialize a multipart upload 

2139 self.offset = 0 

2140 try: 

2141 self._initiate_upload() 

2142 except Exception: 

2143 self.closed = True 

2144 raise 

2145 

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

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

2148 self.buffer = io.BytesIO() 

2149 

2150 def _upload_chunk(self, final=False): 

2151 """Write one part of a multi-block file upload 

2152 

2153 Parameters 

2154 ========== 

2155 final: bool 

2156 This is the last block, so should complete file, if 

2157 self.autocommit is True. 

2158 """ 

2159 # may not yet have been initialized, may need to call _initialize_upload 

2160 

2161 def _initiate_upload(self): 

2162 """Create remote file/upload""" 

2163 pass 

2164 

2165 def _fetch_range(self, start, end): 

2166 """Get the specified set of bytes from remote""" 

2167 return self.fs.cat_file(self.path, start=start, end=end) 

2168 

2169 def read(self, length=-1): 

2170 """ 

2171 Return data from cache, or fetch pieces as necessary 

2172 

2173 Parameters 

2174 ---------- 

2175 length: int (-1) 

2176 Number of bytes to read; if <0, all remaining bytes. 

2177 """ 

2178 length = -1 if length is None else int(length) 

2179 if self.mode != "rb": 

2180 raise ValueError("File not in read mode") 

2181 if length < 0: 

2182 length = self.size - self.loc 

2183 if self.closed: 

2184 raise ValueError("I/O operation on closed file.") 

2185 if length == 0: 

2186 # don't even bother calling fetch 

2187 return b"" 

2188 out = self.cache._fetch(self.loc, self.loc + length) 

2189 

2190 logger.debug( 

2191 "%s read: %i - %i %s", 

2192 self, 

2193 self.loc, 

2194 self.loc + length, 

2195 self.cache._log_stats(), 

2196 ) 

2197 self.loc += len(out) 

2198 return out 

2199 

2200 def readinto(self, b): 

2201 """mirrors builtin file's readinto method 

2202 

2203 https://docs.python.org/3/library/io.html#io.RawIOBase.readinto 

2204 """ 

2205 out = memoryview(b).cast("B") 

2206 data = self.read(out.nbytes) 

2207 out[: len(data)] = data 

2208 return len(data) 

2209 

2210 def readuntil(self, char=b"\n", blocks=None): 

2211 """Return data between current position and first occurrence of char 

2212 

2213 char is included in the output, except if the end of the tile is 

2214 encountered first. 

2215 

2216 Parameters 

2217 ---------- 

2218 char: bytes 

2219 Thing to find 

2220 blocks: None or int 

2221 How much to read in each go. Defaults to file blocksize - which may 

2222 mean a new read on every call. 

2223 """ 

2224 out = [] 

2225 while True: 

2226 start = self.tell() 

2227 part = self.read(blocks or self.blocksize) 

2228 if len(part) == 0: 

2229 break 

2230 found = part.find(char) 

2231 if found > -1: 

2232 out.append(part[: found + len(char)]) 

2233 self.seek(start + found + len(char)) 

2234 break 

2235 out.append(part) 

2236 return b"".join(out) 

2237 

2238 def readline(self): 

2239 """Read until and including the first occurrence of newline character 

2240 

2241 Note that, because of character encoding, this is not necessarily a 

2242 true line ending. 

2243 """ 

2244 return self.readuntil(b"\n") 

2245 

2246 def __next__(self): 

2247 out = self.readline() 

2248 if out: 

2249 return out 

2250 raise StopIteration 

2251 

2252 def __iter__(self): 

2253 return self 

2254 

2255 def readlines(self): 

2256 """Return all data, split by the newline character, including the newline character""" 

2257 data = self.read() 

2258 lines = data.split(b"\n") 

2259 out = [l + b"\n" for l in lines[:-1]] 

2260 if data.endswith(b"\n"): 

2261 return out 

2262 else: 

2263 return out + [lines[-1]] 

2264 # return list(self) ??? 

2265 

2266 def readinto1(self, b): 

2267 return self.readinto(b) 

2268 

2269 def close(self): 

2270 """Close file 

2271 

2272 Finalizes writes, discards cache 

2273 """ 

2274 if getattr(self, "_unclosable", False): 

2275 return 

2276 if self.closed: 

2277 return 

2278 try: 

2279 if self.mode == "rb": 

2280 cache = getattr(self, "cache", None) 

2281 if cache is not None: 

2282 close = getattr(cache, "close", None) 

2283 if callable(close): 

2284 close() 

2285 self.cache = None 

2286 else: 

2287 if not getattr(self, "forced", True): 

2288 self.flush(force=True) 

2289 

2290 if self.fs is not None: 

2291 self.fs.invalidate_cache(self.path) 

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

2293 finally: 

2294 self.closed = True 

2295 

2296 def readable(self): 

2297 """Whether opened for reading""" 

2298 return "r" in self.mode and not self.closed 

2299 

2300 def seekable(self): 

2301 """Whether is seekable (only in read mode)""" 

2302 return self.readable() 

2303 

2304 def writable(self): 

2305 """Whether opened for writing""" 

2306 return self.mode in {"wb", "ab", "xb"} and not self.closed 

2307 

2308 def __reduce__(self): 

2309 if self.mode != "rb": 

2310 raise RuntimeError("Pickling a writeable file is not supported") 

2311 

2312 return reopen, ( 

2313 self.fs, 

2314 self.path, 

2315 self.mode, 

2316 self.blocksize, 

2317 self.loc, 

2318 self.size, 

2319 self.autocommit, 

2320 self.cache.name if self.cache else "none", 

2321 self.kwargs, 

2322 ) 

2323 

2324 def __del__(self): 

2325 if not self.closed: 

2326 self.close() 

2327 

2328 def __str__(self): 

2329 return f"<File-like object {type(self.fs).__name__}, {self.path}>" 

2330 

2331 __repr__ = __str__ 

2332 

2333 def __enter__(self): 

2334 return self 

2335 

2336 def __exit__(self, *args): 

2337 self.close() 

2338 

2339 

2340def reopen(fs, path, mode, blocksize, loc, size, autocommit, cache_type, kwargs): 

2341 file = fs.open( 

2342 path, 

2343 mode=mode, 

2344 block_size=blocksize, 

2345 autocommit=autocommit, 

2346 cache_type=cache_type, 

2347 size=size, 

2348 **kwargs, 

2349 ) 

2350 if loc > 0: 

2351 file.seek(loc) 

2352 return file