-
Notifications
You must be signed in to change notification settings - Fork 1.2k
/
_callback.py
225 lines (179 loc) · 6.1 KB
/
_callback.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
from contextlib import ExitStack
from functools import wraps
from typing import IO, TYPE_CHECKING, Any, Dict, Optional, TypeVar, cast
import fsspec
from funcy import cached_property
if TYPE_CHECKING:
from typing import Callable
from typing_extensions import ParamSpec
from dvc.progress import Tqdm
from dvc.ui._rich_progress import RichTransferProgress
_P = ParamSpec("_P")
_R = TypeVar("_R")
class FsspecCallback(fsspec.Callback):
"""FsspecCallback usable as a context manager, and a few helper methods."""
def wrap_attr(self, fobj: IO, method: str = "read") -> IO:
from tqdm.utils import CallbackIOWrapper
wrapped = CallbackIOWrapper(self.relative_update, fobj, method)
return cast(IO, wrapped)
def wrap_fn(self, fn: "Callable[_P, _R]") -> "Callable[_P, _R]":
@wraps(fn)
def wrapped(*args: "_P.args", **kwargs: "_P.kwargs") -> "_R":
res = fn(*args, **kwargs)
self.relative_update()
return res
return wrapped
def wrap_and_branch(self, fn: "Callable") -> "Callable":
"""
Wraps a function, and pass a new child callback to it.
When the function completes, we increment the parent callback by 1.
"""
wrapped = self.wrap_fn(fn)
@wraps(fn)
def func(path1: str, path2: str):
kw: Dict[str, Any] = {}
with self.branch(path1, path2, kw):
return wrapped(path1, path2, **kw)
return func
def __enter__(self):
return self
def __exit__(self, *exc_args):
self.close()
def close(self):
"""Handle here on exit."""
def relative_update(self, inc: int = 1) -> None:
inc = inc if inc is not None else 0
return super().relative_update(inc)
def absolute_update(self, value: int) -> None:
value = value if value is not None else self.value
return super().absolute_update(value)
@classmethod
def as_callback(
cls, maybe_callback: Optional["FsspecCallback"] = None
) -> "FsspecCallback":
if maybe_callback is None:
return DEFAULT_CALLBACK
return maybe_callback
@classmethod
def as_tqdm_callback(
cls,
callback: Optional["FsspecCallback"] = None,
**tqdm_kwargs: Any,
) -> "FsspecCallback":
return callback or TqdmCallback(**tqdm_kwargs)
@classmethod
def as_rich_callback(
cls, callback: Optional["FsspecCallback"] = None, **rich_kwargs
):
return callback or RichCallback(**rich_kwargs)
def branch(
self,
path_1: str,
path_2: str,
kwargs: Dict[str, Any],
child: "FsspecCallback" = None,
) -> "FsspecCallback":
child = kwargs["callback"] = child or DEFAULT_CALLBACK
return child
class NoOpCallback(FsspecCallback, fsspec.callbacks.NoOpCallback):
pass
class TqdmCallback(FsspecCallback):
def __init__(
self,
size: Optional[int] = None,
value: int = 0,
progress_bar: "Tqdm" = None,
**tqdm_kwargs,
):
tqdm_kwargs["total"] = size or -1
self._tqdm_kwargs = tqdm_kwargs
self._progress_bar = progress_bar
self._stack = ExitStack()
super().__init__(size=size, value=value)
@cached_property
def progress_bar(self):
from dvc.progress import Tqdm
progress_bar = (
self._progress_bar
if self._progress_bar is not None
else Tqdm(**self._tqdm_kwargs)
)
return self._stack.enter_context(progress_bar)
def __enter__(self):
return self
def close(self):
self._stack.close()
def set_size(self, size):
# Tqdm tries to be smart when to refresh,
# so we try to force it to re-render.
super().set_size(size)
self.progress_bar.refresh()
def call(self, hook_name=None, **kwargs):
self.progress_bar.update_to(self.value, total=self.size)
def branch(
self,
path_1: str,
path_2: str,
kwargs,
child: Optional[FsspecCallback] = None,
):
child = child or TqdmCallback(bytes=True, desc=path_1)
return super().branch(path_1, path_2, kwargs, child=child)
class RichCallback(FsspecCallback):
def __init__(
self,
size: Optional[int] = None,
value: int = 0,
progress: "RichTransferProgress" = None,
desc: str = None,
bytes: bool = False, # pylint: disable=redefined-builtin
unit: str = None,
disable: bool = False,
) -> None:
self._progress = progress
self.disable = disable
self._task_kwargs = {
"description": desc or "",
"bytes": bytes,
"unit": unit,
"total": size or 0,
"visible": False,
"progress_type": None if bytes else "summary",
}
self._stack = ExitStack()
super().__init__(size=size, value=value)
@cached_property
def progress(self):
from dvc.ui import ui
from dvc.ui._rich_progress import RichTransferProgress
if self._progress is not None:
return self._progress
progress = RichTransferProgress(
transient=True,
disable=self.disable,
console=ui.error_console,
)
return self._stack.enter_context(progress)
@cached_property
def task(self):
return self.progress.add_task(**self._task_kwargs)
def __enter__(self):
return self
def close(self):
self.progress.clear_task(self.task)
self._stack.close()
def call(self, hook_name=None, **kwargs):
self.progress.update(
self.task,
completed=self.value,
total=self.size,
visible=not self.disable,
)
def branch(
self, path_1, path_2, kwargs, child: Optional[FsspecCallback] = None
):
child = child or RichCallback(
progress=self.progress, desc=path_1, bytes=True
)
return super().branch(path_1, path_2, kwargs, child=child)
DEFAULT_CALLBACK = NoOpCallback()