Skip to content

Commit cdf0310

Browse files
Reduce writer contention during vector repair discovery (#206)
* fix(memory): make revisions, scoped evidence and index repair authoritative * test(eval): define source-bound coding and capacity acceptance protocols * docs(eval): refresh measured payload evidence and rendered figures * feat(ui): center project memory, atomic edits and shared history * docs: record rework acceptance, migration and remaining release gates * fix(docs): keep PyPI links and benchmark alternatives consistent * fix(memory): preserve approval identity and validate lineage metadata * fix(memory): replay completed session transitions after closure * fix(index): prioritize queued erasures over blocked publications * fix(history): retain promoted workspace versions * fix(memory): replay committed edits before embedding * fix(history): authorize broader roots within project reads * fix(dashboard): keep history requests in the selected project * perf: classify vector repair candidates before reserving writer * fix: match store workspace visibility during repair discovery
1 parent 8015713 commit cdf0310

8 files changed

Lines changed: 640 additions & 13 deletions

‎CHANGELOG.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
# Changelog
22

3-
All notable changes to Engraphis are documented here. Format loosely follows
3+
All notable changes to Engraphis are documented here. Format loosely follows
44
[Keep a Changelog](https://keepachangelog.com/); versions use SemVer.
55

66
## [Unreleased]

‎docs/INDEX_REPAIR_MAINTENANCE.md‎

Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,108 @@
1+
# Repair discovery and writer occupancy
2+
3+
This incremental change depends on the reliability candidate at
4+
`31da32c06a5ade989608504472e00dfca6f13f94` (PR #203). It changes candidate
5+
discovery for separate vector indexes. Public entrypoints, repair return fields,
6+
ranking defaults, schema 18, and transaction-sharing native indexes are unchanged.
7+
8+
## Reproduced problem and acceptance
9+
10+
Repair prioritizes erasure/quarantine cleanup before vector updates. Previously,
11+
finding one erasure behind 1,000 queued updates acquired SQLite's writer 1,002
12+
times: registration, 1,000 skipped updates, and one deletion. An independent
13+
connection could not finish a write while classification held that reservation.
14+
15+
Discovery now reads 100-row pages containing IDs, generation, canonical existence,
16+
and provenance/metadata. It avoids memory text and vector payloads. Canonical JSON
17+
decoding and quarantine rules remain shared with the store. Each page is fetched
18+
before yielding; no read transaction or reader lease spans publication. The keyset
19+
advances past the last scanned row even if the page yields no matching candidate.
20+
The store's canonical workspace predicate applies to the memory join, with temporal
21+
filtering disabled. Out-of-binding and missing records remain cleanup candidates;
22+
allowed historical records remain indexable. This matches publication's `get_memory`
23+
view without exposing another workspace's metadata during discovery.
24+
25+
Classification is a hint. Publication still reserves the writer, verifies the
26+
selected generation, rereads current canonical existence, eligibility and vector
27+
identity, applies the provider operation, and acknowledges only that generation.
28+
Stale hints leave recoverable debt. Once the provider-attempt budget is spent,
29+
iteration stops before requesting another filtered candidate, including after a
30+
failed deletion. Otherwise the iterator could scan an irrelevant tail after the
31+
last permitted attempt.
32+
33+
`tests/test_vector_repair_discovery.py` exercises real sync/store sequences for:
34+
35+
- Erasure and quarantine after 105 and 1,000 queued updates; one deletion needs at
36+
most two writer reservations, independent of the stable update backlog.
37+
- An independent writer completing while discovery is deliberately paused.
38+
- Erasure, same-generation quarantine, vector replacement and restoration between
39+
discovery and publication, without stale publication or lost repair debt.
40+
- Canonical decoding of malformed and legacy metadata/provenance.
41+
- Workspace-bound cleanup, multiple allowed workspaces and retained historical
42+
canonical records; external cleanup does not erase the canonical memory.
43+
- Early cleanup, both successful and failing, without scanning newer updates after
44+
the provider budget is exhausted.
45+
46+
Existing sync and storage tests retain coverage of delayed publication, newer
47+
generations during provider callbacks, process/restart recovery and native rollback.
48+
49+
## Reproducible measurement
50+
51+
Run the repair-only probe from the checkout being measured:
52+
53+
```sh
54+
python -m eval.repair_discovery --backlog 1000 --repetition 1 --output repair-1000-1.json
55+
python -m eval.repair_discovery --backlog 10000 --repetition 1 --output repair-10000-1.json
56+
```
57+
58+
For a baseline that predates the driver, run the same driver file with `runpy` from
59+
the baseline checkout. The imported Engraphis package identifies the measured
60+
source; every report records that source's revision and file hashes independently
61+
of the driver hash. Use the same Python/dependencies, machine, storage location and
62+
driver for both sources. Alternate their order across five independent process
63+
repetitions per backlog. Retain every raw result, including failures.
64+
The CLI records failure type and attempted configuration with source/driver identity
65+
and exits nonzero if setup, repair or verification fails; failed work has no success
66+
measurement. An unwritable output location or forced process termination requires
67+
the invoking runner to retain its own exit-status/log record.
68+
69+
The synthetic dataset has fixed 32-dimensional vectors and one erasure after the
70+
declared update backlog. Bulk setup uses canonical store APIs and queue triggers
71+
on a disposable file-backed SQLite database. Setup is excluded from timing. The
72+
measured invocation includes discovery, writer acquisition, synchronous fixture
73+
publication and pending counts. Writer timing includes commit/release overhead;
74+
nested acknowledgement does not acquire or count a second reservation. The same
75+
timing instrumentation applies to both sources.
76+
77+
The adapter makes no network calls. No embedding, recall, tokenizer or answer
78+
generation is measured. This probe does not establish the 100k operating target,
79+
mixed-workload contention, semantic quality or a provider latency guarantee. The
80+
complete-engine protocol and hardware gates remain in
81+
[ENGINE_CAPACITY_PROTOCOL.md](ENGINE_CAPACITY_PROTOCOL.md).
82+
83+
## Compatibility, backout and next dependencies
84+
85+
There is no new migration, policy, service or default. Backout restores the previous
86+
discovery implementation while retaining the canonical database, durable queue and
87+
generation checks. Never delete pending repair work to recover availability.
88+
89+
The following remain separate work:
90+
91+
1. Discovery can scan the whole queue, and repeated calls can repeat that scan.
92+
Pages bound row count, not metadata bytes or total time. Pending counts also
93+
traverse the queue. Measure these costs before adding indexed scheduling state.
94+
2. A permanently failing oldest deletion can consume repeated small attempt
95+
budgets. Durable fairness/backoff and coordination must preserve cleanup priority
96+
without acknowledging unapplied work or fabricating canonical generations.
97+
3. Provider calls still occupy the writer. Moving an arbitrary provider outside it
98+
lets a delayed old upsert recreate an erased vector, even if acknowledgement is
99+
rejected. Hard deadlines need an adapter-level cancellation/fencing contract;
100+
current tests do not prove remote completion safety after a process dies.
101+
4. Legacy resource imports still perform filesystem/extraction/embedding preparation
102+
inside a service writer boundary. A prepared batch must preserve whole-batch
103+
rollback, per-file outcomes, provenance, and caller-owned transactions before
104+
replacing that path. Removing its transaction decorator alone is insufficient.
105+
106+
An optional background repair worker must coordinate ownership and shut down its
107+
own connections cleanly. The dependency-light offline library continues to work
108+
without one. This change does not introduce background scheduling.

‎docs/REWORK_EXECUTION.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ The original checkout's active graph/layout changes remain separate.
1818
| --- | --- | --- | --- |
1919
| P1 reproduced defect | Delayed sync publication restored an erased or outdated external vector. | `core/vector_repair.py` publishes current canonical state under the writer reservation and acknowledges the applied generation. `tests/test_sync_index_repair.py` covers delayed publication, erasure, newer updates, provider failures and native rollback. | Arbitrary synchronous providers can still occupy the writer while publishing. |
2020
| P1 reproduced defect | A blocked vector update prevented later queued erasures from being repaired. | Repair traversal prioritizes canonical deletions and makes bounded progress past deferred updates. Focused regressions cover a one-operation budget, blocked embedding spaces, provider failures and later erasure. | A provider that cannot delete still leaves durable repair debt; deletion is not falsely acknowledged. |
21+
| P2 reproduced contention | Finding one erasure behind 1,000 updates acquired 1,002 writer reservations. | [Repair discovery](INDEX_REPAIR_MAINTENANCE.md) classifies paged canonical headers before reserving the writer, then revalidates inside it. Tests check independent writer progress, stale hints, and immediate stopping after the attempt budget. | Total discovery, repeated scans, failed-deletion fairness and provider latency remain separate scheduling work. |
2122
| P1 reproduced defect | Separate engines accepted multiple governed successors of one record. | `core/mutations.py` validates prepared versions and source claims inside the transaction; schema 18 retains content-free command receipts. `tests/test_governed_concurrency.py` exercises corrections, approvals, promotions and merges through independent engines and spawned processes. | Receipts coordinate processes sharing the canonical database; they are not a new distributed multi-database transaction protocol. |
2223
| P2 reproduced defect | A completed promotion or merge could not be retried after its session closed. | Existing receipts replay before transient active-session and embedding requirements. Tests reopen the engine, disable embedding, replay the result and reject removed successors; new writes still recheck session activity under the writer. | Changed requests are new operations and remain subject to current session and source guards. |
2324
| P1 reproduced defect | A shared claim key or consolidation lineage collapsed distinct repository facts. | Packing deduplicates repeated canonical IDs, preserves full ownership attribution and budgets it. `tests/test_context_scope_grounding.py` retains distinct repositories, values, conditions and title-bound subjects. | Stronger semantic compression remains an experiment; no ranking default changed. |

‎engraphis/core/vector_repair.py‎

Lines changed: 40 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@
1212
vector_index_shares_store_transaction,
1313
)
1414
from engraphis.core.poisoning import inspection_eligible
15-
from engraphis.core.store import _is_memory_database_path
15+
from engraphis.core.store import _is_memory_database_path, _loads
1616

1717
if TYPE_CHECKING:
1818
from engraphis.core.store import Store
@@ -57,29 +57,50 @@ def canonical_search_required(index, store: "Store", *,
5757

5858

5959
def _repair_candidates(store: "Store", target: str, memory_id: Optional[str],
60-
ceiling: tuple[int, str]) -> Iterator[tuple[str, int]]:
61-
"""Page queue identities without loading vectors or revisiting failed work."""
60+
ceiling: tuple[int, str], *,
61+
cleanup_only: bool) -> Iterator[tuple[str, int]]:
62+
"""Read bounded header pages; classification is only a publication hint.
63+
64+
Materialize each page before yielding, without retaining a read transaction.
65+
No vector payload or memory text is needed to skip work for the other phase.
66+
The publisher still revalidates current canonical state under the writer.
67+
"""
6268
after: Optional[tuple[int, str]] = None
6369
while True:
70+
# Match get_memory's instance boundary, including for historical rows.
71+
# Keep the predicate on the LEFT JOIN so hidden/orphaned queue entries
72+
# remain cleanup candidates instead of disappearing from discovery.
73+
scope_where, scope_params = store._where(None, include_invalid=True, alias="m")
74+
memory_join = " AND ".join(["m.id=r.memory_id", *scope_where])
6475
sql = (
65-
"SELECT memory_id,generation FROM vector_index_repairs WHERE identity=? "
66-
"AND (generation,memory_id)<=(?,?)"
76+
"SELECT r.memory_id,r.generation,m.id AS canonical_id,v.id AS vector_id,"
77+
"m.provenance,m.metadata FROM vector_index_repairs r "
78+
f"LEFT JOIN memories m ON {memory_join} "
79+
"LEFT JOIN mem_vectors v ON v.id=r.memory_id "
80+
"WHERE r.identity=? AND (r.generation,r.memory_id)<=(?,?)"
6781
)
68-
params: list[Any] = [target, *ceiling]
82+
params: list[Any] = [*scope_params, target, *ceiling]
6983
if memory_id is not None:
70-
sql += " AND memory_id=?"
84+
sql += " AND r.memory_id=?"
7185
params.append(memory_id)
7286
if after is not None:
73-
sql += " AND (generation,memory_id)>(?,?)"
87+
sql += " AND (r.generation,r.memory_id)>(?,?)"
7488
params.extend(after)
7589
rows = store.conn.execute(
76-
sql + " ORDER BY generation,memory_id LIMIT 100", params,
90+
sql + " ORDER BY r.generation,r.memory_id LIMIT 100", params,
7791
).fetchall()
7892
if not rows:
7993
return
8094
after = (int(rows[-1]["generation"]), str(rows[-1]["memory_id"]))
8195
for row in rows:
82-
yield str(row["memory_id"]), int(row["generation"])
96+
needs_upsert = (
97+
row["canonical_id"] is not None and row["vector_id"] is not None
98+
and inspection_eligible(
99+
_loads(row["provenance"], {}), _loads(row["metadata"], {}),
100+
)
101+
)
102+
if needs_upsert != cleanup_only:
103+
yield str(row["memory_id"]), int(row["generation"])
83104

84105

85106
def repair_vector_index(store: "Store", index: Any, *, embedding_space: str,
@@ -93,7 +114,8 @@ def repair_vector_index(store: "Store", index: Any, *, embedding_space: str,
93114
public engine's compatibility adapter without coupling this coordinator to it.
94115
Cleanup precedes upserts, including when ``limit=1``. The limit bounds provider
95116
attempts; finding cleanup may inspect the whole pending queue in 100-row
96-
pages. Repeated calls can rescan pending upserts; this is not a latency bound.
117+
read-only header pages. Skipped candidates do not acquire writer reservations.
118+
Repeated calls can rescan pending upserts; this is not a latency bound.
97119
"""
98120
if isinstance(limit, bool) or not isinstance(limit, int) or not 1 <= limit <= 1000:
99121
raise ValueError("repair limit must be an integer between 1 and 1000")
@@ -125,7 +147,9 @@ def repair_vector_index(store: "Store", index: Any, *, embedding_space: str,
125147
for cleanup_only in (True, False):
126148
if attempted >= limit or (not cleanup_only and not vector_writes_ready):
127149
break
128-
for selected_id, generation in _repair_candidates(store, target, memory_id, ceiling):
150+
for selected_id, generation in _repair_candidates(
151+
store, target, memory_id, ceiling, cleanup_only=cleanup_only,
152+
):
129153
if attempted >= limit:
130154
break
131155
operation = "delete" if cleanup_only else "upsert"
@@ -176,5 +200,9 @@ def repair_vector_index(store: "Store", index: Any, *, embedding_space: str,
176200
# Cleanup has already had its turn; retain fail-fast publication
177201
# during an upsert outage instead of repeatedly calling the provider.
178202
break
203+
# A filtered iterator may scan a long tail before yielding again. Stop
204+
# here after success or failure, before asking for another candidate.
205+
if attempted >= limit:
206+
break
179207
return {"attempted": attempted, "repaired": repaired,
180208
"pending": store.vector_index_pending(target) or 0}

0 commit comments

Comments
 (0)