Spaces:
Sleeping
Sleeping
| """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) | |