"""analyze_cluster — discover groups that aren't already known (§0.9). Fourth ML analytical tool. Scales the selected numeric columns, fits KMeans for every k in a small range, keeps the k with the best silhouette, and returns one profile row per cluster plus a scatter coloured by cluster ("sepaket"). **Cluster vs segmentation — the distinction the 2026-09-15 discussion drew, and the reason both tools exist.** Segmentation is *by rules*: the criteria are already known ("split spend into tiers"), so the buckets are deterministic. Clustering is for when *"kita enggak tahu sama sekali"* — the groups are discovered from the data. Picking the wrong one of the two is the failure mode the tool descriptions have to prevent, which is why `dont_use_when` names the other explicitly on both sides. **k is chosen, not guessed.** Silhouette is computed for each candidate k and the best wins; the score is returned so the answer can say how well-separated the clusters actually are, and a weak score is a caveat rather than a hidden footnote. Clusters are descriptive — they describe how rows group, not why. Follows docs/tools/ANALYTICAL_TOOLS_CONVENTIONS.md: compute is synchronous except for the `await run_fit` offload; `data` is never self-fetched (Pattern A); errors escape to the invoker's never-throw seam; output is deterministic (fixed random_state). Mirrors analyze_anomaly's shape so the ML tools stay uniform. """ from __future__ import annotations import numpy as np import pandas as pd from pydantic import BaseModel, Field from src.tools.analytics.charts import cluster_scatter_chart from src.tools.analytics.descriptive import ColumnNotFoundError from src.tools.analytics.temporal import _clean from src.tools.ml.executor import run_fit from src.tools.spec_builder import build_description # KMeans on a handful of rows describes noise, not structure. _MIN_ROWS = 20 # Candidate k range (§0.9): small enough to stay interpretable in a narrative. _K_MIN = 2 _K_MAX = 8 # Feature cap (§0.9 R-5), consistent with analyze_importance. _MAX_FEATURES = 20 # Fixed so the same data always yields the same clusters (family-B determinism). _RANDOM_STATE = 42 # A silhouette at or below this means the groups barely separate — say so. _WEAK_SILHOUETTE = 0.25 # Chart column name (internal; never collides with user data). _CLUSTER_COL = "cluster" class NotEnoughDataError(ValueError): """Too few complete rows to cluster (error_code NOT_ENOUGH_DATA).""" class EmptyColumnSelectionError(ValueError): """column_ids was empty (error_code EMPTY_COLUMN_SELECTION).""" class ClusterInput(BaseModel): # `data` FIRST — property order follows field order (planner-prompt readability). data: str = Field(..., description="Placeholder ${t} of the upstream table (Pattern A).") column_ids: list[str] = Field( ..., description="NUMERIC columns describing each row — the traits to group on." ) k: int | None = Field( None, description=f"Number of clusters; omit to choose automatically between " f"{_K_MIN} and {_K_MAX} by silhouette score.", ) DESCRIPTION = build_description( summary="Discover natural groups in the rows and profile each one.", use_when="the user asks what kinds/groups/segments exist WITHOUT saying how to " "split them — 'ada berapa kelompok pelanggan', 'what types of units do we have', " "'kelompokkan berdasarkan pola pemakaian'.", dont_use_when=[ "the user already states the split rule ('tier by spend', 'bucket by age') " "-> analyze_segmentation (that is by rules; this discovers the groups)", "the user wants unusual individual rows -> analyze_anomaly", "the user wants totals per an EXISTING category column -> analyze_aggregate", ], output="the chosen k with its silhouette score, one profile row per cluster " "(size, share, per-column means), and a scatter coloured by cluster.", examples=[ "ada berapa kelompok pelanggan berdasarkan demografi?", "what natural groups are in this data?", "kelompokkan unit berdasarkan pola pemakaiannya", ], ) def _fit_kmeans(x: np.ndarray, k_values: list[int], random_state: int): """Scale once, fit KMeans for each candidate k, return the best by silhouette. Imported lazily so app startup doesn't pay for scikit-learn unless a clustering analysis runs. Runs inside run_fit's thread — this is the whole search, so the per-fit timeout covers every candidate k, not just one. Returns (labels, best_k, best_score, {k: score}). """ from sklearn.cluster import KMeans from sklearn.metrics import silhouette_score from sklearn.preprocessing import StandardScaler # Unsupervised: there is no holdout, so scaling on all rows leaks nothing. x_scaled = StandardScaler().fit_transform(x) scores: dict[int, float] = {} best = (None, 0, -1.0) # (labels, k, score) for k in k_values: labels = KMeans(n_clusters=k, n_init=10, random_state=random_state).fit_predict(x_scaled) if len(set(labels)) < 2: # degenerate fit — silhouette is undefined continue score = float(silhouette_score(x_scaled, labels)) scores[k] = round(score, 4) if score > best[2]: best = (labels, k, score) if best[0] is None: # Every candidate collapsed to one cluster: fall back to the smallest k so the # caller still gets a grouping, with the weak-separation caveat attached. labels = KMeans(n_clusters=_K_MIN, n_init=10, random_state=random_state).fit_predict( x_scaled ) return labels, _K_MIN, 0.0, scores return best[0], best[1], best[2], scores def _profiles(frame: pd.DataFrame, labels: np.ndarray, column_ids: list[str]) -> list[dict]: """One row per cluster: size, share of rows, and the mean of each column.""" total = len(labels) out: list[dict] = [] for label in sorted(set(labels.tolist())): members = frame[labels == label] profile: dict[str, object] = { "cluster": int(label), "size": int(len(members)), "share": round(len(members) / total, 4) if total else 0.0, } for col in column_ids: profile[col] = _clean(members[col].mean()) out.append(profile) return out async def analyze_cluster( df: pd.DataFrame, column_ids: list[str], k: int | None = None, ) -> dict[str, object]: """Group rows by `column_ids` with KMeans, choosing k by silhouette when unset. Async ONLY so the fit search can be offloaded via `run_fit`; everything else is synchronous compute. `data` is materialised upstream (Pattern A) and arrives as `df`. Returns a dict with: method, column_ids, n_rows, k, requested_k (what the caller asked for; differs from k only when a fallback overrode an explicit k — D-10), `quality` (silhouette + floor + weak flag), silhouette_by_k, `clusters` (one profile per group), `caveats`, and `charts` (the reserved auto-chart key). Raises (all wrapped by the invoker's never-throw seam): ColumnNotFoundError, EmptyColumnSelectionError, NotEnoughDataError. """ if not column_ids: raise EmptyColumnSelectionError("column_ids must name at least one numeric column") missing = [c for c in column_ids if c not in df.columns] if missing: raise ColumnNotFoundError(f"columns not found: {missing}") capped = column_ids[:_MAX_FEATURES] features = df[capped].apply(pd.to_numeric, errors="coerce") frame = features.dropna() if frame.shape[0] < _MIN_ROWS: raise NotEnoughDataError( f"need >= {_MIN_ROWS} rows with numeric values in all of {capped} " f"to cluster, got {frame.shape[0]}" ) # Never ask for more clusters than rows, and honour an explicit k as-is. ceiling = min(_K_MAX, frame.shape[0] - 1) k_values = [k] if k else list(range(_K_MIN, max(_K_MIN, ceiling) + 1)) labels, chosen_k, silhouette, by_k = await run_fit( _fit_kmeans, frame.to_numpy(dtype=float), k_values, _RANDOM_STATE, label="analyze_cluster" ) # D-10: an explicit k that every candidate collapsed on is silently replaced by the # fallback's _K_MIN. Detect it (explicit k, but a different k came back) and say so, # instead of the old line that claimed the fallback's k "was given". overridden = k is not None and chosen_k != k caveats = [ "Clusters are descriptive — they show how rows group together, not why, and " "they are not business categories until someone names them.", ] if overridden: caveats.append( f"Requested k={k} could not form {k} separable clusters (every candidate " f"collapsed to a single group), so it fell back to k={chosen_k}; " f"silhouette = {silhouette:.3f}." ) else: caveats.append( f"k={chosen_k} " + ("was given" if k else "chosen by silhouette") + f"; silhouette = {silhouette:.3f} (1.0 = perfectly separated)." ) if silhouette <= _WEAK_SILHOUETTE: caveats.append( "Separation is weak — these groups overlap substantially, so treat the " "split as tentative." ) if len(column_ids) > _MAX_FEATURES: caveats.append( f"Only the first {_MAX_FEATURES} of {len(column_ids)} columns were used." ) chart_df = frame.assign(**{_CLUSTER_COL: labels}) # Two real columns, not PCA components: a user can read "PA vs MTTR" and cannot # read "PC1 vs PC2". With a single column, plot it against the row index. if len(capped) >= 2: chart = cluster_scatter_chart(chart_df, capped[0], capped[1], _CLUSTER_COL) else: chart_df = chart_df.assign(__row__=range(len(chart_df))) chart = cluster_scatter_chart(chart_df, "__row__", capped[0], _CLUSTER_COL) return { "method": "kmeans", "column_ids": capped, "n_rows": int(frame.shape[0]), "k": int(chosen_k), # What the caller asked for (None = auto) vs what `k` ended up — they differ only # when the fallback overrode an explicit k (D-10). "requested_k": k, # Reserved `quality` block (conventions §3) — see analyze_importance. "quality": { "metric": "silhouette", "value": round(float(silhouette), 4), "floor": _WEAK_SILHOUETTE, "weak": silhouette <= _WEAK_SILHOUETTE, }, "silhouette_by_k": by_k, "clusters": _profiles(frame, labels, capped), "caveats": caveats, "charts": [chart], }