Files
cleveragents-core/benchmarks/subplan_execution_bench.py
T
brent.edwards 01b6eb1804
CI / build (push) Successful in 17s
CI / helm (push) Successful in 22s
CI / lint (push) Successful in 28s
CI / typecheck (push) Successful in 47s
CI / benchmark-regression (push) Has been skipped
CI / quality (push) Successful in 3m49s
CI / security (push) Successful in 4m11s
CI / unit_tests (push) Successful in 9m19s
CI / docker (push) Successful in 1m22s
CI / coverage (push) Successful in 12m35s
CI / e2e_tests (push) Successful in 16m17s
CI / integration_tests (push) Successful in 25m5s
CI / status-check (push) Successful in 2s
CI / benchmark-publish (push) Successful in 28m31s
feat(autonomy): parallel execution scales to 10+ concurrent subplans (#1201)
## Summary

Add M6 parallel-scaling coverage for 10+ concurrent subplans:

- **15-subplan parallel scenario** with explicit peak-concurrency bound checks (`max_parallel=10`) and thread-safe concurrency tracking via `_build_executor()`.
- **Deep hierarchical decomposition** coverage (4+ levels) with adjusted leaf condition that only stops early when hitting `max_depth` or when the workset is trivially small (`min_files_per_subplan`).
- **Non-progress guard** in `_build_hierarchy` to prevent pathological recursion when clustering cannot meaningfully split the file set.
- **Small-project regression test** (< 50 files) verifying decomposition depth does not increase unexpectedly with the relaxed leaf condition.
- **ASV benchmark** for 15-subplan parallel execution with `max_parallel=10` to track scaling behavior.

### Removed from this PR

The `_build_hierarchy` child-linkage correctness fix (returning `node_id` from recursive calls instead of using `nodes[-1].node_id`) has been **removed** per review feedback — it is a separate bug fix and will be submitted as an independent issue/PR per CONTRIBUTING.md §Atomic Commits.

## Approach

- **Concurrency tracking:** The `_build_executor()` closure in step definitions detects `context.concurrency_counter` / `context.concurrency_lock` and performs thread-safe peak tracking in a try/finally block.
- **Leaf condition:** Replaced the `max_files_per_subplan` / `max_tokens_per_subplan` leaf check with a `min_files_per_subplan` check to allow deeper decomposition for large projects. Added a non-progress guard so clustering that cannot split the file set terminates immediately rather than recursing to `max_depth`.
- **Deterministic IDs:** `_ids_for_count()` preserves legacy fixed IDs for the first 5 subplans and generates additional deterministic IDs for scale scenarios.

## Validation

### Passing
- `nox -s lint` — all checks passed
- `nox -s typecheck` — 0 errors, 0 warnings
- `nox -s unit_tests` — 12,988 scenarios passed, 0 failed
- `nox -s coverage_report` — 97% (passes `--fail-under=97`)

Closes #855

Reviewed-on: #1201
Co-authored-by: Brent E. Edwards <brent.edwards@cleverthis.com>
Co-committed-by: Brent E. Edwards <brent.edwards@cleverthis.com>
2026-03-31 23:57:39 +00:00

171 lines
6.3 KiB
Python

"""Airspeed Velocity benchmarks for subplan execution scheduler overhead.
Measures scheduling, merge, and integrated execution+merge performance for
SubplanExecutionService and SubplanMergeService across all execution modes
and merge strategies.
"""
from __future__ import annotations
from cleveragents.application.services.subplan_execution_service import (
SubplanExecutionOutput,
SubplanExecutionService,
)
from cleveragents.application.services.subplan_merge_service import (
SubplanMergeService,
)
from cleveragents.domain.models.core.plan import (
ExecutionMode,
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/bench-sub",
)
def _noop_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,
)
class SubplanExecutionSchedulerSuite:
"""Benchmark SubplanExecutionService scheduling overhead."""
def setup(self) -> None:
"""Prepare fixtures for scheduler benchmarks."""
self.statuses = [_make_status(_S1), _make_status(_S2), _make_status(_S3)]
self.statuses_15 = [
_make_status(f"01HGZ6FE0AQDYTR4BXVQ{i:06d}") for i in range(15)
]
self.seq_config = SubplanConfig(execution_mode=ExecutionMode.SEQUENTIAL)
self.par_config = SubplanConfig(
execution_mode=ExecutionMode.PARALLEL, max_parallel=3
)
self.par_scale_config = SubplanConfig(
execution_mode=ExecutionMode.PARALLEL,
max_parallel=10,
)
self.dep_config = SubplanConfig(execution_mode=ExecutionMode.DEPENDENCY_ORDERED)
self.dep_graph = {_S1: [], _S2: [_S1], _S3: [_S2]}
def time_sequential_execution(self) -> None:
"""Time sequential execution of 3 subplans."""
service = SubplanExecutionService(
config=self.seq_config, executor_fn=_noop_executor
)
service.execute_all(subplan_statuses=self.statuses, base_files={})
def time_parallel_execution(self) -> None:
"""Time parallel execution of 3 subplans."""
service = SubplanExecutionService(
config=self.par_config, executor_fn=_noop_executor
)
service.execute_all(subplan_statuses=self.statuses, base_files={})
def time_parallel_execution_15_subplans(self) -> None:
"""Time parallel execution of 15 subplans with max_parallel=10."""
service = SubplanExecutionService(
config=self.par_scale_config,
executor_fn=_noop_executor,
)
service.execute_all(subplan_statuses=self.statuses_15, base_files={})
def time_dependency_ordered_execution(self) -> None:
"""Time dependency-ordered execution of 3 subplans."""
service = SubplanExecutionService(
config=self.dep_config, executor_fn=_noop_executor
)
service.execute_all(
subplan_statuses=self.statuses,
base_files={},
dependency_graph=self.dep_graph,
)
def time_service_construction(self) -> None:
"""Time constructing SubplanExecutionService."""
SubplanExecutionService(config=self.seq_config, executor_fn=_noop_executor)
class SubplanMergeStrategySuite:
"""Benchmark SubplanMergeService strategy overhead."""
def setup(self) -> None:
"""Prepare fixtures for merge benchmarks."""
self.base_files = {"src/main.py": "line1\nline2\nline3\n"}
self.outputs = [
(_S1, {"src/main.py": "line1\nline2\nline3\nnew_a\n"}),
(_S2, {"src/main.py": "new_b\nline1\nline2\nline3\n"}),
]
self.overlapping_outputs = [
(_S1, {"src/main.py": "first version\n"}),
(_S2, {"src/main.py": "second version\n"}),
]
def time_git_three_way_merge(self) -> None:
"""Time git three-way merge of non-overlapping changes."""
service = SubplanMergeService(SubplanMergeStrategy.GIT_THREE_WAY)
service.merge(self.base_files, self.outputs)
def time_sequential_apply_merge(self) -> None:
"""Time sequential apply merge."""
service = SubplanMergeService(SubplanMergeStrategy.SEQUENTIAL_APPLY)
service.merge(self.base_files, self.overlapping_outputs)
def time_last_wins_merge(self) -> None:
"""Time last-wins merge."""
service = SubplanMergeService(SubplanMergeStrategy.LAST_WINS)
service.merge(self.base_files, self.overlapping_outputs)
def time_merge_service_construction(self) -> None:
"""Time constructing SubplanMergeService."""
SubplanMergeService(SubplanMergeStrategy.GIT_THREE_WAY)
def time_single_file_merge(self) -> None:
"""Time merging a single file output."""
service = SubplanMergeService(SubplanMergeStrategy.LAST_WINS)
service.merge(
{"src/main.py": "base\n"},
[(_S1, {"src/main.py": "modified\n"})],
)
class SubplanIntegrationSuite:
"""Benchmark integrated execution+merge overhead."""
def setup(self) -> None:
"""Prepare fixtures for integration benchmarks."""
self.statuses = [_make_status(_S1), _make_status(_S2)]
def time_sequential_with_last_wins(self) -> None:
"""Time sequential execution with last-wins merge."""
config = SubplanConfig(
execution_mode=ExecutionMode.SEQUENTIAL,
merge_strategy=SubplanMergeStrategy.LAST_WINS,
)
service = SubplanExecutionService(config=config, executor_fn=_noop_executor)
service.execute_all(subplan_statuses=self.statuses, base_files={})
def time_parallel_with_git_merge(self) -> None:
"""Time parallel execution with git three-way merge."""
config = SubplanConfig(
execution_mode=ExecutionMode.PARALLEL,
max_parallel=2,
merge_strategy=SubplanMergeStrategy.GIT_THREE_WAY,
)
service = SubplanExecutionService(config=config, executor_fn=_noop_executor)
service.execute_all(subplan_statuses=self.statuses, base_files={})