|
| 1 | +"""Create collisions and company merge/undo on PostgreSQL with FORCE RLS. |
| 2 | +
|
| 3 | +Set TEST_DATABASE_URL to a disposable PostgreSQL database owned by a role that |
| 4 | +can migrate and create the restricted test login (see pg_rls_support). |
| 5 | +""" |
| 6 | + |
| 7 | +import os |
| 8 | +import threading |
| 9 | +from concurrent.futures import ThreadPoolExecutor |
| 10 | +from uuid import uuid4 |
| 11 | + |
| 12 | +import pytest |
| 13 | +from sqlalchemy import event |
| 14 | + |
| 15 | +TEST_DATABASE_URL = os.getenv("TEST_DATABASE_URL") |
| 16 | +pytestmark = [ |
| 17 | + pytest.mark.postgres, |
| 18 | + pytest.mark.skipif(not TEST_DATABASE_URL, reason="TEST_DATABASE_URL not set"), |
| 19 | +] |
| 20 | + |
| 21 | + |
| 22 | +@pytest.fixture(scope="module") |
| 23 | +def app_session(): |
| 24 | + from tests.pg_rls_support import rls_app_session |
| 25 | + |
| 26 | + factory, dispose = rls_app_session(TEST_DATABASE_URL, pool_size=20) |
| 27 | + try: |
| 28 | + yield factory |
| 29 | + finally: |
| 30 | + dispose() |
| 31 | + |
| 32 | + |
| 33 | +def test_concurrent_creates_have_one_winner_without_contact_replacement(app_session, monkeypatch): |
| 34 | + from apps.api.core.tenancy import workspace_scope |
| 35 | + from apps.api.services.leadgen import store as store_module |
| 36 | + from apps.api.services.leadgen.models import Lead, LeadAlreadyExistsError |
| 37 | + |
| 38 | + monkeypatch.setattr(store_module, "SessionLocal", app_session) |
| 39 | + ws = f"create_{uuid4().hex}" |
| 40 | + store = store_module.PgLeadStore(ws) |
| 41 | + writers = 8 |
| 42 | + barrier = threading.Barrier(writers) |
| 43 | + local = threading.local() |
| 44 | + engine = app_session.kw["bind"] |
| 45 | + |
| 46 | + # Hold every writer just before its INSERT: all have observed the absent |
| 47 | + # lead, so this exercises the savepoint/unique-constraint recovery path. |
| 48 | + def before_insert(conn, cursor, statement, parameters, context, executemany): |
| 49 | + if statement.lstrip().upper().startswith("INSERT INTO LEADS ") and not getattr(local, "waited", False): |
| 50 | + local.waited = True |
| 51 | + barrier.wait(timeout=20) |
| 52 | + |
| 53 | + event.listen(engine, "before_cursor_execute", before_insert) |
| 54 | + def create(i): |
| 55 | + with workspace_scope(ws): |
| 56 | + try: |
| 57 | + lead_id = store.upsert_lead(Lead(company="Acme", city="Austin", |
| 58 | + contact_person=f"Person {i}", email=f"person{i}@acme.example"), create_only=True) |
| 59 | + return "created", i, lead_id |
| 60 | + except LeadAlreadyExistsError as exc: |
| 61 | + return "collision", i, exc.lead_id |
| 62 | + |
| 63 | + try: |
| 64 | + with ThreadPoolExecutor(max_workers=writers) as pool: |
| 65 | + outcomes = list(pool.map(create, range(writers))) |
| 66 | + finally: |
| 67 | + event.remove(engine, "before_cursor_execute", before_insert) |
| 68 | + winners = [result for result in outcomes if result[0] == "created"] |
| 69 | + assert len(winners) == 1, outcomes |
| 70 | + _, winner, lead_id = winners[0] |
| 71 | + assert sum(result[0] == "collision" for result in outcomes) == writers - 1 |
| 72 | + assert {result[2] for result in outcomes} == {lead_id} |
| 73 | + with workspace_scope(ws): |
| 74 | + saved = store.get_lead(lead_id) |
| 75 | + assert (saved.contact_person, saved.email) == (f"Person {winner}", f"person{winner}@acme.example") |
| 76 | + # A matching company/city in a different tenant is still independent. |
| 77 | + other = store_module.PgLeadStore(f"other_{uuid4().hex}") |
| 78 | + assert other.upsert_lead(Lead(company="Acme", city="Austin"), create_only=True) != lead_id |
| 79 | + assert other.get_lead(lead_id) is None |
| 80 | + |
| 81 | + |
| 82 | +@pytest.mark.parametrize("observed_side", ["kept", "merged"]) |
| 83 | +def test_merge_refreshes_committed_evidence_and_undo_restores_people(app_session, observed_side): |
| 84 | + from apps.api.core.tenancy import workspace_scope |
| 85 | + from apps.api.services.entities.graph import merge_entities, resolve_company, split_entity |
| 86 | + from apps.api.services.entities.models import CompanyEntity, EntityMergeLog, PersonEmployment, PersonEntity |
| 87 | + from apps.api.services.entities.people import resolve_person |
| 88 | + |
| 89 | + ws = f"merge_{uuid4().hex}" |
| 90 | + with workspace_scope(ws), app_session() as db: |
| 91 | + kept, _ = resolve_company(db, {"company": "Keep Original", "website": "keep.example"}, "a", workspace_id=ws) |
| 92 | + other, _ = resolve_company(db, {"company": "Separate Original", "website": "separate.example"}, "b", workspace_id=ws) |
| 93 | + person, _ = resolve_person(db, workspace_id=ws, name="Synthetic Person", company="Separate Original", |
| 94 | + company_domain="separate.example", linkedin_url="linkedin.com/in/fixture", source="fixture") |
| 95 | + db.commit() |
| 96 | + k, m, p = kept.id, other.id, person.id |
| 97 | + |
| 98 | + with workspace_scope(ws), app_session() as stale: |
| 99 | + # Keep both objects alive in the identity map before the separate commit. |
| 100 | + cached = [stale.get(CompanyEntity, entity_id) for entity_id in (k, m)] |
| 101 | + assert [entity.observation_count for entity in cached] == [1, 1] |
| 102 | + with app_session() as writer: |
| 103 | + domain = "keep.example" if observed_side == "kept" else "separate.example" |
| 104 | + resolve_company(writer, {"company": "Updated observation", "website": domain, |
| 105 | + "phone": "5550001234"}, "c", workspace_id=ws) |
| 106 | + writer.commit() |
| 107 | + merge_entities(stale, k, m, workspace_id=ws) |
| 108 | + |
| 109 | + with workspace_scope(ws), app_session() as db: |
| 110 | + survivor = db.get(CompanyEntity, k) |
| 111 | + assert survivor.observation_count == 3 |
| 112 | + assert set(survivor.sources) == {"a", "b", "c"} |
| 113 | + assert survivor.primary_phone == "5550001234" |
| 114 | + assert db.get(CompanyEntity, m) is None |
| 115 | + assert db.get(PersonEntity, p).company_entity_id == k |
| 116 | + assert db.query(PersonEmployment).filter_by(person_id=p).one().company_entity_id == k |
| 117 | + # Preserve an independent observation accepted after merging, too. |
| 118 | + resolve_company(db, {"company": "Keep Original", "website": "keep.example", |
| 119 | + "email": "later@keep.example"}, "later", workspace_id=ws) |
| 120 | + db.commit() |
| 121 | + log = db.query(EntityMergeLog).filter_by(workspace_id=ws, kept_id=k, merged_id=m).one() |
| 122 | + assert split_entity(db, log.id, workspace_id=ws)["restored_id"] == m |
| 123 | + |
| 124 | + with workspace_scope(ws), app_session() as db: |
| 125 | + kept, other = db.get(CompanyEntity, k), db.get(CompanyEntity, m) |
| 126 | + assert kept.observation_count == (3 if observed_side == "kept" else 2) |
| 127 | + assert other.observation_count == (2 if observed_side == "merged" else 1) |
| 128 | + assert kept.primary_email == "later@keep.example" |
| 129 | + assert set(kept.sources) == ({"a", "c", "later"} if observed_side == "kept" else {"a", "later"}) |
| 130 | + assert set(other.sources) == ({"b", "c"} if observed_side == "merged" else {"b"}) |
| 131 | + assert db.get(PersonEntity, p).company_entity_id == m |
| 132 | + assert db.query(PersonEmployment).filter_by(person_id=p).one().company_entity_id == m |
0 commit comments