"""Helper script for subplan execution Robot Framework smoke tests. Exercises SubplanExecutionService, SubplanMergeService, and integrated execution+merge workflows without requiring the full service layer. """ from __future__ import annotations import sys import threading from cleveragents.application.services.subplan_execution_service import ( SubplanExecutionOutput, SubplanExecutionService, ) from cleveragents.application.services.subplan_merge_service import ( MergeConflictError, SubplanMergeService, ) from cleveragents.domain.models.core.plan import ( ExecutionMode, ProcessingState, SubplanConfig, SubplanMergeStrategy, SubplanStatus, ) _S1 = "01HGZ6FE0AQDYTR4BXVQZ6EA00" _S2 = "01HGZ6FE0AQDYTR4BXVQZ6EB00" _S3 = "01HGZ6FE0AQDYTR4BXVQZ6EC00" def _make_status(subplan_id: str) -> SubplanStatus: return SubplanStatus( subplan_id=subplan_id, action_name="local/robot-sub-action", ) def _success_executor(status: SubplanStatus) -> SubplanExecutionOutput: return SubplanExecutionOutput( subplan_id=status.subplan_id, success=True, files={f"src/{status.subplan_id[-4:]}.py": f"# {status.subplan_id}\n"}, files_changed=1, changeset_summary="Robot test output", ) def _sequential_all() -> None: """Verify sequential execution completes all subplans.""" statuses = [_make_status(_S1), _make_status(_S2), _make_status(_S3)] config = SubplanConfig(execution_mode=ExecutionMode.SEQUENTIAL) service = SubplanExecutionService(config=config, executor_fn=_success_executor) result = service.execute_all(subplan_statuses=statuses, base_files={}) assert result.all_succeeded, f"Not all succeeded: {result.failed_subplan_ids}" completed = [s for s in result.statuses if s.status == ProcessingState.COMPLETE] assert len(completed) == 3, f"Expected 3, got {len(completed)}" print("sequential-all-ok") def _sequential_failfast() -> None: """Verify sequential fail_fast stops after first failure.""" fail_set = {_S2} def executor(status: SubplanStatus) -> SubplanExecutionOutput: if status.subplan_id in fail_set: return SubplanExecutionOutput( subplan_id=status.subplan_id, success=False, error="ValidationError: bad input", ) return _success_executor(status) statuses = [_make_status(_S1), _make_status(_S2), _make_status(_S3)] config = SubplanConfig( execution_mode=ExecutionMode.SEQUENTIAL, fail_fast=True, ) service = SubplanExecutionService(config=config, executor_fn=executor) result = service.execute_all(subplan_statuses=statuses, base_files={}) assert result.statuses[0].status == ProcessingState.COMPLETE assert result.statuses[1].status == ProcessingState.ERRORED assert result.statuses[2].status == ProcessingState.CANCELLED print("sequential-failfast-ok") def _parallel_all() -> None: """Verify parallel execution completes all subplans.""" statuses = [_make_status(_S1), _make_status(_S2), _make_status(_S3)] config = SubplanConfig(execution_mode=ExecutionMode.PARALLEL, max_parallel=3) service = SubplanExecutionService(config=config, executor_fn=_success_executor) result = service.execute_all(subplan_statuses=statuses, base_files={}) assert result.all_succeeded print("parallel-all-ok") def _dep_ordered() -> None: """Verify dependency-ordered execution follows topological order.""" order: list[str] = [] lock = threading.Lock() def executor(status: SubplanStatus) -> SubplanExecutionOutput: with lock: order.append(status.subplan_id) return _success_executor(status) statuses = [_make_status(_S1), _make_status(_S2), _make_status(_S3)] config = SubplanConfig(execution_mode=ExecutionMode.DEPENDENCY_ORDERED) dep_graph = {_S1: [], _S2: [_S1], _S3: [_S2]} service = SubplanExecutionService(config=config, executor_fn=executor) result = service.execute_all( subplan_statuses=statuses, base_files={}, dependency_graph=dep_graph, ) assert result.all_succeeded # Verify order: first occurrences should be A, B, C seen: list[str] = [] for sid in order: if sid not in seen: seen.append(sid) assert seen == [_S1, _S2, _S3], f"Expected [A,B,C], got {seen}" print("dep-ordered-ok") def _merge_git_clean() -> None: """Verify git three-way merge with non-overlapping changes.""" service = SubplanMergeService(SubplanMergeStrategy.GIT_THREE_WAY) base = {"src/main.py": "line1\nline2\nline3\n"} outputs = [ (_S1, {"src/main.py": "line1\nline2\nline3\nnew_a\n"}), (_S2, {"src/main.py": "new_b\nline1\nline2\nline3\n"}), ] result = service.merge(base, outputs) assert result.success, f"Merge failed: {result.conflict_files}" print("merge-git-clean-ok") def _merge_sequential() -> None: """Verify sequential apply merge.""" service = SubplanMergeService(SubplanMergeStrategy.SEQUENTIAL_APPLY) base = {"src/main.py": "original\n"} outputs = [ (_S1, {"src/main.py": "first\n"}), (_S2, {"src/main.py": "second\n"}), ] result = service.merge(base, outputs) assert result.success assert "second" in result.merged_files[0].content print("merge-sequential-ok") def _merge_last_wins() -> None: """Verify last-wins merge takes final output.""" service = SubplanMergeService(SubplanMergeStrategy.LAST_WINS) base = {"src/main.py": "original\n"} outputs = [ (_S1, {"src/main.py": "first\n"}), (_S2, {"src/main.py": "second\n"}), ] result = service.merge(base, outputs) assert result.success assert result.merged_files[0].content == "second\n" print("merge-last-wins-ok") def _merge_fail_conflict() -> None: """Verify fail-on-conflict raises MergeConflictError.""" service = SubplanMergeService(SubplanMergeStrategy.FAIL_ON_CONFLICT) base = {"src/main.py": "base line\n"} outputs = [ (_S1, {"src/main.py": "alpha\n"}), (_S2, {"src/main.py": "beta\n"}), ] try: service.merge(base, outputs) raise AssertionError("Expected MergeConflictError") except MergeConflictError: pass print("merge-fail-conflict-ok") def _exec_merge_integration() -> None: """Verify end-to-end execution + merge integration.""" statuses = [_make_status(_S1), _make_status(_S2)] config = SubplanConfig( execution_mode=ExecutionMode.SEQUENTIAL, merge_strategy=SubplanMergeStrategy.LAST_WINS, ) service = SubplanExecutionService(config=config, executor_fn=_success_executor) result = service.execute_all(subplan_statuses=statuses, base_files={}) assert result.all_succeeded assert result.merge_result is not None assert result.merge_result.success print("exec-merge-integration-ok") def _validation_guards() -> None: """Verify service validation guards.""" # None config try: SubplanExecutionService( config=None, # type: ignore[arg-type] executor_fn=_success_executor, ) raise AssertionError("Expected ValueError for None config") except ValueError: pass # None executor try: SubplanExecutionService( config=SubplanConfig(), executor_fn=None, # type: ignore[arg-type] ) raise AssertionError("Expected ValueError for None executor") except ValueError: pass # None strategy try: SubplanMergeService(strategy=None) # type: ignore[arg-type] raise AssertionError("Expected ValueError for None strategy") except ValueError: pass # Empty statuses svc = SubplanExecutionService(config=SubplanConfig(), executor_fn=_success_executor) try: svc.execute_all(subplan_statuses=[], base_files={}) raise AssertionError("Expected ValueError for empty statuses") except ValueError: pass print("validation-guards-ok") def _exec_summary() -> None: """Verify execution result summary report.""" statuses = [_make_status(_S1), _make_status(_S2)] config = SubplanConfig(execution_mode=ExecutionMode.SEQUENTIAL) service = SubplanExecutionService(config=config, executor_fn=_success_executor) result = service.execute_all(subplan_statuses=statuses, base_files={}) assert result.all_succeeded assert result.total_duration_ms >= 0 assert len(result.statuses) == 2 assert len(result.failed_subplan_ids) == 0 print("exec-summary-ok") def _retry_non_retriable() -> None: """Verify non-retriable error is not retried.""" def executor(status: SubplanStatus) -> SubplanExecutionOutput: return SubplanExecutionOutput( subplan_id=status.subplan_id, success=False, error="ConfigurationError: invalid config", ) config = SubplanConfig(retry_failed=True, max_retries=2) service = SubplanExecutionService(config=config, executor_fn=executor) result = service.execute_all(subplan_statuses=[_make_status(_S1)], base_files={}) assert not result.all_succeeded assert result.statuses[0].attempt_number == 1 print("retry-non-retriable-ok") def _retry_exception_then_succeed() -> None: """Verify executor exception triggers retry and then succeeds.""" call_count: list[int] = [0] def executor(status: SubplanStatus) -> SubplanExecutionOutput: call_count[0] += 1 if call_count[0] <= 1: raise RuntimeError("TimeoutError: simulated exception") return SubplanExecutionOutput( subplan_id=status.subplan_id, success=True, files={"src/main.py": "# recovered\n"}, files_changed=1, ) config = SubplanConfig(retry_failed=True, max_retries=2) service = SubplanExecutionService(config=config, executor_fn=executor) result = service.execute_all(subplan_statuses=[_make_status(_S1)], base_files={}) assert result.all_succeeded assert result.statuses[0].attempt_number == 2 print("retry-exception-ok") def _merge_conflict_exec() -> None: """Verify merge conflict during execution marks not all succeeded.""" def executor(status: SubplanStatus) -> SubplanExecutionOutput: return SubplanExecutionOutput( subplan_id=status.subplan_id, success=True, files={"src/main.py": f"conflict-{status.subplan_id}\n"}, files_changed=1, ) config = SubplanConfig(merge_strategy=SubplanMergeStrategy.FAIL_ON_CONFLICT) service = SubplanExecutionService(config=config, executor_fn=executor) result = service.execute_all( subplan_statuses=[_make_status(_S1), _make_status(_S2)], base_files={"src/main.py": "base\n"}, ) assert not result.all_succeeded print("merge-conflict-exec-ok") def main() -> None: """Dispatch command from sys.argv.""" if len(sys.argv) < 2: raise SystemExit("Expected command argument") command = sys.argv[1] commands: dict[str, object] = { "sequential-all": _sequential_all, "sequential-failfast": _sequential_failfast, "parallel-all": _parallel_all, "dep-ordered": _dep_ordered, "merge-git-clean": _merge_git_clean, "merge-sequential": _merge_sequential, "merge-last-wins": _merge_last_wins, "merge-fail-conflict": _merge_fail_conflict, "exec-merge-integration": _exec_merge_integration, "validation-guards": _validation_guards, "exec-summary": _exec_summary, "retry-non-retriable": _retry_non_retriable, "retry-exception": _retry_exception_then_succeed, "merge-conflict-exec": _merge_conflict_exec, } if command not in commands: raise SystemExit(f"Unknown command: {command}") func = commands[command] if callable(func): func() if __name__ == "__main__": main()