from __future__ import annotations from pathlib import Path from typing import Any import json from .schemas import ExperimentReport, Hypothesis, Insight, PipelineConfiguration, ValidationReport class SharedKnowledgeSpace: """File-based bridge between V1 and T1.""" def __init__(self, root: str | Path): self.root = Path(root) self.reports_dir = self.root / "experiment_reports" self.insights_dir = self.root / "insights" self.configs_dir = self.root / "pipeline_configs" self.hypotheses_dir = self.root / "hypotheses" self.validation_reports_dir = self.root / "validation_reports" self.summary_stats_file = self.root / "summary_statistics.json" for d in ( self.reports_dir, self.insights_dir, self.configs_dir, self.hypotheses_dir, self.validation_reports_dir, ): d.mkdir(parents=True, exist_ok=True) @staticmethod def _write_json(path: Path, payload: dict[str, Any]) -> None: path.write_text(json.dumps(payload, indent=2, ensure_ascii=True), encoding="utf-8") @staticmethod def _read_json(path: Path) -> dict[str, Any]: return json.loads(path.read_text(encoding="utf-8")) def save_report(self, report: ExperimentReport) -> Path: fp = self.reports_dir / f"{report.run_id}.json" self._write_json(fp, report.to_dict()) return fp def save_insight(self, insight: Insight) -> Path: fp = self.insights_dir / f"{insight.insight_id}.json" self._write_json(fp, insight.to_dict()) return fp def save_pipeline_config(self, config: PipelineConfiguration) -> Path: fp = self.configs_dir / f"{config.config_id}.json" self._write_json(fp, config.to_dict()) return fp def save_hypothesis(self, hypothesis: Hypothesis) -> Path: fp = self.hypotheses_dir / f"{hypothesis.hypothesis_id}.json" self._write_json(fp, hypothesis.to_dict()) return fp def save_validation_report(self, report: ValidationReport) -> Path: fp = self.validation_reports_dir / f"{report.validation_id}.json" self._write_json(fp, report.to_dict()) return fp def list_reports(self) -> list[dict[str, Any]]: items: list[dict[str, Any]] = [] for fp in sorted(self.reports_dir.glob("*.json")): items.append(self._read_json(fp)) return items def list_insights(self) -> list[dict[str, Any]]: items: list[dict[str, Any]] = [] for fp in sorted(self.insights_dir.glob("*.json")): items.append(self._read_json(fp)) return items def list_configs(self) -> list[dict[str, Any]]: items: list[dict[str, Any]] = [] for fp in sorted(self.configs_dir.glob("*.json")): items.append(self._read_json(fp)) return items def list_hypotheses(self) -> list[dict[str, Any]]: items: list[dict[str, Any]] = [] for fp in sorted(self.hypotheses_dir.glob("*.json")): items.append(self._read_json(fp)) return items def list_validation_reports(self) -> list[dict[str, Any]]: items: list[dict[str, Any]] = [] for fp in sorted(self.validation_reports_dir.glob("*.json")): items.append(self._read_json(fp)) return items def latest_config_for_task(self, task_scope: str) -> dict[str, Any] | None: matched = [cfg for cfg in self.list_configs() if cfg.get("task_scope") == task_scope] return matched[-1] if matched else None def get_hypothesis_context(self, domain: str, top_k: int = 5) -> list[dict[str, Any]]: domain_l = domain.strip().lower() hypotheses = self.list_hypotheses() validations = self.list_validation_reports() by_hid = {} for vr in validations: hid = vr.get("hypothesis_id", "") if not hid: continue by_hid.setdefault(hid, []).append(vr) scored: list[dict[str, Any]] = [] for h in hypotheses: h_domain = str(h.get("domain", "")).lower() if domain_l and domain_l not in h_domain and domain_l not in " ".join(h.get("tags", [])).lower(): continue reports = by_hid.get(h.get("hypothesis_id", ""), []) if reports: avg_score = sum(float(r.get("score", 0.0)) for r in reports) / len(reports) avg_conf = sum(float(r.get("confidence", 0.0)) for r in reports) / len(reports) else: avg_score = 0.0 avg_conf = 0.0 scored.append( { "hypothesis": h, "validation_reports": reports, "avg_score": avg_score, "avg_confidence": avg_conf, } ) scored.sort(key=lambda x: (x["avg_score"], x["avg_confidence"]), reverse=True) return scored[: max(0, top_k)] def update_summary_statistics(self) -> dict[str, Any]: validations = self.list_validation_reports() tag_stats: dict[str, dict[str, float]] = {} failure_reason_counter: dict[str, int] = {} success_scores: list[float] = [] hypotheses = {h.get("hypothesis_id", ""): h for h in self.list_hypotheses()} for r in validations: status = r.get("status", "inconclusive") score = float(r.get("score", 0.0)) if status == "success": success_scores.append(score) reason = str(r.get("failure_reason", "")).strip() if status == "failed" and reason: failure_reason_counter[reason] = failure_reason_counter.get(reason, 0) + 1 hid = r.get("hypothesis_id", "") h = hypotheses.get(hid, {}) for tag in h.get("tags", []): s = tag_stats.setdefault(tag, {"count": 0.0, "success": 0.0, "score_sum": 0.0}) s["count"] += 1 if status == "success": s["success"] += 1 s["score_sum"] += score tag_summary = {} for tag, s in tag_stats.items(): count = max(1.0, s["count"]) tag_summary[tag] = { "count": int(s["count"]), "success_rate": s["success"] / count, "avg_score": s["score_sum"] / count, } payload = { "total_validation_reports": len(validations), "overall_avg_success_score": (sum(success_scores) / len(success_scores)) if success_scores else 0.0, "tag_summary": tag_summary, "failure_reason_counter": failure_reason_counter, } self._write_json(self.summary_stats_file, payload) return payload def get_summary_statistics(self) -> dict[str, Any]: if not self.summary_stats_file.exists(): return self.update_summary_statistics() return self._read_json(self.summary_stats_file)