diff --git a/engraphis/service.py b/engraphis/service.py index d2b439b..28dd4ad 100644 --- a/engraphis/service.py +++ b/engraphis/service.py @@ -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) ────────── @@ -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).