feat: 增加任务批量勾选和重新上报 (#89)
This commit is contained in:
+195
-1
@@ -50,6 +50,7 @@ from .purchase_reconcile_service import PurchaseReconcileFactory
|
||||
from .selected_android_device_service import SelectedAndroidDeviceService
|
||||
from .settings_repository import SettingsRepository
|
||||
from .task_models import (
|
||||
OutboxStatus,
|
||||
TaskFilters,
|
||||
TaskStatus,
|
||||
TaskSummary,
|
||||
@@ -211,6 +212,82 @@ class ClaimTaskWorker(QObject):
|
||||
self.completed.emit()
|
||||
|
||||
|
||||
class ResultResubmitWorker(QObject):
|
||||
"""在后台逐条重发既有结果 Outbox,不执行任何手机操作。"""
|
||||
|
||||
progress = pyqtSignal(int, int, str)
|
||||
finished = pyqtSignal(object, object, object)
|
||||
completed = pyqtSignal()
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
gateway: AdminGateway,
|
||||
repository: TaskRepository,
|
||||
task_ids: tuple[str, ...],
|
||||
) -> None:
|
||||
super().__init__()
|
||||
self._gateway = gateway
|
||||
self._repository = repository
|
||||
self._task_ids = task_ids
|
||||
self._cancelled = False
|
||||
|
||||
def cancel(self) -> None:
|
||||
"""当前网络请求结束后停止处理后续任务。"""
|
||||
|
||||
self._cancelled = True
|
||||
|
||||
@pyqtSlot()
|
||||
def run(self) -> None:
|
||||
succeeded: list[str] = []
|
||||
failed: list[str] = []
|
||||
skipped: list[str] = []
|
||||
total = len(self._task_ids)
|
||||
try:
|
||||
for current, task_id in enumerate(self._task_ids, 1):
|
||||
if self._cancelled:
|
||||
break
|
||||
self.progress.emit(current, total, task_id)
|
||||
try:
|
||||
event = self._repository.latest_result_outbox(task_id)
|
||||
except Exception:
|
||||
skipped.append(task_id)
|
||||
continue
|
||||
if event is None or event.status is OutboxStatus.SENDING:
|
||||
skipped.append(task_id)
|
||||
continue
|
||||
|
||||
self._repository.mark_outbox_sending(event.id)
|
||||
try:
|
||||
receipt = self._gateway.submit_result(
|
||||
task_id,
|
||||
event.idempotency_key,
|
||||
event.payload_json,
|
||||
)
|
||||
if not receipt.accepted:
|
||||
raise AdminGatewayError(
|
||||
"ADMIN_RESULT_NOT_ACCEPTED",
|
||||
"Admin 未确认接收任务结果",
|
||||
False,
|
||||
)
|
||||
except AdminGatewayError as exc:
|
||||
if exc.retryable:
|
||||
self._repository.mark_outbox_retry(event.id, str(exc))
|
||||
else:
|
||||
self._repository.mark_outbox_failed(event.id, str(exc))
|
||||
failed.append(task_id)
|
||||
continue
|
||||
except Exception as exc:
|
||||
self._repository.mark_outbox_retry(event.id, str(exc))
|
||||
failed.append(task_id)
|
||||
continue
|
||||
|
||||
self._repository.mark_outbox_sent(event.id)
|
||||
succeeded.append(task_id)
|
||||
self.finished.emit(succeeded, failed, skipped)
|
||||
finally:
|
||||
self.completed.emit()
|
||||
|
||||
|
||||
class PDDTaskPageEvent(QObject):
|
||||
"""把 PDD 页面只读操作连接到本地任务 Repository。"""
|
||||
|
||||
@@ -242,6 +319,9 @@ class PDDTaskPageEvent(QObject):
|
||||
self._claim_feedback: Optional[InfoBar] = None
|
||||
self._device_feedback: Optional[InfoBar] = None
|
||||
self._auto_fetch_running = False
|
||||
self._resubmit_busy = False
|
||||
self._resubmit_thread: Optional[QThread] = None
|
||||
self._resubmit_worker: Optional[ResultResubmitWorker] = None
|
||||
self._stop_requested = False
|
||||
self._cycle_next_delay_ms: Optional[int] = None
|
||||
self._stop_status = "自动获取:已停止 · 当前没有执行中的任务"
|
||||
@@ -295,6 +375,7 @@ class PDDTaskPageEvent(QObject):
|
||||
page.refreshRequested.connect(self.refresh_tasks)
|
||||
page.rerunRequested.connect(self.request_rerun)
|
||||
page.rerunCancelRequested.connect(self.request_cancel_rerun)
|
||||
page.resubmitRequested.connect(self.request_resubmit)
|
||||
page.autoFetchRequested.connect(self._request_claim_task)
|
||||
page.detailRequested.connect(self.show_task_detail)
|
||||
page.taskModel.loadMoreRequested.connect(self._load_page)
|
||||
@@ -330,7 +411,7 @@ class PDDTaskPageEvent(QObject):
|
||||
|
||||
if self._closing or not task_id:
|
||||
return
|
||||
if self._auto_fetch_running or self._claim_busy:
|
||||
if self._auto_fetch_running or self._claim_busy or self._resubmit_busy:
|
||||
self._show_rerun_warning(
|
||||
"暂时不能重新执行",
|
||||
"自动获取或其他采集正在运行,请停止并等待当前任务结束。",
|
||||
@@ -386,6 +467,100 @@ class PDDTaskPageEvent(QObject):
|
||||
|
||||
self._start_rerun_worker(task_id)
|
||||
|
||||
@pyqtSlot(object)
|
||||
def request_resubmit(self, task_ids) -> None:
|
||||
"""确认后在后台重发勾选任务的原始结果 Outbox。"""
|
||||
|
||||
stable_ids = tuple(
|
||||
dict.fromkeys(str(value) for value in task_ids if value)
|
||||
)
|
||||
if self._closing or not stable_ids:
|
||||
return
|
||||
if self._auto_fetch_running or self._claim_busy or self._resubmit_busy:
|
||||
self._show_rerun_warning(
|
||||
"暂时不能重新上报",
|
||||
"自动获取、重新执行或另一批上报正在运行,请等待当前操作结束。",
|
||||
)
|
||||
return
|
||||
if self._claim_gateway_error or self._claim_gateway is None:
|
||||
self._show_rerun_warning(
|
||||
"重新上报不可用",
|
||||
self._claim_gateway_error or "Admin 提交服务未初始化",
|
||||
)
|
||||
return
|
||||
|
||||
dialog = MessageBox(
|
||||
f"重新上报 {len(stable_ids)} 条任务结果?",
|
||||
"只会重新提交本地已经保存的原始结果,"
|
||||
"不会重新采集、采购或操作 Android 手机,也不会修改本地结果。",
|
||||
self._page.window(),
|
||||
)
|
||||
dialog.yesButton.setText(f"上报 {len(stable_ids)} 条")
|
||||
dialog.cancelButton.setText("暂不上报")
|
||||
dialog.cancelButton.setFocus()
|
||||
if not dialog.exec():
|
||||
return
|
||||
|
||||
self._start_resubmit_worker(stable_ids)
|
||||
|
||||
def _start_resubmit_worker(self, task_ids: tuple[str, ...]) -> None:
|
||||
"""启动只处理指定结果 Outbox 的工作线程。"""
|
||||
|
||||
assert self._claim_gateway is not None
|
||||
self._resubmit_busy = True
|
||||
self._page.set_resubmit_running(True)
|
||||
self._page.set_engine_status(
|
||||
f"正在准备重新上报 {len(task_ids)} 条任务结果…"
|
||||
)
|
||||
|
||||
thread = QThread(self)
|
||||
worker = ResultResubmitWorker(
|
||||
self._claim_gateway,
|
||||
self._repository,
|
||||
task_ids,
|
||||
)
|
||||
worker.moveToThread(thread)
|
||||
thread.started.connect(worker.run)
|
||||
worker.progress.connect(self._on_resubmit_progress)
|
||||
worker.finished.connect(self._on_resubmit_finished)
|
||||
worker.completed.connect(thread.quit)
|
||||
worker.completed.connect(worker.deleteLater)
|
||||
thread.finished.connect(thread.deleteLater)
|
||||
thread.finished.connect(self._on_resubmit_thread_finished)
|
||||
self._resubmit_thread = thread
|
||||
self._resubmit_worker = worker
|
||||
thread.start()
|
||||
|
||||
@pyqtSlot(int, int, str)
|
||||
def _on_resubmit_progress(self, current: int, total: int, task_id: str) -> None:
|
||||
if not self._closing:
|
||||
self._page.set_engine_status(
|
||||
f"正在重新上报 {current}/{total}:{task_id}"
|
||||
)
|
||||
|
||||
@pyqtSlot(object, object, object)
|
||||
def _on_resubmit_finished(self, succeeded, failed, skipped) -> None:
|
||||
if self._closing:
|
||||
return
|
||||
succeeded_ids = tuple(succeeded)
|
||||
self._page.taskModel.uncheck_tasks(succeeded_ids)
|
||||
message = (
|
||||
f"重新上报完成:成功 {len(succeeded_ids)} 条,"
|
||||
f"失败 {len(failed)} 条,跳过 {len(skipped)} 条"
|
||||
)
|
||||
self._page.set_engine_status(message)
|
||||
if failed or skipped:
|
||||
self._show_claim_error("重新上报未全部完成", message)
|
||||
self._reload()
|
||||
|
||||
@pyqtSlot()
|
||||
def _on_resubmit_thread_finished(self) -> None:
|
||||
self._resubmit_worker = None
|
||||
self._resubmit_thread = None
|
||||
self._resubmit_busy = False
|
||||
if not self._closing:
|
||||
self._page.set_resubmit_running(False)
|
||||
|
||||
def _start_rerun_worker(self, task_id: str) -> None:
|
||||
"""启动只处理指定任务的工作线程。"""
|
||||
|
||||
@@ -557,6 +732,12 @@ class PDDTaskPageEvent(QObject):
|
||||
|
||||
if self._closing:
|
||||
return
|
||||
if self._resubmit_busy:
|
||||
self._show_rerun_warning(
|
||||
"暂时不能获取任务",
|
||||
"任务结果正在重新上报,请等待当前操作结束。",
|
||||
)
|
||||
return
|
||||
if self._auto_fetch_running:
|
||||
self._request_stop_auto_fetch()
|
||||
return
|
||||
@@ -1045,6 +1226,19 @@ class PDDTaskPageEvent(QObject):
|
||||
# 避免窗口销毁时出现 "QThread destroyed while running"。
|
||||
thread.wait(60_000)
|
||||
|
||||
resubmit_worker = self._resubmit_worker
|
||||
resubmit_thread = self._resubmit_thread
|
||||
if resubmit_worker is not None:
|
||||
try:
|
||||
resubmit_worker.cancel()
|
||||
resubmit_worker.progress.disconnect(self._on_resubmit_progress)
|
||||
resubmit_worker.finished.disconnect(self._on_resubmit_finished)
|
||||
except (TypeError, RuntimeError):
|
||||
pass
|
||||
if resubmit_thread is not None and resubmit_thread.isRunning():
|
||||
resubmit_thread.quit()
|
||||
resubmit_thread.wait(60_000)
|
||||
|
||||
|
||||
def summary_to_row(summary: TaskSummary) -> TaskRow:
|
||||
"""把领域摘要转换成只供表格显示的轻量行。"""
|
||||
|
||||
Reference in New Issue
Block a user