226 lines
8.7 KiB
Python
226 lines
8.7 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
import logging
|
|
import os
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
from pathlib import Path
|
|
|
|
from .job_manager import JobManager
|
|
from .log_config import get_logger
|
|
from .schemas import JobPaths
|
|
|
|
log = get_logger("service.runner")
|
|
|
|
|
|
STAGE_RULES: list[tuple[str, str, int]] = [
|
|
("--- 1. 读取数据 ---", "reading_data", 10),
|
|
("--- 2. 智能画幅计算 ---", "sizing_canvas", 20),
|
|
("--- 3. 生成掩膜", "building_mask", 35),
|
|
("--- 4. 计算权重 ---", "computing_weights", 45),
|
|
("--- 5. 启动生成", "placing_words", 65),
|
|
("--- 6. 高清渲染 ---", "rendering", 85),
|
|
("已保存:", "writing_outputs", 92),
|
|
("✅ 完成", "writing_outputs", 99),
|
|
]
|
|
|
|
|
|
class JobRunner:
|
|
def __init__(self, project_root: Path, manager: JobManager) -> None:
|
|
self.project_root = project_root
|
|
self.manager = manager
|
|
self.script_path = self.project_root / "wordcloud_generate_hybrid.py"
|
|
|
|
def _parse_stage(self, line: str, current_stage: str, current_progress: int) -> tuple[str, int]:
|
|
for token, stage, progress in STAGE_RULES:
|
|
if token in line:
|
|
return stage, progress
|
|
|
|
if "尝试 #" in line or "尺度" in line or "二分重试" in line:
|
|
progress = max(current_progress, 70)
|
|
return "placing_words", min(progress + 1, 84)
|
|
|
|
return current_stage, current_progress
|
|
|
|
def run(self, job_id: str, paths: JobPaths, config: dict) -> None:
|
|
t_start = time.perf_counter()
|
|
log.info("=" * 50)
|
|
log.info("[Runner] 任务启动 job_id=%s", job_id)
|
|
log.info(" config_path = %s", paths.config_path)
|
|
log.info(" output_dir = %s", paths.output_dir)
|
|
log.info(" 子进程 python = %s", sys.executable)
|
|
self.manager.set_status(job_id, status="running", stage="starting", progress_percent=1, message="任务启动")
|
|
|
|
with paths.config_path.open("w", encoding="utf-8") as f:
|
|
json.dump(config, f, ensure_ascii=False, indent=2)
|
|
log.info(" 配置文件已写入")
|
|
|
|
cmd = [
|
|
sys.executable,
|
|
str(self.script_path),
|
|
"--config",
|
|
str(paths.config_path),
|
|
]
|
|
|
|
env = os.environ.copy()
|
|
# Force unbuffered stdout so every print() flushes immediately and the
|
|
# frontend SSE log view shows each step in real time. Without this,
|
|
# Python block-buffers stdout when it is a pipe, so lines pile up and
|
|
# only arrive in bursts after the buffer fills or the process exits.
|
|
env["PYTHONUNBUFFERED"] = "1"
|
|
process = subprocess.Popen(
|
|
cmd,
|
|
cwd=str(self.project_root),
|
|
stdout=subprocess.PIPE,
|
|
stderr=subprocess.STDOUT,
|
|
text=True,
|
|
bufsize=1,
|
|
env=env,
|
|
)
|
|
|
|
stage = "starting"
|
|
progress = 1
|
|
output_lines: list[str] = [] # collect all output lines for error reporting
|
|
preview_sent = False
|
|
|
|
assert process.stdout is not None
|
|
for raw in process.stdout:
|
|
line = raw.rstrip("\n")
|
|
output_lines.append(line)
|
|
log.info("[Pipeline] %s", line)
|
|
stage, progress = self._parse_stage(line, stage, progress)
|
|
elapsed = time.perf_counter() - t_start
|
|
self.manager.add_event(
|
|
job_id,
|
|
kind="log",
|
|
stage=stage,
|
|
progress_percent=progress,
|
|
message=line,
|
|
elapsed_seconds=round(elapsed, 3),
|
|
)
|
|
# The PNG is complete before SVG/DB export starts. Publish it as a
|
|
# preview so the frontend does not wait for the slower artifacts.
|
|
if not preview_sent and line.startswith("已保存:"):
|
|
candidate = Path(line.split(":", 1)[1].strip())
|
|
if candidate.suffix.lower() == ".png" and candidate.exists():
|
|
self.manager.set_artifacts(job_id, {"png": str(candidate)})
|
|
self.manager.add_event(
|
|
job_id,
|
|
kind="status",
|
|
stage="preview_ready",
|
|
progress_percent=max(progress, 94),
|
|
message="预览已生成,后台继续导出其余文件",
|
|
elapsed_seconds=round(elapsed, 3),
|
|
)
|
|
preview_sent = True
|
|
|
|
ret = process.wait()
|
|
elapsed = time.perf_counter() - t_start
|
|
log.info("[Runner] 子进程退出 code=%d 耗时=%.2fs", ret, elapsed)
|
|
|
|
# ── 错误时:截取最后 30 行输出作为详细错误信息 ─────────────
|
|
error_detail = None
|
|
if ret != 0:
|
|
# 找到 "生成失败" 或 "错误" 或 traceback 之后的内容
|
|
error_lines = []
|
|
captured = False
|
|
for line in reversed(output_lines):
|
|
if not captured:
|
|
error_lines.append(line)
|
|
if any(kw in line for kw in ("生成失败", "error", "Error", "Traceback", "traceback", "未满足", "放置")):
|
|
captured = True
|
|
elif len(error_lines) < 30:
|
|
error_lines.append(line)
|
|
else:
|
|
break
|
|
error_lines.reverse()
|
|
error_detail = "\n".join(error_lines) if error_lines else f"script exited with code {ret}"
|
|
log.info("[Runner] 错误详情:\n%s", error_detail)
|
|
|
|
png = next(paths.output_dir.glob("*.png"), None)
|
|
# NOTE: do NOT use "*[!_stroke].svg" — in glob, [!...] is a character class,
|
|
# so filenames ending with "e.svg" (e.g. AutoResize.svg) are incorrectly skipped.
|
|
svg = next(
|
|
(p for p in sorted(paths.output_dir.glob("*.svg")) if not p.name.endswith("_stroke.svg")),
|
|
None,
|
|
)
|
|
svg_stroke = next(paths.output_dir.glob("*_stroke.svg"), None)
|
|
db = next(paths.output_dir.glob("*.db"), None)
|
|
metrics = next(paths.output_dir.glob("*metrics*.json"), None)
|
|
|
|
log.info("[Runner] 产物扫描:")
|
|
log.info(" png = %s", png)
|
|
log.info(" svg = %s", svg)
|
|
log.info(" svg_stroke = %s", svg_stroke)
|
|
log.info(" db = %s", db)
|
|
log.info(" metrics = %s", metrics)
|
|
|
|
artifacts = {
|
|
"png": str(png) if png else "",
|
|
"svg": str(svg) if svg else "",
|
|
"svg_stroke": str(svg_stroke) if svg_stroke else "",
|
|
"db": str(db) if db else "",
|
|
"metrics": str(metrics) if metrics else "",
|
|
}
|
|
self.manager.set_artifacts(job_id, artifacts)
|
|
|
|
if ret == 0 and not png:
|
|
log.error("[Runner] 任务失败:退出码=0 但未找到输出图片")
|
|
self.manager.add_event(
|
|
job_id,
|
|
kind="status",
|
|
stage="failed",
|
|
progress_percent=100,
|
|
message="任务失败:未找到输出图片",
|
|
)
|
|
self.manager.set_status(
|
|
job_id,
|
|
status="failed",
|
|
stage="failed",
|
|
progress_percent=100,
|
|
message="任务失败",
|
|
error="missing png artifact",
|
|
)
|
|
return
|
|
|
|
if ret == 0:
|
|
done_message = f"任务完成,用时 {elapsed:.2f} 秒"
|
|
log.info("[Runner] ✅ 任务完成 job_id=%s 总耗时=%.2fs", job_id, elapsed)
|
|
self.manager.add_event(
|
|
job_id,
|
|
kind="status",
|
|
stage="completed",
|
|
progress_percent=100,
|
|
message=done_message,
|
|
elapsed_seconds=round(elapsed, 3),
|
|
)
|
|
self.manager.set_status(
|
|
job_id,
|
|
status="success",
|
|
stage="completed",
|
|
progress_percent=100,
|
|
message=done_message,
|
|
elapsed_seconds=round(elapsed, 3),
|
|
)
|
|
else:
|
|
log.error("[Runner] ❌ 任务失败 job_id=%s exit_code=%d 耗时=%.2fs", job_id, ret, elapsed)
|
|
self.manager.add_event(
|
|
job_id,
|
|
kind="status",
|
|
stage="failed",
|
|
progress_percent=100,
|
|
message=f"任务失败,退出码: {ret},用时 {elapsed:.2f} 秒",
|
|
elapsed_seconds=round(elapsed, 3),
|
|
)
|
|
self.manager.set_status(
|
|
job_id,
|
|
status="failed",
|
|
stage="failed",
|
|
progress_percent=100,
|
|
message="任务失败",
|
|
error=error_detail or f"script exited with code {ret}",
|
|
elapsed_seconds=round(elapsed, 3),
|
|
)
|