fix(software): 修复持久工程首次建单失败
This commit is contained in:
parent
6eb382fbda
commit
72f7bd9107
|
|
@ -8,6 +8,8 @@
|
|||
|
||||
## Unreleased
|
||||
|
||||
- 修复专业软件持久工程首次建单失败后反复重试的问题;新任务会先建立工程空间再提交执行,并在数据库完整性异常时返回可诊断的错误。
|
||||
|
||||
- 专业软件任务失败或取消后会在原对话留下持久通知,页面离线、刷新或对话暂时忙碌时也不会错过结果。
|
||||
|
||||
- 查询专业软件是否可用时不再误提交真实任务;ANSYS 静力分析会要求使用几何中已有的 Named Selection,并在输入缺少命名选择时返回明确错误。
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
Loading…
Reference in New Issue