mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-09-07 00:16:57 +08:00
fix: 收口监控测试线程并修复整理历史 lint (#6304)
This commit is contained in:
+32
-33
@@ -4037,43 +4037,42 @@ class TransferChain(ChainBase, ConfigReloadMixin, metaclass=Singleton):
|
|||||||
file_modify_time=file_item.modify_time,
|
file_modify_time=file_item.modify_time,
|
||||||
fileid=file_item.fileid,
|
fileid=file_item.fileid,
|
||||||
)
|
)
|
||||||
|
if not manual:
|
||||||
|
# 自动路径(目录监控、下载器轮询)与监控分发共用同一套判定,
|
||||||
|
# 否则监控层刚放行的失败重试与升级请求会在这里被全额收回
|
||||||
|
gate_action = evaluate_history_gate(
|
||||||
|
transferd,
|
||||||
|
file_size=file_item.size,
|
||||||
|
file_modify_time=file_item.modify_time,
|
||||||
|
fileid=file_item.fileid,
|
||||||
|
)
|
||||||
|
if not is_skip_action(gate_action):
|
||||||
|
logger.info(
|
||||||
|
f"{file_item.path} 命中"
|
||||||
|
f"{history_description}"
|
||||||
|
f",重新送入整理"
|
||||||
|
)
|
||||||
|
transferd = None
|
||||||
|
|
||||||
if transferd and not manual:
|
if transferd:
|
||||||
# 自动路径(目录监控、下载器轮询)与监控分发共用同一套判定,
|
skipped_history_count += 1
|
||||||
# 否则监控层刚放行的失败重试与升级请求会在这里被全额收回
|
if not transferd.status:
|
||||||
gate_action = evaluate_history_gate(
|
all_success = False
|
||||||
transferd,
|
# 失败记录能走到这里说明重试次数已用尽,此时同样要打已整理标签让种子
|
||||||
file_size=file_item.size,
|
# 退出轮询,否则下载器每一轮都会重新扫描并在这里被拦一次,空转且刷屏
|
||||||
file_modify_time=file_item.modify_time,
|
candidate_hash = download_hash or transferd.download_hash
|
||||||
fileid=file_item.fileid,
|
candidate_downloader = downloader or transferd.downloader
|
||||||
)
|
if candidate_hash and candidate_downloader:
|
||||||
if not is_skip_action(gate_action):
|
skipped_torrents.add(
|
||||||
|
(candidate_hash, candidate_downloader)
|
||||||
|
)
|
||||||
logger.info(
|
logger.info(
|
||||||
f"{file_item.path} 命中"
|
f"{file_item.path} 已整理过("
|
||||||
f"{history_description}"
|
f"{history_description}"
|
||||||
f",重新送入整理"
|
f"),如需重新处理,请删除整理记录。"
|
||||||
)
|
)
|
||||||
transferd = None
|
err_msgs.append(f"{file_item.name} 已整理过")
|
||||||
|
continue
|
||||||
if transferd:
|
|
||||||
skipped_history_count += 1
|
|
||||||
if not transferd.status:
|
|
||||||
all_success = False
|
|
||||||
# 失败记录能走到这里说明重试次数已用尽,此时同样要打已整理标签让种子
|
|
||||||
# 退出轮询,否则下载器每一轮都会重新扫描并在这里被拦一次,空转且刷屏
|
|
||||||
candidate_hash = download_hash or transferd.download_hash
|
|
||||||
candidate_downloader = downloader or transferd.downloader
|
|
||||||
if candidate_hash and candidate_downloader:
|
|
||||||
skipped_torrents.add(
|
|
||||||
(candidate_hash, candidate_downloader)
|
|
||||||
)
|
|
||||||
logger.info(
|
|
||||||
f"{file_item.path} 已整理过("
|
|
||||||
f"{history_description}"
|
|
||||||
f"),如需重新处理,请删除整理记录。"
|
|
||||||
)
|
|
||||||
err_msgs.append(f"{file_item.name} 已整理过")
|
|
||||||
continue
|
|
||||||
|
|
||||||
# 提前获取下载历史,以便获取自定义识别词
|
# 提前获取下载历史,以便获取自定义识别词
|
||||||
downloadhis = DownloadHistoryOper()
|
downloadhis = DownloadHistoryOper()
|
||||||
|
|||||||
@@ -81,8 +81,48 @@ def _run_watchdog(monitor, timeout=10.0):
|
|||||||
finally:
|
finally:
|
||||||
done.set()
|
done.set()
|
||||||
|
|
||||||
Thread(target=runner, daemon=True, name="test-watchdog").start()
|
thread = Thread(target=runner, daemon=True, name="test-watchdog")
|
||||||
return done.wait(timeout=timeout)
|
thread.start()
|
||||||
|
finished = done.wait(timeout=timeout)
|
||||||
|
if finished:
|
||||||
|
thread.join(timeout=timeout)
|
||||||
|
return finished and not thread.is_alive()
|
||||||
|
|
||||||
|
|
||||||
|
def _release_and_join(release, threads, timeout=10.0):
|
||||||
|
"""
|
||||||
|
释放测试注入的阻塞,并等待受控工作线程结束。
|
||||||
|
|
||||||
|
真实故障下的恢复线程允许随进程退出,但测试主动解除阻塞后必须完成回收,
|
||||||
|
避免仍在执行应用代码的守护线程与解释器关闭过程并发。
|
||||||
|
:param release: 控制阻塞动作的事件
|
||||||
|
:param threads: 需要等待的线程集合
|
||||||
|
:param timeout: 每个线程的最长等待秒数
|
||||||
|
"""
|
||||||
|
release.set()
|
||||||
|
for thread in threads:
|
||||||
|
thread.join(timeout=timeout)
|
||||||
|
assert not thread.is_alive(), f"受控线程未能结束: {thread.name}"
|
||||||
|
|
||||||
|
|
||||||
|
def _stop_monitor_threads(monitor, timeout=10.0):
|
||||||
|
"""
|
||||||
|
停止测试期间重建的目录监控,并等待其派生任务结束。
|
||||||
|
|
||||||
|
目录 watcher 会进入 watchfiles 原生扩展,必须在解释器关闭前退出;补偿扫描
|
||||||
|
由重建流程派生,也要在测试边界内完成。
|
||||||
|
:param monitor: Monitor 骨架
|
||||||
|
:param timeout: 每个线程的最长等待秒数
|
||||||
|
"""
|
||||||
|
for watcher in monitor._watchers:
|
||||||
|
if isinstance(watcher, LocalDirectoryWatcher):
|
||||||
|
watcher.stop()
|
||||||
|
watcher.join(timeout=timeout)
|
||||||
|
assert not watcher.is_alive(), f"目录监控线程未能结束: {watcher.watch_path}"
|
||||||
|
for thread in threading.enumerate():
|
||||||
|
if thread.name.startswith("MoviePilot-MonitorCompensation-"):
|
||||||
|
thread.join(timeout=timeout)
|
||||||
|
assert not thread.is_alive(), f"补偿扫描线程未能结束: {thread.name}"
|
||||||
|
|
||||||
|
|
||||||
# --------------------------------------------------------------------------- #
|
# --------------------------------------------------------------------------- #
|
||||||
@@ -115,7 +155,7 @@ def test_watchdog_returns_while_rebuild_blocks_forever(tmp_path, monkeypatch):
|
|||||||
assert _run_watchdog(monitor), "看门狗被重建动作冻死,未能在限定时间内返回"
|
assert _run_watchdog(monitor), "看门狗被重建动作冻死,未能在限定时间内返回"
|
||||||
assert entered.is_set(), "重建动作没有被真正发起"
|
assert entered.is_set(), "重建动作没有被真正发起"
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, monitor._recovery._running.values())
|
||||||
|
|
||||||
|
|
||||||
def test_watchdog_returns_while_pending_retry_blocks_forever(tmp_path, monkeypatch):
|
def test_watchdog_returns_while_pending_retry_blocks_forever(tmp_path, monkeypatch):
|
||||||
@@ -130,7 +170,7 @@ def test_watchdog_returns_while_pending_retry_blocks_forever(tmp_path, monkeypat
|
|||||||
try:
|
try:
|
||||||
assert _run_watchdog(monitor), "看门狗被整理重试驱动冻死,未能在限定时间内返回"
|
assert _run_watchdog(monitor), "看门狗被整理重试驱动冻死,未能在限定时间内返回"
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, monitor._recovery._running.values())
|
||||||
|
|
||||||
|
|
||||||
def test_watchdog_returns_while_local_retry_blocks_forever(tmp_path, monkeypatch):
|
def test_watchdog_returns_while_local_retry_blocks_forever(tmp_path, monkeypatch):
|
||||||
@@ -155,7 +195,7 @@ def test_watchdog_returns_while_local_retry_blocks_forever(tmp_path, monkeypatch
|
|||||||
try:
|
try:
|
||||||
assert _run_watchdog(monitor), "看门狗被本地监控重试冻死,未能在限定时间内返回"
|
assert _run_watchdog(monitor), "看门狗被本地监控重试冻死,未能在限定时间内返回"
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, monitor._recovery._running.values())
|
||||||
|
|
||||||
|
|
||||||
def test_check_watchers_performs_no_filesystem_access(tmp_path, monkeypatch):
|
def test_check_watchers_performs_no_filesystem_access(tmp_path, monkeypatch):
|
||||||
@@ -221,7 +261,8 @@ def test_watchdog_survives_blocking_exists_in_real_rebuild(tmp_path, monkeypatch
|
|||||||
assert _run_watchdog(monitor), "看门狗冻死在真实的重建调用链上"
|
assert _run_watchdog(monitor), "看门狗冻死在真实的重建调用链上"
|
||||||
assert str(tmp_path) in monitor._isolated, "重建无响应后目录没有转入隔离"
|
assert str(tmp_path) in monitor._isolated, "重建无响应后目录没有转入隔离"
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, monitor._recovery._running.values())
|
||||||
|
_stop_monitor_threads(monitor)
|
||||||
|
|
||||||
|
|
||||||
# --------------------------------------------------------------------------- #
|
# --------------------------------------------------------------------------- #
|
||||||
@@ -259,7 +300,7 @@ def test_rebuild_timeout_isolates_directory(tmp_path, monkeypatch):
|
|||||||
assert _run_watchdog(monitor)
|
assert _run_watchdog(monitor)
|
||||||
assert len(calls) == 1, "隔离中的目录仍在被反复重建,会持续泄漏冻死的线程"
|
assert len(calls) == 1, "隔离中的目录仍在被反复重建,会持续泄漏冻死的线程"
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, monitor._recovery._running.values())
|
||||||
|
|
||||||
|
|
||||||
def test_isolated_directory_recovers_after_probe_succeeds(tmp_path, monkeypatch):
|
def test_isolated_directory_recovers_after_probe_succeeds(tmp_path, monkeypatch):
|
||||||
@@ -392,17 +433,19 @@ def test_stuck_transfer_does_not_block_other_files(monkeypatch):
|
|||||||
"""
|
"""
|
||||||
dispatcher.handle_file(storage="local", event_path=Path(f"/mnt/cd2/{name}"), file_size=1)
|
dispatcher.handle_file(storage="local", event_path=Path(f"/mnt/cd2/{name}"), file_size=1)
|
||||||
|
|
||||||
Thread(target=feed, args=("stuck.mkv",), daemon=True).start()
|
stuck_thread = Thread(target=feed, args=("stuck.mkv",), daemon=True)
|
||||||
|
stuck_thread.start()
|
||||||
assert stuck_entered.wait(timeout=5), "卡死的整理没有真正进入"
|
assert stuck_entered.wait(timeout=5), "卡死的整理没有真正进入"
|
||||||
|
|
||||||
done = threading.Event()
|
done = threading.Event()
|
||||||
Thread(target=lambda: (feed("other.mkv"), done.set()), daemon=True).start()
|
other_thread = Thread(target=lambda: (feed("other.mkv"), done.set()), daemon=True)
|
||||||
|
other_thread.start()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
assert done.wait(timeout=5), "一个文件卡在死挂载上就锁死了整条整理链"
|
assert done.wait(timeout=5), "一个文件卡在死挂载上就锁死了整条整理链"
|
||||||
assert "/mnt/cd2/other.mkv" in transferred
|
assert "/mnt/cd2/other.mkv" in transferred
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, (stuck_thread, other_thread))
|
||||||
|
|
||||||
|
|
||||||
# --------------------------------------------------------------------------- #
|
# --------------------------------------------------------------------------- #
|
||||||
@@ -424,7 +467,7 @@ def test_recovery_executor_reports_timeout_without_blocking():
|
|||||||
assert results == {"stuck": RecoveryState.TIMEOUT}
|
assert results == {"stuck": RecoveryState.TIMEOUT}
|
||||||
assert elapsed < 3.0, "执行器没有在超时后放弃冻死的动作"
|
assert elapsed < 3.0, "执行器没有在超时后放弃冻死的动作"
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, executor._running.values())
|
||||||
|
|
||||||
|
|
||||||
def test_recovery_executor_skips_key_with_running_task():
|
def test_recovery_executor_skips_key_with_running_task():
|
||||||
@@ -448,7 +491,7 @@ def test_recovery_executor_skips_key_with_running_task():
|
|||||||
assert executor.run({"k": stuck}, timeout=0.2) == {"k": RecoveryState.BUSY}
|
assert executor.run({"k": stuck}, timeout=0.2) == {"k": RecoveryState.BUSY}
|
||||||
assert len(calls) == 1
|
assert len(calls) == 1
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, executor._running.values())
|
||||||
|
|
||||||
|
|
||||||
def test_recovery_executor_runs_actions_concurrently():
|
def test_recovery_executor_runs_actions_concurrently():
|
||||||
@@ -467,7 +510,7 @@ def test_recovery_executor_runs_actions_concurrently():
|
|||||||
assert set(results.values()) == {RecoveryState.TIMEOUT}
|
assert set(results.values()) == {RecoveryState.TIMEOUT}
|
||||||
assert elapsed < 1.5, "恢复动作是串行等待的,总耗时随目录数增长"
|
assert elapsed < 1.5, "恢复动作是串行等待的,总耗时随目录数增长"
|
||||||
finally:
|
finally:
|
||||||
release.set()
|
_release_and_join(release, executor._running.values())
|
||||||
|
|
||||||
|
|
||||||
def test_recovery_executor_reports_completion_and_swallows_errors():
|
def test_recovery_executor_reports_completion_and_swallows_errors():
|
||||||
|
|||||||
Reference in New Issue
Block a user