-
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtransfer_queue.py
More file actions
305 lines (274 loc) · 11.9 KB
/
Copy pathtransfer_queue.py
File metadata and controls
305 lines (274 loc) · 11.9 KB
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
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
"""
Thread-safe in-memory FIFO transfer queue.
Pure logic, no paramiko or network dependency. A single lock guards every
state change so multiple workers can claim items concurrently without ever
grabbing the same one; claim() is what makes the worker pool in
simple_sftp_client.py (2 workers by default, up to 5 for a batch of many
small files) safe.
"""
import threading
from dataclasses import dataclass
# ───────────── states ─────────────
WAITING = "waiting"
ACTIVE = "active"
COMPLETED = "completed"
FAILED = "failed"
CANCELLED = "cancelled"
SKIPPED = "skipped"
ALL_STATES = (WAITING, ACTIVE, COMPLETED, FAILED, CANCELLED, SKIPPED)
TERMINAL_STATES = (COMPLETED, FAILED, CANCELLED, SKIPPED)
# Most-recent completed/skipped items kept as objects in self._items; once a
# batch has more than this many finished successfully, the oldest ones
# collapse into the self._pruned counters instead of staying in memory. This
# keeps a normal small transfer fully inspectable while still bounding memory
# for a huge job.
RETAIN_FINISHED = 200
@dataclass
class TransferItem:
id: int
direction: str # "upload" or "download"
local_path: str
remote_path: str
name: str
size: int = 0
state: str = WAITING
error: str = ""
cancel_requested: bool = False
on_conflict: str = "overwrite"
is_dir: bool = False
class TransferQueue:
"""FIFO queue of TransferItem, guarded by a single lock for atomic claims."""
def __init__(self):
self._lock = threading.Lock()
self._items = [] # FIFO order, append order preserved
self._by_id = {} # item.id -> item, kept in sync with self._items for O(1) _find
self._next_id = 1
self._paused = False
# Once more than RETAIN_FINISHED items have completed/skipped, the
# oldest ones are dropped from self._items (see _finish) so a huge
# finished batch doesn't keep every item object in memory and get
# copied on every UI poll. These tallies keep counts() accurate for
# the ones that aged out and no longer exist as objects.
self._pruned = {COMPLETED: 0, SKIPPED: 0}
# Running count of live (not yet pruned) COMPLETED+SKIPPED items in
# self._items, kept in sync in _finish and clear_finished so _finish
# never has to re-scan self._items to know whether it is over
# RETAIN_FINISHED.
self._finished_live = 0
def pause(self):
"""Stop claim() from handing out new items. Items already ACTIVE are
left alone and run to completion."""
with self._lock:
self._paused = True
def resume(self):
"""Let claim() hand out WAITING items again."""
with self._lock:
self._paused = False
def is_paused(self):
with self._lock:
return self._paused
def append(self, direction, local_path, remote_path, name, size=0, on_conflict="overwrite",
is_dir=False):
with self._lock:
item = TransferItem(
id=self._next_id,
direction=direction,
local_path=local_path,
remote_path=remote_path,
name=name,
size=size,
on_conflict=on_conflict,
is_dir=is_dir,
)
self._next_id += 1
self._items.append(item)
self._by_id[item.id] = item
return item.id
def claim(self):
"""Atomically find the oldest WAITING item, mark it ACTIVE, and return it."""
with self._lock:
# Workers retire on a None claim, so this is what makes the pool
# wind down while paused instead of picking up more WAITING work.
if self._paused:
return None
for item in self._items:
if item.state == WAITING:
item.state = ACTIVE
return item
return None
def _finish(self, item_id, new_state, error=""):
"""Move an ACTIVE item to a terminal state. No-op (returns False) otherwise.
COMPLETED and SKIPPED items stay visible in self._items up to
RETAIN_FINISHED of them; beyond that cap the oldest ones collapse
into the self._pruned counter instead of staying in memory forever,
which is what keeps a million-file job bounded. FAILED and CANCELLED
are never pruned, so they remain visible and retryable."""
with self._lock:
item = self._find(item_id)
if item is None or item.state != ACTIVE:
return False
item.state = new_state
if error:
item.error = error
if new_state in (COMPLETED, SKIPPED):
# Only one item can cross into COMPLETED/SKIPPED per call, so
# the count can be at most one over the cap: drop the single
# oldest finished item (FIFO, scan from the front) if so.
self._finished_live += 1
if self._finished_live > RETAIN_FINISHED:
for i in self._items:
if i.state in (COMPLETED, SKIPPED):
self._pruned[i.state] += 1
self._items.remove(i)
del self._by_id[i.id]
self._finished_live -= 1
break
return True
def mark_completed(self, item_id):
return self._finish(item_id, COMPLETED)
def mark_failed(self, item_id, error):
return self._finish(item_id, FAILED, error=error)
def mark_skipped(self, item_id):
return self._finish(item_id, SKIPPED)
def mark_cancelled(self, item_id):
"""Move an ACTIVE item to CANCELLED (terminal). False if not ACTIVE or unknown.
This is how the worker finalizes an active transfer that was cancelled
mid-flight; cancel() alone only flags an active item, it does not move it."""
return self._finish(item_id, CANCELLED)
def cancel(self, item_id):
"""WAITING items cancel immediately. ACTIVE items are flagged for the
worker to observe and finish via a terminal method. Returns False if
the item is unknown or already terminal."""
with self._lock:
item = self._find(item_id)
if item is None:
return False
if item.state == WAITING:
item.state = CANCELLED
return True
if item.state == ACTIVE:
item.cancel_requested = True
return True
return False
def requeue(self, item_id):
"""How a failed or user-cancelled item gets put back in line for the
worker pool: reset it to WAITING with a clean error and cancel flag.
False (no change) for any other state or an unknown id."""
with self._lock:
item = self._find(item_id)
if item is None or item.state not in (FAILED, CANCELLED):
return False
item.state = WAITING
item.error = ""
item.cancel_requested = False
return True
def retry_all_failed(self):
"""Put every FAILED item back to WAITING, clearing its error and
cancel flag, same as requeue() does for a single item. CANCELLED
items are left alone: those were intentional and keep the existing
per-item requeue(). Returns how many items were requeued."""
with self._lock:
requeued = 0
for item in self._items:
if item.state == FAILED:
item.state = WAITING
item.error = ""
item.cancel_requested = False
requeued += 1
return requeued
def cancel_all(self):
with self._lock:
for item in self._items:
if item.state == WAITING:
item.state = CANCELLED
elif item.state == ACTIVE:
item.cancel_requested = True
def fail_waiting(self, error):
"""Move every still-WAITING item to FAILED with the given reason, and
return how many were failed. Used when the worker pool cannot open any
transfer session: without this, queued items would sit as WAITING
forever with no worker left to drain them, showing as pending with no
visible failure. ACTIVE and terminal items are left untouched."""
with self._lock:
failed = 0
for item in self._items:
if item.state == WAITING:
item.state = FAILED
item.error = error
failed += 1
return failed
def snapshot(self):
"""Plain dicts (id, direction, name, state, error, is_dir) in FIFO order, copies only."""
with self._lock:
return [
{
"id": item.id,
"direction": item.direction,
"name": item.name,
"state": item.state,
"error": item.error,
"is_dir": item.is_dir,
}
for item in self._items
]
def counts(self):
with self._lock:
result = {state: 0 for state in ALL_STATES}
for item in self._items:
result[item.state] += 1
# Live COMPLETED/SKIPPED items are already counted from
# self._items above. self._pruned only covers the ones that
# aged out past RETAIN_FINISHED and no longer exist as objects,
# so adding it in here cannot double count.
result[COMPLETED] += self._pruned[COMPLETED]
result[SKIPPED] += self._pruned[SKIPPED]
return result
def pending(self):
with self._lock:
return sum(1 for item in self._items if item.state in (WAITING, ACTIVE))
def waiting(self):
"""Count of items still WAITING to be claimed (not counting ACTIVE).
Used by the worker pool to decide whether to start another worker or
let one retire; pending() includes ACTIVE items too, so it is not the
right check there (an active item on another worker would wrongly
keep an idle worker's decision looking like there is more to do)."""
with self._lock:
return sum(1 for item in self._items if item.state == WAITING)
def snapshot_and_pending(self):
"""Same data as snapshot() and pending(), read under one lock grab so
the two numbers can never disagree by one item mid-transfer, the way
two separate lock grabs could."""
with self._lock:
items = [
{
"id": item.id,
"direction": item.direction,
"name": item.name,
"state": item.state,
"error": item.error,
"is_dir": item.is_dir,
}
for item in self._items
]
pending = sum(1 for item in self._items if item.state in (WAITING, ACTIVE))
return items, pending
def clear_finished(self):
"""Remove items in a terminal state, keep waiting/active items. Also
resets the pruned tallies to zero so the next batch's summary starts
clean. Returns the number of item objects removed (the tally reset
is not counted)."""
with self._lock:
before = len(self._items)
kept = []
for item in self._items:
if item.state in TERMINAL_STATES:
del self._by_id[item.id]
else:
kept.append(item)
self._items = kept
self._pruned = {COMPLETED: 0, SKIPPED: 0}
self._finished_live = 0
return before - len(self._items)
def _find(self, item_id):
"""Caller must hold self._lock."""
return self._by_id.get(item_id)