"""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()