feat: 支持持续自动获取任务 (#45)
This commit is contained in:
+166
-15
@@ -22,7 +22,14 @@
|
||||
|
||||
from typing import Dict, Optional
|
||||
|
||||
from PyQt5.QtCore import QCoreApplication, QObject, QThread, pyqtSignal, pyqtSlot
|
||||
from PyQt5.QtCore import (
|
||||
QCoreApplication,
|
||||
QObject,
|
||||
QThread,
|
||||
QTimer,
|
||||
pyqtSignal,
|
||||
pyqtSlot,
|
||||
)
|
||||
from qfluentwidgets import InfoBar, InfoBarPosition
|
||||
|
||||
from .admin_gateway import (
|
||||
@@ -122,6 +129,7 @@ class ClaimTaskWorker(QObject):
|
||||
taskSaved = pyqtSignal(str)
|
||||
duplicateTask = pyqtSignal(str)
|
||||
localSaveFailed = pyqtSignal(str, str)
|
||||
retryableFailed = pyqtSignal(str)
|
||||
failed = pyqtSignal(str)
|
||||
outcome = pyqtSignal(str, str, str)
|
||||
completed = pyqtSignal()
|
||||
@@ -176,7 +184,11 @@ class ClaimTaskWorker(QObject):
|
||||
request_hint = (
|
||||
f",请求编号:{exc.request_id}" if exc.request_id else ""
|
||||
)
|
||||
self.failed.emit(f"{exc}{request_hint}")
|
||||
message = f"{exc}{request_hint}"
|
||||
if exc.retryable:
|
||||
self.retryableFailed.emit(message)
|
||||
else:
|
||||
self.failed.emit(message)
|
||||
except Exception as exc:
|
||||
if not self._cancelled:
|
||||
self.failed.emit(f"执行采集任务失败:{exc}")
|
||||
@@ -195,6 +207,9 @@ class PDDTaskPageEvent(QObject):
|
||||
claim_gateway: Optional[AdminGateway] = None,
|
||||
settings_repository: Optional[SettingsRepository] = None,
|
||||
collect_service_factory: Optional[CollectServiceFactory] = None,
|
||||
next_task_delay_ms: int = 500,
|
||||
no_task_delay_ms: int = 5_000,
|
||||
retry_delays_ms: tuple[int, ...] = (5_000, 10_000, 20_000, 30_000),
|
||||
):
|
||||
super().__init__(parent or page)
|
||||
self._page = page
|
||||
@@ -204,8 +219,23 @@ class PDDTaskPageEvent(QObject):
|
||||
self._claim_busy = False
|
||||
self._claim_thread: Optional[QThread] = None
|
||||
self._claim_worker: Optional[ClaimTaskWorker] = None
|
||||
self._auto_fetch_running = False
|
||||
self._stop_requested = False
|
||||
self._cycle_next_delay_ms: Optional[int] = None
|
||||
self._stop_status = "自动获取:已停止 · 当前没有执行中的任务"
|
||||
self._retry_count = 0
|
||||
self._detail_windows: Dict[str, TaskDetailWindow] = {}
|
||||
self._collect_service_factory = collect_service_factory
|
||||
if next_task_delay_ms < 0 or no_task_delay_ms <= 0:
|
||||
raise ValueError("自动获取等待时间配置无效")
|
||||
if not retry_delays_ms or any(value <= 0 for value in retry_delays_ms):
|
||||
raise ValueError("自动获取重试等待时间必须全部大于 0")
|
||||
self._next_task_delay_ms = next_task_delay_ms
|
||||
self._no_task_delay_ms = no_task_delay_ms
|
||||
self._retry_delays_ms = retry_delays_ms
|
||||
self._next_cycle_timer = QTimer(self)
|
||||
self._next_cycle_timer.setSingleShot(True)
|
||||
self._next_cycle_timer.timeout.connect(self._start_claim_cycle)
|
||||
|
||||
try:
|
||||
self._repository.recover_interrupted_work()
|
||||
@@ -238,7 +268,7 @@ class PDDTaskPageEvent(QObject):
|
||||
page.autoFetchRequested.connect(self._request_claim_task)
|
||||
page.detailRequested.connect(self.show_task_detail)
|
||||
page.taskModel.loadMoreRequested.connect(self._load_page)
|
||||
page.autoFetchButton.setText("获取任务")
|
||||
page.set_auto_fetch_state("stopped")
|
||||
page.destroyed.connect(self.shutdown)
|
||||
application = QCoreApplication.instance()
|
||||
if application is not None:
|
||||
@@ -266,9 +296,12 @@ class PDDTaskPageEvent(QObject):
|
||||
|
||||
@pyqtSlot()
|
||||
def _request_claim_task(self) -> None:
|
||||
"""启动一次后台领取;重复点击和关闭期间直接忽略。"""
|
||||
"""切换持续自动获取;每一轮仍只启动一个后台 Worker。"""
|
||||
|
||||
if self._closing or self._claim_busy:
|
||||
if self._closing:
|
||||
return
|
||||
if self._auto_fetch_running:
|
||||
self._request_stop_auto_fetch()
|
||||
return
|
||||
if self._claim_gateway_error or self._claim_gateway is None:
|
||||
message = self._claim_gateway_error or "Admin 领取服务未初始化"
|
||||
@@ -276,9 +309,28 @@ class PDDTaskPageEvent(QObject):
|
||||
self._show_claim_error("领取任务失败", message)
|
||||
return
|
||||
|
||||
self._auto_fetch_running = True
|
||||
self._stop_requested = False
|
||||
self._retry_count = 0
|
||||
self._page.set_auto_fetch_state("starting")
|
||||
self._page.set_engine_status("自动获取:正在启动…")
|
||||
self._start_claim_cycle()
|
||||
|
||||
@pyqtSlot()
|
||||
def _start_claim_cycle(self) -> None:
|
||||
"""开始一轮任务处理;只允许上一轮线程完全结束后进入。"""
|
||||
|
||||
if (
|
||||
self._closing
|
||||
or not self._auto_fetch_running
|
||||
or self._stop_requested
|
||||
or self._claim_busy
|
||||
):
|
||||
return
|
||||
self._claim_busy = True
|
||||
self._page.autoFetchButton.setEnabled(False)
|
||||
self._page.set_engine_status("正在处理一条采集任务,请稍候…")
|
||||
self._cycle_next_delay_ms = None
|
||||
self._page.set_auto_fetch_state("running")
|
||||
self._page.set_engine_status("自动获取:运行中 · 正在处理一条采集任务…")
|
||||
|
||||
thread = QThread(self)
|
||||
worker = ClaimTaskWorker(
|
||||
@@ -294,6 +346,7 @@ class PDDTaskPageEvent(QObject):
|
||||
worker.taskSaved.connect(self._on_claimed_task_saved)
|
||||
worker.duplicateTask.connect(self._on_duplicate_claimed_task)
|
||||
worker.localSaveFailed.connect(self._on_claimed_task_save_failed)
|
||||
worker.retryableFailed.connect(self._on_claim_retryable_failed)
|
||||
worker.failed.connect(self._on_claim_failed)
|
||||
worker.outcome.connect(self._on_collect_outcome)
|
||||
worker.completed.connect(thread.quit)
|
||||
@@ -304,10 +357,57 @@ class PDDTaskPageEvent(QObject):
|
||||
self._claim_worker = worker
|
||||
thread.start()
|
||||
|
||||
def _request_stop_auto_fetch(self) -> None:
|
||||
"""取消下一轮,并请求当前 Worker 在安全点停止。"""
|
||||
|
||||
self._auto_fetch_running = False
|
||||
self._stop_requested = True
|
||||
self._next_cycle_timer.stop()
|
||||
self._stop_status = "自动获取:已停止 · 当前没有执行中的任务"
|
||||
worker = self._claim_worker
|
||||
if worker is not None:
|
||||
self._page.set_auto_fetch_state("stopping")
|
||||
self._page.set_engine_status(
|
||||
"自动获取:正在停止 · 等待当前任务安全结束"
|
||||
)
|
||||
try:
|
||||
worker.cancel()
|
||||
except RuntimeError:
|
||||
pass
|
||||
return
|
||||
self._finish_auto_fetch_stop()
|
||||
|
||||
def _finish_auto_fetch_stop(self) -> None:
|
||||
"""把按钮和状态区恢复为已停止。"""
|
||||
|
||||
self._page.set_auto_fetch_state("stopped")
|
||||
self._page.set_engine_status(self._stop_status)
|
||||
self._stop_requested = False
|
||||
|
||||
def _continue_after(self, delay_ms: int, status: str) -> None:
|
||||
"""记录本轮结束后的下一次调度,不在 Worker 结束前启动。"""
|
||||
|
||||
self._cycle_next_delay_ms = delay_ms
|
||||
self._page.set_engine_status(status)
|
||||
|
||||
def _stop_after_current(self, status: str) -> None:
|
||||
"""遇到不可自动处理的问题时,完成线程收尾后停止。"""
|
||||
|
||||
self._auto_fetch_running = False
|
||||
self._stop_requested = True
|
||||
self._stop_status = status
|
||||
self._next_cycle_timer.stop()
|
||||
self._page.set_auto_fetch_state("stopping")
|
||||
|
||||
@pyqtSlot()
|
||||
def _on_no_claimed_task(self) -> None:
|
||||
if not self._closing:
|
||||
self._page.set_engine_status("暂无可领取的采集任务")
|
||||
self._retry_count = 0
|
||||
seconds = self._no_task_delay_ms / 1000
|
||||
self._continue_after(
|
||||
self._no_task_delay_ms,
|
||||
f"自动获取:运行中 · 暂无任务,{seconds:g} 秒后再次领取",
|
||||
)
|
||||
|
||||
@pyqtSlot(str)
|
||||
def _on_claimed_task_saved(self, task_id: str) -> None:
|
||||
@@ -316,13 +416,15 @@ class PDDTaskPageEvent(QObject):
|
||||
self._page.set_engine_status(
|
||||
f"已领取任务 {task_id},已保存到本地任务列表"
|
||||
)
|
||||
self._cycle_next_delay_ms = self._next_task_delay_ms
|
||||
self._reload()
|
||||
|
||||
@pyqtSlot(str)
|
||||
def _on_duplicate_claimed_task(self, task_id: str) -> None:
|
||||
if not self._closing:
|
||||
self._page.set_engine_status(
|
||||
f"任务 {task_id} 本地已有,未重复保存"
|
||||
self._continue_after(
|
||||
self._next_task_delay_ms,
|
||||
f"任务 {task_id} 本地已有,未重复保存",
|
||||
)
|
||||
|
||||
@pyqtSlot(str, str)
|
||||
@@ -334,6 +436,22 @@ class PDDTaskPageEvent(QObject):
|
||||
)
|
||||
self._page.set_engine_status(content)
|
||||
self._show_claim_error("本地任务保存失败", content)
|
||||
self._stop_after_current(content)
|
||||
|
||||
@pyqtSlot(str)
|
||||
def _on_claim_retryable_failed(self, message: str) -> None:
|
||||
if self._closing:
|
||||
return
|
||||
index = min(self._retry_count, len(self._retry_delays_ms) - 1)
|
||||
delay_ms = self._retry_delays_ms[index]
|
||||
self._retry_count += 1
|
||||
seconds = delay_ms / 1000
|
||||
status = (
|
||||
f"自动获取:退避等待 · {message},{seconds:g} 秒后重试"
|
||||
)
|
||||
self._continue_after(delay_ms, status)
|
||||
if self._retry_count == 1:
|
||||
self._show_claim_error("Admin 暂时不可用", status)
|
||||
|
||||
@pyqtSlot(str)
|
||||
def _on_claim_failed(self, message: str) -> None:
|
||||
@@ -341,15 +459,32 @@ class PDDTaskPageEvent(QObject):
|
||||
content = message or "领取任务失败,请稍后重试"
|
||||
self._page.set_engine_status(content)
|
||||
self._show_claim_error("领取任务失败", content)
|
||||
self._stop_after_current(content)
|
||||
|
||||
@pyqtSlot(str, str, str)
|
||||
def _on_collect_outcome(self, kind: str, message: str, _task_id: str) -> None:
|
||||
if self._closing:
|
||||
return
|
||||
self._page.set_engine_status(message)
|
||||
self._reload()
|
||||
if kind in {"result_pending", "manual_review", "failed"}:
|
||||
if kind == "no_task":
|
||||
self._on_no_claimed_task()
|
||||
elif kind == "succeeded":
|
||||
self._retry_count = 0
|
||||
self._continue_after(
|
||||
self._next_task_delay_ms,
|
||||
f"自动获取:运行中 · {message}",
|
||||
)
|
||||
elif kind == "result_pending":
|
||||
self._on_claim_retryable_failed(message)
|
||||
elif kind in {"manual_review", "failed"}:
|
||||
self._show_claim_error("采集任务需要处理", message)
|
||||
self._stop_after_current(message)
|
||||
elif kind == "cancelled":
|
||||
self._stop_after_current(message or "自动获取已停止")
|
||||
else:
|
||||
content = message or f"任务返回未知结果:{kind}"
|
||||
self._show_claim_error("自动获取已停止", content)
|
||||
self._stop_after_current(content)
|
||||
|
||||
def _show_claim_error(self, title: str, content: str) -> None:
|
||||
"""显示不会自动消失的可恢复错误,同时保留底部状态文字。"""
|
||||
@@ -368,9 +503,18 @@ class PDDTaskPageEvent(QObject):
|
||||
self._claim_worker = None
|
||||
self._claim_thread = None
|
||||
self._claim_busy = False
|
||||
if not self._closing:
|
||||
self._page.autoFetchButton.setText("获取任务")
|
||||
self._page.autoFetchButton.setEnabled(True)
|
||||
if self._closing:
|
||||
return
|
||||
if not self._auto_fetch_running or self._stop_requested:
|
||||
self._finish_auto_fetch_stop()
|
||||
return
|
||||
delay_ms = self._cycle_next_delay_ms
|
||||
if delay_ms is None:
|
||||
self._stop_after_current("自动获取:已停止 · 本轮没有返回有效结果")
|
||||
self._finish_auto_fetch_stop()
|
||||
return
|
||||
self._page.set_auto_fetch_state("running")
|
||||
self._next_cycle_timer.start(delay_ms)
|
||||
|
||||
def _reload(self) -> None:
|
||||
self._page.begin_task_reload()
|
||||
@@ -446,6 +590,9 @@ class PDDTaskPageEvent(QObject):
|
||||
if self._closing:
|
||||
return
|
||||
self._closing = True
|
||||
self._auto_fetch_running = False
|
||||
self._stop_requested = True
|
||||
self._next_cycle_timer.stop()
|
||||
|
||||
for window in list(self._detail_windows.values()):
|
||||
window.close()
|
||||
@@ -465,6 +612,10 @@ class PDDTaskPageEvent(QObject):
|
||||
self._on_claimed_task_save_failed,
|
||||
),
|
||||
(worker.failed, self._on_claim_failed),
|
||||
(
|
||||
worker.retryableFailed,
|
||||
self._on_claim_retryable_failed,
|
||||
),
|
||||
(worker.outcome, self._on_collect_outcome),
|
||||
)
|
||||
except RuntimeError:
|
||||
|
||||
Reference in New Issue
Block a user