feat: auto-prepare chrome for collection

This commit is contained in:
chengma
2026-07-08 17:25:17 +08:00
parent 574074ebfb
commit 95170ec4aa
8 changed files with 325 additions and 53 deletions
+94 -29
View File
@@ -1218,6 +1218,9 @@ class CollectWorker(BaseWorker):
skipped = 0
failed = 0
done = 0
login_skip_reasons = {}
login_required_accounts = {}
preflight_info = {}
self._run_id = self._create_run_log(eligible, batch_ids)
self._log_run_event(
@@ -1225,7 +1228,7 @@ class CollectWorker(BaseWorker):
)
if self.preflight:
blocked = self._preflight_block(eligible, account_rows, account_by_alias)
blocked, preflight_info = self._preflight_prepare(eligible, account_rows, account_by_alias)
if blocked:
self._log_preflight_blocked(blocked)
summary = self._summary(
@@ -1241,7 +1244,13 @@ class CollectWorker(BaseWorker):
)
self._finish_run_log("blocked", summary)
return summary
self._log_run_event("step=preflight result=success detail=账号检查通过")
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=测试模式跳过采集前检查",
@@ -1270,15 +1279,37 @@ class CollectWorker(BaseWorker):
self._emit_progress(done, total, collected, skipped, failed)
continue
status = self._login_status(account)
if not status.get("logged_in"):
alias = str(task.alias).strip()
if alias in login_skip_reasons:
skipped += 1
done += 1
reason = self._login_skip_reason(status)
reason = login_skip_reasons[alias]
db.mark_skipped(task.id, reason, path=self.db_path)
self.row_updated.emit(task.id, {"status": "skipped", "last_error": reason})
self._log_run_event(
"step=preflight result=skipped detail=任务 {task_id} 商品 {item_id} {reason}".format(
"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_progress(done, total, collected, skipped, failed)
continue
status = self._login_status(account)
if not status.get("logged_in"):
alias = str(task.alias).strip()
skipped += 1
done += 1
reason = self._midrun_login_skip_reason(status)
login_skip_reasons[alias] = reason
login_required_accounts[alias] = self._account_payload(account, reason)
db.mark_skipped(task.id, reason, path=self.db_path)
self.row_updated.emit(task.id, {"status": "skipped", "last_error": 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,
@@ -1394,16 +1425,23 @@ class CollectWorker(BaseWorker):
skipped=skipped,
failed=failed,
batch_ids=batch_ids,
extra={
**preflight_info,
"login_required_accounts": list(login_required_accounts.values()),
},
)
self._finish_run_log("cancelled" if self.should_cancel() else "done", summary)
return summary
def _preflight_block(self, eligible, account_rows, account_by_alias):
def _preflight_prepare(self, eligible, account_rows, account_by_alias):
if not account_rows:
return {
"reason": "NO_ACCOUNTS",
"no_accounts": True,
}
return (
{
"reason": "NO_ACCOUNTS",
"no_accounts": True,
},
{},
)
required_accounts = []
seen_aliases = set()
for task in eligible:
@@ -1412,24 +1450,40 @@ class CollectWorker(BaseWorker):
if account is not None and alias not in seen_aliases:
required_accounts.append(account)
seen_aliases.add(alias)
not_running = []
launch_failed = []
logged_out = []
launched = []
reused = []
for account in required_accounts:
self._log_run_event(
f"step=check_chrome result=start detail=账号 {account.alias} debug_port={account.debug_port}",
f"step=ensure_chrome result=start detail=账号 {account.alias} debug_port={account.debug_port}",
level="info",
)
if not chrome.is_running(account.debug_port):
if chrome.is_running(account.debug_port):
self._log_run_event(
f"step=check_chrome result=blocked detail=账号 {account.alias} CDP 端口未响应 debug_port={account.debug_port}",
level="warning",
f"step=ensure_chrome result=reused detail=账号 {account.alias} Chrome 已打开,复用现有窗口 debug_port={account.debug_port}",
level="info",
)
not_running.append(self._account_payload(account, "CDP 端口未响应"))
continue
self._log_run_event(
f"step=check_chrome result=success detail=账号 {account.alias} debug_port={account.debug_port}",
level="info",
)
reused.append(self._account_payload(account, "已复用"))
else:
try:
result = accounts.launch_for_login(account, path=self.db_path, config=self.config)
except Exception as exc:
reason = diagnostics.redact_log_text(str(exc) or exc.__class__.__name__)
self._log_run_event(
f"step=ensure_chrome result=blocked detail=账号 {account.alias} Chrome 启动失败: {reason}",
level="error",
)
launch_failed.append(self._account_payload(account, f"Chrome 启动失败: {reason}"))
continue
action = "launched" if result.get("launched") else "reused"
detail = "已启动" if action == "launched" else "已复用"
self._log_run_event(
f"step=ensure_chrome result={action} detail=账号 {account.alias} {detail} debug_port={account.debug_port}",
level="info",
)
target = launched if result.get("launched") else reused
target.append(self._account_payload(account, detail))
self._log_run_event(
f"step=login_check result=start detail=账号 {account.alias}",
level="info",
@@ -1449,13 +1503,20 @@ class CollectWorker(BaseWorker):
f"step=login_check result=success detail=账号 {account.alias}",
level="info",
)
if not_running or logged_out:
return {
"reason": "ACCOUNT_NOT_READY",
"not_running": not_running,
"logged_out": logged_out,
}
return None
info = {
"launched_accounts": launched,
"reused_accounts": reused,
"logged_out": logged_out,
}
if launch_failed:
return (
{
"reason": "CHROME_LAUNCH_FAILED",
"launch_failed": launch_failed,
},
info,
)
return None, info
def _account_payload(self, account, reason=None):
payload = {
@@ -1491,6 +1552,10 @@ class CollectWorker(BaseWorker):
reason = status.get("reason")
return f"账号未登录: {reason}" if reason else "账号未登录"
def _midrun_login_skip_reason(self, status):
reason = self._login_skip_reason(status)
return f"采集中途掉登录: {reason}"
def _old_cover_path(self, account, task):
image_root = appconfig.image_dir(self.config)
return image_paths.task_image_path(image_root, task, account, "old")