Skip to content

Commit 92f4801

Browse files
committed
Merge remote-tracking branch 'review-pr/31' into review/integration
2 parents 41d00e5 + 1ffe595 commit 92f4801

2 files changed

Lines changed: 78 additions & 16 deletions

File tree

‎apps/api/services/poller/engine.py‎

Lines changed: 13 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -502,14 +502,14 @@ def _poll_one_source(store, watch_id, workspace_id, src, fire_key, lead_id, back
502502
if src in ("feed", "job_change", "account_group"):
503503
bill_src = src
504504
else:
505-
bill_src = "funding" if src in ("funding", "exec") else "hiring"
505+
bill_src = "funding" if src in ("funding", "exec", "funding_exec") else "hiring"
506506
if not _debit_source(db, workspace_id, bill_src, fire_key):
507507
watch.last_error = "insufficient_credits"
508508
return False
509509

510-
if src == "funding" or src == "exec":
510+
if src in ("funding", "exec", "funding_exec"):
511511
events, patch = sources.fetch_funding_and_exec(
512-
watch, want_funding=(src == "funding"), want_exec=(src == "exec"),
512+
watch, want_funding=(src != "exec"), want_exec=(src != "funding"),
513513
backfill=backfill,
514514
)
515515
elif src == "hiring" or src == "tech":
@@ -577,7 +577,8 @@ def _poll_one_source(store, watch_id, workspace_id, src, fire_key, lead_id, back
577577

578578
def _source_set(watch) -> list:
579579
"""Sources to fan in for a watch kind. company → funding+exec+hiring+tech.
580-
Each entry is one independent-txn source."""
580+
Each entry is one independent-txn source. SEC event types share a watermark
581+
and must consume their filing batch together."""
581582
kind = watch.kind
582583
types = set(watch.signal_types or [])
583584

@@ -591,13 +592,13 @@ def wants(*sts):
591592
settings, "TECH_STACK_WEBSITE_FETCH_ENABLED", False
592593
) else "tech"
593594

595+
# Both SEC detectors consume sec_last_accession. Separate transactions let
596+
# the first detector advance it before the second sees the same filings.
597+
sec_sources = (["funding_exec"] if wants("company_funded") and wants("executive_hired")
598+
else ["funding"] if wants("company_funded")
599+
else ["exec"] if wants("executive_hired") else [])
594600
if kind == "funding":
595-
out = []
596-
if wants("company_funded"):
597-
out.append("funding")
598-
if wants("executive_hired"):
599-
out.append("exec")
600-
return out or ["funding"]
601+
return sec_sources or ["funding"]
601602
if kind == "hiring":
602603
out = []
603604
if wants("hiring_surge"):
@@ -612,16 +613,12 @@ def wants(*sts):
612613
if kind == "account_group":
613614
return ["account_group"]
614615
if kind == "company":
615-
out = []
616-
if wants("company_funded"):
617-
out.append("funding")
618-
if wants("executive_hired"):
619-
out.append("exec")
616+
out = list(sec_sources)
620617
if wants("hiring_surge"):
621618
out.append("hiring")
622619
if wants("new_tech_adopted"):
623620
out.append(tech_src)
624-
return out or ["funding", "exec", "hiring", tech_src]
621+
return out or ["funding_exec", "hiring", tech_src]
625622
return []
626623

627624

‎tests/test_poller_sec_batch.py‎

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
"""Funding and executive detectors must consume their shared filing batch together."""
2+
from types import SimpleNamespace
3+
from uuid import uuid4
4+
5+
import pytest
6+
7+
from apps.api.core.tenancy import current_workspace_var
8+
from apps.api.database import SessionLocal
9+
from apps.api.services.leadgen.enrichment.providers import sec_edgar
10+
from apps.api.services.leadgen.orm_models import LeadRow, SignalRow
11+
from apps.api.services.poller import engine
12+
from apps.api.services.poller.models import WatchSubscription
13+
from apps.api.services.signals.store import SignalStore
14+
15+
16+
@pytest.mark.parametrize("kind,types,expected", [
17+
("company", ["company_funded", "executive_hired"], {"company_funded", "executive_hired"}),
18+
("funding", ["company_funded", "executive_hired"], {"company_funded", "executive_hired"}),
19+
("funding", ["company_funded"], {"company_funded"}),
20+
("company", ["executive_hired"], {"executive_hired"}),
21+
])
22+
def test_sec_poll_preserves_enabled_event_types_and_replay(kind, types, expected, monkeypatch):
23+
workspace, watch_id = f"sec-batch-{uuid4()}", str(uuid4())
24+
filing = SimpleNamespace(cik="0001234567", accession="acc101", filing_date="2026-01-01",
25+
fields={"funding_amount": 5000000},
26+
related_persons=[{"name": "Synthetic Person", "title": "Chief Executive Officer"}])
27+
calls = []
28+
class Provider:
29+
async def list_form_d_since(self, target, since):
30+
calls.append(since)
31+
return [filing] if not since or filing.accession > since else []
32+
monkeypatch.setattr(sec_edgar, "SecEdgarProvider", Provider)
33+
monkeypatch.setattr(engine, "_debit_source", lambda *args: True)
34+
with SessionLocal() as db:
35+
lead = LeadRow(workspace_id=workspace, company="Synthetic SEC Company")
36+
db.add(lead)
37+
db.flush()
38+
lead_id = lead.id
39+
db.add(WatchSubscription(id=watch_id, workspace_id=workspace, kind=kind,
40+
target=lead.company, lead_id=lead_id, interval="daily", enabled=True,
41+
signal_types=types, cursor={"bootstrapped": True, "sec_last_accession": "acc100"}))
42+
db.commit()
43+
store = SignalStore(workspace)
44+
token = current_workspace_var.set(workspace)
45+
try:
46+
for attempt in ("first", "replay"):
47+
with SessionLocal() as db:
48+
watch = db.get(WatchSubscription, watch_id)
49+
kinds = engine._source_set(watch)
50+
for source in kinds:
51+
engine._poll_one_source(store, watch_id, workspace, source, attempt, lead_id, False)
52+
with SessionLocal() as db:
53+
signals = db.query(SignalRow).filter_by(workspace_id=workspace).all()
54+
assert {signal.signal_type for signal in signals} == expected
55+
assert len(signals) == len(expected)
56+
assert {signal.lead_id for signal in signals} == {lead_id}
57+
assert db.get(WatchSubscription, watch_id).cursor["sec_last_accession"] == "acc101"
58+
assert calls == ["acc100", "acc101"]
59+
finally:
60+
with SessionLocal() as db:
61+
db.query(SignalRow).filter_by(workspace_id=workspace).delete()
62+
db.query(WatchSubscription).filter_by(id=watch_id).delete()
63+
db.query(LeadRow).filter_by(workspace_id=workspace).delete()
64+
db.commit()
65+
current_workspace_var.reset(token)

0 commit comments

Comments
 (0)