feat(status): default legacy batches to normal

This commit is contained in:
chengma
2026-07-20 09:41:55 +08:00
parent 00c695d035
commit d687c833e0
17 changed files with 205 additions and 991 deletions
+29
View File
@@ -16,6 +16,8 @@ from .config import make_slug
DEFAULT_BUSY_TIMEOUT_MS = 5000
PRODUCT_STATUS_FEATURE_INTRODUCED_AT = "2026-07-18T16:34:37"
LEGACY_PRODUCT_STATUS_DEFAULT_NOTE = "历史批次默认按架上商品处理(未实时检测)"
VALID_BATCH_FIELDS = {
"source_files_json",
"status",
@@ -493,6 +495,7 @@ def init_db(path=None, conn=None) -> None:
_ensure_task_image_task_columns(database)
_ensure_task_cover_reset_columns(database)
_ensure_task_product_status_columns(database)
_migrate_legacy_product_status_defaults(database)
_ensure_image_studio_project_suite_columns(database)
_ensure_image_studio_project_draft_columns(database)
_ensure_image_studio_job_recovery_columns(database)
@@ -534,6 +537,32 @@ def _ensure_task_product_status_columns(database):
database.execute("ALTER TABLE tasks ADD COLUMN product_status_at TEXT")
def _migrate_legacy_product_status_defaults(database):
"""Treat pre-status-feature active batches as historically on-shelf once."""
database.execute(
"""
UPDATE tasks
SET product_status = 'normal',
product_status_note = ?,
product_status_at = NULL,
updated_at = ?
WHERE (product_status IS NULL OR TRIM(product_status) = '')
AND batch_id IN (
SELECT id
FROM batches
WHERE deleted_at IS NULL
AND created_at < ?
)
""",
(
LEGACY_PRODUCT_STATUS_DEFAULT_NOTE,
_now(),
PRODUCT_STATUS_FEATURE_INTRODUCED_AT,
),
)
def _ensure_image_studio_project_suite_columns(database):
columns = {
row["name"]
-21
View File
@@ -1207,27 +1207,6 @@ def collect(account, task, on_step=None) -> dict:
return result
def recheck_product_status(account, task, on_step=None) -> dict:
"""Read one product status snapshot without changing collected content."""
item_id = _item_id(task)
cdp = open_product(account, item_id, on_step=on_step, bring_to_front=False)
result = None
try:
_notify_collect_step(on_step, "read_product_status")
result = read_product_status(cdp)
if result.get("product_status_error"):
raise EditorError(
"商品状态读取失败:{error}".format(
error=result["product_status_error"]
)
)
finally:
close_target_confirmed = _close_collected_product(cdp)
result["close_target_confirmed"] = close_target_confirmed
return result
def read_product_status(cdp) -> dict:
"""Read the current product's warning state without retaining page HTML."""
-1
View File
@@ -24,7 +24,6 @@ if QT_IMPORT_ERROR is None:
CMHubSettingsWorker,
ApplyWorker,
CollectWorker,
StatusRecheckWorker,
GenerateWorker,
ImageStudioDownloadOriginalWorker,
ImageStudioExportWorker,
+1 -5
View File
@@ -651,10 +651,7 @@ class ApplyTab(QWidget):
if status_text:
lines.append("当前筛选结果没有状态正常且可更新的商品。")
lines.append(status_text)
if update_plan.get("unverified"):
lines.append("请先到①导入采集点击「重新检测商品状态」。")
else:
lines.append("请先回到①导入采集重新确认商品状态。")
lines.append("请先回到①导入采集重新确认商品状态。")
content_error = self._update_content_error(update_plan, update_mode)
if content_error:
if lines:
@@ -730,7 +727,6 @@ class ApplyTab(QWidget):
if not update_plan:
return ""
rows = [
("尚未检测商品状态", update_plan.get("unverified", [])),
("未上架", update_plan.get("unlisted", [])),
("审核中", update_plan.get("reviewing", [])),
("状态未知", update_plan.get("unknown", [])),
+21 -230
View File
@@ -11,26 +11,13 @@ from ...collect_skip import (
from ... import product_status
from ..models import TaskTableModel
from ..widgets import *
from ..workers import (
CollectWorker as _RealCollectWorker,
StatusRecheckWorker as _RealStatusRecheckWorker,
WriteBackWorker as _RealWriteBackWorker,
)
from ..workers import CollectWorker as _RealCollectWorker, WriteBackWorker as _RealWriteBackWorker
def CollectWorker(*args, **kwargs):
return _call_package_attr("CollectWorker", _RealCollectWorker, *args, **kwargs)
def StatusRecheckWorker(*args, **kwargs):
return _call_package_attr(
"StatusRecheckWorker",
_RealStatusRecheckWorker,
*args,
**kwargs,
)
def WriteBackWorker(*args, **kwargs):
return _call_package_attr("WriteBackWorker", _RealWriteBackWorker, *args, **kwargs)
@@ -97,19 +84,15 @@ class CollectTab(QWidget):
self.last_import_stats = None
self.collect_worker = None
self.collect_thread = None
self.status_recheck_worker = None
self.status_recheck_thread = None
self.write_back_worker = None
self.write_back_thread = None
self.last_collect_run_id = None
self.last_status_recheck_run_id = None
self._collect_run_started_at = None
self._collect_task_started_at = None
self._collect_task_elapsed_seconds = 0
self._collect_activity_payload = {}
self._collect_terminal_text = ""
self._collect_stop_requested = False
self._collect_operation_label = "采集"
self._collect_elapsed_timer = QTimer(self)
self._collect_elapsed_timer.setInterval(1000)
self._collect_elapsed_timer.timeout.connect(self._refresh_collect_activity)
@@ -117,8 +100,6 @@ class CollectTab(QWidget):
self.import_button = QPushButton("导入 Excel...")
self.refresh_button = QPushButton("刷新")
self.collect_button = QPushButton("采集旧标题/旧封面")
self.status_recheck_button = QPushButton("重新检测商品状态")
self.status_recheck_button.setObjectName("statusRecheckButton")
self.stop_collect_button = QPushButton("停止")
self.write_back_button = QPushButton("回写旧数据到 Excel")
self.stop_collect_button.setEnabled(False)
@@ -142,7 +123,6 @@ class CollectTab(QWidget):
toolbar.addWidget(self.import_button)
toolbar.addWidget(self.refresh_button)
toolbar.addWidget(self.collect_button)
toolbar.addWidget(self.status_recheck_button)
toolbar.addWidget(self.stop_collect_button)
toolbar.addWidget(self.write_back_button)
toolbar.addStretch(1)
@@ -164,7 +144,7 @@ class CollectTab(QWidget):
self.collect_activity_label = QLabel("")
self.collect_activity_label.setObjectName("collectActivityLabel")
self.collect_activity_label.setAlignment(Qt.AlignRight | Qt.AlignVCenter)
activity_sample = "正在商品状态重检 999/999 · 等待商品页加载 · 本条 99:59"
activity_sample = "正在采集 999/999 · 等待商品页加载 · 本条 99:59"
activity_width = self.collect_activity_label.fontMetrics().horizontalAdvance(activity_sample) + 24
self.collect_activity_label.setFixedWidth(activity_width)
self.collect_activity_label.setVisible(False)
@@ -231,7 +211,6 @@ class CollectTab(QWidget):
self.status_filter.currentIndexChanged.connect(self.refresh_tasks)
self.delete_batch_button.clicked.connect(self.delete_current_batch)
self.collect_button.clicked.connect(self.collect_old_data)
self.status_recheck_button.clicked.connect(self.recheck_product_status)
self.stop_collect_button.clicked.connect(self.stop_collect)
self.write_back_button.clicked.connect(self.write_back_old_data)
self.show_all_button.clicked.connect(self.show_all_tasks)
@@ -253,9 +232,8 @@ class CollectTab(QWidget):
"}"
)
def _start_collect_activity(self, operation_label="采集"):
def _start_collect_activity(self):
now = time.monotonic()
self._collect_operation_label = str(operation_label or "采集")
self._collect_run_started_at = now
self._collect_task_started_at = None
self._collect_task_elapsed_seconds = 0
@@ -328,7 +306,7 @@ class CollectTab(QWidget):
level = "warning"
elif state in {"task_started", "task_step"}:
text = (
f"正在{self._collect_operation_label} {progress} · {step_label} · "
f"正在采集 {progress} · {step_label} · "
f"本条 {_format_collect_elapsed(task_elapsed)}"
)
level = "info"
@@ -344,10 +322,7 @@ class CollectTab(QWidget):
)
level = "danger" if event.get("result") == "failed" else "muted"
else:
if self._collect_operation_label == "采集":
text = f"正在检查账号 · {_format_collect_elapsed(run_elapsed)}"
else:
text = f"正在{self._collect_operation_label}前检查账号 · {_format_collect_elapsed(run_elapsed)}"
text = f"正在检查账号 · {_format_collect_elapsed(run_elapsed)}"
level = "info"
tooltip_parts = []
@@ -376,24 +351,21 @@ class CollectTab(QWidget):
self._collect_stop_requested = False
if outcome == "blocked":
text = f"{self._collect_operation_label}未开始 · 检查未通过"
text = "采集未开始 · 检查未通过"
level = "warning"
tooltip = f"{self._collect_operation_label}前检查未通过,请按弹窗提示处理"
tooltip = "采集前检查未通过,请按弹窗提示处理"
elif outcome == "cancelled":
text = f"{self._collect_operation_label}已停止 · 总用时 {_format_collect_elapsed(total_elapsed)}"
text = f"采集已停止 · 总用时 {_format_collect_elapsed(total_elapsed)}"
level = "warning"
tooltip = f"本轮{self._collect_operation_label}已停止"
tooltip = "本轮采集已停止"
elif outcome == "error":
text = f"{self._collect_operation_label}已结束 · 请查看运行日志"
text = "采集已结束 · 请查看运行日志"
level = "danger"
if self._collect_operation_label == "采集":
tooltip = "采集异常结束,请查看下方采集运行日志"
else:
tooltip = f"{self._collect_operation_label}异常结束,请查看下方运行日志"
tooltip = "采集异常结束,请查看下方采集运行日志"
else:
text = f"{self._collect_operation_label}完成 · 总用时 {_format_collect_elapsed(total_elapsed)}"
text = f"采集完成 · 总用时 {_format_collect_elapsed(total_elapsed)}"
level = "success"
tooltip = f"本轮{self._collect_operation_label}已经完成"
tooltip = "本轮采集已经完成"
self._collect_terminal_text = text
self.collect_activity_label.setToolTip(tooltip)
self._set_collect_activity_style(level)
@@ -407,9 +379,9 @@ class CollectTab(QWidget):
def _append_collect_log(self, message):
self.run_log_view.appendPlainText(str(message))
def _load_latest_collect_run_log(self, run_type="collect"):
def _load_latest_collect_run_log(self):
try:
logs = db.list_run_logs(limit=1, run_type=run_type, path=self.db_path)
logs = db.list_run_logs(limit=1, run_type="collect", path=self.db_path)
if not logs:
return
events = db.list_run_log_events(logs[0].id, limit=30, path=self.db_path)
@@ -436,10 +408,10 @@ class CollectTab(QWidget):
self._set_status(text, level="danger")
QTimer.singleShot(0, lambda: QMessageBox.warning(self, "导入采集", text))
def _show_account_guide(self, message, operation_label="采集"):
def _show_account_guide(self, message):
full_message = (
f"{message}\n\n"
f"本轮{operation_label}已中止。\n"
"本轮采集已中止。\n"
"请先到「账号管理」检查账号配置、Chrome 路径和登录状态。"
)
QMessageBox.warning(self, "账号未就绪", full_message)
@@ -691,11 +663,7 @@ class CollectTab(QWidget):
return self.batch_filter.currentData()
def _update_delete_batch_button(self):
running = bool(
self.collect_thread
or self.status_recheck_thread
or self.write_back_thread
)
running = bool(self.collect_thread or self.write_back_thread)
self.delete_batch_button.setEnabled((not running) and bool(self._selected_batch_id()))
def delete_current_batch(self, checked=False):
@@ -748,9 +716,6 @@ class CollectTab(QWidget):
QMessageBox.information(self, "删除批次", message)
def collect_old_data(self, checked=False):
if self._collect_operation_running() or self.write_back_thread is not None:
self._set_status("当前批处理尚未结束,请稍后再试")
return
tasks = list(self.model.tasks)
if not tasks:
self._set_status("没有可采集任务")
@@ -784,51 +749,6 @@ class CollectTab(QWidget):
self._start_collect_activity()
thread.start()
def recheck_product_status(self, checked=False):
if self._collect_operation_running() or self.write_back_thread is not None:
self._set_status("当前批处理尚未结束,请稍后再试")
return
tasks = list(self.model.tasks)
if not tasks:
self._set_status("当前筛选结果没有可重新检测商品状态的任务")
return
answer = QMessageBox.question(
self,
"重新检测商品状态",
"将重新打开当前筛选结果中的 {count} 条商品详情页,只读取并保存商品状态。\n\n"
"不会读取或覆盖旧标题、旧封面、新标题、新封面;不会改变任务阶段、结果、线上提交标记或 Excel。\n\n"
"确定开始重新检测吗?".format(count=len(tasks)),
QMessageBox.Yes | QMessageBox.No,
QMessageBox.No,
)
if answer != QMessageBox.Yes:
self._set_status("已取消重新检测商品状态")
return
worker = StatusRecheckWorker(
tasks,
db_path=self.db_path,
config=self.config,
diagnostic_log_dir=diagnostics.DEFAULT_LOG_DIR,
)
activity_signal = getattr(worker, "activity", None)
if activity_signal is not None:
activity_signal.connect(self._on_collect_activity)
worker.progress.connect(self._on_status_recheck_progress)
worker.row_updated.connect(self._on_status_recheck_row_updated)
worker.log.connect(self._on_collect_log)
worker.failed.connect(self._on_status_recheck_failed)
worker.finished.connect(self._on_status_recheck_finished)
worker.cancelled.connect(self._on_status_recheck_cancelled)
self.run_log_view.clear()
thread = run_worker(worker, thread_name="StatusRecheckWorker", start=False)
thread.finished.connect(lambda: self._forget_status_recheck_thread(thread))
self.status_recheck_worker = worker
self.status_recheck_thread = thread
self._set_collect_running(True)
self._start_collect_activity("商品状态重检")
self._set_status(f"正在重新检测 {len(tasks)} 条商品状态...")
thread.start()
def _choose_collect_scope(self):
box = ProductStatusScopeDialog(
title="选择采集范围",
@@ -852,12 +772,11 @@ class CollectTab(QWidget):
return box.choice()
def stop_collect(self, checked=False):
worker = self.collect_worker or self.status_recheck_worker
if worker is not None:
worker.cancel()
if self.collect_worker is not None:
self.collect_worker.cancel()
self._collect_stop_requested = True
self._refresh_collect_activity()
self._set_status(f"正在停止{self._collect_operation_label}...")
self._set_status("正在停止采集...")
def write_back_old_data(self, checked=False):
batch_id = self._active_batch_id()
@@ -867,9 +786,6 @@ class CollectTab(QWidget):
self._start_write_back(batch_id)
def _start_write_back(self, batch_id, auto=False):
if self._collect_operation_running():
self._set_status("当前批处理尚未结束,暂不能回写 Excel")
return False
if self.write_back_thread is not None:
self._set_status("Excel 回写正在进行...")
return False
@@ -916,7 +832,6 @@ class CollectTab(QWidget):
self.import_button.setEnabled(not running)
self.refresh_button.setEnabled(not running)
self.collect_button.setEnabled(not running)
self.status_recheck_button.setEnabled(not running)
self.write_back_button.setEnabled(not running)
self.stop_collect_button.setEnabled(running)
self.batch_filter.setEnabled(not running)
@@ -929,7 +844,6 @@ class CollectTab(QWidget):
self.import_button.setEnabled(not running)
self.refresh_button.setEnabled(not running)
self.collect_button.setEnabled(not running)
self.status_recheck_button.setEnabled(not running)
self.write_back_button.setEnabled(not running)
self.batch_filter.setEnabled(not running)
self.shop_filter.setEnabled(not running)
@@ -944,16 +858,6 @@ class CollectTab(QWidget):
self.collect_thread = None
self.collect_worker = None
def _forget_status_recheck_thread(self, thread):
if self.status_recheck_thread is thread:
if self._collect_elapsed_timer.isActive():
self._finish_collect_activity("error")
self.status_recheck_thread = None
self.status_recheck_worker = None
def _collect_operation_running(self):
return bool(self.collect_thread or self.status_recheck_thread)
def _forget_write_back_thread(self, thread):
if self.write_back_thread is thread:
self.write_back_thread = None
@@ -976,25 +880,6 @@ class CollectTab(QWidget):
def _on_collect_failed(self, task_id, error):
self._set_status(f"任务 {task_id} 采集失败:{error}")
def _on_status_recheck_progress(self, payload):
self._set_status(
"商品状态重检进度:{done}/{total},成功{rechecked},略过{skipped},失败{failed}".format(
done=payload.get("done", 0),
total=payload.get("total", 0),
rechecked=payload.get("rechecked", 0),
skipped=payload.get("skipped", 0),
failed=payload.get("failed", 0),
)
)
def _on_status_recheck_row_updated(self, task_id, fields):
self.refresh_tasks()
if self.refresh_workflow_callback is not None:
self.refresh_workflow_callback()
def _on_status_recheck_failed(self, task_id, error):
self._set_status(f"任务 {task_id} 商品状态重检失败:{error}")
def _on_collect_finished(self, payload):
self._set_collect_running(False)
if payload.get("blocked"):
@@ -1043,40 +928,6 @@ class CollectTab(QWidget):
return
self._set_status(message)
def _on_status_recheck_finished(self, payload):
self._set_collect_running(False)
if payload.get("blocked"):
self._finish_collect_activity("blocked", payload)
elif payload.get("error"):
self._finish_collect_activity("error", payload)
else:
self._finish_collect_activity("finished", payload)
self.last_status_recheck_run_id = (
payload.get("run_id") or self.last_status_recheck_run_id
)
self.refresh_tasks()
if self.refresh_workflow_callback is not None:
self.refresh_workflow_callback()
self._load_latest_collect_run_log("status_recheck")
if payload.get("blocked"):
self._show_status_recheck_blocked(payload)
return
message = "商品状态重检完成:成功{rechecked},略过{skipped},失败{failed}".format(
rechecked=payload.get("rechecked", 0),
skipped=payload.get("skipped", 0),
failed=payload.get("failed", 0),
)
counts = payload.get("product_status_counts") or {}
if counts:
message += ";正常{normal},未上架{unlisted},审核中{reviewing},状态未知{unknown}".format(
normal=counts.get("normal", 0),
unlisted=counts.get("unlisted", 0),
reviewing=counts.get("reviewing", 0),
unknown=counts.get("unknown", 0),
)
self._show_status_recheck_account_summary(payload, message)
self._set_status(message)
def _show_collect_blocked(self, payload):
lines = ["采集前检查未通过。"]
if payload.get("no_accounts"):
@@ -1101,18 +952,6 @@ class CollectTab(QWidget):
)
self._show_account_guide("\n".join(lines))
def _show_status_recheck_blocked(self, payload):
lines = ["商品状态重检前检查未通过。"]
if payload.get("no_accounts"):
lines.append("当前没有配置账号。")
launch_failed = payload.get("launch_failed") or []
if launch_failed:
lines.append(
"以下账号 Chrome 启动失败:"
+ "、".join(self._account_label(item) for item in launch_failed)
)
self._show_account_guide("\n".join(lines), operation_label="商品状态重检")
def _show_collect_account_summary(self, payload, message):
launched = payload.get("launched_accounts") or []
reused = payload.get("reused_accounts") or []
@@ -1163,41 +1002,6 @@ class CollectTab(QWidget):
else:
QMessageBox.information(self, "采集完成", text)
def _show_status_recheck_account_summary(self, payload, message):
launched = payload.get("launched_accounts") or []
reused = payload.get("reused_accounts") or []
login_required = payload.get("login_required_accounts") or []
skip_counts = normalize_skip_reason_counts(
payload.get("skip_reason_counts"),
skipped_total=int(payload.get("skipped", 0) or 0),
)
lines = [message]
skip_summary = format_skip_reason_summary(
skip_counts,
skipped_total=int(payload.get("skipped", 0) or 0),
)
if skip_summary:
lines.append(skip_summary)
if launched:
lines.append(
"本轮已自动启动账号 Chrome:"
+ "、".join(self._account_label(item) for item in launched)
)
if login_required:
lines.append(
"以下账号需要补登录:"
+ "、".join(self._account_label(item) for item in login_required)
)
if login_required or skip_counts[LOGIN_REQUIRED] > 0:
lines.append("请到账号管理完成对应账号登录后,再重新检测商品状态。")
if launched or reused:
lines.append("检测结束后不会自动关闭账号 Chrome,请按需自行关闭。")
text = "\n".join(lines)
if login_required or skip_counts[LOGIN_REQUIRED] > 0:
QMessageBox.warning(self, "商品状态重检完成", text)
else:
QMessageBox.information(self, "商品状态重检完成", text)
def _account_label(self, item):
if isinstance(item, dict):
name = item.get("account_name") or item.get("alias") or ""
@@ -1221,19 +1025,6 @@ class CollectTab(QWidget):
)
)
def _on_status_recheck_cancelled(self, payload):
self._set_collect_running(False)
self._finish_collect_activity("cancelled", payload)
self.refresh_tasks()
if self.refresh_workflow_callback is not None:
self.refresh_workflow_callback()
self._set_status(
"商品状态重检已停止:完成{done}/{total}".format(
done=payload.get("done", 0),
total=payload.get("total", 0),
)
)
def _on_write_back_failed(self, task_id, error, auto=False):
message = f"Excel {'自动' if auto else ''}回写失败:{error}"
if "被占用" in str(error):
-377
View File
@@ -2738,383 +2738,6 @@ class CollectWorker(BaseWorker):
def _elapsed_ms(self, started):
return int((time.monotonic() - started) * 1000)
class StatusRecheckWorker(CollectWorker):
"""Re-read product status snapshots without changing task workflow fields."""
def __init__(
self,
tasks,
db_path=None,
config=None,
preflight=True,
diagnostic_log_dir=None,
):
super().__init__(
tasks,
db_path=db_path,
config=config,
preflight=preflight,
diagnostic_log_dir=diagnostic_log_dir,
collect_scope=product_status.COLLECT_SCOPE_ALL,
)
def execute(self):
account_rows = accounts.list_accounts(path=self.db_path, config=self.config)
account_by_alias = {
str(account.alias).strip(): account
for account in account_rows
if str(account.alias).strip()
}
eligible = list(self.tasks)
batch_ids = self._batch_ids(eligible)
total = len(eligible)
rechecked = 0
skipped = 0
failed = 0
done = 0
login_skip_reasons = {}
login_required_accounts = {}
preflight_info = {}
skip_reason_counts = empty_skip_reason_counts()
product_status_counts = {
status: 0 for status in product_status.VALID_PRODUCT_STATUSES
}
self._run_id = self._create_run_log(eligible, batch_ids)
self._emit_activity("preflight_started", total=total, step="preflight")
self._log_run_event(
f"step=preflight result=start detail=商品状态重检开始 total={total}"
)
if self.preflight:
blocked, preflight_info = self._preflight_prepare(
eligible,
account_rows,
account_by_alias,
)
if blocked:
self._log_preflight_blocked(blocked)
summary = self._summary(
ok=False,
total=total,
done=done,
collected=rechecked,
skipped=skipped,
failed=failed,
batch_ids=batch_ids,
blocked=True,
extra={
**blocked,
"rechecked": rechecked,
"skip_reason_counts": dict(skip_reason_counts),
"product_status_counts": dict(product_status_counts),
},
)
self._finish_run_log("blocked", summary)
return summary
for item in preflight_info.get("logged_out") or []:
alias = str(item.get("alias") or "").strip()
reason = item.get("reason") or "账号未登录"
if alias:
login_skip_reasons[alias] = reason
login_required_accounts[alias] = item
self._log_run_event("step=preflight result=success detail=账号就绪检查完成")
else:
self._log_run_event(
"step=preflight result=skipped detail=测试模式跳过商品状态重检前检查",
level="warning",
)
for index, task in enumerate(eligible, start=1):
if self.should_cancel():
break
self._emit_activity(
"task_started",
task=task,
index=index,
total=total,
step="match_account",
)
account = account_by_alias.get(str(task.alias).strip())
if account is None:
skipped += 1
done += 1
skip_reason_counts[ALIAS_UNMATCHED] += 1
reason = "别名未匹配账号"
self._log_run_event(
"step=match_account result=skipped detail=任务 {task_id} 商品 {item_id} {reason}".format(
task_id=task.id,
item_id=task.item_id,
reason=reason,
),
task=task,
level="warning",
)
self._emit_activity(
"task_finished",
task=task,
index=index,
total=total,
step="match_account",
result="skipped",
)
self._emit_recheck_progress(done, total, rechecked, skipped, failed)
continue
alias = str(task.alias).strip()
if alias in login_skip_reasons:
skipped += 1
done += 1
skip_reason_counts[LOGIN_REQUIRED] += 1
reason = login_skip_reasons[alias]
self._log_run_event(
"step=login_check result=skipped detail=任务 {task_id} 商品 {item_id} {reason}".format(
task_id=task.id,
item_id=task.item_id,
reason=reason,
),
task=task,
level="warning",
)
self._emit_activity(
"task_finished",
task=task,
index=index,
total=total,
step="check_login",
result="skipped",
)
self._emit_recheck_progress(done, total, rechecked, skipped, failed)
continue
self._emit_activity(
"task_step",
task=task,
index=index,
total=total,
step="check_login",
)
status = self._confirmed_login_status(account, context="status_recheck", task=task)
if self._is_definitive_logged_out(status):
skipped += 1
done += 1
skip_reason_counts[LOGIN_REQUIRED] += 1
reason = self._midrun_login_skip_reason(status)
login_skip_reasons[alias] = reason
login_required_accounts[alias] = self._account_payload(account, reason)
self._log_run_event(
"step=login_check result=skipped detail=任务 {task_id} 商品 {item_id} {reason}".format(
task_id=task.id,
item_id=task.item_id,
reason=reason,
),
task=task,
level="warning",
)
self._emit_activity(
"task_finished",
task=task,
index=index,
total=total,
step="check_login",
result="skipped",
)
self._emit_recheck_progress(done, total, rechecked, skipped, failed)
continue
if not status.get("logged_in"):
self._log_run_event(
"step=login_check result=uncertain detail=任务 {task_id} 商品 {item_id} 登录状态检测暂时不稳定,继续尝试读取商品状态: {detail}".format(
task_id=task.id,
item_id=task.item_id,
detail=self._login_status_detail(status),
),
task=task,
level="warning",
)
started = time.monotonic()
current_step = "read_product_status"
activity_result = "success"
def on_step(step):
nonlocal current_step
current_step = str(step)
self._emit_activity(
"task_step",
task=task,
index=index,
total=total,
step=current_step,
)
self._log_run_event(
"step={step} result=start detail=任务 {task_id} 商品 {item_id}".format(
step=current_step,
task_id=task.id,
item_id=task.item_id,
),
task=task,
)
try:
result = editor.recheck_product_status(
account,
{"item_id": task.item_id},
on_step=on_step,
)
detected_status = product_status.normalize_status(
result.get("product_status")
)
product_status_counts[detected_status] += 1
current_step = "db_write"
self._emit_activity(
"task_step",
task=task,
index=index,
total=total,
step="save_result",
)
self._log_run_event(
"step=db_write result=start detail=任务 {task_id} 商品 {item_id} 保存商品状态".format(
task_id=task.id,
item_id=task.item_id,
),
task=task,
)
db.set_product_status(
task.id,
detected_status,
result.get("product_status_note"),
path=self.db_path,
)
rechecked += 1
elapsed_ms = self._elapsed_ms(started)
self.row_updated.emit(
task.id,
{
"product_status": detected_status,
"product_status_note": result.get("product_status_note"),
},
)
self._log_run_event(
"step=db_write result=success detail=任务 {task_id} 商品 {item_id} 商品状态重检成功 elapsed_ms={elapsed_ms}".format(
task_id=task.id,
item_id=task.item_id,
elapsed_ms=elapsed_ms,
),
task=task,
)
if result.get("close_target_confirmed") is False:
self._log_run_event(
"step=close_product result=uncertain detail=任务 {task_id} 商品 {item_id} 商品页已请求关闭,但未在短时间内确认关闭;状态结果已保存,继续处理后续任务".format(
task_id=task.id,
item_id=task.item_id,
),
task=task,
level="warning",
)
except Exception as exc:
activity_result = "failed"
failed += 1
error = str(exc) or exc.__class__.__name__
safe_error = diagnostics.redact_log_text(error)
display_error = db.format_failure_error(safe_error, current_step)
elapsed_ms = self._elapsed_ms(started)
self.failed.emit(task.id, display_error)
self._log_run_event(
"step={step} result=failed detail=商品状态重检失败: {error} elapsed_ms={elapsed_ms}".format(
step=current_step,
error=safe_error,
elapsed_ms=elapsed_ms,
),
task=task,
level="error",
)
self._write_diagnostic_log(
"商品状态重检失败",
level="ERROR",
step=current_step,
task=task,
elapsed_ms=elapsed_ms,
payload={"error": safe_error},
exc=exc,
)
finally:
done += 1
self._emit_activity(
"task_finished",
task=task,
index=index,
total=total,
step=current_step,
result=activity_result,
)
self._emit_recheck_progress(done, total, rechecked, skipped, failed)
summary = self._summary(
ok=failed == 0,
total=total,
done=done,
collected=rechecked,
skipped=skipped,
failed=failed,
batch_ids=batch_ids,
extra={
**preflight_info,
"rechecked": rechecked,
"login_required_accounts": list(login_required_accounts.values()),
"skip_reason_counts": dict(skip_reason_counts),
"product_status_counts": dict(product_status_counts),
},
)
self._finish_run_log("cancelled" if self.should_cancel() else "done", summary)
return summary
def _emit_recheck_progress(self, done, total, rechecked, skipped, failed):
self.progress.emit(
{
"done": done,
"total": total,
"rechecked": rechecked,
"skipped": skipped,
"failed": failed,
}
)
def _create_run_log(self, eligible, batch_ids):
try:
return db.create_run_log(
"status_recheck",
dry_run=False,
total=len(eligible),
options={
"batch_ids": batch_ids,
"preflight": self.preflight,
"scope": "current_filters",
},
path=self.db_path,
)
except Exception:
return None
def _finish_run_log(self, status, summary):
if self._run_id is None:
return
try:
db.finish_run_log(
self._run_id,
status=status,
done=summary.get("done", 0),
success_count=summary.get("rechecked", 0),
skipped_count=summary.get("skipped", 0),
failed_count=summary.get("failed", 0),
summary_json=summary,
path=self.db_path,
)
except Exception:
return
class WriteBackWorker(BaseWorker):
"""Write Excel fields back in a background thread."""
-8
View File
@@ -138,7 +138,6 @@ def build_apply_plan(tasks, update_mode) -> dict:
STATUS_UNLISTED: [],
STATUS_REVIEWING: [],
STATUS_UNKNOWN: [],
"unverified": [],
}
normal_candidates = []
for task in base_candidates:
@@ -146,8 +145,6 @@ def build_apply_plan(tasks, update_mode) -> dict:
status_counts[status] += 1
if status == STATUS_NORMAL:
normal_candidates.append(task)
elif status == STATUS_UNKNOWN and not _has_status_snapshot(task):
status_excluded["unverified"].append(task)
else:
status_excluded[status].append(task)
@@ -177,7 +174,6 @@ def build_apply_plan(tasks, update_mode) -> dict:
"unlisted": status_excluded[STATUS_UNLISTED],
"reviewing": status_excluded[STATUS_REVIEWING],
"unknown": status_excluded[STATUS_UNKNOWN],
"unverified": status_excluded["unverified"],
"missing_title": missing_title,
"missing_cover": missing_cover,
"status_scope_excluded": status_excluded_count,
@@ -243,10 +239,6 @@ def _task_status(task) -> str:
return normalize_status(_task_value(task, "product_status"))
def _has_status_snapshot(task) -> bool:
return bool(str(_task_value(task, "product_status") or "").strip())
def _task_value(task, name, default=None):
if isinstance(task, dict):
return task.get(name, default)