Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
135 changes: 134 additions & 1 deletion core/stores/memberships.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@

import sqlite3
from collections.abc import Iterable, Sequence
from dataclasses import dataclass
from dataclasses import dataclass, field
from datetime import UTC, datetime
from pathlib import Path

Expand Down Expand Up @@ -343,6 +343,65 @@ def n_occ(self, content_id: str, *, current_only: bool = True) -> int:
f"WHERE content_id = ? AND tombstoned = 0 {clause}", [content_id]).fetchone()
return int(row[0]) if row else 0

def n_doc_counts(self, *, layer: str | None = None,
current_only: bool = False) -> dict[str, int]:
"""`n_doc(v)` for EVERY atom at once — one GROUP BY instead of |V| point queries.

Defaults to the LIFETIME reading (`current_only=False`), because that is the variant the
rank-frequency histogram is defined over (D6): the corpus's vocabulary shape is a property
of everything it has ever held, not of one cut. Atoms with no occupancy do not appear —
their `n_doc` is 0 by definition, and materializing a zero per orphan would make the
histogram's tail an artifact of the ledger rather than of the corpus."""
clause = "AND current = 1" if current_only else ""
lane = "AND layer = ?" if layer else ""
params = [layer] if layer else []
return {str(r["content_id"]): int(r["n"]) for r in self._conn.execute(
f"SELECT content_id, count(DISTINCT path) AS n FROM memberships "
f"WHERE tombstoned = 0 {clause} {lane} GROUP BY content_id", params).fetchall()}

def rank_frequency(self, *, layer: str | None = None,
current_only: bool = False) -> list[int]:
"""The rank-frequency histogram of lifetime `n_doc(v)`: frequencies sorted DESCENDING, so
index *i* is rank *i+1* (D6, §4).

Returned as plain data rather than a plot: Zipf conformance is a falsifiable corpus
property and must be CHECKED, never assumed (T2/T4), and a caller that wants the check
needs the numbers. The gauge earns its keep only if a shape anomaly localizes something
real — boilerplate consolidation, vocabulary flux; if it never does, D6 says cut it."""
return sorted(self.n_doc_counts(layer=layer, current_only=current_only).values(),
reverse=True)

def occupied_atoms(self, *, layer: str | None = None) -> int:
"""Distinct atoms with at least one (un-tombstoned) occupancy — the denominator that makes
`embeds_avoided` honest. NOT the same as `|V|`: the plane may hold an orphan atom that no
fiber references (D8's crash window), and counting it as occupied would understate the
reuse this store is here to measure."""
lane = "AND layer = ?" if layer else ""
params = [layer] if layer else []
row = self._conn.execute(
f"SELECT count(DISTINCT content_id) FROM memberships WHERE tombstoned = 0 {lane}",
params).fetchone()
return int(row[0]) if row else 0

def occupancy_count(self, *, layer: str | None = None) -> int:
"""`|M|` restricted to one lane (or the whole relation) — the numerator of `|M|/|V|`."""
lane = "WHERE layer = ?" if layer else ""
params = [layer] if layer else []
row = self._conn.execute(
f"SELECT count(*) FROM memberships {lane}", params).fetchone()
return int(row[0]) if row else 0

def lane_gauges(self) -> dict[str, LaneGauge]:
"""Per-layer `(|M|, atoms)` in ONE aggregate query — the per-lane half of `|M|/|V|` (D6).

Kept as a store method rather than a query written at the gauge site so the two counts are
taken at the same cut, from the same scan: computing them separately is how a `|M|` from
after a landing gets divided by a `|V|` from before it."""
return {str(r["layer"]): LaneGauge(occupancies=int(r["n"]), atoms=int(r["v"]))
for r in self._conn.execute(
"SELECT layer, count(*) AS n, count(DISTINCT content_id) AS v "
"FROM memberships WHERE tombstoned = 0 GROUP BY layer").fetchall()}

def atom_ids_of_path(self, path: str) -> set[str]:
return {str(r["content_id"]) for r in self._conn.execute(
"SELECT DISTINCT content_id FROM memberships WHERE path = ?", [path]).fetchall()}
Expand Down Expand Up @@ -484,6 +543,80 @@ def resolve_occupancies(memberships: MembershipStore, hits: Sequence[dict[str, o
return out


@dataclass(frozen=True)
class LaneGauge:
"""One lane's frequency-plane reading."""

occupancies: int = 0 # |M| restricted to this lane
atoms: int = 0 # distinct atoms holding an occupancy here

@property
def dedup_factor(self) -> float:
"""Occupancies per atom — how much reuse this lane is actually buying."""
return (self.occupancies / self.atoms) if self.atoms else 0.0

@property
def embeds_avoided(self) -> int:
"""Occupancies past the first for each atom: the embeds the duplicated model would have
paid and this one does not."""
return max(0, self.occupancies - self.atoms)


@dataclass(frozen=True)
class FrequencyGauges:
"""The D6 standing gauges: `|M|`, `|V|`, the dedup factor, and embeds-avoided — per lane and
over the whole plane.

⚑ **`dedup_factor` IS the D7 falsifier, kept observable forever rather than measured once**
(the S5 amendment). It fails its keep by sitting at ≈1.0 after a full rebuild: that would mean
the membership model bought nothing and D7's economics are false. Reading ≈1.0 is therefore not
"a low number", it is the design being wrong, and the gauge exists to say so out loud.

`plane_atoms` is `|V|` — every atom row in the plane — while `atoms` counts only the atoms some
fiber references. They differ by exactly the orphans (D8's crash window), and keeping them
separate is what stops a repair-pass bug from quietly moving the dedup factor."""

occupancies: int = 0
atoms: int = 0
plane_atoms: int = 0
per_layer: dict[str, LaneGauge] = field(default_factory=dict)

@property
def dedup_factor(self) -> float:
"""`|M|/|V|` (D6) — occupancies per atom in the plane."""
return (self.occupancies / self.plane_atoms) if self.plane_atoms else 0.0

@property
def embeds_avoided(self) -> int:
return max(0, self.occupancies - self.atoms)

@property
def orphans(self) -> int:
return max(0, self.plane_atoms - self.atoms)

def __str__(self) -> str:
lanes = " · ".join(f"{k} {v.dedup_factor:.2f}×" for k, v in sorted(self.per_layer.items()))
return (f"|M|={self.occupancies} |V|={self.plane_atoms} "
f"dedup={self.dedup_factor:.2f}× ({lanes}) "
f"embeds_avoided={self.embeds_avoided} orphans={self.orphans}")


def frequency_gauges(vectors: VectorStore, memberships: MembershipStore) -> FrequencyGauges:
"""Read the D6 gauges. Cheap and read-only: three aggregate SQL queries plus one server-side
row count — no vector crosses into Python, which is what makes this registrable on a cadence
beside the drift-gauge family rather than an occasional investigation.

Every figure here is a QUERY. That is the issue #28 defect class stated as a rule: this repo
already carries docstrings quoting an edge count 8.4× off the live store, and a gauge whose
numbers were baked in at authoring time is that same defect with a dashboard on it."""
return FrequencyGauges(
occupancies=memberships.count(),
atoms=memberships.occupied_atoms(),
plane_atoms=vectors.atom_row_count(),
per_layer=memberships.lane_gauges(),
)


def repair_current_any(vectors: VectorStore, memberships: MembershipStore) -> tuple[int, int]:
"""Rebuild the `current_any` cache from membership truth (D8/R3). Returns (raised, lowered).

Expand Down
132 changes: 131 additions & 1 deletion core/stores/vectorstore.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,9 @@

from __future__ import annotations

from collections.abc import Iterable
from collections.abc import Iterable, Sequence
from dataclasses import dataclass
from datetime import timedelta
from pathlib import Path
from typing import Any

Expand Down Expand Up @@ -90,6 +91,36 @@ def _schema(dim: int) -> pa.Schema:
}


@dataclass(frozen=True)
class CompactionReport:
"""What `VectorStore.compact` physically did (§3).

Both a BEFORE and an AFTER row count are carried on purpose. Compaction's whole contract is
that it is semantically invisible, and the only way to assert invisibility without asserting it
vacuously is to have the two numbers side by side and require them equal WHILE the version
count fell — a compaction that no-ops leaves every number equal, including the one that is
supposed to move."""

rows_before: int = 0
rows_after: int = 0
versions_before: int = 0
versions_after: int = 0

@property
def rows_preserved(self) -> bool:
"""The invariant: compaction removes no logical row."""
return self.rows_before == self.rows_after

@property
def versions_dropped(self) -> int:
return max(0, self.versions_before - self.versions_after)

def __str__(self) -> str:
return (f"rows {self.rows_before}->{self.rows_after} "
f"versions {self.versions_before}->{self.versions_after} "
f"({self.versions_dropped} dropped)")


def is_code_atom_row(row: dict[str, Any]) -> bool:
"""Is this row a shed CODE **atom** row (D1) rather than a source-object chunk row?

Expand Down Expand Up @@ -236,6 +267,51 @@ def rows_for_source(self, source_path: str) -> list[dict[str, Any]]:
self._table().scan().where(f"source_path = {_sql_str(source_path)}")
.limit(0).to_list()]

def project(self, columns: Sequence[str], *, where: str | None = None,
limit: int = 0) -> list[dict[str, Any]]:
"""A projected, predicate-pushed read — the named columns and nothing else.

The generalization of `rows_for_source`' shape, and it exists for one measurable reason:
`vector` is 2560 floats per row, so a scan that does not name it costs a small fraction of
one that does. The rebuild's baseline (`ops/code_rebuild.py`) reads `id`/`layer`/`text`
over the whole code lane to re-derive the dedup economics and must never pay for geometry
it does not read; the carry-forward seed names `vector` precisely because copying it is the
point. `limit(0)` means UNLIMITED (verified empirically against the installed 0.33.0 and
pinned by the shim ratchet, exactly as `rows_for_source` documents).

`where` is a raw LanceDB predicate and is the CALLER's to build — use `_sql_str` for any
value that is not a literal this module wrote itself."""
if TABLE not in self._db.list_tables().tables:
return []
q = self._table().scan().select(list(columns))
if where is not None:
q = q.where(where)
return [dict(r) for r in q.limit(limit).to_list()]

def supersede_legacy_code_rows(self) -> int:
"""Flip every PRE-D1 duplicated code row to `current=false`, retaining it. Returns rows
flipped (0 when the store holds none — so a second call is a no-op).

The predicate is the exact complement of `is_code_atom_row` within the code lane: a row is
legacy iff it is CODE and still carries the occupancy coordinates an atom row sheds. The
rebuild lands the atom plane into the same table as the rows it replaces, so without this
both models answer the default current-view search and the dedup is bought but not served.

This is keep-and-link (`supersede_source`'s mechanism, D2) pointed at the retired ROW MODEL
rather than at a superseded version: one pushed-down predicate, a filtered count, one
in-place `update`. **Nothing is deleted** — `|V|` cannot decrease here, and D5's "purge is
the ONE removal" is untouched — and because the rows remain, the step is reversible by the
same update in the other direction."""
if TABLE not in self._db.list_tables().tables:
return 0
table = self._table()
where = (f"provenance = {_sql_str(Provenance.CODE.value)} "
"AND source_path <> '' AND current = true")
flipped = table.count_rows(where) # portable: do NOT rely on UpdateResult (bp-103 §11)
if flipped:
table.update(where, {"current": False})
return flipped

def delete_source(self, source_path: str) -> None:
"""Drop every derived row for one source document, by `source_path` (the stable doc identity
an amendment replaces a projection under — §4). Idempotent.
Expand Down Expand Up @@ -367,6 +443,20 @@ def atom_rows(self) -> list[dict[str, Any]]:
logged purge)."""
return [r for r in self.all_rows(provenances={Provenance.CODE}) if is_code_atom_row(r)]

def atom_row_count(self) -> int:
"""`|V|` — how many shed CODE-atom rows the plane holds, as a SERVER-SIDE count.

The same number `len(atom_rows())` gives, without the scan: `atom_rows` materializes every
row INCLUDING its 2560-float vector, which is the right cost for the repair pass (it reads
the flag on each row) and entirely the wrong cost for a gauge that runs on a cadence. The
predicate is `is_code_atom_row` written as SQL — the one place those two spellings must
agree, which is why the Python predicate is the documented reading and this cites it."""
if TABLE not in self._db.list_tables().tables:
return 0
return self._table().count_rows(
f"provenance = {_sql_str(Provenance.CODE.value)} "
"AND source_path = '' AND digest = ''")

def set_current_any(self, ids: Iterable[str], value: bool) -> int:
"""Set `current_any` on exactly the named atom rows (D2 step 5). Returns rows written.

Expand Down Expand Up @@ -407,6 +497,46 @@ def delete_atom(self, content_id: str) -> int:
table.delete(where)
return n

def dataset_versions(self) -> int:
"""How many dataset versions the table still retains (§3). One per write batch, and D2 makes
`current_any` flips routine — so this is the number compaction is measured against."""
if TABLE not in self._db.list_tables().tables:
return 0
return len(self._table().list_versions())

def compact(self, *, older_than: timedelta | None = None) -> CompactionReport:
"""Physical maintenance: compact fragments, then drop old dataset versions (§3).

[cross-ref: extension] No compaction path existed in this module (the note verified it), and
§3 makes it part of the store's semantics rather than an operational afterthought: the lance
dataset accumulates a version per write batch — 298 versions / 245 MB measured 2026-07-27
against ~232 MB of raw payload — and `current_any` flips are updates that rewrite fragments.
The rebuild ends with this, and housekeeping runs the cleanup half on cadence.

**Compaction is PHYSICAL, never logical.** It removes no row and changes no vector, so row
count and search results are invariant across it — the acceptance asserts BOTH, because
"nothing changed" is also what a compaction that silently did nothing produces. The version
count dropping is the other half, and it is the half that makes the assertion non-vacuous.

`older_than=None` takes the package's own retention default. Passing `timedelta(0)` reclaims
every superseded version immediately, which is what the rebuild wants (it has just written
thousands of batches) and what a test needs to observe a drop at all — but it forfeits
time-travel to any earlier version, so it is the caller's explicit choice, never the
default. `current_any` stays in lance throughout: the ANN prefilter needs it (D1/§3)."""
if TABLE not in self._db.list_tables().tables:
return CompactionReport()
table = self._table()
before_rows, before_versions = table.count_rows(None), len(table.list_versions())
# ONE call does both halves. The deprecated `compact_files`/`cleanup_old_versions` pair
# routes through `Table.to_lance()` and raises ImportError without the optional `pylance`
# package (found by running it, not by reading about it) — see the shim's note.
table.optimize(cleanup_older_than=older_than, delete_unverified=False)
after_rows, after_versions = table.count_rows(None), len(table.list_versions())
return CompactionReport(
rows_before=before_rows, rows_after=after_rows,
versions_before=before_versions, versions_after=after_versions,
)

def search(self, vector: list[float], *, k: int = 5,
provenances: Iterable[Provenance] | None = None,
include_superseded: bool = False) -> list[dict[str, Any]]:
Expand Down
Loading