877723ab4a
体验稿(8823)逐条过审拍板后动工,不造新机器只做可见性:
件① 三件套(spec §7):
- 1a 卡片已产因子清单:列表项内联 factors(registry hypothesis
血统反查:名/表达式/origin/落成日);单因子名 factorId 前端退役
- 1b 分解批次历史:GET /pipeline/hypotheses/{id}/decompose-jobs
台账新→旧(QA tasks/list 端点形状照抄),卡片折叠区按需拉取
- 1c 轮次进度:run_decompose 每轮 on_progress 回调→内存 job
progress(currentRound/totalRounds/passed/regen;真轮数非 QA 摆设),
台账仍只记起跑+终态;worker 协议改收 job_id
件② 工厂来源列:factors 端点带 origin/hypothesis 血统,
⚡decomposer/✍manual 徽标+回链假设池
件③ 琥珀待办:todos 聚合加 queued/data_check 卡(已确认未分解),
分解转 building 即消行——人工卡点②显性化
件④ 结果持久落卡:前端一次性弹层退役,终态摘要+清单+批次全在卡
测试:后端 TestD7Visibility 六件(历史端点/内联清单/进度中飞可见/
终态不带过程态/琥珀消行/工厂血统)+前端三件套三测;8 目录 2780 绿
120 lines
4.9 KiB
Python
120 lines
4.9 KiB
Python
"""分解 job 台账(10-09)——QA frontend-v2/backend/app.py 任务字典范式移植.
|
|
|
|
10-09 定谳:分解服务端早已成功(registry 落册)而浏览器先超时断开,uvicorn
|
|
对断连请求不写访问日志——响应孤儿化=当日全部「静默卡死」同源.分钟级 LLM
|
|
任务不架在同步 HTTP 上(调超时永远治不了,已三轮验证):POST 起后台任务
|
|
立即返回,前端轮询卡片 decomposeJob 到终态(QA mining start/status 同款).
|
|
|
|
两本账:内存 _JOBS=在飞真相(QA tasks 同款);pipeline_store.decompose_jobs
|
|
表=跨重启台账(QA 没有,我们补——run 台账先于页面,服务重启丢内存后
|
|
台账 running 无主→读侧判 failed(中断),结果永不孤儿化).
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
import uuid
|
|
from datetime import datetime
|
|
from typing import Any, Awaitable, Callable
|
|
|
|
from .llm import LLMError
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
_JOBS: dict[str, dict[str, Any]] = {} # job_id -> job 视图(camelCase)
|
|
_LATEST: dict[str, str] = {} # hyp_id -> 最新 job_id
|
|
_TASKS: set[asyncio.Task] = set() # 持引用防 GC(create_task 裸弃=可能被收)
|
|
|
|
_INTERRUPTED = "服务重启,任务中断——请重试"
|
|
|
|
|
|
class JobFailure(Exception):
|
|
"""worker 内业务性失败(面向用户的中文信息,如卡片状态已变),区别于上游错误."""
|
|
|
|
|
|
def _now() -> str:
|
|
return datetime.now().isoformat(timespec="seconds")
|
|
|
|
|
|
def is_running(hyp_id: str) -> bool:
|
|
job_id = _LATEST.get(hyp_id)
|
|
return job_id is not None and _JOBS.get(job_id, {}).get("status") == "running"
|
|
|
|
|
|
def start(hyp_id: str, worker: Callable[[str], Awaitable[dict]],
|
|
db_path: str) -> dict[str, Any]:
|
|
"""登记 running 行(内存+台账)并起后台任务;同卡在飞由调用方前置拦(409).
|
|
|
|
worker 收 job_id(D7 件①1c:体内每轮回调 update_progress 写在飞进度).
|
|
"""
|
|
job: dict[str, Any] = {
|
|
"jobId": f"job-{datetime.now().strftime('%Y%m%d%H%M%S')}"
|
|
f"-{uuid.uuid4().hex[:8]}",
|
|
"hypId": hyp_id, "status": "running", "rounds": None,
|
|
"registered": [], "failed": [], "error": None,
|
|
"startedAt": _now(), "finishedAt": None}
|
|
_JOBS[job["jobId"]] = job
|
|
_LATEST[hyp_id] = job["jobId"]
|
|
from sanguo_portfolio import pipeline_store
|
|
pipeline_store.upsert_decompose_job(db_path, job)
|
|
task = asyncio.create_task(_run(job["jobId"], worker, db_path))
|
|
_TASKS.add(task)
|
|
task.add_done_callback(_TASKS.discard)
|
|
return dict(job)
|
|
|
|
|
|
def update_progress(job_id: str, *, round_no: int, total_rounds: int,
|
|
passed: int, regen: int, message: str) -> None:
|
|
"""在飞轮次进度(D7 件①1c,QA task.progress 同位):只写内存 job 字典,
|
|
台账仍只记起跑+终态(过程态不落库,与 QA 同纪律)."""
|
|
job = _JOBS.get(job_id)
|
|
if job is None or job.get("status") != "running":
|
|
return
|
|
job["progress"] = {"currentRound": round_no, "totalRounds": total_rounds,
|
|
"passed": passed, "regen": regen, "message": message}
|
|
|
|
|
|
async def _run(job_id: str, worker: Callable[[str], Awaitable[dict]],
|
|
db_path: str) -> None:
|
|
job = _JOBS[job_id]
|
|
try:
|
|
result = await worker(job_id)
|
|
except LLMError as e:
|
|
# P3-8 同款不泄露:真因只进服务端日志,job.error 固定文案
|
|
logger.warning("llm decompose upstream error: %s", e)
|
|
job.update(status="failed", error="LLM 上游返回异常")
|
|
except JobFailure as e:
|
|
job.update(status="failed", error=str(e))
|
|
except ValueError as e:
|
|
job.update(status="failed", error=f"LLM 草稿不合契约: {e}")
|
|
except Exception: # noqa: BLE001 后台任务裸异常无人接=永久 running
|
|
logger.exception("decompose job %s crashed", job_id)
|
|
job.update(status="failed", error="分解任务异常(服务端日志留痕)")
|
|
else:
|
|
job.update(status="completed", rounds=result["rounds"],
|
|
registered=result["registered"], failed=result["failed"])
|
|
job.pop("progress", None) # 终态不再带过程态
|
|
job["finishedAt"] = _now()
|
|
from sanguo_portfolio import pipeline_store
|
|
pipeline_store.upsert_decompose_job(db_path, job)
|
|
|
|
|
|
def view(hyp_id: str, db_path: str) -> dict[str, Any] | None:
|
|
"""最新 job 视图:内存(在飞/本轮终态)→台账(跨重启/刷新).
|
|
|
|
台账 running 而内存无主=服务重启打断,读侧判 failed 并回写台账
|
|
(一处真相;不回写则台账永远躺着僵尸 running).
|
|
"""
|
|
from sanguo_portfolio import pipeline_store
|
|
job_id = _LATEST.get(hyp_id)
|
|
if job_id is not None and job_id in _JOBS:
|
|
return dict(_JOBS[job_id])
|
|
row = pipeline_store.latest_decompose_job(db_path, hyp_id)
|
|
if row is None:
|
|
return None
|
|
if row["status"] == "running":
|
|
row = {**row, "status": "failed", "error": _INTERRUPTED,
|
|
"finishedAt": row["startedAt"]}
|
|
pipeline_store.upsert_decompose_job(db_path, row)
|
|
return row
|