File size: 41,396 Bytes
b221afb
 
 
 
 
 
 
 
 
aaa5147
b221afb
aaa5147
b221afb
 
 
aaa5147
 
b221afb
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
aaa5147
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
b221afb
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
247e2ce
 
 
 
 
 
b221afb
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
d89add7
 
b221afb
 
 
 
 
 
 
 
aaa5147
b221afb
 
 
aaa5147
 
 
 
b221afb
 
 
aaa5147
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
b221afb
aaa5147
b221afb
 
 
aaa5147
b221afb
 
aaa5147
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
b221afb
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
aaa5147
 
 
 
 
b221afb
 
aaa5147
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
"""
agent.py β€” BERTopic Thematic Analysis Agent
Braun & Clarke (2006) six-phase methodology implemented as a ReAct agent
using LangGraph, ChatMistralAI, and MemorySaver.
"""

from __future__ import annotations

import json
import logging
import re
import time
from pathlib import Path
from typing import Any, Generator

logger = logging.getLogger(__name__)

from langchain_mistralai import ChatMistralAI
from langgraph.checkpoint.memory import MemorySaver
from langgraph.prebuilt import create_react_agent

from tools import (
    load_scopus_csv,
    run_bertopic_discovery,
    label_topics_with_llm,
    consolidate_into_themes,
    compare_with_taxonomy,
    generate_comparison_csv,
    export_narrative,
)

# ──────────────────────────────────────────────────────────────────────────────
# Artifact paths (shared across phases)
# ──────────────────────────────────────────────────────────────────────────────

ARTIFACTS_DIR   = Path("artifacts")
ARTIFACTS_DIR.mkdir(exist_ok=True)

_LOADED_DATA    = str(ARTIFACTS_DIR / "loaded_data.json")
_SUMMARIES      = str(ARTIFACTS_DIR / "summaries.json")
_EMB            = str(ARTIFACTS_DIR / "emb.npy")
_LABELS         = str(ARTIFACTS_DIR / "topic_labels.json")
_THEMES         = str(ARTIFACTS_DIR / "themes.json")
_TAXONOMY       = str(ARTIFACTS_DIR / "taxonomy_mapping.json")
_COMPARISON_CSV = str(ARTIFACTS_DIR / "abstract_vs_title_comparison.csv")
_NARRATIVE      = str(ARTIFACTS_DIR / "section7_narrative.txt")

# ──────────────────────────────────────────────────────────────────────────────
# Rate-limit retry config
# ──────────────────────────────────────────────────────────────────────────────

_RL_MAX_RETRIES  = 4                        # max automatic retries on 429
_RL_BACKOFF_SECS = [15, 30, 60, 120]        # wait before each retry attempt


def _is_rate_limit(exc: Exception) -> bool:
    """Return True when *exc* is an HTTP 429 rate-limit error from any client."""
    # httpx.HTTPStatusError carries a .response attribute
    resp = getattr(exc, "response", None)
    if resp is not None and getattr(resp, "status_code", None) == 429:
        return True
    # Fallback: inspect the string representation (handles wrapped exceptions)
    s = str(exc).lower()
    return "429" in s and ("rate limit" in s or "rate_limit" in s or "rate_limited" in s)


# ──────────────────────────────────────────────────────────────────────────────
# System Prompt
# ──────────────────────────────────────────────────────────────────────────────

SYSTEM_PROMPT = """
╔══════════════════════════════════════════════════════════════════════════════╗
β•‘          COMPUTATIONAL THEMATIC ANALYSIS AGENT  β€”  SYSTEM PROMPT           β•‘
β•šβ•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
ROLE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
You are a computational thematic analysis expert trained in the Braun & Clarke
(2006) six-phase framework for rigorous qualitative and mixed-methods research.
You specialise in applying BERTopic-based semantic clustering to academic
literature corpora, with deep expertise in:

  β€’ Systematic literature review methodology
  β€’ Sentence-level semantic embedding and agglomerative clustering
  β€’ LLM-assisted topic labelling and theme consolidation
  β€’ PAJAIS (Pacific-Asia Journal of the Association for Information Systems)
    25-category research taxonomy alignment
  β€’ Transparent, reproducible, human-in-the-loop analytical pipelines

Your outputs are used in peer-reviewed academic research. Precision,
methodological rigour, and faithful adherence to the B&C (2006) phases are
non-negotiable.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
CRITICAL RULES  (must be followed without exception)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1.  ONE PHASE PER MESSAGE. Complete exactly one B&C phase per conversational
    turn. Never skip ahead or combine phases in a single response.

2.  ALL APPROVALS VIA REVIEW TABLE β€” NEVER VIA CHAT. You must NEVER ask the
    user to approve, reject, or rename topics in free-text chat. Every approval
    workflow must go through the Gradio review table. After populating the table
    you must STOP and wait for the user to click "Submit Review".

3.  STOP GATES ARE MANDATORY. At the end of Phases 2, 3, 4, and 5.5 you must
    output the exact STOP phrase:
        ⏸ STOP GATE β€” awaiting your review table submission to continue.
    Do not proceed until the user's next message contains review data.

4.  NEVER HALLUCINATE TOOL RESULTS. If a tool call fails, report the exact
    error verbatim and ask the user how to proceed. Do not invent file paths,
    cluster counts, or topic labels.

5.  COLUMN DISCIPLINE. Never include "Author Keywords" in any clustering run.
    Use only the columns specified in RUN_CONFIGS: Abstract (abstract run) or
    Title (title run).

6.  ARTEFACT HYGIENE. Every tool saves files to the artifacts/ directory.
    Always pass the exact saved_path returned by a previous tool to the next
    tool. Never guess or construct file paths manually.

7.  STREAMING DISCIPLINE. Yield one streamed chunk per reasoning step so the
    Gradio UI can update the phase progress bar in real time.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
TOOLS  (7 available)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1.  load_scopus_csv(csv_path, run_mode)
      Load a Scopus-exported CSV. Counts papers and sentences. Applies
      boilerplate regex filter. Saves loaded_data.json. Use in Phase 1.

2.  run_bertopic_discovery(loaded_data_path)
      Embeds sentences with all-MiniLM-L6-v2 (normalize_embeddings=True).
      Clusters with AgglomerativeClustering(metric=cosine, threshold=0.7).
      No UMAP. Finds 5 nearest centroid sentences per cluster. Generates
      4 Plotly charts. Saves summaries.json + emb.npy. Use in Phase 2.

3.  label_topics_with_llm(summaries_path, top_n)
      Sends top-N topics (max 100) to Mistral via PromptTemplate +
      JsonOutputParser. Returns short labels and descriptions. Saves
      topic_labels.json. Use in Phase 2 after discovery.

4.  consolidate_into_themes(labels_path, summaries_path, emb_path, approved_groups)
      Merges approved topic groups into named themes. Recomputes centroids.
      Saves themes.json. Use in Phase 3 after the review table is submitted.

5.  compare_with_taxonomy(themes_path)
      Maps consolidated themes to PAJAIS 25 categories via Mistral.
      Returns confidence scores and rationale. Saves taxonomy_mapping.json.
      Use in Phase 5.5.

6.  generate_comparison_csv(csv_path, taxonomy_path)
      Produces abstract vs title side-by-side CSV with PAJAIS categories
      and confidence scores. Saves abstract_vs_title_comparison.csv.
      Use in Phase 6.

7.  export_narrative(taxonomy_path)
      Generates a ~500-word Section 7 (Discussion & Implications) as
      flowing academic prose via Mistral. Saves section7_narrative.txt.
      Use in Phase 6.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
BRAUN & CLARKE (2006) SIX-PHASE PROTOCOL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚  PHASE 1 β€” FAMILIARISATION WITH THE DATA                                   β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Objective: Immerse in the corpus. Understand its scope, structure, and quality.

Instructions:
  a. Call load_scopus_csv(csv_path=<user_provided>, run_mode=<"abstract"|"title">).
  b. Display the returned statistics in a clear summary:
       β€’ Total papers loaded
       β€’ Total sentences extracted
       β€’ Sentences remaining after boilerplate filtering
       β€’ Column(s) used
       β€’ Run mode (abstract / title)
  c. Comment briefly on data quality: density, likely noise level, any
     column mapping issues detected.
  d. STOP. Do not proceed to Phase 2 until the user explicitly confirms
     they are satisfied with the loaded data.

Output template:
  πŸ“‚ Phase 1 Complete β€” Familiarisation
  ─────────────────────────────────────
  Papers:              {N}
  Sentences extracted: {S}
  After filtering:     {F}
  Column used:         {C}
  Run mode:            {M}
  [Quality commentary]
  βœ… Ready for Phase 2. Reply "proceed" to start Initial Coding.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚  PHASE 2 β€” GENERATING INITIAL CODES                                        β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Objective: Produce a full set of atomic semantic codes from the corpus.

Instructions:
  a. Call run_bertopic_discovery(loaded_data_path=artifacts/loaded_data.json).
     Report: number of clusters found, total sentences clustered, chart paths.
  b. Call label_topics_with_llm(summaries_path=artifacts/summaries.json, top_n=100).
     Report: number of topics labelled.
  c. Populate the Gradio review table with ALL labelled topics. Each row must
     contain:
       β€’ #           β€” topic_id (integer)
       β€’ Topic Label β€” LLM-generated label
       β€’ Top Evidence β€” first centroid sentence (truncated to 120 chars)
       β€’ Sentences   β€” cluster size
       β€’ Papers      β€” estimated paper count (size Γ· avg sentences per paper)
       β€’ Approve     β€” default True
       β€’ Rename To   β€” empty (user fills)
       β€’ Reasoning   β€” empty (user fills)
  d. Present the 4 Plotly charts by referencing their file paths.
  e. Explain to the user:
       β€’ Check "Approve" for topics to keep; uncheck to discard.
       β€’ Fill "Rename To" with a preferred label; leave blank to keep LLM label.
       β€’ Optionally note merging intentions in "Reasoning".
       β€’ Topics with the same "Reasoning" group tag will be merged in Phase 3.

⏸ STOP GATE β€” awaiting your review table submission to continue.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚  PHASE 3 β€” SEARCHING FOR THEMES                                            β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Objective: Collate approved codes into candidate themes.

Instructions:
  a. Parse the submitted review table. Extract:
       β€’ Approved topic IDs (Approve == True)
       β€’ Rename mappings (Rename To != "")
       β€’ Merge groups (topics sharing the same Reasoning tag)
  b. Construct approved_groups: a JSON list of lists, where each inner list
     contains the topic_ids belonging to one theme. Topics with a shared
     Reasoning tag form one group. Approved topics with no Reasoning tag
     each form a singleton group.
  c. Call consolidate_into_themes(
         labels_path=artifacts/topic_labels.json,
         summaries_path=artifacts/summaries.json,
         emb_path=artifacts/emb.npy,
         approved_groups=<constructed JSON string>
     ).
  d. Display a theme summary table:
       Theme # | Theme Label | Topics Merged | Total Sentences

⏸ STOP GATE β€” awaiting your review table submission to continue.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚  PHASE 4 β€” REVIEWING THEMES / SATURATION CHECK                             β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Objective: Assess whether themes are internally coherent and collectively
exhaustive. Check corpus coverage.

Instructions:
  a. Load artifacts/themes.json (already created by Phase 3).
  b. Compute and display a saturation report:
       β€’ Total sentences covered by all themes vs. total corpus sentences
       β€’ Coverage percentage
       β€’ Theme coherence flag: warn if any theme covers < 1 % of corpus
       β€’ Overlap flag: warn if any two themes share > 30 % vocabulary
       (Vocabulary overlap is approximated by comparing top-evidence sentences
        using word-level Jaccard similarity β€” compute in Python, no tool call.)
  c. Populate the review table again with the THEME list (not topic list):
       β€’ #           β€” theme_id
       β€’ Topic Label β€” current theme_label
       β€’ Top Evidence β€” first top_evidence sentence (120 chars)
       β€’ Sentences   β€” total_size
       β€’ Papers      β€” estimated
       β€’ Approve     β€” default True
       β€’ Rename To   β€” user may provide final name
       β€’ Reasoning   β€” any split/merge instructions
  d. Ask the user to confirm themes, request splits/merges, or rename.

⏸ STOP GATE β€” awaiting your review table submission to continue.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚  PHASE 5 β€” DEFINING AND NAMING THEMES                                      β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Objective: Produce final, publication-ready theme names and definitions.

Instructions:
  a. Apply all renames from the Phase 4 review table to artifacts/themes.json
     in memory (update theme_label field for each theme_id where Rename To
     is non-empty).
  b. For each finalised theme, generate a two-sentence academic definition
     grounded in the top_evidence sentences. Output this as a numbered list.
  c. Confirm the final theme set to the user in a clean summary:
       Theme # | Final Name | Definition (2 sentences) | Sentence Count
  d. STOP. Ask the user to confirm the final names before PAJAIS mapping.
     βœ… Reply "proceed to taxonomy" to continue to Phase 5.5.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚  PHASE 5.5 β€” PAJAIS TAXONOMY ALIGNMENT                                     β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Objective: Map each finalised theme to the PAJAIS 25-category taxonomy.

Instructions:
  a. Call compare_with_taxonomy(themes_path=artifacts/themes.json).
  b. Display the mapping results in a structured table:
       Theme Name | PAJAIS Category | Confidence | Rationale
  c. Highlight any themes with confidence < 0.5 as requiring manual review.
  d. Note any PAJAIS categories not covered by the corpus (research gaps).
  e. Populate the review table with the mapping results:
       β€’ #           β€” theme_id
       β€’ Topic Label β€” theme_label β†’ PAJAIS category
       β€’ Top Evidence β€” rationale (truncated)
       β€’ Approve     β€” default True (uncheck to override mapping)
       β€’ Rename To   β€” alternative PAJAIS category if user disagrees
       β€’ Reasoning   β€” free notes

⏸ STOP GATE β€” awaiting your review table submission to continue.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚  PHASE 6 β€” PRODUCING THE REPORT                                            β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Objective: Generate all final deliverables.

Instructions:
  a. Apply any PAJAIS category overrides from the Phase 5.5 review table
     to artifacts/taxonomy_mapping.json in memory.
  b. Call generate_comparison_csv(
         csv_path=<original CSV path>,
         taxonomy_path=artifacts/taxonomy_mapping.json
     ).
     Report: row count, file path.
  c. Call export_narrative(taxonomy_path=artifacts/taxonomy_mapping.json).
     Report: word count, file path, first 150 chars of preview.
  d. Present a final deliverables checklist:
       βœ… artifacts/summaries.json            β€” raw cluster summaries
       βœ… artifacts/topic_labels.json         β€” LLM-generated labels
       βœ… artifacts/themes.json               β€” consolidated themes
       βœ… artifacts/taxonomy_mapping.json     β€” PAJAIS alignment
       βœ… artifacts/abstract_vs_title_comparison.csv
       βœ… artifacts/section7_narrative.txt    β€” ~500-word Section 7
       βœ… artifacts/chart_cluster_sizes.html
       βœ… artifacts/chart_pca_scatter.html
       βœ… artifacts/chart_top10_pie.html
       βœ… artifacts/chart_centroid_heatmap.html
  e. Congratulate the user and offer to re-run in title mode for comparison.

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
END OF SYSTEM PROMPT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
""".strip()

# ──────────────────────────────────────────────────────────────────────────────
# Tool registry
# ──────────────────────────────────────────────────────────────────────────────

TOOLS = [
    load_scopus_csv,
    run_bertopic_discovery,
    label_topics_with_llm,
    consolidate_into_themes,
    compare_with_taxonomy,
    generate_comparison_csv,
    export_narrative,
]

# ──────────────────────────────────────────────────────────────────────────────
# Phase detection helpers
# ──────────────────────────────────────────────────────────────────────────────

_PHASE_PATTERNS = {
    "loading":    re.compile(r"load_scopus_csv|loaded_data", re.I),
    "embedding":  re.compile(r"run_bertopic_discovery|embedding", re.I),
    "clustering": re.compile(r"summaries\.json|n_topics|clusters found", re.I),
    "labelling":  re.compile(r"label_topics_with_llm|topic_labels", re.I),
    "review":     re.compile(r"STOP GATE|review table|submit review", re.I),
    "done":       re.compile(r"section7_narrative|deliverables checklist", re.I),
}


def _detect_phase(text: str) -> str:
    """Return the most specific pipeline phase detectable from agent output."""
    matched = list(filter(
        lambda kv: kv[1].search(text),
        _PHASE_PATTERNS.items(),
    ))
    return matched[-1][0] if matched else "idle"


# ──────────────────────────────────────────────────────────────────────────────
# Review-table row builder
# ──────────────────────────────────────────────────────────────────────────────

_TABLE_COLS = ["#", "Topic Label", "Top Evidence", "Sentences", "Papers", "Approve", "Rename To", "Reasoning"]


def _topic_to_row(topic: dict, papers_per_sent: float = 0.2) -> list:
    """Convert a topic/theme dict to a review-table row."""
    evidence = (topic.get("top_evidence") or [""])[0]
    return [
        topic.get("topic_id", topic.get("theme_id", 0)),
        topic.get("label", topic.get("theme_label", "")),
        evidence[:120],
        topic.get("size", topic.get("total_size", 0)),
        round(topic.get("size", topic.get("total_size", 0)) * papers_per_sent),
        True,
        "",
        "",
    ]


def _build_review_rows(path: str, id_key: str = "topic_id") -> list[list]:
    """Load a JSON artefact and convert every entry to a review-table row."""
    records = json.loads(Path(path).read_text())
    return list(map(_topic_to_row, records))


# ──────────────────────────────────────────────────────────────────────────────
# Approved-groups extractor  (called in handle_review for Phase 2 β†’ 3)
# ──────────────────────────────────────────────────────────────────────────────

def _extract_approved_groups(rows: list[list]) -> str:
    """
    Parse review-table rows into approved_groups JSON string.

    Groups are formed by the Reasoning field value:
      β€’ Rows sharing a non-empty Reasoning tag β†’ merged into one group
      β€’ Approved rows with empty Reasoning    β†’ singleton group each
      β€’ Unapproved rows (Approve == False)    β†’ discarded
    """
    approved = list(filter(lambda r: r[5] is True or r[5] == "True" or r[5] == 1, rows))

    tagged    = list(filter(lambda r: str(r[7]).strip(), approved))
    untagged  = list(filter(lambda r: not str(r[7]).strip(), approved))

    # Group tagged rows by their Reasoning value
    reasoning_vals = list(set(map(lambda r: str(r[7]).strip(), tagged)))

    tagged_groups = list(map(
        lambda tag: list(map(
            lambda r: int(r[0]),
            filter(lambda r: str(r[7]).strip() == tag, tagged),
        )),
        reasoning_vals,
    ))

    singleton_groups = list(map(lambda r: [int(r[0])], untagged))

    all_groups = tagged_groups + singleton_groups
    return json.dumps(all_groups)


# ──────────────────────────────────────────────────────────────────────────────
# BERTopicAgent  β€” the class consumed by app.py
# ──────────────────────────────────────────────────────────────────────────────

class BERTopicAgent:
    """
    Wraps a LangGraph ReAct agent and exposes the two generator methods
    expected by the Gradio front-end:

        handle_message(message, history, csv_path) β†’ yields 5-tuple
        handle_review(table_rows, history)          β†’ yields 4-tuple
    """

    # Gradio 5-tuple: (history, phase, charts_dict, downloads_list, topic_rows)
    # Gradio 4-tuple: (history, phase, charts_dict, downloads_list)

    def __init__(self) -> None:
        self.phase       = "idle"
        self._charts:    dict[str, str] = {}
        self._downloads: list[str]      = []
        self._csv_path:  str | None     = None

        self._llm       = ChatMistralAI(
            model="mistral-large-latest",
            temperature=0.2,
            streaming=True,
        )
        self._memory    = MemorySaver()

        # handle_tool_error is no longer a @tool() decorator argument in
        # newer LangChain versions β€” set it directly on each tool object.
        for t in TOOLS:
            t.handle_tool_error = True

        self._graph     = create_react_agent(
            model=self._llm,
            tools=TOOLS,
            checkpointer=self._memory,
            prompt=SYSTEM_PROMPT,
        )
        self._thread_id = "bc2006-session-1"

    # ── internal ──────────────────────────────────────────────────────────────

    def _config(self) -> dict:
        return {"configurable": {"thread_id": self._thread_id}}

    def _update_charts(self, text: str) -> None:
        """Scan agent text for chart file paths and register them."""
        found = re.findall(r'artifacts/chart_[a-z_]+\.html', text)
        label_map = {
            "chart_cluster_sizes":    "Cluster Sizes",
            "chart_pca_scatter":      "PCA Scatter",
            "chart_top10_pie":        "Top-10 Pie",
            "chart_centroid_heatmap": "Centroid Heatmap",
        }
        list(map(
            lambda p: self._charts.__setitem__(
                label_map.get(Path(p).stem, Path(p).stem), p
            ),
            found,
        ))

    def _update_downloads(self, text: str) -> None:
        """Scan agent text for downloadable artefact paths."""
        # Added ?: to make it a non-capturing group
        found = re.findall(r'artifacts/[\w_]+\.(?:json|csv|txt|npy|html)', text)
        new   = list(filter(lambda p: p not in self._downloads, found))
        self._downloads.extend(new)

    def _accumulate_stream(
        self,
        stream: Any,
        history: list,
        user_msg: str,
        _attempt: int = 0,
    ) -> Generator[tuple, None, None]:
        """
        Consume a LangGraph stream, yielding Gradio 5-tuples incrementally.

        Automatically retries on HTTP 429 (rate-limit) errors with exponential
        back-off up to _RL_MAX_RETRIES times.  All other exceptions surface a
        friendly error message in the chat rather than crashing the generator.
        """
        accumulated = ""

        try:
            for chunk in stream:
                # LangGraph yields dicts keyed by node name
                node_output = (
                    chunk.get("agent") or
                    chunk.get("tools") or
                    {}
                )
                messages = node_output.get("messages", [])

                text_delta = "".join(list(map(
                    lambda m: getattr(m, "content", "") if hasattr(m, "content") else "",
                    messages,
                )))

                accumulated += text_delta
                self._update_charts(accumulated)
                self._update_downloads(accumulated)
                self.phase = _detect_phase(accumulated)

                updated_history = history + [[user_msg, accumulated]] if accumulated else history

                yield (
                    updated_history,
                    self.phase,
                    dict(self._charts),
                    list(self._downloads),
                    [],          # topic_rows populated in final yield
                )

            # ── Success: final yield with review-table rows ───────────────────
            topic_rows = self._latest_review_rows()
            yield (
                history + [[user_msg, accumulated]],
                self.phase,
                dict(self._charts),
                list(self._downloads),
                topic_rows,
            )

        except Exception as exc:  # noqa: BLE001
            if _is_rate_limit(exc) and _attempt < _RL_MAX_RETRIES:
                # ── Rate-limit: back off then restart the stream ──────────────
                wait = _RL_BACKOFF_SECS[_attempt]
                notice = (
                    f"\n\n⏳ **Mistral rate limit hit** β€” waiting **{wait}s** "
                    f"then retrying automatically "
                    f"(attempt {_attempt + 1}/{_RL_MAX_RETRIES})…"
                )
                logger.warning("Rate limit 429 on attempt %d; sleeping %ds", _attempt, wait)
                yield (
                    history + [[user_msg, accumulated + notice]],
                    "idle",
                    dict(self._charts),
                    list(self._downloads),
                    [],
                )
                time.sleep(wait)

                # Rebuild the stream β€” MemorySaver resumes from last checkpoint
                new_stream = self._graph.stream(
                    {"messages": [{"role": "user", "content": user_msg}]},
                    config=self._config(),
                    stream_mode="updates",
                )
                yield from self._accumulate_stream(
                    new_stream, history, user_msg, _attempt=_attempt + 1
                )

            else:
                # ── Non-retryable error: surface gracefully in chat ───────────
                if _is_rate_limit(exc):
                    err_header = (
                        f"❌ **Rate limit persists after {_RL_MAX_RETRIES} retries.**\n"
                        "Please wait a few minutes before sending another message."
                    )
                else:
                    err_header = f"❌ **API / tool error:** `{type(exc).__name__}: {exc}`"

                logger.exception("Unhandled error in _accumulate_stream (attempt %d)", _attempt)
                self.phase = "idle"
                yield (
                    history + [[user_msg, accumulated + f"\n\n{err_header}"]],
                    "idle",
                    dict(self._charts),
                    list(self._downloads),
                    self._latest_review_rows(),   # keep existing table intact
                )

    def _latest_review_rows(self) -> list[list]:
        """Return review rows from the most recently produced artefact."""
        candidates = [
            (_THEMES,   "theme_id"),
            (_LABELS,   "topic_id"),
            (_SUMMARIES,"topic_id"),
        ]
        existing = list(filter(lambda t: Path(t[0]).exists(), candidates))
        return _build_review_rows(*existing[0]) if existing else []

    # ── public API ────────────────────────────────────────────────────────────

    def handle_message(
        self,
        message:  str,
        history:  list,
        csv_path: str | None = None,
    ) -> Generator[tuple, None, None]:
        """
        Send a user message to the ReAct agent and stream back Gradio 5-tuples.

        Yields: (history, phase, charts_dict, downloads_list, topic_rows)
        """
        self._csv_path = csv_path or self._csv_path

        # Inject CSV path into message so the agent can reference it
        enriched = (
            f"{message}\n\n[SYSTEM CONTEXT] CSV path: {self._csv_path}"
            if self._csv_path and "csv" not in message.lower()
            else message
        )

        stream = self._graph.stream(
            {"messages": [{"role": "user", "content": enriched}]},
            config=self._config(),
            stream_mode="updates",
        )

        yield from self._accumulate_stream(stream, history, message)

    def handle_review(
        self,
        table_data: list,
        history:    list,
    ) -> Generator[tuple, None, None]:
        """
        Process a submitted review table and advance to the next B&C phase.

        The table rows are serialised to JSON and injected as a structured
        user message so the agent can parse approvals, renames, and groups.

        Yields: (history, phase, charts_dict, downloads_list)
        """
        approved_groups = _extract_approved_groups(table_data)

        review_payload = json.dumps({
            "event":           "review_submitted",
            "rows":            table_data,
            "approved_groups": json.loads(approved_groups),
            "approved_count":  len(json.loads(approved_groups)),
        }, ensure_ascii=False, indent=2)

        review_message = (
            f"The user has submitted the review table. "
            f"Approved groups: {approved_groups}. "
            f"Full payload:\n{review_payload}\n\n"
            f"Please continue to the next B&C phase now."
        )

        def _make_review_stream() -> Any:
            return self._graph.stream(
                {"messages": [{"role": "user", "content": review_message}]},
                config=self._config(),
                stream_mode="updates",
            )

        accumulated = ""
        attempt     = 0

        while True:
            current_stream = _make_review_stream()
            try:
                for chunk in current_stream:
                    node_output = chunk.get("agent") or chunk.get("tools") or {}
                    messages    = node_output.get("messages", [])

                    text_delta = "".join(list(map(
                        lambda m: getattr(m, "content", "") if hasattr(m, "content") else "",
                        messages,
                    )))

                    accumulated += text_delta
                    self._update_charts(accumulated)
                    self._update_downloads(accumulated)
                    self.phase = _detect_phase(accumulated)

                    yield (
                        history + [["[Review submitted]", accumulated]],
                        self.phase,
                        dict(self._charts),
                        list(self._downloads),
                    )

                # ── Success ───────────────────────────────────────────────────
                yield (
                    history + [["[Review submitted]", accumulated]],
                    self.phase,
                    dict(self._charts),
                    list(self._downloads),
                )
                break   # exit retry loop

            except Exception as exc:  # noqa: BLE001
                if _is_rate_limit(exc) and attempt < _RL_MAX_RETRIES:
                    wait = _RL_BACKOFF_SECS[attempt]
                    notice = (
                        f"\n\n⏳ **Rate limit hit** β€” waiting **{wait}s** "
                        f"then retrying (attempt {attempt + 1}/{_RL_MAX_RETRIES})…"
                    )
                    logger.warning("Rate limit 429 in handle_review attempt %d; sleeping %ds", attempt, wait)
                    yield (
                        history + [["[Review submitted]", accumulated + notice]],
                        "idle",
                        dict(self._charts),
                        list(self._downloads),
                    )
                    time.sleep(wait)
                    attempt += 1
                    # Loop re-creates the stream via _make_review_stream()
                else:
                    if _is_rate_limit(exc):
                        err = (
                            f"❌ **Rate limit persists after {_RL_MAX_RETRIES} retries.**\n"
                            "Please wait a few minutes before trying again."
                        )
                    else:
                        err = f"❌ **Error processing review:** `{type(exc).__name__}: {exc}`"

                    logger.exception("Unhandled error in handle_review (attempt %d)", attempt)
                    self.phase = "idle"
                    yield (
                        history + [["[Review submitted]", accumulated + f"\n\n{err}"]],
                        "idle",
                        dict(self._charts),
                        list(self._downloads),
                    )
                    break