Skip to content

Commit fd1b08a

Browse files
authored
Merge pull request #28 from rudycelekli/fix/company-merge-refresh-20261002
fix(entities): refresh locked company evidence before merging
2 parents 384235b + 515ee18 commit fd1b08a

2 files changed

Lines changed: 53 additions & 0 deletions

File tree

‎apps/api/services/entities/graph.py‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -479,6 +479,16 @@ def merge_entities(
479479
if workspace_id is not None and kept.workspace_id != workspace_id:
480480
return {"error": "entity_not_found"}
481481

482+
# Acquire in stable ID order and reload both rows. The session may have
483+
# cached either entity before a different importer committed new evidence.
484+
db.flush()
485+
for entity in sorted((kept, merged), key=lambda item: item.id):
486+
_lock_for_update(db, entity)
487+
if kept.workspace_id != merged.workspace_id:
488+
return {"error": "cross_workspace_merge_forbidden"}
489+
if workspace_id is not None and kept.workspace_id != workspace_id:
490+
return {"error": "entity_not_found"}
491+
482492
# Snapshot for split()
483493
merged_keys = [k for (k,) in db.query(EntityBlockingKey.key)
484494
.filter(
Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
"""A merge uses committed evidence, never a session's stale identity-map copy."""
2+
import pytest
3+
from sqlalchemy import create_engine, event
4+
from sqlalchemy.orm import Session
5+
6+
from apps.api.database import Base
7+
from apps.api.services.entities.graph import merge_entities, resolve_company
8+
from apps.api.services.entities.models import CompanyEntity, EntityMergeLog
9+
from apps.api.services.poller.models import WatchSubscription # register watch table
10+
11+
12+
@pytest.mark.parametrize("observed_side", ["kept", "merged"])
13+
def test_merge_preserves_observations_committed_by_another_session(tmp_path, observed_side):
14+
engine = create_engine(f"sqlite:///{tmp_path / 'merge.db'}")
15+
@event.listens_for(engine, "connect")
16+
def foreign_keys(connection, _):
17+
connection.execute("PRAGMA foreign_keys=ON")
18+
Base.metadata.create_all(engine)
19+
with Session(engine, autoflush=False) as seed:
20+
kept, _ = resolve_company(seed, {"company": "Keep Original", "website": "keep.example"}, "source-a", workspace_id="ws")
21+
merged, _ = resolve_company(seed, {"company": "Separate Original", "website": "separate.example"}, "source-b", workspace_id="ws")
22+
seed.commit()
23+
k, m = kept.id, merged.id
24+
with Session(engine, autoflush=False) as stale:
25+
# Keep references alive: both rows are cached before the other commit.
26+
cached_kept, cached_merged = stale.get(CompanyEntity, k), stale.get(CompanyEntity, m)
27+
assert cached_kept.observation_count == cached_merged.observation_count == 1
28+
with Session(engine, autoflush=False) as writer:
29+
domain = "keep.example" if observed_side == "kept" else "separate.example"
30+
fresh, created = resolve_company(writer, {"company": "Later Observation", "website": domain,
31+
"email": "new@fixture.example"}, "source-c", workspace_id="ws")
32+
assert not created and fresh.id == (k if observed_side == "kept" else m)
33+
writer.commit()
34+
merge_entities(stale, k, m, workspace_id="ws")
35+
with Session(engine) as db:
36+
survivor = db.get(CompanyEntity, k)
37+
assert survivor.observation_count == 3
38+
assert sorted(survivor.sources) == ["source-a", "source-b", "source-c"]
39+
assert survivor.primary_email == "new@fixture.example"
40+
assert any(o["source"] == "source-c" for o in survivor.fields["company"])
41+
log = db.query(EntityMergeLog).one()
42+
assert log.snapshot["entity"]["observation_count"] == (2 if observed_side == "merged" else 1)
43+
engine.dispose()

0 commit comments

Comments
 (0)