Skip to content
Closed
Changes from all commits
Commits
Show all changes
14 commits
Select commit Hold shift + click to select a range
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
79 changes: 72 additions & 7 deletions engraphis/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
from engraphis.core.graph_layers import normalize_graph_layer
from engraphis.core.ids import new_id as make_id
from engraphis.core.interfaces import Edge, GraphLayer, MemoryType, Node, Scope, SearchFilter
from engraphis.core.store import normalize_entity_name
from engraphis.core.store import _loads, _merge_edge_provenance, normalize_entity_name
from engraphis.graphdata import build_graph_payload, empty_graph

# ── validation limits (memory-poisoning / resource-exhaustion guards) ──────────
Expand Down Expand Up @@ -2307,14 +2307,79 @@ def _new_repo(old_repo_id):
)

# 3) Edges: relabel workspace/repo, remapping any entity ids folded in step 2.
# When a live source edge collides with an existing live target edge (same
# src/dst/relation/layer/repo), merge metadata instead of violating the
# partial unique index — mirrors Store._deduplicate_live_edges().
src_edges = [dict(x) for x in c.execute(
"SELECT id, repo_id, src, dst FROM edges WHERE workspace_id=?", (wid_src,))]
"SELECT id, repo_id, src, dst, relation, layer, weight, provenance, "
"valid_from, ingested_at, valid_to, expired_at "
"FROM edges WHERE workspace_id=?", (wid_src,))]
for ed in src_edges:
c.execute(
"UPDATE edges SET workspace_id=?, repo_id=?, src=?, dst=? WHERE id=?",
(wid_dst, _new_repo(ed["repo_id"]),
entity_remap.get(ed["src"], ed["src"]), entity_remap.get(ed["dst"], ed["dst"]),
ed["id"]))
new_src = entity_remap.get(ed["src"], ed["src"])
new_dst = entity_remap.get(ed["dst"], ed["dst"])
new_repo = _new_repo(ed["repo_id"])
is_live = ed["valid_to"] is None and ed["expired_at"] is None
target = None
if is_live:
# Check for a live target edge with the same identity.
if new_repo is not None:
target = c.execute(
"SELECT id, weight, provenance, valid_from, ingested_at "
"FROM edges WHERE workspace_id=? AND repo_id=? AND src=? "
"AND dst=? AND relation=? AND layer=? "
"AND valid_to IS NULL AND expired_at IS NULL LIMIT 1",
(wid_dst, new_repo, new_src, new_dst,
ed["relation"], ed["layer"]),
).fetchone()
else:
target = c.execute(
"SELECT id, weight, provenance, valid_from, ingested_at "
"FROM edges WHERE workspace_id=? AND repo_id IS NULL "
"AND src=? AND dst=? AND relation=? AND layer=? "
"AND valid_to IS NULL AND expired_at IS NULL LIMIT 1",
(wid_dst, new_src, new_dst,
ed["relation"], ed["layer"]),
).fetchone()
if target:
# Merge: keep target as survivor, retire source edge.
closed_at = time.time()
src_prov = _loads(ed["provenance"], {})
tgt_prov = _loads(target["provenance"], {})
merged_prov = _merge_edge_provenance(
[tgt_prov, src_prov], merged_ids=[ed["id"]])
merged_weight = max(
float(ed["weight"] or 0.0),
float(target["weight"] or 0.0))
valid_vals = [v for v in (target["valid_from"], ed["valid_from"])
if v is not None]
ingested_vals = [v for v in
(target["ingested_at"], ed["ingested_at"])
if v is not None]
c.execute(
"UPDATE edges SET weight=?, provenance=?, "
"valid_from=?, ingested_at=? WHERE id=?",
(merged_weight,
json.dumps(merged_prov, ensure_ascii=False),
min(valid_vals) if valid_vals else None,
min(ingested_vals) if ingested_vals else None,
target["id"]))
# Move live edge_supports from source to target.
c.execute(
"UPDATE edge_supports SET edge_id=? WHERE edge_id=? "
"AND valid_to IS NULL AND expired_at IS NULL",
(target["id"], ed["id"]))
# Bi-temporally close the source edge.
src_prov["canonical_deduplicated_into"] = target["id"]
c.execute(
"UPDATE edges SET valid_to=?, expired_at=?, "
"provenance=? WHERE id=?",
(closed_at, closed_at,
json.dumps(src_prov, ensure_ascii=False), ed["id"]))
else:
c.execute(
"UPDATE edges SET workspace_id=?, repo_id=?, src=?, dst=? "
"WHERE id=?",
(wid_dst, new_repo, new_src, new_dst, ed["id"]))

# 4) Memories / sessions / events: relabel workspace/repo per distinct repo_id
# bucket (ids, content and history are untouched).
Expand Down