czty's picture
Add files using upload-large-folder tool
a9e46a4 verified
Raw
History Blame Contribute Delete
6.18 kB
from __future__ import annotations
import os
from pathlib import Path
from typing import Any
from .engines import BioinfoE1Validator, BioinfoT1Consultant, BioinfoV1Executor
from .shared_memory import SharedKnowledgeSpace
class DualModeAgentSystem:
"""High-level orchestrator for V1 execution and T1 reflection."""
def __init__(self, project_root: str | Path, execution_backend: str | None = None):
self.project_root = Path(project_root)
backend = execution_backend or os.getenv("BIOCLAW_EXECUTION_BACKEND", "docker")
self.memory = SharedKnowledgeSpace(self.project_root / "shared_knowledge")
self.v1 = BioinfoV1Executor(
memory=self.memory,
default_results_root=self.project_root / "results" / "dual_mode_runs",
project_root=self.project_root,
execution_backend=backend,
)
self.t1 = BioinfoT1Consultant(memory=self.memory)
biomni_root = os.getenv(
"BIOCLAW_BIOMNI_ROOT",
"/225040511/project/BioScientist/agent_system/engines/v1_executor_backup",
)
mcp_config = os.getenv(
"BIOCLAW_BIOMNI_MCP_CONFIG",
str(Path(biomni_root) / "mcp_config_bioscientist_generated.yaml"),
)
self.e1 = BioinfoE1Validator(
memory=self.memory,
biomni_root=biomni_root,
mcp_config_path=mcp_config,
project_root=self.project_root,
)
def execute(
self,
task: str,
input_manifest: dict[str, Any],
task_scope: str,
) -> dict[str, Any]:
cfg = self.memory.latest_config_for_task(task_scope)
return self.v1.execute_task(
task=task,
input_manifest=input_manifest,
task_scope=task_scope,
pipeline_config=cfg,
)
def reflect(self) -> list[dict[str, Any]]:
return self.t1.review_reports()
def consult(self, user_goal: str, task_scope: str) -> dict[str, Any]:
return self.t1.consult(user_goal=user_goal, task_scope=task_scope)
def propose_config(
self,
task_scope: str,
strategy_name: str,
tools: list[str],
parameters: dict[str, Any],
rationale: str,
) -> dict[str, Any]:
return self.t1.emit_pipeline_config(
task_scope=task_scope,
strategy_name=strategy_name,
tools=tools,
parameters=parameters,
rationale=rationale,
)
def autopilot(
self,
user_goal: str,
data_dir: str,
task_scope: str = "first_pipeline",
manifest_overrides: dict[str, Any] | None = None,
) -> dict[str, Any]:
cfg = self.memory.latest_config_for_task(task_scope)
return self.v1.execute_autopilot(
user_goal=user_goal,
data_dir=data_dir,
task_scope=task_scope,
pipeline_config=cfg,
manifest_overrides=manifest_overrides,
)
def register_mcp_servers(self, dry_run: bool = False) -> dict[str, Any]:
return self.e1.ensure_mcp_registered(dry_run=dry_run)
@staticmethod
def _stage_log(stage: str, message: str) -> None:
print(f"[STAGE:{stage}] {message}", flush=True)
def propose_hypotheses(
self,
user_query: str,
task_scope: str,
n: int = 10,
top_k: int = 5,
) -> dict[str, Any]:
self._stage_log("HYPOTHESIS_GENERATION", f"start domain={task_scope} n={n} top_k={top_k}")
generated = self.t1.generate_hypotheses(
user_query=user_query,
domain=task_scope,
n=n,
top_k=top_k,
)
self._stage_log("HYPOTHESIS_GENERATION", f"generated={len(generated)}")
self._stage_log("HYPOTHESIS_RANKING", f"start candidates={len(generated)}")
ranked = self.t1.rank_hypotheses_by_success_proxy(generated, domain=task_scope)
self._stage_log("HYPOTHESIS_RANKING", f"done ranked={len(ranked)}")
return {
"task_scope": task_scope,
"user_query": user_query,
"generated_count": len(generated),
"hypothesis_generation": dict(getattr(self.t1, "last_generation_meta", {})),
"hypotheses_ranked": ranked,
}
def hypothesis_loop(
self,
user_query: str,
task_scope: str,
n: int = 10,
top_k: int = 5,
validate_top_m: int = 3,
validation_level: str = "L1",
register_mcp: bool = True,
) -> dict[str, Any]:
self._stage_log("PIPELINE", f"hypothesis-loop start domain={task_scope} validation_level={validation_level}")
registration = None
if register_mcp:
self._stage_log("MCP_REGISTRATION", "start")
registration = self.e1.ensure_mcp_registered(dry_run=False)
self._stage_log(
"MCP_REGISTRATION",
f"done ok={bool(registration and registration.get('ok', False))}",
)
proposal = self.propose_hypotheses(user_query=user_query, task_scope=task_scope, n=n, top_k=top_k)
ranked = proposal["hypotheses_ranked"]
self._stage_log(
"VALIDATION",
f"start level={validation_level.upper()} top_m={max(0, validate_top_m)} from_ranked={len(ranked)}",
)
validations = self.e1.validate_top_hypotheses(
ranked_hypotheses=ranked,
top_m=validate_top_m,
level=validation_level,
)
self._stage_log("VALIDATION", f"done reports={len(validations)}")
self._stage_log("PIPELINE", "hypothesis-loop finished")
return {
"task_scope": task_scope,
"user_query": user_query,
"registration": registration,
"generated_count": proposal["generated_count"],
"hypothesis_generation": proposal.get("hypothesis_generation", {}),
"validated_count": len(validations),
"ranked_hypotheses": ranked,
"validation_reports": validations,
"summary_statistics": self.memory.get_summary_statistics(),
}