"""Deterministic event rules: subtypes, wording, importance/confidence, idempotency and cross-surface clustering.""" from __future__ import annotations from datetime import UTC, datetime, timedelta import pytest from factories import intel_db # noqa: F401 — registers the fixture from companyatlas.services.events import derive_events, safe_wording, scale_importance from companyatlas.taxonomy import FORBIDDEN_WORDING, ChangeKind COMPANY = {"id": "co_test", "slug": "ztest-acme", "display_name": "Acme", "country": "CA", "industries": []} def _change(surface: str, *, significance: float = 0.6, delta: dict | None = None, diff: dict | None = None, kind: str | None = None) -> dict: from companyatlas.taxonomy import change_kind return {"id": "chg_test", "sensor_id": "sen_test", "company_id": "co_test", "surface": surface, "significance": significance, "kind": kind or str(change_kind(significance)), "structured_delta": delta or {}, "diff": diff or {}, "detected_at": datetime(2026, 9, 12, 12, tzinfo=UTC), "snapshot_before": "snap_a", "snapshot_after": "snap_b", "blocks_added": 0, "blocks_removed": 0, "blocks_modified": 0, "text_delta_ratio": 0.2} def _sensor(surface: str, connector: str = "generic-html-v1") -> dict: return {"id": "sen_test", "surface": surface, "connector_id": connector, "url": f"https://acme.example/{surface}"} def _by_subtype(derived): # type: ignore[no-untyped-def] out: dict[str, list] = {} for d in derived.events: out.setdefault(d.subtype, []).append(d) return out def _assert_clean(derived) -> None: # type: ignore[no-untyped-def] for d in derived.events: low = (d.title + " " + (d.summary or "")).lower() assert not any(bad in low for bad in FORBIDDEN_WORDING), d.title # ------------------------------------------------------------------------------------------------------------ hiring def test_job_count_increase_aggregate_and_ai(): added = [{"title": f"Engineer {i}", "country": "CA", "department": "Engineering"} for i in range(12)] added[0]["title"] = "Senior Machine Learning Engineer" added[1]["title"] = "AI Product Manager" delta = {"jobs": {"added": added, "removed": [], "open_before": 40, "open_after": 52}} d = derive_events(_change("careers", delta=delta), COMPANY, _sensor("careers")) by = _by_subtype(d) assert "JOB_COUNT_INCREASE" in by and "NEW_JOB" not in by # > 5 added → aggregate only ev = by["JOB_COUNT_INCREASE"][0] assert ev.title == "12 new positions detected on careers page" assert ev.old_value == "40" and ev.new_value == "52" assert len(ev.entities["jobs"]) == 12 assert "AI_HIRING" in by and by["AI_HIRING"][0].title == "2 AI-related positions detected on careers page" assert ev.confidence == 0.8 # HTML extraction assert 0 < ev.importance <= 1 assert d.needs_classification is False _assert_clean(d) def test_per_job_events_when_few_added_and_ats_confidence(): delta = {"jobs": {"added": [{"title": "Data Scientist", "location_text": "Toronto, CA"}, {"title": "Account Executive", "remote": True}], "removed": [], "open_before": 10, "open_after": 12}} d = derive_events(_change("jobs_board", delta=delta), COMPANY, _sensor("jobs_board", "greenhouse-v1")) by = _by_subtype(d) assert len(by["NEW_JOB"]) == 2 assert by["NEW_JOB"][0].title == "New position listed: Data Scientist (Toronto, CA)" assert by["NEW_JOB"][1].title == "New position listed: Account Executive (Remote)" assert all(e.confidence == 0.95 for e in d.events) # structured ATS JSON def test_job_count_decrease_wording_and_freeze_signal(): removed = [{"title": f"Role {i}"} for i in range(30)] delta = {"jobs": {"added": [], "removed": removed, "open_before": 40, "open_after": 10}} d = derive_events(_change("careers", significance=0.7, delta=delta), COMPANY, _sensor("careers")) by = _by_subtype(d) dec = by["JOB_COUNT_DECREASE"][0] assert dec.title == "30 monitored job listings no longer visible on careers page" assert "HIRING_FREEZE_SIGNAL" in by assert "no longer visible" in by["HIRING_FREEZE_SIGNAL"][0].title assert by["HIRING_FREEZE_SIGNAL"][0].review == "unexpected_activity" _assert_clean(d) def test_hiring_surge_uses_company_baseline(): added = [{"title": f"Role {i}"} for i in range(9)] delta = {"jobs": {"added": added, "removed": [], "open_before": 100, "open_after": 109}} no_baseline = derive_events(_change("careers", delta=delta), COMPANY, _sensor("careers")) assert "HIRING_SURGE" not in _by_subtype(no_baseline) # 9 < fallback threshold of 10 baseline = {"jobs_new_weekly": {"mean": 1.0, "stddev": 1.0, "samples": 8}} with_baseline = derive_events(_change("careers", delta=delta), COMPANY, _sensor("careers"), baseline=baseline) surge = _by_subtype(with_baseline)["HIRING_SURGE"][0] assert "baseline ≈ 1.0 new/week" in surge.title # ------------------------------------------------------------------------------------------------------------ pricing def test_price_increase_and_tier_changes(): delta = {"plans": {"price_changed": [{"plan_name": "Pro", "before": 49, "after": 59, "currency": "USD", "billing_period": "month"}], "added": [{"plan_name": "Enterprise", "contact_sales": True}], "removed": [{"plan_name": "Starter", "price": 9, "currency": "USD"}]}} d = derive_events(_change("pricing", significance=0.7, delta=delta), COMPANY, _sensor("pricing")) by = _by_subtype(d) inc = by["PRICE_INCREASE"][0] assert inc.title == "Pro plan price observed at $59 (was $49)" assert inc.old_value == "$49" and inc.new_value == "$59" assert inc.payload["pct"] == pytest.approx(20.4, abs=0.1) assert by["NEW_PRICING_TIER"][0].title == "New pricing tier listed: Enterprise (contact sales)" assert "enterprise" in by["NEW_PRICING_TIER"][0].tags assert by["PRICING_TIER_REMOVED"][0].title == "Pricing tier no longer listed: Starter" assert inc.importance > by["PRICING_TIER_REMOVED"][0].importance def test_price_decrease_eur(): delta = {"plans": {"price_changed": [{"plan_name": "Team", "before": 30, "after": 24, "currency": "EUR", "billing_period": "month", "pct": -20}]}} d = derive_events(_change("pricing", delta=delta), COMPANY, _sensor("pricing")) ev = _by_subtype(d)["PRICE_DECREASE"][0] assert ev.title == "Team plan price observed at €24 (was €30)" assert ev.summary == "Decrease of 20.0%, billed per month." def test_generic_pricing_change_from_text_diff(): diff = {"modified": [{"path": "Pricing > Pro", "before": "a", "after": "b"}, {"path": "Pricing > FAQ", "before": "c", "after": "d"}], "added": [], "removed": [], "counts": {"added": 0, "removed": 0, "modified": 2}, "text_delta_ratio": 0.3} d = derive_events(_change("pricing", diff=diff), COMPANY, _sensor("pricing")) ev = _by_subtype(d)["PRICING_CHANGE"][0] assert ev.title == "Pricing page materially updated (2 blocks changed)" assert ev.confidence == 0.7 # text-diff only assert ev.payload["sections"] == ["Pro", "FAQ"] # ------------------------------------------------------------------------------------------------------------ leadership def test_leadership_events_wording(): delta = {"people": {"added": [{"name": "Jane Doe", "title": "Chief Financial Officer", "role_category": "cfo", "is_executive": True}], "removed": [{"name": "John Roe", "title": "Chief Technology Officer", "role_category": "cto", "is_executive": True}], "title_changed": [{"name": "Ann Lee", "before": "VP Operations", "after": "COO"}]}} d = derive_events(_change("leadership", significance=0.7, delta=delta), COMPANY, _sensor("leadership")) by = _by_subtype(d) assert by["NEW_EXECUTIVE"][0].title == "Jane Doe listed as Chief Financial Officer on leadership page" assert by["EXECUTIVE_NO_LONGER_LISTED"][0].title == "John Roe no longer listed on leadership page" assert by["EXECUTIVE_TITLE_CHANGE"][0].title == "Ann Lee now listed as COO (was VP Operations)" assert by["LEADERSHIP_CHANGE"][0].title == "Leadership page updated: 1 added, 1 no longer listed, 1 title change" _assert_clean(d) # ------------------------------------------------------------------------------------------------------------ products / locations def test_product_and_location_rules(): delta = {"products": {"added": [{"name": "Atlas Copilot", "url": "https://acme.example/copilot"}], "removed": [{"name": "Atlas Lite"}]}, "locations": {"added": [{"name": "Toronto office", "city": "Toronto", "country": "CA", "kind": "office"}, {"name": "Tokyo", "city": "Tokyo", "country": "JP", "kind": "office"}], "removed": [{"name": "Berlin", "city": "Berlin", "country": "DE", "kind": "office"}], "new_countries": ["JP"]}} d = derive_events(_change("locations", delta=delta), COMPANY, _sensor("locations"), country_names={"JP": "Japan"}) by = _by_subtype(d) assert by["NEW_PRODUCT"][0].title == "New product listed: Atlas Copilot" assert "ai" in by["NEW_PRODUCT"][0].tags assert by["PRODUCT_REMOVED"][0].title == "Product no longer listed: Atlas Lite" assert by["NEW_LOCATION"][0].title == "New office listed: Toronto, CA" assert by["OFFICE_REMOVED"][0].title == "Office no longer listed: Berlin, DE" assert by["COUNTRY_EXPANSION"][0].title == "New country presence listed: Japan (Tokyo)" assert by["COUNTRY_EXPANSION"][0].importance > by["NEW_LOCATION"][0].importance # ------------------------------------------------------------------------------------------------------------ news / legal / homepage def test_news_subtype_heuristics(): delta = {"news": {"added": [{"title": "Acme announces Q3 2026 earnings results", "url": "https://acme.example/ir/q3", "category": "press"}, {"title": "How we built our new AI search", "url": "https://acme.example/blog/ai-search", "category": "blog"}, {"title": "v2.4 — webhooks and SDK updates", "url": "https://acme.example/changelog/2-4", "category": "changelog"}]}} diff = {"added": [{"path": "Feed", "before": None, "after": "Acme announces Q3 2026 earnings results — revenue up…"}], "modified": [], "removed": [], "counts": {"added": 1}} d = derive_events(_change("feed", delta=delta, diff=diff), COMPANY, _sensor("feed", "rss-feed-v1")) by = _by_subtype(d) assert "EARNINGS_RELEASE" in by and by["EARNINGS_RELEASE"][0].title.startswith("Earnings release: ") assert "BLOG_POST" in by and "ai" in by["BLOG_POST"][0].tags assert "CHANGELOG_ENTRY" in by assert all(e.confidence == 0.9 for e in d.events) # feed / JSON-LD grade evidence assert "EARNINGS_RELEASE" in d.summarize def test_terms_change_sections_and_review(): diff = {"modified": [{"path": "Terms > 7. Termination", "before": "x", "after": "y"}, {"path": "Terms > 12. Governing law", "before": "x", "after": "y"}], "added": [{"path": "Terms > 14. Arbitration", "before": None, "after": "New section"}], "removed": [], "counts": {"added": 1, "removed": 0, "modified": 2}, "text_delta_ratio": 0.12} d = derive_events(_change("legal_terms", significance=0.55, diff=diff), COMPANY, _sensor("legal_terms")) ev = _by_subtype(d)["TERMS_CHANGE"][0] assert ev.title == "Terms of service page materially updated (3 sections changed)" assert ev.payload["sections"] == ["7. Termination", "12. Governing law", "14. Arbitration"] assert ev.review == "legal_sensitive" and ev.confidence == 0.7 assert "TERMS_CHANGE" in d.summarize def test_homepage_redesign_vs_change_and_messaging(): diff = {"modified": [{"path": "Hero", "before": "Old", "after": "New"}], "added": [], "removed": [], "counts": {"added": 0, "removed": 0, "modified": 1}, "text_delta_ratio": 0.6} meta = {"meta": {"title_changed": {"before": "Acme — Payments", "after": "Acme — AI Payments Platform"}}} major = derive_events(_change("homepage", significance=0.7, diff=diff, delta=meta), COMPANY, _sensor("homepage")) by = _by_subtype(major) assert "HOMEPAGE_REDESIGN" in by and "MESSAGING_CHANGE" in by assert by["MESSAGING_CHANGE"][0].old_value == "Acme — Payments" meaningful = derive_events(_change("homepage", significance=0.5, diff=diff), COMPANY, _sensor("homepage")) assert list(_by_subtype(meaningful)) == ["WEBSITE_CHANGE"] assert meaningful.needs_classification is True and meaningful.classification_reason == "ambiguous_surface" def test_docs_api_changelog_text_rules(): diff = {"added": [{"path": "Changelog", "before": None, "after": "2026-09-10: Added bulk export endpoint\nDetails…"}], "modified": [], "removed": [], "counts": {"added": 1, "removed": 0, "modified": 0}, "text_delta_ratio": 0.15} assert _by_subtype(derive_events(_change("changelog", diff=diff), COMPANY, _sensor("changelog")))["CHANGELOG_ENTRY"][0].title == \ "Changelog entry detected: 2026-09-10: Added bulk export endpoint" assert "API_CHANGE" in _by_subtype(derive_events(_change("api", diff=diff), COMPANY, _sensor("api"))) assert "DOC_CHANGE" in _by_subtype(derive_events(_change("docs", diff=diff), COMPANY, _sensor("docs"))) # ------------------------------------------------------------------------------------------------------------ guards @pytest.mark.parametrize("significance", [0.05, 0.3]) def test_noise_and_minor_never_emit(significance): delta = {"jobs": {"added": [{"title": "X"}] * 20, "removed": [], "open_before": 1, "open_after": 21}} d = derive_events(_change("careers", significance=significance, delta=delta), COMPANY, _sensor("careers")) assert d.events == [] and d.needs_classification is False def test_unknown_surface_without_structure_requests_classification(): diff = {"modified": [{"path": "Partners", "before": "a", "after": "b"}], "added": [], "removed": [], "counts": {"modified": 1}} d = derive_events(_change("partners", diff=diff), COMPANY, _sensor("partners")) assert d.events == [] and d.needs_classification and d.classification_reason == "no_deterministic_event" def test_critical_changes_go_to_review_and_importance_scales(): delta = {"plans": {"price_changed": [{"plan_name": "Pro", "before": 10, "after": 12, "currency": "USD"}]}} crit = derive_events(_change("pricing", significance=0.9, delta=delta), COMPANY, _sensor("pricing")) mid = derive_events(_change("pricing", significance=0.5, delta=delta), COMPANY, _sensor("pricing")) assert crit.events[0].review == "major_event" and crit.events[0].importance > mid.events[0].importance assert scale_importance(0.8, 0.5) == pytest.approx(0.8, abs=1e-6) assert scale_importance(0.8, 1.0, 1.0) == 1.0 def test_safe_wording_rewrites_forbidden_phrases(): assert "no longer listed" in safe_wording("72 employees laid off") assert "shut down" not in safe_wording("office shut down").lower() # ============================================================================================================ database flow @pytest.mark.usefixtures("intel_db") async def test_process_pending_is_idempotent_and_clusters_across_surfaces(): from factories import cleanup, make_change, make_company, make_sensor from companyatlas.db import fetch_all, fetch_one, fetch_val, transaction from companyatlas.services.events import process_pending_changes, reprocess_events try: async with transaction() as conn: co = await make_company(conn) newsroom = await make_sensor(conn, co, "newsroom") feed = await make_sensor(conn, co, "feed", path="feed.xml") careers = await make_sensor(conn, co, "careers") item = {"title": "Acme launches Atlas Copilot for enterprises", "url": f"https://{co['canonical_domain']}/news/copilot", "category": "press"} await make_change(conn, newsroom, significance=0.6, structured_delta={"news": {"added": [item]}}) await make_change(conn, feed, significance=0.6, structured_delta={"news": {"added": [{**item, "url": item["url"] + "?utm=rss"}]}}) await make_change(conn, careers, significance=0.6, structured_delta={"jobs": {"added": [{"title": "ML Engineer", "country": "CA"}], "removed": [], "open_before": 5, "open_after": 6}}) await make_change(conn, careers, significance=0.1, structured_delta={}) # noise → archived, no event stats = await process_pending_changes(limit=50) assert stats["changes"] == 3 and stats["archived"] == 1 assert stats["events"] == 5 # NEWS×2 (one duplicate) + JOB_COUNT_INCREASE + AI_HIRING + NEW_JOB assert stats["duplicates"] == 1 async with transaction() as conn: events = await fetch_all(conn, "select * from events where company_id = :c order by created_at", c=co["id"]) news = [e for e in events if e["event_subtype"] == "NEWS_RELEASE"] assert len(news) == 2 and {e["status"] for e in news} == {"active", "duplicate"} canonical = next(e for e in news if e["status"] == "active") dup = next(e for e in news if e["status"] == "duplicate") assert canonical["cluster_id"] == dup["cluster_id"] cluster = await fetch_one(conn, "select * from event_clusters where id = :id", id=canonical["cluster_id"]) assert cluster["source_count"] == 2 and set(cluster["surfaces"]) == {"newsroom", "feed"} assert canonical["confidence"] == pytest.approx(0.93, abs=0.011) # max(0.8 html, 0.9 feed) + 0.03 corroboration assert canonical["payload"]["corroborations"] == 1 sources = await fetch_all(conn, "select * from event_sources where event_id = :e", e=canonical["id"]) assert {s["kind"] for s in sources} == {"primary", "corroboration"} assert await fetch_val(conn, "select count(*) from changes where company_id = :c and status = 'pending'", c=co["id"]) == 0 assert await fetch_val(conn, "select event_count from sensors where id = :s", s=careers["id"]) == 3 assert await fetch_val(conn, "select last_event_at from companies where id = :c", c=co["id"]) is not None assert await fetch_val(conn, "select count(*) from llm_jobs where company_id = :c", c=co["id"]) >= 0 # re-run: nothing pending, nothing new again = await process_pending_changes(limit=50) assert again["events"] == 0 and again["changes"] == 0 # reprocess over processed changes: dedupe keys make it a no-op rep = await reprocess_events(datetime.now(UTC) - timedelta(days=1), company_id=co["id"]) assert rep["changes"] == 3 and rep["events"] == 0 async with transaction() as conn: assert await fetch_val(conn, "select count(*) from events where company_id = :c", c=co["id"]) == 5 finally: await cleanup() @pytest.mark.usefixtures("intel_db") async def test_review_queue_and_llm_jobs_for_legal_change(monkeypatch): from factories import cleanup, make_change, make_company, make_sensor from companyatlas.config import settings from companyatlas.db import fetch_all, transaction from companyatlas.services.events import process_pending_changes monkeypatch.setattr(settings, "llm_base_url", "https://llm.example/v1") monkeypatch.setattr(settings, "llm_enabled", True) try: async with transaction() as conn: co = await make_company(conn) legal = await make_sensor(conn, co, "legal_terms") other = await make_sensor(conn, co, "partners") diff = {"modified": [{"path": "Terms > 7. Termination", "before": "old text", "after": "new text"}], "added": [], "removed": [], "counts": {"modified": 1}, "text_delta_ratio": 0.2} await make_change(conn, legal, significance=0.55, diff=diff) await make_change(conn, other, significance=0.5, diff=diff) stats = await process_pending_changes() assert stats["events"] == 1 and stats["llm_jobs"] == 2 async with transaction() as conn: reviews = await fetch_all(conn, "select kind from review_queue where company_id = :c", c=co["id"]) assert {r["kind"] for r in reviews} == {"legal_sensitive"} jobs = await fetch_all(conn, "select kind from llm_jobs where company_id = :c", c=co["id"]) assert sorted(j["kind"] for j in jobs) == ["classify_change", "summarize_event"] finally: await cleanup() def test_kind_enum_values_used_by_rules(): assert ChangeKind.NOISE == "noise" and ChangeKind.CRITICAL == "critical" def test_news_bursts_become_one_aggregate_event() -> None: from companyatlas.services.events import derive_events items = [{"title": f"Release {i}: quarterly update", "url": f"https://x.com/news/{i}", "published_at": None, "category": "press"} for i in range(9)] change = {"id": "chg_x", "sensor_id": "sen_x", "surface": "newsroom", "significance": 0.55, "kind": "meaningful", "diff": {}, "structured_delta": {"news": {"added": items}}} drafts = derive_events(change, {"id": "co_x", "display_name": "X"}, {"id": "sen_x", "surface": "newsroom", "connector_id": "generic-html-v1", "url": "https://x.com/news"}).events news = [d for d in drafts if d.subtype == "NEWS_RELEASE"] assert len(news) == 1 and news[0].title.startswith("9 new news releases published on") and len(news[0].entities["news"]) == 9 small = dict(change, structured_delta={"news": {"added": items[:2]}}) drafts = derive_events(small, {"id": "co_x", "display_name": "X"}, {"id": "sen_x", "surface": "newsroom", "connector_id": "generic-html-v1", "url": "https://x.com/news"}).events assert len([d for d in drafts if d.subtype == "NEWS_RELEASE"]) == 2