fix(software): 将失败和取消结果通知原对话

This commit is contained in:
caoqianming 2026-08-19 11:07:58 +08:00
parent fac63f6ab5
commit d03a743756
5 changed files with 124 additions and 10 deletions

View File

@ -8,6 +8,8 @@
## Unreleased
- 专业软件任务失败或取消后会在原对话留下持久通知,页面离线、刷新或对话暂时忙碌时也不会错过结果。
- 查询专业软件是否可用时不再误提交真实任务ANSYS 静力分析会要求使用几何中已有的 Named Selection并在输入缺少命名选择时返回明确错误。
- Origin 专业软件任务现在会在原 Windows Node 上保留可继续加工的工程;后续调整直接打开上一版工程,过程只提供轻量预览,工程和其他大文件仅在明确要求下载或交付时才传回并登记为正式产物。

View File

@ -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
if terminal_status == "succeeded":
job.followup_status = "pending"
if terminal_status == "succeeded":
if getattr(job, "workspace_id", None) is not None:
workspace = session.execute(
select(SoftwareWorkspace).where(

View File

@ -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()

View File

@ -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

View File

@ -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=(
[]
if job.status != "succeeded"
else (
_artifact_refs(job.artifact_manifest)
if job.artifact_manifest
else _preview_refs(job.preview_manifest)
)
),
kind="software_job_report",
))