"""Helper utilities for plan persistence E2E Robot integration tests. Covers: - Full lifecycle persistence (all phase transitions) - Restart persistence (close + reopen DB) - Concurrent session access - Action CRUD - Plan tree hierarchy """ from __future__ import annotations import sys import tempfile from datetime import UTC, datetime from pathlib import Path from sqlalchemy import create_engine from sqlalchemy.orm import Session from cleveragents.domain.models.core.action import ( Action, ActionArgument, ActionState, ArgumentRequirement, ArgumentType, ) from cleveragents.domain.models.core.plan import ( NamespacedName, Plan, PlanIdentity, PlanPhase, PlanTimestamps, ProcessingState, ProjectLink, ) from cleveragents.infrastructure.database.models import ( init_database, ) from cleveragents.infrastructure.database.repositories import ( ActionRepository, LifecyclePlanRepository, ) def _create_temp_db() -> tuple[str, Session]: """Create a temporary SQLite database and return (path, session).""" tmp = tempfile.mktemp(suffix=".db") db_url = f"sqlite:///{tmp}" engine = init_database(db_url) session = Session(bind=engine) return tmp, session def _make_action(name: str = "local/e2e-action") -> Action: """Create a minimal action for FK satisfaction.""" return Action( namespaced_name=NamespacedName.parse(name), description="E2E test action", definition_of_done="Done", strategy_actor="local/s", execution_actor="local/e", state=ActionState.AVAILABLE, created_at=datetime(2026, 1, 1, tzinfo=UTC), updated_at=datetime(2026, 1, 1, tzinfo=UTC), ) def _make_plan( plan_id: str, action_name: str = "local/e2e-action", phase: PlanPhase = PlanPhase.ACTION, state: ProcessingState = ProcessingState.QUEUED, parent_plan_id: str | None = None, root_plan_id: str | None = None, ) -> Plan: """Create a Plan domain object.""" now = datetime(2026, 3, 1, tzinfo=UTC) return Plan( identity=PlanIdentity( plan_id=plan_id, parent_plan_id=parent_plan_id, root_plan_id=root_plan_id, attempt=1, ), namespaced_name=NamespacedName(namespace="local", name="e2e-plan"), action_name=action_name, description="E2E test plan", definition_of_done="E2E assertions pass", phase=phase, processing_state=state, strategy_actor="local/s", execution_actor="local/e", project_links=[ProjectLink(project_name="local/proj-e2e")], timestamps=PlanTimestamps(created_at=now, updated_at=now), created_by="e2e-test", tags=["e2e"], reusable=True, read_only=False, ) # --------------------------------------------------------------------------- # Full lifecycle # --------------------------------------------------------------------------- def _full_lifecycle() -> None: """E2E: Create plan, transition through all phases to applied.""" tmp, session = _create_temp_db() factory = lambda: session # noqa: E731 try: # Setup action action_repo = ActionRepository(session_factory=factory) action_repo.create(_make_action()) session.commit() plan_repo = LifecyclePlanRepository(session_factory=factory) # Phase 1: Create plan in action phase plan = _make_plan("01HV00000000000000E2E0FC01") plan_repo.create(plan) session.commit() retrieved = plan_repo.get("01HV00000000000000E2E0FC01") assert retrieved is not None assert retrieved.phase == PlanPhase.ACTION # Phase 2: Transition to strategize retrieved.phase = PlanPhase.STRATEGIZE retrieved.processing_state = ProcessingState.PROCESSING retrieved.timestamps.updated_at = datetime.now(tz=UTC) plan_repo.update(retrieved) session.commit() retrieved = plan_repo.get("01HV00000000000000E2E0FC01") assert retrieved is not None assert retrieved.phase == PlanPhase.STRATEGIZE assert retrieved.processing_state == ProcessingState.PROCESSING # Phase 3: Complete strategize -> execute retrieved.phase = PlanPhase.EXECUTE retrieved.processing_state = ProcessingState.QUEUED retrieved.timestamps.updated_at = datetime.now(tz=UTC) plan_repo.update(retrieved) session.commit() retrieved = plan_repo.get("01HV00000000000000E2E0FC01") assert retrieved is not None assert retrieved.phase == PlanPhase.EXECUTE # Phase 4: Execute -> apply -> applied retrieved.phase = PlanPhase.APPLY retrieved.processing_state = ProcessingState.APPLIED retrieved.timestamps.updated_at = datetime.now(tz=UTC) plan_repo.update(retrieved) session.commit() retrieved = plan_repo.get("01HV00000000000000E2E0FC01") assert retrieved is not None assert retrieved.phase == PlanPhase.APPLY assert retrieved.processing_state == ProcessingState.APPLIED assert retrieved.is_terminal print("full-lifecycle-ok") finally: session.close() Path(tmp).unlink(missing_ok=True) # --------------------------------------------------------------------------- # Restart persistence # --------------------------------------------------------------------------- def _restart_persistence() -> None: """E2E: Verify plan survives DB close and reopen.""" tmp, session = _create_temp_db() factory = lambda: session # noqa: E731 try: action_repo = ActionRepository(session_factory=factory) action_repo.create(_make_action()) session.commit() plan_repo = LifecyclePlanRepository(session_factory=factory) plan = _make_plan( "01HV00000000000000E2ERST01", phase=PlanPhase.EXECUTE, state=ProcessingState.PROCESSING, ) plan_repo.create(plan) session.commit() # Close session (simulate process restart) session.close() # Reopen engine2 = create_engine(f"sqlite:///{tmp}", echo=False) session2 = Session(bind=engine2) factory2 = lambda: session2 # noqa: E731 plan_repo2 = LifecyclePlanRepository(session_factory=factory2) retrieved = plan_repo2.get("01HV00000000000000E2ERST01") assert retrieved is not None assert retrieved.phase == PlanPhase.EXECUTE assert retrieved.processing_state == ProcessingState.PROCESSING assert retrieved.description == "E2E test plan" assert len(retrieved.project_links) == 1 session2.close() engine2.dispose() print("restart-persistence-ok") finally: Path(tmp).unlink(missing_ok=True) # --------------------------------------------------------------------------- # Concurrent access # --------------------------------------------------------------------------- def _concurrent_access() -> None: """E2E: Two sessions see each other's committed data.""" tmp, session1 = _create_temp_db() factory1 = lambda: session1 # noqa: E731 try: action_repo = ActionRepository(session_factory=factory1) action_repo.create(_make_action()) session1.commit() plan_repo1 = LifecyclePlanRepository(session_factory=factory1) plan = _make_plan("01HV00000000000000E2ECNC01") plan_repo1.create(plan) session1.commit() # Open second session engine2 = create_engine(f"sqlite:///{tmp}", echo=False) session2 = Session(bind=engine2) factory2 = lambda: session2 # noqa: E731 plan_repo2 = LifecyclePlanRepository(session_factory=factory2) retrieved = plan_repo2.get("01HV00000000000000E2ECNC01") assert retrieved is not None assert retrieved.description == "E2E test plan" session2.close() engine2.dispose() print("concurrent-access-ok") finally: session1.close() Path(tmp).unlink(missing_ok=True) # --------------------------------------------------------------------------- # Action CRUD # --------------------------------------------------------------------------- def _action_crud() -> None: """E2E: Create, read, update, delete an action.""" tmp, session = _create_temp_db() factory = lambda: session # noqa: E731 try: repo = ActionRepository(session_factory=factory) # Create action = Action( namespaced_name=NamespacedName.parse("local/crud-action"), description="CRUD test", definition_of_done="Done", strategy_actor="local/s", execution_actor="local/e", arguments=[ ActionArgument( name="target", arg_type=ArgumentType.STRING, requirement=ArgumentRequirement.REQUIRED, description="Target path", ), ], invariants=["No deletions"], state=ActionState.AVAILABLE, created_at=datetime(2026, 1, 1, tzinfo=UTC), updated_at=datetime(2026, 1, 1, tzinfo=UTC), ) repo.create(action) session.commit() # Read fetched = repo.get_by_name("local/crud-action") assert fetched is not None assert fetched.description == "CRUD test" assert len(fetched.arguments) == 1 assert len(fetched.invariants) == 1 # Update fetched.description = "Updated CRUD test" fetched.updated_at = datetime.now(tz=UTC) repo.update(fetched) session.commit() refetched = repo.get_by_name("local/crud-action") assert refetched is not None assert refetched.description == "Updated CRUD test" # Delete deleted = repo.delete("local/crud-action") assert deleted is True session.commit() gone = repo.get_by_name("local/crud-action") assert gone is None print("action-crud-ok") finally: session.close() Path(tmp).unlink(missing_ok=True) # --------------------------------------------------------------------------- # Plan tree # --------------------------------------------------------------------------- def _plan_tree() -> None: """E2E: Parent-child plan hierarchy persistence.""" tmp, session = _create_temp_db() factory = lambda: session # noqa: E731 try: action_repo = ActionRepository(session_factory=factory) action_repo.create(_make_action()) session.commit() plan_repo = LifecyclePlanRepository(session_factory=factory) # Parent parent = _make_plan("01HV00000000000000E2ETRP01") plan_repo.create(parent) session.commit() # Child child = _make_plan( "01HV00000000000000E2ETRC01", parent_plan_id="01HV00000000000000E2ETRP01", root_plan_id="01HV00000000000000E2ETRP01", ) plan_repo.create(child) session.commit() loaded_child = plan_repo.get("01HV00000000000000E2ETRC01") assert loaded_child is not None assert loaded_child.identity.parent_plan_id == "01HV00000000000000E2ETRP01" assert loaded_child.identity.root_plan_id == "01HV00000000000000E2ETRP01" assert loaded_child.is_subplan loaded_parent = plan_repo.get("01HV00000000000000E2ETRP01") assert loaded_parent is not None assert not loaded_parent.is_subplan print("plan-tree-ok") finally: session.close() Path(tmp).unlink(missing_ok=True) # --------------------------------------------------------------------------- # Main dispatcher # --------------------------------------------------------------------------- def main() -> None: if len(sys.argv) < 2: raise SystemExit("Expected command argument") command = sys.argv[1] commands = { "full-lifecycle": _full_lifecycle, "restart-persistence": _restart_persistence, "concurrent-access": _concurrent_access, "action-crud": _action_crud, "plan-tree": _plan_tree, } if command not in commands: raise SystemExit(f"Unknown command: {command}") commands[command]() if __name__ == "__main__": main()