diff --git a/scripts/run_behave_parallel.py b/scripts/run_behave_parallel.py new file mode 100644 index 000000000..3a8bd5d94 --- /dev/null +++ b/scripts/run_behave_parallel.py @@ -0,0 +1,332 @@ +"""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()