from __future__ import annotations import importlib.util from pathlib import Path from typing import Any import json import os import re import shlex import shutil import subprocess import traceback import requests from ..mcp_router import MCPToolRouter from ..schemas import ExperimentReport, FailureMode from ..shared_memory import SharedKnowledgeSpace class BioinfoV1Executor: """Execution engine: plan task -> call MCP servers -> publish structured report.""" def __init__( self, memory: SharedKnowledgeSpace, default_results_root: str | Path, project_root: str | Path, execution_backend: str = "docker", ): self.memory = memory self.default_results_root = Path(default_results_root) self.default_results_root.mkdir(parents=True, exist_ok=True) self.project_root = Path(project_root) self.router = MCPToolRouter(project_root=project_root) self.execution_backend = execution_backend @staticmethod def _append_log(log_file: Path, message: str) -> None: with log_file.open("a", encoding="utf-8") as lf: lf.write(message.rstrip() + "\n") @staticmethod def _jsonable(value: Any) -> Any: if isinstance(value, Path): return str(value) if isinstance(value, dict): return {str(k): BioinfoV1Executor._jsonable(v) for k, v in value.items()} if isinstance(value, list): return [BioinfoV1Executor._jsonable(v) for v in value] if isinstance(value, tuple): return [BioinfoV1Executor._jsonable(v) for v in value] return value @staticmethod def _read_log_tail(log_file: Path, max_lines: int = 200) -> str: try: lines = log_file.read_text(encoding="utf-8", errors="replace").splitlines() return "\n".join(lines[-max_lines:]) except Exception: return "" @staticmethod def _pick_trimmed_pair(trim_dir: Path) -> tuple[Path | None, Path | None]: r1_patterns = [ "*_val_1.fq.gz", "*_val_1.fq", "*_R1_val_1.fq.gz", "*_1_val_1.fq.gz", "*_fastp_R1.fastq", "*_fastp_R1.fastq.gz", ] r2_patterns = [ "*_val_2.fq.gz", "*_val_2.fq", "*_R2_val_2.fq.gz", "*_2_val_2.fq.gz", "*_fastp_R2.fastq", "*_fastp_R2.fastq.gz", ] r1_candidates: list[Path] = [] r2_candidates: list[Path] = [] for pattern in r1_patterns: r1_candidates.extend(sorted(trim_dir.glob(pattern))) for pattern in r2_patterns: r2_candidates.extend(sorted(trim_dir.glob(pattern))) return (r1_candidates[-1] if r1_candidates else None, r2_candidates[-1] if r2_candidates else None) def _plan_task( self, task: str, task_scope: str, pipeline_config: dict[str, Any] | None, ) -> list[dict[str, Any]]: """ Plan execution in a Biomni-like style: 1) understand scope 2) choose candidate tools 3) execute step-by-step with artifacts passed forward """ cfg_tools = (pipeline_config or {}).get("tools") or [] if task_scope == "first_pipeline": task_lower = task.lower() wants_alignment = any( kw in task_lower for kw in ("align", "alignment", "map", "mapping", "variant", "snp", "bam") ) benchmark_repro_mode = any( kw in task_lower for kw in ("code repository", "original open-source code", "reproduce paper", "reproducibility") ) base_plan = [ { "name": "qc_raw", "description": "Quality control on raw reads", "candidates": cfg_tools or ["fastqc", "fastp"], }, { "name": "trim", "description": "Adapter/quality trimming", "candidates": cfg_tools or ["trim_galore", "cutadapt", "trimmomatic", "fastp"], }, { "name": "align", "description": "Read alignment to reference index", "candidates": cfg_tools or ["bowtie2", "bwa", "hisat2", "star", "minimap2"], }, { "name": "qc_trimmed", "description": "Quality control after trimming", "candidates": cfg_tools or ["fastqc", "qualimap"], }, { "name": "aggregate", "description": "Aggregate reports", "candidates": cfg_tools or ["multiqc"], }, ] if not wants_alignment: base_plan = [s for s in base_plan if s["name"] != "align"] if benchmark_repro_mode: base_plan.extend( [ { "name": "source_clone", "description": "Clone original paper source code repository", "local_executor": "source_clone", "candidates": [], }, { "name": "source_execute", "description": "Run reproducibility-oriented source code execution attempts", "local_executor": "source_execute", "candidates": [], }, { "name": "summarize_repro", "description": "Summarize reproducibility with deviations and evidence", "local_executor": "summarize_repro", "candidates": [], }, ] ) return base_plan # generic mode: one tool per configured entry if cfg_tools: return [ { "name": f"step_{i+1}", "description": f"Configured execution step for {tool}", "candidates": [tool], } for i, tool in enumerate(cfg_tools) ] # fallback return [ { "name": "generic_step", "description": f"Generic execution for task: {task}", "candidates": ["fastqc"], } ] def _route_or_generate(self, tool_name: str, log_file: Path): if self.router.has_tool(tool_name): return self.router.resolve(tool_name) self._append_log(log_file, f"[planner] MCP tool missing: {tool_name}, invoking converter...") self._try_generate_mcp_server(tool_name, log_file) self.router.refresh() if self.router.has_tool(tool_name): return self.router.resolve(tool_name) return None def _try_generate_mcp_server(self, tool_name: str, log_file: Path) -> None: converter_path = self.project_root / "src" / "bioinfomcp_converter.py" if not converter_path.exists(): self._append_log(log_file, f"[converter] converter not found: {converter_path}") return try: spec = importlib.util.spec_from_file_location("bioinfomcp_converter", str(converter_path)) if spec is None or spec.loader is None: raise RuntimeError("Failed to load converter module spec.") module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) converter_cls = getattr(module, "BioinfoMCP", None) if converter_cls is None: raise RuntimeError("BioinfoMCP class not found in converter.") converter = converter_cls(model="openai") ok, error_msg, code = converter.autogenerate_mcp_tool( tool_name=tool_name, manual="--help", run_help_command=True, ) if not ok or not code: raise RuntimeError(f"Converter failed: {error_msg}") mcp_dir = self.project_root / "mcp-servers" / f"mcp_{tool_name}" / "app" mcp_dir.mkdir(parents=True, exist_ok=True) server_file = mcp_dir / f"{tool_name}_server.py" server_file.write_text(self._wrap_generated_mcp_code(code), encoding="utf-8") self._append_log(log_file, f"[converter] generated MCP server: {server_file}") except Exception as exc: # pragma: no cover - defensive path self._append_log(log_file, f"[converter] generation failed for {tool_name}: {exc}") @staticmethod def _wrap_generated_mcp_code(code: str) -> str: if "FastMCP" in code and "mcp = FastMCP()" in code: if "if __name__ == '__main__':" in code or "if __name__ == \"__main__\":" in code: return code return f"{code}\n\nif __name__ == '__main__':\n mcp.run()\n" return ( "from fastmcp import FastMCP\n" "mcp = FastMCP()\n\n" f"{code}\n\n" "if __name__ == '__main__':\n" " mcp.run()\n" ) @staticmethod def _load_tool_callable(server_script: Path, function_name: str): mod_name = f"mcp_module_{server_script.stem}_{abs(hash(str(server_script)))}" spec = importlib.util.spec_from_file_location(mod_name, str(server_script)) if spec is None or spec.loader is None: raise RuntimeError(f"Cannot load module for {server_script}") module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) func = getattr(module, function_name, None) if func is None: raise RuntimeError(f"Function '{function_name}' not found in {server_script}") return func def _invoke_mcp_tool( self, route, kwargs: dict[str, Any], log_file: Path, ) -> dict[str, Any]: self._append_log( log_file, f"[toolcall] {route.tool_name}.{route.function_name} script={route.server_script} kwargs={json.dumps({k: str(v) for k, v in kwargs.items()}, ensure_ascii=True)}", ) if self.execution_backend == "docker": return self._invoke_mcp_tool_docker(route, kwargs=kwargs, log_file=log_file) try: tool_fn = self._load_tool_callable(route.server_script, route.function_name) result = tool_fn(**kwargs) if not isinstance(result, dict): result = {"raw_result": result} ok = "error" not in result self._append_log(log_file, f"[toolcall] status={'ok' if ok else 'error'}") return {"ok": ok, "result": result} except FileNotFoundError as exc: tb = traceback.format_exc() self._append_log(log_file, f"[toolcall] FileNotFoundError: {exc}\n{tb}") return {"ok": False, "result": {"error": f"FileNotFoundError: {exc}", "traceback": tb}} except Exception as exc: tb = traceback.format_exc() self._append_log(log_file, f"[toolcall] exception: {exc}\n{tb}") return {"ok": False, "result": {"error": str(exc), "traceback": tb}} def _invoke_mcp_tool_docker(self, route, kwargs: dict[str, Any], log_file: Path) -> dict[str, Any]: if shutil.which("docker") is None: return { "ok": False, "result": {"error": "docker_not_found", "hint": "Install Docker or switch execution_backend=python"}, } run_dir = log_file.parent payload_file = run_dir / f"_mcp_payload_{route.function_name}.json" runner_file = run_dir / "_mcp_docker_runner.py" payload = { "server_script": str(route.server_script), "function_name": route.function_name, "kwargs": self._jsonable(kwargs), } payload_file.write_text(json.dumps(payload, ensure_ascii=True), encoding="utf-8") runner_file.write_text(self._docker_runner_script(), encoding="utf-8") mount_root = self.project_root.parent docker_cmd = [ "docker", "run", "--rm", "-v", f"{mount_root}:{mount_root}", "-w", str(mount_root), route.image_name, "python", str(runner_file), str(payload_file), ] self._append_log(log_file, f"[toolcall-docker] cmd={' '.join(docker_cmd)}") completed = subprocess.run(docker_cmd, capture_output=True, text=True) stdout = completed.stdout or "" stderr = completed.stderr or "" if completed.returncode != 0: self._append_log(log_file, f"[toolcall-docker] failed rc={completed.returncode}\n{stderr}") return { "ok": False, "result": { "error": f"docker_run_failed rc={completed.returncode}", "stdout": stdout, "stderr": stderr, "image": route.image_name, "hint": ( f"Ensure image '{route.image_name}' exists/builds, " f"or run with execution_backend=python." ), }, } try: result = json.loads(stdout.strip() or "{}") if not isinstance(result, dict): result = {"raw_result": result} except Exception: result = {"raw_stdout": stdout, "raw_stderr": stderr} ok = "error" not in result return {"ok": ok, "result": result} @staticmethod def _docker_runner_script() -> str: return """from __future__ import annotations import importlib.util import inspect import json import sys from pathlib import Path from typing import Any, get_args, get_origin, Union def to_jsonable(v: Any): if isinstance(v, Path): return str(v) if isinstance(v, dict): return {str(k): to_jsonable(val) for k, val in v.items()} if isinstance(v, list): return [to_jsonable(x) for x in v] if isinstance(v, tuple): return [to_jsonable(x) for x in v] return v def is_path_type(tp: Any) -> bool: return tp is Path def convert_value(value: Any, ann: Any) -> Any: if ann is inspect._empty: return value origin = get_origin(ann) if origin in (list, tuple): args = get_args(ann) inner = args[0] if args else Any if isinstance(value, list): return [convert_value(v, inner) for v in value] return value if origin is Union: for a in get_args(ann): if a is type(None): continue try: return convert_value(value, a) except Exception: pass return value if is_path_type(ann): return Path(value) if value is not None else value return value def main(): payload_path = Path(sys.argv[1]) payload = json.loads(payload_path.read_text(encoding="utf-8")) server_script = payload["server_script"] function_name = payload["function_name"] kwargs = payload.get("kwargs", {}) spec = importlib.util.spec_from_file_location("mcp_runtime_mod", server_script) if spec is None or spec.loader is None: print(json.dumps({"error": f"cannot_load_module:{server_script}"})) sys.exit(0) mod = importlib.util.module_from_spec(spec) spec.loader.exec_module(mod) fn = getattr(mod, function_name, None) if fn is None: print(json.dumps({"error": f"function_not_found:{function_name}"})) sys.exit(0) sig = inspect.signature(fn) call_kwargs = {} for k, v in kwargs.items(): if k in sig.parameters: call_kwargs[k] = convert_value(v, sig.parameters[k].annotation) else: call_kwargs[k] = v try: result = fn(**call_kwargs) print(json.dumps(to_jsonable(result), ensure_ascii=True)) except Exception as exc: import traceback print(json.dumps({"error": str(exc), "traceback": traceback.format_exc()}, ensure_ascii=True)) if __name__ == "__main__": main() """ def _build_step_kwargs( self, step_name: str, route_function: str, context: dict[str, Any], ) -> dict[str, Any]: run_dir: Path = context["run_dir"] r1: Path = context["r1"] r2: Path = context["r2"] threads: int = context["threads"] quality_cutoff: int = context["quality_cutoff"] index_base = context["index_base"] artifacts = context["artifacts"] if step_name == "qc_raw": outdir = run_dir / "01_fastqc_raw" outdir.mkdir(parents=True, exist_ok=True) return {"input_files": [r1, r2], "outdir": outdir, "threads": threads} if step_name == "trim": outdir = run_dir / "02_trim" outdir.mkdir(parents=True, exist_ok=True) if route_function == "fastp": out1 = outdir / f"{r1.stem}_fastp_R1.fastq" out2 = outdir / f"{r2.stem}_fastp_R2.fastq" json_report = str(outdir / "fastp.json") html_report = str(outdir / "fastp.html") return { "in1": r1, "in2": r2, "out1": out1, "out2": out2, "qualified_quality_phred": quality_cutoff, "cut_mean_quality": quality_cutoff, "cut_right": True, "thread": threads, "json": json_report, "html": html_report, "report_title": "fastp report (BioClawMCP)", } return { "input_files": [r1, r2], "paired": True, "quality": quality_cutoff, "output_dir": outdir, "cores": threads, } if step_name == "align": align_dir = run_dir / "03_align" align_dir.mkdir(parents=True, exist_ok=True) sam_out = align_dir / "alignment.sam" artifacts["sam"] = sam_out trimmed_r1 = artifacts.get("trimmed_r1", r1) trimmed_r2 = artifacts.get("trimmed_r2", r2) if route_function == "bowtie2_align": return { "index_base": str(index_base), "mate1_files": str(trimmed_r1), "mate2_files": str(trimmed_r2), "sam_output": sam_out, "threads": threads, } # generic aligner fall back shape return { "input_files": [trimmed_r1, trimmed_r2], "reference_index_base": str(index_base), "output_dir": align_dir, "threads": threads, } if step_name == "qc_trimmed": outdir = run_dir / "04_fastqc_trimmed" outdir.mkdir(parents=True, exist_ok=True) trimmed_r1 = artifacts.get("trimmed_r1", r1) trimmed_r2 = artifacts.get("trimmed_r2", r2) return {"input_files": [trimmed_r1, trimmed_r2], "outdir": outdir, "threads": threads} if step_name == "aggregate": outdir = run_dir / "05_multiqc" outdir.mkdir(parents=True, exist_ok=True) return { "analysis_directory": run_dir, "outdir": outdir, "filename": "multiqc_report.html", "force": True, } # generic fallback for configured custom step return { "analysis_directory": run_dir, "outdir": run_dir / f"{step_name}_output", } @staticmethod def _extract_code_repo(task: str, input_manifest: dict[str, Any]) -> str: manifest_repo = str(input_manifest.get("code_repo") or "").strip() if manifest_repo: return manifest_repo m = re.search(r"Code repo:\s*(https?://\S+)", task, re.I) if m: return m.group(1).rstrip(").,;") return "" def _run_local_command(self, cmd: list[str], log_file: Path, cwd: Path | None = None, timeout_s: int = 900) -> dict[str, Any]: try: self._append_log(log_file, f"[local-cmd] {' '.join(cmd)} cwd={cwd or self.project_root}") completed = subprocess.run( cmd, cwd=str(cwd) if cwd else None, capture_output=True, text=True, timeout=timeout_s, ) return { "command_executed": " ".join(cmd), "return_code": completed.returncode, "stdout": completed.stdout or "", "stderr": completed.stderr or "", } except Exception as exc: tb = traceback.format_exc() return { "command_executed": " ".join(cmd), "return_code": -1, "error": str(exc), "traceback": tb, } @staticmethod def _extract_readme_run_commands(source_dir: Path, max_commands: int = 6) -> list[list[str]]: candidates: list[list[str]] = [] readme_files = [ source_dir / "README.md", source_dir / "readme.md", source_dir / "README.rst", ] text = "" for fp in readme_files: if fp.exists(): try: text = fp.read_text(encoding="utf-8", errors="replace") break except Exception: continue if not text: return candidates # extract fenced code blocks first blocks = re.findall(r"```(?:bash|sh|shell)?\n([\s\S]*?)```", text, re.I) lines: list[str] = [] for b in blocks: lines.extend(b.splitlines()) if not lines: lines = text.splitlines() allowed_prefix = ( "python ", "python3 ", "bash ", "sh ", "./", "Rscript ", "R -e ", "R --vanilla ", "R ", "snakemake", "nextflow run", ) forbidden = ("sudo ", "rm -rf", "docker system", "shutdown", "reboot") for raw in lines: line = raw.strip() if line.startswith("$ "): line = line[2:].strip() if not line or line.startswith("#"): continue if any(bad in line for bad in forbidden): continue if line.startswith(allowed_prefix): try: cmd = shlex.split(line) except Exception: continue if cmd: candidates.append(cmd) if len(candidates) >= max_commands: break return candidates @staticmethod def _is_safe_run_command(cmd: list[str]) -> bool: if not cmd: return False joined = " ".join(cmd).lower() forbidden = [ "rm -rf", "sudo ", "shutdown", "reboot", ":(){:|:&};:", "mkfs", "dd if=", "curl ", "wget ", "scp ", "ssh ", "docker system", ] if any(x in joined for x in forbidden): return False allowed_heads = {"python", "python3", "bash", "sh", "rscript", "r", "snakemake", "nextflow", "make"} head = cmd[0].lower() return head in allowed_heads or head.startswith("./") def _plan_source_commands_with_gemini( self, *, task: str, source_dir: Path, input_manifest: dict[str, Any], log_file: Path, max_commands: int = 5, ) -> list[list[str]]: api_key = os.getenv("GEMINI_API_KEY", "").strip() print("api_key:"+api_key) if not api_key: self._append_log(log_file, "[gemini] planner skipped: GEMINI_API_KEY missing") return [] readme_text = "" for fp in [source_dir / "README.md", source_dir / "readme.md", source_dir / "README.rst"]: if fp.exists(): readme_text = fp.read_text(encoding="utf-8", errors="replace") break if len(readme_text) > 12000: readme_text = readme_text[:12000] print("readme_text:"+readme_text) endpoint = ( "https://generativelanguage.googleapis.com/v1beta/models/" f"gemini-2.5-flash-lite:generateContent?key={api_key}" ) planner_prompt = ( "You are a bioinformatics reproducibility execution planner.\n" "Given repository README and task context, propose up to 5 LOCAL runnable commands.\n" "Output JSON with key `commands`, where each command is either:\n" '1) ["python","script.py","--arg"] or 2) "python script.py --arg"\n' "Rules:\n" "- Only local run commands: python/python3/bash/sh/Rscript/R/snakemake/nextflow/make\n" "- No network download commands (curl/wget/git clone), no sudo, no destructive commands.\n" "- Prefer commands that execute original method in the repository.\n" "- Return JSON only.\n" ) payload = { "contents": [ { "role": "user", "parts": [ { "text": planner_prompt + "\n" + json.dumps( { "task": task, "repo_path": str(source_dir), "data_type": input_manifest.get("data_type", ""), "repro_steps": input_manifest.get("repro_steps", []), "readme_excerpt": readme_text, }, ensure_ascii=True, ) } ], } ], "generationConfig": {"temperature": 0.1}, } print("payload:"+json.dumps(payload, ensure_ascii=False)) try: resp = requests.post(endpoint, json=payload, timeout=60) if resp.status_code >= 400: self._append_log( log_file, f"[gemini] planner failed status={resp.status_code} body={(resp.text or '')[:500]}", ) return self._extract_readme_run_commands(source_dir, max_commands=max_commands) data = resp.json() text = "" for part in (data.get("candidates", [{}])[0].get("content", {}).get("parts", []) or []): if "text" in part: text += part["text"] obj = {} try: obj = json.loads(text) if text.strip() else {} except Exception: m = re.search(r"```json\s*([\s\S]*?)```", text, re.I) if m: obj = json.loads(m.group(1)) else: m = re.search(r"(\{[\s\S]*\})", text) if m: obj = json.loads(m.group(1)) commands = obj.get("commands", []) # Tolerate model output as a plain list print("obj:"+json.dumps(obj, ensure_ascii=False)) if not commands and isinstance(obj, list): commands = obj # Tolerate inline text commands if not commands and text.strip(): candidate_lines = [ln.strip() for ln in text.splitlines() if ln.strip()] commands = [ ln.lstrip("- ").strip() for ln in candidate_lines if ln.lower().startswith(("python ", "python3 ", "bash ", "sh ", "rscript ", "r ", "snakemake", "nextflow", "make ")) ] safe_cmds: list[list[str]] = [] for item in commands[:max_commands]: if isinstance(item, str): cmd = shlex.split(item) elif isinstance(item, list): cmd = [str(x) for x in item] else: continue if self._is_safe_run_command(cmd): safe_cmds.append(cmd) if not safe_cmds: self._append_log(log_file, f"[gemini] empty command plan; raw={(text or '')[:500]}") return self._extract_readme_run_commands(source_dir, max_commands=max_commands) self._append_log(log_file, f"[gemini] planned_commands={json.dumps(safe_cmds, ensure_ascii=True)}") return safe_cmds except Exception as exc: self._append_log(log_file, f"[gemini] planner exception: {exc}") return self._extract_readme_run_commands(source_dir, max_commands=max_commands) def _execute_local_step( self, step_name: str, *, task: str, input_manifest: dict[str, Any], context: dict[str, Any], log_file: Path, ) -> dict[str, Any]: run_dir: Path = context["run_dir"] artifacts: dict[str, Any] = context["artifacts"] code_repo = self._extract_code_repo(task=task, input_manifest=input_manifest) source_dir = run_dir / "06_source_code" artifacts["source_repo"] = code_repo artifacts["source_dir"] = source_dir if step_name == "source_clone": if not code_repo or "unknown_repo" in code_repo: return {"ok": False, "skipped": True, "result": {"warning": "source_repo_missing"}} source_dir.parent.mkdir(parents=True, exist_ok=True) if source_dir.exists(): result = self._run_local_command(["git", "-C", str(source_dir), "pull"], log_file=log_file) else: result = self._run_local_command(["git", "clone", "--depth", "1", code_repo, str(source_dir)], log_file=log_file) ok = result.get("return_code", 1) == 0 return {"ok": ok, "result": result} if step_name == "source_execute": if not source_dir.exists(): return {"ok": False, "skipped": True, "result": {"warning": "source_dir_not_ready"}} attempts: list[dict[str, Any]] = [] run_attempt_success = False has_r = shutil.which("R") is not None # Attempt 1: install dependencies if common file exists req = source_dir / "requirements.txt" if req.exists(): setup = self._run_local_command(["python", "-m", "pip", "install", "-r", str(req)], cwd=source_dir, log_file=log_file) setup["attempt_type"] = "setup" attempts.append(setup) r_renv = source_dir / "renv.lock" r_desc = source_dir / "DESCRIPTION" if r_renv.exists(): if has_r: setup = self._run_local_command( [ "R", "--vanilla", "-e", "if (!requireNamespace('renv', quietly=TRUE)) install.packages('renv', repos='https://cloud.r-project.org'); renv::restore(prompt=FALSE)", ], cwd=source_dir, log_file=log_file, timeout_s=1800, ) setup["attempt_type"] = "setup" attempts.append(setup) else: attempts.append( { "attempt_type": "setup", "return_code": -1, "warning": "R_not_found_for_renv_restore", "command_executed": "R --vanilla -e ", } ) elif r_desc.exists() and has_r: setup = self._run_local_command( [ "R", "--vanilla", "-e", "if (!requireNamespace('remotes', quietly=TRUE)) install.packages('remotes', repos='https://cloud.r-project.org'); remotes::install_local('.', upgrade='never')", ], cwd=source_dir, log_file=log_file, timeout_s=1800, ) setup["attempt_type"] = "setup" attempts.append(setup) # Attempt 2: common reproducibility entry scripts candidate_scripts = ["reproduce.sh", "run.sh", "scripts/reproduce.sh"] executed = False for rel in candidate_scripts: fp = source_dir / rel if fp.exists(): run = self._run_local_command(["bash", str(fp)], cwd=source_dir, log_file=log_file) run["attempt_type"] = "run" attempts.append(run) run_attempt_success = run_attempt_success or (run.get("return_code", 1) == 0) executed = True break # Attempt 3: Gemini agent plans and drives repo-specific runnable commands gemini_cmds: list[list[str]] = [] if not executed: gemini_cmds = self._plan_source_commands_with_gemini( task=task, source_dir=source_dir, input_manifest=input_manifest, log_file=log_file, ) print(gemini_cmds) for cmd in gemini_cmds: run = self._run_local_command(cmd, cwd=source_dir, log_file=log_file) run["attempt_type"] = "run_gemini" attempts.append(run) if run.get("return_code", 1) == 0: run_attempt_success = True executed = True break if run_attempt_success: return {"ok": True, "result": {"attempts": attempts, "gemini_planned_commands": gemini_cmds}} warning = "source_execution_not_reproduced" if not gemini_cmds and not executed: warning = "gemini_no_executable_command" return { "ok": False, "skipped": True, "result": { "warning": warning, "gemini_planned_commands": gemini_cmds, "attempts": attempts, }, } if step_name == "summarize_repro": summary = { "source_repo": code_repo or "unknown_repo", "data_sources": input_manifest.get("data_sources", []), "repro_steps": input_manifest.get("repro_steps", []), "notes": ( "Summary generated from benchmark prompt execution. " "Review source execution attempts and MCP preprocessing outputs for final conclusion." ), } return {"ok": True, "result": summary} return {"ok": False, "result": {"error": f"unsupported_local_step:{step_name}"}} @staticmethod def _bowtie2_index_exists(index_base: str | Path) -> bool: base = Path(index_base) bt2_suffixes = [".1.bt2", ".2.bt2", ".3.bt2", ".4.bt2", ".rev.1.bt2", ".rev.2.bt2"] bt2l_suffixes = [s.replace(".bt2", ".bt2l") for s in bt2_suffixes] return any((base.parent / (base.name + s)).exists() for s in bt2_suffixes) or any( (base.parent / (base.name + s)).exists() for s in bt2l_suffixes ) def execute_task( self, task: str, input_manifest: dict[str, Any], task_scope: str, pipeline_config: dict[str, Any] | None = None, ) -> dict[str, Any]: report = ExperimentReport.new( task=task, input_manifest=input_manifest, pipeline_config_id=(pipeline_config or {}).get("config_id"), ) run_dir = self.default_results_root / report.run_id run_dir.mkdir(parents=True, exist_ok=True) log_path = run_dir / "pipeline.log" log_path.write_text("", encoding="utf-8") plan = self._plan_task(task=task, task_scope=task_scope, pipeline_config=pipeline_config) route_snapshot: dict[str, Any] = {} self._append_log(log_path, f"[planner] generated_plan={json.dumps(plan, ensure_ascii=True)}") try: params = (pipeline_config or {}).get("parameters", {}) threads = int(params.get("threads", 4)) quality_cutoff = int(params.get("quality_cutoff", 20)) index_base = input_manifest.get("reference_index_base") or params.get("reference_index_base") r1 = Path(input_manifest["r1"]).resolve() r2 = Path(input_manifest["r2"]).resolve() if task_scope != "first_pipeline": self._append_log( log_path, "[planner] non-first_pipeline scope detected; running generic configured steps plan", ) if not r1.exists() or not r2.exists(): raise FileNotFoundError("Input FASTQ files not found for r1/r2.") index_ready = bool(index_base) and self._bowtie2_index_exists(index_base) if not index_ready: before = len(plan) plan = [s for s in plan if s.get("name") != "align"] self._append_log( log_path, "[planner] align_step_skipped reason=missing_or_invalid_reference_index", ) self._append_log(log_path, f"[planner] plan_rewritten from_steps={before} to_steps={len(plan)}") step_results: dict[str, Any] = {} artifacts: dict[str, Any] = {} warnings: list[str] = [] skippable_steps = {"aggregate"} context = { "run_dir": run_dir, "r1": r1, "r2": r2, "threads": threads, "quality_cutoff": quality_cutoff, "index_base": index_base, "artifacts": artifacts, } for step in plan: step_name = step["name"] if step.get("local_executor"): local_result = self._execute_local_step( step_name=step_name, task=task, input_manifest=input_manifest, context=context, log_file=log_path, ) step_results[step_name] = local_result if not local_result.get("ok"): warn = f"skip_or_fail_local_step:{step_name}" warnings.append(warn) self._append_log(log_path, f"[planner] {warn}") continue candidates = step["candidates"] chosen_route = None for candidate in candidates: chosen_route = self._route_or_generate(candidate, log_file=log_path) if chosen_route is not None: if self.router.binary_available(chosen_route.function_name): break self._append_log( log_path, f"[planner] skip_candidate={candidate} reason=binary_unavailable binary={chosen_route.binary_name}", ) chosen_route = None if chosen_route is None: if step_name in skippable_steps: warn = f"skip_step:{step_name} reason=no_available_mcp_route" warnings.append(warn) self._append_log(log_path, f"[planner] {warn}") step_results[step_name] = { "ok": False, "skipped": True, "result": {"warning": warn}, } continue raise RuntimeError(f"no_mcp_route_found_for_step:{step_name}") route_snapshot[step_name] = { "requested_candidates": candidates, "selected_tool": chosen_route.tool_name, "selected_function": chosen_route.function_name, "binary": chosen_route.binary_name, "binary_available": self.router.binary_available(chosen_route.function_name), "server_script": str(chosen_route.server_script), "server_exists": chosen_route.server_script.exists(), } kwargs = self._build_step_kwargs( step_name=step_name, route_function=chosen_route.function_name, context=context, ) result = self._invoke_mcp_tool(chosen_route, kwargs=kwargs, log_file=log_path) step_results[step_name] = result if not result["ok"]: if step_name in skippable_steps: warn = f"skip_step:{step_name} reason=tool_execution_failed" warnings.append(warn) self._append_log(log_path, f"[planner] {warn}") step_results[step_name] = { "ok": False, "skipped": True, "result": { **self._jsonable(result["result"]), "warning": warn, }, } continue raise RuntimeError(f"{step_name}_failed") # capture key artifacts output_files = result["result"].get("output_files") or [] if step_name == "trim": trim_dir = run_dir / "02_trim" # best-effort: prefer explicit output_files when available if output_files: r1_candidates = [Path(p) for p in output_files if ("_val_1" in p or "_fastp_R1" in p)] r2_candidates = [Path(p) for p in output_files if ("_val_2" in p or "_fastp_R2" in p)] if r1_candidates and r2_candidates: artifacts["trimmed_r1"] = r1_candidates[-1] artifacts["trimmed_r2"] = r2_candidates[-1] if "trimmed_r1" not in artifacts or "trimmed_r2" not in artifacts: trimmed_r1, trimmed_r2 = self._pick_trimmed_pair(trim_dir) if trimmed_r1 and trimmed_r2: artifacts["trimmed_r1"] = trimmed_r1 artifacts["trimmed_r2"] = trimmed_r2 if "trimmed_r1" not in artifacts or "trimmed_r2" not in artifacts: raise RuntimeError("trimmed_files_not_found") multiqc_report = run_dir / "05_multiqc" / "multiqc_report.html" if (step_results.get("aggregate") or {}).get("ok") and multiqc_report.exists(): artifacts["multiqc_report"] = multiqc_report self._collect_fastqc_artifacts(run_dir=run_dir, artifacts=artifacts) summary = { "status": "completed_with_warnings" if warnings else "completed", "task_scope": task_scope, "task": task, "plan": self._jsonable(plan), "steps": self._jsonable({k: {"ok": v["ok"], "result": v["result"]} for k, v in step_results.items()}), "warnings": warnings, "artifacts": { "trimmed_r1": str(artifacts.get("trimmed_r1", "")), "trimmed_r2": str(artifacts.get("trimmed_r2", "")), "sam": str(artifacts.get("sam", "")), "multiqc_report": str(artifacts.get("multiqc_report", multiqc_report)), "fastqc_raw_html": self._jsonable(artifacts.get("fastqc_raw_html", [])), "fastqc_trimmed_html": self._jsonable(artifacts.get("fastqc_trimmed_html", [])), }, "router": self._jsonable(route_snapshot), } (run_dir / "summary.json").write_text( json.dumps(summary, indent=2, ensure_ascii=True), encoding="utf-8", ) analysis_report_path = self._write_analysis_report( run_dir=run_dir, task=task, task_scope=task_scope, status=summary["status"], step_results=step_results, warnings=warnings, artifacts=artifacts, ) report.status = "completed_with_warnings" if warnings else "completed" report.execution_log = str(log_path) report.output_artifacts = [str(run_dir / "summary.json"), str(analysis_report_path)] if artifacts.get("multiqc_report"): report.output_artifacts.append(str(artifacts["multiqc_report"])) report.metrics = { "steps": len(plan), "has_pipeline_config": bool(pipeline_config), "task_scope": task_scope, "plan": self._jsonable(plan), "router": self._jsonable(route_snapshot), "step_results": self._jsonable(step_results), "warnings": warnings, "log_tail": self._read_log_tail(log_path), } if warnings: report.notes = ( "Executed with warnings via MCP routing plan. " f"Warnings: {'; '.join(warnings)}" ) else: report.notes = "Executed via MCP routing plan with dynamic tool selection and converter fallback." except Exception as exc: # pragma: no cover - defensive path report.status = "failed" report.execution_log = str(log_path) report.failures.append( FailureMode( step_name="executor", error_type=type(exc).__name__, error_message=str(exc), hint="Inspect pipeline.log and parameters.", ) ) self._append_log(log_path, f"[V1] status=failed error={exc}") failure_summary = { "status": "failed", "task_scope": task_scope, "task": task, "plan": self._jsonable(plan), "router": self._jsonable(route_snapshot), "partial_steps": self._jsonable(step_results), "artifacts": self._jsonable(artifacts), "error": {"type": type(exc).__name__, "message": str(exc)}, "log_tail": self._read_log_tail(log_path), } (run_dir / "summary.json").write_text( json.dumps(failure_summary, indent=2, ensure_ascii=True), encoding="utf-8", ) report.output_artifacts = [str(run_dir / "summary.json")] report.metrics = { "task_scope": task_scope, "router": self._jsonable(route_snapshot), "plan": self._jsonable(plan), "partial_steps": self._jsonable(step_results), "log_tail": self._read_log_tail(log_path), } self.memory.save_report(report) return report.to_dict() def execute_autopilot( self, user_goal: str, data_dir: str | Path, task_scope: str = "first_pipeline", pipeline_config: dict[str, Any] | None = None, manifest_overrides: dict[str, Any] | None = None, ) -> dict[str, Any]: r1, r2 = self._discover_paired_fastq(Path(data_dir)) manifest: dict[str, Any] = {"r1": str(r1), "r2": str(r2)} if manifest_overrides: manifest.update(manifest_overrides) # reference index is optional; execute_task 会根据有效性决定是否执行 align if pipeline_config and pipeline_config.get("parameters", {}).get("reference_index_base"): manifest["reference_index_base"] = pipeline_config["parameters"]["reference_index_base"] return self.execute_task( task=user_goal, input_manifest=manifest, task_scope=task_scope, pipeline_config=pipeline_config, ) @staticmethod def _discover_paired_fastq(data_dir: Path) -> tuple[Path, Path]: if not data_dir.exists(): raise FileNotFoundError(f"data_dir not found: {data_dir}") files = sorted( [ p for p in data_dir.iterdir() if p.is_file() and p.name.lower().endswith((".fastq", ".fq", ".fastq.gz", ".fq.gz")) ] ) if len(files) < 2: raise ValueError("Need at least two FASTQ files in data_dir.") # 优先找 R1/R2 命名对 for f in files: name = f.name r2_name = ( name.replace("_R1", "_R2") .replace(".R1.", ".R2.") .replace("_1.", "_2.") ) for g in files: if g.name == r2_name: return f, g # fallback: 前两个 FASTQ return files[0], files[1] @staticmethod def _collect_fastqc_artifacts(run_dir: Path, artifacts: dict[str, Any]) -> None: raw_dir = run_dir / "01_fastqc_raw" trim_dir = run_dir / "04_fastqc_trimmed" if raw_dir.exists(): artifacts["fastqc_raw_html"] = [str(p) for p in sorted(raw_dir.glob("*_fastqc.html"))] if trim_dir.exists(): artifacts["fastqc_trimmed_html"] = [str(p) for p in sorted(trim_dir.glob("*_fastqc.html"))] @staticmethod def _extract_fastp_metrics(step_results: dict[str, Any]) -> dict[str, Any]: trim = step_results.get("trim", {}) result = trim.get("result", {}) if isinstance(trim, dict) else {} output_files = result.get("output_files") or [] json_fp = None for fp in output_files: if str(fp).endswith(".json") and "fastp" in str(fp): json_fp = Path(fp) break if not json_fp or not json_fp.exists(): return {} try: payload = json.loads(json_fp.read_text(encoding="utf-8")) summary = payload.get("summary", {}) before = summary.get("before_filtering", {}) after = summary.get("after_filtering", {}) filtering = summary.get("filtering_result", {}) return { "before_total_reads": before.get("total_reads"), "after_total_reads": after.get("total_reads"), "q30_rate_after": after.get("q30_rate"), "passed_filter_reads": filtering.get("passed_filter_reads"), "low_quality_reads": filtering.get("low_quality_reads"), } except Exception: return {} def _write_analysis_report( self, run_dir: Path, task: str, task_scope: str, status: str, step_results: dict[str, Any], warnings: list[str], artifacts: dict[str, Any], ) -> Path: metrics = self._extract_fastp_metrics(step_results) lines: list[str] = [] lines.append("# BioClawMCP Analysis Report") lines.append("") lines.append(f"- Task: {task}") lines.append(f"- Task Scope: {task_scope}") lines.append(f"- Status: {status}") lines.append("") if warnings: lines.append("## Warnings") for w in warnings: lines.append(f"- {w}") lines.append("") lines.append("## Pipeline Execution Summary") lines.append("| Step | Status |") lines.append("|---|---|") for step_name, step_payload in step_results.items(): ok = step_payload.get("ok", False) skipped = step_payload.get("skipped", False) status_text = "skipped" if skipped else ("ok" if ok else "failed") lines.append(f"| {step_name} | {status_text} |") lines.append("") if metrics: lines.append("## fastp Key Metrics") lines.append(f"- Before total reads: {metrics.get('before_total_reads')}") lines.append(f"- After total reads: {metrics.get('after_total_reads')}") lines.append(f"- Q30 rate after: {metrics.get('q30_rate_after')}") lines.append(f"- Passed filter reads: {metrics.get('passed_filter_reads')}") lines.append(f"- Low quality reads: {metrics.get('low_quality_reads')}") lines.append("") lines.append("## Output Artifacts") if artifacts.get("trimmed_r1"): lines.append(f"- Trimmed R1: {artifacts.get('trimmed_r1')}") if artifacts.get("trimmed_r2"): lines.append(f"- Trimmed R2: {artifacts.get('trimmed_r2')}") if artifacts.get("sam"): lines.append(f"- Alignment SAM: {artifacts.get('sam')}") if artifacts.get("multiqc_report"): lines.append(f"- MultiQC report: {artifacts.get('multiqc_report')}") for fp in artifacts.get("fastqc_raw_html", []): lines.append(f"- FastQC raw HTML: {fp}") for fp in artifacts.get("fastqc_trimmed_html", []): lines.append(f"- FastQC trimmed HTML: {fp}") lines.append("") lines.append("## Per-step Command & Logs") for step_name, step_payload in step_results.items(): lines.append(f"### {step_name}") result = step_payload.get("result", {}) if isinstance(step_payload, dict) else {} cmd = result.get("command_executed", "") if cmd: lines.append("```bash") lines.append(cmd) lines.append("```") stderr = result.get("stderr", "") stdout = result.get("stdout", "") if stderr: lines.append("stderr:") lines.append("```text") lines.append(str(stderr)[:6000]) lines.append("```") if stdout: lines.append("stdout:") lines.append("```text") lines.append(str(stdout)[:6000]) lines.append("```") lines.append("") report_path = run_dir / "analysis_report.md" report_path.write_text("\n".join(lines), encoding="utf-8") return report_path