Skip to content

Commit df587ce

Browse files
Publish sync conflict vectors safely
1 parent 0a2110b commit df587ce

4 files changed

Lines changed: 88 additions & 7 deletions

File tree

‎engraphis/core/sync.py‎

Lines changed: 12 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1350,9 +1350,11 @@ def _apply_one(self, d: dict, rec, report: dict, accepted: dict, known: dict,
13501350
if existing is not None and provenance_is_approved(existing.provenance):
13511351
content_changed = rec.content != existing.content
13521352
self._rehome_external_record(rec, src_device=src_device)
1353-
self._preserve_hlc_conflict(
1353+
conflict_action = self._preserve_hlc_conflict(
13541354
existing, rec, report=report, known=known, dry_run=dry_run,
13551355
)
1356+
if conflict_action is not None:
1357+
pending_index_actions.append(conflict_action)
13561358
if not dry_run and content_changed:
13571359
self.store.audit(
13581360
"sync:%s" % _clamp_str(src_device or "peer", 128),
@@ -1412,9 +1414,11 @@ def _apply_one(self, d: dict, rec, report: dict, accepted: dict, known: dict,
14121414
rec.valid_to_recorded_at = now_ts()
14131415
rec.embedding = None
14141416
if existing is not None:
1415-
self._preserve_hlc_conflict(
1417+
conflict_action = self._preserve_hlc_conflict(
14161418
existing, rec, report=report, known=known, dry_run=dry_run,
14171419
)
1420+
if conflict_action is not None:
1421+
pending_index_actions.append(conflict_action)
14181422
if existing is None:
14191423
if not dry_run:
14201424
index_action = self._write(rec, commit=False)
@@ -1743,11 +1747,11 @@ def _preserve_hlc_conflict(
17431747
report: dict,
17441748
known: dict,
17451749
dry_run: bool,
1746-
) -> None:
1750+
) -> Optional[_VectorIndexAction]:
17471751
"""Keep the losing concurrent edit as one deterministic untrusted successor."""
17481752
conflict = self._hlc_conflict(existing, incoming)
17491753
if conflict is None:
1750-
return
1754+
return None
17511755
physical, logical, existing_hash, incoming_hash = conflict
17521756
winner = (
17531757
existing
@@ -1844,9 +1848,10 @@ def _preserve_hlc_conflict(
18441848
or not _same_sync_payload(already_preserved, preserved)
18451849
):
18461850
raise SyncError("sync conflict identity collision")
1847-
return
1851+
return None
1852+
index_action = None
18481853
if not dry_run:
1849-
self._write(preserved, commit=False)
1854+
index_action = self._write(preserved, commit=False)
18501855
self.store.audit(
18511856
"sync",
18521857
"sync_conflict_preserved",
@@ -1860,6 +1865,7 @@ def _preserve_hlc_conflict(
18601865
)
18611866
known[conflict_id] = preserved
18621867
report["conflicts_preserved"] += 1
1868+
return index_action
18631869

18641870
def _write(
18651871
self, rec: MemoryRecord, *, commit: bool = True,

‎engraphis/service.py‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3162,6 +3162,8 @@ def _apply_postgres_schema_snapshot(
31623162
claim_kind="catalog_snapshot_chunk",
31633163
resolve_conflicts=True,
31643164
))
3165+
if not stored_rows:
3166+
return {"workspace": workspace, "stored": 0, "entities": 0, "relations": 0}
31653167
stored = stored_rows[0]
31663168
wid, rid = self._require_scope(workspace, repo)
31673169
actual_ids: dict[str, str] = {}

‎tests/test_postgres_schema.py‎

Lines changed: 28 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,8 @@
66

77
from engraphis.backends import postgres_schema
88
from engraphis.core.interfaces import SchemaSnapshot, SearchFilter
9-
from engraphis.service import MemoryService
9+
import engraphis.service as service_module
10+
from engraphis.service import MAX_CONTENT_CHARS, MemoryService
1011

1112

1213
class _Cursor:
@@ -360,6 +361,32 @@ def inspect(self, supplied, *, schemas=None):
360361
assert "secret" not in serialized
361362

362363

364+
def test_empty_postgres_chunk_result_returns_without_indexing(monkeypatch):
365+
snapshot = SchemaSnapshot(
366+
title="PostgreSQL schema: empty",
367+
text="x" * (MAX_CONTENT_CHARS + 1),
368+
metadata={"database": "empty", "source_digest": "digest"},
369+
)
370+
371+
class _Introspector:
372+
def inspect(self, supplied, *, schemas=None):
373+
return snapshot
374+
375+
class _EmptyExtractor:
376+
def extract(self, _text):
377+
return []
378+
379+
monkeypatch.setattr(
380+
postgres_schema, "get_postgres_introspector", lambda: _Introspector()
381+
)
382+
monkeypatch.setattr(service_module, "ChunkingExtractor", _EmptyExtractor)
383+
service = MemoryService.create(":memory:")
384+
385+
assert service.import_postgres_schema(
386+
"postgresql://local/empty", workspace="acme"
387+
) == {"workspace": "acme", "stored": 0, "entities": 0, "relations": 0}
388+
389+
363390
def test_large_postgres_snapshot_keeps_every_chunk_distinct(monkeypatch):
364391
snapshot = SchemaSnapshot(
365392
title="PostgreSQL schema: large",

‎tests/test_sync.py‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2672,6 +2672,52 @@ def fail_late(actor, action, target, detail="", *, commit=True):
26722672
assert publications == []
26732673

26742674

2675+
def test_hlc_conflict_variant_publishes_external_vector():
2676+
engine = MemoryEngine.create(":memory:", vector_backend="numpy")
2677+
publications = []
2678+
2679+
class RecordingExternalIndex:
2680+
def upsert(self, ids, _vecs, meta=None, *, commit=True):
2681+
publications.extend(ids)
2682+
2683+
def delete(self, ids, *, commit=True):
2684+
del ids, commit
2685+
2686+
lower_node = f"dev_{'0' * 26}"
2687+
higher_node = f"dev_{'1' * 26}"
2688+
2689+
def bundle(content, node):
2690+
return {
2691+
"format": SYNC_FORMAT,
2692+
"version": 2,
2693+
"device_id": node,
2694+
"workspace_name": "w",
2695+
"repos": {},
2696+
"memories": [{
2697+
"id": "same-hlc-id",
2698+
"content": content,
2699+
"ingested_at": 42.0,
2700+
"valid_from": 42.0,
2701+
"modified_hlc": format_modified_hlc(42, 1, node),
2702+
}],
2703+
"mem_links": [],
2704+
}
2705+
2706+
syncer = SyncEngine(
2707+
engine.store,
2708+
embedder=engine.embedder,
2709+
vector_index=RecordingExternalIndex(),
2710+
)
2711+
syncer.apply_bundle(bundle("lower-node edit", lower_node), into_workspace="w")
2712+
publications.clear()
2713+
syncer.apply_bundle(bundle("higher-node edit", higher_node), into_workspace="w")
2714+
2715+
conflict_id = engine.store.conn.execute(
2716+
"SELECT id FROM memories WHERE id <> 'same-hlc-id'"
2717+
).fetchone()["id"]
2718+
assert conflict_id in publications
2719+
2720+
26752721
def test_sync_configured_embedder_failure_aborts_before_memory_write(caplog):
26762722
engine = MemoryEngine.create(":memory:", vector_backend="numpy")
26772723

0 commit comments

Comments
 (0)