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

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

941 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 # A negative start or end counts back from the file length. Only 

860 # fsspec's own file classes expose `.size`; a plain file (e.g. from 

861 # the caching filesystems) does not, so measure with a seek instead. 

862 needs_size = (start is not None and start < 0) or ( 

863 end is not None and end < 0 

864 ) 

865 size = None 

866 seeked_for_size = False 

867 if needs_size: 

868 size = getattr(f, "size", None) 

869 if size is None: 

870 size = f.seek(0, 2) 

871 seeked_for_size = True 

872 if start is not None: 

873 if start >= 0: 

874 f.seek(start) 

875 else: 

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

877 elif seeked_for_size: 

878 # measuring the size above moved the cursor to the end; with no 

879 # start given, put it back so the read begins at the start. When 

880 # .size was used instead, or nothing was measured, the cursor 

881 # is untouched, as before, so a plain full read still works 

882 # on non-seekable files. 

883 f.seek(0) 

884 if end is not None: 

885 if end < 0: 

886 end = size + end 

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

888 return f.read() 

889 

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

891 """Set the bytes of given file""" 

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

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

894 # not as well supported 

895 raise FileExistsError 

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

897 f.write(value) 

898 

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

900 """Put value into path 

901 

902 (counterpart to ``cat``) 

903 

904 Parameters 

905 ---------- 

906 path: string or dict(str, bytes) 

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

908 a mapping of {path: bytesvalue}. 

909 value: bytes, optional 

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

911 ``path`` is a dict 

912 """ 

913 if isinstance(path, str): 

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

915 elif isinstance(path, dict): 

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

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

918 else: 

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

920 

921 def cat_ranges( 

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

923 ): 

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

925 

926 Parameters 

927 ---------- 

928 paths: list 

929 A list of of filepaths on this filesystems 

930 starts, ends: int or list 

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

932 used to read all the specified files. 

933 """ 

934 if max_gap is not None: 

935 raise NotImplementedError 

936 if not isinstance(paths, list): 

937 raise TypeError 

938 if not isinstance(starts, list): 

939 starts = [starts] * len(paths) 

940 if not isinstance(ends, list): 

941 ends = [ends] * len(paths) 

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

943 raise ValueError 

944 out = [] 

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

946 try: 

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

948 except Exception as e: 

949 if on_error == "return": 

950 out.append(e) 

951 else: 

952 raise 

953 return out 

954 

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

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

957 

958 Parameters 

959 ---------- 

960 recursive: bool 

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

962 contained files 

963 on_error : "raise", "omit", "return" 

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

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

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

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

968 instance. 

969 kwargs: passed to cat_file 

970 

971 Returns 

972 ------- 

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

974 or the path has been otherwise expanded 

975 """ 

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

977 if ( 

978 len(paths) > 1 

979 or isinstance(path, list) 

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

981 ): 

982 out = {} 

983 for path in paths: 

984 try: 

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

986 except Exception as e: 

987 if on_error == "raise": 

988 raise 

989 if on_error == "return": 

990 out[path] = e 

991 return out 

992 else: 

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

994 

995 def get_file( 

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

997 ): 

998 """Copy single remote file to local""" 

999 from .implementations.local import LocalFileSystem 

1000 

1001 if outfile is None and isfilelike(lpath): 

1002 outfile = lpath 

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

1004 os.makedirs(lpath, exist_ok=True) 

1005 return None 

1006 

1007 if outfile is None: 

1008 fs = LocalFileSystem(auto_mkdir=True) 

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

1010 

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

1012 close_outfile = outfile is None 

1013 if close_outfile: 

1014 outfile = open(lpath, "wb") 

1015 

1016 try: 

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

1018 data = True 

1019 while data: 

1020 data = f1.read(self.blocksize) 

1021 segment_len = outfile.write(data) 

1022 if segment_len is None: 

1023 segment_len = len(data) 

1024 callback.relative_update(segment_len) 

1025 finally: 

1026 if close_outfile: 

1027 outfile.close() 

1028 

1029 def get( 

1030 self, 

1031 rpath, 

1032 lpath, 

1033 recursive=False, 

1034 callback=DEFAULT_CALLBACK, 

1035 maxdepth=None, 

1036 **kwargs, 

1037 ): 

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

1039 

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

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

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

1043 and will be expanded. 

1044 

1045 Calls get_file for each source. 

1046 """ 

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

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

1049 # are provided as lists 

1050 rpaths = rpath 

1051 lpaths = lpath 

1052 else: 

1053 from .implementations.local import ( 

1054 LocalFileSystem, 

1055 make_path_posix, 

1056 trailing_sep, 

1057 ) 

1058 

1059 source_is_str = isinstance(rpath, str) 

1060 rpaths = self.expand_path( 

1061 rpath, recursive=recursive, maxdepth=maxdepth, **kwargs 

1062 ) 

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

1064 # Non-recursive glob does not copy directories 

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

1066 if not rpaths: 

1067 return 

1068 

1069 if isinstance(lpath, str): 

1070 lpath = make_path_posix(lpath) 

1071 

1072 source_is_file = len(rpaths) == 1 

1073 dest_is_dir = isinstance(lpath, str) and ( 

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

1075 ) 

1076 

1077 exists = source_is_str and ( 

1078 (has_magic(rpath) and source_is_file) 

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

1080 ) 

1081 lpaths = other_paths( 

1082 rpaths, 

1083 lpath, 

1084 exists=exists, 

1085 flatten=not source_is_str, 

1086 ) 

1087 if isinstance(lpath, str): 

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

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

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

1091 check_contained(lpath, lpaths) 

1092 

1093 callback.set_size(len(lpaths)) 

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

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

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

1097 

1098 def put_file( 

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

1100 ): 

1101 """Copy single file to remote""" 

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

1103 raise FileExistsError 

1104 if os.path.isdir(lpath): 

1105 self.makedirs(rpath, exist_ok=True) 

1106 return None 

1107 

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

1109 size = f1.seek(0, 2) 

1110 callback.set_size(size) 

1111 f1.seek(0) 

1112 

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

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

1115 while f1.tell() < size: 

1116 data = f1.read(self.blocksize) 

1117 segment_len = f2.write(data) 

1118 if segment_len is None: 

1119 segment_len = len(data) 

1120 callback.relative_update(segment_len) 

1121 

1122 def put( 

1123 self, 

1124 lpath, 

1125 rpath, 

1126 recursive=False, 

1127 callback=DEFAULT_CALLBACK, 

1128 maxdepth=None, 

1129 **kwargs, 

1130 ): 

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

1132 

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

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

1135 will go within. 

1136 

1137 Calls put_file for each source. 

1138 """ 

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

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

1141 # are provided as lists 

1142 rpaths = rpath 

1143 lpaths = lpath 

1144 else: 

1145 from .implementations.local import ( 

1146 LocalFileSystem, 

1147 make_path_posix, 

1148 trailing_sep, 

1149 ) 

1150 

1151 source_is_str = isinstance(lpath, str) 

1152 if source_is_str: 

1153 lpath = make_path_posix(lpath) 

1154 fs = LocalFileSystem() 

1155 lpaths = fs.expand_path( 

1156 lpath, recursive=recursive, maxdepth=maxdepth, **kwargs 

1157 ) 

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

1159 # Non-recursive glob does not copy directories 

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

1161 if not lpaths: 

1162 return 

1163 

1164 source_is_file = len(lpaths) == 1 

1165 dest_is_dir = isinstance(rpath, str) and ( 

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

1167 ) 

1168 

1169 rpath = ( 

1170 self._strip_protocol(rpath) 

1171 if isinstance(rpath, str) 

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

1173 ) 

1174 exists = source_is_str and ( 

1175 (has_magic(lpath) and source_is_file) 

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

1177 ) 

1178 rpaths = other_paths( 

1179 lpaths, 

1180 rpath, 

1181 exists=exists, 

1182 flatten=not source_is_str, 

1183 ) 

1184 

1185 callback.set_size(len(rpaths)) 

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

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

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

1189 

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

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

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

1193 return f.read(size) 

1194 

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

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

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

1198 f.seek(0, 2) 

1199 f.seek(max(f.tell() - size, 0)) 

1200 return f.read() 

1201 

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

1203 raise NotImplementedError 

1204 

1205 def copy( 

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

1207 ): 

1208 """Copy within two locations in the filesystem 

1209 

1210 on_error : "raise", "ignore" 

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

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

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

1214 """ 

1215 if on_error is None and recursive: 

1216 on_error = "ignore" 

1217 elif on_error is None: 

1218 on_error = "raise" 

1219 

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

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

1222 # are provided as lists 

1223 paths1 = path1 

1224 paths2 = path2 

1225 else: 

1226 from .implementations.local import trailing_sep 

1227 

1228 source_is_str = isinstance(path1, str) 

1229 paths1 = self.expand_path( 

1230 path1, recursive=recursive, maxdepth=maxdepth, **kwargs 

1231 ) 

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

1233 # Non-recursive glob does not copy directories 

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

1235 if not paths1: 

1236 return 

1237 

1238 source_is_file = len(paths1) == 1 

1239 dest_is_dir = isinstance(path2, str) and ( 

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

1241 ) 

1242 

1243 exists = source_is_str and ( 

1244 (has_magic(path1) and source_is_file) 

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

1246 ) 

1247 paths2 = other_paths( 

1248 paths1, 

1249 path2, 

1250 exists=exists, 

1251 flatten=not source_is_str, 

1252 ) 

1253 

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

1255 try: 

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

1257 except FileNotFoundError: 

1258 if on_error == "raise": 

1259 raise 

1260 

1261 def expand_path( 

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

1263 ): 

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

1265 to files or directories. 

1266 

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

1268 """ 

1269 

1270 if maxdepth is not None and maxdepth < 1: 

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

1272 

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

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

1275 else: 

1276 out = set() 

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

1278 for p in path: 

1279 if not assume_literal and has_magic(p): 

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

1281 out |= bit 

1282 if recursive: 

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

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

1285 # after decrementing then avoid expand_path call. 

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

1287 continue 

1288 out |= set( 

1289 self.expand_path( 

1290 list(bit), 

1291 recursive=recursive, 

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

1293 assume_literal=True, 

1294 **kwargs, 

1295 ) 

1296 ) 

1297 continue 

1298 elif recursive: 

1299 rec = set( 

1300 self.find( 

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

1302 ) 

1303 ) 

1304 out |= rec 

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

1306 # should only check once, for the root 

1307 out.add(p) 

1308 if not out: 

1309 raise FileNotFoundError(path) 

1310 return sorted(out) 

1311 

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

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

1314 if path1 == path2: 

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

1316 else: 

1317 # explicitly raise exception to prevent data corruption 

1318 self.copy( 

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

1320 ) 

1321 self.rm(path1, recursive=recursive) 

1322 

1323 def rm_file(self, path): 

1324 """Delete a file""" 

1325 self._rm(path) 

1326 

1327 def _rm(self, path): 

1328 """Delete one file""" 

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

1330 raise NotImplementedError 

1331 

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

1333 """Delete files or directories. 

1334 

1335 Parameters 

1336 ---------- 

1337 path: str or list of str 

1338 Files or directories to delete. 

1339 recursive: bool 

1340 If True, recursively delete directories and their contents. 

1341 maxdepth: int or None 

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

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

1344 possible. 

1345 """ 

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

1347 for p in reversed(path): 

1348 self.rm_file(p) 

1349 

1350 @classmethod 

1351 def _parent(cls, path): 

1352 path = cls._strip_protocol(path) 

1353 if "/" in path: 

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

1355 return cls.root_marker + parent 

1356 else: 

1357 return cls.root_marker 

1358 

1359 def _open( 

1360 self, 

1361 path, 

1362 mode="rb", 

1363 block_size=None, 

1364 autocommit=True, 

1365 cache_options=None, 

1366 **kwargs, 

1367 ): 

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

1369 return AbstractBufferedFile( 

1370 self, 

1371 path, 

1372 mode, 

1373 block_size, 

1374 autocommit, 

1375 cache_options=cache_options, 

1376 **kwargs, 

1377 ) 

1378 

1379 def open( 

1380 self, 

1381 path, 

1382 mode="rb", 

1383 block_size=None, 

1384 cache_options=None, 

1385 compression=None, 

1386 **kwargs, 

1387 ): 

1388 """ 

1389 Return a file-like object from the filesystem 

1390 

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

1392 block. 

1393 

1394 Parameters 

1395 ---------- 

1396 path: str 

1397 Target file 

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

1399 See builtin ``open()`` 

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

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

1402 atomic is implementation-dependent. 

1403 block_size: int 

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

1405 cache_options : dict, optional 

1406 Extra arguments to pass through to the cache. 

1407 compression: string or None 

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

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

1410 compression from the filename suffix. 

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

1412 """ 

1413 import io 

1414 

1415 path = self._strip_protocol(path) 

1416 if "b" not in mode: 

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

1418 

1419 text_kwargs = { 

1420 k: kwargs.pop(k) 

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

1422 if k in kwargs 

1423 } 

1424 return io.TextIOWrapper( 

1425 self.open( 

1426 path, 

1427 mode, 

1428 block_size=block_size, 

1429 cache_options=cache_options, 

1430 compression=compression, 

1431 **kwargs, 

1432 ), 

1433 **text_kwargs, 

1434 ) 

1435 else: 

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

1437 f = self._open( 

1438 path, 

1439 mode=mode, 

1440 block_size=block_size, 

1441 autocommit=ac, 

1442 cache_options=cache_options, 

1443 **kwargs, 

1444 ) 

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

1446 self.transaction.files.append(f) 

1447 if compression is not None: 

1448 from fsspec.compression import compr 

1449 from fsspec.core import get_compression 

1450 

1451 compression = get_compression(path, compression) 

1452 compress = compr[compression] 

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

1454 return f 

1455 

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

1457 """Create empty file, or update timestamp 

1458 

1459 Parameters 

1460 ---------- 

1461 path: str 

1462 file location 

1463 truncate: bool 

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

1465 leave file unchanged, if backend allows this 

1466 """ 

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

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

1469 pass 

1470 else: 

1471 raise NotImplementedError # update timestamp, if possible 

1472 

1473 def ukey(self, path): 

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

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

1476 

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

1478 """Read a block of bytes from 

1479 

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

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

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

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

1484 bytestring returned WILL include the end delimiter string. 

1485 

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

1487 

1488 Parameters 

1489 ---------- 

1490 fn: string 

1491 Path to filename 

1492 offset: int 

1493 Byte offset to start read 

1494 length: int 

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

1496 delimiter: bytes (optional) 

1497 Ensure reading starts and stops at delimiter bytestring 

1498 

1499 Examples 

1500 -------- 

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

1502 b'Alice, 100\\nBo' 

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

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

1505 

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

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

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

1509 

1510 See Also 

1511 -------- 

1512 :func:`fsspec.utils.read_block` 

1513 """ 

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

1515 # Only fsspec's own file classes expose `.size`; a plain file (e.g. 

1516 # from the caching filesystems) does not, so fall back to a seek. 

1517 size = getattr(f, "size", None) 

1518 if size is None: 

1519 size = f.seek(0, 2) 

1520 f.seek(0) 

1521 if length is None: 

1522 length = size 

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

1524 length = size - offset 

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

1526 

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

1528 """ 

1529 JSON representation of this filesystem instance. 

1530 

1531 Parameters 

1532 ---------- 

1533 include_password: bool, default True 

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

1535 

1536 Returns 

1537 ------- 

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

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

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

1541 keyword arguments as their own keys. 

1542 

1543 Warnings 

1544 -------- 

1545 Serialized filesystems may contain sensitive information which have been 

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

1547 store and send them in a secure environment! 

1548 """ 

1549 from .json import FilesystemJSONEncoder 

1550 

1551 return json.dumps( 

1552 self, 

1553 cls=type( 

1554 "_FilesystemJSONEncoder", 

1555 (FilesystemJSONEncoder,), 

1556 {"include_password": include_password}, 

1557 ), 

1558 ) 

1559 

1560 @staticmethod 

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

1562 """ 

1563 Recreate a filesystem instance from JSON representation. 

1564 

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

1566 

1567 Parameters 

1568 ---------- 

1569 blob: str 

1570 

1571 Returns 

1572 ------- 

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

1574 

1575 Warnings 

1576 -------- 

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

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

1579 at import time. 

1580 """ 

1581 from .json import FilesystemJSONDecoder 

1582 

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

1584 

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

1586 """ 

1587 JSON-serializable dictionary representation of this filesystem instance. 

1588 

1589 Parameters 

1590 ---------- 

1591 include_password: bool, default True 

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

1593 

1594 Returns 

1595 ------- 

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

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

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

1599 keyword arguments as their own keys. 

1600 

1601 Warnings 

1602 -------- 

1603 Serialized filesystems may contain sensitive information which have been 

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

1605 store and send them in a secure environment! 

1606 """ 

1607 from .json import FilesystemJSONEncoder 

1608 

1609 json_encoder = FilesystemJSONEncoder() 

1610 

1611 cls = type(self) 

1612 proto = self.protocol 

1613 

1614 storage_options = dict(self.storage_options) 

1615 if not include_password: 

1616 storage_options.pop("password", None) 

1617 

1618 return dict( 

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

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

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

1622 **json_encoder.make_serializable(storage_options), 

1623 ) 

1624 

1625 @staticmethod 

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

1627 """ 

1628 Recreate a filesystem instance from dictionary representation. 

1629 

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

1631 

1632 Parameters 

1633 ---------- 

1634 dct: Dict[str, Any] 

1635 

1636 Returns 

1637 ------- 

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

1639 

1640 Warnings 

1641 -------- 

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

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

1644 at import time. 

1645 """ 

1646 from .json import FilesystemJSONDecoder 

1647 

1648 json_decoder = FilesystemJSONDecoder() 

1649 

1650 dct = dict(dct) # Defensive copy 

1651 

1652 cls = FilesystemJSONDecoder.try_resolve_fs_cls(dct) 

1653 if cls is None: 

1654 raise ValueError("Not a serialized AbstractFileSystem") 

1655 

1656 dct.pop("cls", None) 

1657 dct.pop("protocol", None) 

1658 

1659 return cls( 

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

1661 **json_decoder.unmake_serializable(dct), 

1662 ) 

1663 

1664 def _get_pyarrow_filesystem(self): 

1665 """ 

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

1667 """ 

1668 # all instances already also derive from pyarrow 

1669 return self 

1670 

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

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

1673 

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

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

1676 """ 

1677 from .mapping import FSMap 

1678 

1679 return FSMap( 

1680 root, 

1681 self, 

1682 check=check, 

1683 create=create, 

1684 missing_exceptions=missing_exceptions, 

1685 ) 

1686 

1687 @classmethod 

1688 def clear_instance_cache(cls): 

1689 """ 

1690 Clear the cache of filesystem instances. 

1691 

1692 Notes 

1693 ----- 

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

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

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

1697 since the instances refcount will not drop to zero until 

1698 ``clear_instance_cache`` is called. 

1699 """ 

1700 cls._cache.clear() 

1701 

1702 def created(self, path): 

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

1704 raise NotImplementedError 

1705 

1706 def modified(self, path): 

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

1708 raise NotImplementedError 

1709 

1710 def tree( 

1711 self, 

1712 path: str = "/", 

1713 recursion_limit: int = 2, 

1714 max_display: int = 25, 

1715 display_size: bool = False, 

1716 prefix: str = "", 

1717 is_last: bool = True, 

1718 first: bool = True, 

1719 indent_size: int = 4, 

1720 ) -> str: 

1721 """ 

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

1723 

1724 Parameters 

1725 ---------- 

1726 path: Root path to start traversal from 

1727 recursion_limit: Maximum depth of directory traversal 

1728 max_display: Maximum number of items to display per directory 

1729 display_size: Whether to display file sizes 

1730 prefix: Current line prefix for visual tree structure 

1731 is_last: Whether current item is last in its level 

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

1733 indent_size: Number of spaces by indent 

1734 

1735 Returns 

1736 ------- 

1737 str: A string representing the tree structure. 

1738 

1739 Example 

1740 ------- 

1741 >>> from fsspec import filesystem 

1742 

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

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

1745 >>> print(tree) 

1746 """ 

1747 

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

1749 """Format bytes as text.""" 

1750 for prefix, k in ( 

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

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

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

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

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

1756 ): 

1757 if n >= 0.9 * k: 

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

1759 return f"{n}B" 

1760 

1761 result = [] 

1762 

1763 if first: 

1764 result.append(path) 

1765 

1766 if recursion_limit: 

1767 indent = " " * indent_size 

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

1769 contents.sort( 

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

1771 ) 

1772 

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

1774 displayed_contents = contents[:max_display] 

1775 remaining_count = len(contents) - max_display 

1776 else: 

1777 displayed_contents = contents 

1778 remaining_count = 0 

1779 

1780 for i, item in enumerate(displayed_contents): 

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

1782 remaining_count == 0 

1783 ) 

1784 

1785 branch = ( 

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

1787 if is_last_item 

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

1789 ) 

1790 branch += " " 

1791 new_prefix = prefix + ( 

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

1793 ) 

1794 

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

1796 

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

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

1799 num_files = sum( 

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

1801 ) 

1802 num_folders = sum( 

1803 1 

1804 for sub_item in sub_contents 

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

1806 ) 

1807 

1808 if num_files == 0 and num_folders == 0: 

1809 size = " (empty folder)" 

1810 elif num_files == 0: 

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

1812 elif num_folders == 0: 

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

1814 else: 

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

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

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

1818 else: 

1819 size = "" 

1820 

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

1822 

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

1824 result.append( 

1825 self.tree( 

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

1827 recursion_limit=recursion_limit - 1, 

1828 max_display=max_display, 

1829 display_size=display_size, 

1830 prefix=new_prefix, 

1831 is_last=is_last_item, 

1832 first=False, 

1833 indent_size=indent_size, 

1834 ) 

1835 ) 

1836 

1837 if remaining_count > 0: 

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

1839 result.append( 

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

1841 ) 

1842 

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

1844 

1845 # ------------------------------------------------------------------------ 

1846 # Aliases 

1847 

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

1849 """Alias of `AbstractFileSystem.cat_file`.""" 

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

1851 

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

1853 """Alias of `AbstractFileSystem.pipe_file`.""" 

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

1855 

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

1857 """Alias of `AbstractFileSystem.mkdir`.""" 

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

1859 

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

1861 """Alias of `AbstractFileSystem.makedirs`.""" 

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

1863 

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

1865 """Alias of `AbstractFileSystem.ls`.""" 

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

1867 

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

1869 """Alias of `AbstractFileSystem.copy`.""" 

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

1871 

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

1873 """Alias of `AbstractFileSystem.mv`.""" 

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

1875 

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

1877 """Alias of `AbstractFileSystem.info`.""" 

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

1879 

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

1881 """Alias of `AbstractFileSystem.du`.""" 

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

1883 

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

1885 """Alias of `AbstractFileSystem.mv`.""" 

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

1887 

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

1889 """Alias of `AbstractFileSystem.rm`.""" 

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

1891 

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

1893 """Alias of `AbstractFileSystem.put`.""" 

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

1895 

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

1897 """Alias of `AbstractFileSystem.get`.""" 

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

1899 

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

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

1902 

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

1904 way of delegating credentials. 

1905 

1906 Parameters 

1907 ---------- 

1908 path : str 

1909 The path on the filesystem 

1910 expiration : int 

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

1912 

1913 Returns 

1914 ------- 

1915 URL : str 

1916 The signed URL 

1917 

1918 Raises 

1919 ------ 

1920 NotImplementedError : if method is not implemented for a filesystem 

1921 """ 

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

1923 

1924 def _isfilestore(self): 

1925 # Originally inherited from pyarrow DaskFileSystem. Keeping this 

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

1927 # legacy fsspec-compatible filesystems and thus accepts fsspec 

1928 # filesystems as well 

1929 return False 

1930 

1931 

1932class AbstractBufferedFile(io.IOBase): 

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

1934 

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

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

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

1938 ``_initiate_upload`` and ``_fetch_range``. 

1939 """ 

1940 

1941 DEFAULT_BLOCK_SIZE = 5 * 2**20 

1942 _details = None 

1943 

1944 def __init__( 

1945 self, 

1946 fs, 

1947 path, 

1948 mode="rb", 

1949 block_size="default", 

1950 autocommit=True, 

1951 cache_type="readahead", 

1952 cache_options=None, 

1953 size=None, 

1954 **kwargs, 

1955 ): 

1956 """ 

1957 Template for files with buffered reading and writing 

1958 

1959 Parameters 

1960 ---------- 

1961 fs: instance of FileSystem 

1962 path: str 

1963 location in file-system 

1964 mode: str 

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

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

1967 block_size: int 

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

1969 autocommit: bool 

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

1971 happens when file is being closed. 

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

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

1974 cache_options : dict 

1975 Additional options passed to the constructor for the cache specified 

1976 by `cache_type`. 

1977 size: int 

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

1979 kwargs: 

1980 Gets stored as self.kwargs 

1981 """ 

1982 from .core import caches 

1983 

1984 self.path = path 

1985 self.fs = fs 

1986 self.mode = mode 

1987 self.blocksize = ( 

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

1989 ) 

1990 self.loc = 0 

1991 self.autocommit = autocommit 

1992 self.end = None 

1993 self.start = None 

1994 self.closed = False 

1995 

1996 if cache_options is None: 

1997 cache_options = {} 

1998 

1999 if "trim" in kwargs: 

2000 warnings.warn( 

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

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

2003 FutureWarning, 

2004 ) 

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

2006 

2007 self.kwargs = kwargs 

2008 

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

2010 raise NotImplementedError("File mode not supported") 

2011 if mode == "rb": 

2012 if size is not None: 

2013 self.size = size 

2014 else: 

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

2016 self.cache = caches[cache_type]( 

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

2018 ) 

2019 else: 

2020 self.buffer = io.BytesIO() 

2021 self.offset = None 

2022 self.forced = False 

2023 self.location = None 

2024 

2025 @property 

2026 def details(self): 

2027 if self._details is None: 

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

2029 return self._details 

2030 

2031 @details.setter 

2032 def details(self, value): 

2033 self._details = value 

2034 self.size = value["size"] 

2035 

2036 @property 

2037 def full_name(self): 

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

2039 

2040 @property 

2041 def closed(self): 

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

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

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

2045 

2046 @closed.setter 

2047 def closed(self, c): 

2048 self._closed = c 

2049 

2050 def __hash__(self): 

2051 if "w" in self.mode: 

2052 return id(self) 

2053 else: 

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

2055 

2056 def __eq__(self, other): 

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

2058 if self is other: 

2059 return True 

2060 return ( 

2061 isinstance(other, type(self)) 

2062 and self.mode == "rb" 

2063 and other.mode == "rb" 

2064 and hash(self) == hash(other) 

2065 ) 

2066 

2067 def commit(self): 

2068 """Move from temp to final destination""" 

2069 

2070 def discard(self): 

2071 """Throw away temporary file""" 

2072 

2073 def info(self): 

2074 """File information about this path""" 

2075 if self.readable(): 

2076 return self.details 

2077 else: 

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

2079 

2080 def tell(self): 

2081 """Current file location""" 

2082 return self.loc 

2083 

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

2085 """Set current file location 

2086 

2087 Parameters 

2088 ---------- 

2089 loc: int 

2090 byte location 

2091 whence: {0, 1, 2} 

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

2093 """ 

2094 loc = int(loc) 

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

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

2097 if whence == 0: 

2098 nloc = loc 

2099 elif whence == 1: 

2100 nloc = self.loc + loc 

2101 elif whence == 2: 

2102 nloc = self.size + loc 

2103 else: 

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

2105 if nloc < 0: 

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

2107 self.loc = nloc 

2108 return self.loc 

2109 

2110 def write(self, data): 

2111 """ 

2112 Write data to buffer. 

2113 

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

2115 or equal to blocksize. 

2116 

2117 Parameters 

2118 ---------- 

2119 data: bytes 

2120 Set of bytes to be written. 

2121 """ 

2122 if not self.writable(): 

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

2124 if self.closed: 

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

2126 if self.forced: 

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

2128 out = self.buffer.write(data) 

2129 self.loc += out 

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

2131 self.flush() 

2132 return out 

2133 

2134 def flush(self, force=False): 

2135 """ 

2136 Write buffered data to backend store. 

2137 

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

2139 the file is being closed. 

2140 

2141 Parameters 

2142 ---------- 

2143 force: bool 

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

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

2146 """ 

2147 

2148 if self.closed: 

2149 raise ValueError("Flush on closed file") 

2150 if force and self.forced: 

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

2152 if force: 

2153 self.forced = True 

2154 

2155 if self.readable(): 

2156 # no-op to flush on read-mode 

2157 return 

2158 

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

2160 # Defer write on small block 

2161 return 

2162 

2163 if self.offset is None: 

2164 # Initialize a multipart upload 

2165 self.offset = 0 

2166 try: 

2167 self._initiate_upload() 

2168 except Exception: 

2169 self.closed = True 

2170 raise 

2171 

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

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

2174 self.buffer = io.BytesIO() 

2175 

2176 def _upload_chunk(self, final=False): 

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

2178 

2179 Parameters 

2180 ========== 

2181 final: bool 

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

2183 self.autocommit is True. 

2184 """ 

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

2186 

2187 def _initiate_upload(self): 

2188 """Create remote file/upload""" 

2189 pass 

2190 

2191 def _fetch_range(self, start, end): 

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

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

2194 

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

2196 """ 

2197 Return data from cache, or fetch pieces as necessary 

2198 

2199 Parameters 

2200 ---------- 

2201 length: int (-1) 

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

2203 """ 

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

2205 if self.mode != "rb": 

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

2207 if length < 0: 

2208 length = self.size - self.loc 

2209 if self.closed: 

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

2211 if length == 0: 

2212 # don't even bother calling fetch 

2213 return b"" 

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

2215 

2216 logger.debug( 

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

2218 self, 

2219 self.loc, 

2220 self.loc + length, 

2221 self.cache._log_stats(), 

2222 ) 

2223 self.loc += len(out) 

2224 return out 

2225 

2226 def readinto(self, b): 

2227 """mirrors builtin file's readinto method 

2228 

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

2230 """ 

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

2232 data = self.read(out.nbytes) 

2233 out[: len(data)] = data 

2234 return len(data) 

2235 

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

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

2238 

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

2240 encountered first. 

2241 

2242 Parameters 

2243 ---------- 

2244 char: bytes 

2245 Thing to find 

2246 blocks: None or int 

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

2248 mean a new read on every call. 

2249 """ 

2250 out = [] 

2251 while True: 

2252 start = self.tell() 

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

2254 if len(part) == 0: 

2255 break 

2256 found = part.find(char) 

2257 if found > -1: 

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

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

2260 break 

2261 out.append(part) 

2262 return b"".join(out) 

2263 

2264 def readline(self): 

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

2266 

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

2268 true line ending. 

2269 """ 

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

2271 

2272 def __next__(self): 

2273 out = self.readline() 

2274 if out: 

2275 return out 

2276 raise StopIteration 

2277 

2278 def __iter__(self): 

2279 return self 

2280 

2281 def readlines(self): 

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

2283 data = self.read() 

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

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

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

2287 return out 

2288 else: 

2289 return out + [lines[-1]] 

2290 # return list(self) ??? 

2291 

2292 def readinto1(self, b): 

2293 return self.readinto(b) 

2294 

2295 def close(self): 

2296 """Close file 

2297 

2298 Finalizes writes, discards cache 

2299 """ 

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

2301 return 

2302 if self.closed: 

2303 return 

2304 try: 

2305 if self.mode == "rb": 

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

2307 if cache is not None: 

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

2309 if callable(close): 

2310 close() 

2311 self.cache = None 

2312 else: 

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

2314 self.flush(force=True) 

2315 

2316 if self.fs is not None: 

2317 self.fs.invalidate_cache(self.path) 

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

2319 finally: 

2320 self.closed = True 

2321 

2322 def readable(self): 

2323 """Whether opened for reading""" 

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

2325 

2326 def seekable(self): 

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

2328 return self.readable() 

2329 

2330 def writable(self): 

2331 """Whether opened for writing""" 

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

2333 

2334 def __reduce__(self): 

2335 if self.mode != "rb": 

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

2337 

2338 return reopen, ( 

2339 self.fs, 

2340 self.path, 

2341 self.mode, 

2342 self.blocksize, 

2343 self.loc, 

2344 self.size, 

2345 self.autocommit, 

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

2347 self.kwargs, 

2348 ) 

2349 

2350 def __del__(self): 

2351 if not self.closed: 

2352 self.close() 

2353 

2354 def __str__(self): 

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

2356 

2357 __repr__ = __str__ 

2358 

2359 def __enter__(self): 

2360 return self 

2361 

2362 def __exit__(self, *args): 

2363 self.close() 

2364 

2365 

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

2367 file = fs.open( 

2368 path, 

2369 mode=mode, 

2370 block_size=blocksize, 

2371 autocommit=autocommit, 

2372 cache_type=cache_type, 

2373 size=size, 

2374 **kwargs, 

2375 ) 

2376 if loc > 0: 

2377 file.seek(loc) 

2378 return file