From 264f3a8f442ba9b6ac5b2219fe2af09750401591 Mon Sep 17 00:00:00 2001 From: jraedisch Date: Mon, 13 Apr 2026 17:29:18 +0200 Subject: [PATCH] for_you_v2: add account affinity, credibility scoring, and z-score ranking Adds medoid-based account affinity, profile credibility heuristics, per-cluster z-score normalization, and cluster diversity malus. Single best post per window instead of top-N. Co-Authored-By: Claude Opus 4.6 (1M context) --- experiments/for_you_v2.py | 214 ++++++++++++++++++++++++++++++++------ 1 file changed, 185 insertions(+), 29 deletions(-) diff --git a/experiments/for_you_v2.py b/experiments/for_you_v2.py index 45ca664..75675e9 100644 --- a/experiments/for_you_v2.py +++ b/experiments/for_you_v2.py @@ -9,7 +9,7 @@ ensures distinctive interests score higher than catch-all clusters. Works with or without bearer token (128d vs 768d embeddings). Usage: - python3 for_you_v2.py [token] --window 60 [--top 3] + python3 for_you_v2.py [token] --window 60 """ import argparse @@ -53,6 +53,75 @@ def divepool_search(token, query, limit=100, did=None, cluster=False, return json.loads(resp.read()) +DIVEPOOL_MEDOIDS = "https://divepool.social/api/v1/medoids" + + +def fetch_account_medoids(token, dids): + """Batch-fetch up to 3 cluster medoids per account for up to 10 DIDs.""" + payload = {"dids": dids[:25]} + data = json.dumps(payload).encode() + headers = {"Content-Type": "application/json"} + if token: + headers["Authorization"] = f"Bearer {token}" + req = urllib.request.Request(DIVEPOOL_MEDOIDS, data=data, headers=headers, + method="POST") + try: + with urllib.request.urlopen(req, timeout=10) as resp: + return json.loads(resp.read()) + except Exception: + return {"accounts": {}} + + +def fetch_bsky_profile(did): + """Fetch public profile stats. Returns (followers, following, posts) or None.""" + url = f"{BSKY_API}/app.bsky.actor.getProfile?actor={urllib.parse.quote(did)}" + req = urllib.request.Request(url, headers={"Accept": "application/json"}) + try: + with urllib.request.urlopen(req, timeout=5) as resp: + data = json.loads(resp.read()) + return ( + data.get("followersCount", 0), + data.get("followsCount", 0), + data.get("postsCount", 0), + ) + except Exception: + return None + + +def account_credibility(followers, following, posts): + """Score from -1 (spam) through 0 (neutral/big) to +1 (small & legit). + Penalizes spam bots; boosts approachable small accounts; neutral for large ones.""" + if posts == 0: + return -0.5 # zero-post accounts are almost always fake + + ratio = followers / posts + + # Base spam signal: followers/posts ratio + # ratio 0.08 → -0.6, ratio 0.3 → 0.0, ratio 0.75 → +0.5, ratio 1+ → +0.6 + base = (ratio - 0.3) / (ratio + 0.2) + + # Suspicious following patterns (follow-back bots, fake follower farms) + if followers > 0: + ff_ratio = following / followers + if ff_ratio > 10: + base = min(base, -0.5) + elif ff_ratio > 5: + base *= 0.5 + elif following > 100: + base = min(base, -0.5) + + # Small & legit bonus: accounts under ~200 followers with healthy ratios + # are the approachable long-tail we want to surface + if ratio >= 0.3 and followers < 200: + base += 0.3 # boost small legit accounts + + # Big accounts: cap positive score — they don't need help + if followers > 1000: + base = min(base, 0.1) + + return max(-1.0, min(1.0, base)) + + def resolve_post_text(did, collection, rkey): uri = f"at://{did}/{collection}/{rkey}" url = f"{BSKY_API}/app.bsky.feed.getPostThread?uri={urllib.parse.quote(uri)}&depth=0" @@ -307,8 +376,6 @@ def main(): parser.add_argument("did", help="Your AT Protocol DID") parser.add_argument("--window", type=int, required=True, help="Selection window in seconds") - parser.add_argument("--top", type=int, default=3, - help="Max posts to show per window (default: 3)") args = parser.parse_args() token, user_did = args.token, args.did @@ -378,8 +445,18 @@ def main(): print(f" {label}: weight {w.mean():.2f} (min={w.min():.2f}, max={w.max():.2f})", file=sys.stderr, flush=True) + # ── User cluster centroids (for account affinity scoring) ──────────── + user_centroids = [] + for cid in sorted(cluster_id_to_label): + mask = [i for i, c in enumerate(ref_cluster_ids) if c == cid] + if mask: + centroid = ref_matrix[mask].mean(axis=0) + centroid = centroid / np.linalg.norm(centroid) + user_centroids.append(centroid) + user_centroid_matrix = np.stack(user_centroids) # (K, dim) + print(f"\n{len(ref_vecs)} refs, {n_clusters} clusters, " - f"window={args.window}s, top={args.top}", + f"window={args.window}s", file=sys.stderr, flush=True) # ── Step 3: Stream firehose, score every post, pick best per window ── @@ -400,8 +477,34 @@ def main(): shown_scores = [] window_start = time.time() - # Each candidate: (comp_score, best_sim, best_label, other_labels, did, col, rkey) - candidates = [] + candidates = [] # (z, comp_score, best_label, other_labels, did, col, rkey) + cluster_shown_count = defaultdict(int) # cumulative malus: how many times each cluster won + MALUS = 0.5 # z-score penalty per prior win — after 2 wins, cluster needs z 1.0 higher + AFFINITY_BONUS = 1.0 # max z-score bonus for account affinity (scaled by similarity) + CREDIBILITY_WEIGHT = 1.5 # scales credibility score (now -1 to +1) + account_affinity_cache = {} # did -> affinity score + profile_cache = {} # did -> (followers, following, posts) + + # Per-cluster running stats for z-score normalization (Welford's online algo) + cluster_stats = defaultdict(lambda: [0, 0.0, 0.0]) # [count, mean, M2] + total_scored = 0 + WARMUP = 100 # min scored posts before z-scores kick in + + def update_cluster_stats(label, score): + s = cluster_stats[label] + s[0] += 1 + delta = score - s[1] + s[1] += delta / s[0] + s[2] += delta * (score - s[1]) + + def z_score(label, score): + s = cluster_stats[label] + if s[0] < 5: + return score # not enough data, fall back to raw + std = (s[2] / s[0]) ** 0.5 + if std < 1e-6: + return 0.0 + return (score - s[1]) / std def is_spam(did, text): if did_freq[did] > 5: @@ -414,20 +517,24 @@ def main(): return True return False - def show(handle, text, rkey, comp_score, best_sim, label): + def show(handle, text, rkey, adj_z, comp_score, label, affinity=0.0, cred=0.5): link = bsky_link(handle, rkey) - gold = "" - if len(shown_scores) >= 10: - mean = np.mean(shown_scores) - std = np.std(shown_scores) - if std > 0 and comp_score >= mean + 1.5 * std: - gold = " *" - shown_scores.append(comp_score) - print(f"[{comp_score:.2f}{gold} {label}]\n{text}\n{link} — @{handle}\n", flush=True) - stats.record_shown(best_sim) + aff_str = f" aff={affinity:.2f}" if affinity > 0 else "" + cred_str = f" cred={cred:+.2f}" if abs(cred) > 0.1 else "" + age_sec = time.time() - tid_to_timestamp(rkey) + if age_sec < 60: + age_str = f"{age_sec:.0f}s" + elif age_sec < 3600: + age_str = f"{age_sec / 60:.0f}m" + else: + age_str = f"{age_sec / 3600:.1f}h" + print(f"[z={adj_z:.1f} {comp_score:.2f}{aff_str}{cred_str} {age_str} {label}]\n" + f"{text}\n{link} — @{handle}\n", + flush=True) + stats.record_shown(comp_score) def flush_window(): - """Pick the best candidates from the window, resolve and show them.""" + """Pick the best candidate by z-score, resolve and show.""" nonlocal window_start, candidates stats.windows_elapsed += 1 @@ -437,11 +544,53 @@ def main(): candidates = [] return - candidates.sort(key=lambda c: -c[0]) - picks = candidates[:args.top] - - shown_in_window = 0 - for comp_score, best_sim, best_label, other_labels, did, col, rkey in picks: + # Fetch account medoids + profiles for top candidates + candidate_dids = list(dict.fromkeys( + d for _, _, _, _, d, _, _ in sorted(candidates, key=lambda c: -c[0])[:20] + )) + uncached_medoids = [d for d in candidate_dids[:25] if d not in account_affinity_cache] + if uncached_medoids: + resp = fetch_account_medoids(token, uncached_medoids) + for did_key, acct in resp.get("accounts", {}).items(): + medoids = acct.get("medoids", []) + if not medoids: + account_affinity_cache[did_key] = 0.0 + continue + med_vecs = [] + for m in medoids: + e = m.get("embedding", []) + if e: + med_vecs.append(normalize(np.array(e, dtype=np.float32))) + if not med_vecs: + account_affinity_cache[did_key] = 0.0 + continue + med_matrix = np.stack(med_vecs) # (up to 3, dim) + sims = med_matrix @ user_centroid_matrix.T # (3, K) + account_affinity_cache[did_key] = float(sims.max()) + for d in uncached_medoids: + if d not in account_affinity_cache: + account_affinity_cache[d] = 0.0 + + uncached_profiles = [d for d in candidate_dids[:10] if d not in profile_cache] + for d in uncached_profiles: + prof = fetch_bsky_profile(d) + profile_cache[d] = prof if prof else (0, 0, 0) + + # Apply cumulative malus + account affinity + credibility + adjusted = [] + for z, cs, bl, ol, d, c, rk in candidates: + affinity = account_affinity_cache.get(d, 0.0) + prof = profile_cache.get(d) + cred = account_credibility(*prof) if prof else 0.5 + adj = (z + - MALUS * cluster_shown_count[bl] + + AFFINITY_BONUS * affinity + + CREDIBILITY_WEIGHT * cred) + adjusted.append((adj, z, cs, bl, ol, d, c, rk, cred)) + adjusted.sort(key=lambda c: -c[0]) # sort by adjusted z-score + + shown = False + for adj_z, z, comp_score, best_label, other_labels, did, col, rkey, cred in adjusted: handle, text = resolve_post_text(did, col, rkey) if not handle or not text: stats.fail_resolve += 1 @@ -455,10 +604,13 @@ def main(): else: label = best_label stats.cluster_hits[best_label] += 1 - show(handle, text, rkey, comp_score, best_sim, label) - shown_in_window += 1 + cluster_shown_count[best_label] += 1 + affinity = account_affinity_cache.get(did, 0.0) + show(handle, text, rkey, adj_z, comp_score, label, affinity, cred) + shown = True + break # only show top 1 - if shown_in_window == 0: + if not shown: stats.windows_empty += 1 # One-line summary after flush @@ -507,17 +659,21 @@ def main(): stats.record_score(best_sim) best_label = ref_labels[best_idx] + update_cluster_stats(best_label, comp_score) + total_scored += 1 + + z = z_score(best_label, comp_score) if total_scored >= WARMUP else comp_score + other_labels = [cluster_id_to_label.get(cid, "?") for cid in list(hit_clusters)[:2]] candidates.append(( - comp_score, best_sim, best_label, + z, comp_score, best_label, other_labels, did, col, rkey, )) # Keep candidate list bounded - max_keep = args.top * 5 - if len(candidates) > max_keep * 2: + if len(candidates) > 20: candidates.sort(key=lambda c: -c[0]) - candidates = candidates[:max_keep] + candidates = candidates[:10] # Check if window is up now = time.time() -- 2.51.2