Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions evalbench/evalproto/eval_agent.proto
Original file line number Diff line number Diff line change
Expand Up @@ -89,8 +89,17 @@ message ReporterSpec {
float timeout_seconds = 3;
}

message ReportingContext {
string job_id = 1;
string run_time = 2; // Timestamp of evaluation run
string store_type = 3; // "CONFIGS", "EVALS", "SCORES", "SUMMARY"
string results_json = 4; // Serialized DataFrame records as JSON string (orient="records")
string database = 5; // Optional target database name
}

message ReportingRequest {
ReporterSpec reporter = 1;
ReportingContext context = 2;
}

message ReporterResult {
Expand Down
49 changes: 35 additions & 14 deletions evalbench/reporting/remote_reporter.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ def __init__(
self.reporter_name = reporter_name
self.config = dict(reporting_config) if isinstance(reporting_config, dict) else {}
self.timeout_seconds = float(self.config.get("timeout_seconds", 300.0))
self._reported = False
self.database = str(self.config.get("database", ""))
logger.info(
"Initialized RemoteReporter: name=%s, config=%s",
self.reporter_name,
Expand All @@ -36,18 +36,14 @@ def __init__(

def store(self, results: pd.DataFrame, store_type: Any) -> None:
type_name = getattr(store_type, "name", str(store_type))
# Delegate once during evaluation results processing
if type_name != "EVALS":
return
if self._reported:
return

session_id = rpc_id_var.get()
if session_id not in AGENT_GRPC_PROXY_QUEUES:
logger.warning(
"RemoteReporter: session_id %s not in AGENT_GRPC_PROXY_QUEUES, "
"skipping delegated reporting",
"skipping delegated reporting for %s",
session_id,
type_name,
)
return

Expand All @@ -62,15 +58,38 @@ def store(self, results: pd.DataFrame, store_type: Any) -> None:
timeout_seconds=self.timeout_seconds,
)

results_json = ""
if isinstance(results, pd.DataFrame) and not results.empty:
try:
results_json = results.to_json(orient="records", date_format="iso")
except Exception as e:
logger.warning(
"RemoteReporter: Failed to serialize DataFrame to JSON: %s",
e,
)
results_json = ""

reporting_context = eval_agent_pb2.ReportingContext(
job_id=str(self.job_id or ""),
run_time=str(self.run_time or ""),
store_type=type_name,
results_json=results_json,
database=self.database,
)

msg = eval_agent_pb2.AgentStreamMessage(
session_id=session_id,
correlation_id=correlation_id,
reporting_request=eval_agent_pb2.ReportingRequest(reporter=reporter_spec),
reporting_request=eval_agent_pb2.ReportingRequest(
reporter=reporter_spec,
context=reporting_context,
),
)

logger.info(
"[REVERSE_REPORTER] Dispatching ReportingRequest for '%s' (correlation_id=%s)",
"[REMOTE_REPORTER] Dispatching ReportingRequest for '%s' type=%s (correlation_id=%s)",
self.reporter_name,
type_name,
correlation_id,
)
out_queue.put(msg)
Expand All @@ -81,26 +100,28 @@ def store(self, results: pd.DataFrame, store_type: Any) -> None:
res = resp_msg.reporting_response.result
if res.success:
logger.info(
"[REVERSE_REPORTER] Delegated reporter '%s' completed successfully: %s",
"[REMOTE_REPORTER] Delegated reporter '%s' (%s) completed successfully: %s",
self.reporter_name,
type_name,
res.result_json,
)
else:
logger.error(
"[REVERSE_REPORTER] Delegated reporter '%s' failed: %s",
"[REMOTE_REPORTER] Delegated reporter '%s' (%s) failed: %s",
self.reporter_name,
type_name,
res.error_message,
)
self._reported = True
else:
logger.warning(
"[REVERSE_REPORTER] Received unexpected response on stream: %s",
"[REMOTE_REPORTER] Received unexpected response on stream: %s",
resp_msg.WhichOneof("payload"),
)
except queue.Empty:
logger.error(
"[REVERSE_REPORTER] Timed out waiting for ReportingResponse for '%s'",
"[REMOTE_REPORTER] Timed out waiting for ReportingResponse for '%s' (%s)",
self.reporter_name,
type_name,
)
finally:
inboxes.pop(correlation_id, None)
99 changes: 97 additions & 2 deletions evalbench/test/agent_grpc_proxy_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,7 @@ def answer_scorer():
def test_remote_reporter_success(self):
reporter = RemoteReporter(
"gcs_artifacts",
{"bucket": "test-bucket", "path_prefix": "runs"},
{"bucket": "test-bucket", "path_prefix": "runs", "database": "bigquery"},
job_id="job_123",
run_time="2026-08-18",
)
Expand All @@ -131,6 +131,14 @@ def answer_reporter():
rep_spec = msg.reporting_request.reporter
self.assertEqual(rep_spec.reporter_name, "gcs_artifacts")

ctx = msg.reporting_request.context
self.assertEqual(ctx.job_id, "job_123")
self.assertEqual(ctx.run_time, "2026-08-18")
self.assertEqual(ctx.store_type, "EVALS")
self.assertEqual(ctx.database, "bigquery")
parsed_data = json.loads(ctx.results_json)
self.assertEqual(parsed_data, [{"eval_id": "1"}])

reply = eval_agent_pb2.AgentStreamMessage(
session_id=self.session_id,
correlation_id=corr_id,
Expand All @@ -150,7 +158,94 @@ def answer_reporter():
df = pd.DataFrame({"eval_id": ["1"]})
reporter.store(df, STORETYPE.EVALS)
t.join()
self.assertTrue(reporter._reported)

def test_remote_reporter_all_store_types(self):
reporter = RemoteReporter(
"csv",
{"output_directory": "/tmp/results"},
job_id="job_456",
run_time="2026-08-18",
)

dispatched_types = []

def answer_reporter():
for _ in range(4):
msg = self.out_queue.get(timeout=2.0)
corr_id = msg.correlation_id
ctx = msg.reporting_request.context
dispatched_types.append(ctx.store_type)

reply = eval_agent_pb2.AgentStreamMessage(
session_id=self.session_id,
correlation_id=corr_id,
reporting_response=eval_agent_pb2.ReportingResponse(
result=eval_agent_pb2.ReporterResult(
reporter_name="csv",
success=True,
result_json=json.dumps({"written": ctx.store_type}),
)
),
)
self.inboxes[corr_id].put(reply)

t = threading.Thread(target=answer_reporter)
t.start()

df = pd.DataFrame({"key": ["val"]})
reporter.store(df, STORETYPE.CONFIGS)
reporter.store(df, STORETYPE.EVALS)
reporter.store(df, STORETYPE.SCORES)
reporter.store(df, STORETYPE.SUMMARY)
t.join()

self.assertEqual(
dispatched_types,
["CONFIGS", "EVALS", "SCORES", "SUMMARY"],
)

def test_remote_reporter_timeout(self):
reporter = RemoteReporter(
"slow_reporter",
{"timeout_seconds": 0.1},
job_id="job_slow",
run_time="2026-08-18",
)
df = pd.DataFrame({"eval_id": ["1"]})
# Should complete without raising exception
reporter.store(df, STORETYPE.EVALS)

def test_remote_reporter_failure_response(self):
reporter = RemoteReporter(
"failing_reporter",
{"timeout_seconds": 5.0},
job_id="job_fail",
run_time="2026-08-18",
)

def answer_failure():
msg = self.out_queue.get(timeout=2.0)
corr_id = msg.correlation_id

reply = eval_agent_pb2.AgentStreamMessage(
session_id=self.session_id,
correlation_id=corr_id,
reporting_response=eval_agent_pb2.ReportingResponse(
result=eval_agent_pb2.ReporterResult(
reporter_name="failing_reporter",
success=False,
error_message="GCS bucket not accessible",
)
),
)
self.inboxes[corr_id].put(reply)

t = threading.Thread(target=answer_failure)
t.start()

df = pd.DataFrame({"eval_id": ["1"]})
reporter.store(df, STORETYPE.EVALS)
t.join()


if __name__ == "__main__":
Expand Down
Loading