test_phase_three_ownership.py 2.8 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980
  1. from __future__ import annotations
  2. from pathlib import Path
  3. import pytest
  4. from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker
  5. from script_build_host.domain.errors import (
  6. MissionAlreadyOwned,
  7. MissionFencingTokenStale,
  8. MissionStopRequested,
  9. )
  10. from script_build_host.domain.records import BuildStatus
  11. from script_build_host.infrastructure.ownership import FencedCommandGate, OwnerLease
  12. from script_build_host.repositories.legacy_state import SqlAlchemyLegacyBuildStateRepository
  13. from script_build_host.repositories.sqlalchemy import SqlAlchemyMissionBindingRepository
  14. async def _bound_build(sessions: async_sessionmaker[AsyncSession]) -> int:
  15. state = SqlAlchemyLegacyBuildStateRepository(sessions)
  16. build_id = await state.create(
  17. execution_id=1,
  18. topic_build_id=2,
  19. topic_id=3,
  20. agent_type="AigcAgent",
  21. agent_config={},
  22. data_source_url=None,
  23. strategies_config={},
  24. root_trace_id="root-owner-test",
  25. )
  26. await SqlAlchemyMissionBindingRepository(sessions).create(
  27. script_build_id=build_id,
  28. root_trace_id="root-owner-test",
  29. input_snapshot_id=1,
  30. engine_version="test",
  31. schema_version="test/v1",
  32. )
  33. return build_id
  34. @pytest.mark.asyncio
  35. async def test_owner_lease_excludes_second_owner_and_grows_epoch(
  36. database: tuple[AsyncEngine, async_sessionmaker[AsyncSession]], tmp_path: Path
  37. ) -> None:
  38. _, sessions = database
  39. build_id = await _bound_build(sessions)
  40. first = OwnerLease(sessions, tmp_path / "leases", owner_instance_id="owner-a")
  41. second = OwnerLease(sessions, tmp_path / "leases", owner_instance_id="owner-b")
  42. first_token = await first.acquire(build_id)
  43. with pytest.raises(MissionAlreadyOwned):
  44. await second.acquire(build_id)
  45. await first.release()
  46. second_token = await second.acquire(build_id)
  47. assert second_token.owner_epoch == first_token.owner_epoch + 1
  48. with pytest.raises(MissionFencingTokenStale):
  49. await FencedCommandGate(sessions).verify(first_token)
  50. await second.release()
  51. @pytest.mark.asyncio
  52. async def test_stop_epoch_fences_current_owner_without_taking_process_lease(
  53. database: tuple[AsyncEngine, async_sessionmaker[AsyncSession]], tmp_path: Path
  54. ) -> None:
  55. _, sessions = database
  56. build_id = await _bound_build(sessions)
  57. lease = OwnerLease(sessions, tmp_path / "leases", owner_instance_id="owner")
  58. token = await lease.acquire(build_id)
  59. gate = FencedCommandGate(sessions)
  60. await gate.verify(token)
  61. assert await gate.request_stop(build_id) == token.stop_epoch + 1
  62. with pytest.raises(MissionStopRequested):
  63. await gate.verify(token)
  64. assert (
  65. await SqlAlchemyLegacyBuildStateRepository(sessions).get_status(build_id)
  66. is BuildStatus.STOPPING
  67. )
  68. await lease.release()