Skip to content

Commit f33df23

Browse files
committed
fix(storage): harden bloat diagnostic evidence
1 parent 6f243a7 commit f33df23

4 files changed

Lines changed: 443 additions & 24 deletions

File tree

routes/diagnostics_routes.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
"""Diagnostics routes — /api/db/stats, /api/rag/stats, /api/test/youtube, /api/test-research."""
22

3+
import asyncio
34
import logging
45
import os
56
from typing import Dict, Any
@@ -58,7 +59,8 @@ async def get_diagnostics_logs(request: Request, limit: int = 200) -> Dict[str,
5859
async def get_storage_bloat_diagnostics(request: Request) -> Dict[str, Any]:
5960
require_admin(request)
6061
try:
61-
return collect_configured_storage_bloat_diagnostics(
62+
return await asyncio.to_thread(
63+
collect_configured_storage_bloat_diagnostics,
6264
default_db_path=APP_DB,
6365
upload_dir=UPLOAD_DIR,
6466
)

src/storage_diagnostics.py

Lines changed: 60 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,28 @@
4949
"extra",
5050
"attachments",
5151
)
52+
REFERENCE_COVERED_SOURCES = (
53+
"chat_messages.content",
54+
"chat_messages.attachment_metadata",
55+
"documents.current_content.pdf_source_marker",
56+
"document_versions.content.pdf_source_marker",
57+
)
58+
REFERENCE_UNCOVERED_SOURCES = (
59+
"gallery_images.filename",
60+
"gallery_images.file_hash",
61+
"notes.image_url",
62+
"notes.color",
63+
"notes.content",
64+
"notes.items",
65+
"calendars.color",
66+
"calendar_events.color",
67+
"calendar_events.description",
68+
"calendar_events.location",
69+
)
70+
DURABLE_REFERENCE_SOURCES_INCOMPLETE = "durable_reference_sources_incomplete"
71+
INCOMPLETE_REFERENCE_SAMPLE_SCOPE = (
72+
"candidates_observed_within_incomplete_reference_scope"
73+
)
5274
REFERENCE_CONTAINER_KEYS = {
5375
"attachment",
5476
"attachments",
@@ -96,7 +118,8 @@ def collect_storage_bloat_diagnostics(
96118
upload_report = _collect_upload_report(
97119
Path(upload_dir or UPLOAD_DIR),
98120
live_upload_ids,
99-
db_report["upload_references"]["complete"],
121+
db_report["upload_references"]["attempted_scan_complete"],
122+
db_report["upload_references"]["reference_coverage"]["complete"],
100123
warnings,
101124
)
102125
return {
@@ -201,6 +224,15 @@ def _connect_readonly(db_path: Path) -> sqlite3.Connection:
201224
return conn
202225

203226

227+
def _durable_reference_coverage_report() -> dict[str, Any]:
228+
"""Describe implemented reference scans without exposing persisted values."""
229+
return {
230+
"complete": not REFERENCE_UNCOVERED_SOURCES,
231+
"covered_sources": list(REFERENCE_COVERED_SOURCES),
232+
"uncovered_sources": list(REFERENCE_UNCOVERED_SOURCES),
233+
}
234+
235+
204236
def _collect_db_report(
205237
db_path: Path | None,
206238
warnings: list[str],
@@ -221,7 +253,9 @@ def _collect_db_report(
221253
"scan_limit_scope": "per_source",
222254
"value_char_scan_limit": MAX_REFERENCE_CONTENT_CHARS,
223255
"json_node_scan_limit": MAX_REFERENCE_JSON_NODES,
256+
"attempted_scan_complete": False,
224257
"complete": False,
258+
"reference_coverage": _durable_reference_coverage_report(),
225259
},
226260
}
227261
live_upload_ids: set[str] = set()
@@ -491,7 +525,9 @@ def _live_upload_reference_report(
491525
"scan_limit_scope": "per_source",
492526
"value_char_scan_limit": MAX_REFERENCE_CONTENT_CHARS,
493527
"json_node_scan_limit": MAX_REFERENCE_JSON_NODES,
528+
"attempted_scan_complete": False,
494529
"complete": False,
530+
"reference_coverage": _durable_reference_coverage_report(),
495531
}
496532
live_ids: set[str] = set()
497533
scans_complete: list[bool] = []
@@ -575,7 +611,11 @@ def _live_upload_reference_report(
575611
)
576612

577613
report["live_reference_count"] = len(live_ids)
578-
report["complete"] = bool(scans_complete) and all(scans_complete)
614+
attempted_scan_complete = bool(scans_complete) and all(scans_complete)
615+
report["attempted_scan_complete"] = attempted_scan_complete
616+
report["complete"] = (
617+
attempted_scan_complete and report["reference_coverage"]["complete"]
618+
)
579619
return report, live_ids
580620

581621

@@ -788,6 +828,7 @@ def _collect_upload_report(
788828
upload_dir: Path,
789829
live_upload_ids: set[str],
790830
reference_scan_complete: bool,
831+
durable_reference_coverage_complete: bool,
791832
warnings: list[str],
792833
) -> dict[str, Any]:
793834
report: dict[str, Any] = {
@@ -813,6 +854,9 @@ def _collect_upload_report(
813854
observed_count=0,
814855
sample=[],
815856
reference_scan_complete=reference_scan_complete,
857+
durable_reference_coverage_complete=(
858+
durable_reference_coverage_complete
859+
),
816860
upload_traversal_complete=False,
817861
),
818862
}
@@ -831,6 +875,7 @@ def _collect_upload_report(
831875
manifest_ids,
832876
live_upload_ids,
833877
reference_scan_complete,
878+
durable_reference_coverage_complete,
834879
warnings,
835880
)
836881
report.update(file_report)
@@ -845,6 +890,7 @@ def _upload_file_report(
845890
manifest_ids: set[str],
846891
live_upload_ids: set[str],
847892
reference_scan_complete: bool,
893+
durable_reference_coverage_complete: bool,
848894
warnings: list[str],
849895
) -> dict[str, Any]:
850896
state = UploadTraversalState(visited_dirs=1)
@@ -861,7 +907,7 @@ def _upload_file_report(
861907
orphan_count += 1
862908
if len(orphan_sample) < MAX_ORPHAN_ROWS:
863909
orphan_sample.append(
864-
_upload_sample(record, upload_dir, manifest_ids, live_upload_ids)
910+
_upload_sample(record, upload_dir, manifest_ids)
865911
)
866912

867913
traversal_unreadable = state.skipped_unreadable_count > 0
@@ -884,6 +930,9 @@ def _upload_file_report(
884930
observed_count=orphan_count,
885931
sample=orphan_sample,
886932
reference_scan_complete=reference_scan_complete,
933+
durable_reference_coverage_complete=(
934+
durable_reference_coverage_complete
935+
),
887936
upload_traversal_complete=not traversal_incomplete,
888937
upload_traversal_truncated=state.truncated,
889938
upload_traversal_unreadable=traversal_unreadable,
@@ -896,13 +945,16 @@ def _suspected_orphan_report(
896945
observed_count: int,
897946
sample: list[dict[str, Any]],
898947
reference_scan_complete: bool,
948+
durable_reference_coverage_complete: bool,
899949
upload_traversal_complete: bool,
900950
upload_traversal_truncated: bool = False,
901951
upload_traversal_unreadable: bool = False,
902952
) -> dict[str, Any]:
903953
reasons: list[str] = []
904954
if not reference_scan_complete:
905955
reasons.append("db_reference_scan_incomplete")
956+
if not durable_reference_coverage_complete:
957+
reasons.append(DURABLE_REFERENCE_SOURCES_INCOMPLETE)
906958
if upload_traversal_truncated or (
907959
not upload_traversal_complete and not upload_traversal_unreadable
908960
):
@@ -917,6 +969,11 @@ def _suspected_orphan_report(
917969
"count": observed_count if complete else None,
918970
"observed_count": observed_count,
919971
"reason_incomplete": reasons or None,
972+
"sample_scope": (
973+
"suspected_orphan_candidates"
974+
if complete
975+
else INCOMPLETE_REFERENCE_SAMPLE_SCOPE
976+
),
920977
"sample": sample,
921978
}
922979

@@ -1024,14 +1081,12 @@ def _upload_sample(
10241081
record: UploadFileRecord,
10251082
upload_dir: Path,
10261083
manifest_ids: set[str],
1027-
live_upload_ids: set[str],
10281084
) -> dict[str, Any]:
10291085
return {
10301086
"path_hash": _hash_upload_identifier(_relative_upload_path(upload_dir, record.path)),
10311087
"relative_dir_bucket": _safe_relative_dir_bucket(upload_dir, record.path),
10321088
"size_bytes": record.size_bytes,
10331089
"manifest_present": _upload_file_in_manifest(record.path, manifest_ids),
1034-
"live_db_reference_found": _upload_file_is_live(record.path, live_upload_ids),
10351090
}
10361091

10371092

tests/test_diagnostics_service_route.py

Lines changed: 96 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
"""Route-level regression tests for GET /api/diagnostics/services.
1+
"""Route-level regression tests for admin diagnostics endpoints.
22
33
The reviewer asked for explicit coverage of unauthenticated / non-admin / admin
44
access to this admin diagnostics route, beyond the unit tests for the collector.
@@ -8,6 +8,9 @@
88
fail, so the suite stays green in minimal environments; CI installs
99
requirements, so the tests run there.
1010
"""
11+
import asyncio
12+
import threading
13+
1114
import pytest
1215

1316
fastapi = pytest.importorskip("fastapi")
@@ -66,3 +69,95 @@ def gate(_request: Request):
6669
body = r.json()
6770
assert set(body) == {"overall", "services", "timestamp"}
6871
assert body["overall"] == "ok"
72+
73+
74+
def _storage_bloat_endpoint():
75+
router = diag.setup_diagnostics_routes(
76+
rag_manager=None,
77+
rag_available=False,
78+
research_handler=None,
79+
memory_vector=None,
80+
)
81+
return next(
82+
route.endpoint
83+
for route in router.routes
84+
if route.path == "/api/diagnostics/storage-bloat"
85+
)
86+
87+
88+
@pytest.mark.asyncio
89+
async def test_storage_bloat_collection_leaves_event_loop_responsive(monkeypatch):
90+
loop = asyncio.get_running_loop()
91+
event_loop_thread = threading.get_ident()
92+
collector_started = asyncio.Event()
93+
collector_release = threading.Event()
94+
collector_threads = []
95+
expected = {
96+
"status": "success",
97+
"database": {},
98+
"uploads": {},
99+
"warnings": [],
100+
}
101+
102+
def blocked_collector(*, default_db_path, upload_dir):
103+
assert default_db_path == diag.APP_DB
104+
assert upload_dir == diag.UPLOAD_DIR
105+
collector_threads.append(threading.get_ident())
106+
loop.call_soon_threadsafe(collector_started.set)
107+
if not collector_release.wait(timeout=2):
108+
raise AssertionError("collector was not released by the event loop")
109+
return expected
110+
111+
monkeypatch.setattr(diag, "require_admin", lambda _request: None)
112+
monkeypatch.setattr(
113+
diag,
114+
"collect_configured_storage_bloat_diagnostics",
115+
blocked_collector,
116+
)
117+
endpoint = _storage_bloat_endpoint()
118+
route_task = asyncio.create_task(endpoint(object()))
119+
120+
try:
121+
await asyncio.wait_for(collector_started.wait(), timeout=2)
122+
heartbeat = asyncio.Event()
123+
loop.call_soon(heartbeat.set)
124+
await asyncio.wait_for(heartbeat.wait(), timeout=2)
125+
126+
assert route_task.done() is False
127+
assert len(collector_threads) == 1
128+
assert collector_threads[0] != event_loop_thread
129+
finally:
130+
collector_release.set()
131+
132+
assert await asyncio.wait_for(route_task, timeout=2) == expected
133+
134+
135+
def test_storage_bloat_collector_error_returns_generic_500(monkeypatch):
136+
def failing_collector(*, default_db_path, upload_dir):
137+
raise RuntimeError("private database path and credentials")
138+
139+
monkeypatch.setattr(diag, "require_admin", lambda _request: None)
140+
monkeypatch.setattr(
141+
diag,
142+
"collect_configured_storage_bloat_diagnostics",
143+
failing_collector,
144+
)
145+
app = FastAPI()
146+
app.include_router(
147+
diag.setup_diagnostics_routes(
148+
rag_manager=None,
149+
rag_available=False,
150+
research_handler=None,
151+
memory_vector=None,
152+
)
153+
)
154+
client = TestClient(app, raise_server_exceptions=False)
155+
156+
response = client.get("/api/diagnostics/storage-bloat")
157+
158+
assert response.status_code == 500
159+
assert response.json() == {
160+
"detail": "Failed to retrieve storage bloat diagnostics"
161+
}
162+
assert "private database path" not in response.text
163+
assert "credentials" not in response.text

0 commit comments

Comments
 (0)