Make run_db_checks an orchestration list

At 327 lines it interleaved reading the tables, deriving the lookups every
check needs, and the bodies of ten checks. Six checks were already
functions; the rest were inline, so none of them could be read or tested
without the others around them.

Each one is now a function taking exactly what it needs:
_check_document_meta_parity, _check_orphaned_chunks, _check_orphaned_items,
_check_documents_without_items, _check_dangling_item_refs,
_check_vector_dimension, _check_unembedded_chunks, _check_picture_data,
_check_settings_row and _check_pending_migrations. run_db_checks reads the
tables, then appends results.

_document_centroids takes the vector reduction. Passing the matrix as a
parameter keeps it a local of run_db_checks, so the del before clustering
still drops the last reference — measured at 6.2 MB allocated to reduce a
102 MB matrix, no second copy. There is no snapshot object: one holding
vectors would keep the largest allocation alive past the del.

The reduction also rebound doc_ids from the document-id set to the centroid
id list halfway through the function. The centroid ids have their own name
now.

No test changes: the 80 doctor tests cover these through run_db_checks and
pass unchanged.
This commit is contained in:
Yiorgis Gozadinos 2026-08-20 14:22:29 +03:00
parent d429ac4996
commit 0592206f16
No known key found for this signature in database

View file

@ -433,6 +433,227 @@ def _check_duplicate_documents(
)
def _check_document_meta_parity(
doc_ids: set[str], meta_doc_ids: set[str]
) -> CheckResult:
"""documents <-> document_meta must be 1:1."""
orphan_docs = doc_ids - meta_doc_ids
orphan_meta = meta_doc_ids - doc_ids
if not (orphan_docs or orphan_meta):
return CheckResult(
name="document_meta_parity",
severity=Severity.OK,
message="documents and document_meta are consistent.",
)
details = [f"document with no meta: {d}" for d in _sample(sorted(orphan_docs))]
details += [f"meta with no document: {d}" for d in _sample(sorted(orphan_meta))]
return CheckResult(
name="document_meta_parity",
severity=Severity.FAIL,
message="documents and document_meta are out of sync.",
remediation="haiku-rag rebuild",
details=details,
)
def _check_orphaned_chunks(chunk_doc_ids: set[str], doc_ids: set[str]) -> CheckResult:
"""Chunks referencing a document that no longer exists."""
orphans = chunk_doc_ids - doc_ids
return CheckResult(
name="orphaned_chunks",
severity=Severity.FAIL if orphans else Severity.OK,
message=(
"Chunks reference missing documents." if orphans else "No orphaned chunks."
),
remediation="haiku-rag rebuild" if orphans else None,
details=_sample(sorted(orphans)),
)
def _check_orphaned_items(item_doc_ids: set[str], doc_ids: set[str]) -> CheckResult:
"""Document items referencing a document that no longer exists."""
orphans = item_doc_ids - doc_ids
return CheckResult(
name="orphaned_document_items",
severity=Severity.FAIL if orphans else Severity.OK,
message=(
"Document items reference missing documents."
if orphans
else "No orphaned document items."
),
remediation="haiku-rag rebuild" if orphans else None,
details=_sample(sorted(orphans)),
)
def _check_documents_without_items(
doc_ids: set[str], chunk_doc_ids: set[str], item_doc_ids: set[str]
) -> CheckResult:
"""A chunked document must have items; one without them is corrupt. Empty
documents legitimately have neither, so only chunked ones are flagged."""
missing = (doc_ids & chunk_doc_ids) - item_doc_ids
return CheckResult(
name="documents_without_items",
severity=Severity.WARN if missing else Severity.OK,
message=(
f"{len(missing)} chunked document(s) have no document items."
if missing
else "Every chunked document has document items."
),
remediation="haiku-rag rebuild" if missing else None,
details=_sample(sorted(missing)),
)
def _check_dangling_item_refs(
chunk_rows: list[dict], self_refs_by_doc: dict[str, set[str]]
) -> CheckResult:
"""Chunk metadata may reference self_refs that do not exist for that document."""
dangling: list[str] = []
for row in chunk_rows:
refs = json.loads(row.get("metadata") or "{}").get("doc_item_refs") or []
known = self_refs_by_doc.get(row["document_id"], set())
if any(ref not in known for ref in refs):
dangling.append(row["id"])
return CheckResult(
name="dangling_doc_item_refs",
severity=Severity.FAIL if dangling else Severity.OK,
message=(
f"{len(dangling)} chunk(s) reference missing document items."
if dangling
else "All chunk doc_item_refs resolve."
),
remediation="haiku-rag rebuild" if dangling else None,
details=_sample(dangling),
)
def _check_vector_dimension(stored_dim: int | None, actual_dim: int) -> CheckResult:
if stored_dim and stored_dim != actual_dim:
return CheckResult(
name="vector_dimension",
severity=Severity.FAIL,
message=(
f"Chunk vector size {actual_dim} does not match stored "
f"vector_dim {stored_dim}."
),
remediation="haiku-rag rebuild",
)
return CheckResult(
name="vector_dimension",
severity=Severity.OK,
message=f"Chunk vectors are {actual_dim}-dimensional.",
)
def _check_unembedded_chunks(id_column, embedded: "np.ndarray") -> CheckResult:
"""All-zero vectors, reported as a count with a few sampled ids so a large
corpus never materializes every chunk id."""
zero_rows = np.nonzero(~embedded)[0]
zero_count = int(zero_rows.size)
sample = [id_column[int(i)].as_py() for i in zero_rows[:_SAMPLE_LIMIT]]
if zero_count > _SAMPLE_LIMIT:
sample.append(f"... (+{zero_count - _SAMPLE_LIMIT} more)")
return CheckResult(
name="unembedded_chunks",
severity=Severity.WARN if zero_count else Severity.OK,
message=(
f"{zero_count} chunk(s) have an all-zero (unembedded) vector."
if zero_count
else "All chunks are embedded."
),
remediation="haiku-rag rebuild --embed-only" if zero_count else None,
details=sample,
)
def _document_centroids(
document_id_column, vectors: "np.ndarray", embedded: "np.ndarray", dim: int
) -> tuple[list[str], "np.ndarray", "np.ndarray"]:
"""Reduce each document's chunk vectors to one summed centroid.
Dictionary-encode the document ids into integer codes, then sum each
document's embedded rows in a single pass per document — no second full copy
of the vector matrix. Returns (document ids, summed centroids, chunk counts);
the caller normalizes.
"""
encoded = document_id_column.combine_chunks().dictionary_encode()
ids = encoded.dictionary.to_pylist()
codes = encoded.indices.to_numpy(zero_copy_only=False)
centroids = np.zeros((len(ids), dim), dtype=np.float32)
counts = np.zeros(len(ids), dtype=np.int64)
order = np.argsort(codes, kind="stable")
bounds = np.searchsorted(codes, np.arange(len(ids) + 1), sorter=order)
for d in range(len(ids)):
rows = order[bounds[d] : bounds[d + 1]]
rows = rows[embedded[rows]]
counts[d] = rows.size
if rows.size:
centroids[d] = vectors[rows].sum(axis=0)
return ids, centroids, counts
def _check_picture_data(
missing_picture_docs: list[str], content_type_by_doc: dict[str, str]
) -> CheckResult:
"""Pictures from image/PDF sources should carry raster bytes. Pictures that
are external image references in a text document (markdown, HTML) have no
embedded bytes by nature, so a missing raster there is expected."""
real_missing = [
doc_id
for doc_id in missing_picture_docs
if not content_type_by_doc.get(doc_id, "").startswith("text/")
]
return CheckResult(
name="picture_data",
severity=Severity.WARN if real_missing else Severity.OK,
message=(
f"{len(real_missing)} picture item(s) in image/PDF documents "
"have no image data."
if real_missing
else "Pictures that should carry image data have it."
),
remediation="haiku-rag rebuild" if real_missing else None,
details=_sample(sorted(set(real_missing))),
)
def _check_settings_row(total_settings: int, canonical: int) -> CheckResult:
"""Settings must hold exactly one canonical row."""
if total_settings == 0 or canonical != 1:
return CheckResult(
name="settings_row",
severity=Severity.FAIL,
message=(
f"Expected exactly one 'settings' row, found {canonical} "
f"(of {total_settings} total)."
),
remediation="haiku-rag migrate",
)
return CheckResult(
name="settings_row",
severity=Severity.OK,
message="Settings row is present.",
)
def _check_pending_migrations(stored_version: str) -> CheckResult:
pending = (
get_pending_upgrades(stored_version) if stored_version != "unknown" else []
)
return CheckResult(
name="pending_migrations",
severity=Severity.WARN if pending else Severity.OK,
message=(
f"{len(pending)} migration(s) pending (db version {stored_version})."
if pending
else f"Database is up to date (version {stored_version})."
),
remediation="haiku-rag migrate" if pending else None,
details=[f"{step.version}: {step.description or ''}" for step in pending],
)
async def run_db_checks(
store: Store,
config: AppConfig,
@ -484,110 +705,22 @@ async def run_db_checks(
labels_by_doc.setdefault(row["document_id"], set()).add(row["label"])
notify("Checking referential integrity")
# documents <-> document_meta must be 1:1.
orphan_docs = doc_ids - meta_doc_ids
orphan_meta = meta_doc_ids - doc_ids
if orphan_docs or orphan_meta:
details = [f"document with no meta: {d}" for d in _sample(sorted(orphan_docs))]
details += [f"meta with no document: {d}" for d in _sample(sorted(orphan_meta))]
results.append(
CheckResult(
name="document_meta_parity",
severity=Severity.FAIL,
message="documents and document_meta are out of sync.",
remediation="haiku-rag rebuild",
details=details,
)
)
else:
results.append(
CheckResult(
name="document_meta_parity",
severity=Severity.OK,
message="documents and document_meta are consistent.",
)
)
# Orphaned chunks / items reference a document that no longer exists.
orphan_chunk_docs = chunk_doc_ids - doc_ids
results.append(
CheckResult(
name="orphaned_chunks",
severity=Severity.FAIL if orphan_chunk_docs else Severity.OK,
message=(
"Chunks reference missing documents."
if orphan_chunk_docs
else "No orphaned chunks."
),
remediation="haiku-rag rebuild" if orphan_chunk_docs else None,
details=_sample(sorted(orphan_chunk_docs)),
)
)
orphan_item_docs = item_doc_ids - doc_ids
results.append(
CheckResult(
name="orphaned_document_items",
severity=Severity.FAIL if orphan_item_docs else Severity.OK,
message=(
"Document items reference missing documents."
if orphan_item_docs
else "No orphaned document items."
),
remediation="haiku-rag rebuild" if orphan_item_docs else None,
details=_sample(sorted(orphan_item_docs)),
)
)
results.append(_check_document_meta_parity(doc_ids, meta_doc_ids))
results.append(_check_orphaned_chunks(chunk_doc_ids, doc_ids))
results.append(_check_orphaned_items(item_doc_ids, doc_ids))
notify("Checking document chunking")
# Documents with no chunks, classified by what they contain and whether the
# embedder can index images.
results += _classify_unchunked(
doc_ids - chunk_doc_ids, labels_by_doc, store.embedder.supports_images
)
# A chunked document must have items; one without them is corrupt. Empty
# documents legitimately have neither, so only flag the chunked ones.
docs_missing_items = (doc_ids & chunk_doc_ids) - item_doc_ids
results.append(
CheckResult(
name="documents_without_items",
severity=Severity.WARN if docs_missing_items else Severity.OK,
message=(
f"{len(docs_missing_items)} chunked document(s) have no document items."
if docs_missing_items
else "Every chunked document has document items."
),
remediation="haiku-rag rebuild" if docs_missing_items else None,
details=_sample(sorted(docs_missing_items)),
)
)
results.append(_check_documents_without_items(doc_ids, chunk_doc_ids, item_doc_ids))
notify("Checking chunk references")
# Chunk metadata may reference self_refs that do not exist for that document.
dangling: list[str] = []
for row in chunk_rows:
refs = json.loads(row.get("metadata") or "{}").get("doc_item_refs") or []
known = self_refs_by_doc.get(row["document_id"], set())
if any(ref not in known for ref in refs):
dangling.append(row["id"])
results.append(
CheckResult(
name="dangling_doc_item_refs",
severity=Severity.FAIL if dangling else Severity.OK,
message=(
f"{len(dangling)} chunk(s) reference missing document items."
if dangling
else "All chunk doc_item_refs resolve."
),
remediation="haiku-rag rebuild" if dangling else None,
details=_sample(dangling),
)
)
results.append(_check_dangling_item_refs(chunk_rows, self_refs_by_doc))
notify("Scanning chunk vectors")
# Vector dimension consistency and unembedded (all-zero) vectors share one
# scan of the vector column — the heaviest check on large corpora.
# Vector dimension, unembedded vectors and duplicate detection share one
# scan of the vector column — the heaviest read on large corpora.
arrow = (
await store.chunks_table.query()
.select(["id", "vector", "document_id"])
@ -596,26 +729,7 @@ async def run_db_checks(
stored = await SettingsRepository(store).get_current_settings()
stored_dim = stored.get("embeddings", {}).get("model", {}).get("vector_dim")
actual_dim = arrow.schema.field("vector").type.list_size
if stored_dim and stored_dim != actual_dim:
results.append(
CheckResult(
name="vector_dimension",
severity=Severity.FAIL,
message=(
f"Chunk vector size {actual_dim} does not match stored "
f"vector_dim {stored_dim}."
),
remediation="haiku-rag rebuild",
)
)
else:
results.append(
CheckResult(
name="vector_dimension",
severity=Severity.OK,
message=f"Chunk vectors are {actual_dim}-dimensional.",
)
)
results.append(_check_vector_dimension(stored_dim, actual_dim))
# Reshape the Arrow fixed-size-list child buffer directly into an (N, dim)
# float32 matrix. Going through to_pylist() would box N*dim Python floats
@ -625,50 +739,17 @@ async def run_db_checks(
vectors = vec_col.values.to_numpy(zero_copy_only=False).reshape(-1, actual_dim)
embedded = vectors.any(axis=1) if vectors.size else np.zeros(0, dtype=bool)
# Unembedded (all-zero) chunks: report a count and a few sampled ids without
# materializing every chunk id.
zero_rows = np.nonzero(~embedded)[0]
zero_count = int(zero_rows.size)
id_col = arrow.column("id")
zero_sample = [id_col[int(i)].as_py() for i in zero_rows[:_SAMPLE_LIMIT]]
if zero_count > _SAMPLE_LIMIT:
zero_sample.append(f"... (+{zero_count - _SAMPLE_LIMIT} more)")
results.append(
CheckResult(
name="unembedded_chunks",
severity=Severity.WARN if zero_count else Severity.OK,
message=(
f"{zero_count} chunk(s) have an all-zero (unembedded) vector."
if zero_count
else "All chunks are embedded."
),
remediation="haiku-rag rebuild --embed-only" if zero_count else None,
details=zero_sample,
)
)
results.append(_check_unembedded_chunks(arrow.column("id"), embedded))
notify("Detecting near-duplicate documents")
# Near-identical documents (centroid cosine). Reduce each document's chunk
# vectors to one summed centroid during the scan: dictionary-encode the
# document ids into integer codes, then sum each document's embedded rows in
# a single pass per document — no second full copy of the vector matrix.
encoded = arrow.column("document_id").combine_chunks().dictionary_encode()
doc_ids = encoded.dictionary.to_pylist()
codes = encoded.indices.to_numpy(zero_copy_only=False)
centroids = np.zeros((len(doc_ids), actual_dim), dtype=np.float32)
counts = np.zeros(len(doc_ids), dtype=np.int64)
order = np.argsort(codes, kind="stable")
bounds = np.searchsorted(codes, np.arange(len(doc_ids) + 1), sorter=order)
for d in range(len(doc_ids)):
rows = order[bounds[d] : bounds[d + 1]]
rows = rows[embedded[rows]]
counts[d] = rows.size
if rows.size:
centroids[d] = vectors[rows].sum(axis=0)
centroid_doc_ids, centroids, counts = _document_centroids(
arrow.column("document_id"), vectors, embedded, actual_dim
)
# The matrix is the largest object here; drop it before clustering.
del vectors
results.append(
_check_duplicate_documents(
doc_ids,
centroid_doc_ids,
centroids,
counts,
uri_by_doc,
@ -679,9 +760,6 @@ async def run_db_checks(
)
notify("Checking picture data")
# Pictures from image/PDF sources should carry raster bytes. Pictures that
# are external image references in a text document (markdown, HTML) have no
# embedded bytes by nature, so a missing raster there is expected.
missing_picture_docs = [
row["document_id"]
for row in await store.document_items_table.query()
@ -689,72 +767,16 @@ async def run_db_checks(
.where("label = 'picture' AND picture_data IS NULL")
.to_list()
]
real_missing = [
doc_id
for doc_id in missing_picture_docs
if not content_type_by_doc.get(doc_id, "").startswith("text/")
]
results.append(
CheckResult(
name="picture_data",
severity=Severity.WARN if real_missing else Severity.OK,
message=(
f"{len(real_missing)} picture item(s) in image/PDF documents "
"have no image data."
if real_missing
else "Pictures that should carry image data have it."
),
remediation="haiku-rag rebuild" if real_missing else None,
details=_sample(sorted(set(real_missing))),
)
)
results.append(_check_picture_data(missing_picture_docs, content_type_by_doc))
notify("Checking settings and indexes")
# Settings must hold exactly one canonical row.
total_settings = await store.settings_table.count_rows()
canonical = len(
await store.settings_table.query().where("id = 'settings'").to_list()
)
if total_settings == 0 or canonical != 1:
results.append(
CheckResult(
name="settings_row",
severity=Severity.FAIL,
message=(
f"Expected exactly one 'settings' row, found {canonical} "
f"(of {total_settings} total)."
),
remediation="haiku-rag migrate",
)
)
else:
results.append(
CheckResult(
name="settings_row",
severity=Severity.OK,
message="Settings row is present.",
)
)
results.append(_check_settings_row(total_settings, canonical))
results.append(_check_embedding_drift(stored, config))
stored_version = str(stored.get("version", "unknown"))
pending = (
get_pending_upgrades(stored_version) if stored_version != "unknown" else []
)
results.append(
CheckResult(
name="pending_migrations",
severity=Severity.WARN if pending else Severity.OK,
message=(
f"{len(pending)} migration(s) pending (db version {stored_version})."
if pending
else f"Database is up to date (version {stored_version})."
),
remediation="haiku-rag migrate" if pending else None,
details=[f"{step.version}: {step.description or ''}" for step in pending],
)
)
results.append(_check_pending_migrations(str(stored.get("version", "unknown"))))
results.append(_check_vector_index(stats))