from __future__ import annotations import json from datetime import datetime, timezone from pathlib import Path from typing import Any from .bucket import RunPaths, read_json, read_text, write_json from .eval_config import effective_eval_config from .security import redact _TERMINAL = {"success", "failed", "failure", "error", "cancelled", "manual", "stale"} def _now() -> str: return datetime.now(timezone.utc).isoformat() def _safe_segment(value: Any, default: str) -> str: cleaned = "".join(ch if ch.isalnum() or ch in {"-", "_"} else "-" for ch in str(value or "").strip()) return cleaned.strip("-_") or default def _eval_mount(cfg: dict[str, Any] | None = None) -> Path | None: cfg = cfg or effective_eval_config() if not cfg.get("enabled"): return None mount = Path(str(cfg.get("job_mount_path") or "/evals")) return mount def _publish_status_path(paths: RunPaths) -> str: return f"{paths.root}/eval_publish_status.json" def _should_attempt_publish(record: dict[str, Any], bundle_state: dict[str, Any] | None = None) -> bool: status = str(record.get("outcome", {}).get("status") or record.get("phase") or "").lower() state_status = str((bundle_state or {}).get("status") or "").lower() phase = str(record.get("phase") or "").lower() return status in _TERMINAL or state_status in _TERMINAL or phase in {"final", "failure"} def maybe_publish_eval_record( run_id: str, *, bucket_source: str, token: str | None = None, state: dict[str, Any] | None = None, force: bool = False, ) -> dict[str, Any]: """Publish a run-local anonymized eval record to the operator archive. User Jobs write only ``runs//eval_record.json`` in the user's bucket. The ASF Space backend owns the private operator eval bucket mounted at ``/evals`` and copies the already-anonymized record there. This lets evals work for all users without giving their Jobs write access to the private operator bucket. """ cfg = effective_eval_config() if not cfg.get("enabled"): return {"enabled": False, "published": False, "reason": "eval_archive_disabled"} mount = _eval_mount(cfg) if mount is None: return {"enabled": False, "published": False, "reason": "eval_archive_disabled"} if not mount.exists() or not mount.is_dir(): return {"enabled": True, "published": False, "reason": "eval_mount_missing", "mount": str(mount)} if not mount.is_dir(): return {"enabled": True, "published": False, "reason": "eval_mount_not_directory", "mount": str(mount)} paths = RunPaths(run_id, bucket_source=bucket_source) status_path = _publish_status_path(paths) if not force: prev = read_json(status_path, token=token) or {} if prev.get("published") and prev.get("schema_version") == "1.0": return {**prev, "enabled": True, "skipped": True, "reason": "already_published"} record = read_json(f"{paths.root}/eval_record.json", token=token) or {} if not record: return {"enabled": True, "published": False, "reason": "record_not_ready"} if not isinstance(record, dict): return {"enabled": True, "published": False, "reason": "invalid_record"} if not force and not _should_attempt_publish(record, state): return {"enabled": True, "published": False, "reason": "run_not_terminal"} privacy = record.get("privacy") if isinstance(record.get("privacy"), dict) else {} if privacy.get("generated_code_stored") or privacy.get("raw_prompts_stored") or privacy.get("tokens_stored"): return {"enabled": True, "published": False, "reason": "privacy_flags_rejected"} finished = str(record.get("finished_at") or record.get("started_at") or _now()) day = finished[:10] if len(finished) >= 10 else _now()[:10] try: yyyy, mm, dd = day.split("-") except Exception: yyyy, mm, dd = _now()[:10].split("-") anon_run_id = _safe_segment(record.get("anonymous_run_id") or record.get("run_id_hash"), "unknown-run") dest_dir = mount / yyyy / mm / dd / anon_run_id dest_dir.mkdir(parents=True, exist_ok=True) dest = dest_dir / "eval_record.json" dest.write_text(json.dumps(record, indent=2, ensure_ascii=False, sort_keys=True) + "\n", encoding="utf-8") compact = read_text(f"{paths.root}/events_compact.jsonl", token=token) if not compact: timeline = record.get("timeline") if isinstance(record.get("timeline"), list) else [] compact = "".join(json.dumps(item, ensure_ascii=False) + "\n" for item in timeline) if compact: (dest_dir / "events_compact.jsonl").write_text(compact, encoding="utf-8") report = { "schema_version": "1.0", "published": True, "published_at": _now(), "publish_mode": "backend", "archive_path": str(dest), "archive_relative_path": str(dest.relative_to(mount)), "eval_bucket_source": cfg.get("bucket_source") or "", "eval_bucket_path": cfg.get("bucket_path") or "evals", "mount": str(mount), "anonymous_run_id": anon_run_id, } try: write_json(status_path, report, token=token) except Exception as exc: # noqa: BLE001 report["status_write_warning"] = redact(str(exc))[:500] return {"enabled": True, **report}