Team Ai
Apppublic

CodeChamp95/topic_modelling

sourceHugging Faceupdated 4mo agoView on Hugging Face
0likes
agent.py781 linesDownload Raw Back to root
1"""2agent.py — BERTopic Thematic Analysis Agent3Braun & Clarke (2006) six-phase methodology implemented as a ReAct agent4using LangGraph, ChatMistralAI, and MemorySaver.5"""6 7from __future__ import annotations8 9import json10import logging11import re12import time13from pathlib import Path14from typing import Any, Generator15 16logger = logging.getLogger(__name__)17 18from langchain_mistralai import ChatMistralAI19from langgraph.checkpoint.memory import MemorySaver20from langgraph.prebuilt import create_react_agent21 22from tools import (23    load_scopus_csv,24    run_bertopic_discovery,25    label_topics_with_llm,26    consolidate_into_themes,27    compare_with_taxonomy,28    generate_comparison_csv,29    export_narrative,30)31 32# ──────────────────────────────────────────────────────────────────────────────33# Artifact paths (shared across phases)34# ──────────────────────────────────────────────────────────────────────────────35 36ARTIFACTS_DIR   = Path("artifacts")37ARTIFACTS_DIR.mkdir(exist_ok=True)38 39_LOADED_DATA    = str(ARTIFACTS_DIR / "loaded_data.json")40_SUMMARIES      = str(ARTIFACTS_DIR / "summaries.json")41_EMB            = str(ARTIFACTS_DIR / "emb.npy")42_LABELS         = str(ARTIFACTS_DIR / "topic_labels.json")43_THEMES         = str(ARTIFACTS_DIR / "themes.json")44_TAXONOMY       = str(ARTIFACTS_DIR / "taxonomy_mapping.json")45_COMPARISON_CSV = str(ARTIFACTS_DIR / "abstract_vs_title_comparison.csv")46_NARRATIVE      = str(ARTIFACTS_DIR / "section7_narrative.txt")47 48# ──────────────────────────────────────────────────────────────────────────────49# Rate-limit retry config50# ──────────────────────────────────────────────────────────────────────────────51 52_RL_MAX_RETRIES  = 4                        # max automatic retries on 42953_RL_BACKOFF_SECS = [15, 30, 60, 120]        # wait before each retry attempt54 55 56def _is_rate_limit(exc: Exception) -> bool:57    """Return True when *exc* is an HTTP 429 rate-limit error from any client."""58    # httpx.HTTPStatusError carries a .response attribute59    resp = getattr(exc, "response", None)60    if resp is not None and getattr(resp, "status_code", None) == 429:61        return True62    # Fallback: inspect the string representation (handles wrapped exceptions)63    s = str(exc).lower()64    return "429" in s and ("rate limit" in s or "rate_limit" in s or "rate_limited" in s)65 66 67# ──────────────────────────────────────────────────────────────────────────────68# System Prompt69# ──────────────────────────────────────────────────────────────────────────────70 71SYSTEM_PROMPT = """72╔══════════════════════════════════════════════════════════════════════════════╗73║          COMPUTATIONAL THEMATIC ANALYSIS AGENT  —  SYSTEM PROMPT           ║74╚══════════════════════════════════════════════════════════════════════════════╝75 76━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━77ROLE78━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━79You are a computational thematic analysis expert trained in the Braun & Clarke80(2006) six-phase framework for rigorous qualitative and mixed-methods research.81You specialise in applying BERTopic-based semantic clustering to academic82literature corpora, with deep expertise in:83 84  • Systematic literature review methodology85  • Sentence-level semantic embedding and agglomerative clustering86  • LLM-assisted topic labelling and theme consolidation87  • PAJAIS (Pacific-Asia Journal of the Association for Information Systems)88    25-category research taxonomy alignment89  • Transparent, reproducible, human-in-the-loop analytical pipelines90 91Your outputs are used in peer-reviewed academic research. Precision,92methodological rigour, and faithful adherence to the B&C (2006) phases are93non-negotiable.94 95━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━96CRITICAL RULES  (must be followed without exception)97━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━981.  ONE PHASE PER MESSAGE. Complete exactly one B&C phase per conversational99    turn. Never skip ahead or combine phases in a single response.100 1012.  ALL APPROVALS VIA REVIEW TABLE — NEVER VIA CHAT. You must NEVER ask the102    user to approve, reject, or rename topics in free-text chat. Every approval103    workflow must go through the Gradio review table. After populating the table104    you must STOP and wait for the user to click "Submit Review".105 1063.  STOP GATES ARE MANDATORY. At the end of Phases 2, 3, 4, and 5.5 you must107    output the exact STOP phrase:108        ⏸ STOP GATE — awaiting your review table submission to continue.109    Do not proceed until the user's next message contains review data.110 1114.  NEVER HALLUCINATE TOOL RESULTS. If a tool call fails, report the exact112    error verbatim and ask the user how to proceed. Do not invent file paths,113    cluster counts, or topic labels.114 1155.  COLUMN DISCIPLINE. Never include "Author Keywords" in any clustering run.116    Use only the columns specified in RUN_CONFIGS: Abstract (abstract run) or117    Title (title run).118 1196.  ARTEFACT HYGIENE. Every tool saves files to the artifacts/ directory.120    Always pass the exact saved_path returned by a previous tool to the next121    tool. Never guess or construct file paths manually.122 1237.  STREAMING DISCIPLINE. Yield one streamed chunk per reasoning step so the124    Gradio UI can update the phase progress bar in real time.125 126━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━127TOOLS  (7 available)128━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━1291.  load_scopus_csv(csv_path, run_mode)130      Load a Scopus-exported CSV. Counts papers and sentences. Applies131      boilerplate regex filter. Saves loaded_data.json. Use in Phase 1.132 1332.  run_bertopic_discovery(loaded_data_path)134      Embeds sentences with all-MiniLM-L6-v2 (normalize_embeddings=True).135      Clusters with AgglomerativeClustering(metric=cosine, threshold=0.7).136      No UMAP. Finds 5 nearest centroid sentences per cluster. Generates137      4 Plotly charts. Saves summaries.json + emb.npy. Use in Phase 2.138 1393.  label_topics_with_llm(summaries_path, top_n)140      Sends top-N topics (max 100) to Mistral via PromptTemplate +141      JsonOutputParser. Returns short labels and descriptions. Saves142      topic_labels.json. Use in Phase 2 after discovery.143 1444.  consolidate_into_themes(labels_path, summaries_path, emb_path, approved_groups)145      Merges approved topic groups into named themes. Recomputes centroids.146      Saves themes.json. Use in Phase 3 after the review table is submitted.147 1485.  compare_with_taxonomy(themes_path)149      Maps consolidated themes to PAJAIS 25 categories via Mistral.150      Returns confidence scores and rationale. Saves taxonomy_mapping.json.151      Use in Phase 5.5.152 1536.  generate_comparison_csv(csv_path, taxonomy_path)154      Produces abstract vs title side-by-side CSV with PAJAIS categories155      and confidence scores. Saves abstract_vs_title_comparison.csv.156      Use in Phase 6.157 1587.  export_narrative(taxonomy_path)159      Generates a ~500-word Section 7 (Discussion & Implications) as160      flowing academic prose via Mistral. Saves section7_narrative.txt.161      Use in Phase 6.162 163━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━164BRAUN & CLARKE (2006) SIX-PHASE PROTOCOL165━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━166 167┌─────────────────────────────────────────────────────────────────────────────┐168│  PHASE 1 — FAMILIARISATION WITH THE DATA                                   │169└─────────────────────────────────────────────────────────────────────────────┘170Objective: Immerse in the corpus. Understand its scope, structure, and quality.171 172Instructions:173  a. Call load_scopus_csv(csv_path=<user_provided>, run_mode=<"abstract"|"title">).174  b. Display the returned statistics in a clear summary:175       • Total papers loaded176       • Total sentences extracted177       • Sentences remaining after boilerplate filtering178       • Column(s) used179       • Run mode (abstract / title)180  c. Comment briefly on data quality: density, likely noise level, any181     column mapping issues detected.182  d. STOP. Do not proceed to Phase 2 until the user explicitly confirms183     they are satisfied with the loaded data.184 185Output template:186  📂 Phase 1 Complete — Familiarisation187  ─────────────────────────────────────188  Papers:              {N}189  Sentences extracted: {S}190  After filtering:     {F}191  Column used:         {C}192  Run mode:            {M}193  [Quality commentary]194  ✅ Ready for Phase 2. Reply "proceed" to start Initial Coding.195 196━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━197 198┌─────────────────────────────────────────────────────────────────────────────┐199│  PHASE 2 — GENERATING INITIAL CODES                                        │200└─────────────────────────────────────────────────────────────────────────────┘201Objective: Produce a full set of atomic semantic codes from the corpus.202 203Instructions:204  a. Call run_bertopic_discovery(loaded_data_path=artifacts/loaded_data.json).205     Report: number of clusters found, total sentences clustered, chart paths.206  b. Call label_topics_with_llm(summaries_path=artifacts/summaries.json, top_n=100).207     Report: number of topics labelled.208  c. Populate the Gradio review table with ALL labelled topics. Each row must209     contain:210       • #           — topic_id (integer)211       • Topic Label — LLM-generated label212       • Top Evidence — first centroid sentence (truncated to 120 chars)213       • Sentences   — cluster size214       • Papers      — estimated paper count (size ÷ avg sentences per paper)215       • Approve     — default True216       • Rename To   — empty (user fills)217       • Reasoning   — empty (user fills)218  d. Present the 4 Plotly charts by referencing their file paths.219  e. Explain to the user:220       • Check "Approve" for topics to keep; uncheck to discard.221       • Fill "Rename To" with a preferred label; leave blank to keep LLM label.222       • Optionally note merging intentions in "Reasoning".223       • Topics with the same "Reasoning" group tag will be merged in Phase 3.224 225⏸ STOP GATE — awaiting your review table submission to continue.226 227━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━228 229┌─────────────────────────────────────────────────────────────────────────────┐230│  PHASE 3 — SEARCHING FOR THEMES                                            │231└─────────────────────────────────────────────────────────────────────────────┘232Objective: Collate approved codes into candidate themes.233 234Instructions:235  a. Parse the submitted review table. Extract:236       • Approved topic IDs (Approve == True)237       • Rename mappings (Rename To != "")238       • Merge groups (topics sharing the same Reasoning tag)239  b. Construct approved_groups: a JSON list of lists, where each inner list240     contains the topic_ids belonging to one theme. Topics with a shared241     Reasoning tag form one group. Approved topics with no Reasoning tag242     each form a singleton group.243  c. Call consolidate_into_themes(244         labels_path=artifacts/topic_labels.json,245         summaries_path=artifacts/summaries.json,246         emb_path=artifacts/emb.npy,247         approved_groups=<constructed JSON string>248     ).249  d. Display a theme summary table:250       Theme # | Theme Label | Topics Merged | Total Sentences251 252⏸ STOP GATE — awaiting your review table submission to continue.253 254━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━255 256┌─────────────────────────────────────────────────────────────────────────────┐257│  PHASE 4 — REVIEWING THEMES / SATURATION CHECK                             │258└─────────────────────────────────────────────────────────────────────────────┘259Objective: Assess whether themes are internally coherent and collectively260exhaustive. Check corpus coverage.261 262Instructions:263  a. Load artifacts/themes.json (already created by Phase 3).264  b. Compute and display a saturation report:265       • Total sentences covered by all themes vs. total corpus sentences266       • Coverage percentage267       • Theme coherence flag: warn if any theme covers < 1 % of corpus268       • Overlap flag: warn if any two themes share > 30 % vocabulary269       (Vocabulary overlap is approximated by comparing top-evidence sentences270        using word-level Jaccard similarity — compute in Python, no tool call.)271  c. Populate the review table again with the THEME list (not topic list):272       • #           — theme_id273       • Topic Label — current theme_label274       • Top Evidence — first top_evidence sentence (120 chars)275       • Sentences   — total_size276       • Papers      — estimated277       • Approve     — default True278       • Rename To   — user may provide final name279       • Reasoning   — any split/merge instructions280  d. Ask the user to confirm themes, request splits/merges, or rename.281 282⏸ STOP GATE — awaiting your review table submission to continue.283 284━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━285 286┌─────────────────────────────────────────────────────────────────────────────┐287│  PHASE 5 — DEFINING AND NAMING THEMES                                      │288└─────────────────────────────────────────────────────────────────────────────┘289Objective: Produce final, publication-ready theme names and definitions.290 291Instructions:292  a. Apply all renames from the Phase 4 review table to artifacts/themes.json293     in memory (update theme_label field for each theme_id where Rename To294     is non-empty).295  b. For each finalised theme, generate a two-sentence academic definition296     grounded in the top_evidence sentences. Output this as a numbered list.297  c. Confirm the final theme set to the user in a clean summary:298       Theme # | Final Name | Definition (2 sentences) | Sentence Count299  d. STOP. Ask the user to confirm the final names before PAJAIS mapping.300     ✅ Reply "proceed to taxonomy" to continue to Phase 5.5.301 302━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━303 304┌─────────────────────────────────────────────────────────────────────────────┐305│  PHASE 5.5 — PAJAIS TAXONOMY ALIGNMENT                                     │306└─────────────────────────────────────────────────────────────────────────────┘307Objective: Map each finalised theme to the PAJAIS 25-category taxonomy.308 309Instructions:310  a. Call compare_with_taxonomy(themes_path=artifacts/themes.json).311  b. Display the mapping results in a structured table:312       Theme Name | PAJAIS Category | Confidence | Rationale313  c. Highlight any themes with confidence < 0.5 as requiring manual review.314  d. Note any PAJAIS categories not covered by the corpus (research gaps).315  e. Populate the review table with the mapping results:316       • #           — theme_id317       • Topic Label — theme_label → PAJAIS category318       • Top Evidence — rationale (truncated)319       • Approve     — default True (uncheck to override mapping)320       • Rename To   — alternative PAJAIS category if user disagrees321       • Reasoning   — free notes322 323⏸ STOP GATE — awaiting your review table submission to continue.324 325━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━326 327┌─────────────────────────────────────────────────────────────────────────────┐328│  PHASE 6 — PRODUCING THE REPORT                                            │329└─────────────────────────────────────────────────────────────────────────────┘330Objective: Generate all final deliverables.331 332Instructions:333  a. Apply any PAJAIS category overrides from the Phase 5.5 review table334     to artifacts/taxonomy_mapping.json in memory.335  b. Call generate_comparison_csv(336         csv_path=<original CSV path>,337         taxonomy_path=artifacts/taxonomy_mapping.json338     ).339     Report: row count, file path.340  c. Call export_narrative(taxonomy_path=artifacts/taxonomy_mapping.json).341     Report: word count, file path, first 150 chars of preview.342  d. Present a final deliverables checklist:343       ✅ artifacts/summaries.json            — raw cluster summaries344       ✅ artifacts/topic_labels.json         — LLM-generated labels345       ✅ artifacts/themes.json               — consolidated themes346       ✅ artifacts/taxonomy_mapping.json     — PAJAIS alignment347       ✅ artifacts/abstract_vs_title_comparison.csv348       ✅ artifacts/section7_narrative.txt    — ~500-word Section 7349       ✅ artifacts/chart_cluster_sizes.html350       ✅ artifacts/chart_pca_scatter.html351       ✅ artifacts/chart_top10_pie.html352       ✅ artifacts/chart_centroid_heatmap.html353  e. Congratulate the user and offer to re-run in title mode for comparison.354 355━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━356END OF SYSTEM PROMPT357━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━358""".strip()359 360# ──────────────────────────────────────────────────────────────────────────────361# Tool registry362# ──────────────────────────────────────────────────────────────────────────────363 364TOOLS = [365    load_scopus_csv,366    run_bertopic_discovery,367    label_topics_with_llm,368    consolidate_into_themes,369    compare_with_taxonomy,370    generate_comparison_csv,371    export_narrative,372]373 374# ──────────────────────────────────────────────────────────────────────────────375# Phase detection helpers376# ──────────────────────────────────────────────────────────────────────────────377 378_PHASE_PATTERNS = {379    "loading":    re.compile(r"load_scopus_csv|loaded_data", re.I),380    "embedding":  re.compile(r"run_bertopic_discovery|embedding", re.I),381    "clustering": re.compile(r"summaries\.json|n_topics|clusters found", re.I),382    "labelling":  re.compile(r"label_topics_with_llm|topic_labels", re.I),383    "review":     re.compile(r"STOP GATE|review table|submit review", re.I),384    "done":       re.compile(r"section7_narrative|deliverables checklist", re.I),385}386 387 388def _detect_phase(text: str) -> str:389    """Return the most specific pipeline phase detectable from agent output."""390    matched = list(filter(391        lambda kv: kv[1].search(text),392        _PHASE_PATTERNS.items(),393    ))394    return matched[-1][0] if matched else "idle"395 396 397# ──────────────────────────────────────────────────────────────────────────────398# Review-table row builder399# ──────────────────────────────────────────────────────────────────────────────400 401_TABLE_COLS = ["#", "Topic Label", "Top Evidence", "Sentences", "Papers", "Approve", "Rename To", "Reasoning"]402 403 404def _topic_to_row(topic: dict, papers_per_sent: float = 0.2) -> list:405    """Convert a topic/theme dict to a review-table row."""406    evidence = (topic.get("top_evidence") or [""])[0]407    return [408        topic.get("topic_id", topic.get("theme_id", 0)),409        topic.get("label", topic.get("theme_label", "")),410        evidence[:120],411        topic.get("size", topic.get("total_size", 0)),412        round(topic.get("size", topic.get("total_size", 0)) * papers_per_sent),413        True,414        "",415        "",416    ]417 418 419def _build_review_rows(path: str, id_key: str = "topic_id") -> list[list]:420    """Load a JSON artefact and convert every entry to a review-table row."""421    records = json.loads(Path(path).read_text())422    return list(map(_topic_to_row, records))423 424 425# ──────────────────────────────────────────────────────────────────────────────426# Approved-groups extractor  (called in handle_review for Phase 2 → 3)427# ──────────────────────────────────────────────────────────────────────────────428 429def _extract_approved_groups(rows: list[list]) -> str:430    """431    Parse review-table rows into approved_groups JSON string.432 433    Groups are formed by the Reasoning field value:434      • Rows sharing a non-empty Reasoning tag → merged into one group435      • Approved rows with empty Reasoning    → singleton group each436      • Unapproved rows (Approve == False)    → discarded437    """438    approved = list(filter(lambda r: r[5] is True or r[5] == "True" or r[5] == 1, rows))439 440    tagged    = list(filter(lambda r: str(r[7]).strip(), approved))441    untagged  = list(filter(lambda r: not str(r[7]).strip(), approved))442 443    # Group tagged rows by their Reasoning value444    reasoning_vals = list(set(map(lambda r: str(r[7]).strip(), tagged)))445 446    tagged_groups = list(map(447        lambda tag: list(map(448            lambda r: int(r[0]),449            filter(lambda r: str(r[7]).strip() == tag, tagged),450        )),451        reasoning_vals,452    ))453 454    singleton_groups = list(map(lambda r: [int(r[0])], untagged))455 456    all_groups = tagged_groups + singleton_groups457    return json.dumps(all_groups)458 459 460# ──────────────────────────────────────────────────────────────────────────────461# BERTopicAgent  — the class consumed by app.py462# ──────────────────────────────────────────────────────────────────────────────463 464class BERTopicAgent:465    """466    Wraps a LangGraph ReAct agent and exposes the two generator methods467    expected by the Gradio front-end:468 469        handle_message(message, history, csv_path) → yields 5-tuple470        handle_review(table_rows, history)          → yields 4-tuple471    """472 473    # Gradio 5-tuple: (history, phase, charts_dict, downloads_list, topic_rows)474    # Gradio 4-tuple: (history, phase, charts_dict, downloads_list)475 476    def __init__(self) -> None:477        self.phase       = "idle"478        self._charts:    dict[str, str] = {}479        self._downloads: list[str]      = []480        self._csv_path:  str | None     = None481 482        self._llm       = ChatMistralAI(483            model="mistral-large-latest",484            temperature=0.2,485            streaming=True,486        )487        self._memory    = MemorySaver()488 489        # handle_tool_error is no longer a @tool() decorator argument in490        # newer LangChain versions — set it directly on each tool object.491        for t in TOOLS:492            t.handle_tool_error = True493 494        self._graph     = create_react_agent(495            model=self._llm,496            tools=TOOLS,497            checkpointer=self._memory,498            prompt=SYSTEM_PROMPT,499        )500        self._thread_id = "bc2006-session-1"501 502    # ── internal ──────────────────────────────────────────────────────────────503 504    def _config(self) -> dict:505        return {"configurable": {"thread_id": self._thread_id}}506 507    def _update_charts(self, text: str) -> None:508        """Scan agent text for chart file paths and register them."""509        found = re.findall(r'artifacts/chart_[a-z_]+\.html', text)510        label_map = {511            "chart_cluster_sizes":    "Cluster Sizes",512            "chart_pca_scatter":      "PCA Scatter",513            "chart_top10_pie":        "Top-10 Pie",514            "chart_centroid_heatmap": "Centroid Heatmap",515        }516        list(map(517            lambda p: self._charts.__setitem__(518                label_map.get(Path(p).stem, Path(p).stem), p519            ),520            found,521        ))522 523    def _update_downloads(self, text: str) -> None:524        """Scan agent text for downloadable artefact paths."""525        # Added ?: to make it a non-capturing group526        found = re.findall(r'artifacts/[\w_]+\.(?:json|csv|txt|npy|html)', text)527        new   = list(filter(lambda p: p not in self._downloads, found))528        self._downloads.extend(new)529 530    def _accumulate_stream(531        self,532        stream: Any,533        history: list,534        user_msg: str,535        _attempt: int = 0,536    ) -> Generator[tuple, None, None]:537        """538        Consume a LangGraph stream, yielding Gradio 5-tuples incrementally.539 540        Automatically retries on HTTP 429 (rate-limit) errors with exponential541        back-off up to _RL_MAX_RETRIES times.  All other exceptions surface a542        friendly error message in the chat rather than crashing the generator.543        """544        accumulated = ""545 546        try:547            for chunk in stream:548                # LangGraph yields dicts keyed by node name549                node_output = (550                    chunk.get("agent") or551                    chunk.get("tools") or552                    {}553                )554                messages = node_output.get("messages", [])555 556                text_delta = "".join(list(map(557                    lambda m: getattr(m, "content", "") if hasattr(m, "content") else "",558                    messages,559                )))560 561                accumulated += text_delta562                self._update_charts(accumulated)563                self._update_downloads(accumulated)564                self.phase = _detect_phase(accumulated)565 566                updated_history = history + [[user_msg, accumulated]] if accumulated else history567 568                yield (569                    updated_history,570                    self.phase,571                    dict(self._charts),572                    list(self._downloads),573                    [],          # topic_rows populated in final yield574                )575 576            # ── Success: final yield with review-table rows ───────────────────577            topic_rows = self._latest_review_rows()578            yield (579                history + [[user_msg, accumulated]],580                self.phase,581                dict(self._charts),582                list(self._downloads),583                topic_rows,584            )585 586        except Exception as exc:  # noqa: BLE001587            if _is_rate_limit(exc) and _attempt < _RL_MAX_RETRIES:588                # ── Rate-limit: back off then restart the stream ──────────────589                wait = _RL_BACKOFF_SECS[_attempt]590                notice = (591                    f"\n\n⏳ **Mistral rate limit hit** — waiting **{wait}s** "592                    f"then retrying automatically "593                    f"(attempt {_attempt + 1}/{_RL_MAX_RETRIES})…"594                )595                logger.warning("Rate limit 429 on attempt %d; sleeping %ds", _attempt, wait)596                yield (597                    history + [[user_msg, accumulated + notice]],598                    "idle",599                    dict(self._charts),600                    list(self._downloads),601                    [],602                )603                time.sleep(wait)604 605                # Rebuild the stream — MemorySaver resumes from last checkpoint606                new_stream = self._graph.stream(607                    {"messages": [{"role": "user", "content": user_msg}]},608                    config=self._config(),609                    stream_mode="updates",610                )611                yield from self._accumulate_stream(612                    new_stream, history, user_msg, _attempt=_attempt + 1613                )614 615            else:616                # ── Non-retryable error: surface gracefully in chat ───────────617                if _is_rate_limit(exc):618                    err_header = (619                        f"❌ **Rate limit persists after {_RL_MAX_RETRIES} retries.**\n"620                        "Please wait a few minutes before sending another message."621                    )622                else:623                    err_header = f"❌ **API / tool error:** `{type(exc).__name__}: {exc}`"624 625                logger.exception("Unhandled error in _accumulate_stream (attempt %d)", _attempt)626                self.phase = "idle"627                yield (628                    history + [[user_msg, accumulated + f"\n\n{err_header}"]],629                    "idle",630                    dict(self._charts),631                    list(self._downloads),632                    self._latest_review_rows(),   # keep existing table intact633                )634 635    def _latest_review_rows(self) -> list[list]:636        """Return review rows from the most recently produced artefact."""637        candidates = [638            (_THEMES,   "theme_id"),639            (_LABELS,   "topic_id"),640            (_SUMMARIES,"topic_id"),641        ]642        existing = list(filter(lambda t: Path(t[0]).exists(), candidates))643        return _build_review_rows(*existing[0]) if existing else []644 645    # ── public API ────────────────────────────────────────────────────────────646 647    def handle_message(648        self,649        message:  str,650        history:  list,651        csv_path: str | None = None,652    ) -> Generator[tuple, None, None]:653        """654        Send a user message to the ReAct agent and stream back Gradio 5-tuples.655 656        Yields: (history, phase, charts_dict, downloads_list, topic_rows)657        """658        self._csv_path = csv_path or self._csv_path659 660        # Inject CSV path into message so the agent can reference it661        enriched = (662            f"{message}\n\n[SYSTEM CONTEXT] CSV path: {self._csv_path}"663            if self._csv_path and "csv" not in message.lower()664            else message665        )666 667        stream = self._graph.stream(668            {"messages": [{"role": "user", "content": enriched}]},669            config=self._config(),670            stream_mode="updates",671        )672 673        yield from self._accumulate_stream(stream, history, message)674 675    def handle_review(676        self,677        table_data: list,678        history:    list,679    ) -> Generator[tuple, None, None]:680        """681        Process a submitted review table and advance to the next B&C phase.682 683        The table rows are serialised to JSON and injected as a structured684        user message so the agent can parse approvals, renames, and groups.685 686        Yields: (history, phase, charts_dict, downloads_list)687        """688        approved_groups = _extract_approved_groups(table_data)689 690        review_payload = json.dumps({691            "event":           "review_submitted",692            "rows":            table_data,693            "approved_groups": json.loads(approved_groups),694            "approved_count":  len(json.loads(approved_groups)),695        }, ensure_ascii=False, indent=2)696 697        review_message = (698            f"The user has submitted the review table. "699            f"Approved groups: {approved_groups}. "700            f"Full payload:\n{review_payload}\n\n"701            f"Please continue to the next B&C phase now."702        )703 704        def _make_review_stream() -> Any:705            return self._graph.stream(706                {"messages": [{"role": "user", "content": review_message}]},707                config=self._config(),708                stream_mode="updates",709            )710 711        accumulated = ""712        attempt     = 0713 714        while True:715            current_stream = _make_review_stream()716            try:717                for chunk in current_stream:718                    node_output = chunk.get("agent") or chunk.get("tools") or {}719                    messages    = node_output.get("messages", [])720 721                    text_delta = "".join(list(map(722                        lambda m: getattr(m, "content", "") if hasattr(m, "content") else "",723                        messages,724                    )))725 726                    accumulated += text_delta727                    self._update_charts(accumulated)728                    self._update_downloads(accumulated)729                    self.phase = _detect_phase(accumulated)730 731                    yield (732                        history + [["[Review submitted]", accumulated]],733                        self.phase,734                        dict(self._charts),735                        list(self._downloads),736                    )737 738                # ── Success ───────────────────────────────────────────────────739                yield (740                    history + [["[Review submitted]", accumulated]],741                    self.phase,742                    dict(self._charts),743                    list(self._downloads),744                )745                break   # exit retry loop746 747            except Exception as exc:  # noqa: BLE001748                if _is_rate_limit(exc) and attempt < _RL_MAX_RETRIES:749                    wait = _RL_BACKOFF_SECS[attempt]750                    notice = (751                        f"\n\n⏳ **Rate limit hit** — waiting **{wait}s** "752                        f"then retrying (attempt {attempt + 1}/{_RL_MAX_RETRIES})…"753                    )754                    logger.warning("Rate limit 429 in handle_review attempt %d; sleeping %ds", attempt, wait)755                    yield (756                        history + [["[Review submitted]", accumulated + notice]],757                        "idle",758                        dict(self._charts),759                        list(self._downloads),760                    )761                    time.sleep(wait)762                    attempt += 1763                    # Loop re-creates the stream via _make_review_stream()764                else:765                    if _is_rate_limit(exc):766                        err = (767                            f"❌ **Rate limit persists after {_RL_MAX_RETRIES} retries.**\n"768                            "Please wait a few minutes before trying again."769                        )770                    else:771                        err = f"❌ **Error processing review:** `{type(exc).__name__}: {exc}`"772 773                    logger.exception("Unhandled error in handle_review (attempt %d)", attempt)774                    self.phase = "idle"775                    yield (776                        history + [["[Review submitted]", accumulated + f"\n\n{err}"]],777                        "idle",778                        dict(self._charts),779                        list(self._downloads),780                    )781                    break