czty's picture
Add files using upload-large-folder tool
c6856c9 verified
Raw
History Blame Contribute Delete
25.6 kB
#!/usr/bin/env python3
from __future__ import annotations
import argparse
import ast
import csv
import importlib.util
import inspect
import json
import multiprocessing as mp
import os
import random
import re
import shutil
import subprocess
import sys
import tempfile
import traceback
from concurrent.futures import ThreadPoolExecutor, as_completed
from collections import Counter, defaultdict
from pathlib import Path
from typing import Any
PROJECT_ROOT = Path(__file__).resolve().parents[2]
DEFAULT_MCP_ROOT = PROJECT_ROOT / "biomni_web" / "backend" / "data" / "mcp_generated"
DEFAULT_CATALOG = PROJECT_ROOT / "graph_outputs" / "mcp_generated_graph_all_20260522_124110" / "server_catalog.json"
DEFAULT_OUTPUT_DIR = PROJECT_ROOT / "experiments" / "execution_audit" / "results"
def normalize_server_name(server_dir: Path) -> str:
name = server_dir.name
return name[4:] if name.startswith("mcp_") else name
def safe_name(name: str) -> str:
return re.sub(r"[^0-9a-zA-Z_]", "_", name)
def classify_error(stderr: str, stdout: str = "", returncode: int | None = None) -> str:
text = f"{stderr}\n{stdout}".lower()
if returncode == 124 or "timeout" in text or "timed out" in text:
return "timeout"
if "no module named" in text or "modulenotfounderror" in text or "importerror" in text:
return "missing_python_dependency"
if "no such file or directory" in text or "not found" in text or "filenotfounderror" in text:
return "missing_executable_or_file"
if "syntaxerror" in text or "indentationerror" in text:
return "syntax_error"
if "validation error" in text or "schema" in text or "missing required" in text:
return "schema_or_argument_error"
if "calledprocesserror" in text or "non-zero exit" in text or "exit status" in text:
return "runtime_nonzero_exit"
if "traceback" in text or "exception" in text or "error" in text:
return "runtime_exception"
if returncode not in (None, 0):
return "nonzero_exit"
return "none"
def run_cmd(cmd: list[str], *, timeout: float, cwd: Path | None = None) -> dict[str, Any]:
try:
proc = subprocess.run(
cmd,
cwd=str(cwd) if cwd else None,
text=True,
capture_output=True,
timeout=timeout,
)
return {
"returncode": proc.returncode,
"stdout": proc.stdout[-4000:],
"stderr": proc.stderr[-4000:],
"timed_out": False,
}
except subprocess.TimeoutExpired as exc:
return {
"returncode": 124,
"stdout": (exc.stdout or "")[-4000:] if isinstance(exc.stdout, str) else "",
"stderr": (exc.stderr or "")[-4000:] if isinstance(exc.stderr, str) else "timeout",
"timed_out": True,
}
def server_files(server_dir: Path) -> tuple[Path | None, Path | None]:
app_dir = server_dir / "app"
shim = next(iter(sorted(app_dir.glob("*_shim_server.py"))), None)
raw = next(
iter(sorted(p for p in app_dir.glob("*_server.py") if not p.name.endswith("_shim_server.py"))),
None,
)
return shim, raw
def audit_syntax(path: Path, python: str, timeout: float) -> dict[str, Any]:
result = run_cmd([python, "-m", "py_compile", str(path)], timeout=timeout)
ok = result["returncode"] == 0
return {
"ok": ok,
"error_type": "none" if ok else classify_error(result["stderr"], result["stdout"], result["returncode"]),
"stderr": result["stderr"],
}
def _import_worker(path: str, queue: mp.Queue) -> None:
try:
module_name = f"_mcp_audit_{safe_name(Path(path).stem)}_{os.getpid()}"
spec = importlib.util.spec_from_file_location(module_name, path)
if spec is None or spec.loader is None:
raise RuntimeError("cannot build import spec")
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
queue.put({"ok": True, "error": ""})
except BaseException:
queue.put({"ok": False, "error": traceback.format_exc()[-4000:]})
def run_import(path: Path, timeout: float) -> dict[str, Any]:
ctx = mp.get_context("spawn")
queue: mp.Queue = ctx.Queue()
proc = ctx.Process(target=_import_worker, args=(str(path), queue))
proc.start()
proc.join(timeout)
if proc.is_alive():
proc.terminate()
proc.join(1)
return {"ok": False, "error_type": "timeout", "stderr": "import timed out"}
payload = queue.get() if not queue.empty() else {"ok": proc.exitcode == 0, "error": ""}
ok = bool(payload.get("ok"))
err = payload.get("error", "")
return {
"ok": ok,
"error_type": "none" if ok else classify_error(err, returncode=proc.exitcode),
"stderr": err,
}
def audit_import_subprocess(path: Path, python: str, timeout: float) -> dict[str, Any]:
code = (
"import importlib.util, pathlib; "
f"p=pathlib.Path({str(path)!r}); "
"spec=importlib.util.spec_from_file_location('_mcp_audit_mod', str(p)); "
"m=importlib.util.module_from_spec(spec); "
"spec.loader.exec_module(m)"
)
result = run_cmd([python, "-c", code], timeout=timeout)
ok = result["returncode"] == 0
return {
"ok": ok,
"error_type": "none" if ok else classify_error(result["stderr"], result["stdout"], result["returncode"]),
"stderr": result["stderr"],
}
def audit_startup(path: Path, python: str, timeout: float, grace: float) -> dict[str, Any]:
proc = subprocess.Popen(
[python, str(path)],
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
)
try:
stdout, stderr = proc.communicate(timeout=grace)
ok = proc.returncode == 0
return {
"ok": ok,
"returncode": proc.returncode,
"error_type": "none" if ok else classify_error(stderr, stdout, proc.returncode),
"stderr": (stderr or "")[-4000:],
}
except subprocess.TimeoutExpired:
proc.terminate()
try:
stdout, stderr = proc.communicate(timeout=timeout)
except subprocess.TimeoutExpired:
proc.kill()
stdout, stderr = proc.communicate(timeout=1)
# A stdio MCP server should keep running while it waits for requests.
if not (stderr or "").strip():
return {"ok": True, "returncode": 124, "error_type": "none", "stderr": ""}
return {
"ok": False,
"returncode": 124,
"error_type": classify_error(stderr, stdout, 124),
"stderr": (stderr or "")[-4000:],
}
def extract_command_literals(source_path: Path) -> list[str]:
try:
text = source_path.read_text(encoding="utf-8")
except OSError:
return []
candidates: list[str] = []
for match in re.finditer(r"(?:command|cmd)\s*(?::[^=]+)?=\s*\[\s*['\"]([^'\"]+)['\"]", text):
candidates.append(match.group(1))
for match in re.finditer(r"subprocess\.(?:run|Popen|check_output|check_call)\(\s*\[\s*['\"]([^'\"]+)['\"]", text):
candidates.append(match.group(1))
cleaned = []
for item in candidates:
base = Path(item).name
if base and base not in {"python", "python3", "Rscript", "bash", "sh"}:
cleaned.append(base)
return sorted(set(cleaned))
def audit_help_version(commands: list[str], timeout: float) -> dict[str, Any]:
if not commands:
return {"probed": False, "ok": False, "command": "", "flag": "", "error_type": "no_command_literal"}
for command in commands:
if shutil.which(command) is None:
continue
for flag in ("--help", "-h", "--version", "-version", "-v"):
result = run_cmd([command, flag], timeout=timeout)
combined = f"{result['stdout']}\n{result['stderr']}".strip()
if result["returncode"] == 0 or combined:
return {
"probed": True,
"ok": True,
"command": command,
"flag": flag,
"returncode": result["returncode"],
"error_type": "none",
}
first = commands[0]
return {
"probed": True,
"ok": False,
"command": first,
"flag": "",
"error_type": "missing_executable_or_file",
}
def load_catalog(path: Path) -> list[dict[str, Any]]:
return json.loads(path.read_text(encoding="utf-8"))
def stratified_tool_sample(catalog: list[dict[str, Any]], sample_size: int, seed: int) -> list[dict[str, Any]]:
by_category: dict[str, list[dict[str, Any]]] = defaultdict(list)
for server in catalog:
for tool in server.get("tools", []):
by_category[server.get("category", "unknown")].append(
{
"server": server["name"],
"category": server.get("category", "unknown"),
"tool": tool.get("name"),
"schema": tool.get("inputSchema", {}),
}
)
rng = random.Random(seed)
total = sum(len(items) for items in by_category.values())
selected: list[dict[str, Any]] = []
for category, items in sorted(by_category.items()):
quota = max(1, round(sample_size * len(items) / total))
rng.shuffle(items)
selected.extend(items[:quota])
if len(selected) > sample_size:
selected = selected[:sample_size]
elif len(selected) < sample_size:
seen = {(x["server"], x["tool"]) for x in selected}
rest = [
item
for items in by_category.values()
for item in items
if (item["server"], item["tool"]) not in seen
]
rng.shuffle(rest)
selected.extend(rest[: sample_size - len(selected)])
return selected
def make_fixture_files(tmp: Path) -> dict[str, Path]:
files = {
"fasta": tmp / "tiny.fa",
"fastq": tmp / "tiny.fastq",
"csv": tmp / "tiny.csv",
"tsv": tmp / "tiny.tsv",
"txt": tmp / "tiny.txt",
"vcf": tmp / "tiny.vcf",
"bam": tmp / "tiny.bam",
"gff": tmp / "tiny.gff",
"json": tmp / "tiny.json",
}
files["fasta"].write_text(">seq1\nACGTACGTACGTACGTACGT\n", encoding="utf-8")
files["fastq"].write_text("@seq1\nACGTACGTACGT\n+\nFFFFFFFFFFFF\n", encoding="utf-8")
files["csv"].write_text("gene,count\nA,1\nB,2\n", encoding="utf-8")
files["tsv"].write_text("gene\tcount\nA\t1\nB\t2\n", encoding="utf-8")
files["txt"].write_text("tiny fixture\n", encoding="utf-8")
files["vcf"].write_text("##fileformat=VCFv4.2\n#CHROM\tPOS\tID\tREF\tALT\tQUAL\tFILTER\tINFO\n1\t1\t.\tA\tC\t.\t.\t.\n", encoding="utf-8")
files["bam"].write_bytes(b"BAM\1")
files["gff"].write_text("##gff-version 3\nseq1\t.\tgene\t1\t10\t.\t+\t.\tID=g1\n", encoding="utf-8")
files["json"].write_text('{"items": [1, 2, 3]}\n', encoding="utf-8")
return files
def value_for_arg(name: str, schema: dict[str, Any], fixtures: dict[str, Path], tmp: Path) -> Any:
lname = name.lower()
typ = (schema or {}).get("type", "string")
if any(k in lname for k in ("out", "output", "report", "result", "prefix", "log")):
return str(tmp / f"{safe_name(name)}.out")
if "fastq" in lname or "fq" in lname:
return str(fixtures["fastq"])
if "fasta" in lname or "fa" == lname or lname.endswith("_fa") or "reference" in lname or "query" in lname:
return str(fixtures["fasta"])
if "vcf" in lname:
return str(fixtures["vcf"])
if "bam" in lname or "sam" in lname:
return str(fixtures["bam"])
if "gff" in lname or "gtf" in lname or "annotation" in lname:
return str(fixtures["gff"])
if "json" in lname:
return str(fixtures["json"])
if "csv" in lname or "table" in lname or "matrix" in lname or "count" in lname:
return str(fixtures["csv"])
if "tsv" in lname:
return str(fixtures["tsv"])
if "file" in lname or "path" in lname or "input" in lname or "contig" in lname or "reads" in lname:
return str(fixtures["fasta"])
if typ == "boolean":
return False
if typ == "number" or typ == "integer":
return 1
if typ == "array":
return [str(fixtures["fasta"])]
if typ == "object":
return {}
return "tiny"
def find_source_for_server(mcp_root: Path, server: str) -> Path | None:
server_dir = mcp_root / f"mcp_{server}"
shim, raw = server_files(server_dir)
return shim or raw
def _tool_worker(source: str, tool_name: str, kwargs: dict[str, Any], queue: mp.Queue) -> None:
try:
path = Path(source)
module_name = f"_mcp_tool_audit_{safe_name(path.stem)}_{os.getpid()}"
spec = importlib.util.spec_from_file_location(module_name, str(path))
if spec is None or spec.loader is None:
raise RuntimeError("cannot build import spec")
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
fn = getattr(module, tool_name, None)
if fn is None and hasattr(module, "_load_functions"):
for candidate in module._load_functions():
if getattr(candidate, "__name__", "") == tool_name:
fn = candidate
break
if fn is None:
raise AttributeError(f"tool function not found: {tool_name}")
signature = inspect.signature(fn)
coerced_kwargs = {}
for key, value in kwargs.items():
annotation = signature.parameters.get(key).annotation if key in signature.parameters else inspect._empty
annotation_text = str(annotation)
if value is not None and ("Path" in annotation_text or annotation is Path):
coerced_kwargs[key] = Path(value)
else:
coerced_kwargs[key] = value
kwargs = coerced_kwargs
result = fn(**kwargs)
queue.put({"ok": True, "result_type": type(result).__name__, "error": ""})
except BaseException:
queue.put({"ok": False, "result_type": "", "error": traceback.format_exc()[-4000:]})
def audit_tool_call(mcp_root: Path, item: dict[str, Any], timeout: float) -> dict[str, Any]:
source = find_source_for_server(mcp_root, item["server"])
if source is None:
return {**item, "ok": False, "error_type": "missing_executable_or_file", "stderr": "server source not found"}
schema = item.get("schema") or {}
properties = schema.get("properties") or {}
required = schema.get("required") or []
with tempfile.TemporaryDirectory(prefix="mcp_tool_audit_") as tmp_s:
tmp = Path(tmp_s)
fixtures = make_fixture_files(tmp)
kwargs = {name: value_for_arg(name, properties.get(name, {}), fixtures, tmp) for name in required}
ctx = mp.get_context("spawn")
queue: mp.Queue = ctx.Queue()
proc = ctx.Process(target=_tool_worker, args=(str(source), item["tool"], kwargs, queue))
proc.start()
proc.join(timeout)
if proc.is_alive():
proc.terminate()
proc.join(1)
return {**item, "ok": False, "error_type": "timeout", "stderr": "tool call timed out", "kwargs": kwargs}
payload = queue.get() if not queue.empty() else {"ok": proc.exitcode == 0, "error": ""}
ok = bool(payload.get("ok"))
err = payload.get("error", "")
return {
**item,
"ok": ok,
"error_type": "none" if ok else classify_error(err, returncode=proc.exitcode),
"stderr": err,
"kwargs": kwargs,
"result_type": payload.get("result_type", ""),
}
def audit_one_server(payload: tuple[str, str, float, float, float, float, float]) -> dict[str, Any]:
server_dir_s, python, syntax_timeout, import_timeout, startup_timeout, startup_grace, help_timeout = payload
server_dir = Path(server_dir_s)
server = normalize_server_name(server_dir)
shim, raw = server_files(server_dir)
target = shim or raw
if target is None:
return {"server": server, "has_source": False}
syntax = audit_syntax(target, python, syntax_timeout)
imported = (
audit_import_subprocess(target, python, import_timeout)
if syntax["ok"]
else {"ok": False, "error_type": syntax["error_type"], "stderr": syntax["stderr"]}
)
startup = (
audit_startup(target, python, startup_timeout, startup_grace)
if syntax["ok"]
else {"ok": False, "error_type": syntax["error_type"], "stderr": syntax["stderr"]}
)
commands = extract_command_literals(raw or target)
help_version = audit_help_version(commands, help_timeout)
return {
"server": server,
"source": str(target),
"has_source": True,
"syntax_ok": syntax["ok"],
"import_ok": imported["ok"],
"startup_ok": startup["ok"],
"help_version_probed": help_version["probed"],
"help_version_ok": help_version["ok"],
"help_version_command": help_version.get("command", ""),
"help_version_flag": help_version.get("flag", ""),
"syntax_error_type": syntax["error_type"],
"import_error_type": imported["error_type"],
"startup_error_type": startup["error_type"],
"help_version_error_type": help_version["error_type"],
}
def pct(n: int, d: int) -> str:
return f"{(100 * n / d):.1f}%" if d else "n/a"
def write_csv(path: Path, rows: list[dict[str, Any]]) -> None:
if not rows:
path.write_text("", encoding="utf-8")
return
fieldnames = sorted({key for row in rows for key in row})
with path.open("w", newline="", encoding="utf-8") as handle:
writer = csv.DictWriter(handle, fieldnames=fieldnames)
writer.writeheader()
writer.writerows(rows)
def main() -> None:
parser = argparse.ArgumentParser(description="Execution audit for generated BioinfoMCP servers.")
parser.add_argument("--mcp-root", type=Path, default=DEFAULT_MCP_ROOT)
parser.add_argument("--catalog", type=Path, default=DEFAULT_CATALOG)
parser.add_argument("--output-dir", type=Path, default=DEFAULT_OUTPUT_DIR)
parser.add_argument("--python", default=sys.executable)
parser.add_argument("--sample-tools", type=int, default=300)
parser.add_argument("--seed", type=int, default=13)
parser.add_argument("--syntax-timeout", type=float, default=8)
parser.add_argument("--import-timeout", type=float, default=8)
parser.add_argument("--startup-grace", type=float, default=2)
parser.add_argument("--startup-timeout", type=float, default=1)
parser.add_argument("--help-timeout", type=float, default=4)
parser.add_argument("--tool-timeout", type=float, default=8)
parser.add_argument("--server-workers", type=int, default=16)
parser.add_argument("--tool-workers", type=int, default=8)
args = parser.parse_args()
args.output_dir.mkdir(parents=True, exist_ok=True)
server_dirs = sorted(p for p in args.mcp_root.glob("mcp_*") if p.is_dir())
server_rows: list[dict[str, Any]] = []
server_payloads = [
(
str(server_dir),
args.python,
args.syntax_timeout,
args.import_timeout,
args.startup_timeout,
args.startup_grace,
args.help_timeout,
)
for server_dir in server_dirs
]
with mp.Pool(processes=max(1, args.server_workers)) as pool:
for index, row in enumerate(pool.imap_unordered(audit_one_server, server_payloads), start=1):
server_rows.append(row)
if index % 100 == 0:
print(f"audited {index}/{len(server_dirs)} servers", flush=True)
server_rows.sort(key=lambda row: row.get("server", ""))
catalog = load_catalog(args.catalog)
sample = stratified_tool_sample(catalog, args.sample_tools, args.seed)
tool_rows: list[dict[str, Any]] = []
with ThreadPoolExecutor(max_workers=max(1, args.tool_workers)) as executor:
futures = [executor.submit(audit_tool_call, args.mcp_root, item, args.tool_timeout) for item in sample]
for index, future in enumerate(as_completed(futures), start=1):
tool_rows.append(future.result())
if index % 25 == 0:
print(f"audited {index}/{len(sample)} sampled tools", flush=True)
tool_rows.sort(key=lambda row: (row.get("category", ""), row.get("server", ""), row.get("tool", "")))
server_total = len(server_rows)
source_present_total = sum(1 for r in server_rows if r.get("has_source"))
tool_total_catalog = sum(len(s.get("tools", [])) for s in catalog)
summary_rows = [
{
"audit_layer": "Generated server entrypoint present",
"unit": "mcp_generated server directories",
"n_tested": server_total,
"n_success": source_present_total,
"success_rate": pct(source_present_total, server_total),
"notes": "server directory contains *_shim_server.py or *_server.py under app/",
},
{
"audit_layer": "Server syntax/compile",
"unit": "MCP servers",
"n_tested": source_present_total,
"n_success": sum(1 for r in server_rows if r.get("syntax_ok")),
"success_rate": pct(sum(1 for r in server_rows if r.get("syntax_ok")), source_present_total),
"notes": "python -m py_compile on generated shim/raw server",
},
{
"audit_layer": "Server import",
"unit": "MCP servers",
"n_tested": sum(1 for r in server_rows if r.get("syntax_ok")),
"n_success": sum(1 for r in server_rows if r.get("import_ok")),
"success_rate": pct(
sum(1 for r in server_rows if r.get("import_ok")),
sum(1 for r in server_rows if r.get("syntax_ok")),
),
"notes": "import generated server module in isolated process",
},
{
"audit_layer": "MCP startup",
"unit": "MCP servers",
"n_tested": sum(1 for r in server_rows if r.get("syntax_ok")),
"n_success": sum(1 for r in server_rows if r.get("startup_ok")),
"success_rate": pct(
sum(1 for r in server_rows if r.get("startup_ok")),
sum(1 for r in server_rows if r.get("syntax_ok")),
),
"notes": "stdio server starts and does not immediately traceback; waiting servers count as success",
},
{
"audit_layer": "Help/version probe",
"unit": "servers with inferred underlying CLI",
"n_tested": sum(1 for r in server_rows if r.get("help_version_probed")),
"n_success": sum(1 for r in server_rows if r.get("help_version_ok")),
"success_rate": pct(
sum(1 for r in server_rows if r.get("help_version_ok")),
sum(1 for r in server_rows if r.get("help_version_probed")),
),
"notes": "run inferred command with --help/-h/--version when command literal is found",
},
{
"audit_layer": "Tiny fixture execution",
"unit": "stratified sampled tools",
"n_tested": len(tool_rows),
"n_success": sum(1 for r in tool_rows if r.get("ok")),
"success_rate": pct(sum(1 for r in tool_rows if r.get("ok")), len(tool_rows)),
"notes": f"sampled from {tool_total_catalog} catalog tools; required args filled with tiny FASTA/FASTQ/CSV/VCF fixtures",
},
]
server_error_counts = Counter()
for row in server_rows:
if not row.get("has_source"):
server_error_counts["missing_server_entrypoint"] += 1
for key in ("syntax_error_type", "import_error_type", "startup_error_type", "help_version_error_type"):
value = row.get(key)
if value and value != "none" and value != "no_command_literal":
server_error_counts[value] += 1
tool_error_counts = Counter(r.get("error_type", "none") for r in tool_rows if not r.get("ok"))
error_rows = []
for name, count in server_error_counts.most_common():
error_rows.append({"scope": "server", "error_type": name, "count": count})
for name, count in tool_error_counts.most_common():
error_rows.append({"scope": "sampled_tool", "error_type": name, "count": count})
write_csv(args.output_dir / "server_audit.csv", server_rows)
write_csv(args.output_dir / "tool_fixture_audit.csv", tool_rows)
write_csv(args.output_dir / "summary_table.csv", summary_rows)
write_csv(args.output_dir / "error_taxonomy.csv", error_rows)
(args.output_dir / "summary.json").write_text(
json.dumps(
{
"mcp_root": str(args.mcp_root),
"catalog": str(args.catalog),
"server_count": server_total,
"catalog_tool_count": tool_total_catalog,
"sampled_tool_count": len(tool_rows),
"summary_table": summary_rows,
"error_taxonomy": error_rows,
},
ensure_ascii=False,
indent=2,
),
encoding="utf-8",
)
print(json.dumps(summary_rows, indent=2), flush=True)
print(f"wrote audit results to {args.output_dir}", flush=True)
if __name__ == "__main__":
main()