Files
aetherbound-guild/tools/run_test_attempt.py
T

477 lines
17 KiB
Python
Executable File

#!/usr/bin/env python3
"""Run one source-bound test attempt without leaking temporary work roots."""
from __future__ import annotations
import argparse
import hashlib
import json
import os
from pathlib import Path
import shutil
import signal
import subprocess
import sys
import tarfile
import tempfile
import threading
import time
from typing import Any
RUNNER_ID = "aetherbound-test-attempt-v1"
OWNER_FILE = ".abg-attempt-owner.json"
RESULT_FILE = "result.json"
KEEP_FILE = ".keep"
STRICT_DIAGNOSTICS = (
"SCRIPT ERROR:",
"Parse Error:",
"ERROR:",
"ObjectDB instances leaked",
"resources still in use at exit",
"RID allocations of type",
"RID of type",
)
_child: subprocess.Popen[bytes] | None = None
_interrupted_signal: int | None = None
HANDLED_SIGNALS = (signal.SIGINT, signal.SIGTERM, signal.SIGHUP)
PROCESS_GROUP_TERM_TIMEOUT = 1.0
TREE_REMOVE_ATTEMPTS = 20
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--label", required=True, help="Stable family name used for retention")
parser.add_argument("--cwd", type=Path, default=Path.cwd())
parser.add_argument("--git-ref", help="Run from a fresh git archive of this revision")
parser.add_argument("--cwd-subdir", default=".", help="Working directory inside an archived source")
parser.add_argument("--workspace-root", type=Path, default=Path("/private/tmp/aetherbound-guild-test-work"))
parser.add_argument("--artifact-root", type=Path, default=Path("/private/tmp/aetherbound-guild-test-evidence"))
parser.add_argument("--keep-failed", type=int, default=3)
parser.add_argument("--keep-passed", type=int, default=3)
parser.add_argument("--timeout", type=float, default=900.0)
parser.add_argument("--require-marker", action="append", default=[])
parser.add_argument("--allow-diagnostics", action="store_true")
parser.add_argument("command", nargs=argparse.REMAINDER)
args = parser.parse_args()
if args.command and args.command[0] == "--":
args.command = args.command[1:]
if not args.command:
parser.error("a command is required after --")
if args.keep_failed < 0 or args.keep_passed < 0:
parser.error("retention counts must be non-negative")
if args.timeout <= 0:
parser.error("timeout must be positive")
return args
def safe_label(value: str) -> str:
label = "".join(character if character.isalnum() or character in "._-" else "-" for character in value)
label = label.strip("-")
if not label:
raise ValueError("label must contain an alphanumeric character")
return label
def write_json(path: Path, payload: dict[str, Any]) -> None:
path.write_text(json.dumps(payload, ensure_ascii=True, indent=2, sort_keys=True) + "\n", encoding="utf-8")
def process_alive(pid: int) -> bool:
try:
os.kill(pid, 0)
except ProcessLookupError:
return False
except PermissionError:
return True
return True
def process_group_alive(pgid: int) -> bool:
try:
os.killpg(pgid, 0)
except ProcessLookupError:
return False
except PermissionError:
return True
return True
def remove_tree(path: Path) -> None:
for attempt in range(TREE_REMOVE_ATTEMPTS):
try:
shutil.rmtree(path)
return
except FileNotFoundError:
return
except OSError:
if attempt + 1 == TREE_REMOVE_ATTEMPTS:
raise
time.sleep(0.1)
def cleanup_stale_workspaces(workspace_root: Path) -> list[str]:
removed: list[str] = []
if not workspace_root.is_dir():
return removed
for candidate in sorted(workspace_root.iterdir()):
owner_path = candidate / OWNER_FILE
if not candidate.is_dir() or not owner_path.is_file():
continue
try:
owner = json.loads(owner_path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
continue
if owner.get("created_by") != RUNNER_ID:
continue
pid = int(owner.get("pid", 0))
if pid > 0 and process_alive(pid):
continue
remove_tree(candidate)
removed.append(str(candidate))
return removed
def sha256_file(path: Path) -> str:
digest = hashlib.sha256()
with path.open("rb") as source:
for chunk in iter(lambda: source.read(1024 * 1024), b""):
digest.update(chunk)
return digest.hexdigest()
def archive_evidence(source: Path, destination: Path) -> tuple[int, str]:
files = sorted(path for path in source.rglob("*") if path.is_file())
if not files:
return 0, ""
with tarfile.open(destination, "w:gz", compresslevel=6) as archive:
for path in files:
archive.add(path, arcname=Path("evidence") / path.relative_to(source), recursive=False)
return len(files), sha256_file(destination)
def referenced_by_git(repository: Path, candidate: Path) -> bool:
completed = subprocess.run(
["git", "grep", "-F", "-q", "--", str(candidate)],
cwd=repository,
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
check=False,
)
return completed.returncode == 0
def prune_artifacts(family_root: Path, repository: Path, outcomes: set[str], keep: int) -> list[str]:
candidates: list[tuple[str, Path]] = []
for path in family_root.iterdir():
result_path = path / RESULT_FILE
if not path.is_dir() or not result_path.is_file() or (path / KEEP_FILE).exists():
continue
try:
result = json.loads(result_path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
continue
if result.get("created_by") != RUNNER_ID or result.get("outcome") not in outcomes:
continue
if referenced_by_git(repository, path):
continue
candidates.append((str(result.get("started_at", "")), path))
candidates.sort(reverse=True)
removed: list[str] = []
for _started_at, path in candidates[keep:]:
remove_tree(path)
removed.append(str(path))
return removed
def hash_inventory(root: Path) -> list[str]:
rows: list[str] = []
for path in sorted(path for path in root.rglob("*") if path.is_file() and path.name != "inventory.sha256"):
digest = sha256_file(path)
rows.append(f"{digest} {path.relative_to(root)}")
return rows
def inspect_terminal(log_path: Path, required_markers: list[str], allow_diagnostics: bool) -> tuple[list[str], dict[str, int]]:
text = log_path.read_text(encoding="utf-8", errors="replace")
diagnostics: list[str] = []
if not allow_diagnostics:
for line in text.splitlines():
if any(marker in line for marker in STRICT_DIAGNOSTICS):
diagnostics.append(line.strip())
marker_counts = {marker: text.count(marker) for marker in required_markers}
return diagnostics, marker_counts
def extract_source(repository: Path, git_ref: str, destination: Path, log) -> str:
archive_path = destination.parent / "source.tar"
subprocess.run(
["git", "archive", "--format=tar", f"--output={archive_path}", git_ref],
cwd=repository,
stdout=log,
stderr=subprocess.STDOUT,
check=True,
)
source_revision = subprocess.check_output(["git", "rev-parse", git_ref], cwd=repository, text=True).strip()
destination.mkdir(parents=True)
subprocess.run(["/usr/bin/tar", "-xf", str(archive_path), "-C", str(destination)], stdout=log, stderr=subprocess.STDOUT, check=True)
archive_path.unlink()
return source_revision
def pump_output(stream, log) -> None:
while True:
chunk = stream.read(65536)
if not chunk:
return
log.write(chunk)
log.flush()
sys.stdout.buffer.write(chunk)
sys.stdout.buffer.flush()
def terminate_child(sig: int = signal.SIGTERM) -> bool:
global _child
if _child is None:
return True
try:
os.killpg(_child.pid, sig)
except ProcessLookupError:
return True
except PermissionError:
return False
return True
def stop_child_process_group() -> None:
if _child is None:
return
pgid = _child.pid
terminate_child(signal.SIGTERM)
deadline = time.monotonic() + PROCESS_GROUP_TERM_TIMEOUT
while process_group_alive(pgid) and time.monotonic() < deadline:
time.sleep(0.05)
if process_group_alive(pgid):
terminate_child(signal.SIGKILL)
deadline = time.monotonic() + PROCESS_GROUP_TERM_TIMEOUT
while process_group_alive(pgid) and time.monotonic() < deadline:
time.sleep(0.05)
def handle_signal(sig: int, _frame) -> None:
global _interrupted_signal
_interrupted_signal = sig
terminate_child(signal.SIGTERM)
raise KeyboardInterrupt
def main() -> int:
global _child
previous_signal_mask = None
attempt_root: Path | None = None
artifact_dir: Path | None = None
if hasattr(signal, "pthread_sigmask"):
previous_signal_mask = signal.pthread_sigmask(signal.SIG_BLOCK, HANDLED_SIGNALS)
try:
args = parse_args()
label = safe_label(args.label)
repository = args.cwd.resolve()
workspace_root = args.workspace_root.resolve()
artifact_root = args.artifact_root.resolve()
workspace_root.mkdir(parents=True, exist_ok=True)
family_root = artifact_root / label
family_root.mkdir(parents=True, exist_ok=True)
stale_removed = cleanup_stale_workspaces(workspace_root)
attempt_root = Path(tempfile.mkdtemp(prefix=f"{label}-", dir=workspace_root))
started_epoch = time.time()
started_at = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(started_epoch))
attempt_id = time.strftime("%Y%m%dT%H%M%SZ", time.gmtime(started_epoch)) + f"-{os.getpid()}"
artifact_dir = family_root / attempt_id
artifact_dir.mkdir()
command_log = artifact_dir / "command.log"
evidence_dir = attempt_root / "evidence"
evidence_dir.mkdir()
owner = {
"created_by": RUNNER_ID,
"pid": os.getpid(),
"label": label,
"started_at": started_at,
"artifact_dir": str(artifact_dir),
}
write_json(attempt_root / OWNER_FILE, owner)
except BaseException:
if attempt_root is not None:
remove_tree(attempt_root)
if artifact_dir is not None:
remove_tree(artifact_dir)
if previous_signal_mask is not None:
signal_mask = previous_signal_mask
previous_signal_mask = None
signal.pthread_sigmask(signal.SIG_SETMASK, signal_mask)
raise
source_revision = "WORKTREE"
timed_out = False
return_code = 1
archive_count = 0
archive_sha256 = ""
outcome = "FAIL"
error = ""
diagnostics: list[str] = []
marker_counts: dict[str, int] = {}
process_group_removed = True
try:
if previous_signal_mask is not None:
signal_mask = previous_signal_mask
previous_signal_mask = None
signal.pthread_sigmask(signal.SIG_SETMASK, signal_mask)
with command_log.open("wb") as log:
if args.git_ref:
source_root = attempt_root / "source"
source_revision = extract_source(repository, args.git_ref, source_root, log)
run_cwd = (source_root / args.cwd_subdir).resolve()
if source_root not in run_cwd.parents and run_cwd != source_root:
raise ValueError("cwd-subdir escapes the archived source")
else:
run_cwd = (repository / args.cwd_subdir).resolve()
if not run_cwd.is_dir():
raise FileNotFoundError(f"working directory does not exist: {run_cwd}")
env = os.environ.copy()
for name in ("home", "cache", "tmp", "browser"):
(attempt_root / name).mkdir()
env.update({
"HOME": str(attempt_root / "home"),
"ABG_ORIGINAL_HOME": os.environ.get("HOME", ""),
"XDG_CACHE_HOME": str(attempt_root / "cache"),
"TMPDIR": str(attempt_root / "tmp"),
"TMP": str(attempt_root / "tmp"),
"TEMP": str(attempt_root / "tmp"),
"ABG_BROWSER_TEMP_PARENT": str(attempt_root / "browser"),
"ABG_ATTEMPT_ROOT": str(attempt_root),
"ABG_EVIDENCE_DIR": str(evidence_dir),
"ABG_SOURCE_REVISION": source_revision,
})
_child = subprocess.Popen(
args.command,
cwd=run_cwd,
env=env,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
start_new_session=True,
)
assert _child.stdout is not None
pump = threading.Thread(target=pump_output, args=(_child.stdout, log), daemon=True)
pump.start()
try:
return_code = _child.wait(timeout=args.timeout)
except subprocess.TimeoutExpired:
timed_out = True
terminate_child(signal.SIGTERM)
try:
_child.wait(timeout=5)
except subprocess.TimeoutExpired:
terminate_child(signal.SIGKILL)
_child.wait(timeout=5)
return_code = 124
pump.join(timeout=5)
diagnostics, marker_counts = inspect_terminal(
command_log, args.require_marker, args.allow_diagnostics
)
marker_failure = any(count != 1 for count in marker_counts.values())
if not timed_out and return_code == 0 and (diagnostics or marker_failure):
return_code = 86
if diagnostics:
error = "strict runtime diagnostics detected"
else:
error = "required marker count is not exactly one"
outcome = "TIMEOUT" if timed_out else ("PASS" if return_code == 0 else "FAIL")
except KeyboardInterrupt:
outcome = "INTERRUPTED"
return_code = 128 + int(_interrupted_signal or signal.SIGINT)
error = f"interrupted by signal {_interrupted_signal or signal.SIGINT}"
except Exception as exc: # Preserve a compact terminal instead of a work-root leak.
outcome = "HARNESS_ERROR"
return_code = 1
error = f"{type(exc).__name__}: {exc}"
finally:
if _child is not None and _child.poll() is None:
terminate_child(signal.SIGTERM)
try:
_child.wait(timeout=5)
except subprocess.TimeoutExpired:
terminate_child(signal.SIGKILL)
_child.wait(timeout=5)
stop_child_process_group()
process_group_removed = _child is None or not process_group_alive(_child.pid)
if not process_group_removed:
outcome = "HARNESS_ERROR"
return_code = 1
error = "owned child process group remained after TERM/KILL cleanup"
archive_path = artifact_dir / "evidence.tar.gz"
try:
archive_count, archive_sha256 = archive_evidence(evidence_dir, archive_path)
if archive_count == 0 and archive_path.exists():
archive_path.unlink()
except Exception as exc:
outcome = "HARNESS_ERROR"
return_code = 1
error = f"evidence archival failed: {type(exc).__name__}: {exc}"
finally:
remove_tree(attempt_root)
finished_at = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())
result = {
"created_by": RUNNER_ID,
"label": label,
"started_at": started_at,
"finished_at": finished_at,
"command": args.command,
"repository": str(repository),
"source_revision": source_revision,
"outcome": outcome,
"return_code": return_code,
"timed_out": timed_out,
"error": error,
"strict_diagnostics": diagnostics,
"required_marker_counts": marker_counts,
"process_group_removed": process_group_removed,
"workspace_removed": not attempt_root.exists(),
"stale_workspaces_removed": stale_removed,
"evidence_files": archive_count,
"evidence_archive_sha256": archive_sha256,
}
write_json(artifact_dir / RESULT_FILE, result)
inventory = hash_inventory(artifact_dir)
(artifact_dir / "inventory.sha256").write_text("\n".join(inventory) + "\n", encoding="utf-8")
with (family_root / "attempts.jsonl").open("a", encoding="utf-8") as journal:
journal.write(json.dumps({**result, "artifact_dir": str(artifact_dir)}, ensure_ascii=True, sort_keys=True) + "\n")
removed = []
removed.extend(prune_artifacts(family_root, repository, {"PASS"}, args.keep_passed))
removed.extend(
prune_artifacts(
family_root,
repository,
{"FAIL", "TIMEOUT", "INTERRUPTED", "HARNESS_ERROR"},
args.keep_failed,
)
)
if removed:
with (family_root / "retention.jsonl").open("a", encoding="utf-8") as retention_log:
retention_log.write(json.dumps({"at": finished_at, "removed": sorted(removed)}, sort_keys=True) + "\n")
print(
f"ABG_TEST_ATTEMPT_{outcome} label={label} rc={return_code} "
f"workspace_removed={str(not attempt_root.exists()).lower()} artifact={artifact_dir}"
)
return return_code
if __name__ == "__main__":
for handled_signal in HANDLED_SIGNALS:
signal.signal(handled_signal, handle_signal)
raise SystemExit(main())