study-buddy / app /observability /retrieval.py
GitHub Actions
deploy d092bea3608b7a29952f16357fda39b7a29e399b
2e818da
Raw
History Blame Contribute Delete
6.99 kB
"""Retrieval-outcome classification and stage/candidate observation helpers.
Shared by the instrumented RAG call sites (ingestion, ``ChromaDBClient``,
``handlers._get_chunks``) so instrumentation logic stays out of the product
modules and out of ``handlers.py``. Per the plan's file-ownership table this
module owns "retrieval outcome classification and stage/candidate
observations".
It emits nothing itself: no span, metric, log, or file write happens here. It
only classifies exceptions into stable, low-cardinality categories and computes
content-free count/distribution attributes (never chunk text, never document
IDs as measures). Callers hand the returned dicts to
``Operation.set(...)`` / ``op.add_count(...)`` which re-sanitize before export.
"""
from __future__ import annotations
import statistics
from typing import Any
# Stable, low-cardinality retrieval-error categories. Callers pass the stage
# that failed; these are the vocabulary the "retrieval_error" logs/attributes
# use so a dashboard can group by cause without unbounded cardinality.
COLLECTION_UNAVAILABLE = "collection_unavailable"
COLLECTION_COUNT_FAILED = "collection_count_failed"
VECTOR_SEARCH_FAILED = "vector_search_failed"
EMBEDDING_FAILED = "embedding_failed"
RESULT_PREPARE_FAILED = "result_prepare_failed"
UNKNOWN_ERROR = "unknown_error"
def error_category(exc: BaseException) -> str:
"""A stable, content-free category for an exception (its class name).
Never includes the exception message, which can embed a filesystem path or
document text. The operation boundary's sanitizer bounds it further.
"""
return type(exc).__name__
def chunk_char_stats(chunks: list[str]) -> dict[str, int]:
"""Count + character-size distribution of a chunk list (no text).
Returns only integer counts/lengths -- min/median/max/total characters and
the chunk count -- so a dashboard can see chunk-size distribution without
any chunk content ever leaving the process.
"""
lengths = [len(c) for c in chunks]
if not lengths:
return {"chunk_count": 0}
return {
"chunk_count": len(lengths),
"chunk_chars_total": sum(lengths),
"chunk_chars_min": min(lengths),
"chunk_chars_median": int(statistics.median(lengths)),
"chunk_chars_max": max(lengths),
}
def candidate_stats(candidates: list[dict[str, Any]]) -> dict[str, int | float]:
"""Content-free retrieval-behavior counts from a candidate dict list.
Emits ``raw_candidate_count`` plus document/section *diversity* (distinct
counts only -- never the ids/labels themselves) and
``duplicate_candidate_ratio`` (a fraction, never chunk text), matching the
canonical "Retrieval-behavior measures" contract. Applies identically
regardless of which pipeline version produced ``candidates`` -- it is a
pure function of the list handed to it, so calling it on a naive top-5
and on an over-fetched-then-deduplicated top-5 uses the exact same
yardstick, which is what makes a v1/v2 comparison meaningful.
"""
docs = {str(c.get("document_id")) for c in candidates if c.get("document_id")}
sources = {str(c.get("source")) for c in candidates if c.get("source")}
return {
"raw_candidate_count": len(candidates),
"document_diversity": len(docs),
"section_diversity": len(sources),
"duplicate_candidate_ratio": duplicate_candidate_ratio(candidates),
}
# --- Near-duplicate detection -------------------------------------------------
#
# Deterministic and content-free *in what it emits*: these helpers read chunk
# text internally (this is product code operating on already-retrieved
# in-process candidates, not telemetry export -- see module docstring), but
# every value that ever reaches `Operation.set(...)`/a span/a metric is a
# count or a ratio, never the text itself.
#
# Similarity definition: whitespace/case-normalized token-set Jaccard
# similarity. Chosen over an embedding-distance or edit-distance definition
# because it directly matches the failure mode Task 7 diagnosed from real
# SigNoz data -- the chunker's 64-character overlap produces adjacent chunks
# that share most of their words verbatim (near-identical, not paraphrased),
# and Jaccard on tokens is deterministic, has no model/network dependency,
# and is trivially the same computation used to both *measure* duplication
# (`duplicate_candidate_ratio`) and *act on it* (the dedup mechanism in
# `ChromaDBClient.query_observed`) -- one shared definition of "duplicate",
# not two that could quietly drift apart.
NEAR_DUPLICATE_JACCARD_THRESHOLD = 0.8
def _normalize_text(text: str) -> str:
return " ".join(text.lower().split())
def _token_set(text: str) -> set[str]:
return set(_normalize_text(text).split())
def _jaccard(a: set[str], b: set[str]) -> float:
if not a and not b:
return 1.0
if not a or not b:
return 0.0
union = len(a | b)
return (len(a & b) / union) if union else 0.0
def is_near_duplicate(
text_a: str, text_b: str, *, threshold: float = NEAR_DUPLICATE_JACCARD_THRESHOLD
) -> bool:
"""True if two chunk texts are exact or near-duplicate of each other.
Exact matches always count (fast path, and correct even for empty
strings, which have Jaccard similarity 1.0 by the ``_jaccard`` convention
above but are a degenerate case worth short-circuiting explicitly).
"""
if text_a == text_b:
return True
return _jaccard(_token_set(text_a), _token_set(text_b)) >= threshold
def find_duplicate_indices(candidates: list[dict[str, Any]]) -> set[int]:
"""Indices of candidates that are a near-duplicate of some *earlier*
candidate in the same list.
First-occurrence-wins: within a group of mutually near-duplicate
candidates, the earliest (highest-ranked) one is kept and every later one
in that group is flagged. Order-dependent by design -- both the
diversity/duplication *measurement* (`candidate_stats`) and the dedup
*mechanism* (`ChromaDBClient.query_observed`) receive candidates in
Chroma's own relevance-ranked order, so "keep the earliest" means "keep
the most relevant member of each near-duplicate group."
"""
kept_texts: list[str] = []
dup_indices: set[int] = set()
for i, c in enumerate(candidates):
text = str(c.get("text", ""))
if any(is_near_duplicate(text, kept) for kept in kept_texts):
dup_indices.add(i)
else:
kept_texts.append(text)
return dup_indices
def duplicate_candidate_ratio(candidates: list[dict[str, Any]]) -> float:
"""Fraction of ``candidates`` that are a near-duplicate of an earlier one.
A pure number (0.0-1.0), never chunk text -- safe to attach to a span or
metric under the sanitization rules this module's docstring describes.
"""
if not candidates:
return 0.0
return len(find_duplicate_indices(candidates)) / len(candidates)