codekingpro/portable-devtools
114k
1from __future__ import annotations2 3import inspect4import logging5import os6import shutil7import uuid8 9from .asyn import AsyncFileSystem, _run_coros_in_chunks, sync_wrapper10from .callbacks import DEFAULT_CALLBACK11from .core import filesystem, get_filesystem_class, split_protocol, url_to_fs12 13_generic_fs = {}14logger = logging.getLogger("fsspec.generic")15 16 17def set_generic_fs(protocol, **storage_options):18 """Populate the dict used for method=="generic" lookups"""19 _generic_fs[protocol] = filesystem(protocol, **storage_options)20 21 22def _resolve_fs(url, method, protocol=None, storage_options=None):23 """Pick instance of backend FS"""24 url = url[0] if isinstance(url, (list, tuple)) else url25 protocol = protocol or split_protocol(url)[0]26 storage_options = storage_options or {}27 if method == "default":28 return filesystem(protocol)29 if method == "generic":30 return _generic_fs[protocol]31 if method == "current":32 cls = get_filesystem_class(protocol)33 return cls.current()34 if method == "options":35 fs, _ = url_to_fs(url, **storage_options.get(protocol, {}))36 return fs37 raise ValueError(f"Unknown FS resolution method: {method}")38 39 40def rsync(41 source,42 destination,43 delete_missing=False,44 source_field="size",45 dest_field="size",46 update_cond="different",47 inst_kwargs=None,48 fs=None,49 **kwargs,50):51 """Sync files between two directory trees52 53 (experimental)54 55 Parameters56 ----------57 source: str58 Root of the directory tree to take files from. This must be a directory, but59 do not include any terminating "/" character60 destination: str61 Root path to copy into. The contents of this location should be62 identical to the contents of ``source`` when done. This will be made a63 directory, and the terminal "/" should not be included.64 delete_missing: bool65 If there are paths in the destination that don't exist in the66 source and this is True, delete them. Otherwise, leave them alone.67 source_field: str | callable68 If ``update_field`` is "different", this is the key in the info69 of source files to consider for difference. Maybe a function of the70 info dict.71 dest_field: str | callable72 If ``update_field`` is "different", this is the key in the info73 of destination files to consider for difference. May be a function of74 the info dict.75 update_cond: "different"|"always"|"never"76 If "always", every file is copied, regardless of whether it exists in77 the destination. If "never", files that exist in the destination are78 not copied again. If "different" (default), only copy if the info79 fields given by ``source_field`` and ``dest_field`` (usually "size")80 are different. Other comparisons may be added in the future.81 inst_kwargs: dict|None82 If ``fs`` is None, use this set of keyword arguments to make a83 GenericFileSystem instance84 fs: GenericFileSystem|None85 Instance to use if explicitly given. The instance defines how to86 to make downstream file system instances from paths.87 88 Returns89 -------90 dict of the copy operations that were performed, {source: destination}91 """92 fs = fs or GenericFileSystem(**(inst_kwargs or {}))93 source = fs._strip_protocol(source)94 destination = fs._strip_protocol(destination)95 allfiles = fs.find(source, withdirs=True, detail=True)96 if not fs.isdir(source):97 raise ValueError("Can only rsync on a directory")98 otherfiles = fs.find(destination, withdirs=True, detail=True)99 dirs = [100 a101 for a, v in allfiles.items()102 if v["type"] == "directory" and a.replace(source, destination) not in otherfiles103 ]104 logger.debug(f"{len(dirs)} directories to create")105 if dirs:106 fs.make_many_dirs(107 [dirn.replace(source, destination) for dirn in dirs], exist_ok=True108 )109 allfiles = {a: v for a, v in allfiles.items() if v["type"] == "file"}110 logger.debug(f"{len(allfiles)} files to consider for copy")111 to_delete = [112 o113 for o, v in otherfiles.items()114 if o.replace(destination, source) not in allfiles and v["type"] == "file"115 ]116 for k, v in allfiles.copy().items():117 otherfile = k.replace(source, destination)118 if otherfile in otherfiles:119 if update_cond == "always":120 allfiles[k] = otherfile121 elif update_cond == "never":122 allfiles.pop(k)123 elif update_cond == "different":124 inf1 = source_field(v) if callable(source_field) else v[source_field]125 v2 = otherfiles[otherfile]126 inf2 = dest_field(v2) if callable(dest_field) else v2[dest_field]127 if inf1 != inf2:128 # details mismatch, make copy129 allfiles[k] = otherfile130 else:131 # details match, don't copy132 allfiles.pop(k)133 else:134 # file not in target yet135 allfiles[k] = otherfile136 logger.debug(f"{len(allfiles)} files to copy")137 if allfiles:138 source_files, target_files = zip(*allfiles.items())139 fs.cp(source_files, target_files, **kwargs)140 logger.debug(f"{len(to_delete)} files to delete")141 if delete_missing and to_delete:142 fs.rm(to_delete)143 return allfiles144 145 146class GenericFileSystem(AsyncFileSystem):147 """Wrapper over all other FS types148 149 <experimental!>150 151 This implementation is a single unified interface to be able to run FS operations152 over generic URLs, and dispatch to the specific implementations using the URL153 protocol prefix.154 155 Note: instances of this FS are always async, even if you never use it with any async156 backend.157 """158 159 protocol = "generic" # there is no real reason to ever use a protocol with this FS160 161 def __init__(self, default_method="default", storage_options=None, **kwargs):162 """163 164 Parameters165 ----------166 default_method: str (optional)167 Defines how to configure backend FS instances. Options are:168 - "default": instantiate like FSClass(), with no169 extra arguments; this is the default instance of that FS, and can be170 configured via the config system171 - "generic": takes instances from the `_generic_fs` dict in this module,172 which you must populate before use. Keys are by protocol173 - "options": expects storage_options, a dict mapping protocol to174 kwargs to use when constructing the filesystem175 - "current": takes the most recently instantiated version of each FS176 """177 self.method = default_method178 self.st_opts = storage_options179 super().__init__(**kwargs)180 181 def _parent(self, path):182 fs = _resolve_fs(path, self.method, storage_options=self.st_opts)183 return fs.unstrip_protocol(fs._parent(path))184 185 def _strip_protocol(self, path):186 # normalization only187 fs = _resolve_fs(path, self.method, storage_options=self.st_opts)188 return fs.unstrip_protocol(fs._strip_protocol(path))189 190 async def _find(self, path, maxdepth=None, withdirs=False, detail=False, **kwargs):191 fs = _resolve_fs(path, self.method, storage_options=self.st_opts)192 if fs.async_impl:193 out = await fs._find(194 path, maxdepth=maxdepth, withdirs=withdirs, detail=True, **kwargs195 )196 else:197 out = fs.find(198 path, maxdepth=maxdepth, withdirs=withdirs, detail=True, **kwargs199 )200 result = {}201 for k, v in out.items():202 v = v.copy() # don't corrupt target FS dircache203 name = fs.unstrip_protocol(k)204 v["name"] = name205 result[name] = v206 if detail:207 return result208 return list(result)209 210 async def _info(self, url, **kwargs):211 fs = _resolve_fs(url, self.method)212 if fs.async_impl:213 out = await fs._info(url, **kwargs)214 else:215 out = fs.info(url, **kwargs)216 out = out.copy() # don't edit originals217 out["name"] = fs.unstrip_protocol(out["name"])218 return out219 220 async def _ls(221 self,222 url,223 detail=True,224 **kwargs,225 ):226 fs = _resolve_fs(url, self.method)227 if fs.async_impl:228 out = await fs._ls(url, detail=True, **kwargs)229 else:230 out = fs.ls(url, detail=True, **kwargs)231 out = [o.copy() for o in out] # don't edit originals232 for o in out:233 o["name"] = fs.unstrip_protocol(o["name"])234 if detail:235 return out236 else:237 return [o["name"] for o in out]238 239 async def _cat_file(240 self,241 url,242 **kwargs,243 ):244 fs = _resolve_fs(url, self.method)245 if fs.async_impl:246 return await fs._cat_file(url, **kwargs)247 else:248 return fs.cat_file(url, **kwargs)249 250 async def _pipe_file(251 self,252 path,253 value,254 **kwargs,255 ):256 fs = _resolve_fs(path, self.method, storage_options=self.st_opts)257 if fs.async_impl:258 return await fs._pipe_file(path, value, **kwargs)259 else:260 return fs.pipe_file(path, value, **kwargs)261 262 async def _rm(self, url, **kwargs):263 urls = url264 if isinstance(urls, str):265 urls = [urls]266 fs = _resolve_fs(urls[0], self.method)267 if fs.async_impl:268 await fs._rm(urls, **kwargs)269 else:270 fs.rm(url, **kwargs)271 272 async def _makedirs(self, path, exist_ok=False):273 logger.debug("Make dir %s", path)274 fs = _resolve_fs(path, self.method, storage_options=self.st_opts)275 if fs.async_impl:276 await fs._makedirs(path, exist_ok=exist_ok)277 else:278 fs.makedirs(path, exist_ok=exist_ok)279 280 def rsync(self, source, destination, **kwargs):281 """Sync files between two directory trees282 283 See `func:rsync` for more details.284 """285 rsync(source, destination, fs=self, **kwargs)286 287 async def _cp_file(288 self,289 url,290 url2,291 blocksize=2**20,292 callback=DEFAULT_CALLBACK,293 tempdir: str | None = None,294 **kwargs,295 ):296 fs = _resolve_fs(url, self.method)297 fs2 = _resolve_fs(url2, self.method)298 if fs is fs2:299 # pure remote300 if fs.async_impl:301 return await fs._copy(url, url2, **kwargs)302 else:303 return fs.copy(url, url2, **kwargs)304 await copy_file_op(fs, [url], fs2, [url2], tempdir, 1, on_error="raise")305 306 async def _make_many_dirs(self, urls, exist_ok=True):307 fs = _resolve_fs(urls[0], self.method)308 if fs.async_impl:309 coros = [fs._makedirs(u, exist_ok=exist_ok) for u in urls]310 await _run_coros_in_chunks(coros)311 else:312 for u in urls:313 fs.makedirs(u, exist_ok=exist_ok)314 315 make_many_dirs = sync_wrapper(_make_many_dirs)316 317 async def _copy(318 self,319 path1: list[str],320 path2: list[str],321 recursive: bool = False,322 on_error: str = "ignore",323 maxdepth: int | None = None,324 batch_size: int | None = None,325 tempdir: str | None = None,326 **kwargs,327 ):328 # TODO: special case for one FS being local, which can use get/put329 # TODO: special case for one being memFS, which can use cat/pipe330 if recursive:331 raise NotImplementedError("Please use fsspec.generic.rsync")332 path1 = [path1] if isinstance(path1, str) else path1333 path2 = [path2] if isinstance(path2, str) else path2334 335 fs = _resolve_fs(path1, self.method)336 fs2 = _resolve_fs(path2, self.method)337 338 if fs is fs2:339 if fs.async_impl:340 return await fs._copy(path1, path2, **kwargs)341 else:342 return fs.copy(path1, path2, **kwargs)343 344 await copy_file_op(345 fs, path1, fs2, path2, tempdir, batch_size, on_error=on_error346 )347 348 349async def copy_file_op(350 fs1, url1, fs2, url2, tempdir=None, batch_size=20, on_error="ignore"351):352 import tempfile353 354 tempdir = tempdir or tempfile.mkdtemp()355 try:356 coros = [357 _copy_file_op(358 fs1,359 u1,360 fs2,361 u2,362 os.path.join(tempdir, uuid.uuid4().hex),363 )364 for u1, u2 in zip(url1, url2)365 ]366 out = await _run_coros_in_chunks(367 coros, batch_size=batch_size, return_exceptions=True368 )369 finally:370 shutil.rmtree(tempdir)371 if on_error == "return":372 return out373 elif on_error == "raise":374 for o in out:375 if isinstance(o, Exception):376 raise o377 378 379async def _copy_file_op(fs1, url1, fs2, url2, local, on_error="ignore"):380 if fs1.async_impl:381 await fs1._get_file(url1, local)382 else:383 fs1.get_file(url1, local)384 if fs2.async_impl:385 await fs2._put_file(local, url2)386 else:387 fs2.put_file(local, url2)388 os.unlink(local)389 logger.debug("Copy %s -> %s; done", url1, url2)390 391 392async def maybe_await(cor):393 if inspect.iscoroutine(cor):394 return await cor395 else:396 return cor397 