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
1 change: 1 addition & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ NEWS_GLINER_MODEL=EmergentMethods/gliner_medium_news-v2.1
NEWS_SPACY_MODEL=
ACTIVISM_THRESHOLDS=
CROWDED_TRADE_MIN_MANAGERS=5
SIMILARITY_CROWDING_MIN_SCORE=0.5

# Data-quality and freshness monitoring
DQ_HARVEST_WINDOW_MINUTES=1560
Expand Down
22 changes: 22 additions & 0 deletions alembic/versions/017_manager_similarity_cosine.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
"""Store an embedding-cosine score alongside Jaccard similarity.

Revision ID: 017
Revises: 016
"""

import sqlalchemy as sa

from alembic import op

revision = "017"
down_revision = "016"
branch_labels = None
depends_on = None


def upgrade() -> None:
op.add_column("manager_similarity", sa.Column("cosine", sa.REAL(), nullable=True))


def downgrade() -> None:
op.drop_column("manager_similarity", "cosine")
6 changes: 6 additions & 0 deletions alerts/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,12 @@ def _evaluate_condition(self, condition: dict[str, Any], event: AlertEvent) -> b
if count is None or count < int(expected):
return False
continue
if key == "similar_manager_count_gte":
count = _as_int(payload.get("similar_manager_count"))
threshold = _as_float(expected)
if count is None or threshold is None or count < threshold:
return False
continue
Comment thread
stranske marked this conversation as resolved.
if key == "any_new_filing":
if bool(expected) and event.event_type != "new_filing":
return False
Expand Down
1 change: 1 addition & 0 deletions alerts/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
"min_ownership_pct": (0.0, 100.0),
"min_delta_pct": (0.0, 100.0),
"threshold_crossed": (0.0, 100.0),
"similar_manager_count_gte": (1.0, None),
}


Expand Down
40 changes: 31 additions & 9 deletions api/managers.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import io
import json
import logging
import math
import os
import re
import sqlite3
Expand Down Expand Up @@ -1352,8 +1353,9 @@ async def get_manager(
async def get_similar_managers(
id: int = Path(..., ge=1, description="Manager identifier"),
limit: int = Query(10, ge=1, le=100),
basis: str = Query("jaccard", pattern="^(jaccard|cosine)$"),
):
"""Return the strongest deterministic holding-overlap peers for a manager."""
"""Return the strongest holding-overlap or embedding-cosine peers for a manager."""
Comment thread
stranske marked this conversation as resolved.
conn = None
try:
conn = connect_db()
Expand All @@ -1369,23 +1371,43 @@ async def get_similar_managers(
raise HTTPException(status_code=404, detail="Manager not found")
ensure_manager_similarity_table(conn)
ph = "?" if isinstance(conn, sqlite3.Connection) else "%s"
score_column = "cosine" if basis == "cosine" else "jaccard"
rows = conn.execute(
"SELECT CASE WHEN manager_id_a = "
+ ph
+ " THEN manager_id_b ELSE manager_id_a END, jaccard, overlap_count, union_count "
"FROM manager_similarity WHERE manager_id_a = " + ph + " OR manager_id_b = " + ph + " "
"ORDER BY jaccard DESC, overlap_count DESC LIMIT " + ph,
(id, id, id, limit),
+ " THEN manager_id_b ELSE manager_id_a END, "
+ score_column
+ ", jaccard, cosine, overlap_count, union_count "
"FROM manager_similarity WHERE (manager_id_a = "
+ ph
+ " OR manager_id_b = "
+ ph
+ ") AND "
+ score_column
+ " IS NOT NULL ORDER BY "
+ score_column
+ " DESC, overlap_count DESC",
(id, id, id),
).fetchall()
finite_rows = [
row
for row in rows
if all(
value is None or math.isfinite(float(value)) for value in (row[1], row[2], row[3])
)
]
return {
"items": [
{
"manager_id": int(row[0]),
"jaccard": float(row[1]),
"overlap_count": int(row[2]),
"union_count": int(row[3]),
"basis": basis,
"score": float(row[1]),
"jaccard": float(row[2]),
"cosine": float(row[3]) if row[3] is not None else None,
"overlap_count": int(row[4]),
"union_count": int(row[5]),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
for row in rows
for row in finite_rows[:limit]
]
}
except DB_ERROR_TYPES as exc:
Expand Down
36 changes: 35 additions & 1 deletion etl/conviction_flow.py
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,32 @@ def _resolve_crowded_trade_min_managers(default: int = 3) -> int:
return max(1, parsed)


def _resolve_similarity_crowding_min_score(default: float = 0.5) -> float:
raw = os.getenv("SIMILARITY_CROWDING_MIN_SCORE")
if raw is None or raw.strip() == "":
return default
try:
return parse_finite_float(raw, min_value=0.0, max_value=1.0, allow_none=False)
except ValueError:
logger.warning("Invalid SIMILARITY_CROWDING_MIN_SCORE; using default", extra={"raw": raw})
return default


def _similar_manager_ids(conn: Any, manager_ids: list[int], minimum_score: float) -> list[int]:
"""Return crowd members connected by a Jaccard similarity at the configured floor."""
if len(manager_ids) < 2:
return []
ph = get_placeholder(conn)
placeholders = ", ".join([ph] * len(manager_ids))
rows = conn.execute(
"SELECT manager_id_a, manager_id_b FROM manager_similarity "
f"WHERE manager_id_a IN ({placeholders}) AND manager_id_b IN ({placeholders}) "
f"AND jaccard >= {ph}",
(*manager_ids, *manager_ids, minimum_score),
).fetchall()
return sorted({int(manager_id) for row in rows for manager_id in row})


def _ensure_crowded_trades_table(conn: Any) -> None:
if isinstance(conn, sqlite3.Connection):
conn.execute("""CREATE TABLE IF NOT EXISTS crowded_trades (
Expand Down Expand Up @@ -696,12 +722,20 @@ def dispatch_conviction_alerts(
conn = connect_db()
total_alerts = 0
try:
similarity_floor = _resolve_similarity_crowding_min_score()
for row in _load_crowded_trade_rows(conn, report_date):
manager_ids = json.loads(row[3]) if isinstance(row[3], str) else list(row[3])
similar_manager_ids = _similar_manager_ids(
conn, [int(manager_id) for manager_id in manager_ids], similarity_floor
)
payload = {
"cusip": str(row[0]),
"name_of_issuer": str(row[1]) if row[1] is not None else None,
"manager_count": int(row[2]),
"manager_ids": json.loads(row[3]) if isinstance(row[3], str) else list(row[3]),
"manager_ids": manager_ids,
"similar_manager_ids": similar_manager_ids,
"similar_manager_count": len(similar_manager_ids),
"similarity_floor": similarity_floor,
"total_value_usd": float(row[4]) if row[4] is not None else None,
"avg_conviction_pct": float(row[5]) if row[5] is not None else None,
"max_conviction_pct": float(row[6]) if row[6] is not None else None,
Expand Down
38 changes: 34 additions & 4 deletions etl/manager_similarity_flow.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,24 @@

import sqlite3
from itertools import combinations
from math import isfinite, sqrt
from typing import Any

from adapters.base import get_placeholder, get_table_columns
from embeddings import embed_text


def cosine_similarity(left: list[float], right: list[float]) -> float | None:
"""Return cosine similarity for two finite, equally-sized vectors."""
if len(left) != len(right) or not left or not all(isfinite(value) for value in (*left, *right)):
return None
denominator = sqrt(sum(value * value for value in left)) * sqrt(
sum(value * value for value in right)
)
if denominator == 0 or not isfinite(denominator):
return None
score = sum(a * b for a, b in zip(left, right, strict=True)) / denominator
return score if isfinite(score) else None


def ensure_manager_similarity_table(conn: Any) -> None:
Expand All @@ -15,7 +30,7 @@ def ensure_manager_similarity_table(conn: Any) -> None:
conn.execute("""CREATE TABLE IF NOT EXISTS manager_similarity (
manager_id_a INTEGER NOT NULL REFERENCES managers(id),
manager_id_b INTEGER NOT NULL REFERENCES managers(id),
jaccard REAL NOT NULL, overlap_count INTEGER NOT NULL, union_count INTEGER NOT NULL,
jaccard REAL NOT NULL, cosine REAL, overlap_count INTEGER NOT NULL, union_count INTEGER NOT NULL,
computed_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (manager_id_a, manager_id_b),
CHECK (manager_id_a < manager_id_b)
Expand All @@ -26,6 +41,9 @@ def ensure_manager_similarity_table(conn: Any) -> None:
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_manager_similarity_b ON manager_similarity(manager_id_b)"
)
columns = {row[1] for row in conn.execute("PRAGMA table_info(manager_similarity)")}
if "cosine" not in columns:
conn.execute("ALTER TABLE manager_similarity ADD COLUMN cosine REAL")


def compute_manager_similarity(conn: Any) -> int:
Expand Down Expand Up @@ -53,6 +71,10 @@ def compute_manager_similarity(conn: Any) -> int:
securities = holdings.setdefault(int(manager_id), set())
if security is not None:
securities.add(str(security))
vectors = {
manager_id: embed_text(" ".join(sorted(securities)))
for manager_id, securities in holdings.items()
}

def replace_rows() -> int:
conn.execute("DELETE FROM manager_similarity")
Expand All @@ -64,9 +86,17 @@ def replace_rows() -> int:
if not union:
continue
conn.execute(
"INSERT INTO manager_similarity (manager_id_a, manager_id_b, jaccard, overlap_count, union_count) "
f"VALUES ({ph}, {ph}, {ph}, {ph}, {ph})",
(left, right, len(overlap) / len(union), len(overlap), len(union)),
"INSERT INTO manager_similarity "
"(manager_id_a, manager_id_b, jaccard, cosine, overlap_count, union_count) "
f"VALUES ({ph}, {ph}, {ph}, {ph}, {ph}, {ph})",
(
left,
right,
len(overlap) / len(union),
cosine_similarity(vectors[left], vectors[right]),
len(overlap),
len(union),
),
)
count += 1
return count
Expand Down
1 change: 1 addition & 0 deletions schema.sql
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,7 @@ CREATE TABLE IF NOT EXISTS manager_similarity (
manager_id_a bigint NOT NULL REFERENCES managers(manager_id),
manager_id_b bigint NOT NULL REFERENCES managers(manager_id),
jaccard real NOT NULL,
cosine real,
overlap_count integer NOT NULL,
union_count integer NOT NULL,
computed_at timestamptz NOT NULL DEFAULT now(),
Expand Down
19 changes: 16 additions & 3 deletions tests/test_alert_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,7 @@ def test_alert_engine_supports_news_and_crowded_trade_conditions(tmp_path):
conn,
name="Crowding Rule",
event_type="crowded_trade_change",
condition_json='{"manager_count_gte":8}',
condition_json='{"similar_manager_count_gte":3}',
)
engine = AlertEngine(conn)

Expand All @@ -147,13 +147,26 @@ def test_alert_engine_supports_news_and_crowded_trade_conditions(tmp_path):
AlertEvent(
event_type="crowded_trade_change",
manager_id=1,
payload={"manager_count": 9},
payload={"manager_count": 9, "similar_manager_count": 3},
)
)

assert [alert.rule.name for alert in recent_news] == ["News Surge"]
assert recent_news[0].channels == ["email", "streamlit"]
assert [alert.rule.name for alert in crowded] == ["Crowding Rule"]
_insert_rule(
conn,
name="Fractional Crowding Rule",
event_type="crowded_trade_change",
condition_json='{"similar_manager_count_gte":3.9}',
)
fractional_crowding = engine.evaluate(
AlertEvent(
event_type="crowded_trade_change",
manager_id=1,
payload={"manager_count": 9, "similar_manager_count": 3},
)
)
assert [alert.rule.name for alert in fractional_crowding] == ["Crowding Rule"]
finally:
conn.close()

Expand Down
15 changes: 15 additions & 0 deletions tests/test_conviction_flow.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,21 @@ def _seed_manager_and_filing(
)


def test_similarity_crowding_counts_only_connected_managers():
conn = sqlite3.connect(":memory:")
conn.execute(
"CREATE TABLE manager_similarity (manager_id_a INTEGER, manager_id_b INTEGER, jaccard REAL)"
)
conn.executemany(
"INSERT INTO manager_similarity VALUES (?, ?, ?)",
[(1, 2, 0.8), (1, 3, 0.7), (2, 3, 0.9), (3, 4, 0.2), (1, 4, 0.5)],
)

assert conviction_flow._similar_manager_ids(conn, [1, 2, 3, 4], 0.5) == [1, 2, 3, 4]
assert conviction_flow._similar_manager_ids(conn, [1, 2, 3, 4], 0.8) == [1, 2, 3]
assert conviction_flow._similar_manager_ids(conn, [1, 2, 3, 4], 0.95) == []


def test_compute_conviction_scores_known_portfolio(tmp_path):
db_path = tmp_path / "conviction.db"
conn = sqlite3.connect(db_path)
Expand Down
14 changes: 14 additions & 0 deletions tests/test_crowded_contrarian.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,13 @@ def _setup_db(tmp_path, manager_count: int) -> str:
"computed_at TEXT DEFAULT CURRENT_TIMESTAMP, "
"UNIQUE(cusip, report_date))"
)
conn.execute(
"CREATE TABLE manager_similarity ("
"manager_id_a INTEGER NOT NULL, "
"manager_id_b INTEGER NOT NULL, "
"jaccard REAL NOT NULL, "
"CHECK (manager_id_a < manager_id_b))"
)
conn.execute(
"CREATE TABLE daily_diffs ("
"diff_id INTEGER PRIMARY KEY AUTOINCREMENT, "
Expand Down Expand Up @@ -310,6 +317,10 @@ def test_dispatch_conviction_alerts_emits_events_for_detected_rows(tmp_path, mon
report_date = "2024-05-01"

detect_crowded_trades.fn(report_date, min_managers=3, conn=conn)
conn.executemany(
"INSERT INTO manager_similarity(manager_id_a, manager_id_b, jaccard) VALUES (?, ?, ?)",
[(1, 2, 0.9), (1, 3, 0.9), (1, 4, 0.9), (1, 5, 0.9)],
)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
for manager_id in (1, 2, 3, 4):
conn.execute(
"INSERT INTO daily_diffs(manager_id, report_date, cusip, name_of_issuer, delta_type, "
Expand Down Expand Up @@ -351,6 +362,9 @@ def fake_fire_alerts_for_event_sync(db_conn, event):
"crowded_trade_change",
"contrarian_signal",
]
assert events[0].payload["similar_manager_ids"] == [1, 2, 3, 4, 5]
assert events[0].payload["similar_manager_count"] == 5
assert events[0].payload["similarity_floor"] == 0.5


def test_detect_contrarian_signals_skips_split_consensus(tmp_path):
Expand Down
Loading
Loading