Files
cleveragents-core/scripts/run_behave_parallel.py
T
freemo 37586882b3
CI / benchmark-publish (pull_request) Has been skipped
CI / lint (pull_request) Failing after 17s
CI / build (pull_request) Successful in 21s
CI / helm (pull_request) Successful in 23s
CI / security (pull_request) Failing after 46s
CI / typecheck (pull_request) Failing after 49s
CI / coverage (pull_request) Has been skipped
CI / benchmark-regression (pull_request) Has been skipped
CI / unit_tests (pull_request) Failing after 1m48s
CI / docker (pull_request) Has been skipped
CI / quality (pull_request) Successful in 3m41s
CI / e2e_tests (pull_request) Failing after 16m17s
CI / integration_tests (pull_request) Failing after 21m20s
CI / status-check (pull_request) Failing after 1s
chore(ci): extract behave-parallel runner script from noxfile.py into scripts/
Move the large embedded `_BEHAVE_PARALLEL_CLI_SOURCE` string constant out of
noxfile.py and into a standalone `scripts/run_behave_parallel.py` module.

The `_install_behave_parallel()` helper now reads the script from disk via
`Path(__file__).parent / 'scripts' / 'run_behave_parallel.py'` instead of
embedding the source as a raw string literal. This allows ruff to lint and
type-check the runner independently, and makes noxfile.py significantly
shorter and easier to read.

No functional changes: the installed `behave-parallel` entry point is
identical to the previous embedded version. Parallel and sequential modes,
coverage integration, and the multiprocessing fork model are all preserved.

Fixed two SIM105 lint violations in the extracted script (replaced
try/except/pass with contextlib.suppress).

ISSUES CLOSED: #1538
2026-04-02 23:44:30 +00:00

333 lines
11 KiB
Python

"""In-process parallel behave runner.
Replaces the old subprocess-per-feature model with direct use of
behave's ``Runner`` API. Step definitions and environment hooks are
loaded once per process; feature files are parsed and executed without
Python interpreter startup overhead.
Parallelism modes
-----------------
* **Sequential** (``--processes 1`` or ``BEHAVE_PARALLEL_COVERAGE=1``):
All features run in a single ``Runner.run()`` call.
* **Parallel** (``--processes N``, N > 1, no coverage):
Features are split into *N* equal-size chunks. A
``multiprocessing.Pool`` with the ``fork`` start method dispatches
each chunk to a worker that creates its own ``Runner`` (hooks and
step definitions are re-loaded cheaply because all heavy modules are
already in memory from the parent).
"""
from __future__ import annotations
import argparse
import io
import multiprocessing
import os
import sys
import time
from contextlib import redirect_stderr, redirect_stdout, suppress
from pathlib import Path
from typing import Any
DEFAULT_FEATURE_ROOT = "features/"
# Type alias for summary dictionaries
Summary = dict[str, Any]
# ---------------------------------------------------------------------------
# Summary helpers
# ---------------------------------------------------------------------------
def _empty_summary() -> Summary:
return {
"features": {"passed": 0, "failed": 0, "errors": 0, "skipped": 0},
"scenarios": {"passed": 0, "failed": 0, "errors": 0, "skipped": 0},
"steps": {"passed": 0, "failed": 0, "errors": 0, "skipped": 0},
"duration": 0.0,
}
def _extract_summary(runner: Any) -> Summary:
"""Build a summary dict from a completed behave Runner."""
summary = _empty_summary()
for feature in runner.features:
status_name = feature.status.name
if status_name == "passed":
summary["features"]["passed"] += 1
elif status_name == "skipped":
summary["features"]["skipped"] += 1
else:
summary["features"]["failed"] += 1
summary["duration"] += feature.duration or 0.0
for scenario in feature.walk_scenarios():
sname = scenario.status.name
if sname == "passed":
summary["scenarios"]["passed"] += 1
elif sname == "skipped":
summary["scenarios"]["skipped"] += 1
elif sname == "failed":
summary["scenarios"]["failed"] += 1
else:
summary["scenarios"]["errors"] += 1
for step in scenario.steps:
stname = step.status.name
if stname == "passed":
summary["steps"]["passed"] += 1
elif stname == "skipped":
summary["steps"]["skipped"] += 1
elif stname == "failed":
summary["steps"]["failed"] += 1
else:
summary["steps"]["errors"] += 1
return summary
def _merge_summaries(summaries: list[Summary]) -> Summary:
total = _empty_summary()
for s in summaries:
for bucket in ("features", "scenarios", "steps"):
for field in ("passed", "failed", "errors", "skipped"):
total[bucket][field] += s.get(bucket, {}).get(field, 0)
total["duration"] += float(s.get("duration", 0.0))
return total
def _format_duration(seconds: float) -> str:
minutes = int(seconds // 60)
remainder = seconds % 60
if minutes:
return f"{minutes}m {remainder:.3f}s"
return f"{remainder:.3f}s"
def _print_overall_summary(
total: Summary,
wall_seconds: float | None = None,
) -> None:
print("\nOverall summary:")
for bucket in ("features", "scenarios", "steps"):
b = total[bucket]
print(
f"{b['passed']} {bucket} passed, {b['failed']} failed, "
f"{b['errors']} errored, {b['skipped']} skipped"
)
if total["duration"]:
print(f"Took {_format_duration(float(total['duration']))}")
if wall_seconds is not None:
print(f"Wall time: {_format_duration(wall_seconds)}")
def _has_failures(total: Summary) -> bool:
return (
total["features"]["failed"] > 0
or total["features"]["errors"] > 0
or total["scenarios"]["failed"] > 0
or total["scenarios"]["errors"] > 0
)
def _no_scenarios_ran(total: Summary) -> bool:
"""Return True when no scenario reached a terminal state.
This catches runner-level crashes (e.g. ``before_all`` failure) that
prevent any scenario from executing. Without this guard the summary
would contain all-zero counters, ``_has_failures()`` would return
``False``, and CI would silently pass a broken suite.
"""
s = total["scenarios"]
return s["passed"] + s["failed"] + s["errors"] + s["skipped"] == 0
# ---------------------------------------------------------------------------
# Feature discovery
# ---------------------------------------------------------------------------
def _iter_features(paths: list[str]) -> list[str]:
collected = []
for path in paths:
p = Path(path)
if p.is_dir():
collected.extend(sorted(str(fp) for fp in p.rglob("*.feature")))
else:
collected.append(str(p))
return collected
def _extract_features_and_args(argv: list[str]) -> tuple[list[str], list[str]]:
positional = []
options = []
for arg in argv:
if arg.startswith("-"):
options.append(arg)
else:
positional.append(arg)
if positional:
return positional, options
return [DEFAULT_FEATURE_ROOT], options
# ---------------------------------------------------------------------------
# In-process behave execution
# ---------------------------------------------------------------------------
def _make_runner(feature_paths: list[str], behave_args: list[str]) -> Any:
"""Create a behave Runner with proper configuration defaults.
Mirrors the format-defaulting logic from ``behave.__main__.run_behave``
so that ``-q`` and bare invocations get a sensible formatter instead
of crashing on ``config.format is None``.
"""
from behave.configuration import Configuration
from behave.runner import Runner
from behave.runner_util import reset_runtime
reset_runtime()
args = list(behave_args) + [str(p) for p in feature_paths]
config = Configuration(command_args=args)
if not config.format:
config.format = [config.default_format]
return Runner(config)
def _run_features_inprocess(
feature_paths: list[str], behave_args: list[str]
) -> tuple[bool, Summary]:
"""Run *all* feature_paths in a single behave Runner invocation.
Returns ``(failed: bool, summary: dict)``.
"""
runner = _make_runner(feature_paths, behave_args)
failed = runner.run()
summary = _extract_summary(runner)
return failed, summary
def _worker_run_features(
payload: tuple[list[str], list[str]],
) -> tuple[bool, str, str, Summary]:
"""Entry point for multiprocessing workers.
Runs a chunk of feature files in a forked child process. Heavy
Python modules (cleveragents, behave, SQLAlchemy, ...) are already
loaded in the parent and shared via copy-on-write after ``fork``.
``load_step_definitions()`` and ``load_hooks()`` still execute
inside each worker (they ``exec()`` the step .py files and
environment.py), but every ``import`` they trigger is a cache hit.
"""
feature_paths, behave_args = payload
stdout_buf = io.StringIO()
stderr_buf = io.StringIO()
runner = _make_runner(feature_paths, behave_args)
with redirect_stdout(stdout_buf), redirect_stderr(stderr_buf):
failed = runner.run()
summary = _extract_summary(runner)
return failed, stdout_buf.getvalue(), stderr_buf.getvalue(), summary
# ---------------------------------------------------------------------------
# CLI entry point
# ---------------------------------------------------------------------------
def main(argv: list[str] | None = None) -> None:
argv = list(sys.argv[1:] if argv is None else argv)
parser = argparse.ArgumentParser(add_help=False)
parser.add_argument("--processes", "-j", type=int, default=None)
known, remaining = parser.parse_known_args(argv)
processes = known.processes or os.cpu_count() or 1
feature_args, other_args = _extract_features_and_args(remaining)
# Pass-through --help / --version to behave
if any(flag in remaining for flag in ("-h", "--help", "--version")):
import subprocess
code = subprocess.run(
[sys.executable, "-m", "behave", *other_args, *feature_args]
).returncode
sys.exit(code)
feature_paths = _iter_features(feature_args)
if not feature_paths:
print("No feature files found", file=sys.stderr)
sys.exit(0)
coverage_mode = bool(os.environ.get("BEHAVE_PARALLEL_COVERAGE"))
start = time.monotonic()
if processes <= 1 or coverage_mode or len(feature_paths) == 1:
# ---- sequential in-process mode ----
_, total = _run_features_inprocess(feature_paths, other_args)
else:
# ---- parallel in-process mode (multiprocessing fork) ----
# Pre-import heavy modules so forked children get them for free.
with suppress(ImportError):
import cleveragents # noqa: F401
with suppress(ImportError):
import behave # noqa: F401
# Split features into roughly equal chunks.
chunk_size = max(1, (len(feature_paths) + processes - 1) // processes)
chunks = [
feature_paths[i : i + chunk_size]
for i in range(0, len(feature_paths), chunk_size)
]
ctx = multiprocessing.get_context("fork")
with ctx.Pool(processes=min(processes, len(chunks))) as pool:
results = pool.map(
_worker_run_features,
[(chunk, other_args) for chunk in chunks],
)
summaries = []
for _worker_failed, stdout, stderr, summary in results:
if stdout:
print(stdout, end="")
if stderr:
print(stderr, end="", file=sys.stderr)
summaries.append(summary)
total = _merge_summaries(summaries)
wall = time.monotonic() - start
_print_overall_summary(total, wall_seconds=wall)
# Use the summary-based check rather than the raw runner ``failed``
# boolean. The ``@tdd_expected_fail`` handler in environment.py
# inverts scenario statuses for TDD issue-capture tests, but behave's
# ``runner.run()`` tracks step failures in a local variable that
# cannot be updated by after_scenario hooks. Relying solely on the
# summary (which reflects the corrected scenario statuses) ensures
# that TDD-inverted scenarios do not cause a spurious exit-code 1.
if _has_failures(total):
sys.exit(1)
# Safety net: if features were requested but zero scenarios ran, the
# runner crashed before executing any scenario (e.g. ``before_all``
# failure). Treat this as a failure so CI does not silently pass.
if feature_paths and _no_scenarios_ran(total):
print(
"ERROR: features were requested but no scenarios ran — "
"possible runner-level crash.",
file=sys.stderr,
)
sys.exit(1)
if __name__ == "__main__":
main()