From 72f7bd91075277f69dcbee5532fa00033afdcc48 Mon Sep 17 00:00:00 2001 From: caoqianming Date: Wed, 19 Aug 2026 12:31:05 +0800 Subject: [PATCH] =?UTF-8?q?fix(software):=20=E4=BF=AE=E5=A4=8D=E6=8C=81?= =?UTF-8?q?=E4=B9=85=E5=B7=A5=E7=A8=8B=E9=A6=96=E6=AC=A1=E5=BB=BA=E5=8D=95?= =?UTF-8?q?=E5=A4=B1=E8=B4=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGELOG.md | 2 + core/software_jobs.py | 18 +++- tests/test_software_job_creation.py | 126 ++++++++++++++++++++++++++++ 3 files changed, 143 insertions(+), 3 deletions(-) create mode 100644 tests/test_software_job_creation.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 8a6ae61..d562793 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,8 @@ ## Unreleased +- 修复专业软件持久工程首次建单失败后反复重试的问题;新任务会先建立工程空间再提交执行,并在数据库完整性异常时返回可诊断的错误。 + - 专业软件任务失败或取消后会在原对话留下持久通知,页面离线、刷新或对话暂时忙碌时也不会错过结果。 - 查询专业软件是否可用时不再误提交真实任务;ANSYS 静力分析会要求使用几何中已有的 Named Selection,并在输入缺少命名选择时返回明确错误。 diff --git a/core/software_jobs.py b/core/software_jobs.py index e58524e..71dc832 100644 --- a/core/software_jobs.py +++ b/core/software_jobs.py @@ -438,16 +438,28 @@ def create_job( with session.begin_nested(): if new_workspace is not None: session.add(new_workspace) + # SoftwareJob.workspace_id and SoftwareWorkspace.head_job_id form a + # nullable FK cycle. No ORM relationship exists to give the unit of + # work an insertion dependency, so flush the parent explicitly before + # inserting the first job in a workspace. + session.flush([new_workspace]) session.add(row) - session.flush() + session.flush([row]) return _job_dict(row), True - except IntegrityError: + except IntegrityError as exc: existing = session.execute( select(SoftwareJob).where( SoftwareJob.user_id == user_id, SoftwareJob.idempotency_key == key, ) - ).scalar_one() + ).scalar_one_or_none() + if existing is None: + diagnostic = getattr(getattr(exc, "orig", None), "diag", None) + constraint = getattr(diagnostic, "constraint_name", None) + suffix = f" ({constraint})" if constraint else "" + raise SoftwareJobError( + f"software job database integrity validation failed{suffix}" + ) from exc if ( existing.task_id != task_id or existing.capability != capability diff --git a/tests/test_software_job_creation.py b/tests/test_software_job_creation.py new file mode 100644 index 0000000..0a93e69 --- /dev/null +++ b/tests/test_software_job_creation.py @@ -0,0 +1,126 @@ +from __future__ import annotations + +import unittest +from contextlib import contextmanager, nullcontext +from types import SimpleNamespace +from unittest.mock import MagicMock, call, patch +from uuid import uuid4 + +from sqlalchemy.exc import IntegrityError + +from core.software_jobs import SoftwareJobError, create_job +from core.storage.models import SoftwareJob, SoftwareWorkspace + + +class _Result: + def __init__(self, *, first=None, scalar=None, scalars=None): + self._first = first + self._scalar = scalar + self._scalars = scalars or [] + + def first(self): + return self._first + + def scalar_one_or_none(self): + return self._scalar + + def scalars(self): + return SimpleNamespace(all=lambda: self._scalars) + + +class SoftwareJobCreationTests(unittest.TestCase): + def setUp(self) -> None: + self.user_id = uuid4() + self.task_id = uuid4() + self.artifact_id = uuid4() + self.request = { + "schema_version": 2, + "inputs": [{"key": "data", "artifact_id": str(self.artifact_id)}], + "operation": {"plot": {"type": "line"}}, + "outputs": [{"key": "figure_png", "type": "figure", "format": "png"}], + } + self.contract = SimpleNamespace( + workspace=SimpleNamespace(), + output_namespace="origin", + input_policy={ + "suffixes": [".xlsx"], + "max_bytes": 1024, + "max_total_bytes": 2048, + }, + normalize_request=lambda request: (request, "a" * 64), + input_bindings=lambda request: request["inputs"], + ) + self.artifact = SimpleNamespace( + artifact_id=self.artifact_id, + user_id=self.user_id, + status="active", + current_path="demo.xlsx", + size_bytes=128, + content_sha256="b" * 64, + ) + + def _session(self, *, fail_first_flush: bool = False): + session = MagicMock() + results = iter([ + _Result(first=(self.task_id,)), + _Result(scalar=None), + _Result(scalars=[self.artifact]), + _Result(scalar=None), + _Result(scalar=None), + ]) + session.execute.side_effect = lambda _statement: next(results) + session.begin_nested.return_value = nullcontext() + if fail_first_flush: + session.flush.side_effect = IntegrityError( + "insert workspace", {}, Exception("foreign key violation") + ) + return session + + @staticmethod + def _scope(session): + @contextmanager + def scope(): + yield session + + return scope + + def _create(self, session): + with ( + patch("core.software_jobs.session_scope", self._scope(session)), + patch("core.software_jobs.get_contract", return_value=self.contract), + ): + return create_job( + self.user_id, + self.task_id, + idempotency_key="workspace-first-job", + capability="origin.plot@v2", + request=self.request, + ) + + def test_first_workspace_job_flushes_workspace_before_job(self) -> None: + session = self._session() + + job, created = self._create(session) + + self.assertTrue(created) + self.assertEqual(job["workspace_id"], str(session.add.call_args_list[0].args[0].workspace_id)) + workspace = session.add.call_args_list[0].args[0] + row = session.add.call_args_list[1].args[0] + self.assertIsInstance(workspace, SoftwareWorkspace) + self.assertIsInstance(row, SoftwareJob) + self.assertEqual(row.workspace_id, workspace.workspace_id) + self.assertEqual(session.flush.call_args_list, [call([workspace]), call([row])]) + + def test_non_idempotency_integrity_error_is_not_masked_as_no_result(self) -> None: + session = self._session(fail_first_flush=True) + + with self.assertRaisesRegex( + SoftwareJobError, "database integrity validation failed" + ): + self._create(session) + + self.assertEqual(session.execute.call_count, 5) + + +if __name__ == "__main__": + unittest.main()