Task 002: add shared domain contracts
This commit is contained in:
@@ -5,6 +5,10 @@ ENV PYTHONUNBUFFERED=1
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
COPY apps/runner/requirements.txt ./requirements.txt
|
||||
RUN pip install --no-cache-dir -r requirements.txt
|
||||
|
||||
COPY apps/backend/src/domain/contracts ./backend_contracts/contracts
|
||||
COPY apps/runner/src ./src
|
||||
|
||||
EXPOSE 8010
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
pydantic==2.13.4
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib.util
|
||||
import os
|
||||
import sys
|
||||
from functools import lru_cache
|
||||
from pathlib import Path
|
||||
from typing import Any, Mapping
|
||||
|
||||
|
||||
@lru_cache(maxsize=1)
|
||||
def _backend_contracts_module() -> Any:
|
||||
contracts_dir = _backend_contracts_dir()
|
||||
init_file = contracts_dir / "__init__.py"
|
||||
module_name = "_pipeline_backend_domain_contracts"
|
||||
|
||||
existing = sys.modules.get(module_name)
|
||||
if existing is not None and Path(str(existing.__file__)).resolve() == init_file.resolve():
|
||||
return existing
|
||||
sys.modules.pop(module_name, None)
|
||||
|
||||
spec = importlib.util.spec_from_file_location(
|
||||
module_name,
|
||||
init_file,
|
||||
submodule_search_locations=[str(contracts_dir)],
|
||||
)
|
||||
if spec is None or spec.loader is None:
|
||||
raise RuntimeError(f"Cannot load backend contract module from {init_file}")
|
||||
|
||||
module = importlib.util.module_from_spec(spec)
|
||||
sys.modules[module_name] = module
|
||||
spec.loader.exec_module(module)
|
||||
return module
|
||||
|
||||
|
||||
def _backend_contracts_dir() -> Path:
|
||||
for candidate in _backend_contract_candidates():
|
||||
init_file = candidate / "__init__.py"
|
||||
if init_file.is_file():
|
||||
return candidate
|
||||
|
||||
candidates = ", ".join(str(path) for path in _backend_contract_candidates())
|
||||
raise RuntimeError(f"Cannot find backend contract package. Checked: {candidates}")
|
||||
|
||||
|
||||
def _backend_contract_candidates() -> tuple[Path, ...]:
|
||||
configured = os.environ.get("PIPELINE_BACKEND_CONTRACTS_DIR")
|
||||
candidates: list[Path] = []
|
||||
if configured:
|
||||
candidates.append(Path(configured))
|
||||
|
||||
current_file = Path(__file__).resolve()
|
||||
candidates.append(Path("/app/backend_contracts/contracts"))
|
||||
|
||||
for parent in current_file.parents:
|
||||
candidates.append(parent / "apps" / "backend" / "src" / "domain" / "contracts")
|
||||
|
||||
return tuple(dict.fromkeys(candidates))
|
||||
|
||||
|
||||
def validate_agent_job_output(payload: Mapping[str, Any]) -> Any:
|
||||
contracts = _backend_contracts_module()
|
||||
return contracts.AgentJobOutput.model_validate(payload)
|
||||
@@ -0,0 +1,62 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import sys
|
||||
import os
|
||||
import shutil
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
|
||||
from pydantic import ValidationError
|
||||
|
||||
|
||||
REPO_ROOT = Path(__file__).resolve().parents[3]
|
||||
RUNNER_ROOT = REPO_ROOT / "apps" / "runner"
|
||||
sys.path.insert(0, str(RUNNER_ROOT))
|
||||
|
||||
from src.application.job_output_validation import ( # noqa: E402
|
||||
_backend_contracts_module,
|
||||
validate_agent_job_output,
|
||||
)
|
||||
|
||||
|
||||
class RunnerJobOutputValidationTest(unittest.TestCase):
|
||||
def test_valid_runner_job_output_uses_backend_contracts(self) -> None:
|
||||
output = validate_agent_job_output(
|
||||
{
|
||||
"status": "SUCCEEDED",
|
||||
"output_files": [{"path": "outputs/plan.json"}],
|
||||
"payload": {"artifact": "plan"},
|
||||
}
|
||||
)
|
||||
|
||||
self.assertEqual("SUCCEEDED", output.status.value)
|
||||
|
||||
def test_invalid_runner_job_output_status_is_rejected(self) -> None:
|
||||
with self.assertRaises(ValidationError):
|
||||
validate_agent_job_output({"status": "NOT_A_STATUS"})
|
||||
|
||||
def test_runner_image_layout_can_use_copied_backend_contracts(self) -> None:
|
||||
source_contracts = REPO_ROOT / "apps" / "backend" / "src" / "domain" / "contracts"
|
||||
|
||||
with tempfile.TemporaryDirectory() as temp_dir:
|
||||
copied_contracts = Path(temp_dir) / "backend_contracts" / "contracts"
|
||||
shutil.copytree(source_contracts, copied_contracts)
|
||||
|
||||
previous = os.environ.get("PIPELINE_BACKEND_CONTRACTS_DIR")
|
||||
os.environ["PIPELINE_BACKEND_CONTRACTS_DIR"] = str(copied_contracts)
|
||||
_backend_contracts_module.cache_clear()
|
||||
try:
|
||||
output = validate_agent_job_output({"status": "SUCCEEDED"})
|
||||
finally:
|
||||
_backend_contracts_module.cache_clear()
|
||||
if previous is None:
|
||||
os.environ.pop("PIPELINE_BACKEND_CONTRACTS_DIR", None)
|
||||
else:
|
||||
os.environ["PIPELINE_BACKEND_CONTRACTS_DIR"] = previous
|
||||
|
||||
self.assertEqual("SUCCEEDED", output.status.value)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user