From d03a74375635d62e8130d44aa3d69a11271e2652 Mon Sep 17 00:00:00 2001 From: caoqianming Date: Wed, 19 Aug 2026 11:07:58 +0800 Subject: [PATCH] =?UTF-8?q?fix(software):=20=E5=B0=86=E5=A4=B1=E8=B4=A5?= =?UTF-8?q?=E5=92=8C=E5=8F=96=E6=B6=88=E7=BB=93=E6=9E=9C=E9=80=9A=E7=9F=A5?= =?UTF-8?q?=E5=8E=9F=E5=AF=B9=E8=AF=9D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGELOG.md | 2 + core/software_jobs.py | 3 +- tests/test_software_followups.py | 63 +++++++++++++++++++++++++++++++- tests/test_software_nodes.py | 29 +++++++++++++++ web/software_followups.py | 37 +++++++++++++++---- 5 files changed, 124 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 319c46a..8a6ae61 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,8 @@ ## Unreleased +- 专业软件任务失败或取消后会在原对话留下持久通知,页面离线、刷新或对话暂时忙碌时也不会错过结果。 + - 查询专业软件是否可用时不再误提交真实任务;ANSYS 静力分析会要求使用几何中已有的 Named Selection,并在输入缺少命名选择时返回明确错误。 - Origin 专业软件任务现在会在原 Windows Node 上保留可继续加工的工程;后续调整直接打开上一版工程,过程只提供轻量预览,工程和其他大文件仅在明确要求下载或交付时才传回并登记为正式产物。 diff --git a/core/software_jobs.py b/core/software_jobs.py index 7a3bf32..e58524e 100644 --- a/core/software_jobs.py +++ b/core/software_jobs.py @@ -210,6 +210,7 @@ def request_job_cancel(user_id: UUID, job_id: UUID) -> tuple[dict, dict | None]: job.stage = "terminal" job.error = {"code": "USER_CANCELLED", "detail": "Cancelled before dispatch."} job.terminal_at = now + job.followup_status = "pending" return _job_dict(job), None job.status = "cancelling" job.stage = "cancel_requested" @@ -1158,8 +1159,8 @@ def record_job_terminal(node_id: UUID, payload: dict) -> None: job.local_manifest = local_manifest job.preview_manifest = preview_manifest job.terminal_at = now + job.followup_status = "pending" if terminal_status == "succeeded": - job.followup_status = "pending" if getattr(job, "workspace_id", None) is not None: workspace = session.execute( select(SoftwareWorkspace).where( diff --git a/tests/test_software_followups.py b/tests/test_software_followups.py index 92a7c89..4ba85f2 100644 --- a/tests/test_software_followups.py +++ b/tests/test_software_followups.py @@ -31,7 +31,7 @@ class _Session: self.added.append(row) -def _job(action: str): +def _job(action: str, *, status: str = "succeeded"): job_id = uuid4() task_id = uuid4() return SimpleNamespace( @@ -39,9 +39,10 @@ def _job(action: str): task_id=task_id, user_id=uuid4(), capability="origin.plot@v2", - status="succeeded", + status=status, followup_status="pending", completion_action=action, + error={}, artifact_manifest=[ { "artifact_id": str(uuid4()), @@ -119,6 +120,64 @@ class SoftwareFollowupTests(unittest.TestCase): self.assertEqual(job.followup_status, "pending") self.assertFalse(session.added) + def test_failed_job_persists_fixed_assistant_message_without_artifacts(self): + job = _job("analyze", status="failed") + job.error = {"code": "COLUMN_MISSING", "detail": " Y2\n does not exist "} + task = SimpleNamespace(task_id=job.task_id, run_status="idle") + session = _Session([job.task_id, task, job]) + + @contextmanager + def scope(): + yield session + + with ( + patch("web.software_followups.session_scope", scope), + patch("web.software_followups.allocate_message_idx", return_value=9), + ): + claim = claim_followup(job.job_id) + + self.assertEqual(claim.action, "report") + self.assertEqual(job.followup_status, "completed") + message = session.added[0] + self.assertEqual(message.payload["role"], "assistant") + self.assertIn("专业软件任务执行失败", message.payload["content"]) + self.assertIn("Y2 does not exist", message.payload["content"]) + self.assertEqual(message.artifact_refs, []) + + def test_cancelled_job_persists_fixed_assistant_message(self): + job = _job("analyze", status="cancelled") + task = SimpleNamespace(task_id=job.task_id, run_status="idle") + session = _Session([job.task_id, task, job]) + + @contextmanager + def scope(): + yield session + + with ( + patch("web.software_followups.session_scope", scope), + patch("web.software_followups.allocate_message_idx", return_value=10), + ): + claim = claim_followup(job.job_id) + + self.assertEqual(claim.action, "report") + self.assertIn("专业软件任务已取消", session.added[0].payload["content"]) + + def test_busy_task_leaves_failed_notification_pending(self): + job = _job("report", status="failed") + task = SimpleNamespace(task_id=job.task_id, run_status="running") + session = _Session([job.task_id, task]) + + @contextmanager + def scope(): + yield session + + with patch("web.software_followups.session_scope", scope): + claim = claim_followup(job.job_id) + + self.assertIsNone(claim) + self.assertEqual(job.followup_status, "pending") + self.assertFalse(session.added) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_software_nodes.py b/tests/test_software_nodes.py index 22dc277..886e263 100644 --- a/tests/test_software_nodes.py +++ b/tests/test_software_nodes.py @@ -330,6 +330,7 @@ class SoftwareJobProtocolTests(unittest.TestCase): self.assertEqual(result["status"], "cancelled") self.assertIsNone(node_message) self.assertEqual(job.error["code"], "USER_CANCELLED") + self.assertEqual(job.followup_status, "pending") @patch("core.software_jobs.session_scope") def test_running_job_persists_cancel_before_sending(self, session_scope) -> None: @@ -858,6 +859,34 @@ class SoftwareJobProtocolTests(unittest.TestCase): self.assertEqual(job.status, "succeeded") self.assertEqual(job.followup_status, "pending") + @patch("core.software_jobs.session_scope") + def test_failed_terminal_queues_completion_followup(self, session_scope) -> None: + session = session_scope.return_value.__enter__.return_value + node_id = uuid4() + lease_id = uuid4() + digest = "e" * 64 + job = type("Job", (), {})() + job.job_id = uuid4() + job.node_id = node_id + job.lease_id = lease_id + job.request_digest = digest + job.status = "running" + job.capability = "origin.plot@v2" + job.progress = 40 + session.execute.return_value.scalar_one_or_none.return_value = job + + record_job_terminal(node_id, { + "job_id": str(job.job_id), + "lease_id": str(lease_id), + "request_digest": digest, + "status": "failed", + "error": {"code": "COLUMN_MISSING", "detail": "Y2 does not exist"}, + "artifact_manifest": [], + }) + + self.assertEqual(job.status, "failed") + self.assertEqual(job.followup_status, "pending") + @patch("core.software_jobs.session_scope") def test_disconnect_does_not_requeue_active_jobs(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value diff --git a/web/software_followups.py b/web/software_followups.py index 034b1a4..bc0131b 100644 --- a/web/software_followups.py +++ b/web/software_followups.py @@ -1,4 +1,4 @@ -"""Software Job 完成后的固定报告与 Agent 自动续跑。""" +"""Software Job 终态后的固定报告与 Agent 自动续跑。""" from __future__ import annotations import asyncio @@ -15,6 +15,8 @@ from core.storage.models import Message, SoftwareJob, Task from .common import INSTANCE from .run_lifecycle import RunScheduleError, schedule_claimed_run +TERMINAL_STATUSES = {"succeeded", "failed", "cancelled"} + @dataclass(frozen=True) class FollowupClaim: @@ -56,6 +58,23 @@ def _preview_refs(manifest: list) -> list[dict]: def _report_text(job: SoftwareJob) -> str: contract = get_contract(job.capability) + if job.status == "failed": + error = job.error if isinstance(job.error, dict) else {} + code = str(error.get("code") or "SOFTWARE_JOB_FAILED")[:100] + detail = " ".join(str(error.get("detail") or "未提供具体错误信息").split())[:500] + return ( + f"专业软件任务执行失败:{contract.display_name}\n\n" + f"- Job ID:`{job.job_id}`\n" + f"- 错误代码:`{code}`\n" + f"- 原因:{detail}\n" + "- 结果:任务未生成正式产物;调整输入或运行环境后可重新提交" + ) + if job.status == "cancelled": + return ( + f"专业软件任务已取消:{contract.display_name}\n\n" + f"- Job ID:`{job.job_id}`\n" + "- 结果:任务已停止,未完成的中间输出不会作为正式产物发布" + ) if contract.workspace is not None and getattr(job, "workspace_id", None) is not None: exported = [ str(item.get("filename")) @@ -115,7 +134,7 @@ def pending_followup_ids(limit: int = 20) -> list[UUID]: return list(session.execute( select(SoftwareJob.job_id) .where( - SoftwareJob.status == "succeeded", + SoftwareJob.status.in_(TERMINAL_STATUSES), SoftwareJob.followup_status == "pending", ) .order_by(SoftwareJob.terminal_at, SoftwareJob.job_id) @@ -162,21 +181,25 @@ def claim_followup(job_id: UUID) -> FollowupClaim | None: ).scalar_one_or_none() if ( job is None - or job.status != "succeeded" + or job.status not in TERMINAL_STATUSES or job.followup_status != "pending" ): return None next_idx = allocate_message_idx(session, task.task_id, locked_task=task) - if job.completion_action == "report": + if job.status != "succeeded" or job.completion_action == "report": session.add(Message( task_id=task.task_id, idx=next_idx, payload={"role": "assistant", "content": _report_text(job)}, artifact_refs=( - _artifact_refs(job.artifact_manifest) - if job.artifact_manifest - else _preview_refs(job.preview_manifest) + [] + if job.status != "succeeded" + else ( + _artifact_refs(job.artifact_manifest) + if job.artifact_manifest + else _preview_refs(job.preview_manifest) + ) ), kind="software_job_report", ))