b8732dfc6f
CI / build (pull_request) Successful in 17s
CI / helm (pull_request) Successful in 17s
CI / push-validation (pull_request) Successful in 10s
CI / lint (pull_request) Successful in 39s
CI / quality (pull_request) Successful in 50s
CI / typecheck (pull_request) Successful in 52s
CI / security (pull_request) Successful in 53s
CI / e2e_tests (pull_request) Successful in 2m13s
CI / coverage (pull_request) Successful in 5m35s
CI / integration_tests (pull_request) Successful in 6m40s
CI / unit_tests (pull_request) Successful in 7m38s
CI / docker (pull_request) Successful in 1m45s
CI / status-check (pull_request) Successful in 1s
CI / lint (push) Successful in 16s
CI / quality (push) Successful in 17s
CI / build (push) Successful in 23s
CI / helm (push) Successful in 24s
CI / typecheck (push) Successful in 53s
CI / security (push) Successful in 53s
CI / push-validation (push) Successful in 38s
CI / e2e_tests (push) Successful in 3m14s
CI / unit_tests (push) Successful in 6m37s
CI / integration_tests (push) Successful in 6m39s
CI / docker (push) Successful in 12s
CI / coverage (push) Successful in 10m53s
CI / status-check (push) Successful in 1s
Adjusted test running and file-detection logic to stabilize unit tests in overlayfs environments and improve target feature handling. - Modified scripts/run_behave_parallel.py to run sequentially when there are 2 or fewer feature files, avoiding fork deadlocks on overlayfs and reducing nox-based unit test timeouts for agent_skills_loader and skill_search features. - Updated noxfile.py to correctly detect feature files in posargs, fixing the prior logic that appended the "features/" directory when specific feature files were provided. This ensures precise test selection and avoids unnecessary path expansion. Rationale: These changes address the root causes of flaky unit test timeouts by preventing problematic forking behavior with small feature sets and by ensuring nox respects explicitly provided feature file paths. ISSUES CLOSED: #9374
405 lines
14 KiB
Python
405 lines
14 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
|
|
import traceback
|
|
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
|
|
|
|
|
|
# Aliases for chunk-level checks — identical semantics to the overall-summary
|
|
# helpers; exposed under distinct names for readability at the call site.
|
|
_chunk_has_failures = _has_failures
|
|
_chunk_no_scenarios_ran = _no_scenarios_ran
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 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.
|
|
|
|
If the worker crashes with an unhandled exception, the traceback is
|
|
captured into stderr and a crash summary with ``features.errors = 1``
|
|
is returned so that the parent can detect the crash via
|
|
``_chunk_has_failures`` (and also ``_chunk_no_scenarios_ran``, since no
|
|
scenarios reached a terminal state) and replay the diagnostics.
|
|
"""
|
|
feature_paths, behave_args = payload
|
|
|
|
stdout_buf = io.StringIO()
|
|
stderr_buf = io.StringIO()
|
|
|
|
try:
|
|
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
|
|
except Exception:
|
|
# Surface the full traceback so the parent can replay it for
|
|
# diagnostics. Set features.errors = 1 so that _has_failures()
|
|
# detects the crash in the merged total even when other workers
|
|
# passed (preventing a silent CI green when a partial worker crash
|
|
# occurs alongside passing chunks).
|
|
tb_text = traceback.format_exc()
|
|
stderr_buf.write(f"\nWORKER CRASH:\n{tb_text}")
|
|
crash_summary = _empty_summary()
|
|
crash_summary["features"]["errors"] = 1
|
|
return True, stdout_buf.getvalue(), stderr_buf.getvalue(), crash_summary
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Parallel result aggregation
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _aggregate_worker_results(
|
|
results: list[tuple[bool, str, str, Summary]],
|
|
) -> Summary:
|
|
"""Aggregate worker results, replaying logs only for failed chunks.
|
|
|
|
Iterates over the list of ``(worker_failed, stdout, stderr, summary)``
|
|
tuples returned by the multiprocessing pool. For each chunk whose
|
|
summary-based check indicates failure (failed/errored scenarios or
|
|
features, or a crash that produced no terminal scenario results), the
|
|
captured stdout and stderr are replayed to the parent's streams.
|
|
Passing chunks have their output suppressed to reduce CI noise.
|
|
|
|
Uses summary-based checks rather than the raw ``worker_failed`` boolean
|
|
so that TDD-inverted chunks (``@tdd_expected_fail`` scenarios whose
|
|
status has been corrected to "passed" by the after_scenario hook) do
|
|
not trigger spurious log replay.
|
|
"""
|
|
summaries: list[Summary] = []
|
|
for idx, (_, stdout, stderr, summary) in enumerate(results):
|
|
# Rely on summary-based checks: the @tdd_expected_fail handler in
|
|
# environment.py inverts scenario statuses but cannot update the
|
|
# runner.run() local variable, so the raw worker_failed flag is
|
|
# not trustworthy for TDD-inverted chunks.
|
|
chunk_failed = _chunk_has_failures(summary) or _chunk_no_scenarios_ran(summary)
|
|
if chunk_failed:
|
|
if stdout:
|
|
print(f"\n--- Worker {idx} stdout (chunk failed) ---")
|
|
print(stdout, end="")
|
|
if stderr:
|
|
print(
|
|
f"\n--- Worker {idx} stderr (chunk failed) ---",
|
|
file=sys.stderr,
|
|
)
|
|
print(stderr, end="", file=sys.stderr)
|
|
summaries.append(summary)
|
|
return _merge_summaries(summaries)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 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()
|
|
|
|
# Run sequentially if:
|
|
# - processes <= 1 (explicitly requested)
|
|
# - coverage_mode (slipcover requires single process)
|
|
# - only 1 feature file (no parallelism benefit)
|
|
# - very few feature files relative to processes (fork overhead > benefit)
|
|
# When feature_paths <= 2, sequential is faster and avoids fork deadlocks
|
|
if (
|
|
processes <= 1
|
|
or coverage_mode
|
|
or len(feature_paths) == 1
|
|
or len(feature_paths) <= 2
|
|
):
|
|
# ---- 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],
|
|
)
|
|
|
|
total = _aggregate_worker_results(results)
|
|
|
|
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()
|