"""Resumable local CLI for the Agents Course evaluation and submission API.""" from __future__ import annotations import argparse import hashlib import importlib import json import os import re import shutil import sys from dataclasses import dataclass from datetime import UTC, datetime from pathlib import Path from typing import Any, Iterable import requests from agent_system import ( AgentConfigurationError, AgentSettings, LocalAgentSystem, clean_submission_value, is_placeholder_answer, ) from attachment_processing import AttachmentProcessingError, AttachmentProcessor PROJECT_ROOT = Path(__file__).resolve().parent DEFAULT_API_URL = "https://agents-course-unit4-scoring.hf.space" DEFAULT_SPACE_ID = "BmanClark/Agents_Course_final" DEFAULT_GAIA_REPO_ID = "gaia-benchmark/GAIA" DEFAULT_GAIA_DATA_DIR = "2023/validation" MAX_ATTACHMENT_BYTES = 100 * 1024 * 1024 class EvaluationError(RuntimeError): """Raised for invalid API responses, cache data, or submission state.""" def configure_utf8_stdio() -> None: """Keep agent/tool output Unicode-safe on Windows cp1252 consoles.""" for stream in (sys.stdout, sys.stderr): reconfigure = getattr(stream, "reconfigure", None) if callable(reconfigure): reconfigure(encoding="utf-8", errors="replace") @dataclass(frozen=True) class RunnerSettings: api_url: str username: str space_id: str local_dir: Path gaia_repo_id: str gaia_data_dir: str @classmethod def from_env(cls) -> "RunnerSettings": return cls( api_url=os.getenv("EVALUATION_API_URL", DEFAULT_API_URL).rstrip("/"), username=os.getenv("HF_USERNAME", "").strip(), space_id=os.getenv("SPACE_ID", DEFAULT_SPACE_ID).strip(), local_dir=Path( os.getenv("LOCAL_DATA_DIR", str(PROJECT_ROOT / ".local")) ).resolve(), gaia_repo_id=os.getenv( "GAIA_DATASET_REPO", DEFAULT_GAIA_REPO_ID ).strip(), gaia_data_dir=os.getenv( "GAIA_DATASET_DIR", DEFAULT_GAIA_DATA_DIR ).strip("/"), ) @property def agent_code_url(self) -> str: if "/" not in self.space_id: raise EvaluationError( "SPACE_ID must use the form username/space-name." ) return f"https://huggingface.co/spaces/{self.space_id}/tree/main" class EvaluationClient: def __init__(self, settings: RunnerSettings) -> None: self.settings = settings self.session = requests.Session() self.session.headers.update( {"User-Agent": "BmanClark-agents-course-local-runner/1.0"} ) def questions(self, random_only: bool = False) -> list[dict[str, Any]]: endpoint = "random-question" if random_only else "questions" try: response = self.session.get( f"{self.settings.api_url}/{endpoint}", timeout=30 ) response.raise_for_status() data = response.json() except (requests.RequestException, ValueError) as exc: raise EvaluationError(f"Could not fetch {endpoint}: {exc}") from exc if isinstance(data, dict): data = [data] if not isinstance(data, list) or not data: raise EvaluationError(f"The {endpoint} endpoint returned no tasks.") for item in data: if not isinstance(item, dict) or not item.get("task_id") or not item.get( "question" ): raise EvaluationError(f"Malformed question record: {item!r}") return data def download_attachment(self, question: dict[str, Any]) -> Path | None: file_name = str(question.get("file_name") or "").strip() if not file_name: return None task_id = safe_component(str(question["task_id"])) destination_dir = self.settings.local_dir / "attachments" / task_id destination_dir.mkdir(parents=True, exist_ok=True) destination = destination_dir / safe_filename(file_name) if destination.is_file() and destination.stat().st_size > 0: return destination partial = destination.with_suffix(destination.suffix + ".part") total = 0 try: with self.session.get( f"{self.settings.api_url}/files/{question['task_id']}", timeout=120, stream=True, ) as response: response.raise_for_status() declared_size = int(response.headers.get("content-length", "0") or 0) if declared_size > MAX_ATTACHMENT_BYTES: raise EvaluationError( f"Attachment {file_name} exceeds the 100 MB safety limit." ) with partial.open("wb") as handle: for chunk in response.iter_content(chunk_size=1024 * 1024): if not chunk: continue total += len(chunk) if total > MAX_ATTACHMENT_BYTES: raise EvaluationError( f"Attachment {file_name} exceeds the 100 MB safety limit." ) handle.write(chunk) partial.replace(destination) except requests.HTTPError as exc: partial.unlink(missing_ok=True) if exc.response is not None and exc.response.status_code == 404: return self._download_gaia_attachment(file_name, destination) raise EvaluationError(f"Could not download {file_name}: {exc}") from exc except (requests.RequestException, OSError, ValueError) as exc: partial.unlink(missing_ok=True) raise EvaluationError(f"Could not download {file_name}: {exc}") from exc except EvaluationError: partial.unlink(missing_ok=True) raise return destination def _download_gaia_attachment(self, file_name: str, destination: Path) -> Path: """Fall back to the official gated GAIA repository after a service 404.""" try: from huggingface_hub import hf_hub_download except ImportError as exc: raise EvaluationError( "The course file endpoint returned 404 and huggingface_hub is not " "installed for the official GAIA fallback." ) from exc repository_path = f"{self.settings.gaia_data_dir}/{safe_filename(file_name)}" fallback_dir = self.settings.local_dir / "hf-downloads" try: downloaded = Path( hf_hub_download( repo_id=self.settings.gaia_repo_id, filename=repository_path, repo_type="dataset", local_dir=fallback_dir, ) ) if downloaded.stat().st_size > MAX_ATTACHMENT_BYTES: raise EvaluationError( f"Attachment {file_name} exceeds the 100 MB safety limit." ) destination.parent.mkdir(parents=True, exist_ok=True) shutil.copyfile(downloaded, destination) except EvaluationError: raise except Exception as exc: raise EvaluationError( "The course file endpoint returned 404 and the official gated GAIA " "fallback could not download the attachment. Accept access at " "https://huggingface.co/datasets/gaia-benchmark/GAIA, then run " r".\.venv\Scripts\hf.exe auth login. " f"Underlying error: {exc}" ) from exc return destination def submit(self, answers: list[dict[str, str]]) -> dict[str, Any]: if not self.settings.username: raise EvaluationError( "HF_USERNAME is required for submission. Set it in the shell first." ) payload = { "username": self.settings.username, "agent_code": self.settings.agent_code_url, "answers": answers, } try: response = self.session.post( f"{self.settings.api_url}/submit", json=payload, timeout=120 ) response.raise_for_status() result = response.json() except requests.HTTPError as exc: detail = exc.response.text[:1_000] if exc.response is not None else str(exc) raise EvaluationError(f"Submission was rejected: {detail}") from exc except (requests.RequestException, ValueError) as exc: raise EvaluationError(f"Submission failed: {exc}") from exc if not isinstance(result, dict): raise EvaluationError("Submission response was not a JSON object.") return result class AnswerCache: """Private, atomic local cache keyed by evaluation task ID.""" VERSION = 1 def __init__(self, path: Path) -> None: self.path = path self.data: dict[str, Any] = {"version": self.VERSION, "answers": {}} self.load() def load(self) -> None: if not self.path.exists(): return try: data = json.loads(self.path.read_text(encoding="utf-8")) except (OSError, ValueError) as exc: raise EvaluationError(f"Could not read answer cache {self.path}: {exc}") from exc if data.get("version") != self.VERSION or not isinstance( data.get("answers"), dict ): raise EvaluationError( f"Unsupported or malformed answer cache: {self.path}" ) self.data = data def get_valid(self, question: dict[str, Any]) -> str | None: entry = self.get_entry(question) if entry is None: return None answer = entry.get("answer") return answer if isinstance(answer, str) and answer.strip() else None def get_entry(self, question: dict[str, Any]) -> dict[str, Any] | None: entry = self.data["answers"].get(str(question["task_id"])) if not isinstance(entry, dict): return None if entry.get("question_sha256") != question_digest(str(question["question"])): return None return entry def is_validated(self, question: dict[str, Any]) -> bool: entry = self.get_entry(question) return bool(entry is not None and entry.get("validated", True)) def record( self, question: dict[str, Any], answer: str, agent_signature: str, attachment_name: str | None, validated: bool = True, validation_issues: list[str] | None = None, ) -> None: self.data["answers"][str(question["task_id"])] = { "answer": answer, "question_sha256": question_digest(str(question["question"])), "agent_signature": agent_signature, "attachment_name": attachment_name, "validated": validated, "validation_issues": validation_issues or [], "completed_at": datetime.now(UTC).isoformat(), } self.save() def save(self) -> None: self.path.parent.mkdir(parents=True, exist_ok=True) temporary = self.path.with_suffix(self.path.suffix + ".tmp") temporary.write_text( json.dumps(self.data, indent=2, ensure_ascii=False) + "\n", encoding="utf-8", ) temporary.replace(self.path) def question_digest(question: str) -> str: return hashlib.sha256(question.encode("utf-8")).hexdigest() def safe_component(value: str) -> str: cleaned = re.sub(r"[^A-Za-z0-9._-]", "_", value) if not cleaned or cleaned in {".", ".."}: raise EvaluationError(f"Unsafe path component: {value!r}") return cleaned def safe_filename(value: str) -> str: name = Path(value.replace("\\", "/")).name return safe_component(name) def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser( description="Run and submit the Hugging Face Agents Course evaluation locally." ) commands = parser.add_subparsers(dest="command", required=True) commands.add_parser("check", help="Check dependencies, Ollama, and local models.") commands.add_parser("test", help="Solve and cache one random evaluation task.") run = commands.add_parser("run", help="Solve and cache evaluation tasks.") run.add_argument("--task-id", action="append", help="Only run this task ID.") run.add_argument("--limit", type=int, help="Run at most this many selected tasks.") run.add_argument("--force", action="store_true", help="Ignore valid cached answers.") run.add_argument( "--best-effort", action="store_true", help=( "Cache an explicitly marked best-effort answer when evidence validation " "still fails after all configured research rounds." ), ) commands.add_parser("status", help="Show cache coverage without displaying answers.") submit = commands.add_parser("submit", help="Submit all valid cached answers.") submit.add_argument( "--yes", action="store_true", help="Skip the interactive SUBMIT confirmation." ) return parser def check_environment() -> None: required_modules = [ "av", "requests", "smolagents", "litellm", "openpyxl", "faster_whisper", ] missing = [] for module in required_modules: try: importlib.import_module(module) except ImportError: missing.append(module) if missing: raise EvaluationError( "Missing Python modules: " + ", ".join(missing) + ". Run: python -m pip install -r requirements.txt" ) models = LocalAgentSystem.check_ollama(AgentSettings.from_env()) print(f"Ollama is reachable; {len(models)} local model(s) found.") print("Required text and multimodal models are installed.") def selected_questions( questions: Iterable[dict[str, Any]], task_ids: list[str] | None, limit: int | None ) -> list[dict[str, Any]]: selected = list(questions) if task_ids: wanted = set(task_ids) selected = [q for q in selected if str(q["task_id"]) in wanted] found = {str(q["task_id"]) for q in selected} missing = sorted(wanted - found) if missing: raise EvaluationError("Unknown task ID(s): " + ", ".join(missing)) if limit is not None: if limit < 1: raise EvaluationError("--limit must be at least 1.") selected = selected[:limit] return selected def solve_tasks( questions: list[dict[str, Any]], client: EvaluationClient, cache: AnswerCache, force: bool, best_effort: bool = False, ) -> int: agent: LocalAgentSystem | None = None processor = AttachmentProcessor() failures = 0 for index, question in enumerate(questions, start=1): task_id = str(question["task_id"]) cached = cache.get_valid(question) if cached is not None and not force: print(f"[{index}/{len(questions)}] {task_id}: cached; skipping") continue print(f"[{index}/{len(questions)}] {task_id}: solving") try: attachment = client.download_attachment(question) evidence = processor.process(attachment, str(question["question"])) if agent is None: LocalAgentSystem.check_ollama(AgentSettings.from_env()) agent = LocalAgentSystem() result = agent.solve_result( task_id, str(question["question"]), evidence, allow_best_effort=best_effort, ) cache.record( question, result.answer, agent.signature, attachment.name if attachment else None, validated=result.validated, validation_issues=result.issues, ) quality = "validated" if result.validated else "BEST EFFORT" print( f"[{index}/{len(questions)}] {task_id}: " f"{quality} answer cached: {result.answer}" ) except ( AgentConfigurationError, AttachmentProcessingError, EvaluationError, ValueError, ) as exc: failures += 1 print(f"[{index}/{len(questions)}] {task_id}: ERROR: {exc}", file=sys.stderr) return failures def print_status(questions: list[dict[str, Any]], cache: AnswerCache) -> int: complete = sum(cache.get_valid(question) is not None for question in questions) validated = sum(cache.is_validated(question) for question in questions) best_effort = complete - validated print( f"Cached answers: {complete}/{len(questions)} " f"({validated} validated, {best_effort} best effort)" ) for question in questions: if cache.get_valid(question) is None: state = "missing" elif cache.is_validated(question): state = "validated" else: state = "best-effort" attachment = str(question.get("file_name") or "none") print(f" {question['task_id']}: {state}; attachment={attachment}") return complete def submit_cached( questions: list[dict[str, Any]], client: EvaluationClient, cache: AnswerCache, assume_yes: bool, ) -> None: answers = [] missing = [] for question in questions: answer = cache.get_valid(question) if answer is None: missing.append(str(question["task_id"])) else: answer = audit_submission_answer(str(question["task_id"]), answer) answers.append( {"task_id": str(question["task_id"]), "submitted_answer": answer} ) if missing: raise EvaluationError( f"Refusing a partial submission: {len(missing)} task(s) are missing." ) print(f"Username: {client.settings.username or ''}") print(f"Agent code: {client.settings.agent_code_url}") print(f"Answers ready: {len(answers)}") validated = sum(cache.is_validated(question) for question in questions) print(f"Evidence-validated: {validated}; best effort: {len(answers) - validated}") if not assume_yes: confirmation = input("Type SUBMIT to send these answers for scoring: ").strip() if confirmation != "SUBMIT": print("Submission cancelled.") return result = client.submit(answers) submission_dir = client.settings.local_dir / "submissions" submission_dir.mkdir(parents=True, exist_ok=True) timestamp = datetime.now(UTC).strftime("%Y%m%dT%H%M%SZ") (submission_dir / f"{timestamp}.json").write_text( json.dumps(result, indent=2, ensure_ascii=False) + "\n", encoding="utf-8" ) print( "Submission successful: " f"{result.get('score', 'N/A')}% " f"({result.get('correct_count', '?')}/{result.get('total_attempted', '?')})" ) if result.get("message"): print(result["message"]) def audit_submission_answer(task_id: str, raw_answer: str) -> str: """Normalize and reject malformed legacy cache entries before submission.""" try: answer = clean_submission_value(raw_answer) except ValueError as error: raise EvaluationError( f"Task {task_id} has an invalid cached answer: {error}" ) from error if is_placeholder_answer(answer): raise EvaluationError( f"Task {task_id} still has a placeholder cached answer: {answer!r}" ) if not any(character.isalnum() for character in answer): raise EvaluationError( f"Task {task_id} has a punctuation-only cached answer: {answer!r}" ) return answer def main(argv: list[str] | None = None) -> int: configure_utf8_stdio() args = build_parser().parse_args(argv) settings = RunnerSettings.from_env() client = EvaluationClient(settings) cache = AnswerCache(settings.local_dir / "answers.json") try: if args.command == "check": check_environment() return 0 if args.command == "test": questions = client.questions(random_only=True) return 1 if solve_tasks(questions, client, cache, force=True) else 0 questions = client.questions() if args.command == "status": print_status(questions, cache) return 0 if args.command == "run": chosen = selected_questions(questions, args.task_id, args.limit) return ( 1 if solve_tasks( chosen, client, cache, args.force, best_effort=args.best_effort, ) else 0 ) if args.command == "submit": submit_cached(questions, client, cache, args.yes) return 0 except (AgentConfigurationError, EvaluationError) as exc: print(f"Error: {exc}", file=sys.stderr) return 2 raise AssertionError(f"Unhandled command: {args.command}")