Aurelius / core /discovery.py
murtaza-2007
Rebuild Aurelius as a multi-source graph intelligence engine
75ce203
Raw
History Blame Contribute Delete
9.97 kB
"""Aurelius core β€” relationship scoring and hidden-connection discovery.
Two query types beyond pathfinding, both source-agnostic (they speak only
the GraphSource protocol, so they work identically on live Wikipedia and
on an ingested protein or finance graph):
relate(a, b) β€” how strongly are two nodes connected, and through what?
Evidence: direct edges, directed 2-hop paths (a→x→b),
co-targets (aβ†’x←b), co-sources (xβ†’a, xβ†’b when backlinks
exist), and embedding similarity. Returns a 0-100 strength
plus the actual intermediaries, not just a number.
discover(a) β€” Swanson ABC literature-based discovery: rank nodes C that
share many intermediaries B with A (A→B→C) but have NO
direct A→C edge. High bridge-count + high embedding
similarity + no direct link = a candidate hidden
connection. On stored sources the similarity half uses
the fused [text ; Ξ±Β·struct] vector, which is what lets a
structurally-close-but-textually-far node surface.
"""
from __future__ import annotations
import asyncio
import numpy as np
from .embedding import EmbeddingCache, cosine_similarity
from .representation import fuse
from .source import GraphSource
from .types import NodeRef
# Caps that keep live-mode (API-backed) discovery bounded.
_RELATE_NEIGHBOR_CAP = 400
_DISCOVER_BRIDGE_CAP = 20 # B nodes expanded (concurrently) per discover
_DISCOVER_CANDIDATES = 800 # C pool cap before ranking
_DISCOVER_EMBED_CAP = 200 # candidates text-embedded in live mode
def _store_of(source: GraphSource):
"""The GraphStore behind an ingested adapter, or None for live ones."""
return getattr(source, "store", None) if hasattr(source, "ingested") else None
async def _fused_embedding(source: GraphSource, ref: NodeRef,
emb_cache: EmbeddingCache) -> np.ndarray | None:
"""Fused vector from the store when available, else live text embed."""
store = _store_of(source)
if store is not None:
t = store.get_embedding(source.name, ref.id, "text")
s = store.get_embedding(source.name, ref.id, "struct")
if t is not None or s is not None:
return fuse(t, s, source.name)
info = await source.node_info(ref, rich=True)
vecs = await emb_cache.embed([ref.key()], [info.text or ref.title])
return vecs[0] if vecs and vecs[0].size else None
# ══════════════════════════════════════════════════════════════
# relate(a, b)
# ══════════════════════════════════════════════════════════════
async def relate(source: GraphSource, a_query: str, b_query: str) -> dict:
a = await source.resolve(a_query)
b = await source.resolve(b_query)
if not a or not b:
missing = a_query if not a else b_query
return {"error": f"Cannot find: '{missing}'"}
nb_a, nb_b = await asyncio.gather(
source.neighbors(a), source.neighbors(b))
out_a = {e.dst.id: e.dst for e in nb_a[:_RELATE_NEIGHBOR_CAP]}
out_b = {e.dst.id: e.dst for e in nb_b[:_RELATE_NEIGHBOR_CAP]}
in_a: dict[str, NodeRef] = {}
in_b: dict[str, NodeRef] = {}
if source.supports_backlinks:
bk_a, bk_b = await asyncio.gather(
source.back_neighbors(a), source.back_neighbors(b))
in_a = {e.src.id: e.src for e in bk_a}
in_b = {e.src.id: e.src for e in bk_b}
direct_ab = b.id in out_a
direct_ba = a.id in out_b
# Directed 2-hop a→x→b: x is an out-neighbor of a AND an in-neighbor
# of b (or, without backlinks, unverifiable β€” skipped).
paths_ab = [out_a[x] for x in (set(out_a) & set(in_b))] if in_b else []
paths_ba = [out_b[x] for x in (set(out_b) & set(in_a))] if in_a else []
co_targets = [out_a[x] for x in (set(out_a) & set(out_b))] # aβ†’x←b
co_sources = [in_a[x] for x in (set(in_a) & set(in_b))] # x→a, x→b
emb_cache = EmbeddingCache()
ea, eb = await asyncio.gather(
_fused_embedding(source, a, emb_cache),
_fused_embedding(source, b, emb_cache))
sim = cosine_similarity(ea, eb)
# Composite strength: direct edges dominate, then 2-hop evidence
# (saturating), then shared-neighbourhood evidence, then similarity.
def _sat(count: int, scale: float) -> float:
return 1.0 - float(np.exp(-count / scale))
strength = (
(0.35 if (direct_ab or direct_ba) else 0.0)
+ 0.30 * _sat(len(paths_ab) + len(paths_ba), 5.0)
+ 0.20 * _sat(len(co_targets) + len(co_sources), 20.0)
+ 0.15 * max(0.0, sim)
)
def _refs(refs: list[NodeRef], cap: int = 12) -> list[dict]:
return [{"id": r.id, "title": r.title} for r in refs[:cap]]
return {
"source": source.name,
"a": {"id": a.id, "title": a.title},
"b": {"id": b.id, "title": b.title},
"direct": {"a_to_b": direct_ab, "b_to_a": direct_ba},
"paths_a_to_b": _refs(paths_ab),
"paths_b_to_a": _refs(paths_ba),
"n_paths": len(paths_ab) + len(paths_ba),
"co_targets": _refs(co_targets),
"n_co_targets": len(co_targets),
"co_sources": _refs(co_sources),
"n_co_sources": len(co_sources),
"similarity": round(sim, 3),
"strength": round(100 * min(1.0, strength), 1),
}
# ══════════════════════════════════════════════════════════════
# discover(a)
# ══════════════════════════════════════════════════════════════
async def discover(source: GraphSource, a_query: str, k: int = 12) -> dict:
a = await source.resolve(a_query)
if not a:
return {"error": f"Cannot find: '{a_query}'"}
emb_cache = EmbeddingCache()
ea = await _fused_embedding(source, a, emb_cache)
nb_a = await source.neighbors(a)
direct: dict[str, NodeRef] = {e.dst.id: e.dst for e in nb_a}
if not direct:
return {"error": f"'{a.title}' has no outbound links to walk."}
# Choose the B set: rank a's neighbours by text similarity to a so the
# bridges we expand are the *relevant* ones, then expand concurrently.
b_refs = list(direct.values())
if len(b_refs) > _DISCOVER_BRIDGE_CAP and ea is not None:
infos = await source.node_infos(b_refs)
keys = [r.key() for r in b_refs]
embs = await emb_cache.embed(
keys, [i.text or r.title for i, r in zip(infos, b_refs)])
# Text-vs-fused dims can differ (stored fused vectors are longer);
# compare on the shared text prefix length.
d = min(ea.shape[0], embs[0].shape[0]) if embs[0].size else 0
sims = [cosine_similarity(e[:d], ea[:d]) if e.size else 0.0
for e in embs]
order = np.argsort(sims)[::-1]
b_refs = [b_refs[i] for i in order[:_DISCOVER_BRIDGE_CAP]]
else:
b_refs = b_refs[:_DISCOVER_BRIDGE_CAP]
results = await asyncio.gather(
*(source.neighbors(b) for b in b_refs), return_exceptions=True)
# Aggregate candidates C with their bridges B (A→B→C, no A→C).
bridges_of: dict[str, list[NodeRef]] = {}
cand_ref: dict[str, NodeRef] = {}
for b_ref, edges in zip(b_refs, results):
if isinstance(edges, BaseException):
continue
for e in edges:
c = e.dst
if c.id == a.id or c.id in direct:
continue
bridges_of.setdefault(c.id, []).append(b_ref)
cand_ref[c.id] = c
if len(cand_ref) >= _DISCOVER_CANDIDATES * 4:
break
if not cand_ref:
return {"a": {"id": a.id, "title": a.title}, "source": source.name,
"candidates": []}
# Rank: bridge support first, then embedding similarity on the top pool.
pool = sorted(cand_ref, key=lambda c: -len(bridges_of[c]))[:_DISCOVER_CANDIDATES]
store = _store_of(source)
sims: dict[str, float] = {}
if store is not None and ea is not None:
for c in pool:
t = store.get_embedding(source.name, c, "text")
s = store.get_embedding(source.name, c, "struct")
ec = fuse(t, s, source.name)
d = min(ea.shape[0], ec.shape[0]) if ec is not None else 0
sims[c] = cosine_similarity(ec[:d], ea[:d]) if d else 0.0
elif ea is not None:
head = pool[:_DISCOVER_EMBED_CAP]
refs = [cand_ref[c] for c in head]
infos = await source.node_infos(refs)
embs = await emb_cache.embed(
[r.key() for r in refs],
[i.text or r.title for i, r in zip(infos, refs)])
d0 = ea.shape[0]
for c, e in zip(head, embs):
d = min(d0, e.shape[0]) if e.size else 0
sims[c] = cosine_similarity(e[:d], ea[:d]) if d else 0.0
max_bridges = max(len(bridges_of[c]) for c in pool)
scored = []
for c in pool:
support = len(bridges_of[c]) / max_bridges
scored.append((0.55 * support + 0.45 * max(0.0, sims.get(c, 0.0)), c))
scored.sort(reverse=True)
candidates = [{
"id": c,
"title": cand_ref[c].title,
"score": round(100 * sc, 1),
"n_bridges": len(bridges_of[c]),
"similarity": round(sims.get(c, 0.0), 3),
"bridges": [{"id": b.id, "title": b.title}
for b in bridges_of[c][:6]],
} for sc, c in scored[:k]]
return {"a": {"id": a.id, "title": a.title}, "source": source.name,
"candidates": candidates}