diff --git a/src/paperbot/api/routes/paperscool.py b/src/paperbot/api/routes/paperscool.py index 5a09f237..57554350 100644 --- a/src/paperbot/api/routes/paperscool.py +++ b/src/paperbot/api/routes/paperscool.py @@ -3,6 +3,7 @@ import copy import os import re +import time from threading import Thread from typing import Any, Dict, List, Optional from urllib.parse import urlparse @@ -41,11 +42,13 @@ from paperbot.infrastructure.stores.paper_store import PaperStore from paperbot.infrastructure.stores.pipeline_session_store import PipelineSessionStore from paperbot.infrastructure.stores.research_store import SqlAlchemyResearchStore +from paperbot.infrastructure.stores.workflow_metric_store import WorkflowMetricStore from paperbot.utils.text_processing import extract_github_url router = APIRouter() _paper_search_service: Optional[PaperSearchService] = None _pipeline_session_store = PipelineSessionStore() +_workflow_metric_store: Optional[WorkflowMetricStore] = None # Test compatibility hook: unit tests can monkeypatch this to inject a fake workflow. PapersCoolTopicSearchWorkflow = None @@ -84,6 +87,24 @@ def _get_paper_search_service() -> PaperSearchService: return _paper_search_service +def _get_workflow_metric_store() -> WorkflowMetricStore: + global _workflow_metric_store + if _workflow_metric_store is None: + _workflow_metric_store = WorkflowMetricStore() + return _workflow_metric_store + + +def _count_report_claims_and_evidence(report: Dict[str, Any]) -> tuple[int, int]: + claims = 0 + evidences = 0 + for query in report.get("queries") or []: + for item in query.get("top_items") or []: + claims += 1 + if item.get("url") or item.get("pdf_url") or item.get("external_url"): + evidences += 1 + return claims, evidences + + async def _run_topic_search( *, queries: List[str], @@ -237,6 +258,10 @@ async def topic_search(req: PapersCoolSearchRequest): async def _dailypaper_stream(req: DailyPaperRequest): """SSE generator for the full DailyPaper pipeline.""" cleaned_queries = [q.strip() for q in req.queries if (q or "").strip()] + started = time.perf_counter() + phase_ms: Dict[str, float] = {} + phase_start = started + metric_store = _get_workflow_metric_store() session = _pipeline_session_store.start_session( workflow="paperscool_daily", @@ -334,6 +359,8 @@ async def _dailypaper_stream(req: DailyPaperRequest): "session_id": session_id, }, ) + phase_ms["search"] = round((time.perf_counter() - phase_start) * 1000.0, 2) + phase_start = time.perf_counter() # Phase 2 — Build Report if req.resume and isinstance(session_state.get("report"), dict): @@ -362,6 +389,8 @@ async def _dailypaper_stream(req: DailyPaperRequest): "session_id": session_id, }, ) + phase_ms["build"] = round((time.perf_counter() - phase_start) * 1000.0, 2) + phase_start = time.perf_counter() query_items: List[Dict[str, Any]] = [] paper_query_map: Dict[int, str] = {} @@ -612,9 +641,12 @@ async def _dailypaper_stream(req: DailyPaperRequest): checkpoint="enriched", state={"search_result": search_result, "report": report}, ) + phase_ms["enrich"] = round((time.perf_counter() - phase_start) * 1000.0, 2) + phase_start = time.perf_counter() if req.require_approval: preview_markdown = render_daily_paper_markdown(report) + claims, evidences = _count_report_claims_and_evidence(report) pending_payload = { "report": report, "markdown": preview_markdown, @@ -641,6 +673,19 @@ async def _dailypaper_stream(req: DailyPaperRequest): }, ) yield StreamEvent(type="result", data=pending_payload) + metric_store.record_metric( + workflow="paperscool_daily", + stage="approval_pending", + status="pending_approval", + claim_count=claims, + evidence_count=evidences, + elapsed_ms=(time.perf_counter() - started) * 1000.0, + detail={ + "session_id": session_id, + "phase_ms": phase_ms, + "resume": bool(req.resume), + }, + ) return # Phase 5 — Persist + Notify @@ -689,6 +734,8 @@ async def _dailypaper_stream(req: DailyPaperRequest): email_to_override=_validate_email_list(req.notify_email_to) or None, ) + phase_ms["persist"] = round((time.perf_counter() - phase_start) * 1000.0, 2) + result_payload = { "report": report, "markdown": markdown, @@ -703,6 +750,22 @@ async def _dailypaper_stream(req: DailyPaperRequest): session_id=session_id, result=result_payload, status="completed" ) + claims, evidences = _count_report_claims_and_evidence(report) + metric_store.record_metric( + workflow="paperscool_daily", + stage="result", + status="completed", + claim_count=claims, + evidence_count=evidences, + elapsed_ms=(time.perf_counter() - started) * 1000.0, + detail={ + "session_id": session_id, + "phase_ms": phase_ms, + "enable_judge": bool(req.enable_judge), + "enable_llm_analysis": bool(req.enable_llm_analysis), + }, + ) + yield StreamEvent(type="result", data=result_payload) @@ -712,6 +775,9 @@ async def generate_daily_report(req: DailyPaperRequest): if not cleaned_queries: raise HTTPException(status_code=400, detail="queries is required") + started = time.perf_counter() + metric_store = _get_workflow_metric_store() + # Fast sync path when no long-running step is requested — avoids SSE overhead if ( not req.enable_llm_analysis @@ -720,7 +786,29 @@ async def generate_daily_report(req: DailyPaperRequest): and not req.resume and not req.session_id ): - return await _sync_daily_report(req, cleaned_queries) + try: + payload = await _sync_daily_report(req, cleaned_queries) + report = payload.report if isinstance(payload, DailyPaperResponse) else {} + claims, evidences = _count_report_claims_and_evidence(report) + metric_store.record_metric( + workflow="paperscool_daily", + stage="sync_result", + status="completed", + claim_count=claims, + evidence_count=evidences, + elapsed_ms=(time.perf_counter() - started) * 1000.0, + detail={"mode": "sync"}, + ) + return payload + except Exception as exc: + metric_store.record_metric( + workflow="paperscool_daily", + stage="sync_result", + status="failed", + elapsed_ms=(time.perf_counter() - started) * 1000.0, + detail={"mode": "sync", "error": str(exc)}, + ) + raise # SSE streaming path for long-running operations return StreamingResponse( @@ -842,6 +930,7 @@ async def approve_daily_session(session_id: str): raise HTTPException(status_code=409, detail="session is not pending approval") final_payload = _finalize_approved_session(session) + claims, evidences = _count_report_claims_and_evidence(final_payload.get("report") or {}) _pipeline_session_store.update_status( session_id=session_id, status="completed", @@ -849,6 +938,14 @@ async def approve_daily_session(session_id: str): state_patch={"approved_at": True}, result=final_payload, ) + _get_workflow_metric_store().record_metric( + workflow="paperscool_daily", + stage="approval_finalize", + status="completed", + claim_count=claims, + evidence_count=evidences, + detail={"session_id": session_id, "mode": "approval"}, + ) updated = _pipeline_session_store.get_session(session_id) return PipelineSessionResponse(session=updated or {}) @@ -878,6 +975,12 @@ async def reject_daily_session(session_id: str, req: ApprovalDecisionRequest): state_patch={"reject_reason": req.reason or ""}, result=rejected_result, ) + _get_workflow_metric_store().record_metric( + workflow="paperscool_daily", + stage="approval_finalize", + status="rejected", + detail={"session_id": session_id, "reason": req.reason or ""}, + ) updated = _pipeline_session_store.get_session(session_id) return PipelineSessionResponse(session=updated or {}) @@ -1196,6 +1299,8 @@ def enrich_papers_with_repo_data(req: PapersCoolReposRequest): async def _paperscool_analyze_stream(req: PapersCoolAnalyzeRequest): + started = time.perf_counter() + metric_store = _get_workflow_metric_store() report = copy.deepcopy(req.report) llm_service = get_llm_service() @@ -1353,6 +1458,20 @@ async def _paperscool_analyze_stream(req: PapersCoolAnalyzeRequest): yield StreamEvent(type="judge_done", data=report["judge"]) markdown = render_daily_paper_markdown(report) + claims, evidences = _count_report_claims_and_evidence(report) + metric_store.record_metric( + workflow="paperscool_analyze", + stage="result", + status="completed", + claim_count=claims, + evidence_count=evidences, + elapsed_ms=(time.perf_counter() - started) * 1000.0, + detail={ + "run_judge": bool(req.run_judge), + "run_trends": bool(req.run_trends), + "run_insight": bool(req.run_insight), + }, + ) yield StreamEvent(type="result", data={"report": report, "markdown": markdown}) diff --git a/src/paperbot/api/routes/research.py b/src/paperbot/api/routes/research.py index 786a616f..9222ef60 100644 --- a/src/paperbot/api/routes/research.py +++ b/src/paperbot/api/routes/research.py @@ -2,6 +2,7 @@ import os import re +import time from collections import Counter from datetime import datetime, timedelta, timezone from typing import Any, Dict, List, Optional, Tuple @@ -15,6 +16,7 @@ from paperbot.infrastructure.api_clients.semantic_scholar import SemanticScholarClient from paperbot.infrastructure.stores.memory_store import SqlAlchemyMemoryStore from paperbot.infrastructure.stores.research_store import SqlAlchemyResearchStore +from paperbot.infrastructure.stores.workflow_metric_store import WorkflowMetricStore from paperbot.memory.eval.collector import MemoryMetricCollector from paperbot.memory.extractor import extract_memories from paperbot.memory.schema import MemoryCandidate, NormalizedMessage @@ -26,6 +28,7 @@ _memory_store = SqlAlchemyMemoryStore() _track_router = TrackRouter(research_store=_research_store, memory_store=_memory_store) _metric_collector: Optional[MemoryMetricCollector] = None +_workflow_metric_store: Optional[WorkflowMetricStore] = None _paper_store: Optional["PaperStore"] = None _paper_search_service: Optional["PaperSearchService"] = None @@ -97,6 +100,13 @@ def _get_metric_collector() -> MemoryMetricCollector: return _metric_collector +def _get_workflow_metric_store() -> WorkflowMetricStore: + global _workflow_metric_store + if _workflow_metric_store is None: + _workflow_metric_store = WorkflowMetricStore() + return _workflow_metric_store + + def _get_paper_store() -> "PaperStore": """Lazy initialization of paper store.""" from paperbot.infrastructure.stores.paper_store import PaperStore @@ -1105,11 +1115,32 @@ class ContextResponse(BaseModel): context_pack: Dict[str, Any] +class WorkflowMetricsResponse(BaseModel): + summary: Dict[str, Any] + + +@router.get("/research/metrics/workflows", response_model=WorkflowMetricsResponse) +def get_workflow_metrics_summary( + days: int = Query(7, ge=1, le=90), + workflow: Optional[str] = None, + track_id: Optional[int] = None, +): + summary = _get_workflow_metric_store().summarize( + days=days, + workflow=workflow, + track_id=track_id, + ) + return WorkflowMetricsResponse(summary=summary) + + @router.post("/research/context", response_model=ContextResponse) async def build_context(req: ContextRequest): set_trace_id() # Initialize trace_id for this request Logger.info("Received build context request", file=LogFiles.HARVEST) + started = time.perf_counter() + metric_store = _get_workflow_metric_store() + if req.activate_track_id is not None: Logger.info("Activating research track", file=LogFiles.HARVEST) activated = _research_store.activate_track( @@ -1162,7 +1193,38 @@ async def build_context(req: ContextRequest): Logger.info( f"Context pack built successfully, found {paper_count} papers", file=LogFiles.HARVEST ) + evidence_count = sum(1 for p in (pack.get("paper_recommendations") or []) if p.get("url")) + metric_store.record_metric( + workflow="research_context", + stage=req.stage, + status="completed", + track_id=req.track_id, + claim_count=paper_count, + evidence_count=evidence_count, + elapsed_ms=(time.perf_counter() - started) * 1000.0, + detail={ + "paper_limit": int(req.paper_limit), + "memory_limit": int(req.memory_limit), + "offline": bool(req.offline), + "sources": list(req.sources or []), + }, + ) return ContextResponse(context_pack=pack) + except Exception as exc: + metric_store.record_metric( + workflow="research_context", + stage=req.stage, + status="failed", + track_id=req.track_id, + elapsed_ms=(time.perf_counter() - started) * 1000.0, + detail={ + "error": str(exc), + "paper_limit": int(req.paper_limit), + "memory_limit": int(req.memory_limit), + "offline": bool(req.offline), + }, + ) + raise finally: await engine.close() diff --git a/src/paperbot/infrastructure/stores/models.py b/src/paperbot/infrastructure/stores/models.py index dbdb4cc8..502ba40d 100644 --- a/src/paperbot/infrastructure/stores/models.py +++ b/src/paperbot/infrastructure/stores/models.py @@ -347,6 +347,27 @@ class MemoryEvalMetricModel(Base): detail_json: Mapped[str] = mapped_column(Text, default="{}") +class WorkflowEvalMetricModel(Base): + """Evaluation metrics for research workflow observability and evidence coverage.""" + + __tablename__ = "workflow_eval_metrics" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + ts: Mapped[datetime] = mapped_column(DateTime(timezone=True), index=True) + + workflow: Mapped[str] = mapped_column(String(64), index=True) + stage: Mapped[str] = mapped_column(String(64), default="", index=True) + status: Mapped[str] = mapped_column(String(32), default="completed", index=True) + track_id: Mapped[Optional[int]] = mapped_column(Integer, nullable=True, index=True) + + claim_count: Mapped[int] = mapped_column(Integer, default=0) + evidence_count: Mapped[int] = mapped_column(Integer, default=0) + coverage_rate: Mapped[float] = mapped_column(Float, default=0.0, index=True) + elapsed_ms: Mapped[float] = mapped_column(Float, default=0.0) + + detail_json: Mapped[str] = mapped_column(Text, default="{}") + + class ResearchTrackModel(Base): """ User research direction / track. diff --git a/src/paperbot/infrastructure/stores/workflow_metric_store.py b/src/paperbot/infrastructure/stores/workflow_metric_store.py new file mode 100644 index 00000000..93d08f2c --- /dev/null +++ b/src/paperbot/infrastructure/stores/workflow_metric_store.py @@ -0,0 +1,237 @@ +from __future__ import annotations + +import json +from collections import defaultdict +from datetime import datetime, timedelta, timezone +from typing import Any, Dict, List, Optional + +from sqlalchemy import select + +from paperbot.infrastructure.stores.models import Base, WorkflowEvalMetricModel +from paperbot.infrastructure.stores.sqlalchemy_db import SessionProvider, get_db_url + + +def _utcnow() -> datetime: + return datetime.now(timezone.utc) + + +class WorkflowMetricStore: + """Persist and aggregate workflow quality metrics (coverage/latency/status).""" + + def __init__(self, db_url: Optional[str] = None, *, auto_create_schema: bool = True): + self.db_url = db_url or get_db_url() + self._provider = SessionProvider(self.db_url) + if auto_create_schema: + Base.metadata.create_all(self._provider.engine) + + def record_metric( + self, + *, + workflow: str, + stage: str = "", + status: str = "completed", + track_id: Optional[int] = None, + claim_count: int = 0, + evidence_count: int = 0, + elapsed_ms: float = 0.0, + detail: Optional[Dict[str, Any]] = None, + ) -> Dict[str, Any]: + claims = max(0, int(claim_count or 0)) + evidences = max(0, int(evidence_count or 0)) + coverage_rate = (float(evidences) / float(claims)) if claims > 0 else 0.0 + + row = WorkflowEvalMetricModel( + ts=_utcnow(), + workflow=(workflow or "unknown")[:64], + stage=(stage or "")[:64], + status=(status or "completed")[:32], + track_id=int(track_id) if track_id is not None else None, + claim_count=claims, + evidence_count=evidences, + coverage_rate=max(0.0, min(1.0, float(coverage_rate))), + elapsed_ms=max(0.0, float(elapsed_ms or 0.0)), + detail_json=json.dumps(detail or {}, ensure_ascii=False), + ) + + with self._provider.session() as session: + session.add(row) + session.commit() + session.refresh(row) + return self._to_dict(row) + + def summarize( + self, + *, + days: int = 7, + workflow: Optional[str] = None, + track_id: Optional[int] = None, + ) -> Dict[str, Any]: + window_days = max(1, min(int(days), 90)) + since = _utcnow() - timedelta(days=window_days) + + stmt = select(WorkflowEvalMetricModel).where(WorkflowEvalMetricModel.ts >= since) + if workflow: + stmt = stmt.where(WorkflowEvalMetricModel.workflow == str(workflow)[:64]) + if track_id is not None: + stmt = stmt.where(WorkflowEvalMetricModel.track_id == int(track_id)) + + with self._provider.session() as session: + rows = session.execute(stmt).scalars().all() + + totals = { + "runs": 0, + "success_runs": 0, + "failed_runs": 0, + "avg_elapsed_ms": 0.0, + "claim_count": 0, + "evidence_count": 0, + "coverage_rate": 0.0, + } + + by_day: Dict[str, Dict[str, Any]] = {} + by_workflow: Dict[str, Dict[str, Any]] = {} + by_track: Dict[str, Dict[str, Any]] = {} + failure_stages: Dict[str, int] = defaultdict(int) + + for row in rows: + totals["runs"] += 1 + if row.status == "completed": + totals["success_runs"] += 1 + else: + totals["failed_runs"] += 1 + failure_stages[row.stage or "unknown"] += 1 + + totals["avg_elapsed_ms"] += float(row.elapsed_ms or 0.0) + totals["claim_count"] += int(row.claim_count or 0) + totals["evidence_count"] += int(row.evidence_count or 0) + + date_key = (row.ts or _utcnow()).date().isoformat() + day = by_day.setdefault( + date_key, + { + "date": date_key, + "runs": 0, + "success_runs": 0, + "failed_runs": 0, + "claim_count": 0, + "evidence_count": 0, + "avg_elapsed_ms": 0.0, + "coverage_rate": 0.0, + }, + ) + self._add_to_bucket(day, row) + + wf = by_workflow.setdefault( + row.workflow or "unknown", + { + "workflow": row.workflow or "unknown", + "runs": 0, + "success_runs": 0, + "failed_runs": 0, + "claim_count": 0, + "evidence_count": 0, + "avg_elapsed_ms": 0.0, + "coverage_rate": 0.0, + }, + ) + self._add_to_bucket(wf, row) + + track_key = str(row.track_id) if row.track_id is not None else "none" + track_bucket = by_track.setdefault( + track_key, + { + "track_id": row.track_id, + "runs": 0, + "success_runs": 0, + "failed_runs": 0, + "claim_count": 0, + "evidence_count": 0, + "avg_elapsed_ms": 0.0, + "coverage_rate": 0.0, + }, + ) + self._add_to_bucket(track_bucket, row) + + if totals["runs"] > 0: + totals["avg_elapsed_ms"] = round(totals["avg_elapsed_ms"] / totals["runs"], 2) + totals["coverage_rate"] = self._coverage(totals["claim_count"], totals["evidence_count"]) + + return { + "window_days": window_days, + "totals": totals, + "by_day": [self._finalize_bucket(by_day[k]) for k in sorted(by_day.keys())], + "by_workflow": sorted( + (self._finalize_bucket(v) for v in by_workflow.values()), + key=lambda x: int(x.get("runs") or 0), + reverse=True, + ), + "by_track": sorted( + (self._finalize_bucket(v) for v in by_track.values()), + key=lambda x: int(x.get("runs") or 0), + reverse=True, + ), + "failure_stages": dict( + sorted(failure_stages.items(), key=lambda kv: kv[1], reverse=True) + ), + } + + def _add_to_bucket(self, bucket: Dict[str, Any], row: WorkflowEvalMetricModel) -> None: + bucket["runs"] = int(bucket.get("runs") or 0) + 1 + if row.status == "completed": + bucket["success_runs"] = int(bucket.get("success_runs") or 0) + 1 + else: + bucket["failed_runs"] = int(bucket.get("failed_runs") or 0) + 1 + bucket["claim_count"] = int(bucket.get("claim_count") or 0) + int(row.claim_count or 0) + bucket["evidence_count"] = int(bucket.get("evidence_count") or 0) + int( + row.evidence_count or 0 + ) + bucket["avg_elapsed_ms"] = float(bucket.get("avg_elapsed_ms") or 0.0) + float( + row.elapsed_ms or 0.0 + ) + + def _finalize_bucket(self, bucket: Dict[str, Any]) -> Dict[str, Any]: + runs = max(0, int(bucket.get("runs") or 0)) + if runs > 0: + bucket["avg_elapsed_ms"] = round(float(bucket.get("avg_elapsed_ms") or 0.0) / runs, 2) + else: + bucket["avg_elapsed_ms"] = 0.0 + bucket["coverage_rate"] = self._coverage( + int(bucket.get("claim_count") or 0), int(bucket.get("evidence_count") or 0) + ) + return bucket + + @staticmethod + def _coverage(claim_count: int, evidence_count: int) -> float: + claims = max(0, int(claim_count or 0)) + evidences = max(0, int(evidence_count or 0)) + if claims <= 0: + return 0.0 + return round(max(0.0, min(1.0, evidences / claims)), 4) + + @staticmethod + def _to_dict(row: WorkflowEvalMetricModel) -> Dict[str, Any]: + try: + detail = json.loads(row.detail_json or "{}") + if not isinstance(detail, dict): + detail = {} + except Exception: + detail = {} + return { + "id": int(row.id), + "ts": row.ts.isoformat() if row.ts else None, + "workflow": row.workflow, + "stage": row.stage, + "status": row.status, + "track_id": row.track_id, + "claim_count": int(row.claim_count or 0), + "evidence_count": int(row.evidence_count or 0), + "coverage_rate": float(row.coverage_rate or 0.0), + "elapsed_ms": float(row.elapsed_ms or 0.0), + "detail": detail, + } + + def close(self) -> None: + try: + self._provider.engine.dispose() + except Exception: + pass diff --git a/tests/unit/test_research_workflow_metrics_route.py b/tests/unit/test_research_workflow_metrics_route.py new file mode 100644 index 00000000..040e2514 --- /dev/null +++ b/tests/unit/test_research_workflow_metrics_route.py @@ -0,0 +1,36 @@ +from __future__ import annotations + +from pathlib import Path + +from fastapi.testclient import TestClient + +from paperbot.api import main as api_main +from paperbot.api.routes import research as research_route +from paperbot.infrastructure.stores.workflow_metric_store import WorkflowMetricStore + + +def test_research_workflow_metrics_summary_route(tmp_path: Path, monkeypatch): + db_url = f"sqlite:///{tmp_path / 'workflow-route.db'}" + store = WorkflowMetricStore(db_url=db_url) + store.record_metric( + workflow="research_context", + stage="auto", + status="completed", + track_id=3, + claim_count=12, + evidence_count=9, + elapsed_ms=222, + ) + + monkeypatch.setattr(research_route, "_workflow_metric_store", store) + + with TestClient(api_main.app) as client: + resp = client.get( + "/api/research/metrics/workflows?days=7&workflow=research_context&track_id=3" + ) + + assert resp.status_code == 200 + summary = resp.json()["summary"] + assert summary["totals"]["runs"] >= 1 + assert summary["totals"]["claim_count"] >= 12 + assert summary["totals"]["evidence_count"] >= 9 diff --git a/tests/unit/test_workflow_metric_store.py b/tests/unit/test_workflow_metric_store.py new file mode 100644 index 00000000..9d4b8896 --- /dev/null +++ b/tests/unit/test_workflow_metric_store.py @@ -0,0 +1,66 @@ +from __future__ import annotations + +from pathlib import Path + +from paperbot.infrastructure.stores.workflow_metric_store import WorkflowMetricStore + + +def test_workflow_metric_store_records_and_summarizes(tmp_path: Path): + db_url = f"sqlite:///{tmp_path / 'workflow-metrics.db'}" + store = WorkflowMetricStore(db_url=db_url) + + store.record_metric( + workflow="paperscool_daily", + stage="result", + status="completed", + claim_count=10, + evidence_count=8, + elapsed_ms=1200, + detail={"source": "test"}, + ) + store.record_metric( + workflow="research_context", + stage="auto", + status="failed", + claim_count=4, + evidence_count=2, + elapsed_ms=600, + detail={"error": "boom"}, + ) + + summary = store.summarize(days=7) + + assert summary["totals"]["runs"] == 2 + assert summary["totals"]["success_runs"] == 1 + assert summary["totals"]["failed_runs"] == 1 + assert summary["totals"]["claim_count"] == 14 + assert summary["totals"]["evidence_count"] == 10 + assert summary["totals"]["coverage_rate"] == round(10 / 14, 4) + assert len(summary["by_workflow"]) == 2 + + +def test_workflow_metric_store_filters_by_workflow_and_track(tmp_path: Path): + db_url = f"sqlite:///{tmp_path / 'workflow-filter.db'}" + store = WorkflowMetricStore(db_url=db_url) + + store.record_metric( + workflow="research_context", + stage="auto", + status="completed", + track_id=1, + claim_count=6, + evidence_count=6, + ) + store.record_metric( + workflow="research_context", + stage="auto", + status="completed", + track_id=2, + claim_count=5, + evidence_count=4, + ) + + scoped = store.summarize(days=7, workflow="research_context", track_id=1) + assert scoped["totals"]["runs"] == 1 + assert scoped["totals"]["claim_count"] == 6 + assert scoped["totals"]["coverage_rate"] == 1.0