Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 29 additions & 3 deletions catalyst_edge_mcp/adapters/gdelt.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

import hashlib
import json
import re
from collections.abc import Mapping
from datetime import datetime, timedelta
from email.utils import parsedate_to_datetime
Expand Down Expand Up @@ -44,6 +45,9 @@
MAX_ARTICLES = 50
MAX_RETRY_AFTER_SECONDS = 300.0
GDELT_GATE = ProviderGate(name="gdelt", concurrency=1, requests_per_second=1 / 6)
URL_PUBLICATION_DATE = re.compile(
r"(?<!\d)(20\d{2})[/-](0[1-9]|1[0-2])[/-](0[1-9]|[12]\d|3[01])(?!\d)"
)


class GdeltAdapter:
Expand Down Expand Up @@ -255,8 +259,11 @@ def _observation(
) -> EventObservation | None:
title = " ".join(str(article.get("title") or "").split())[:240]
url = str(article.get("url") or "").strip()
published_at = self._article_datetime(article.get("seendate"))
if not title or not url or published_at is None:
observed_at = self._article_datetime(article.get("seendate"))
published_at = self._url_publication_datetime(url)
if not title or not url or observed_at is None or published_at is None:
return None
if published_at > observed_at + timedelta(days=1):
return None
parsed = urlsplit(url)
if parsed.scheme.lower() != "https" or not parsed.hostname:
Expand All @@ -271,7 +278,7 @@ def _observation(
canonical_url=url,
title=title,
published_at=published_at,
observed_at=now,
observed_at=observed_at,
retrieved_at=now,
raw_sha256=raw_sha256,
parser_version=PARSER_VERSION,
Expand Down Expand Up @@ -306,6 +313,7 @@ def title_predicate(title: str, published_at: datetime) -> bool:
since,
title_predicate=title_predicate,
)
events = [event for event in events if self._event_is_recent(event, since)]
evidence = [self._evidence(event) for event in events]
effective_status = status or (
SourceStatus.FRESH if evidence else SourceStatus.NO_OBSERVATIONS
Expand Down Expand Up @@ -468,6 +476,24 @@ def _article_datetime(value: object) -> datetime | None:
return None
return parsed.replace(tzinfo=UTC) if parsed.tzinfo is None else parsed.astimezone(UTC)

@staticmethod
def _url_publication_datetime(url: str) -> datetime | None:
match = URL_PUBLICATION_DATE.search(urlsplit(url).path)
if match is None:
return None
try:
return datetime(*(int(value) for value in match.groups()), tzinfo=UTC)
except ValueError:
return None

@classmethod
def _event_is_recent(cls, event: StoredEvent, since: datetime) -> bool:
source = event.primary_source
if source.parser_version.startswith(PARSER_VERSION):
published_at = cls._url_publication_datetime(source.canonical_url)
return published_at is not None and published_at >= since
return event.published_at >= since

@staticmethod
def _retry_after(value: str | None, now: datetime) -> float | None:
if not value:
Expand Down
1 change: 1 addition & 0 deletions catalyst_edge_mcp/research.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ def _reviewable(item: Evidence) -> bool:
and item.confidence >= 0.50
and materiality not in {"discovery_only", "not_material"}
and bool(item.sources)
and bool(_claim_id(item))
)


Expand Down
109 changes: 87 additions & 22 deletions catalyst_edge_mcp/sec_filings.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

from __future__ import annotations

import asyncio
import hashlib
import re
from datetime import datetime, timedelta
Expand Down Expand Up @@ -263,6 +264,16 @@ def resolve_sec_ticker(ticker_to_cik: dict[str, str], ticker: str) -> str | None
return ticker_to_cik.get(ticker) or ticker_to_cik.get(ticker.replace(".", "-"))


def sec_ticker_is_current(payload: dict[str, Any], ticker: str) -> bool | None:
"""Use the issuer's current SEC ticker list when the field is available."""
values = payload.get("tickers")
if not isinstance(values, list):
return None
requested = {ticker, ticker.replace(".", "-"), ticker.replace("-", ".")}
current = {str(value).upper() for value in values}
return bool(requested & current)


async def lookup_cik_via_company_search(client: httpx.AsyncClient, ticker: str) -> str | None:
"""Resolve a ticker SEC's mapping file omits, via EDGAR company search.

Expand All @@ -284,6 +295,55 @@ async def lookup_cik_via_company_search(client: httpx.AsyncClient, ticker: str)
return match.group(1).zfill(10) if match else None


class SecCikResolver:
"""Share bounded SEC ticker resolution across sibling adapters."""

def __init__(self, *, miss_ttl_seconds: float = 60.0) -> None:
self._ticker_to_cik: dict[str, str] | None = None
self._map_lock = asyncio.Lock()
self._ticker_locks: dict[str, asyncio.Lock] = {}
self._fallback_cache: dict[str, tuple[float, str | None]] = {}
self._miss_ttl_seconds = miss_ttl_seconds

async def resolve(self, client: httpx.AsyncClient, ticker: str) -> str | None:
mapping = await self._mapping(client)
mapped = resolve_sec_ticker(mapping, ticker)
if mapped is not None:
return mapped

lock = self._ticker_locks.setdefault(ticker, asyncio.Lock())
async with lock:
now = asyncio.get_running_loop().time()
cached = self._fallback_cache.get(ticker)
if cached is not None and cached[0] > now:
return cached[1]
resolved = await lookup_cik_via_company_search(client, ticker)
self._fallback_cache[ticker] = (now + self._miss_ttl_seconds, resolved)
return resolved

async def _mapping(self, client: httpx.AsyncClient) -> dict[str, str]:
if self._ticker_to_cik is not None:
return self._ticker_to_cik
async with self._map_lock:
if self._ticker_to_cik is None:
async with SEC_GATE.request():
response = await client.get(TICKER_MAP_URL)
response.raise_for_status()
payload = response.json()
fields = payload.get("fields", [])
try:
ticker_index = fields.index("ticker")
cik_index = fields.index("cik")
except ValueError as exc:
raise ValueError("Unexpected SEC ticker mapping schema") from exc
self._ticker_to_cik = {
str(row[ticker_index]).upper(): str(row[cik_index]).zfill(10)
for row in payload.get("data", [])
if len(row) > max(ticker_index, cik_index)
}
return self._ticker_to_cik


class SecFilingsAdapter:
"""Collect recent filing metadata from the official SEC submissions API."""

Expand All @@ -298,13 +358,14 @@ def __init__(
clock=None,
fund_tickers: frozenset[str] = frozenset(),
store_path: str | None = None,
cik_resolver: SecCikResolver | None = None,
) -> None:
if "@" not in user_agent:
raise ValueError("SEC User-Agent must include a contact email address")
self.user_agent = user_agent
self._client = client
self._clock = clock or (lambda: datetime.now(UTC))
self._ticker_to_cik: dict[str, str] | None = None
self._cik_resolver = cik_resolver or SecCikResolver()
self._fund_tickers = fund_tickers
self.store = EvidenceStore(store_path) if store_path else None

Expand Down Expand Up @@ -374,8 +435,29 @@ async def _collect(
response.raise_for_status()
retrieved_at = self._as_utc(self._clock())
cutoff = retrieved_at - timedelta(days=lookback_days)
evidence = self._normalize_recent(response.json(), cik, cutoff, retrieved_at)
warnings = []
payload = response.json()
evidence = self._normalize_recent(payload, cik, cutoff, retrieved_at)
current_identity = sec_ticker_is_current(payload, ticker)
warnings = (
[f"SEC issuer {cik} no longer lists {ticker} as a current ticker."]
if current_identity is False
else []
)
reason_records = (
[
scoped_reason(
ReasonCode.ENTITY_REJECTED,
ReasonScope.EVALUATION,
ticker,
source_id=self.provider,
family=self.family,
observed_at=retrieved_at,
detail="ticker_not_current_for_sec_issuer",
)
]
if current_identity is False
else []
)
for item in evidence:
source = item.sources[0]
accession = source.accession_or_record_id
Expand Down Expand Up @@ -408,6 +490,7 @@ async def _collect(
status=SourceStatus.FRESH if evidence else SourceStatus.NO_OBSERVATIONS,
policy_decision=PolicyDecision.APPROVED,
collected_at=retrieved_at,
reason_records=reason_records,
)

async def _exhibit_links(
Expand Down Expand Up @@ -590,25 +673,7 @@ def _require_archive_url(url: str) -> None:
raise ValueError("SEC primary document URL is outside the official archive")

async def _resolve_cik(self, client: httpx.AsyncClient, ticker: str) -> str | None:
if self._ticker_to_cik is None:
async with SEC_GATE.request():
response = await client.get(TICKER_MAP_URL)
response.raise_for_status()
payload = response.json()
fields = payload.get("fields", [])
try:
ticker_index = fields.index("ticker")
cik_index = fields.index("cik")
except ValueError as exc:
raise ValueError("Unexpected SEC ticker mapping schema") from exc
self._ticker_to_cik = {
str(row[ticker_index]).upper(): str(row[cik_index]).zfill(10)
for row in payload.get("data", [])
if len(row) > max(ticker_index, cik_index)
}
return resolve_sec_ticker(self._ticker_to_cik, ticker) or await (
lookup_cik_via_company_search(client, ticker)
)
return await self._cik_resolver.resolve(client, ticker)

@classmethod
def _normalize_recent(
Expand Down
49 changes: 27 additions & 22 deletions catalyst_edge_mcp/sec_ownership.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,9 +30,8 @@
from catalyst_edge_mcp.sec_filings import (
SEC_GATE,
SUBMISSIONS_URL,
TICKER_MAP_URL,
lookup_cik_via_company_search,
resolve_sec_ticker,
SecCikResolver,
sec_ticker_is_current,
)

OWNERSHIP_FORMS = frozenset({"3", "3/A", "4", "4/A", "5", "5/A"})
Expand Down Expand Up @@ -254,13 +253,14 @@ def __init__(
clock=None,
fund_tickers: frozenset[str] = frozenset(),
store_path: str | None = None,
cik_resolver: SecCikResolver | None = None,
) -> None:
if "@" not in user_agent:
raise ValueError("SEC User-Agent must include a contact email address")
self.user_agent = user_agent
self._client = client
self._clock = clock or (lambda: datetime.now(UTC))
self._ticker_to_cik: dict[str, str] | None = None
self._cik_resolver = cik_resolver or SecCikResolver()
self._fund_tickers = fund_tickers
self.store = EvidenceStore(store_path) if store_path else None

Expand Down Expand Up @@ -322,13 +322,33 @@ async def _collect(
],
)
payload = await self._get_json(client, SUBMISSIONS_URL.format(cik=cik))
current_identity = sec_ticker_is_current(payload, ticker)
recent = payload.get("filings", {}).get("recent", {})
if not isinstance(recent, dict) or not isinstance(recent.get("form"), list):
raise ValueError("Unexpected SEC submissions schema")
cutoff = now - timedelta(days=lookback_days)
records: list[dict[str, Any]] = []
proposed_sales: list[Evidence] = []
warnings: list[str] = []
warnings = (
[f"SEC issuer {cik} no longer lists {ticker} as a current ticker."]
if current_identity is False
else []
)
reason_records = (
[
scoped_reason(
ReasonCode.ENTITY_REJECTED,
ReasonScope.EVALUATION,
ticker,
source_id=self.provider,
family=self.family,
observed_at=now,
detail="ticker_not_current_for_sec_issuer",
)
]
if current_identity is False
else []
)
planned = self._plan_ownership_documents(recent, cutoff, warnings)
if len(planned) > MAX_OWNERSHIP_DOCUMENTS:
# Never truncate silently. Measured p99 is 20 documents and the observed
Expand Down Expand Up @@ -383,7 +403,7 @@ async def _collect(
self._record_grouped_claims(ticker, evidence)
if not evidence:
warnings.append(f"No qualifying direct SEC insider activity found for {ticker}.")
return self._result(evidence, warnings, now)
return self._result(evidence, warnings, now, reason_records=reason_records)

async def _parsed_documents(
self,
Expand Down Expand Up @@ -688,22 +708,7 @@ def _form_144_evidence(
)

async def _resolve_cik(self, client: httpx.AsyncClient, ticker: str) -> str | None:
if self._ticker_to_cik is None:
payload = await self._get_json(client, TICKER_MAP_URL)
fields = payload.get("fields", [])
try:
ticker_index = fields.index("ticker")
cik_index = fields.index("cik")
except ValueError as exc:
raise ValueError("Unexpected SEC ticker mapping schema") from exc
self._ticker_to_cik = {
str(row[ticker_index]).upper(): str(row[cik_index]).zfill(10)
for row in payload.get("data", [])
if len(row) > max(ticker_index, cik_index)
}
return resolve_sec_ticker(self._ticker_to_cik, ticker) or await (
lookup_cik_via_company_search(client, ticker)
)
return await self._cik_resolver.resolve(client, ticker)

async def _get_json(self, client: httpx.AsyncClient, url: str) -> dict[str, Any]:
async with SEC_GATE.request():
Expand Down
5 changes: 4 additions & 1 deletion catalyst_edge_mcp/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@
)
from catalyst_edge_mcp.registry_config import RegistryBundle, load_registry_bundle
from catalyst_edge_mcp.replay.runtime import load_runtime_scorer
from catalyst_edge_mcp.sec_filings import SecFilingsAdapter
from catalyst_edge_mcp.sec_filings import SecCikResolver, SecFilingsAdapter
from catalyst_edge_mcp.sec_funds import SecFundAdapter
from catalyst_edge_mcp.sec_ownership import SecInsiderAdapter
from catalyst_edge_mcp.service import CatalystService
Expand All @@ -49,17 +49,20 @@ def build_service(
adapters = []
if settings.sec_user_agent:
fund_tickers = frozenset(registry.fund_identity_index)
cik_resolver = SecCikResolver()
adapters.extend(
[
SecFilingsAdapter(
settings.sec_user_agent,
fund_tickers=fund_tickers,
store_path=settings.evidence_store_path,
cik_resolver=cik_resolver,
),
SecInsiderAdapter(
settings.sec_user_agent,
fund_tickers=fund_tickers,
store_path=settings.evidence_store_path,
cik_resolver=cik_resolver,
),
SecFundAdapter(
settings.sec_user_agent,
Expand Down
7 changes: 6 additions & 1 deletion catalyst_edge_mcp/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -322,7 +322,12 @@ async def evaluate(self, request: ToolInput) -> CatalystEdgeResponse:
compact = self._compact(scored.evidence)
attributions = source_attributions(used_source_ids)
summary = build_summary(compact, missing, request.risk_mode)
checks = next_checks(compact, request.risk_mode, request.lookback_days)
checks = next_checks(
compact,
request.risk_mode,
request.lookback_days,
reason_records=reason_records,
)
research = build_research_assessment(
compact,
missing_families=missing,
Expand Down
Loading