Spaces:
Sleeping
Sleeping
| """Control Room observability summary. | |
| Projects the live :class:`~app.observability.local_store.LocalObservationStore` | |
| -- the canonical, ordinary-product-traffic observation log built in Tasks 1-4 | |
| -- into the compact shape the Control Room endpoint returns under the | |
| ``observability`` key. This module only reads already-recorded observations | |
| and aggregates them; it never recomputes timings from ``ProjectActivity`` and | |
| it never emits a span, metric, log, or telemetry write of its own. | |
| Status vocabulary (checked in this priority order, matching the plan's | |
| "never silently convert missing data to zero" constraint): | |
| 1. ``not_instrumented`` -- telemetry itself is off (`bootstrap` mode is | |
| ``"disabled"``); there is no store to read at all. | |
| 2. ``not_observed`` -- telemetry is on but nothing has been recorded yet | |
| (store missing or empty). | |
| 3. ``degraded`` -- telemetry itself had trouble (one or more suppressed | |
| export/store failures were counted), even though some data did land. | |
| 4. ``observed`` -- real data present, nothing degraded. | |
| """ | |
| from __future__ import annotations | |
| from typing import Any | |
| from app.observability import bootstrap | |
| # Maps a recorded `operation` name (see call sites across app/rag, | |
| # app/websockets/handlers.py, app/agents/cerebras_client.py, and | |
| # app/services/student_memory.py) to the one canonical stage-duration name | |
| # from the plan preamble's "Operational measures" table. More than one | |
| # operation name can share a canonical stage (e.g. both "embedding.query" and | |
| # "embedding.batch" measure `embedding_ms`, and all three memory operations | |
| # measure `memory_ms`); their durations are pooled before computing | |
| # percentiles so the Control Room projection stays keyed by the one canonical | |
| # name, not by the internal call-site name. | |
| _OPERATION_TO_STAGE: dict[str, str] = { | |
| "pdf.page_extraction": "page_extraction_ms", | |
| "rag.chunking": "chunking_ms", | |
| "rag.ingestion": "ingestion_ms", | |
| "embedding.query": "embedding_ms", | |
| "embedding.batch": "embedding_ms", | |
| "chroma.collection_lookup": "collection_lookup_ms", | |
| "chroma.vector_search": "vector_search_ms", | |
| "rag.retrieval": "retrieval_ms", | |
| "context.assembly": "context_assembly_ms", | |
| "memory.push": "memory_ms", | |
| "memory.flush": "memory_ms", | |
| "memory.recall": "memory_ms", | |
| "llm.generate": "llm_ms", | |
| "citation.attach": "citation_attachment_ms", | |
| } | |
| def _percentile(sorted_values: list[float], fraction: float) -> float: | |
| """Nearest-rank-with-interpolation percentile over an already-sorted list.""" | |
| if len(sorted_values) == 1: | |
| return sorted_values[0] | |
| index = fraction * (len(sorted_values) - 1) | |
| lower = int(index) | |
| upper = min(lower + 1, len(sorted_values) - 1) | |
| weight = index - lower | |
| return sorted_values[lower] + (sorted_values[upper] - sorted_values[lower]) * weight | |
| def _stage_summary(rows: list[dict[str, Any]]) -> dict[str, dict[str, float | int]]: | |
| """p50/p95/count per canonical stage, over rows with a real duration only. | |
| A row whose `duration_ms` is `None` is skipped entirely for that stage -- | |
| never treated as a zero-duration sample. | |
| """ | |
| durations_by_stage: dict[str, list[float]] = {} | |
| for row in rows: | |
| stage = _OPERATION_TO_STAGE.get(row.get("operation")) | |
| if stage is None: | |
| continue | |
| duration = row.get("duration_ms") | |
| if duration is None: | |
| continue | |
| durations_by_stage.setdefault(stage, []).append(float(duration)) | |
| summary: dict[str, dict[str, float | int]] = {} | |
| for stage, durations in durations_by_stage.items(): | |
| ordered = sorted(durations) | |
| summary[stage] = { | |
| "p50": round(_percentile(ordered, 0.50), 3), | |
| "p95": round(_percentile(ordered, 0.95), 3), | |
| "count": len(ordered), | |
| } | |
| return summary | |
| def _retrieval_summary(rows: list[dict[str, Any]]) -> dict[str, int]: | |
| """success/empty/error counts from every row carrying a retrieval outcome. | |
| Only rows where the call site attached a `RetrievalOutcome` (currently | |
| `rag.retrieval`, plus the benchmark harness's `rag.benchmark_query`) carry | |
| a `retrieval` field; every other row is silently skipped, not counted as | |
| a zero. | |
| """ | |
| success = 0 | |
| empty = 0 | |
| error = 0 | |
| for row in rows: | |
| retrieval = row.get("retrieval") | |
| if not retrieval: | |
| continue | |
| status = retrieval.get("status") | |
| if status == "success": | |
| success += 1 | |
| elif status == "success_empty": | |
| empty += 1 | |
| elif status == "error_fallback": | |
| error += 1 | |
| return {"success_count": success, "empty_count": empty, "error_count": error} | |
| def _latest_trace_id(rows: list[dict[str, Any]]) -> str | None: | |
| for row in reversed(rows): | |
| trace_id = row.get("trace_id") | |
| if trace_id: | |
| return trace_id | |
| return None | |
| def _latest_pipeline_version(rows: list[dict[str, Any]]) -> str | None: | |
| for row in reversed(rows): | |
| experiment = row.get("experiment") | |
| if experiment and experiment.get("pipeline_version"): | |
| return experiment["pipeline_version"] | |
| return None | |
| def _empty_projection() -> dict[str, Any]: | |
| return { | |
| "latest_trace_id": None, | |
| "pipeline_version": None, | |
| "stages": {}, | |
| "retrieval": {"success_count": 0, "empty_count": 0, "error_count": 0}, | |
| } | |
| def observability_summary() -> dict[str, Any]: | |
| """Build the Control Room `observability` projection from live state. | |
| Always returns the same top-level shape regardless of status, so the | |
| frontend never has to branch on which keys exist: | |
| ```json | |
| { | |
| "status": "observed", | |
| "latest_trace_id": "...", | |
| "pipeline_version": "rag-naive-v1", | |
| "stages": {"retrieval_ms": {"p50": 1.0, "p95": 2.0, "count": 3}}, | |
| "retrieval": {"success_count": 1, "empty_count": 0, "error_count": 0}, | |
| "exporter": {"mode": "local", "suppressed_failures": 0}, | |
| "signoz_url": null | |
| } | |
| ``` | |
| """ | |
| state = bootstrap.get_state() | |
| suppressed_failures = bootstrap.get_suppressed_failure_count() | |
| exporter = {"mode": state.mode, "suppressed_failures": suppressed_failures} | |
| signoz_url = state.config.signoz_ui_url if state.mode == "full" else None | |
| if state.mode == "disabled": | |
| return { | |
| "status": "not_instrumented", | |
| "exporter": exporter, | |
| "signoz_url": signoz_url, | |
| **_empty_projection(), | |
| } | |
| store = state.local_store | |
| rows = store.all() if store is not None else [] | |
| if not rows: | |
| return { | |
| "status": "not_observed", | |
| "exporter": exporter, | |
| "signoz_url": signoz_url, | |
| **_empty_projection(), | |
| } | |
| status = "degraded" if suppressed_failures > 0 else "observed" | |
| return { | |
| "status": status, | |
| "latest_trace_id": _latest_trace_id(rows), | |
| "pipeline_version": _latest_pipeline_version(rows), | |
| "stages": _stage_summary(rows), | |
| "retrieval": _retrieval_summary(rows), | |
| "exporter": exporter, | |
| "signoz_url": signoz_url, | |
| } | |