From 45cdbfde036d1eefe96c743ba0c2987362a46113 Mon Sep 17 00:00:00 2001 From: Saurabh Maurya Date: Tue, 18 Aug 2026 22:29:12 +0000 Subject: [PATCH 1/3] feat(agent): introduce reverse proxy bidirectional streaming evaluation Implements generalized reverse proxy bidirectional streaming architecture for autonomous multi-turn agent evaluation in Evalbench. Key Changes: - evalproto/eval_agent.proto: Defined typed AgentStreamMessage envelopes for turn requests/responses, in-flight remote scoring, remote workspace archival, and session summary delivery. - evalproto/eval_service.proto: Added rpc AgentInteract streaming endpoint. - evalbench/eval_service.py: Implemented AgentInteract servicer bridging streaming clients with in-memory AGENT_PROXY_QUEUES. - evalbench/generators/models/agentic_reverse_proxy.py: Added AgenticReverseProxyGenerator implementing AgentCliGenerator. - evalbench/scorers/remote_scorer.py: Added RemoteScorerProxy and transparent `remote: true` delegation in score.py. - evalbench/reporting/remote_artifact_reporter.py: Added RemoteArtifactReporter and transparent `reporting.gcs.remote: true` delegation. - evalbench/test/agentic_reverse_proxy_test.py: Added comprehensive unit tests. --- evalbench/eval_service.py | 156 +++++++++++++ evalbench/evalproto/eval_agent.proto | 142 ++++++++++++ evalbench/evalproto/eval_service.proto | 6 + evalbench/generators/models/__init__.py | 2 + .../models/agentic_reverse_proxy.py | 218 ++++++++++++++++++ evalbench/reporting/__init__.py | 14 +- .../reporting/remote_artifact_reporter.py | 68 ++++++ evalbench/scorers/remote_scorer.py | 86 +++++++ evalbench/scorers/score.py | 5 + evalbench/test/agentic_reverse_proxy_test.py | 157 +++++++++++++ 10 files changed, 851 insertions(+), 3 deletions(-) create mode 100644 evalbench/evalproto/eval_agent.proto create mode 100644 evalbench/generators/models/agentic_reverse_proxy.py create mode 100644 evalbench/reporting/remote_artifact_reporter.py create mode 100644 evalbench/scorers/remote_scorer.py create mode 100644 evalbench/test/agentic_reverse_proxy_test.py diff --git a/evalbench/eval_service.py b/evalbench/eval_service.py index bf315503..a7aad5ce 100644 --- a/evalbench/eval_service.py +++ b/evalbench/eval_service.py @@ -29,12 +29,15 @@ eval_request_pb2, eval_response_pb2, eval_service_pb2_grpc, + eval_agent_pb2, + eval_agent_pb2_grpc, ) from util.service import ( load_session_configs, get_dataset_from_request, ) from generators.models.grpc_proxy import PROXY_QUEUES +from generators.models.agentic_reverse_proxy import AGENT_PROXY_QUEUES import threading from util.context import rpc_id_var @@ -453,6 +456,159 @@ async def read_from_client(): PROXY_QUEUES.pop(session_id, None) logging.info(f"Cleaned up proxy queues for session {session_id}") + async def AgentInteract( + self, + request_iterator: AsyncIterator[eval_agent_pb2.AgentStreamMessage], + context: grpc.ServicerContext, + ) -> AsyncGenerator[eval_agent_pb2.AgentStreamMessage, None]: + """Bidirectional stream linking Google3 Autonomous Agents to Evalbench AgentEvaluator.""" + session_id = rpc_id_var.get() + session = SESSIONMANAGER.get_session(session_id) + config, db_configs, model_config, setup_config = load_session_configs(session) + + if config is None: + context.set_code(grpc.StatusCode.FAILED_PRECONDITION) + context.set_details("Session not configured") + return + + logging.info("Starting an AgentInteract bidirectional stream for session %s...", session_id) + config["session_id"] = session_id + + inboxes: dict[str, queue.Queue] = {} # correlation_id -> queue.Queue + out_queue = queue.Queue() # Evalbench -> worker + + config["agent_inboxes"] = inboxes + config["agent_out_queue"] = out_queue + AGENT_PROXY_QUEUES[session_id] = (inboxes, out_queue) + + # Load dataset and instantiate orchestrator + dataset_config_json = config.get("dataset_config") + dataset_dict = load_dataset_from_json(dataset_config_json, config) + + dataset = [] + for _, item_list in dataset_dict.items(): + dataset.extend(item_list) + + num_evals = config.get("num_evals_to_run") + if num_evals and int(num_evals) > 0: + dataset = dataset[:int(num_evals)] + + orchestrator = get_orchestrator(config, db_configs, setup_config, report_progress=True) + loop = asyncio.get_event_loop() + ctx = contextvars.copy_context() + + try: + def _cleanup_on_drop(ctx): + if session_id in AGENT_PROXY_QUEUES: + AGENT_PROXY_QUEUES.pop(session_id, None) + logging.info(f"Cleaned up agent proxy queues for session {session_id}") + + context.add_done_callback(_cleanup_on_drop) + + async def run_eval_and_process(): + await loop.run_in_executor(None, ctx.run, orchestrator.evaluate, dataset) + job_id, run_time, results_tf, scores_tf, multi_trial_scores_tf = orchestrator.process() + reporters = get_reporters(config.get("reporting") or {}, job_id, run_time) + logging.info("Processing agent evaluation results...") + summary = await loop.run_in_executor( + None, + ctx.run, + _process_results, + reporters, + job_id, + run_time, + results_tf, + scores_tf, + multi_trial_scores_tf, + config, + model_config, + db_configs, + ) + return job_id, summary + + eval_task = asyncio.create_task(run_eval_and_process()) + + async def read_from_client(): + async for response in request_iterator: + corr_id = response.correlation_id + logging.info( + "Server-Inbound: Received AgentStreamMessage (correlation_id=%s, payload=%s)", + corr_id, + response.WhichOneof("payload"), + ) + if corr_id in inboxes: + inboxes[corr_id].put(response) + else: + logging.warning( + "Server-Inbound: Orphaned AgentStreamMessage correlation_id '%s' (active inboxes: %s)", + corr_id, + list(inboxes.keys()), + ) + + read_task = asyncio.create_task(read_from_client()) + + # Yield loop: pop from out_queue and yield to worker + job_id = None + summary = None + while True: + if eval_task.done(): + logging.info("Agent Evaluator & Reporting task finished for session %s.", session_id) + try: + job_id, summary = eval_task.result() + except Exception as e: + logging.error("Agent Evaluator & Reporting task failed: %s", e, exc_info=True) + break + + if SESSIONMANAGER.get_session(session_id) is None: + logging.warning(f"Session {session_id} deleted. Terminating stream.") + context.set_code(grpc.StatusCode.NOT_FOUND) + context.set_details("Session deleted") + return + + try: + out_msg: eval_agent_pb2.AgentStreamMessage = await asyncio.to_thread(out_queue.get, True, 0.5) + logging.info( + "Server-Outbound: Yielding AgentStreamMessage (correlation_id=%s, payload=%s)", + out_msg.correlation_id, + out_msg.WhichOneof("payload"), + ) + yield out_msg + except queue.Empty: + continue + except Exception as e: + logging.error("Server-Outbound: Error yielding message: %s", e, exc_info=True) + continue + + # Flush any remaining messages from out_queue before finishing + while not out_queue.empty(): + try: + out_msg = out_queue.get_nowait() + yield out_msg + except Exception: + break + + read_task.cancel() + try: + await read_task + except asyncio.CancelledError: + pass + + if job_id and summary: + logging.info(f"Finished Agent Evaluation Job ID {job_id}. Summary: {summary}") + final_msg = eval_agent_pb2.AgentStreamMessage( + session_id=session_id, + correlation_id="final_summary", + session_summary=eval_agent_pb2.SessionSummaryMessage( + job_id=job_id, + summary_json=json.dumps(summary), + ), + ) + yield final_msg + + finally: + AGENT_PROXY_QUEUES.pop(session_id, None) + logging.info(f"Cleaned up agent proxy queues for session {session_id}") + def _process_results( reporters, job_id, run_time, results_tf, scores_tf, multi_trial_scores_tf, config, model_config, db_configs diff --git a/evalbench/evalproto/eval_agent.proto b/evalbench/evalproto/eval_agent.proto new file mode 100644 index 00000000..d877b58a --- /dev/null +++ b/evalbench/evalproto/eval_agent.proto @@ -0,0 +1,142 @@ +edition = "2023"; + +package cloud_databases_eval_proto; + +message AgentStreamMessage { + string session_id = 1; + string scenario_id = 2; + string correlation_id = 3; // Correlates async requests and responses across threads + + oneof payload { + // 0. Pre-Flight Health & Readiness Probe + HealthCheckRequest health_check_request = 5; + HealthCheckResponse health_check_response = 6; + + // 1. Scenario Lifecycle + LifecycleRequest lifecycle_request = 10; + LifecycleResponse lifecycle_response = 11; + + // 2. Turn Execution + TurnRequest turn_request = 20; + TurnResponse turn_response = 21; + + // 3. In-Flight Scoring Execution (Named Remote Scorers) + ScoringRequest scoring_request = 30; + ScoringResponse scoring_response = 31; + + // 4. Remote Workspace Archival + ArtifactRequest artifact_request = 40; + ArtifactResponse artifact_response = 41; + + // 5. Final Session Completion + SessionSummaryMessage session_summary = 50; + } +} + +// --- 0. Health Check Messages --- +message HealthCheckRequest { + string probe_command = 1; + float timeout_seconds = 2; +} + +message HealthCheckResponse { + bool is_healthy = 1; + bool internet_egress_ok = 2; + int32 active_skills_count = 3; + repeated string available_tools = 4; + string diagnostic_output = 5; + string error_message = 6; +} + +// --- 1. Lifecycle Messages --- +message LifecycleRequest { + enum Stage { STAGE_UNSPECIFIED = 0; SETUP = 1; RESET = 2; TEARDOWN = 3; } + Stage stage = 1; + string command = 2; + map env_vars = 3; +} + +message LifecycleResponse { + int32 exit_code = 1; + string stdout = 2; + string stderr = 3; + string error_message = 4; +} + +// --- 2. Turn Execution Messages --- +message TurnRequest { + int32 turn_index = 1; + string prompt = 2; + string system_instruction = 3; + map env = 4; + string working_dir = 5; + float timeout_seconds = 6; + bool resume = 7; +} + +message ToolCallRecord { + string tool_id = 1; + string tool_name = 2; + string parameters_json = 3; + string output = 4; + string status = 5; // "success", "error" + int64 duration_ms = 6; +} + +message TurnResponse { + int32 turn_index = 1; + string response_text = 2; + string stdout = 3; + string stderr = 4; + int32 exit_code = 5; + repeated ToolCallRecord tool_calls = 6; + map token_stats = 7; + bool execution_completed = 8; + string error_message = 9; +} + +// --- 3. In-Flight Scoring Messages --- +message RemoteScorerSpec { + string name = 1; // Scorer identifier, e.g. "dataform_compile", "dbt_run", "notebook_eval" + string config_json = 2; // Serialized configuration dictionary from YAML + float timeout_seconds = 3; // Timeout per scorer +} + +message ScoringRequest { + repeated RemoteScorerSpec scorers = 1; +} + +message ScoreResult { + string name = 1; + float score = 2; // 100.0 (PASS), 0.0 (FAIL), or continuous score + int32 exit_code = 3; + string stdout = 4; + string stderr = 5; + string logs = 6; + string error_message = 7; +} + +message ScoringResponse { + repeated ScoreResult results = 1; +} + +// --- 4. Artifact Archival Messages --- +message ArtifactRequest { + string target_gcs_bucket = 1; + string target_gcs_prefix = 2; + string export_path = 3; + repeated string exclude_patterns = 4; +} + +message ArtifactResponse { + string gcs_uri = 1; + int64 archive_size_bytes = 2; + repeated string exported_files = 3; + string error_message = 4; +} + +// --- 5. Session Summary --- +message SessionSummaryMessage { + string job_id = 1; + string summary_json = 2; +} diff --git a/evalbench/evalproto/eval_service.proto b/evalbench/evalproto/eval_service.proto index 490ce6d0..4e95a48e 100644 --- a/evalbench/evalproto/eval_service.proto +++ b/evalbench/evalproto/eval_service.proto @@ -6,6 +6,7 @@ import "eval_config.proto"; import "eval_connect.proto"; import "eval_request.proto"; import "eval_response.proto"; +import "eval_agent.proto"; option java_multiple_files = true; @@ -40,6 +41,11 @@ service EvalService { // option deadline = 1800; } + // Autonomous Agent Bidirectional stream. + rpc AgentInteract(stream AgentStreamMessage) returns (stream AgentStreamMessage) { + // option deadline = 1800; + } + // PrepareCodeEvalInputs for NL2Code Evaluation rpc PrepareCodeEvalInputs(EvalCodeInputRequest) returns (stream EvalCodeInputRequest) { // option deadline = 1800; diff --git a/evalbench/generators/models/__init__.py b/evalbench/generators/models/__init__.py index 15b9841d..c4dd0b1e 100644 --- a/evalbench/generators/models/__init__.py +++ b/evalbench/generators/models/__init__.py @@ -15,6 +15,7 @@ from .mcp_tools import McpToolsGenerator from .noop_agent import NoopAgentGenerator from .agent_runtime import AgentRuntimeGenerator +from .agentic_reverse_proxy import AgenticReverseProxyGenerator from util.config import load_yaml_config @@ -34,6 +35,7 @@ def get_generator(global_models, model_config_path: str, db: DB = None): "querydata": lambda: QueryData(config), "query_data_api": lambda: QueryDataAPIGenerator(config), "grpc_proxy": lambda: GrpcProxyModel(config), + "agentic_reverse_proxy": lambda: AgenticReverseProxyGenerator(config), "gemini_cli": lambda: GeminiCliGenerator(config), "claude_code": lambda: ClaudeCodeGenerator(config), "codex_cli": lambda: CodexCliGenerator(config), diff --git a/evalbench/generators/models/agentic_reverse_proxy.py b/evalbench/generators/models/agentic_reverse_proxy.py new file mode 100644 index 00000000..9c4dc8d7 --- /dev/null +++ b/evalbench/generators/models/agentic_reverse_proxy.py @@ -0,0 +1,218 @@ +import json +import logging +import queue +import subprocess +import time +import uuid +from typing import Any, Dict, List, Optional + +from .agent_cli import AgentCliGenerator +from evalproto import eval_agent_pb2 +from util.context import rpc_id_var + +# Global dictionary for active reverse proxy sessions: +# session_id -> (inboxes_dict, out_queue) +AGENT_PROXY_QUEUES: Dict[str, tuple[Dict[str, queue.Queue], queue.Queue]] = {} + + +class CLICommand: + def __init__(self, cli: str, prompt: str, env: dict = None, resume: bool = False, session_id: str = None, cwd: str = None): + self.cli = cli + self.prompt = prompt + self.env = env if env else {} + self.resume = resume + self.session_id = session_id + self.cwd = cwd + + +class AgenticReverseProxyGenerator(AgentCliGenerator): + """Generator proxying multi-turn agent execution across the reverse bidi stream.""" + + def __init__(self, querygenerator_config: Dict[str, Any]): + super().__init__(querygenerator_config) + self.name = "agentic_reverse_proxy" + self.timeout_seconds = float(querygenerator_config.get("timeout_seconds", 300.0)) + self.turn_counter: Dict[str, int] = {} + logging.info("Initialized AgenticReverseProxyGenerator (timeout=%ss)", self.timeout_seconds) + + @property + def version(self) -> str: + return "agentic_reverse_proxy" + + def generate_internal(self, prompt: str) -> str: + res = self.safe_generate(self.create_command(self.name, prompt)) + return res.stdout + + def create_command( + self, + cli: str, + prompt: str, + env: dict = None, + resume: bool = False, + session_id: str = None, + cwd: str = None, + ) -> CLICommand: + return CLICommand( + cli=cli or self.name, + prompt=prompt, + env=env, + resume=resume, + session_id=session_id, + cwd=cwd, + ) + + def safe_generate(self, cli_cmd: CLICommand | dict | str) -> subprocess.CompletedProcess: + if isinstance(cli_cmd, CLICommand): + prompt = cli_cmd.prompt + session_id = cli_cmd.session_id or rpc_id_var.get() + resume = cli_cmd.resume + env = cli_cmd.env + cwd = cli_cmd.cwd + elif isinstance(cli_cmd, dict): + prompt = cli_cmd.get("prompt", "") + session_id = cli_cmd.get("session_id") or rpc_id_var.get() + resume = cli_cmd.get("resume", False) + env = cli_cmd.get("env", {}) + cwd = cli_cmd.get("cwd") + else: + prompt = str(cli_cmd) + session_id = rpc_id_var.get() + resume = False + env = {} + cwd = None + + if session_id not in AGENT_PROXY_QUEUES: + ctx_id = rpc_id_var.get() + if ctx_id in AGENT_PROXY_QUEUES: + session_id = ctx_id + else: + logging.error("AgenticReverseProxy: session_id %s not in AGENT_PROXY_QUEUES (keys: %s)", session_id, list(AGENT_PROXY_QUEUES.keys())) + return subprocess.CompletedProcess( + args=["agentic_reverse_proxy"], + returncode=1, + stdout="", + stderr=f"Session {session_id} not connected to reverse stream", + ) + + inboxes, out_queue = AGENT_PROXY_QUEUES[session_id] + + turn_idx = self.turn_counter.get(session_id, 0) + 1 + self.turn_counter[session_id] = turn_idx + + correlation_id = str(uuid.uuid4()) + inbox: queue.Queue[eval_agent_pb2.AgentStreamMessage] = queue.Queue() + inboxes[correlation_id] = inbox + + turn_req = eval_agent_pb2.TurnRequest( + turn_index=turn_idx, + prompt=prompt, + env={k: str(v) for k, v in env.items()} if env else {}, + working_dir=cwd or "/workspace", + timeout_seconds=self.timeout_seconds, + resume=resume, + ) + + msg = eval_agent_pb2.AgentStreamMessage( + session_id=session_id, + correlation_id=correlation_id, + turn_request=turn_req, + ) + + logging.info("[REVERSE_PROXY] Dispatching TurnRequest turn=%d (correlation_id=%s) to out_queue", turn_idx, correlation_id) + out_queue.put(msg) + + try: + resp_msg = inbox.get(timeout=self.timeout_seconds) + except queue.Empty: + logging.error("[REVERSE_PROXY] Timed out waiting for TurnResponse (correlation_id=%s)", correlation_id) + return subprocess.CompletedProcess( + args=["agentic_reverse_proxy"], + returncode=124, + stdout="", + stderr="Timed out waiting for agent response from reverse stream", + ) + finally: + inboxes.pop(correlation_id, None) + + if not resp_msg.HasField("turn_response"): + err_details = resp_msg.WhichOneof("payload") + logging.error("[REVERSE_PROXY] Unexpected message received on inbox: %s", err_details) + return subprocess.CompletedProcess( + args=["agentic_reverse_proxy"], + returncode=1, + stdout="", + stderr=f"Unexpected payload on stream: {err_details}", + ) + + turn_resp = resp_msg.turn_response + tool_calls_list = [] + tools_by_name = {} + for tc in turn_resp.tool_calls: + t_entry = { + "tool_name": tc.tool_name, + "parameters": tc.parameters_json, + "output": tc.output, + "status": tc.status, + "duration_ms": tc.duration_ms, + } + tool_calls_list.append(t_entry) + if tc.tool_name not in tools_by_name: + tools_by_name[tc.tool_name] = {"parameters": []} + tools_by_name[tc.tool_name]["parameters"].append(tc.parameters_json) + + envelope = { + "session_id": session_id, + "response": turn_resp.response_text, + "stdout": turn_resp.stdout, + "stderr": turn_resp.stderr, + "exit_code": turn_resp.exit_code, + "tool_calls": tool_calls_list, + "stats": { + "tools": { + "totalCalls": len(tool_calls_list), + "byName": tools_by_name, + }, + "tokens": dict(turn_resp.token_stats), + }, + } + + raw_stdout = json.dumps(envelope, indent=2) + return subprocess.CompletedProcess( + args=["agentic_reverse_proxy"], + returncode=turn_resp.exit_code, + stdout=raw_stdout, + stderr=turn_resp.stderr, + ) + + def parse_response(self, stdout: str) -> dict: + if not stdout: + return {} + try: + return json.loads(stdout) + except Exception: + return {"response": stdout} + + def extract_tools(self, stdout: str) -> list[str]: + output_json = self.parse_response(stdout) + if "stats" in output_json and "tools" in output_json["stats"] and "byName" in output_json["stats"]["tools"]: + return list(output_json["stats"]["tools"]["byName"].keys()) + if "tool_calls" in output_json: + return [tc.get("tool_name") for tc in output_json["tool_calls"] if isinstance(tc, dict) and "tool_name" in tc] + return [] + + def extract_skills(self, stdout: str) -> list[str]: + output_json = self.parse_response(stdout) + skills = [] + for tc in output_json.get("tool_calls", []): + if isinstance(tc, dict) and tc.get("tool_name") in ("activate_skill", "Skill", "use_skill"): + params = tc.get("parameters", {}) + if isinstance(params, str): + try: + params = json.loads(params) + except Exception: + params = {} + if isinstance(params, dict): + sname = params.get("skill_name") or params.get("skill") or params.get("name") + if sname and sname not in skills: + skills.append(sname) + return skills diff --git a/evalbench/reporting/__init__.py b/evalbench/reporting/__init__.py index 3bb14b5a..cf7f2480 100644 --- a/evalbench/reporting/__init__.py +++ b/evalbench/reporting/__init__.py @@ -2,6 +2,7 @@ from .bqstore import BigQueryReporter from .report import Reporter from .gcs_artifact import GcsReporter +from .remote_artifact_reporter import RemoteArtifactReporter def get_reporters(reporting_config, job_id, run_time) -> list[Reporter]: @@ -15,7 +16,14 @@ def get_reporters(reporting_config, job_id, run_time) -> list[Reporter]: if "csv" in reporting_config: reporters.append(CsvReporter( reporting_config["csv"], job_id, run_time)) - if "gcs_artifacts" in reporting_config: - reporters.append(GcsReporter( - reporting_config["gcs_artifacts"], job_id, run_time)) + if "remote_artifacts" in reporting_config: + reporters.append(RemoteArtifactReporter( + reporting_config["remote_artifacts"], job_id, run_time)) + else: + gcs_cfg = reporting_config.get("gcs") or reporting_config.get("gcs_artifacts") or reporting_config.get("artifacts") + if gcs_cfg is not None: + if isinstance(gcs_cfg, dict) and gcs_cfg.get("remote", False): + reporters.append(RemoteArtifactReporter(gcs_cfg, job_id, run_time)) + else: + reporters.append(GcsReporter(gcs_cfg, job_id, run_time)) return reporters diff --git a/evalbench/reporting/remote_artifact_reporter.py b/evalbench/reporting/remote_artifact_reporter.py new file mode 100644 index 00000000..51578c6b --- /dev/null +++ b/evalbench/reporting/remote_artifact_reporter.py @@ -0,0 +1,68 @@ +import logging +import queue +import uuid +from typing import Any, Dict + +import pandas as pd + +from evalproto import eval_agent_pb2 +from reporting.report import Reporter, STORETYPE +from util.context import rpc_id_var +from generators.models.agentic_reverse_proxy import AGENT_PROXY_QUEUES + +logger = logging.getLogger(__name__) + + +class RemoteArtifactReporter(Reporter): + """Reporter delegating workspace dump and GCS upload across the reverse bidi stream.""" + + def __init__(self, reporting_config: Dict[str, Any] | None, job_id: str, run_time: Any): + super().__init__(reporting_config, job_id, run_time) + self.bucket = reporting_config.get("bucket", "sobi_dc_share") if reporting_config else "sobi_dc_share" + self.path_prefix = reporting_config.get("path_prefix", "runs") if reporting_config else "runs" + self.export_path = reporting_config.get("export_path", "/workspace") if reporting_config else "/workspace" + self.exclude_patterns = reporting_config.get("exclude_patterns", [".venv", "node_modules", "skills"]) if reporting_config else [".venv", "node_modules", "skills"] + logger.info("Initialized RemoteArtifactReporter: bucket=%s, prefix=%s", self.bucket, self.path_prefix) + + def store(self, results: pd.DataFrame, store_type: Any) -> None: + type_name = getattr(store_type, "name", str(store_type)) + if type_name != "EVALS": + return + + session_id = rpc_id_var.get() + if session_id not in AGENT_PROXY_QUEUES: + logger.warning("RemoteArtifactReporter: session_id %s not in AGENT_PROXY_QUEUES, skipping remote archival", session_id) + return + + inboxes, out_queue = AGENT_PROXY_QUEUES[session_id] + correlation_id = str(uuid.uuid4()) + inbox: queue.Queue[eval_agent_pb2.AgentStreamMessage] = queue.Queue() + inboxes[correlation_id] = inbox + + artifact_req = eval_agent_pb2.ArtifactRequest( + target_gcs_bucket=self.bucket, + target_gcs_prefix=self.path_prefix, + export_path=self.export_path, + exclude_patterns=self.exclude_patterns, + ) + + msg = eval_agent_pb2.AgentStreamMessage( + session_id=session_id, + correlation_id=correlation_id, + artifact_request=artifact_req, + ) + + logger.info("[REVERSE_REPORTER] Dispatching ArtifactRequest to out_queue (correlation_id=%s)", correlation_id) + out_queue.put(msg) + + try: + resp_msg = inbox.get(timeout=300.0) + if resp_msg.HasField("artifact_response"): + art_resp = resp_msg.artifact_response + logger.info("[REVERSE_REPORTER] Remote workspace archived successfully to: %s (%d bytes)", art_resp.gcs_uri, art_resp.archive_size_bytes) + else: + logger.warning("[REVERSE_REPORTER] Received non-artifact response: %s", resp_msg.WhichOneof("payload")) + except queue.Empty: + logger.error("[REVERSE_REPORTER] Timed out waiting for ArtifactResponse") + finally: + inboxes.pop(correlation_id, None) diff --git a/evalbench/scorers/remote_scorer.py b/evalbench/scorers/remote_scorer.py new file mode 100644 index 00000000..822a0a01 --- /dev/null +++ b/evalbench/scorers/remote_scorer.py @@ -0,0 +1,86 @@ +import json +import logging +import queue +import uuid +from typing import Any, Dict, List, Tuple + +from evalproto import eval_agent_pb2 +from scorers.comparator import Comparator +from util.context import rpc_id_var +from generators.models.agentic_reverse_proxy import AGENT_PROXY_QUEUES + + +class RemoteScorerProxy(Comparator): + """Comparator proxying scoring evaluation across the reverse bidi stream.""" + + def __init__(self, name: str, config: Dict[str, Any]): + super().__init__(config or {}) + self.name = name + self.config = dict(config) if isinstance(config, dict) else {} + self.timeout_seconds = float(self.config.get("timeout_seconds", 300.0)) + logging.info("Initialized RemoteScorerProxy: name=%s, config=%s", self.name, self.config) + + def compare( + self, + nl_prompt: str, + golden_sql: str, + query_type: str, + golden_result: str, + golden_eval_results: str, + golden_error: str, + generated_sql: str, + generated_result: str, + eval_results: str, + generated_error: str, + **kwargs: Any, + ) -> Tuple[float, str] | List[Tuple[str, float, str]]: + session_id = rpc_id_var.get() + if session_id not in AGENT_PROXY_QUEUES: + logging.error("RemoteScorerProxy: session_id %s not found in AGENT_PROXY_QUEUES", session_id) + return (0.0, f"Error: session_id '{session_id}' not connected to reverse stream") + + inboxes, out_queue = AGENT_PROXY_QUEUES[session_id] + correlation_id = str(uuid.uuid4()) + inbox: queue.Queue[eval_agent_pb2.AgentStreamMessage] = queue.Queue() + inboxes[correlation_id] = inbox + + scorer_spec = eval_agent_pb2.RemoteScorerSpec( + name=self.name, + config_json=json.dumps(self.config), + timeout_seconds=self.timeout_seconds, + ) + + scoring_req = eval_agent_pb2.ScoringRequest(scorers=[scorer_spec]) + msg = eval_agent_pb2.AgentStreamMessage( + session_id=session_id, + correlation_id=correlation_id, + scoring_request=scoring_req, + ) + + logging.info("[REVERSE_SCORER] Dispatching ScoringRequest for '%s' (correlation_id=%s)", self.name, correlation_id) + out_queue.put(msg) + + try: + resp_msg = inbox.get(timeout=self.timeout_seconds) + except queue.Empty: + logging.error("[REVERSE_SCORER] Timed out waiting for ScoringResponse for '%s' (correlation_id=%s)", self.name, correlation_id) + return (0.0, f"Error: Timed out waiting for remote scorer '{self.name}' response") + finally: + inboxes.pop(correlation_id, None) + + if not resp_msg.HasField("scoring_response"): + err_details = resp_msg.WhichOneof("payload") + logging.error("[REVERSE_SCORER] Unexpected message on stream: %s", err_details) + return (0.0, f"Error: Unexpected payload on stream: {err_details}") + + scoring_resp = resp_msg.scoring_response + results: List[Tuple[str, float, str]] = [] + for r in scoring_resp.results: + log_output = r.logs or r.stdout or r.error_message or f"exit_code={r.exit_code}" + results.append((r.name, float(r.score), log_output)) + + if len(results) == 1 and results[0][0] == self.name: + return (results[0][1], results[0][2]) + elif results: + return results + return (0.0, f"Error: Remote scorer '{self.name}' returned empty results") diff --git a/evalbench/scorers/score.py b/evalbench/scorers/score.py index 80aa7304..8dd4d6cb 100644 --- a/evalbench/scorers/score.py +++ b/evalbench/scorers/score.py @@ -32,6 +32,7 @@ from scorers import trajectorymatcher from scorers import turncount from scorers.dataset_quality.scorer import DatasetQualityScorer +from scorers.remote_scorer import RemoteScorerProxy DEFAULT_SCORERS: dict[str, type[comparator.Comparator]] = { @@ -140,6 +141,10 @@ def get_scorer_instance( if not isinstance(scorer_config, dict): return [] + # If marked remote, delegate across reverse bidi stream + if scorer_config.get("remote", False): + return [RemoteScorerProxy(scorer_name, scorer_config)] + # Direct match in global DEFAULT_SCORERS map if scorer_name in DEFAULT_SCORERS: return _build_instances( diff --git a/evalbench/test/agentic_reverse_proxy_test.py b/evalbench/test/agentic_reverse_proxy_test.py new file mode 100644 index 00000000..66c848fb --- /dev/null +++ b/evalbench/test/agentic_reverse_proxy_test.py @@ -0,0 +1,157 @@ +import queue +import unittest +from unittest.mock import MagicMock + +from evalproto import eval_agent_pb2 +from generators.models.agentic_reverse_proxy import ( + AgenticReverseProxyGenerator, + AGENT_PROXY_QUEUES, +) +from reporting.remote_artifact_reporter import RemoteArtifactReporter +from scorers.remote_scorer import RemoteScorerProxy +from util.context import rpc_id_var + + +class TestAgenticReverseProxy(unittest.TestCase): + + def setUp(self): + self.session_id = "test_session_123" + self.inboxes = {} + self.out_queue = queue.Queue() + AGENT_PROXY_QUEUES[self.session_id] = (self.inboxes, self.out_queue) + rpc_id_var.set(self.session_id) + + def tearDown(self): + AGENT_PROXY_QUEUES.pop(self.session_id, None) + + def test_generator_turn_success(self): + generator = AgenticReverseProxyGenerator({"timeout_seconds": 5.0}) + + def answer(): + msg = self.out_queue.get(timeout=2.0) + corr_id = msg.correlation_id + turn_req = msg.turn_request + self.assertEqual(turn_req.turn_index, 1) + self.assertEqual(turn_req.prompt, "hello agent") + + t1 = eval_agent_pb2.ToolCallRecord( + tool_id="call_1", + tool_name="list_directory", + parameters_json='{"path": "/workspace"}', + output="files", + status="success", + duration_ms=50, + ) + turn_resp = eval_agent_pb2.TurnResponse( + turn_index=1, + response_text="Found files", + stdout="done", + exit_code=0, + tool_calls=[t1], + token_stats={"input_tokens": 100}, + execution_completed=True, + ) + reply = eval_agent_pb2.AgentStreamMessage( + session_id=self.session_id, + correlation_id=corr_id, + turn_response=turn_resp, + ) + self.inboxes[corr_id].put(reply) + + import threading + t = threading.Thread(target=answer) + t.start() + + cmd = generator.create_command("agent", "hello agent", session_id=self.session_id) + res = generator.safe_generate(cmd) + t.join() + + self.assertEqual(res.returncode, 0) + parsed = generator.parse_response(res.stdout) + self.assertEqual(parsed["response"], "Found files") + self.assertEqual(len(parsed["tool_calls"]), 1) + self.assertEqual(parsed["tool_calls"][0]["tool_name"], "list_directory") + + def test_remote_scorer_success(self): + scorer = RemoteScorerProxy("dataform_compile", {"timeout_seconds": 5.0}) + + def answer_scorer(): + msg = self.out_queue.get(timeout=2.0) + corr_id = msg.correlation_id + self.assertEqual(msg.WhichOneof("payload"), "scoring_request") + spec = msg.scoring_request.scorers[0] + self.assertEqual(spec.name, "dataform_compile") + + score_res = eval_agent_pb2.ScoreResult( + name="dataform_compile", + score=100.0, + exit_code=0, + stdout="Compiled 1 action", + logs="Verified", + ) + reply = eval_agent_pb2.AgentStreamMessage( + session_id=self.session_id, + correlation_id=corr_id, + scoring_response=eval_agent_pb2.ScoringResponse(results=[score_res]), + ) + self.inboxes[corr_id].put(reply) + + import threading + t = threading.Thread(target=answer_scorer) + t.start() + + score, logs = scorer.compare( + nl_prompt="build pipeline", + golden_sql="", + query_type="", + golden_result="", + golden_eval_results="", + golden_error="", + generated_sql="", + generated_result="", + eval_results="", + generated_error="", + ) + t.join() + + self.assertEqual(score, 100.0) + self.assertEqual(logs, "Verified") + + def test_remote_artifact_reporter_success(self): + reporter = RemoteArtifactReporter( + {"bucket": "test-bucket", "path_prefix": "runs", "export_path": "/workspace"}, + job_id="job_123", + run_time="2026-08-18", + ) + + def answer_artifact(): + msg = self.out_queue.get(timeout=2.0) + corr_id = msg.correlation_id + self.assertEqual(msg.WhichOneof("payload"), "artifact_request") + art_req = msg.artifact_request + self.assertEqual(art_req.target_gcs_bucket, "test-bucket") + + reply = eval_agent_pb2.AgentStreamMessage( + session_id=self.session_id, + correlation_id=corr_id, + artifact_response=eval_agent_pb2.ArtifactResponse( + gcs_uri="gs://test-bucket/runs/fake_home.zip", + archive_size_bytes=2048, + exported_files=["notes.txt"], + ), + ) + self.inboxes[corr_id].put(reply) + + import threading + t = threading.Thread(target=answer_artifact) + t.start() + + import pandas as pd + from reporting.report import STORETYPE + df = pd.DataFrame({"eval_id": ["1"]}) + reporter.store(df, STORETYPE.EVALS) + t.join() + + +if __name__ == "__main__": + unittest.main() From ee40316e34e7e4f6601db8d5bc0aacd7aa3fbe50 Mon Sep 17 00:00:00 2001 From: Saurabh Maurya Date: Tue, 18 Aug 2026 23:04:23 +0000 Subject: [PATCH 2/3] chore: clean up unused imports and add comment to cancelled exception block --- evalbench/eval_service.py | 1 + evalbench/generators/models/agentic_reverse_proxy.py | 3 +-- evalbench/reporting/remote_artifact_reporter.py | 2 +- evalbench/test/agentic_reverse_proxy_test.py | 1 - 4 files changed, 3 insertions(+), 4 deletions(-) diff --git a/evalbench/eval_service.py b/evalbench/eval_service.py index a7aad5ce..bffb5260 100644 --- a/evalbench/eval_service.py +++ b/evalbench/eval_service.py @@ -591,6 +591,7 @@ async def read_from_client(): try: await read_task except asyncio.CancelledError: + # Expected when cleaning up inbound reader task upon eval completion pass if job_id and summary: diff --git a/evalbench/generators/models/agentic_reverse_proxy.py b/evalbench/generators/models/agentic_reverse_proxy.py index 9c4dc8d7..cbd87bea 100644 --- a/evalbench/generators/models/agentic_reverse_proxy.py +++ b/evalbench/generators/models/agentic_reverse_proxy.py @@ -2,9 +2,8 @@ import logging import queue import subprocess -import time import uuid -from typing import Any, Dict, List, Optional +from typing import Any, Dict from .agent_cli import AgentCliGenerator from evalproto import eval_agent_pb2 diff --git a/evalbench/reporting/remote_artifact_reporter.py b/evalbench/reporting/remote_artifact_reporter.py index 51578c6b..5e15ec0b 100644 --- a/evalbench/reporting/remote_artifact_reporter.py +++ b/evalbench/reporting/remote_artifact_reporter.py @@ -6,7 +6,7 @@ import pandas as pd from evalproto import eval_agent_pb2 -from reporting.report import Reporter, STORETYPE +from reporting.report import Reporter from util.context import rpc_id_var from generators.models.agentic_reverse_proxy import AGENT_PROXY_QUEUES diff --git a/evalbench/test/agentic_reverse_proxy_test.py b/evalbench/test/agentic_reverse_proxy_test.py index 66c848fb..2dbc7c88 100644 --- a/evalbench/test/agentic_reverse_proxy_test.py +++ b/evalbench/test/agentic_reverse_proxy_test.py @@ -1,6 +1,5 @@ import queue import unittest -from unittest.mock import MagicMock from evalproto import eval_agent_pb2 from generators.models.agentic_reverse_proxy import ( From 19644300333824026e5f324e8fee20b809488ea5 Mon Sep 17 00:00:00 2001 From: Saurabh Maurya Date: Tue, 18 Aug 2026 23:06:04 +0000 Subject: [PATCH 3/3] chore: adhere to repository linter standards and logging conventions --- evalbench/eval_service.py | 1 - .../models/agentic_reverse_proxy.py | 39 ++++++++++++------- .../reporting/remote_artifact_reporter.py | 10 ++--- evalbench/scorers/remote_scorer.py | 23 ++++++----- evalbench/test/agentic_reverse_proxy_test.py | 2 +- 5 files changed, 44 insertions(+), 31 deletions(-) diff --git a/evalbench/eval_service.py b/evalbench/eval_service.py index bffb5260..fc226a8a 100644 --- a/evalbench/eval_service.py +++ b/evalbench/eval_service.py @@ -30,7 +30,6 @@ eval_response_pb2, eval_service_pb2_grpc, eval_agent_pb2, - eval_agent_pb2_grpc, ) from util.service import ( load_session_configs, diff --git a/evalbench/generators/models/agentic_reverse_proxy.py b/evalbench/generators/models/agentic_reverse_proxy.py index cbd87bea..72520bba 100644 --- a/evalbench/generators/models/agentic_reverse_proxy.py +++ b/evalbench/generators/models/agentic_reverse_proxy.py @@ -3,19 +3,30 @@ import queue import subprocess import uuid -from typing import Any, Dict +from typing import Any -from .agent_cli import AgentCliGenerator from evalproto import eval_agent_pb2 from util.context import rpc_id_var +from .agent_cli import AgentCliGenerator + +logger = logging.getLogger(__name__) + # Global dictionary for active reverse proxy sessions: # session_id -> (inboxes_dict, out_queue) -AGENT_PROXY_QUEUES: Dict[str, tuple[Dict[str, queue.Queue], queue.Queue]] = {} +AGENT_PROXY_QUEUES: dict[str, tuple[dict[str, queue.Queue], queue.Queue]] = {} class CLICommand: - def __init__(self, cli: str, prompt: str, env: dict = None, resume: bool = False, session_id: str = None, cwd: str = None): + def __init__( + self, + cli: str, + prompt: str, + env: dict | None = None, + resume: bool = False, + session_id: str | None = None, + cwd: str | None = None, + ): self.cli = cli self.prompt = prompt self.env = env if env else {} @@ -27,12 +38,12 @@ def __init__(self, cli: str, prompt: str, env: dict = None, resume: bool = False class AgenticReverseProxyGenerator(AgentCliGenerator): """Generator proxying multi-turn agent execution across the reverse bidi stream.""" - def __init__(self, querygenerator_config: Dict[str, Any]): + def __init__(self, querygenerator_config: dict[str, Any]): super().__init__(querygenerator_config) self.name = "agentic_reverse_proxy" self.timeout_seconds = float(querygenerator_config.get("timeout_seconds", 300.0)) - self.turn_counter: Dict[str, int] = {} - logging.info("Initialized AgenticReverseProxyGenerator (timeout=%ss)", self.timeout_seconds) + self.turn_counter: dict[str, int] = {} + logger.info("Initialized AgenticReverseProxyGenerator (timeout=%ss)", self.timeout_seconds) @property def version(self) -> str: @@ -46,10 +57,10 @@ def create_command( self, cli: str, prompt: str, - env: dict = None, + env: dict | None = None, resume: bool = False, - session_id: str = None, - cwd: str = None, + session_id: str | None = None, + cwd: str | None = None, ) -> CLICommand: return CLICommand( cli=cli or self.name, @@ -85,7 +96,7 @@ def safe_generate(self, cli_cmd: CLICommand | dict | str) -> subprocess.Complete if ctx_id in AGENT_PROXY_QUEUES: session_id = ctx_id else: - logging.error("AgenticReverseProxy: session_id %s not in AGENT_PROXY_QUEUES (keys: %s)", session_id, list(AGENT_PROXY_QUEUES.keys())) + logger.error("AgenticReverseProxy: session_id %s not in AGENT_PROXY_QUEUES (keys: %s)", session_id, list(AGENT_PROXY_QUEUES.keys())) return subprocess.CompletedProcess( args=["agentic_reverse_proxy"], returncode=1, @@ -117,13 +128,13 @@ def safe_generate(self, cli_cmd: CLICommand | dict | str) -> subprocess.Complete turn_request=turn_req, ) - logging.info("[REVERSE_PROXY] Dispatching TurnRequest turn=%d (correlation_id=%s) to out_queue", turn_idx, correlation_id) + logger.info("[REVERSE_PROXY] Dispatching TurnRequest turn=%d (correlation_id=%s) to out_queue", turn_idx, correlation_id) out_queue.put(msg) try: resp_msg = inbox.get(timeout=self.timeout_seconds) except queue.Empty: - logging.error("[REVERSE_PROXY] Timed out waiting for TurnResponse (correlation_id=%s)", correlation_id) + logger.error("[REVERSE_PROXY] Timed out waiting for TurnResponse (correlation_id=%s)", correlation_id) return subprocess.CompletedProcess( args=["agentic_reverse_proxy"], returncode=124, @@ -135,7 +146,7 @@ def safe_generate(self, cli_cmd: CLICommand | dict | str) -> subprocess.Complete if not resp_msg.HasField("turn_response"): err_details = resp_msg.WhichOneof("payload") - logging.error("[REVERSE_PROXY] Unexpected message received on inbox: %s", err_details) + logger.error("[REVERSE_PROXY] Unexpected message received on inbox: %s", err_details) return subprocess.CompletedProcess( args=["agentic_reverse_proxy"], returncode=1, diff --git a/evalbench/reporting/remote_artifact_reporter.py b/evalbench/reporting/remote_artifact_reporter.py index 5e15ec0b..575ec0fd 100644 --- a/evalbench/reporting/remote_artifact_reporter.py +++ b/evalbench/reporting/remote_artifact_reporter.py @@ -1,14 +1,13 @@ import logging import queue import uuid -from typing import Any, Dict +from typing import Any import pandas as pd - from evalproto import eval_agent_pb2 +from generators.models.agentic_reverse_proxy import AGENT_PROXY_QUEUES from reporting.report import Reporter from util.context import rpc_id_var -from generators.models.agentic_reverse_proxy import AGENT_PROXY_QUEUES logger = logging.getLogger(__name__) @@ -16,12 +15,13 @@ class RemoteArtifactReporter(Reporter): """Reporter delegating workspace dump and GCS upload across the reverse bidi stream.""" - def __init__(self, reporting_config: Dict[str, Any] | None, job_id: str, run_time: Any): + def __init__(self, reporting_config: dict[str, Any] | None, job_id: str, run_time: Any): super().__init__(reporting_config, job_id, run_time) self.bucket = reporting_config.get("bucket", "sobi_dc_share") if reporting_config else "sobi_dc_share" self.path_prefix = reporting_config.get("path_prefix", "runs") if reporting_config else "runs" self.export_path = reporting_config.get("export_path", "/workspace") if reporting_config else "/workspace" - self.exclude_patterns = reporting_config.get("exclude_patterns", [".venv", "node_modules", "skills"]) if reporting_config else [".venv", "node_modules", "skills"] + default_excludes = [".venv", "node_modules", "skills"] + self.exclude_patterns = reporting_config.get("exclude_patterns", default_excludes) if reporting_config else default_excludes logger.info("Initialized RemoteArtifactReporter: bucket=%s, prefix=%s", self.bucket, self.path_prefix) def store(self, results: pd.DataFrame, store_type: Any) -> None: diff --git a/evalbench/scorers/remote_scorer.py b/evalbench/scorers/remote_scorer.py index 822a0a01..05ecd69b 100644 --- a/evalbench/scorers/remote_scorer.py +++ b/evalbench/scorers/remote_scorer.py @@ -2,23 +2,26 @@ import logging import queue import uuid -from typing import Any, Dict, List, Tuple +from typing import Any from evalproto import eval_agent_pb2 +from generators.models.agentic_reverse_proxy import AGENT_PROXY_QUEUES from scorers.comparator import Comparator from util.context import rpc_id_var -from generators.models.agentic_reverse_proxy import AGENT_PROXY_QUEUES + + +logger = logging.getLogger(__name__) class RemoteScorerProxy(Comparator): """Comparator proxying scoring evaluation across the reverse bidi stream.""" - def __init__(self, name: str, config: Dict[str, Any]): + def __init__(self, name: str, config: dict[str, Any]): super().__init__(config or {}) self.name = name self.config = dict(config) if isinstance(config, dict) else {} self.timeout_seconds = float(self.config.get("timeout_seconds", 300.0)) - logging.info("Initialized RemoteScorerProxy: name=%s, config=%s", self.name, self.config) + logger.info("Initialized RemoteScorerProxy: name=%s, config=%s", self.name, self.config) def compare( self, @@ -33,10 +36,10 @@ def compare( eval_results: str, generated_error: str, **kwargs: Any, - ) -> Tuple[float, str] | List[Tuple[str, float, str]]: + ) -> tuple[float, str] | list[tuple[str, float, str]]: session_id = rpc_id_var.get() if session_id not in AGENT_PROXY_QUEUES: - logging.error("RemoteScorerProxy: session_id %s not found in AGENT_PROXY_QUEUES", session_id) + logger.error("RemoteScorerProxy: session_id %s not found in AGENT_PROXY_QUEUES", session_id) return (0.0, f"Error: session_id '{session_id}' not connected to reverse stream") inboxes, out_queue = AGENT_PROXY_QUEUES[session_id] @@ -57,24 +60,24 @@ def compare( scoring_request=scoring_req, ) - logging.info("[REVERSE_SCORER] Dispatching ScoringRequest for '%s' (correlation_id=%s)", self.name, correlation_id) + logger.info("[REVERSE_SCORER] Dispatching ScoringRequest for '%s' (correlation_id=%s)", self.name, correlation_id) out_queue.put(msg) try: resp_msg = inbox.get(timeout=self.timeout_seconds) except queue.Empty: - logging.error("[REVERSE_SCORER] Timed out waiting for ScoringResponse for '%s' (correlation_id=%s)", self.name, correlation_id) + logger.error("[REVERSE_SCORER] Timed out waiting for ScoringResponse for '%s' (correlation_id=%s)", self.name, correlation_id) return (0.0, f"Error: Timed out waiting for remote scorer '{self.name}' response") finally: inboxes.pop(correlation_id, None) if not resp_msg.HasField("scoring_response"): err_details = resp_msg.WhichOneof("payload") - logging.error("[REVERSE_SCORER] Unexpected message on stream: %s", err_details) + logger.error("[REVERSE_SCORER] Unexpected message on stream: %s", err_details) return (0.0, f"Error: Unexpected payload on stream: {err_details}") scoring_resp = resp_msg.scoring_response - results: List[Tuple[str, float, str]] = [] + results: list[tuple[str, float, str]] = [] for r in scoring_resp.results: log_output = r.logs or r.stdout or r.error_message or f"exit_code={r.exit_code}" results.append((r.name, float(r.score), log_output)) diff --git a/evalbench/test/agentic_reverse_proxy_test.py b/evalbench/test/agentic_reverse_proxy_test.py index 2dbc7c88..5d5908a3 100644 --- a/evalbench/test/agentic_reverse_proxy_test.py +++ b/evalbench/test/agentic_reverse_proxy_test.py @@ -3,8 +3,8 @@ from evalproto import eval_agent_pb2 from generators.models.agentic_reverse_proxy import ( - AgenticReverseProxyGenerator, AGENT_PROXY_QUEUES, + AgenticReverseProxyGenerator, ) from reporting.remote_artifact_reporter import RemoteArtifactReporter from scorers.remote_scorer import RemoteScorerProxy