Cognitive-rag / backend / run_report.py
run_report.py
Raw
from __future__ import annotations

import json
import re
import uuid
from datetime import datetime, timezone
from pathlib import Path
from typing import Any

from backend import config


INVALID_FOLDER_CHARS = r'[\\/:*?"<>|\)]'


def ensure_run_report_dirs() -> None:
    config.RUN_REPORT_DIR.mkdir(parents=True, exist_ok=True)
    config.RUN_REPORT_UPLOAD_DIR.mkdir(parents=True, exist_ok=True)
    config.RUN_REPORT_PROMPT_DIR.mkdir(parents=True, exist_ok=True)


def sanitize_folder_name(name: str) -> str:
    cleaned = re.sub(INVALID_FOLDER_CHARS, "_", (name or "").strip())
    cleaned = cleaned.replace("\r", "_").replace("\n", "_")
    cleaned = re.sub(r"\s+", "_", cleaned)
    cleaned = re.sub(r"_+", "_", cleaned).strip("._ ")
    return cleaned or "file"


def _norm_path(value: str | None) -> str:
    if not value:
        return ""
    return value.replace("\\", "/").strip("/")


def _record_identity(record: dict[str, Any]) -> tuple[str, str]:
    source = record.get("source", {})
    if not isinstance(source, dict):
        source = {}
    source_path = _norm_path(source.get("path"))
    filename = (source.get("filename") or record.get("document_name") or "").strip()
    return source_path, filename


def _matches_upload(record: dict[str, Any], upload_rel: str) -> bool:
    source_path, filename = _record_identity(record)
    upload_rel_norm = _norm_path(upload_rel)
    upload_basename = Path(upload_rel_norm).name if upload_rel_norm else ""
    upload_has_dir = "/" in upload_rel_norm

    if upload_has_dir:
        return bool(source_path) and source_path == upload_rel_norm

    if source_path and Path(source_path).name == upload_basename:
        return True
    if filename and filename == upload_basename:
        return True
    return False


def _append_jsonl_line(path: Path, line: str) -> None:
    with path.open("a", encoding="utf-8") as handle:
        handle.write(line if line.endswith("\n") else f"{line}\n")


def write_upload_reports(
    *,
    job_id: str,
    uploaded_files: list[str],
    ingestion_jsonl: Path,
    chunking_jsonl: Path,
    embedding_jsonl: Path,
) -> dict[str, Any]:
    """
    Copies per-file stage artifacts into backend/run_report/upload/<sanitized_filename>/.
    """
    ensure_run_report_dirs()

    targets: dict[str, dict[str, Path]] = {}
    for upload_rel in uploaded_files:
        file_name = Path(_norm_path(upload_rel)).name or upload_rel
        target_dir = config.RUN_REPORT_UPLOAD_DIR / sanitize_folder_name(file_name)
        target_dir.mkdir(parents=True, exist_ok=True)
        targets[upload_rel] = {
            "dir": target_dir,
            "ingestion": target_dir / "ingestion.jsonl",
            "chunking": target_dir / "chunking.jsonl",
            "embedding": target_dir / "embedding.jsonl",
        }
        for stage_file in ("ingestion", "chunking", "embedding"):
            targets[upload_rel][stage_file].write_text("", encoding="utf-8")

    def process_stage(stage_name: str, source_path: Path) -> dict[str, int]:
        counts = {upload_rel: 0 for upload_rel in uploaded_files}
        if not source_path.exists():
            return counts
        with source_path.open("r", encoding="utf-8") as handle:
            for line in handle:
                raw = line.strip()
                if not raw:
                    continue
                try:
                    record = json.loads(raw)
                except json.JSONDecodeError:
                    continue
                for upload_rel in uploaded_files:
                    if _matches_upload(record, upload_rel):
                        _append_jsonl_line(targets[upload_rel][stage_name], line)
                        counts[upload_rel] += 1
        return counts

    ingestion_counts = process_stage("ingestion", ingestion_jsonl)
    chunking_counts = process_stage("chunking", chunking_jsonl)
    embedding_counts = process_stage("embedding", embedding_jsonl)

    summary: dict[str, Any] = {"job_id": job_id, "files": []}
    for upload_rel in uploaded_files:
        summary["files"].append(
            {
                "upload_rel": upload_rel,
                "report_dir": str(targets[upload_rel]["dir"]),
                "ingestion_records": ingestion_counts.get(upload_rel, 0),
                "chunking_records": chunking_counts.get(upload_rel, 0),
                "embedding_records": embedding_counts.get(upload_rel, 0),
            }
        )
    return summary


def _prompt_report_dir(query: str) -> Path:
    ensure_run_report_dirs()
    base_name = sanitize_folder_name(query)[:120] or "prompt"
    folder = config.RUN_REPORT_PROMPT_DIR / base_name
    if not folder.exists():
        folder.mkdir(parents=True, exist_ok=True)
        return folder

    # Same question asked multiple times: keep question-based name with numeric suffix.
    i = 2
    while True:
        candidate = config.RUN_REPORT_PROMPT_DIR / f"{base_name}__{i}"
        if not candidate.exists():
            candidate.mkdir(parents=True, exist_ok=True)
            return candidate
        i += 1


def write_prompt_reports(
    *,
    query: str,
    strategy: str,
    context_sent_to_generator: list[dict[str, Any]],
    retrieved_topk: list[dict[str, Any]],
    dense_topk: list[dict[str, Any]],
    sparse_topk: list[dict[str, Any]],
    related_chunks: list[dict[str, Any]],
    memory_chunks: list[dict[str, Any]],
    memory_raw_count: int = 0,
    memory_dropped_count: int = 0,
    answer: str,
    used_chunks: list[dict[str, Any]],
    evidence_block: list[str] | None = None,
    original_query_retrieval: dict[str, Any] | None = None,
    multi_hop_retrievals: list[dict[str, Any]] | None = None,
    multi_hop_triggered: bool = False,
    multi_hop_reason: str | None = None,
    memory_write_debug: dict[str, Any] | None = None,
    memory_retrieval_filter: dict[str, Any] | None = None,
) -> dict[str, str]:
    """
    Writes retrieval_report.json and generation_report.json under backend/run_report/prompt/<id>/.
    """
    folder = _prompt_report_dir(query)

    known_score_keys = (
        "score",
        "dense_score",
        "sparse_score",
        "rrf_score",
        "cross_encoder_score",
        "original_score",
    )

    def compact_chunk(chunk: dict[str, Any]) -> dict[str, Any]:
        item: dict[str, Any] = {}
        for key in (
            "chunk_id",
            "doc_id",
            "text",
            "modality",
            "page",
            "section",
            "timestamp",
            "source_path",
            "retrieval_channels",
        ):
            value = chunk.get(key)
            if value is not None and value != "":
                item[key] = value
        for key in known_score_keys:
            value = chunk.get(key)
            if value is not None and value != "":
                item[key] = value
        return item

    def build_retrieval_section(section: dict[str, Any]) -> dict[str, Any]:
        context_bucket = [compact_chunk(chunk) for chunk in section.get("context_sent_to_generator", [])]
        retrieved_bucket = [compact_chunk(chunk) for chunk in section.get("retrieved_topk", [])]
        dense_bucket = [compact_chunk(chunk) for chunk in section.get("dense_topk", [])]
        sparse_bucket = [compact_chunk(chunk) for chunk in section.get("sparse_topk", [])]
        related_bucket = [compact_chunk(chunk) for chunk in section.get("related_chunks", [])]
        memory_bucket = [compact_chunk(chunk) for chunk in section.get("memory_chunks", [])]

        score_keys_available: list[str] = []
        seen_scores: set[str] = set()
        for bucket in (
            context_bucket,
            retrieved_bucket,
            dense_bucket,
            sparse_bucket,
            related_bucket,
            memory_bucket,
        ):
            for chunk in bucket:
                for key in known_score_keys:
                    if key in chunk and key not in seen_scores:
                        seen_scores.add(key)
                        score_keys_available.append(key)

        return {
            "query": section.get("query", ""),
            "strategy": section.get("strategy", "hybrid"),
            "context_sent_to_generator": context_bucket,
            "retrieved_topk": retrieved_bucket,
            "dense_topk": dense_bucket,
            "sparse_topk": sparse_bucket,
            "related_chunks": related_bucket,
            "memory_chunks": memory_bucket,
            "counts": {
                "context_sent_to_generator": len(context_bucket),
                "retrieved_topk": len(retrieved_bucket),
                "dense_topk": len(dense_bucket),
                "sparse_topk": len(sparse_bucket),
                "related_chunks": len(related_bucket),
                "memory_chunks": len(memory_bucket),
                "memory_raw_count": int(section.get("memory_raw_count", 0)),
                "memory_dropped_count": int(section.get("memory_dropped_count", 0)),
            },
            "score_schema": score_keys_available,
        }

    first_section_input = original_query_retrieval or {
        "query": query,
        "strategy": strategy,
        "context_sent_to_generator": context_sent_to_generator,
        "retrieved_topk": retrieved_topk,
        "dense_topk": dense_topk,
        "sparse_topk": sparse_topk,
        "related_chunks": related_chunks,
        "memory_chunks": memory_chunks,
        "memory_raw_count": memory_raw_count,
        "memory_dropped_count": memory_dropped_count,
    }
    first_section = build_retrieval_section(first_section_input)
    context_bucket = [compact_chunk(chunk) for chunk in context_sent_to_generator]
    used_bucket = [compact_chunk(chunk) for chunk in used_chunks]

    retrieval_payload: dict[str, Any] = {
        "original_query": query,
        "original_query_retrieval": first_section,
        "context_sent_to_generator": context_bucket,
        "multi_hop_triggered": bool(multi_hop_triggered),
    }
    if multi_hop_triggered and multi_hop_reason:
        retrieval_payload["multi_hop_reason"] = multi_hop_reason

    if multi_hop_retrievals:
        if len(multi_hop_retrievals) == 1:
            section = multi_hop_retrievals[0]
            retrieval_payload["multi_hop_query"] = section.get("query", "")
            retrieval_payload["multi_hop_retrieval"] = build_retrieval_section(section)
        else:
            for idx, section in enumerate(multi_hop_retrievals, 1):
                retrieval_payload[f"multi_hop_query_{idx}"] = section.get("query", "")
                retrieval_payload[f"multi_hop_retrieval_{idx}"] = build_retrieval_section(section)
        retrieval_payload["multi_hop_retrievals"] = [
            build_retrieval_section(section) for section in multi_hop_retrievals
        ]
    if memory_write_debug is not None:
        retrieval_payload["memory_write"] = memory_write_debug
    if memory_retrieval_filter is not None:
        retrieval_payload["memory_retrieval_filter"] = memory_retrieval_filter
    generation_payload: dict[str, Any] = {
        "query": query,
        "answer": answer,
        "used_chunks": used_bucket,
    }
    if evidence_block:
        generation_payload["evidence_block"] = evidence_block

    retrieval_file = folder / "retrieval_report.json"
    generation_file = folder / "generation_report.json"
    retrieval_file.write_text(
        json.dumps(retrieval_payload, ensure_ascii=True, indent=2),
        encoding="utf-8",
    )
    generation_file.write_text(
        json.dumps(generation_payload, ensure_ascii=True, indent=2),
        encoding="utf-8",
    )
    return {
        "folder": str(folder),
        "retrieval_report": str(retrieval_file),
        "generation_report": str(generation_file),
    }