Commit 323709a5 authored by Data Governance Dev's avatar Data Governance Dev

feat: 增量结果发布 + 每 tab 三状态徽章 + OverviewStats 必填字段兜底

之前结果只在所有 step 完成后才返回,导致前端看不到中间进度。本轮改:

后端增量发布:
- orchestrator.run_governance_workflow 新增 on_section_update 回调
  回调在每步完成且有 section_key 时触发
- job_manager 新增 _sync_section_update 同步回调到 job.result['sections']
  (job.result 在任务结束时再被整体 result 替换)
- api/routes.py /jobs/{id}/result 取消'必须 completed'限制,
  running/failed/cancelled 状态也允许读取 sections

前端轮询 + tab 状态:
- pollStatus 每 1.5s 同时拉 status 与 result;
  loadResult 改为 merge sections(不整体替换 jobResult)
- TAB_STEP_MAP 映射 tab -> (step, sectionKey)
- tabStatus computed 按状态返回 'idle'|'running'|'done'
- TAB_BADGE 配置:info(非分析对象)/warning(分析中)/success(分析完成)
- 6 个 el-tab-pane 全部改用 #label slot,加 el-tag 徽章

bug 修复:OverviewStats 必填字段在运行中没数据导致 500:
- core/models.py: database/host/db_type 改为带默认值 ''
- api/routes.py: OverviewStats 构造用 try/except 兜底,失败降级到全零值

AST 全部通过;本地 == HTTP 服务返回一致。
parent 9ac8c11b
...@@ -186,7 +186,7 @@ async def cancel_job(job_id: str): ...@@ -186,7 +186,7 @@ async def cancel_job(job_id: str):
return {"ok": True, "job_id": job_id} return {"ok": True, "job_id": job_id}
# ── 任务结果 ── # ── 任务结果(支持增量:运行中也可拿到已完成的 sections) ──
@router.get("/jobs/{job_id}/result", response_model=JobResultResponse) @router.get("/jobs/{job_id}/result", response_model=JobResultResponse)
async def get_job_result(job_id: str): async def get_job_result(job_id: str):
jm = JobManager.get_instance() jm = JobManager.get_instance()
...@@ -194,20 +194,26 @@ async def get_job_result(job_id: str): ...@@ -194,20 +194,26 @@ async def get_job_result(job_id: str):
if not job: if not job:
logger.warning(f"GET /api/jobs/{job_id}/result 任务不存在") logger.warning(f"GET /api/jobs/{job_id}/result 任务不存在")
raise HTTPException(status_code=404, detail="任务不存在") raise HTTPException(status_code=404, detail="任务不存在")
if job.status != "completed": # 不再要求 status == "completed":运行中也允许前端拉取已完成的 sections
logger.warning(f"GET /api/jobs/{job_id}/result 任务未完成(status={job.status})") if job.status not in ("running", "completed", "failed", "cancelled"):
raise HTTPException(status_code=400, detail=f"任务未完成(当前: {job.status})") logger.warning(f"GET /api/jobs/{job_id}/result 任务状态异常(status={job.status})")
raise HTTPException(status_code=400, detail=f"任务状态异常(当前: {job.status})")
sections = job.result.get("sections", {})
overview_dict = sections.get("overview", {}) sections = job.result.get("sections", {}) if isinstance(job.result, dict) else {}
overview = OverviewStats(**{k: v for k, v in overview_dict.items() if k in OverviewStats.model_fields}) overview_dict = sections.get("overview", {}) or {}
logger.info(f"GET /api/jobs/{job_id}/result 成功, sections={list(sections.keys())}") # OverviewStats 必有字段(database/host/db_type)在运行中还没生成,
# 用 try/except 兜底:构造失败时降级到零值
try:
overview = OverviewStats(**{k: v for k, v in overview_dict.items() if k in OverviewStats.model_fields})
except Exception:
overview = OverviewStats() # 全部走默认值
logger.info(f"GET /api/jobs/{job_id}/result 成功(status={job.status}), sections={list(sections.keys())}")
return JobResultResponse( return JobResultResponse(
job_id=job_id, job_id=job_id,
status=job.status, status=job.status,
overview=overview, overview=overview,
sections=sections, sections=sections,
reports=job.result.get("reports", {}), reports=job.result.get("reports", {}) if isinstance(job.result, dict) else {},
) )
......
...@@ -229,6 +229,7 @@ class JobManager: ...@@ -229,6 +229,7 @@ class JobManager:
on_log=lambda entry: self._sync_log(job, entry), on_log=lambda entry: self._sync_log(job, entry),
on_step_start=lambda num, title: self._sync_step_start(job, num, title), on_step_start=lambda num, title: self._sync_step_start(job, num, title),
on_step_done=lambda num, title, ok: self._sync_step_done(job, num, title, ok), on_step_done=lambda num, title, ok: self._sync_step_done(job, num, title, ok),
on_section_update=lambda section_key, data: self._sync_section_update(job, section_key, data),
cancel_event=job._cancel_event, cancel_event=job._cancel_event,
) )
job.result = result job.result = result
...@@ -278,6 +279,12 @@ class JobManager: ...@@ -278,6 +279,12 @@ class JobManager:
else: else:
job.steps_failed.append(num) job.steps_failed.append(num)
def _sync_section_update(self, job: Job, section_key: str, data: dict) -> None:
"""Step 完成时增量落盘:让 API 可以拿到部分 sections"""
if not isinstance(job.result, dict):
job.result = {}
job.result.setdefault("sections", {})[section_key] = data
# ── 查询 ── # ── 查询 ──
def get(self, job_id: str) -> Optional[Job]: def get(self, job_id: str) -> Optional[Job]:
return self.jobs.get(job_id) return self.jobs.get(job_id)
......
...@@ -95,9 +95,9 @@ class StandardsListResponse(BaseModel): ...@@ -95,9 +95,9 @@ class StandardsListResponse(BaseModel):
# ── 治理结果(前端多级表格消费) ── # ── 治理结果(前端多级表格消费) ──
class OverviewStats(BaseModel): class OverviewStats(BaseModel):
database: str database: str = ""
host: str host: str = ""
db_type: str db_type: str = ""
total_tables: int = 0 total_tables: int = 0
total_fields: int = 0 total_fields: int = 0
high_empty_fields: int = 0 high_empty_fields: int = 0
......
...@@ -88,6 +88,7 @@ async def run_governance_workflow( ...@@ -88,6 +88,7 @@ async def run_governance_workflow(
on_log: Callable[[dict], None] | None = None, on_log: Callable[[dict], None] | None = None,
on_step_start: Callable[[int, str], None] | None = None, on_step_start: Callable[[int, str], None] | None = None,
on_step_done: Callable[[int, str, bool], None] | None = None, on_step_done: Callable[[int, str, bool], None] | None = None,
on_section_update: Callable[[str, dict], None] | None = None,
cancel_event: Optional[asyncio.Event] = None, cancel_event: Optional[asyncio.Event] = None,
) -> dict: ) -> dict:
"""主入口:跑一轮治理,返回结构化结果。 """主入口:跑一轮治理,返回结构化结果。
...@@ -267,6 +268,9 @@ async def run_governance_workflow( ...@@ -267,6 +268,9 @@ async def run_governance_workflow(
section_key = output.get("section_key") section_key = output.get("section_key")
if section_key: if section_key:
sections[section_key] = output["data"] sections[section_key] = output["data"]
if on_section_update:
try: on_section_update(section_key, output["data"])
except Exception: pass
if output.get("extras"): if output.get("extras"):
step_outputs[step_name(num)] = output["extras"] step_outputs[step_name(num)] = output["extras"]
......
This diff is collapsed.
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment