diff --git a/noxfile.py b/noxfile.py index c1e69a2b..7cf3e6a9 100644 --- a/noxfile.py +++ b/noxfile.py @@ -95,6 +95,10 @@ def _install_behave_parallel(session: nox.Session) -> None: Parallel execution uses ``multiprocessing.Pool`` with the ``fork`` start method so that already-imported modules are shared (copy-on-write) across workers. + + The runner script is read from ``scripts/run_behave_parallel.py`` so + that the noxfile stays concise and the script can be linted and typed + independently. """ tmp_dir = Path(session.create_tmp()) source_dir = tmp_dir / "behave-parallel-inprocess" @@ -103,7 +107,9 @@ def _install_behave_parallel(session: nox.Session) -> None: pkg_dir = source_dir / "behave_parallel" pkg_dir.mkdir(parents=True, exist_ok=True) (pkg_dir / "__init__.py").write_text("\n") - (pkg_dir / "cli.py").write_text(_BEHAVE_PARALLEL_CLI_SOURCE) + + runner_script = Path(__file__).parent / "scripts" / "run_behave_parallel.py" + (pkg_dir / "cli.py").write_text(runner_script.read_text(encoding="utf-8")) setup_path = source_dir / "setup.py" setup_path.write_text( @@ -124,334 +130,6 @@ def _install_behave_parallel(session: nox.Session) -> None: session.install(str(source_dir)) -# The in-process behave-parallel CLI source is kept as a module-level -# constant so the noxfile itself stays readable and ``ruff`` can still -# lint the surrounding Python without tripping over a giant raw string -# embedded inside a function body. -_BEHAVE_PARALLEL_CLI_SOURCE = r''' -"""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 -from pathlib import Path - -DEFAULT_FEATURE_ROOT = "features/" - - -# --------------------------------------------------------------------------- -# Summary helpers -# --------------------------------------------------------------------------- - -def _empty_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): - """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): - 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): - 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, wall_seconds=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): - 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): - """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): - 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): - 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, behave_args): - """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, behave_args): - """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): - """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=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. - try: - import cleveragents # noqa: F401 - except ImportError: - pass - try: - import behave # noqa: F401 - except ImportError: - pass - - # 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() -''' - - # ============================================================================= diff --git a/scripts/run_behave_parallel.py b/scripts/run_behave_parallel.py new file mode 100644 index 00000000..3a8bd5d9 --- /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()