feat: add diagnostic logs for generation flows

This commit is contained in:
chengma
2026-06-29 17:49:45 +08:00
parent 588b352eac
commit cd1256bc5e
17 changed files with 1455 additions and 104 deletions
+128 -22
View File
@@ -6,12 +6,13 @@ import copy
import json
import mimetypes
import os
import threading
import time
import urllib.error
import urllib.request
import uuid
from . import appconfig, db
from . import appconfig, db, diagnostics
from . import prompts as prompt_module
from .config import make_slug
@@ -34,12 +35,15 @@ def gen_title(
retry=None,
config=None,
models_path=appconfig.AI_MODELS_PATH,
on_step=None,
):
"""Generate a new product title from a prompt and the old title."""
cfg = appconfig.load_config() if config is None else config
ai_cfg = appconfig.ai_config(cfg)
_notify_step(on_step, "load_text_model")
model = _role_model("text", ai_cfg.get("default_text_model"), models_path)
_notify_step(on_step, "title_build_request")
payload = _chat_payload(
model,
[
@@ -51,6 +55,7 @@ def gen_title(
],
)
attempts = _attempt_count(ai_cfg, retry)
_notify_step(on_step, "title_request")
data = _call_with_retry(
model,
payload,
@@ -58,6 +63,7 @@ def gen_title(
attempts,
request_kind="json",
)
_notify_step(on_step, "title_parse_response")
text = _extract_text(data).strip()
if not text:
raise AIError("AI 返回为空标题")
@@ -73,9 +79,11 @@ def gen_cover(
retry=None,
config=None,
models_path=appconfig.AI_MODELS_PATH,
on_step=None,
):
"""Generate a new cover image and save it as a JPEG file."""
_notify_step(on_step, "cover_validate_input")
old_cover_path = os.path.abspath(str(old_cover_path))
if not os.path.exists(old_cover_path):
raise FileNotFoundError("旧封面图片不存在: %s" % old_cover_path)
@@ -84,14 +92,17 @@ def gen_cover(
cfg = appconfig.load_config() if config is None else config
ai_cfg = appconfig.ai_config(cfg)
_notify_step(on_step, "load_image_model")
model = _role_model("image", ai_cfg.get("default_image_model"), models_path)
resolution = str(resolution or ai_cfg.get("resolution", "1k"))
quality = _jpg_quality(jpg_quality if jpg_quality is not None else ai_cfg.get("jpg_quality", 90))
attempts = _attempt_count(ai_cfg, retry)
api_type = model.get("api_type", "auto")
_notify_step(on_step, "cover_build_request")
if api_type == "images_edits":
body, content_type = _image_edit_body(model, cover_prompt, old_cover_path, resolution)
_notify_step(on_step, "cover_request")
data = _call_with_retry(
model,
body,
@@ -102,12 +113,24 @@ def gen_cover(
)
else:
payload = _image_chat_payload(model, cover_prompt, old_cover_path, resolution)
_notify_step(on_step, "cover_request")
data = _call_with_retry(model, payload, cfg, attempts, request_kind="json")
_notify_step(on_step, "cover_parse_response")
image_bytes = _extract_image_bytes(data, model, cfg)
_notify_step(on_step, "cover_save")
return _save_jpeg(image_bytes, out_path, resolution, quality)
def _notify_step(callback, step):
if callback is None:
return
try:
callback(step)
except Exception:
pass
def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=None):
"""Generate titles first, then covers, and persist each successful task."""
@@ -132,6 +155,8 @@ def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=No
image_root = runtime.get("image_dir") or appconfig.image_dir(config)
account_by_alias = runtime.get("account_by_alias") or {}
on_task_update = runtime.get("on_task_update")
on_event = runtime.get("on_event")
on_error = runtime.get("on_error")
title_prompt = _prompt_value(prompts, "title")
cover_prompt = _prompt_value(prompts, "cover")
should_stop = should_stop or (lambda: False)
@@ -149,6 +174,22 @@ def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=No
}
_emit_generation_progress(on_progress, summary)
title_results = {}
step_by_task = {}
step_lock = threading.Lock()
def set_step(task, step):
with step_lock:
step_by_task[getattr(task, "id", None)] = str(step)
def get_step(task, fallback):
with step_lock:
return step_by_task.get(getattr(task, "id", None), fallback)
def step_callback(task, phase):
def callback(step):
set_step(task, step)
_emit_generation_event(on_event, task, phase, step, "start")
return callback
with ThreadPoolExecutor(
max_workers=max(1, int(generation_cfg.get("title_concurrency", 1)))
@@ -158,6 +199,8 @@ def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=No
if should_stop():
summary["cancelled"] = True
break
set_step(task, "title_submit")
_emit_generation_event(on_event, task, "title", "title_submit", "start")
futures[
executor.submit(
gen_title,
@@ -166,6 +209,7 @@ def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=No
retry=generation_cfg.get("retry"),
config=config,
models_path=models_path,
on_step=step_callback(task, "title"),
)
] = task
for future in as_completed(futures):
@@ -176,12 +220,18 @@ def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=No
try:
title_results[task.id] = future.result()
summary["title_done"] += 1
set_step(task, "title_done")
_emit_generation_event(on_event, task, "title", "title_done", "success")
except CancelledError:
summary["cancelled"] = True
_emit_generation_event(on_event, task, "title", get_step(task, "title_request"), "cancelled", level="warning")
except Exception as exc:
summary["failed"] += 1
summary["ok"] = False
_mark_generate_failed(task, exc, db_path, on_task_update)
step = get_step(task, "title_request")
error = _mark_generate_failed(task, exc, db_path, on_task_update)
_emit_generation_event(on_event, task, "title", step, "failed", detail=error, level="error")
_emit_generation_error(on_error, task, "title", step, exc, error)
_emit_generation_progress(on_progress, summary)
cover_tasks = [
@@ -197,23 +247,38 @@ def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=No
summary["cancelled"] = True
break
new_title = title_results[task.id]
rendered_cover_prompt = prompt_module.render_prompt(
cover_prompt,
_prompt_context(task, new_title, account_by_alias),
)
futures[
executor.submit(
gen_cover,
rendered_cover_prompt,
getattr(task, "old_cover_path", "") or "",
_new_cover_path(task, account_by_alias, image_root),
resolution=generation_cfg.get("resolution"),
jpg_quality=generation_cfg.get("jpg_quality"),
retry=generation_cfg.get("retry"),
config=config,
models_path=models_path,
try:
set_step(task, "cover_prompt_render")
_emit_generation_event(on_event, task, "cover", "cover_prompt_render", "start")
rendered_cover_prompt = prompt_module.render_prompt(
cover_prompt,
_prompt_context(task, new_title, account_by_alias),
)
] = (task, new_title)
_emit_generation_event(on_event, task, "cover", "cover_prompt_render", "success")
set_step(task, "cover_submit")
_emit_generation_event(on_event, task, "cover", "cover_submit", "start")
futures[
executor.submit(
gen_cover,
rendered_cover_prompt,
getattr(task, "old_cover_path", "") or "",
_new_cover_path(task, account_by_alias, image_root),
resolution=generation_cfg.get("resolution"),
jpg_quality=generation_cfg.get("jpg_quality"),
retry=generation_cfg.get("retry"),
config=config,
models_path=models_path,
on_step=step_callback(task, "cover"),
)
] = (task, new_title)
except Exception as exc:
summary["failed"] += 1
summary["ok"] = False
step = get_step(task, "cover_prompt_render")
error = _mark_generate_failed(task, exc, db_path, on_task_update)
_emit_generation_event(on_event, task, "cover", step, "failed", detail=error, level="error")
_emit_generation_error(on_error, task, "cover", step, exc, error)
_emit_generation_progress(on_progress, summary)
for future in as_completed(futures):
task, new_title = futures[future]
if should_stop():
@@ -221,6 +286,8 @@ def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=No
_cancel_pending(futures)
try:
new_cover_path = future.result()
set_step(task, "db_write")
_emit_generation_event(on_event, task, "cover", "db_write", "start")
db.set_generated(task.id, new_title, new_cover_path, path=db_path)
summary["cover_done"] += 1
if on_task_update is not None:
@@ -233,19 +300,23 @@ def generate_batch(tasks, prompts, ai_cfg=None, on_progress=None, should_stop=No
"new_cover_path": new_cover_path,
},
)
_emit_generation_event(on_event, task, "cover", "db_write", "success")
except CancelledError:
summary["cancelled"] = True
_emit_generation_event(on_event, task, "cover", get_step(task, "cover_request"), "cancelled", level="warning")
except Exception as exc:
summary["failed"] += 1
summary["ok"] = False
_mark_generate_failed(task, exc, db_path, on_task_update)
step = get_step(task, "cover_request")
error = _mark_generate_failed(task, exc, db_path, on_task_update)
_emit_generation_event(on_event, task, "cover", step, "failed", detail=error, level="error")
_emit_generation_error(on_error, task, "cover", step, exc, error)
_emit_generation_progress(on_progress, summary)
if summary["cancelled"]:
summary["ok"] = False
return summary
def _role_model(category, name, models_path):
if not name:
raise AIError("未配置默认 %s 模型" % category)
@@ -332,10 +403,45 @@ def _new_cover_path(task, account_by_alias, image_root):
def _mark_generate_failed(task, exc, db_path, on_task_update):
error = str(exc) or exc.__class__.__name__
error = diagnostics.redact_log_text(str(exc) or exc.__class__.__name__)
db.mark_failed(task.id, "generate", error, path=db_path)
if on_task_update is not None:
on_task_update(task.id, {"status": "failed", "last_error": error})
return error
def _emit_generation_event(callback, task, phase, step, result, detail=None, level="info"):
if callback is None:
return
payload = {
"task": task,
"phase": phase,
"step": str(step),
"result": result,
"level": level,
}
if detail is not None:
payload["detail"] = diagnostics.redact_log_text(detail)
try:
callback(payload)
except Exception:
return
def _emit_generation_error(callback, task, phase, step, exc, error):
if callback is None:
return
payload = {
"task": task,
"phase": phase,
"step": str(step),
"error": diagnostics.redact_log_text(error),
"exception": exc,
}
try:
callback(payload)
except Exception:
return
def _cancel_pending(futures):
@@ -401,7 +507,7 @@ def _call_once(model, body, config, request_kind, content_type=None):
else:
data = body
request = urllib.request.Request(
model["url"],
appconfig.model_request_url(model),
data=data,
headers=_headers(model, content_type),
method="POST",
+28 -1
View File
@@ -425,6 +425,33 @@ def get_model(name, path=AI_MODELS_PATH) -> dict:
return copy.deepcopy(models[_model_index(models, name)])
def model_request_url(model) -> str:
"""Return the HTTP endpoint used for a configured AI model."""
raw_url = str(model.get("url", "") or "").strip()
api_type = str(model.get("api_type", "auto") or "auto").strip()
if api_type == "images_edits":
return _append_default_endpoint(raw_url, "images/edits")
return _append_default_endpoint(raw_url, "chat/completions")
def _append_default_endpoint(raw_url, endpoint):
if not raw_url:
return raw_url
parts = urllib.parse.urlsplit(raw_url)
path = parts.path.rstrip("/")
lowered = path.lower()
endpoint_path = "/" + endpoint.strip("/")
if lowered.endswith(endpoint_path):
return raw_url
base_markers = ("", "/v1", "/v1beta", "/api/v1", "/api/v1beta")
if lowered in base_markers or lowered.endswith(base_markers[1:]):
path = path + endpoint_path
return urllib.parse.urlunsplit(
(parts.scheme, parts.netloc, path, parts.query, parts.fragment)
)
return raw_url
def _test_request_payload(model):
if model["api_type"] == "images_edits":
payload = {"model": model["model"], "prompt": "ping"}
@@ -455,7 +482,7 @@ def test_ai_model(name, path=AI_MODELS_PATH) -> dict:
"utf-8"
)
request = urllib.request.Request(
model["url"],
model_request_url(model),
data=data,
headers={
"Authorization": "Bearer " + model["api_key"],
+118
View File
@@ -0,0 +1,118 @@
"""Diagnostic logging helpers for local debug logs."""
from __future__ import annotations
import datetime as _dt
import json
import os
import re
import traceback
from . import appconfig
DEFAULT_LOG_DIR = "logs"
DEFAULT_LOG_FILE = "cmshopee.log"
DEFAULT_MAX_BYTES = 2 * 1024 * 1024
DEFAULT_BACKUPS = 3
_SECRET_ASSIGNMENT_RE = re.compile(
r"(?i)\b(api[_-]?key|apikey|token|password|cookie|authorization)\b\s*[:=]\s*([^\s,;]+)"
)
_BEARER_RE = re.compile(r"(?i)\bBearer\s+[A-Za-z0-9._\-]+")
def _redact_text(value):
text = str(value)
text = _SECRET_ASSIGNMENT_RE.sub(lambda match: f"{match.group(1)}=***", text)
return _BEARER_RE.sub("Bearer ***", text)
def _sanitize_log_value(value):
value = appconfig.sanitize_for_log(value)
if isinstance(value, dict):
return {key: _sanitize_log_value(child) for key, child in value.items()}
if isinstance(value, list):
return [_sanitize_log_value(item) for item in value]
if isinstance(value, tuple):
return tuple(_sanitize_log_value(item) for item in value)
if isinstance(value, str):
return _redact_text(value)
return value
def redact_log_text(value):
"""Return free-form diagnostic text with common secret assignments redacted."""
return _redact_text(value)
def diagnostic_log_path(log_dir=None, filename=DEFAULT_LOG_FILE):
directory = os.path.abspath(log_dir or DEFAULT_LOG_DIR)
return os.path.join(directory, filename)
def write_diagnostic_log(
message,
*,
level="INFO",
step=None,
task_id=None,
alias=None,
item_id=None,
elapsed_ms=None,
payload=None,
exc=None,
log_dir=None,
log_path=None,
max_bytes=DEFAULT_MAX_BYTES,
backups=DEFAULT_BACKUPS,
):
"""Append a sanitized diagnostic log entry and return the log path."""
path = os.path.abspath(log_path or diagnostic_log_path(log_dir))
os.makedirs(os.path.dirname(path), exist_ok=True)
_rotate_if_needed(path, max_bytes=max_bytes, backups=backups)
entry = {
"time": _dt.datetime.now().isoformat(timespec="seconds"),
"level": str(level or "INFO").upper(),
"message": _redact_text(message),
}
if step is not None:
entry["step"] = _redact_text(step)
if task_id is not None:
entry["task_id"] = task_id
if alias is not None:
entry["alias"] = _redact_text(alias)
if item_id is not None:
entry["item_id"] = _redact_text(item_id)
if elapsed_ms is not None:
entry["elapsed_ms"] = int(elapsed_ms)
if payload is not None:
entry["payload"] = _sanitize_log_value(payload)
if exc is not None:
entry["exception"] = exc.__class__.__name__
tb = "".join(traceback.format_exception(type(exc), exc, exc.__traceback__))
entry["traceback"] = _redact_text(appconfig.redact_secrets(tb))
safe_entry = _sanitize_log_value(entry)
with open(path, "a", encoding="utf-8") as fh:
fh.write(json.dumps(safe_entry, ensure_ascii=False, sort_keys=True))
fh.write("\n")
return path
def _rotate_if_needed(path, max_bytes=DEFAULT_MAX_BYTES, backups=DEFAULT_BACKUPS):
if max_bytes <= 0 or backups <= 0:
return
if not os.path.exists(path) or os.path.getsize(path) < max_bytes:
return
oldest = f"{path}.{int(backups)}"
if os.path.exists(oldest):
os.remove(oldest)
for index in range(int(backups) - 1, 0, -1):
source = f"{path}.{index}"
target = f"{path}.{index + 1}"
if os.path.exists(source):
os.replace(source, target)
os.replace(path, f"{path}.1")
+17 -3
View File
@@ -305,9 +305,10 @@ def is_logged_in(account) -> bool:
return bool(login_status(account).get("logged_in"))
def open_product(account, item_id) -> CDP:
def open_product(account, item_id, on_step=None) -> CDP:
"""Open or reuse a product edit tab, navigate to a clean edit URL, and wait ready."""
_notify_collect_step(on_step, "open_product")
host = _cdp_host(account)
item_id = str(item_id)
url = _product_url(account, item_id)
@@ -326,6 +327,7 @@ def open_product(account, item_id) -> CDP:
except Exception:
pass
cdp.send("Page.navigate", {"url": url})
_notify_collect_step(on_step, "wait_ready")
_wait_ready(cdp)
return cdp
@@ -366,17 +368,20 @@ def download_cover(src, out_path) -> str:
return out_path
def collect(account, task) -> dict:
def collect(account, task, on_step=None) -> dict:
"""Collect current title and cover snapshot before any edits."""
item_id = _item_id(task)
cdp = open_product(account, item_id)
cdp = open_product(account, item_id, on_step=on_step)
try:
_notify_collect_step(on_step, "read_title")
old_title = read_title(cdp)
_notify_collect_step(on_step, "read_cover")
old_cover_src = read_cover_src(cdp)
out_path = _get(task, "old_cover_path")
if not out_path:
out_path = os.path.join(_image_root(account), f"{item_id}_old.jpg")
_notify_collect_step(on_step, "download_cover")
old_cover_path = download_cover(old_cover_src, out_path)
return {
"old_title": old_title,
@@ -387,6 +392,15 @@ def collect(account, task) -> dict:
_close_collected_product(cdp)
def _notify_collect_step(callback, step):
if callback is None:
return
try:
callback(step)
except Exception:
pass
def _close_collected_product(cdp):
target_id = getattr(cdp, "target_id", None)
created_by_app = bool(getattr(cdp, "created_by_app", False))
+581 -42
View File
@@ -5,6 +5,7 @@ from __future__ import annotations
import os
import sys
import threading
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
try:
@@ -84,7 +85,7 @@ QTabBar::tab:hover:!selected {
if QT_IMPORT_ERROR is None:
from . import accounts, ai, appconfig, chrome, db, editor, excel, prompts
from . import accounts, ai, appconfig, chrome, db, diagnostics, editor, excel, prompts
from . import config as account_config
@@ -477,6 +478,9 @@ if QT_IMPORT_ERROR is None:
self.batch_filter.setObjectName("batchFilter")
self.shop_filter = QComboBox()
self.shop_filter.setObjectName("shopFilter")
self.item_filter = QLineEdit()
self.item_filter.setObjectName("generateItemFilter")
self.item_filter.setPlaceholderText("商品ID")
self.status_filter = QComboBox()
self.status_filter.setObjectName("statusFilter")
for label, value in self.STATUS_FILTERS:
@@ -488,6 +492,8 @@ if QT_IMPORT_ERROR is None:
filter_layout.addWidget(self.batch_filter, 2)
filter_layout.addWidget(QLabel("店铺"))
filter_layout.addWidget(self.shop_filter, 1)
filter_layout.addWidget(QLabel("商品ID"))
filter_layout.addWidget(self.item_filter, 1)
filter_layout.addWidget(QLabel("状态"))
filter_layout.addWidget(self.status_filter, 1)
filter_layout.addWidget(self.refresh_button)
@@ -502,12 +508,20 @@ if QT_IMPORT_ERROR is None:
self.task_table.horizontalHeader().setSectionResizeMode(QHeaderView.Stretch)
self.task_table.verticalHeader().setVisible(False)
self.run_log_view = QPlainTextEdit()
self.run_log_view.setObjectName("generateRunLogView")
self.run_log_view.setReadOnly(True)
self.run_log_view.setMaximumHeight(128)
self.run_log_view.setPlaceholderText("AI生成运行日志")
right_panel = QWidget()
right_layout = QVBoxLayout(right_panel)
right_layout.setContentsMargins(12, 0, 0, 0)
right_layout.addLayout(filter_layout)
right_layout.addWidget(self.summary_label)
right_layout.addWidget(self.task_table, 1)
right_layout.addWidget(QLabel("AI生成运行日志"))
right_layout.addWidget(self.run_log_view)
self.splitter = QSplitter(Qt.Horizontal)
self.splitter.addWidget(left_panel)
@@ -529,6 +543,7 @@ if QT_IMPORT_ERROR is None:
self.batch_filter.currentIndexChanged.connect(self.refresh_tasks)
self.shop_filter.currentIndexChanged.connect(self.refresh_tasks)
self.item_filter.textChanged.connect(self.refresh_tasks)
self.status_filter.currentIndexChanged.connect(self.refresh_tasks)
self.refresh_button.clicked.connect(self.refresh_tasks)
self.save_title_button.clicked.connect(self.save_title_prompt)
@@ -546,11 +561,33 @@ if QT_IMPORT_ERROR is None:
self.refresh_cover_templates()
self.refresh_tasks()
self._load_latest_generate_run_log()
def _set_status(self, message):
if self.status_callback is not None:
self.status_callback(message)
def _on_generate_log(self, message):
self._append_generate_log(message)
self._set_status(message)
def _append_generate_log(self, message):
self.run_log_view.appendPlainText(str(message))
def _load_latest_generate_run_log(self):
try:
logs = db.list_run_logs(limit=1, run_type="generate", path=self.db_path)
if not logs:
return
events = db.list_run_log_events(logs[0].id, limit=40, path=self.db_path)
except Exception:
return
lines = [
f"{event.created_at} [{event.level}] {event.message}"
for event in events
]
self.run_log_view.setPlainText("\n".join(lines))
def save_title_prompt(self, checked=False):
try:
prompts.save_title_prompt(
@@ -711,13 +748,15 @@ if QT_IMPORT_ERROR is None:
prompt_values,
db_path=self.db_path,
config=self.config,
diagnostic_log_dir=diagnostics.DEFAULT_LOG_DIR,
)
worker.progress.connect(self._on_generate_progress)
worker.row_updated.connect(self._on_generate_row_updated)
worker.log.connect(self._set_status)
worker.log.connect(self._on_generate_log)
worker.failed.connect(self._on_generate_failed)
worker.finished.connect(self._on_generate_finished)
worker.cancelled.connect(self._on_generate_cancelled)
self.run_log_view.clear()
thread = run_worker(worker, thread_name="GenerateWorker", start=False)
thread.finished.connect(lambda: self._forget_generate_thread(thread))
self.generate_worker = worker
@@ -784,6 +823,7 @@ if QT_IMPORT_ERROR is None:
self.refresh_button.setEnabled(not running)
self.batch_filter.setEnabled(not running)
self.shop_filter.setEnabled(not running)
self.item_filter.setEnabled(not running)
self.status_filter.setEnabled(not running)
self.save_title_button.setEnabled(not running)
self.save_cover_template_button.setEnabled(not running)
@@ -809,6 +849,7 @@ if QT_IMPORT_ERROR is None:
def _on_generate_finished(self, payload):
self._set_generate_running(False)
self.refresh_tasks()
self._load_latest_generate_run_log()
self._update_generate_progress(payload)
if payload.get("error"):
self._set_status(f"AI 生成失败:{payload.get('error')}")
@@ -818,6 +859,7 @@ if QT_IMPORT_ERROR is None:
def _on_generate_cancelled(self, payload):
self._set_generate_running(False)
self.refresh_tasks()
self._load_latest_generate_run_log()
self._update_generate_progress(payload)
self._set_status("AI 生成已停止:" + self._generate_progress_text(payload))
@@ -874,6 +916,7 @@ if QT_IMPORT_ERROR is None:
selected_batch = self.batch_filter.currentData()
selected_shop = self.shop_filter.currentData()
selected_status = self.status_filter.currentData() or "all"
item_query = self.item_filter.text().strip()
self._populate_batch_filter(batches, selected_batch)
selected_batch = self.batch_filter.currentData()
batch_tasks = db.list_tasks(batch_id=selected_batch, path=self.db_path)
@@ -882,6 +925,7 @@ if QT_IMPORT_ERROR is None:
filtered_tasks = [
task for task in batch_tasks
if self._matches_shop(task, selected_shop)
and self._matches_item(task, item_query)
and self._matches_status(task, selected_status)
]
except Exception as exc:
@@ -940,6 +984,11 @@ if QT_IMPORT_ERROR is None:
def _matches_shop(self, task, selected_shop):
return selected_shop is None or str(task.alias).strip() == selected_shop
def _matches_item(self, task, item_query):
if not item_query:
return True
return item_query in str(getattr(task, "item_id", ""))
def _matches_status(self, task, selected_status):
if selected_status in (None, "all"):
return True
@@ -990,6 +1039,9 @@ if QT_IMPORT_ERROR is None:
self.batch_filter.setObjectName("applyBatchFilter")
self.shop_filter = QComboBox()
self.shop_filter.setObjectName("applyShopFilter")
self.item_filter = QLineEdit()
self.item_filter.setObjectName("applyItemFilter")
self.item_filter.setPlaceholderText("商品ID")
self.status_filter = QComboBox()
self.status_filter.setObjectName("applyStatusFilter")
for label, value in self.STATUS_FILTERS:
@@ -1001,6 +1053,8 @@ if QT_IMPORT_ERROR is None:
filter_layout.addWidget(self.batch_filter, 2)
filter_layout.addWidget(QLabel("店铺"))
filter_layout.addWidget(self.shop_filter, 1)
filter_layout.addWidget(QLabel("商品ID"))
filter_layout.addWidget(self.item_filter, 1)
filter_layout.addWidget(QLabel("状态"))
filter_layout.addWidget(self.status_filter, 1)
filter_layout.addWidget(self.refresh_button)
@@ -1045,6 +1099,7 @@ if QT_IMPORT_ERROR is None:
self.batch_filter.currentIndexChanged.connect(self.refresh_tasks)
self.shop_filter.currentIndexChanged.connect(self.refresh_tasks)
self.item_filter.textChanged.connect(self.refresh_tasks)
self.status_filter.currentIndexChanged.connect(self.refresh_tasks)
self.refresh_button.clicked.connect(self.refresh_tasks)
self.start_update_button.clicked.connect(self.start_update)
@@ -1066,6 +1121,7 @@ if QT_IMPORT_ERROR is None:
selected_batch = self.batch_filter.currentData()
selected_shop = self.shop_filter.currentData()
selected_status = self.status_filter.currentData() or "generated"
item_query = self.item_filter.text().strip()
self._populate_batch_filter(batches, selected_batch)
selected_batch = self.batch_filter.currentData()
batch_tasks = [
@@ -1077,6 +1133,7 @@ if QT_IMPORT_ERROR is None:
filtered_tasks = [
task for task in batch_tasks
if self._matches_shop(task, selected_shop)
and self._matches_item(task, item_query)
and self._matches_status(task, selected_status)
]
except Exception as exc:
@@ -1219,6 +1276,9 @@ if QT_IMPORT_ERROR is None:
def _shop_filter_label(self):
return self.shop_filter.currentText() or "全部店铺"
def _item_filter_label(self):
return self.item_filter.text().strip() or "全部商品"
def _is_update_task(self, task):
if task.stage in {"generated", "applied"}:
return True
@@ -1227,6 +1287,11 @@ if QT_IMPORT_ERROR is None:
def _matches_shop(self, task, selected_shop):
return selected_shop is None or str(task.alias).strip() == selected_shop
def _matches_item(self, task, item_query):
if not item_query:
return True
return item_query in str(getattr(task, "item_id", ""))
def _matches_status(self, task, selected_status):
if selected_status in (None, "all"):
return True
@@ -1254,6 +1319,7 @@ if QT_IMPORT_ERROR is None:
"即将按当前筛选结果开始更新 Shopee 线上商品。\n\n"
f"批次:{self._batch_filter_label()}\n"
f"店铺:{self._shop_filter_label()}\n"
f"商品ID:{self._item_filter_label()}\n"
f"状态:{self._status_label()}\n"
f"任务数:{len(tasks)}\n\n"
"安全设置:"
@@ -1315,6 +1381,7 @@ if QT_IMPORT_ERROR is None:
self.refresh_button.setEnabled(not running)
self.batch_filter.setEnabled(not running)
self.shop_filter.setEnabled(not running)
self.item_filter.setEnabled(not running)
self.status_filter.setEnabled(not running)
self._update_write_back_button()
@@ -1323,6 +1390,7 @@ if QT_IMPORT_ERROR is None:
self.refresh_button.setEnabled(not running)
self.batch_filter.setEnabled(not running)
self.shop_filter.setEnabled(not running)
self.item_filter.setEnabled(not running)
self.status_filter.setEnabled(not running)
self.write_back_button.setEnabled(False if running else bool(self._active_batch_ids()))
@@ -1596,6 +1664,7 @@ if QT_IMPORT_ERROR is None:
self.collect_thread = None
self.write_back_worker = None
self.write_back_thread = None
self.last_collect_run_id = None
self.import_button = QPushButton("导入 Excel...")
self.refresh_button = QPushButton("刷新")
@@ -1632,6 +1701,12 @@ if QT_IMPORT_ERROR is None:
self.table.horizontalHeader().setSectionResizeMode(QHeaderView.Stretch)
self.table.verticalHeader().setVisible(False)
self.run_log_view = QPlainTextEdit()
self.run_log_view.setObjectName("collectRunLogView")
self.run_log_view.setReadOnly(True)
self.run_log_view.setMaximumHeight(128)
self.run_log_view.setPlaceholderText("采集运行日志")
self.empty_label = QLabel("")
layout = QVBoxLayout(self)
@@ -1640,6 +1715,8 @@ if QT_IMPORT_ERROR is None:
layout.addLayout(summary_layout)
layout.addWidget(self.match_detail_label)
layout.addWidget(self.table, 1)
layout.addWidget(QLabel("采集运行日志"))
layout.addWidget(self.run_log_view)
layout.addWidget(self.empty_label)
self.import_button.clicked.connect(self.import_excel)
@@ -1651,11 +1728,41 @@ if QT_IMPORT_ERROR is None:
self.show_unmatched_button.clicked.connect(self.show_unmatched_tasks)
self.refresh_tasks()
self._load_latest_collect_run_log()
def _set_status(self, message):
if self.status_callback is not None:
self.status_callback(message)
def _on_collect_log(self, message):
self._append_collect_log(message)
self._set_status(message)
def _append_collect_log(self, message):
self.run_log_view.appendPlainText(str(message))
def _load_latest_collect_run_log(self):
try:
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)
except Exception:
return
lines = [
f"{event.created_at} [{event.level}] {event.message}"
for event in events
]
self.run_log_view.setPlainText("\n".join(lines))
def _log_collect_run_event(self, run_id, message, level="info"):
safe_message = diagnostics.redact_log_text(message)
try:
db.add_run_log_event(run_id, safe_message, level=level, path=self.db_path)
except Exception:
return
self._append_collect_log(safe_message)
def _show_error(self, message):
QMessageBox.warning(self, "导入采集", str(message))
self._set_status(str(message))
@@ -1723,13 +1830,19 @@ if QT_IMPORT_ERROR is None:
if not tasks:
self._set_status("没有可采集任务")
return
worker = CollectWorker(tasks, db_path=self.db_path, config=self.config)
worker = CollectWorker(
tasks,
db_path=self.db_path,
config=self.config,
diagnostic_log_dir=diagnostics.DEFAULT_LOG_DIR,
)
worker.progress.connect(self._on_collect_progress)
worker.row_updated.connect(self._on_collect_row_updated)
worker.log.connect(self._set_status)
worker.log.connect(self._on_collect_log)
worker.failed.connect(self._on_collect_failed)
worker.finished.connect(self._on_collect_finished)
worker.cancelled.connect(self._on_collect_cancelled)
self.run_log_view.clear()
thread = run_worker(worker, thread_name="CollectWorker", start=False)
thread.finished.connect(lambda: self._forget_collect_thread(thread))
self.collect_worker = worker
@@ -1830,7 +1943,9 @@ if QT_IMPORT_ERROR is None:
def _on_collect_finished(self, payload):
self._set_collect_running(False)
self.last_collect_run_id = payload.get("run_id") or self.last_collect_run_id
self.refresh_tasks()
self._load_latest_collect_run_log()
if payload.get("blocked"):
self._show_collect_blocked(payload)
return
@@ -1842,6 +1957,11 @@ if QT_IMPORT_ERROR is None:
if payload.get("collected", 0) > 0:
batch_id = self._active_batch_id()
if batch_id and self._start_write_back(batch_id, auto=True):
if self.last_collect_run_id:
self._log_collect_run_event(
self.last_collect_run_id,
"step=excel_write_back result=start detail=采集成功后自动回写旧数据到 Excel",
)
self._set_status(f"{message},正在自动回写 Excel...")
return
if not batch_id:
@@ -1900,6 +2020,12 @@ if QT_IMPORT_ERROR is None:
message += "\n请关闭原 Excel 后重试;SQLite 已保留采集结果,也可另存副本。"
QMessageBox.warning(self, "回写旧数据", message)
self._set_status(message.replace("\n", " "))
if auto and self.last_collect_run_id:
self._log_collect_run_event(
self.last_collect_run_id,
f"step=excel_write_back result=failed detail={error}",
level="error",
)
def _on_write_back_finished(self, payload, auto=False):
self._set_write_back_running(False)
@@ -1907,6 +2033,12 @@ if QT_IMPORT_ERROR is None:
error = payload.get("error") or "未知错误"
retry_hint = ",可点击「回写旧数据到 Excel」手动重试" if auto else ""
self._set_status(f"Excel {'自动' if auto else ''}回写失败:{error}{retry_hint}")
if auto and self.last_collect_run_id:
self._log_collect_run_event(
self.last_collect_run_id,
f"step=excel_write_back result=failed detail={error}",
level="error",
)
return
self.refresh_tasks()
self._set_status(
@@ -1916,6 +2048,11 @@ if QT_IMPORT_ERROR is None:
rows=payload.get("rows", 0),
)
)
if auto and self.last_collect_run_id:
self._log_collect_run_event(
self.last_collect_run_id,
"step=excel_write_back result=success detail=旧数据已回写 Excel",
)
def show_all_tasks(self, checked=False):
self.model.set_filter_mode("all")
@@ -2072,12 +2209,21 @@ if QT_IMPORT_ERROR is None:
class GenerateWorker(BaseWorker):
"""Generate titles and covers for collected tasks."""
def __init__(self, tasks, prompt_values, db_path=None, config=None):
def __init__(
self,
tasks,
prompt_values,
db_path=None,
config=None,
diagnostic_log_dir=None,
):
super().__init__()
self.tasks = list(tasks)
self.prompt_values = dict(prompt_values or {})
self.db_path = db_path
self.config = config
self.diagnostic_log_dir = diagnostic_log_dir
self._run_id = None
def execute(self):
account_rows = accounts.list_accounts(path=self.db_path, config=self.config)
@@ -2086,23 +2232,174 @@ if QT_IMPORT_ERROR is None:
for account in account_rows
if str(account.alias).strip()
}
return ai.generate_batch(
self.tasks,
self.prompt_values,
ai_cfg={
"config": self.config,
"db_path": self.db_path,
"image_dir": appconfig.image_dir(self.config),
"account_by_alias": account_by_alias,
"on_task_update": self._emit_row_update,
},
on_progress=self.progress.emit,
should_stop=self.should_cancel,
eligible = [
task for task in self.tasks
if getattr(task, "stage", None) == "collected"
]
batch_ids = self._batch_ids(eligible)
self._run_id = self._create_run_log(eligible, batch_ids)
self._log_run_event(
f"phase=preflight step=start result=start detail=AI生成开始 total={len(eligible)}"
)
try:
summary = ai.generate_batch(
self.tasks,
self.prompt_values,
ai_cfg={
"config": self.config,
"db_path": self.db_path,
"image_dir": appconfig.image_dir(self.config),
"account_by_alias": account_by_alias,
"on_task_update": self._emit_row_update,
"on_event": self._on_generation_event,
"on_error": self._on_generation_error,
},
on_progress=self.progress.emit,
should_stop=self.should_cancel,
)
except Exception as exc:
error = diagnostics.redact_log_text(str(exc) or exc.__class__.__name__)
summary = {
"ok": False,
"error": error,
"total": len(eligible),
"title_done": 0,
"cover_done": 0,
"failed": len(eligible),
"cancelled": self.should_cancel(),
}
self._log_run_event(
f"phase=worker step=execute result=failed detail={error}",
level="error",
)
self._write_diagnostic_log(
"AI生成运行失败",
level="ERROR",
step="execute",
payload={"error": error},
exc=exc,
)
summary["run_id"] = self._run_id
summary["batch_ids"] = batch_ids
status = "cancelled" if summary.get("cancelled") else "done"
self._finish_run_log(status, summary)
return summary
def _emit_row_update(self, task_id, fields):
self.row_updated.emit(int(task_id), dict(fields or {}))
def _on_generation_event(self, payload):
task = payload.get("task")
phase = payload.get("phase") or "generate"
step = payload.get("step") or "unknown"
result = payload.get("result") or "start"
detail = payload.get("detail")
message = f"phase={phase} step={step} result={result}"
if detail:
message += f" detail={detail}"
self._log_run_event(message, task=task, level=payload.get("level") or "info")
def _on_generation_error(self, payload):
task = payload.get("task")
phase = payload.get("phase") or "generate"
step = payload.get("step") or "unknown"
error = diagnostics.redact_log_text(payload.get("error") or "未知错误")
self._write_diagnostic_log(
"AI生成任务失败",
level="ERROR",
step=step,
task=task,
payload={"phase": phase, "error": error},
exc=payload.get("exception"),
)
def _batch_ids(self, tasks):
batch_ids = []
for task in tasks:
batch_id = getattr(task, "batch_id", None)
if batch_id and batch_id not in batch_ids:
batch_ids.append(batch_id)
return batch_ids
def _create_run_log(self, eligible, batch_ids):
try:
ai_cfg = appconfig.ai_config(self.config)
return db.create_run_log(
"generate",
dry_run=False,
total=len(eligible),
options={
"batch_ids": batch_ids,
"default_text_model": ai_cfg.get("default_text_model"),
"default_image_model": ai_cfg.get("default_image_model"),
"resolution": ai_cfg.get("resolution"),
"title_concurrency": ai_cfg.get("title_concurrency"),
"image_concurrency": ai_cfg.get("image_concurrency"),
},
path=self.db_path,
)
except Exception:
return None
def _finish_run_log(self, status, summary):
if self._run_id is None:
return
try:
done = int(summary.get("cover_done", 0) or 0) + int(summary.get("failed", 0) or 0)
db.finish_run_log(
self._run_id,
status=status,
done=done,
success_count=summary.get("cover_done", 0),
skipped_count=0,
failed_count=summary.get("failed", 0),
summary_json=summary,
path=self.db_path,
)
except Exception:
return
def _log_run_event(self, message, task=None, level="info"):
safe_message = diagnostics.redact_log_text(message)
self.log.emit(str(safe_message))
if self._run_id is None:
return
try:
db.add_run_log_event(
self._run_id,
safe_message,
task_id=getattr(task, "id", None),
alias=getattr(task, "alias", None),
item_id=getattr(task, "item_id", None),
level=level,
path=self.db_path,
)
except Exception:
return
def _write_diagnostic_log(
self,
message,
level="INFO",
step=None,
task=None,
payload=None,
exc=None,
):
try:
diagnostics.write_diagnostic_log(
message,
level=level,
step=step,
task_id=getattr(task, "id", None),
alias=getattr(task, "alias", None),
item_id=getattr(task, "item_id", None),
payload=payload,
exc=exc,
log_dir=self.diagnostic_log_dir,
)
except Exception:
return
class ApplyWorker(BaseWorker):
"""Apply generated title/cover changes, optionally previewing or grouping by account."""
@@ -2522,13 +2819,14 @@ if QT_IMPORT_ERROR is None:
return
def _log_run_event(self, message, task=None, level="info"):
self.log.emit(str(message))
safe_message = diagnostics.redact_log_text(message)
self.log.emit(str(safe_message))
if self._run_id is None:
return
try:
db.add_run_log_event(
self._run_id,
message,
safe_message,
task_id=getattr(task, "id", None),
alias=getattr(task, "alias", None),
item_id=getattr(task, "item_id", None),
@@ -2542,12 +2840,21 @@ if QT_IMPORT_ERROR is None:
class CollectWorker(BaseWorker):
"""Collect old title and cover for imported tasks."""
def __init__(self, tasks, db_path=None, config=None, preflight=True):
def __init__(
self,
tasks,
db_path=None,
config=None,
preflight=True,
diagnostic_log_dir=None,
):
super().__init__()
self.tasks = list(tasks)
self.db_path = db_path
self.config = config
self.preflight = preflight
self.diagnostic_log_dir = diagnostic_log_dir
self._run_id = None
def execute(self):
account_rows = accounts.list_accounts(path=self.db_path, config=self.config)
@@ -2560,27 +2867,41 @@ if QT_IMPORT_ERROR is None:
task for task in self.tasks
if getattr(task, "stage", None) == "imported"
]
batch_ids = self._batch_ids(eligible)
total = len(eligible)
collected = 0
skipped = 0
failed = 0
done = 0
self._run_id = self._create_run_log(eligible, batch_ids)
self._log_run_event(
f"step=preflight result=start detail=采集运行开始 total={total}"
)
if self.preflight:
blocked = self._preflight_block(eligible, account_rows, account_by_alias)
if blocked:
blocked.update(
{
"ok": False,
"blocked": True,
"total": total,
"done": 0,
"collected": 0,
"skipped": 0,
"failed": 0,
}
self._log_preflight_blocked(blocked)
summary = self._summary(
ok=False,
total=total,
done=done,
collected=collected,
skipped=skipped,
failed=failed,
batch_ids=batch_ids,
blocked=True,
extra=blocked,
)
return blocked
self._finish_run_log("blocked", summary)
return summary
self._log_run_event("step=preflight result=success detail=账号检查通过")
else:
self._log_run_event(
"step=preflight result=skipped detail=测试模式跳过采集前检查",
level="warning",
)
for task in eligible:
if self.should_cancel():
@@ -2592,6 +2913,15 @@ if QT_IMPORT_ERROR is None:
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=preflight 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
@@ -2602,10 +2932,41 @@ if QT_IMPORT_ERROR is None:
reason = self._login_skip_reason(status)
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(
task_id=task.id,
item_id=task.item_id,
reason=reason,
),
task=task,
level="warning",
)
self._emit_progress(done, total, collected, skipped, failed)
continue
started = time.monotonic()
current_step = "db_write"
def on_step(step):
nonlocal current_step
current_step = str(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:
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.mark_running(task.id, "collect", path=self.db_path)
self.row_updated.emit(task.id, {"status": "running"})
result = editor.collect(
@@ -2614,6 +2975,15 @@ if QT_IMPORT_ERROR is None:
"item_id": task.item_id,
"old_cover_path": self._old_cover_path(account, task),
},
on_step=on_step,
)
current_step = "db_write"
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_collected(
task.id,
@@ -2622,6 +2992,7 @@ if QT_IMPORT_ERROR is None:
path=self.db_path,
)
collected += 1
elapsed_ms = self._elapsed_ms(started)
self.row_updated.emit(
task.id,
{
@@ -2631,24 +3002,55 @@ if QT_IMPORT_ERROR is None:
"old_cover_path": result.get("old_cover_path", ""),
},
)
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,
)
except Exception as exc:
failed += 1
error = str(exc) or exc.__class__.__name__
db.mark_failed(task.id, "collect", error, path=self.db_path)
self.failed.emit(task.id, error)
self.row_updated.emit(task.id, {"status": "failed", "last_error": error})
safe_error = diagnostics.redact_log_text(error)
elapsed_ms = self._elapsed_ms(started)
db.mark_failed(task.id, "collect", safe_error, path=self.db_path)
self.failed.emit(task.id, safe_error)
self.row_updated.emit(task.id, {"status": "failed", "last_error": safe_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_progress(done, total, collected, skipped, failed)
return {
"ok": failed == 0,
"total": total,
"done": done,
"collected": collected,
"skipped": skipped,
"failed": failed,
}
summary = self._summary(
ok=failed == 0,
total=total,
done=done,
collected=collected,
skipped=skipped,
failed=failed,
batch_ids=batch_ids,
)
self._finish_run_log("cancelled" if self.should_cancel() else "done", summary)
return summary
def _preflight_block(self, eligible, account_rows, account_by_alias):
if not account_rows:
@@ -2727,6 +3129,143 @@ if QT_IMPORT_ERROR is None:
)
)
def _batch_ids(self, tasks):
batch_ids = []
for task in tasks:
batch_id = getattr(task, "batch_id", None)
if batch_id and batch_id not in batch_ids:
batch_ids.append(batch_id)
return batch_ids
def _summary(
self,
ok,
total,
done,
collected,
skipped,
failed,
batch_ids,
blocked=False,
extra=None,
):
summary = {
"ok": ok,
"total": total,
"done": done,
"collected": collected,
"skipped": skipped,
"failed": failed,
"batch_ids": batch_ids,
"run_id": self._run_id,
}
if blocked:
summary["blocked"] = True
if extra:
summary.update(extra)
return summary
def _create_run_log(self, eligible, batch_ids):
try:
return db.create_run_log(
"collect",
dry_run=False,
total=len(eligible),
options={
"batch_ids": batch_ids,
"preflight": self.preflight,
},
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("collected", 0),
skipped_count=summary.get("skipped", 0),
failed_count=summary.get("failed", 0),
summary_json=summary,
path=self.db_path,
)
except Exception:
return
def _log_run_event(self, message, task=None, level="info"):
safe_message = diagnostics.redact_log_text(message)
self.log.emit(str(safe_message))
if self._run_id is None:
return
try:
db.add_run_log_event(
self._run_id,
safe_message,
task_id=getattr(task, "id", None),
alias=getattr(task, "alias", None),
item_id=getattr(task, "item_id", None),
level=level,
path=self.db_path,
)
except Exception:
return
def _log_preflight_blocked(self, blocked):
if blocked.get("no_accounts"):
self._log_run_event(
"step=preflight result=blocked detail=当前没有配置账号",
level="warning",
)
for item in blocked.get("not_running") or []:
self._log_run_event(
"step=preflight result=blocked detail=账号 {alias} Chrome 未启动或调试端口不可访问: {reason}".format(
alias=item.get("alias") or "",
reason=item.get("reason") or "",
),
level="warning",
)
for item in blocked.get("logged_out") or []:
self._log_run_event(
"step=preflight result=blocked detail=账号 {alias} 未登录 Shopee: {reason}".format(
alias=item.get("alias") or "",
reason=item.get("reason") or "",
),
level="warning",
)
def _write_diagnostic_log(
self,
message,
level="INFO",
step=None,
task=None,
elapsed_ms=None,
payload=None,
exc=None,
):
try:
diagnostics.write_diagnostic_log(
message,
level=level,
step=step,
task_id=getattr(task, "id", None),
alias=getattr(task, "alias", None),
item_id=getattr(task, "item_id", None),
elapsed_ms=elapsed_ms,
payload=payload,
exc=exc,
log_dir=self.diagnostic_log_dir,
)
except Exception:
return
def _elapsed_ms(self, started):
return int((time.monotonic() - started) * 1000)
class WriteBackWorker(BaseWorker):
"""Write Excel fields back in a background thread."""