Skip to content

Commit 1ffe595

Browse files
committed
fix(poller): consume funding and executive filings in one batch
Signed-off-by: Rudy Celekli <47457359+rudycelekli@users.noreply.github.com>
1 parent 9fe0b57 commit 1ffe595

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":
@@ -570,7 +570,8 @@ def _poll_one_source(store, watch_id, workspace_id, src, fire_key, lead_id, back
570570

571571
def _source_set(watch) -> list:
572572
"""Sources to fan in for a watch kind. company → funding+exec+hiring+tech.
573-
Each entry is one independent-txn source."""
573+
Each entry is one independent-txn source. SEC event types share a watermark
574+
and must consume their filing batch together."""
574575
kind = watch.kind
575576
types = set(watch.signal_types or [])
576577

@@ -584,13 +585,13 @@ def wants(*sts):
584585
settings, "TECH_STACK_WEBSITE_FETCH_ENABLED", False
585586
) else "tech"
586587

588+
# Both SEC detectors consume sec_last_accession. Separate transactions let
589+
# the first detector advance it before the second sees the same filings.
590+
sec_sources = (["funding_exec"] if wants("company_funded") and wants("executive_hired")
591+
else ["funding"] if wants("company_funded")
592+
else ["exec"] if wants("executive_hired") else [])
587593
if kind == "funding":
588-
out = []
589-
if wants("company_funded"):
590-
out.append("funding")
591-
if wants("executive_hired"):
592-
out.append("exec")
593-
return out or ["funding"]
594+
return sec_sources or ["funding"]
594595
if kind == "hiring":
595596
out = []
596597
if wants("hiring_surge"):
@@ -605,16 +606,12 @@ def wants(*sts):
605606
if kind == "account_group":
606607
return ["account_group"]
607608
if kind == "company":
608-
out = []
609-
if wants("company_funded"):
610-
out.append("funding")
611-
if wants("executive_hired"):
612-
out.append("exec")
609+
out = list(sec_sources)
613610
if wants("hiring_surge"):
614611
out.append("hiring")
615612
if wants("new_tech_adopted"):
616613
out.append(tech_src)
617-
return out or ["funding", "exec", "hiring", tech_src]
614+
return out or ["funding_exec", "hiring", tech_src]
618615
return []
619616

620617

‎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)