CodeChamp95/topic_modelling
0
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