ae52db3147
All three cost real release cycles, and none is fixed by widening a budget. daemon_application_cancels_physical_job_only_after_final_session waited for the SUBSCRIBER COUNT to reach 2, then cancelled both sessions and asserted the physical job had started exactly once. The job starts asynchronously after subscription, so both cancels could land first and leave starts == 0 -- arguably the correct outcome. It now waits for the state the assertions actually require. Verified 47/47 on Windows, the only platform it ever failed. The parallel harness refused outright when the suite leader had already exited, because taskkill /T cannot walk a tree from a dead PID. But the leader can exit between the timeout decision and that call, so the harness itself lost a race: a natural exit at the wrong moment failed the whole wave. It now proves cleanup the only way still available -- nothing parented to that PID -- and its contract asserts the PROPERTY rather than the phrase "tree cleanup" it used to grep for. That string pin is what broke when the guard was reworded while behaving correctly; the contract now checks rc==2 AND that the descendant really did survive, which would also catch a guard that claims to fail closed while leaking. extract_wide_flat_file_is_linear took ONE sample per size, so the ratio carried the noise of both. On a loaded Windows VM linear code measured 51x against a 40x bound (184ms -> 9387ms). Best-of-N instead: timing noise only ever adds time, so the minimum is the cheapest good estimate of the noise-free cost. The bound is deliberately unchanged -- it sits where linear (~20x) and quadratic (~128x) are each >=2x away, so raising it would move the test toward the very signal it exists to catch. Now measures 19.1x on Windows, 21.4x on macOS. Also: the smoke's `cli` helper redirected stderr to a file and discarded it, so any of the 10 bare `VAR=$(cli ...)` assignments could kill the run under `set -euo pipefail` printing NOTHING. One such abort cost a full Windows cycle just to locate and still could not be attributed. It now surfaces the command and its stderr. Neutral wording on purpose: one call site expects a non-zero exit and must not read as a failure. Signed-off-by: Martin Vogel <martin.vogel.tech@gmail.com>
447 lines
15 KiB
Python
Executable File
447 lines
15 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""Run one wave of C test suites without nested shell worker processes.
|
|
|
|
The caller owns suite selection, sharding, and final union/count checks. This
|
|
helper owns native child processes directly, writes one result for every suite,
|
|
and bounds a child that never exits. Keeping accounting in this single parent
|
|
avoids an MSYS2 failure mode where a completed native child left its `bash -c`
|
|
worker permanently stuck before the result append.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import os
|
|
import pathlib
|
|
import re
|
|
import signal
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
from dataclasses import dataclass
|
|
|
|
|
|
SUITE_NAME = re.compile(r"^[a-z0-9_]+$")
|
|
SUMMARY = re.compile(r"^ (?P<passed>[0-9]+) passed")
|
|
FAILED = re.compile(r"(?:^|, )(?P<failed>[0-9]+) failed")
|
|
SKIPPED = re.compile(r"(?:^|, )(?P<skipped>[0-9]+) skipped")
|
|
SLOW_SUITES = frozenset(("incremental", "store_arch", "daemon_runtime"))
|
|
POLL_SECONDS = 0.05
|
|
|
|
|
|
@dataclass
|
|
class ActiveSuite:
|
|
name: str
|
|
process: subprocess.Popen[bytes]
|
|
log_path: pathlib.Path
|
|
log_file: object
|
|
started: float
|
|
timeout: int
|
|
|
|
|
|
def parse_args() -> argparse.Namespace:
|
|
parser = argparse.ArgumentParser(
|
|
description="Run a bounded parallel wave of test-runner suites"
|
|
)
|
|
parser.add_argument("--suite-file", required=True, type=pathlib.Path)
|
|
parser.add_argument("--log-dir", required=True, type=pathlib.Path)
|
|
parser.add_argument("--results-file", required=True, type=pathlib.Path)
|
|
parser.add_argument("--jobs", required=True, type=int)
|
|
parser.add_argument("--timeout", required=True, type=int)
|
|
parser.add_argument("--slow-timeout", required=True, type=int)
|
|
parser.add_argument("--kill-grace", required=True, type=int)
|
|
parser.add_argument(
|
|
"--test-post-exit-barrier-dir",
|
|
type=pathlib.Path,
|
|
help=argparse.SUPPRESS,
|
|
)
|
|
parser.add_argument(
|
|
"--test-pre-terminate-barrier-dir",
|
|
type=pathlib.Path,
|
|
help=argparse.SUPPRESS,
|
|
)
|
|
parser.add_argument(
|
|
"runner_command",
|
|
nargs="+",
|
|
help="runner executable and any fixed arguments; suite name is appended",
|
|
)
|
|
args = parser.parse_args()
|
|
for name in ("jobs", "timeout", "slow_timeout", "kill_grace"):
|
|
if getattr(args, name) < 1:
|
|
parser.error(f"--{name.replace('_', '-')} must be at least 1")
|
|
return args
|
|
|
|
|
|
def read_suites(path: pathlib.Path) -> list[str]:
|
|
try:
|
|
suites = path.read_text(encoding="utf-8").splitlines()
|
|
except OSError as exc:
|
|
raise RuntimeError(f"cannot read suite file {path}: {exc}") from exc
|
|
malformed = [suite for suite in suites if SUITE_NAME.fullmatch(suite) is None]
|
|
if malformed:
|
|
raise RuntimeError(f"malformed suite name in {path}: {malformed[0]!r}")
|
|
if len(set(suites)) != len(suites):
|
|
raise RuntimeError(f"duplicate suite name in {path}")
|
|
return suites
|
|
|
|
|
|
def append_log(path: pathlib.Path, message: str) -> None:
|
|
with path.open("a", encoding="utf-8", newline="\n") as stream:
|
|
stream.write(message)
|
|
stream.write("\n")
|
|
|
|
|
|
def start_suite(
|
|
suite: str,
|
|
runner_command: list[str],
|
|
log_dir: pathlib.Path,
|
|
timeout: int,
|
|
) -> ActiveSuite | None:
|
|
log_path = log_dir / f"{suite}.log"
|
|
log_file = log_path.open("wb")
|
|
popen_args: dict[str, object] = {
|
|
"stdout": log_file,
|
|
"stderr": subprocess.STDOUT,
|
|
}
|
|
if os.name == "nt":
|
|
popen_args["creationflags"] = subprocess.CREATE_NEW_PROCESS_GROUP
|
|
else:
|
|
popen_args["start_new_session"] = True
|
|
try:
|
|
process = subprocess.Popen(runner_command + [suite], **popen_args)
|
|
except OSError as exc:
|
|
log_file.close()
|
|
append_log(log_path, f" FAIL: could not start suite {suite!r}: {exc}")
|
|
return None
|
|
return ActiveSuite(
|
|
name=suite,
|
|
process=process,
|
|
log_path=log_path,
|
|
log_file=log_file,
|
|
started=time.monotonic(),
|
|
timeout=timeout,
|
|
)
|
|
|
|
|
|
def windows_descendants(pid: int, timeout: int) -> bool:
|
|
"""True if any live process still claims `pid` as its parent.
|
|
|
|
Used only when the suite leader has already exited: `taskkill /T` cannot
|
|
walk a tree from a dead PID, so cleanup is proven by asking whether anything
|
|
is still parented to it. One level deep on purpose -- Windows does not
|
|
reparent orphans, so a grandchild keeps pointing at its own (dead) parent
|
|
and would not be found here. That is a weaker proof than taskkill /T, which
|
|
is why it is reserved for the case where the strong proof is impossible.
|
|
"""
|
|
try:
|
|
completed = subprocess.run(
|
|
[
|
|
"powershell.exe",
|
|
"-NoProfile",
|
|
"-NonInteractive",
|
|
"-Command",
|
|
"@(Get-CimInstance Win32_Process -Filter "
|
|
f"'ParentProcessId={pid}').Count",
|
|
],
|
|
check=False,
|
|
stdin=subprocess.DEVNULL,
|
|
capture_output=True,
|
|
text=True,
|
|
timeout=timeout,
|
|
)
|
|
except (OSError, subprocess.TimeoutExpired):
|
|
return True # cannot prove absence -> assume the worst
|
|
if completed.returncode != 0:
|
|
return True
|
|
return (completed.stdout or "").strip() not in ("0", "")
|
|
|
|
|
|
def terminate_process_tree(active: ActiveSuite, kill_grace: int) -> None:
|
|
process = active.process
|
|
leader_exited = process.poll() is not None
|
|
if os.name == "nt":
|
|
if leader_exited:
|
|
# The leader can exit on its own between the timeout decision and
|
|
# this call. Refusing outright made the harness itself lose a race:
|
|
# a natural exit at the wrong moment failed the whole wave, which is
|
|
# how a deliberately-hanging fixture suite reddened a release run.
|
|
# taskkill /T cannot walk a tree from a dead PID, so prove cleanup
|
|
# the only way still available -- nothing is parented to it.
|
|
if windows_descendants(process.pid, kill_grace):
|
|
raise RuntimeError(
|
|
f"suite {active.name!r} leader exited leaving live descendants"
|
|
)
|
|
return
|
|
try:
|
|
completed = subprocess.run(
|
|
[
|
|
"taskkill.exe",
|
|
"/PID",
|
|
str(process.pid),
|
|
"/T",
|
|
"/F",
|
|
],
|
|
check=False,
|
|
stdin=subprocess.DEVNULL,
|
|
stdout=subprocess.DEVNULL,
|
|
stderr=subprocess.DEVNULL,
|
|
timeout=kill_grace,
|
|
)
|
|
except (OSError, subprocess.TimeoutExpired):
|
|
completed = None
|
|
if completed is None or completed.returncode != 0:
|
|
if process.poll() is None:
|
|
process.kill()
|
|
try:
|
|
process.wait(timeout=kill_grace)
|
|
except subprocess.TimeoutExpired:
|
|
pass
|
|
raise RuntimeError(
|
|
f"suite {active.name!r} taskkill could not prove process-tree cleanup"
|
|
)
|
|
try:
|
|
process.wait(timeout=kill_grace)
|
|
except subprocess.TimeoutExpired as exc:
|
|
process.kill()
|
|
raise RuntimeError(
|
|
f"suite {active.name!r} process tree resisted forced termination"
|
|
) from exc
|
|
return
|
|
|
|
def group_active() -> bool:
|
|
try:
|
|
os.killpg(process.pid, 0)
|
|
return True
|
|
except ProcessLookupError:
|
|
return False
|
|
except PermissionError:
|
|
return True
|
|
|
|
def wait_for_group_exit(deadline: float) -> bool:
|
|
while time.monotonic() < deadline:
|
|
process.poll()
|
|
if not group_active():
|
|
return True
|
|
time.sleep(POLL_SECONDS)
|
|
process.poll()
|
|
return not group_active()
|
|
|
|
try:
|
|
os.killpg(process.pid, signal.SIGTERM)
|
|
except ProcessLookupError:
|
|
return
|
|
if wait_for_group_exit(time.monotonic() + kill_grace):
|
|
return
|
|
try:
|
|
os.killpg(process.pid, signal.SIGKILL)
|
|
except ProcessLookupError:
|
|
return
|
|
if wait_for_group_exit(time.monotonic() + kill_grace):
|
|
return
|
|
raise RuntimeError(
|
|
f"suite {active.name!r} process group persisted after forced termination"
|
|
)
|
|
|
|
|
|
def wait_for_test_pre_terminate_barrier(
|
|
barrier_dir: pathlib.Path | None,
|
|
active: ActiveSuite,
|
|
) -> None:
|
|
if barrier_dir is None:
|
|
return
|
|
hold = barrier_dir / f"{active.name}.hold"
|
|
if not hold.exists():
|
|
return
|
|
ready = barrier_dir / f"{active.name}.ready"
|
|
leader_exited = barrier_dir / f"{active.name}.leader-exited"
|
|
release = barrier_dir / f"{active.name}.release"
|
|
ready.write_text(f"{active.process.pid}\n", encoding="utf-8")
|
|
deadline = time.monotonic() + 10
|
|
while not release.exists():
|
|
returncode = active.process.poll()
|
|
if returncode is not None and not leader_exited.exists():
|
|
leader_exited.write_text(f"{returncode}\n", encoding="utf-8")
|
|
if time.monotonic() >= deadline:
|
|
raise RuntimeError(
|
|
f"test pre-terminate barrier for {active.name!r} was not released"
|
|
)
|
|
time.sleep(POLL_SECONDS)
|
|
|
|
|
|
def wait_for_test_post_exit_barrier(
|
|
barrier_dir: pathlib.Path | None,
|
|
suite: str,
|
|
) -> None:
|
|
if barrier_dir is None:
|
|
return
|
|
hold = barrier_dir / f"{suite}.hold"
|
|
if not hold.exists():
|
|
return
|
|
ready = barrier_dir / f"{suite}.ready"
|
|
release = barrier_dir / f"{suite}.release"
|
|
ready.write_text("child exited; result intentionally not recorded\n", encoding="utf-8")
|
|
deadline = time.monotonic() + 10
|
|
while not release.exists():
|
|
if time.monotonic() >= deadline:
|
|
raise RuntimeError(f"test post-exit barrier for {suite!r} was not released")
|
|
time.sleep(POLL_SECONDS)
|
|
|
|
|
|
def parse_summary(log_path: pathlib.Path) -> tuple[int, int, int] | None:
|
|
try:
|
|
stream = log_path.open(encoding="utf-8", errors="replace")
|
|
except OSError as exc:
|
|
raise RuntimeError(f"cannot read suite log {log_path}: {exc}") from exc
|
|
last_summary = None
|
|
with stream:
|
|
for line in stream:
|
|
match = SUMMARY.match(line)
|
|
if match is None:
|
|
continue
|
|
failed = FAILED.search(line)
|
|
skipped = SKIPPED.search(line)
|
|
last_summary = (
|
|
int(match.group("passed")),
|
|
int(failed.group("failed")) if failed is not None else 0,
|
|
int(skipped.group("skipped")) if skipped is not None else 0,
|
|
)
|
|
return last_summary
|
|
|
|
|
|
def record_result(
|
|
active: ActiveSuite,
|
|
returncode: int,
|
|
results_file: pathlib.Path,
|
|
timed_out: bool,
|
|
) -> None:
|
|
active.log_file.close()
|
|
elapsed = max(0, int(time.monotonic() - active.started))
|
|
if timed_out:
|
|
returncode = 124
|
|
append_log(
|
|
active.log_path,
|
|
f" FAIL: suite {active.name!r} exceeded {active.timeout}s wall clock "
|
|
"(killed as hung)",
|
|
)
|
|
summary = parse_summary(active.log_path)
|
|
if returncode == 0 and summary is None:
|
|
returncode = 97
|
|
append_log(
|
|
active.log_path,
|
|
f" FAIL: suite {active.name!r} exited 0 without a completion summary "
|
|
"(ran nothing?)",
|
|
)
|
|
passed, failed, skipped = summary or (0, 0, 0)
|
|
result = (
|
|
f"{active.name} rc={returncode} pass={passed} fail={failed} "
|
|
f"skip={skipped} secs={elapsed}"
|
|
)
|
|
with results_file.open("a", encoding="utf-8", newline="\n") as stream:
|
|
stream.write(result)
|
|
stream.write("\n")
|
|
print(f" {result}", flush=True)
|
|
|
|
|
|
def record_start_failure(
|
|
suite: str,
|
|
log_dir: pathlib.Path,
|
|
results_file: pathlib.Path,
|
|
) -> None:
|
|
result = f"{suite} rc=98 pass=0 fail=0 skip=0 secs=0"
|
|
with results_file.open("a", encoding="utf-8", newline="\n") as stream:
|
|
stream.write(result)
|
|
stream.write("\n")
|
|
print(f" {result}", flush=True)
|
|
if not (log_dir / f"{suite}.log").exists():
|
|
append_log(log_dir / f"{suite}.log", f" FAIL: suite {suite!r} did not start")
|
|
|
|
|
|
def run_wave(args: argparse.Namespace) -> None:
|
|
suites = read_suites(args.suite_file)
|
|
args.log_dir.mkdir(parents=True, exist_ok=True)
|
|
args.results_file.parent.mkdir(parents=True, exist_ok=True)
|
|
args.results_file.touch(exist_ok=True)
|
|
pending = list(suites)
|
|
active: dict[str, ActiveSuite] = {}
|
|
try:
|
|
while pending or active:
|
|
while pending and len(active) < args.jobs:
|
|
suite = pending.pop(0)
|
|
timeout = (
|
|
args.slow_timeout if suite in SLOW_SUITES else args.timeout
|
|
)
|
|
started = start_suite(
|
|
suite,
|
|
list(args.runner_command),
|
|
args.log_dir,
|
|
timeout,
|
|
)
|
|
if started is None:
|
|
record_start_failure(suite, args.log_dir, args.results_file)
|
|
else:
|
|
active[suite] = started
|
|
|
|
made_progress = False
|
|
now = time.monotonic()
|
|
for suite, running in list(active.items()):
|
|
returncode = running.process.poll()
|
|
timed_out = returncode is None and now - running.started >= running.timeout
|
|
if returncode is None and not timed_out:
|
|
continue
|
|
if timed_out:
|
|
wait_for_test_pre_terminate_barrier(
|
|
args.test_pre_terminate_barrier_dir,
|
|
running,
|
|
)
|
|
terminate_process_tree(running, args.kill_grace)
|
|
returncode = running.process.returncode
|
|
wait_for_test_post_exit_barrier(
|
|
args.test_post_exit_barrier_dir,
|
|
suite,
|
|
)
|
|
record_result(
|
|
running,
|
|
int(returncode if returncode is not None else 124),
|
|
args.results_file,
|
|
timed_out,
|
|
)
|
|
del active[suite]
|
|
made_progress = True
|
|
if active and not made_progress:
|
|
time.sleep(POLL_SECONDS)
|
|
finally:
|
|
cleanup_errors: list[str] = []
|
|
for running in active.values():
|
|
try:
|
|
terminate_process_tree(running, args.kill_grace)
|
|
except (OSError, RuntimeError) as exc:
|
|
cleanup_errors.append(f"{running.name}: {exc}")
|
|
finally:
|
|
try:
|
|
running.log_file.close()
|
|
except OSError as exc:
|
|
cleanup_errors.append(f"{running.name} log close: {exc}")
|
|
if cleanup_errors:
|
|
raise RuntimeError(
|
|
"parallel scheduler cleanup failed: " + "; ".join(cleanup_errors)
|
|
)
|
|
|
|
|
|
def main() -> int:
|
|
args = parse_args()
|
|
try:
|
|
if os.environ.get("MSYSTEM") and os.name != "nt":
|
|
raise RuntimeError(
|
|
"Windows/MSYS test runs require the native MinGW Python "
|
|
"(os.name must be 'nt')"
|
|
)
|
|
run_wave(args)
|
|
except (OSError, RuntimeError) as exc:
|
|
print(f"FAIL: parallel scheduler infrastructure error: {exc}", file=sys.stderr)
|
|
return 2
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|