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
121 changes: 120 additions & 1 deletion src/paperbot/api/routes/paperscool.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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],
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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] = {}
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The calculation of claims and evidences is also performed earlier in the if req.require_approval: block (line 649). To avoid code duplication and improve maintainability, you could calculate these values once after the report is finalized (e.g., after the enrichment phase around line 645) and before the approval check, then reuse the variables in both places.

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)


Expand All @@ -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
Expand All @@ -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(
Expand Down Expand Up @@ -842,13 +930,22 @@ 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",
checkpoint="result",
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 {})

Expand Down Expand Up @@ -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 {})

Expand Down Expand Up @@ -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()

Expand Down Expand Up @@ -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})


Expand Down
62 changes: 62 additions & 0 deletions src/paperbot/api/routes/research.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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()

Expand Down
21 changes: 21 additions & 0 deletions src/paperbot/infrastructure/stores/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading