Skip to content

Commit 02c34cf

Browse files
committed
feat(syncer): trees 重试 + 可取消 + 逐文件进度
- git/trees 关键调用套重试退避(_fetch_json),避免单次抖动整段失败 - 全链路加 cancel Event:阶段间 + as_completed 循环内即退, 取消时 shutdown(wait=False, cancel_futures=True) 不阻塞,已落盘 manifest 保留续传 - 接活死掉的 progress_cb:sync_maps 逐文件回调「地图 N/M」,sync_all 回调阶段名 UI:同步按钮运行中复用为取消入口(✕ 取消同步),进度进日志面板
1 parent 5be2364 commit 02c34cf

2 files changed

Lines changed: 180 additions & 30 deletions

File tree

aao/resources/syncer.py

Lines changed: 148 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -15,8 +15,11 @@
1515
from __future__ import annotations
1616

1717
import json
18+
import threading
19+
import time
1820
import urllib.parse
1921
import urllib.request
22+
from collections.abc import Callable
2023
from concurrent.futures import ThreadPoolExecutor, as_completed
2124
from dataclasses import dataclass
2225
from pathlib import Path
@@ -37,10 +40,16 @@ class SyncResult:
3740
failed: int = 0 # 显式失败条目数
3841
skipped: int = 0 # 已是最新被跳过的条目数
3942
error: str = "" # 致命错误描述(fatal 时填)
43+
cancelled: bool = False # 被外部取消
4044

4145
@property
4246
def message(self) -> str:
4347
"""供 UI 进度回调使用的人话汇报。"""
48+
if self.cancelled:
49+
base = f"{self.name}已取消"
50+
if self.done > 0:
51+
base += f"(已下载 {self.done} 条)"
52+
return base
4453
if self.error:
4554
return f"{self.name}更新失败:{self.error}"
4655
if self.failed == 0:
@@ -52,6 +61,11 @@ def message(self) -> str:
5261
return f"{self.name}部分失败(成功 {self.done},失败 {self.failed})"
5362

5463

64+
def _cancelled(cancel: threading.Event | None) -> bool:
65+
"""统一的取消信号检查(cancel=None 视为永不取消)。"""
66+
return cancel is not None and cancel.is_set()
67+
68+
5569
# 主源:GitHub dev-v2 分支(main 分支改名而来)
5670
_REPO = "MaaAssistantArknights/MaaAssistantArknights"
5771
_BRANCH = "dev-v2"
@@ -112,20 +126,68 @@ def _make_request(url: str, token: str | None = None) -> urllib.request.Request:
112126
return req
113127

114128

129+
_DOWNLOAD_RETRIES = 3
130+
_DOWNLOAD_BACKOFF_SEC = 0.5
131+
132+
133+
def _fetch_json(
134+
opener: urllib.request.OpenerDirector,
135+
url: str,
136+
*,
137+
token: str | None = None,
138+
timeout: int = 60,
139+
) -> dict[str, Any] | None:
140+
"""带重试的 GET + JSON 解析,用于关键 API 调用(git/trees)。
141+
142+
单文件下载用 _download(落盘);这里只拉一份 JSON 进内存。
143+
关键 API 抖动一次就整段同步失败不划算,套同样 3 次重试退避。
144+
"""
145+
last_err: Exception | None = None
146+
for attempt in range(1, _DOWNLOAD_RETRIES + 1):
147+
try:
148+
with opener.open(_make_request(url, token), timeout=timeout) as resp:
149+
return json.loads(resp.read())
150+
except Exception as e: # noqa: BLE001
151+
last_err = e
152+
if attempt < _DOWNLOAD_RETRIES:
153+
time.sleep(_DOWNLOAD_BACKOFF_SEC * attempt)
154+
continue
155+
logger.error("拉取失败(%d 次重试均失败) %s: %s", _DOWNLOAD_RETRIES, url, e)
156+
return None
157+
assert last_err is not None # 循环必 return,防御
158+
return None
159+
160+
115161
def _download(
116162
opener: urllib.request.OpenerDirector, url: str, dest: Path, token: str | None = None
117163
) -> bool:
118-
try:
119-
with opener.open(_make_request(url, token), timeout=30) as resp:
120-
data = resp.read()
121-
dest.parent.mkdir(parents=True, exist_ok=True)
122-
tmp = dest.with_suffix(dest.suffix + ".tmp")
123-
tmp.write_bytes(data)
124-
tmp.replace(dest)
125-
return True
126-
except Exception as e: # noqa: BLE001
127-
logger.warning("下载失败 %s: %s", url, e)
128-
return False
164+
"""带重试的单文件下载。
165+
166+
CDN 暂时抖动 / TLS timeout 不至于让一个文件永远漏掉:
167+
失败重试最多 3 次,间隔 500ms 退避。
168+
成功即刻返回;3 次全失败才返回 False,由上层记录 manifest 不动、下次再补。
169+
"""
170+
last_err: Exception | None = None
171+
for attempt in range(1, _DOWNLOAD_RETRIES + 1):
172+
try:
173+
with opener.open(_make_request(url, token), timeout=30) as resp:
174+
data = resp.read()
175+
dest.parent.mkdir(parents=True, exist_ok=True)
176+
tmp = dest.with_suffix(dest.suffix + ".tmp")
177+
tmp.write_bytes(data)
178+
tmp.replace(dest)
179+
return True
180+
except Exception as e: # noqa: BLE001
181+
last_err = e
182+
if attempt < _DOWNLOAD_RETRIES:
183+
time.sleep(_DOWNLOAD_BACKOFF_SEC * attempt)
184+
continue
185+
logger.warning("下载失败(%d 次重试均失败) %s: %s", _DOWNLOAD_RETRIES, url, e)
186+
return False
187+
# 理论不可达(循环必 return),但满足类型检查 + 防御。
188+
assert last_err is not None
189+
logger.warning("下载失败 %s: %s", url, last_err)
190+
return False
129191

130192

131193
def _list_remote_tiles(
@@ -143,11 +205,8 @@ def _list_remote_tiles(
143205
blob sha 与 Contents API sha 一致,可直接沿用同一份 manifest。
144206
"""
145207
logger.info("从 GitHub git/trees 拉取 dev-v2 完整树...")
146-
try:
147-
with opener.open(_make_request(_TREES_API, token), timeout=60) as resp:
148-
data = json.loads(resp.read())
149-
except Exception as e: # noqa: BLE001
150-
logger.error("拉取 git/trees 失败: %s", e)
208+
data = _fetch_json(opener, _TREES_API, token=token, timeout=60)
209+
if data is None:
151210
return None
152211

153212
if data.get("truncated"):
@@ -198,10 +257,13 @@ def _get_battle_data(
198257
tmp.unlink(missing_ok=True)
199258

200259

201-
def sync_operators(proxy: str | None = None) -> SyncResult:
260+
def sync_operators(proxy: str | None = None, cancel: threading.Event | None = None) -> SyncResult:
202261
data_dir = project_root() / "data"
203262
data_dir.mkdir(parents=True, exist_ok=True)
204263

264+
if _cancelled(cancel):
265+
return SyncResult(ok=False, name="干员名", cancelled=True)
266+
205267
opener = _make_opener(proxy or _settings_proxy())
206268
token = _settings_github_token()
207269
raw = _get_battle_data(opener, token=token)
@@ -283,19 +345,28 @@ def _rebuild_level_codes(map_dir: Path) -> int:
283345
return len(codes)
284346

285347

286-
def sync_maps(proxy: str | None = None) -> SyncResult:
348+
def sync_maps(
349+
proxy: str | None = None,
350+
progress_cb: Callable[[str], None] | None = None,
351+
cancel: threading.Event | None = None,
352+
) -> SyncResult:
287353
data_dir = project_root() / "data"
288354
map_dir = data_dir / "map"
289355
map_dir.mkdir(parents=True, exist_ok=True)
290356

357+
if _cancelled(cancel):
358+
return SyncResult(ok=False, name="地图", cancelled=True)
359+
291360
# 1. 先扫盘生成 level_codes.json —— 即便下载从未成功,运行期也能用已落盘的旧地图。
292361
pre_codes = _rebuild_level_codes(map_dir)
293362
logger.info("level_codes.json 初始刷新: %d 关", pre_codes)
294363

295364
# 2. 增量下载(基于 git/trees blob sha)。
296365
opener = _make_opener(proxy or _settings_proxy())
297366
token = _settings_github_token()
298-
result = _download_maps_remote(map_dir, opener, token=token)
367+
result = _download_maps_remote(
368+
map_dir, opener, token=token, progress_cb=progress_cb, cancel=cancel
369+
)
299370

300371
# 3. 下载完成后再扫一次,把新关卡也写进 level_codes.json。
301372
post_codes = _rebuild_level_codes(map_dir)
@@ -314,6 +385,8 @@ def sync_maps(proxy: str | None = None) -> SyncResult:
314385
failed,
315386
post_codes,
316387
)
388+
if _cancelled(cancel):
389+
return SyncResult(ok=False, name="地图", cancelled=True, done=done, skipped=skipped)
317390
return SyncResult(
318391
ok=failed == 0,
319392
name="地图",
@@ -334,7 +407,11 @@ def _load_manifest(path: Path) -> dict[str, str]:
334407

335408

336409
def _download_maps_remote(
337-
dst_dir: Path, opener: urllib.request.OpenerDirector, token: str | None = None
410+
dst_dir: Path,
411+
opener: urllib.request.OpenerDirector,
412+
token: str | None = None,
413+
progress_cb: Callable[[str], None] | None = None,
414+
cancel: threading.Event | None = None,
338415
) -> tuple[int, int, int] | None:
339416
"""从 GitHub 下载地图文件(增量:基于 git/trees blob sha)。
340417
@@ -344,6 +421,11 @@ def _download_maps_remote(
344421
- sha 变化或缺失 → 并发下载 raw URL
345422
失败的不更新 manifest sha,下次重试。
346423
424+
cancel 被置位时:停止派发新完成的 future、立即返回(已提交的小文件下载
425+
让池自然跑完,不阻塞调用方),已落盘的 manifest 进度保留供下次续传。
426+
427+
progress_cb 在下载期间被周期性回调("地图 N/M")。
428+
347429
Returns:
348430
(ready, done, failed) 或 None(git/trees 拉取失败)。
349431
ready = 下载/已是最新成功的累计。
@@ -383,13 +465,21 @@ def _commit_manifest() -> None:
383465
done = 0
384466
failed = 0
385467
if tasks:
386-
logger.info("地图增量下载: %d 个需更新 / %d 个总文件", len(tasks), ready + len(tasks))
387-
with ThreadPoolExecutor(max_workers=_MAX_WORKERS) as pool:
468+
total = len(tasks)
469+
logger.info("地图增量下载: %d 个需更新 / %d 个总文件", total, ready + total)
470+
if progress_cb:
471+
progress_cb(f"地图 0/{total}")
472+
# 不用 with:取消时 shutdown(wait=False, cancel_futures=True) 立即返回,
473+
# 不阻塞调用方等已派发的小文件下载跑完。
474+
pool = ThreadPoolExecutor(max_workers=_MAX_WORKERS)
475+
try:
388476
futs = {
389477
pool.submit(_download, opener, url, dst, token): (name, sha)
390478
for name, url, dst, sha in tasks
391479
}
392480
for fut in as_completed(futs):
481+
if _cancelled(cancel):
482+
break
393483
name, sha = futs[fut]
394484
if fut.result():
395485
ready += 1
@@ -398,22 +488,54 @@ def _commit_manifest() -> None:
398488
# chunked commit:每 100 个落盘一次,被中断也不丢进度。
399489
if done % 100 == 0:
400490
_commit_manifest()
401-
logger.info("地图下载进度: %d/%d", done, len(tasks))
491+
if progress_cb and (done % 50 == 0 or done == total):
492+
progress_cb(f"地图 {done}/{total}")
402493
else:
403494
failed += 1
495+
finally:
496+
# 取消:丢弃排队未启动的 future,不等正在跑的(小文件下载让其自然结束)。
497+
# 正常结束:等所有 future 跑完再退出,避免遗留线程。
498+
if _cancelled(cancel):
499+
pool.shutdown(wait=False, cancel_futures=True)
500+
else:
501+
pool.shutdown(wait=True)
404502
else:
405503
logger.info("地图数据已是最新(%d 文件)", ready)
406504

407505
_commit_manifest()
408506
return ready, done, failed
409507

410508

411-
def sync_all(proxy: str | None = None) -> list[SyncResult]:
412-
"""同步干员名 + 地图。返回每一步的 SyncResult,调用方可据此汇报。"""
509+
def sync_all(
510+
proxy: str | None = None,
511+
progress_cb: Callable[[str], None] | None = None,
512+
cancel: threading.Event | None = None,
513+
) -> list[SyncResult]:
514+
"""同步干员名 + 地图。返回每一步的 SyncResult,调用方可据此汇报。
515+
516+
progress_cb 用于阶段 + 地图逐文件进度;cancel 置位后尽快在阶段/文件间退出。
517+
"""
413518
logger.info("资源同步(远程)...")
414-
results = [sync_operators(proxy=proxy), sync_maps(proxy=proxy)]
519+
results: list[SyncResult] = []
520+
521+
if _cancelled(cancel):
522+
results.append(SyncResult(ok=False, name="干员名", cancelled=True))
523+
else:
524+
if progress_cb:
525+
progress_cb("干员名…")
526+
results.append(sync_operators(proxy=proxy, cancel=cancel))
527+
528+
if _cancelled(cancel):
529+
results.append(SyncResult(ok=False, name="地图", cancelled=True))
530+
else:
531+
if progress_cb:
532+
progress_cb("地图…")
533+
results.append(sync_maps(proxy=proxy, progress_cb=progress_cb, cancel=cancel))
534+
415535
if all(r.ok for r in results):
416536
logger.info("资源同步完成")
537+
elif any(r.cancelled for r in results):
538+
logger.info("资源同步已取消")
417539
else:
418540
failed_steps = [r.name for r in results if not r.ok]
419541
logger.warning("资源同步部分失败: %s", ", ".join(failed_steps))

aao/ui/settings_page.py

Lines changed: 32 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
from __future__ import annotations
88

99
import json
10+
import threading
1011
from typing import TYPE_CHECKING, Any
1112

1213
from PySide6.QtCore import QObject, Qt, QThread, QTimer, Signal
@@ -71,17 +72,26 @@ class _ResourceWorker(QObject):
7172
def __init__(self, mode: str):
7273
super().__init__()
7374
self._mode = mode # "sync" | "check_update" | "update_all"
75+
# 取消信号:跨线程 Event,主线程 set() 即可让 sync_all 在阶段/文件间退出。
76+
self._cancel = threading.Event()
77+
78+
def cancel(self) -> None:
79+
"""主线程调用:请求中止正在进行的同步。"""
80+
self._cancel.set()
7481

7582
def run(self) -> None:
7683
try:
7784
if self._mode == "sync":
7885
self.log.emit("开始同步资源(干员名 + 地图)...")
79-
results = sync_all()
86+
results = sync_all(
87+
progress_cb=lambda m: self.log.emit(m),
88+
cancel=self._cancel,
89+
)
8090
summary = ";".join(r.message for r in results)
8191
if all(r.ok for r in results):
8292
self.finished_ok.emit(f"资源同步完成({summary})")
8393
else:
84-
# 部分失败也走 finished_ok——结果消息里已包含失败信息
94+
# 部分失败/取消也走 finished_ok——结果消息里已包含失败/取消信息
8595
# 这里不弹 failed 信号以免 UI 显示成" fatal exception"。
8696
self.finished_ok.emit(summary)
8797
elif self._mode == "check_update":
@@ -149,6 +159,7 @@ def __init__(self):
149159
super().__init__()
150160
self._worker: _ResourceWorker | None = None
151161
self._thread: QThread | None = None
162+
self._res_mode: str | None = None # 当前后台资源操作类型(用于取消按钮判断)
152163
self._auto_preview_done = False
153164
self._collapsibles: dict[str, CollapsibleBox] = {}
154165
self._build_ui()
@@ -316,7 +327,7 @@ def _build_ui(self) -> None:
316327
self.btn_bg_pick.clicked.connect(self._on_bg_pick)
317328
self.btn_bg_clear.clicked.connect(self._on_bg_clear)
318329
self.slider_bg.valueChanged.connect(self._on_bg_opacity)
319-
self.btn_sync.clicked.connect(lambda: self._run_resource("sync"))
330+
self.btn_sync.clicked.connect(self._on_sync_clicked)
320331
self.btn_check.clicked.connect(lambda: self._run_resource("check_update"))
321332
self.btn_refresh_win.clicked.connect(self._refresh_windows)
322333
self.btn_preview.clicked.connect(self._preview_window)
@@ -568,11 +579,20 @@ def _on_save(self) -> None:
568579
self.lbl_op.setText("设置已保存(端口/profile/代理/Token 变更需重启或下次同步生效)")
569580
self.settings_changed.emit()
570581

582+
def _on_sync_clicked(self) -> None:
583+
"""同步按钮:未在跑 → 启动同步;正在同步 → 请求取消。"""
584+
if self._thread is not None and self._res_mode == "sync" and self._worker is not None:
585+
self._worker.cancel()
586+
self.btn_sync.setText("取消中…")
587+
self.btn_sync.setEnabled(False)
588+
return
589+
self._run_resource("sync")
590+
571591
def _run_resource(self, mode: str) -> None:
572592
if self._thread is not None:
573593
self.lbl_op.setText("上一次操作还在进行中…")
574594
return
575-
self._set_res_buttons(False)
595+
self._res_mode = mode
576596
self.lbl_op.setText("进行中…")
577597
self._worker = _ResourceWorker(mode)
578598
self._thread = QThread()
@@ -584,11 +604,18 @@ def _run_resource(self, mode: str) -> None:
584604
self._worker.finished_ok.connect(self._thread.quit)
585605
self._worker.failed.connect(self._thread.quit)
586606
self._thread.finished.connect(self._cleanup_thread)
607+
if mode == "sync":
608+
# 同步可取消:检查按钮禁用,同步按钮复用为取消入口。
609+
self.btn_check.setEnabled(False)
610+
self.btn_sync.setText("✕ 取消同步")
611+
else:
612+
self._set_res_buttons(False)
587613
self._thread.start()
588614

589615
def _on_res_done(self, msg: str) -> None:
590616
self.lbl_op.setText(msg)
591617
self._set_res_buttons(True)
618+
self.btn_sync.setText("🔄 同步资源")
592619
self.lbl_res_status.setText(self._res_status_text())
593620

594621
def _set_res_buttons(self, enabled: bool) -> None:
@@ -600,3 +627,4 @@ def _cleanup_thread(self) -> None:
600627
self._thread.wait()
601628
self._worker = None
602629
self._thread = None
630+
self._res_mode = None

0 commit comments

Comments
 (0)