"""Persistent identity for complaint clusters. This is the memory. The problem this solves: every run re-groups the feedback from scratch, so this week's "ads" cluster is technically a brand new object that merely resembles last week's. Without a way to tell they are the same problem, you can never say whether anything is getting better or worse. The fix: give every problem a permanent ID and a history. After each run, every cluster gets a fingerprint saved to disk alongside its ID. On the next run, new clusters are matched against those saved fingerprints. A close enough match means it is the same problem, so it inherits the ID and this run gets appended to its history. No match means it is genuinely new. Also tracks the awkward cases honestly: problems that split in two, problems that merge, and problems that stop being mentioned at all. """ import json from datetime import date from pathlib import Path import numpy as np REGISTRY_VERSION = 1 MATCH_THRESHOLD = 0.82 # how similar two fingerprints must be to count as "same" DORMANT_AFTER = 3 # runs without a mention before a problem is called dormant def _fingerprint_text(cluster): """The text that represents a cluster's identity.""" parts = [cluster["name"]] parts += [e["text"] for e in cluster.get("examples", [])[:3]] return " ".join(parts) HASH_DIMS = 512 def _hashed_bag(texts, dims=HASH_DIMS): """Fixed-size word fingerprint that is stable across separate runs. A plain bag of words cannot be used here: the vocabulary would be rebuilt from each run's data, so the vectors would have different lengths and meanings every time and nothing would ever match its past self. Hashing each word into a fixed number of slots sidesteps that. crc32 is used rather than Python's built-in hash() because the latter is randomized per process, which would silently break matching between runs. """ import zlib matrix = np.zeros((len(texts), dims)) for row, text in enumerate(texts): for word in text.lower().split(): if len(word) > 3: matrix[row, zlib.crc32(word.encode()) % dims] += 1 norms = np.linalg.norm(matrix, axis=1, keepdims=True) norms[norms == 0] = 1 return matrix / norms def fingerprints(clusters): """Vector per cluster, plus which method produced it. Fingerprints from different methods are never compared, since their numbers mean different things. """ texts = [_fingerprint_text(c) for c in clusters] try: import analyze_semantic return np.asarray(analyze_semantic.embed(texts)), "embedding" except Exception: return _hashed_bag(texts), "hashed" def load(path): path = Path(path) if not path.exists(): return {"version": REGISTRY_VERSION, "next_id": 1, "clusters": {}} return json.loads(path.read_text()) def save(registry, path): path = Path(path) path.parent.mkdir(parents=True, exist_ok=True) path.write_text(json.dumps(registry, indent=2)) def _new_id(registry): cluster_id = f"C{registry['next_id']:03d}" registry["next_id"] += 1 return cluster_id def reconcile(clusters, registry, run_date=None, threshold=MATCH_THRESHOLD): """Match this run's clusters to known problems. Returns (clusters, events). Matching is greedy and one to one: the strongest pair is matched first, and neither side can be matched twice. That keeps a single known problem from absorbing three of this run's clusters silently. """ run_date = run_date or date.today().isoformat() vectors, method = fingerprints(clusters) known_ids = list(registry["clusters"].keys()) known_vectors = ( np.asarray([registry["clusters"][k]["fingerprint"] for k in known_ids]) if known_ids else np.zeros((0, vectors.shape[1])) ) # Similarity only comparable if the fingerprints were built the same way. comparable = [ i for i, k in enumerate(known_ids) if registry["clusters"][k].get("method") == method and len(registry["clusters"][k]["fingerprint"]) == vectors.shape[1] ] pairs = [] if comparable: similarity = vectors @ known_vectors[comparable].T for new_index in range(len(clusters)): for col, known_index in enumerate(comparable): pairs.append( (float(similarity[new_index, col]), new_index, known_ids[known_index]) ) pairs.sort(reverse=True) matched_new, matched_known, assignment = set(), set(), {} for score, new_index, known_id in pairs: if score < threshold: break if new_index in matched_new or known_id in matched_known: continue matched_new.add(new_index) matched_known.add(known_id) assignment[new_index] = (known_id, score) events = {"new": [], "continuing": [], "dormant": [], "reawakened": []} for i, cluster in enumerate(clusters): if i in assignment: cluster_id, score = assignment[i] entry = registry["clusters"][cluster_id] was_absent = entry.get("absent_runs", 0) >= DORMANT_AFTER entry["absent_runs"] = 0 entry["label"] = cluster["name"] # labels drift; keep the latest entry["fingerprint"] = vectors[i].tolist() entry["last_seen"] = run_date cluster["match_score"] = round(score, 3) (events["reawakened"] if was_absent else events["continuing"]).append(cluster_id) else: cluster_id = _new_id(registry) registry["clusters"][cluster_id] = { "id": cluster_id, "label": cluster["name"], "fingerprint": vectors[i].tolist(), "method": method, "first_seen": run_date, "last_seen": run_date, "absent_runs": 0, "history": [], } cluster["match_score"] = None events["new"].append(cluster_id) entry = registry["clusters"][cluster_id] point = { "date": run_date, "count": cluster["count"], "share": cluster["share"], "sources": cluster["source_spread"], } # Re-running on the same day replaces that day rather than double counting. if entry["history"] and entry["history"][-1]["date"] == run_date: entry["history"][-1] = point else: entry["history"].append(point) cluster["cluster_id"] = cluster_id cluster["first_seen"] = entry["first_seen"] cluster["history"] = entry["history"] cluster["trend"] = trend_of(entry["history"]) # Anything known but not seen this run seen = {c["cluster_id"] for c in clusters} for cluster_id, entry in registry["clusters"].items(): if cluster_id not in seen: entry["absent_runs"] = entry.get("absent_runs", 0) + 1 if entry["absent_runs"] >= DORMANT_AFTER: events["dormant"].append(cluster_id) return clusters, events def trend_of(history, min_points=2): """Describe how a problem's volume is moving. Honest about small samples.""" if len(history) < min_points: # Same keys every time, so callers never have to check which shape # they got back. return { "direction": "new", "change_pct": None, "step_pct": None, "runs": len(history), } first_share = history[0]["share"] or 0.01 last_share = history[-1]["share"] change = 100 * (last_share - first_share) / first_share previous, latest = history[-2]["share"], history[-1]["share"] step = 100 * (latest - previous) / (previous or 0.01) if abs(step) < 10: direction = "steady" elif step > 0: direction = "growing" else: direction = "shrinking" return { "direction": direction, "change_pct": round(change, 1), "step_pct": round(step, 1), "runs": len(history), } def summarize(events, registry): lines = [] if events["new"]: lines.append(f" {len(events['new'])} new problem(s) first seen this run") if events["continuing"]: lines.append(f" {len(events['continuing'])} known problem(s) still active") if events["reawakened"]: labels = ", ".join( registry["clusters"][c]["label"][:35] for c in events["reawakened"] ) lines.append(f" reawakened after going quiet: {labels}") if events["dormant"]: labels = ", ".join( registry["clusters"][c]["label"][:35] for c in events["dormant"][:3] ) lines.append(f" {len(events['dormant'])} gone quiet ({labels})") return lines