From ab7ee18498236c0c3b71f983fc50bc44dfcf966b Mon Sep 17 00:00:00 2001 From: Tim Stranske Date: Fri, 31 Jul 2026 08:14:03 -0500 Subject: [PATCH 1/6] feat: add cosine similarity crowding alerts --- .../versions/016_manager_similarity_cosine.py | 22 +++++++++ alerts/engine.py | 5 ++ alerts/models.py | 1 + api/managers.py | 29 +++++++++--- etl/conviction_flow.py | 36 +++++++++++++- etl/manager_similarity_flow.py | 37 +++++++++++++-- schema.sql | 1 + tests/test_alert_engine.py | 4 +- tests/test_conviction_flow.py | 14 ++++++ tests/test_manager_similarity_flow.py | 47 +++++++++++++++---- 10 files changed, 174 insertions(+), 22 deletions(-) create mode 100644 alembic/versions/016_manager_similarity_cosine.py diff --git a/alembic/versions/016_manager_similarity_cosine.py b/alembic/versions/016_manager_similarity_cosine.py new file mode 100644 index 00000000..f6c3d186 --- /dev/null +++ b/alembic/versions/016_manager_similarity_cosine.py @@ -0,0 +1,22 @@ +"""Store an embedding-cosine score alongside Jaccard similarity. + +Revision ID: 016 +Revises: 015 +""" + +import sqlalchemy as sa + +from alembic import op + +revision = "016" +down_revision = "015" +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") diff --git a/alerts/engine.py b/alerts/engine.py index 207dfbc2..04f51779 100644 --- a/alerts/engine.py +++ b/alerts/engine.py @@ -83,6 +83,11 @@ 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")) + if count is None or count < int(expected): + return False + continue if key == "any_new_filing": if bool(expected) and event.event_type != "new_filing": return False diff --git a/alerts/models.py b/alerts/models.py index c0f0c5c7..990f857f 100644 --- a/alerts/models.py +++ b/alerts/models.py @@ -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), } diff --git a/api/managers.py b/api/managers.py index eaaee784..6a07a83c 100644 --- a/api/managers.py +++ b/api/managers.py @@ -1352,8 +1352,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.""" conn = None try: conn = connect_db() @@ -1369,21 +1370,35 @@ 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, + + " 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 LIMIT " + + ph, (id, id, id, limit), ).fetchall() 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]), } for row in rows ] diff --git a/etl/conviction_flow.py b/etl/conviction_flow.py index d63c3aa5..0ab68550 100644 --- a/etl/conviction_flow.py +++ b/etl/conviction_flow.py @@ -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 ( @@ -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, diff --git a/etl/manager_similarity_flow.py b/etl/manager_similarity_flow.py index cb6c48bc..03671ba2 100644 --- a/etl/manager_similarity_flow.py +++ b/etl/manager_similarity_flow.py @@ -4,9 +4,23 @@ import sqlite3 from itertools import combinations +from math import 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: + return None + denominator = sqrt(sum(value * value for value in left)) * sqrt( + sum(value * value for value in right) + ) + if denominator == 0: + return None + return sum(a * b for a, b in zip(left, right, strict=True)) / denominator def ensure_manager_similarity_table(conn: Any) -> None: @@ -15,7 +29,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) @@ -26,6 +40,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: @@ -53,6 +70,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") @@ -64,9 +85,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 diff --git a/schema.sql b/schema.sql index 255aeb5d..06da28c3 100644 --- a/schema.sql +++ b/schema.sql @@ -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(), diff --git a/tests/test_alert_engine.py b/tests/test_alert_engine.py index 10c1c1fc..eddcbb98 100644 --- a/tests/test_alert_engine.py +++ b/tests/test_alert_engine.py @@ -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) @@ -147,7 +147,7 @@ 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}, ) ) diff --git a/tests/test_conviction_flow.py b/tests/test_conviction_flow.py index 5e3346de..ad25123f 100644 --- a/tests/test_conviction_flow.py +++ b/tests/test_conviction_flow.py @@ -46,6 +46,20 @@ 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)], + ) + + assert conviction_flow._similar_manager_ids(conn, [1, 2, 3, 4], 0.5) == [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) diff --git a/tests/test_manager_similarity_flow.py b/tests/test_manager_similarity_flow.py index c13b090b..857bf437 100644 --- a/tests/test_manager_similarity_flow.py +++ b/tests/test_manager_similarity_flow.py @@ -4,10 +4,11 @@ from typing import Any, cast import httpx +import pytest from api import managers as managers_module from api.chat import app -from etl.manager_similarity_flow import compute_manager_similarity +from etl.manager_similarity_flow import compute_manager_similarity, cosine_similarity def test_similarity_uses_union_denominator_and_latest_filings(): @@ -38,6 +39,11 @@ def test_similarity_uses_union_denominator_and_latest_filings(): } +def test_cosine_similarity_matches_hand_computed_vectors(): + assert cosine_similarity([1.0, 1.0], [1.0, 0.0]) == pytest.approx(1 / 2**0.5) + assert cosine_similarity([0.0, 0.0], [1.0, 0.0]) is None + + def test_similarity_rebuild_rolls_back_on_insert_failure(): conn = sqlite3.connect(":memory:") conn.executescript("""PRAGMA foreign_keys = ON; @@ -60,14 +66,16 @@ def test_similarity_rebuild_rolls_back_on_insert_failure(): assert conn.execute("SELECT * FROM manager_similarity").fetchall() == before -async def _get_similar_manager(manager_id: int, limit: int): +async def _get_similar_manager(manager_id: int, limit: int, basis: str = "jaccard"): await cast(Any, app.router).startup() try: transport = httpx.ASGITransport(app=cast(Any, app)) async with httpx.AsyncClient( transport=transport, base_url="http://test", timeout=5.0 ) as client: - return await client.get(f"/managers/{manager_id}/similar", params={"limit": limit}) + return await client.get( + f"/managers/{manager_id}/similar", params={"limit": limit, "basis": basis} + ) finally: await cast(Any, app.router).shutdown() @@ -81,11 +89,11 @@ def test_similar_manager_endpoint_orders_canonical_pairs_and_returns_404(tmp_pat "INSERT INTO managers (id, name) VALUES (?, ?)", [(1, "One"), (2, "Two"), (3, "Three")] ) conn.execute( - "CREATE TABLE manager_similarity (manager_id_a INTEGER, manager_id_b INTEGER, jaccard REAL, overlap_count INTEGER, union_count INTEGER)" + "CREATE TABLE manager_similarity (manager_id_a INTEGER, manager_id_b INTEGER, jaccard REAL, cosine REAL, overlap_count INTEGER, union_count INTEGER)" ) conn.executemany( - "INSERT INTO manager_similarity VALUES (?, ?, ?, ?, ?)", - [(1, 2, 0.5, 2, 4), (1, 3, 0.75, 3, 4)], + "INSERT INTO manager_similarity VALUES (?, ?, ?, ?, ?, ?)", + [(1, 2, 0.5, 0.8, 2, 4), (1, 3, 0.75, 0.6, 3, 4)], ) conn.commit() conn.close() @@ -93,11 +101,34 @@ def test_similar_manager_endpoint_orders_canonical_pairs_and_returns_404(tmp_pat response = asyncio.run(_get_similar_manager(1, 1)) assert response.status_code == 200 assert response.json() == { - "items": [{"manager_id": 3, "jaccard": 0.75, "overlap_count": 3, "union_count": 4}] + "items": [ + { + "manager_id": 3, + "basis": "jaccard", + "score": 0.75, + "jaccard": 0.75, + "cosine": 0.6, + "overlap_count": 3, + "union_count": 4, + } + ] } reverse_response = asyncio.run(_get_similar_manager(2, 10)) assert reverse_response.status_code == 200 assert reverse_response.json() == { - "items": [{"manager_id": 1, "jaccard": 0.5, "overlap_count": 2, "union_count": 4}] + "items": [ + { + "manager_id": 1, + "basis": "jaccard", + "score": 0.5, + "jaccard": 0.5, + "cosine": 0.8, + "overlap_count": 2, + "union_count": 4, + } + ] } + cosine_response = asyncio.run(_get_similar_manager(1, 10, basis="cosine")) + assert cosine_response.json()["items"][0]["manager_id"] == 2 + assert cosine_response.json()["items"][0]["score"] == 0.8 assert asyncio.run(_get_similar_manager(999, 10)).status_code == 404 From 466a893111bdbcddde6bedd12f3e2df6db6efb0c Mon Sep 17 00:00:00 2001 From: Tim Stranske Date: Fri, 31 Jul 2026 08:26:39 -0500 Subject: [PATCH 2/6] fix: cover cosine crowding dependencies --- .env.example | 1 + tests/test_crowded_contrarian.py | 11 +++++++++++ 2 files changed, 12 insertions(+) diff --git a/.env.example b/.env.example index 288bd322..3f62b5bd 100644 --- a/.env.example +++ b/.env.example @@ -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 diff --git a/tests/test_crowded_contrarian.py b/tests/test_crowded_contrarian.py index 9f8fca0f..ba957a46 100644 --- a/tests/test_crowded_contrarian.py +++ b/tests/test_crowded_contrarian.py @@ -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, " @@ -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)], + ) for manager_id in (1, 2, 3, 4): conn.execute( "INSERT INTO daily_diffs(manager_id, report_date, cusip, name_of_issuer, delta_type, " From a11170bf2cb61c0bcc1b88a5238d27c257f85629 Mon Sep 17 00:00:00 2001 From: stranske Date: Fri, 31 Jul 2026 08:43:17 -0500 Subject: [PATCH 3/6] fix: renumber cosine migration to 017 after activism 016 merge The activism campaign migration merged to main as revision 016, so this branch's cosine migration collided on the same identifier and Alembic reported two heads named 016. --- ...ilarity_cosine.py => 017_manager_similarity_cosine.py} | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) rename alembic/versions/{016_manager_similarity_cosine.py => 017_manager_similarity_cosine.py} (83%) diff --git a/alembic/versions/016_manager_similarity_cosine.py b/alembic/versions/017_manager_similarity_cosine.py similarity index 83% rename from alembic/versions/016_manager_similarity_cosine.py rename to alembic/versions/017_manager_similarity_cosine.py index f6c3d186..2639ad2a 100644 --- a/alembic/versions/016_manager_similarity_cosine.py +++ b/alembic/versions/017_manager_similarity_cosine.py @@ -1,15 +1,15 @@ """Store an embedding-cosine score alongside Jaccard similarity. -Revision ID: 016 -Revises: 015 +Revision ID: 017 +Revises: 016 """ import sqlalchemy as sa from alembic import op -revision = "016" -down_revision = "015" +revision = "017" +down_revision = "016" branch_labels = None depends_on = None From 4e574c4953fe6780666a18fc44ab9e0f1da15086 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Fri, 31 Jul 2026 14:05:56 +0000 Subject: [PATCH 4/6] chore(codex-keepalive): apply updates (PR #1497) --- langsmith-fleet-worker-attempt.json | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) create mode 100644 langsmith-fleet-worker-attempt.json diff --git a/langsmith-fleet-worker-attempt.json b/langsmith-fleet-worker-attempt.json new file mode 100644 index 00000000..dc5c0634 --- /dev/null +++ b/langsmith-fleet-worker-attempt.json @@ -0,0 +1,16 @@ +{ + "agent": "codex", + "cli_version": "0.125.0", + "emitted_at": "2026-07-31T14:05:55.307761Z", + "execution_profile": "codex-default", + "fallback_models": [ + "gpt-5.4" + ], + "operation_role": "worker", + "pr_number": "1497", + "requested_model": "gpt-5.5", + "runner": "reusable-codex-run", + "schema": "langsmith-fleet/v1", + "selected_model": "", + "selection_reason": "" +} From 530304585a853db51e7cbe1cce28b71777d3c954 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Fri, 31 Jul 2026 14:11:02 +0000 Subject: [PATCH 5/6] chore(codex-keepalive): apply updates (PR #1497) --- langsmith-fleet-worker-attempt.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/langsmith-fleet-worker-attempt.json b/langsmith-fleet-worker-attempt.json index dc5c0634..6742b81e 100644 --- a/langsmith-fleet-worker-attempt.json +++ b/langsmith-fleet-worker-attempt.json @@ -1,7 +1,7 @@ { "agent": "codex", "cli_version": "0.125.0", - "emitted_at": "2026-07-31T14:05:55.307761Z", + "emitted_at": "2026-07-31T14:11:00.524790Z", "execution_profile": "codex-default", "fallback_models": [ "gpt-5.4" From f6a0b97410c077f4c4aa388654483fb47e38586d Mon Sep 17 00:00:00 2001 From: Tim Stranske Date: Fri, 31 Jul 2026 09:30:56 -0500 Subject: [PATCH 6/6] fix: address cosine similarity review findings --- alerts/engine.py | 3 ++- api/managers.py | 15 +++++++++++---- etl/manager_similarity_flow.py | 9 +++++---- langsmith-fleet-worker-attempt.json | 16 ---------------- tests/test_alert_engine.py | 15 ++++++++++++++- tests/test_conviction_flow.py | 5 +++-- tests/test_crowded_contrarian.py | 3 +++ tests/test_manager_similarity_flow.py | 2 ++ 8 files changed, 40 insertions(+), 28 deletions(-) delete mode 100644 langsmith-fleet-worker-attempt.json diff --git a/alerts/engine.py b/alerts/engine.py index 04f51779..ac6d2e6b 100644 --- a/alerts/engine.py +++ b/alerts/engine.py @@ -85,7 +85,8 @@ def _evaluate_condition(self, condition: dict[str, Any], event: AlertEvent) -> b continue if key == "similar_manager_count_gte": count = _as_int(payload.get("similar_manager_count")) - if count is None or count < int(expected): + threshold = _as_float(expected) + if count is None or threshold is None or count < threshold: return False continue if key == "any_new_filing": diff --git a/api/managers.py b/api/managers.py index 6a07a83c..6adc69cf 100644 --- a/api/managers.py +++ b/api/managers.py @@ -6,6 +6,7 @@ import io import json import logging +import math import os import re import sqlite3 @@ -1385,10 +1386,16 @@ async def get_similar_managers( + score_column + " IS NOT NULL ORDER BY " + score_column - + " DESC, overlap_count DESC LIMIT " - + ph, - (id, id, id, limit), + + " 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": [ { @@ -1400,7 +1407,7 @@ async def get_similar_managers( "overlap_count": int(row[4]), "union_count": int(row[5]), } - for row in rows + for row in finite_rows[:limit] ] } except DB_ERROR_TYPES as exc: diff --git a/etl/manager_similarity_flow.py b/etl/manager_similarity_flow.py index 03671ba2..684e4ce6 100644 --- a/etl/manager_similarity_flow.py +++ b/etl/manager_similarity_flow.py @@ -4,7 +4,7 @@ import sqlite3 from itertools import combinations -from math import sqrt +from math import isfinite, sqrt from typing import Any from adapters.base import get_placeholder, get_table_columns @@ -13,14 +13,15 @@ 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: + 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: + if denominator == 0 or not isfinite(denominator): return None - return sum(a * b for a, b in zip(left, right, strict=True)) / denominator + 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: diff --git a/langsmith-fleet-worker-attempt.json b/langsmith-fleet-worker-attempt.json deleted file mode 100644 index 6742b81e..00000000 --- a/langsmith-fleet-worker-attempt.json +++ /dev/null @@ -1,16 +0,0 @@ -{ - "agent": "codex", - "cli_version": "0.125.0", - "emitted_at": "2026-07-31T14:11:00.524790Z", - "execution_profile": "codex-default", - "fallback_models": [ - "gpt-5.4" - ], - "operation_role": "worker", - "pr_number": "1497", - "requested_model": "gpt-5.5", - "runner": "reusable-codex-run", - "schema": "langsmith-fleet/v1", - "selected_model": "", - "selection_reason": "" -} diff --git a/tests/test_alert_engine.py b/tests/test_alert_engine.py index eddcbb98..d1f5987a 100644 --- a/tests/test_alert_engine.py +++ b/tests/test_alert_engine.py @@ -150,10 +150,23 @@ def test_alert_engine_supports_news_and_crowded_trade_conditions(tmp_path): 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() diff --git a/tests/test_conviction_flow.py b/tests/test_conviction_flow.py index ad25123f..6ecbc144 100644 --- a/tests/test_conviction_flow.py +++ b/tests/test_conviction_flow.py @@ -53,10 +53,11 @@ def test_similarity_crowding_counts_only_connected_managers(): ) conn.executemany( "INSERT INTO manager_similarity VALUES (?, ?, ?)", - [(1, 2, 0.8), (1, 3, 0.7), (2, 3, 0.9), (3, 4, 0.2)], + [(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] + 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) == [] diff --git a/tests/test_crowded_contrarian.py b/tests/test_crowded_contrarian.py index ba957a46..0fd95239 100644 --- a/tests/test_crowded_contrarian.py +++ b/tests/test_crowded_contrarian.py @@ -362,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): diff --git a/tests/test_manager_similarity_flow.py b/tests/test_manager_similarity_flow.py index 857bf437..bcfdfb4d 100644 --- a/tests/test_manager_similarity_flow.py +++ b/tests/test_manager_similarity_flow.py @@ -42,6 +42,7 @@ def test_similarity_uses_union_denominator_and_latest_filings(): def test_cosine_similarity_matches_hand_computed_vectors(): assert cosine_similarity([1.0, 1.0], [1.0, 0.0]) == pytest.approx(1 / 2**0.5) assert cosine_similarity([0.0, 0.0], [1.0, 0.0]) is None + assert cosine_similarity([float("nan")], [1.0]) is None def test_similarity_rebuild_rolls_back_on_insert_failure(): @@ -131,4 +132,5 @@ def test_similar_manager_endpoint_orders_canonical_pairs_and_returns_404(tmp_pat cosine_response = asyncio.run(_get_similar_manager(1, 10, basis="cosine")) assert cosine_response.json()["items"][0]["manager_id"] == 2 assert cosine_response.json()["items"][0]["score"] == 0.8 + assert asyncio.run(_get_similar_manager(1, 10, basis="invalid")).status_code == 400 assert asyncio.run(_get_similar_manager(999, 10)).status_code == 404