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
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
333 lines
11 KiB
Python
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()
|