Skip to main content

nexus_core/app/
swarm.rs

1//! `/swarm`: a per-session roster of persona+model rows that, when swarm
2//! mode is on, turns a normal chat message into a live conversation between
3//! personas. Conversation proceeds in rounds: every persona gets a chance to
4//! respond before the moderator can decide whether the panel has converged.
5//! A synthesis reply is then written as the discussion's canonical assistant
6//! message.
7
8use std::collections::HashSet;
9use std::sync::Arc;
10
11use anyhow::Result;
12use tokio::sync::mpsc;
13use tokio::time::{Duration, timeout};
14
15use super::App;
16use crate::app::backends::Backends;
17use crate::db::Persona;
18use crate::provider::openrouter::OpenRouter;
19use crate::provider::{ChatMessage, ChatParams, StreamEvent};
20use crate::tools::ToolExecutor;
21
22/// Maximum complete panel rounds before synthesis. Every persona present at
23/// the start of a round gets one response opportunity in that round.
24const MAX_ROUNDS: usize = 4;
25/// Tool round-trips available to each persona response.
26const SWARM_PERSONA_MAX_TOOL_ITERS: usize = 8;
27
28/// Hard ceiling on how many personas a conversation can grow to via the
29/// moderator adding new voices mid-run.
30const MAX_PERSONAS: usize = 6;
31
32/// A `/swarm` turn update tagged with the session it belongs to (the turn
33/// runs in the background, so the viewer may have navigated away by the
34/// time an update lands).
35pub type SwarmMsg = (String, SwarmUpdate);
36
37#[derive(Clone)]
38pub enum SwarmUpdate {
39    /// The roster was empty, so one was suggested — persist it before the
40    /// conversation's first turn.
41    RosterSuggested(Vec<Persona>),
42    /// Status-line progress text.
43    Progress(String),
44    /// One persona's reply for the turn just run.
45    Reply {
46        persona: String,
47        model: String,
48        content: String,
49    },
50    /// The moderator decided the conversation needed a new voice — persist
51    /// it to the roster so it shows up in the popup too.
52    PersonaJoined(Persona),
53    /// The turn's final synthesis reply — the canonical assistant message.
54    Synthesis(String),
55    Error(String),
56}
57
58impl App {
59    /// Flip swarm mode for the active session.
60    pub fn toggle_swarm_mode(&mut self) -> Result<()> {
61        let Some(session) = &mut self.session else {
62            return Ok(());
63        };
64        let on = !session.swarm_mode;
65        session.swarm_mode = on;
66        self.db.set_session_swarm_mode(&session.id, on)?;
67        self.push_status(format!("swarm mode: {}", if on { "ON" } else { "OFF" }));
68        Ok(())
69    }
70
71    /// Start a swarm turn for the just-sent message (already pushed to
72    /// `self.messages`/the db by `send_message`). No-op if one's running.
73    pub fn start_swarm_turn(&mut self) {
74        let Some(session) = self.session.clone() else {
75            return;
76        };
77        if self.swarm_rx.is_some() {
78            self.push_status("a swarm turn is already running".to_string());
79            return;
80        }
81        let default_model = self.current_model.clone().unwrap_or_default();
82        let Some((default_provider, raw_default_model)) =
83            self.resolve_model_backend(&default_model)
84        else {
85            self.pending_events
86                .push_back(super::AppEvent::OpenLoginPopup);
87            return;
88        };
89        // The swarm meta agent has no separate model setting — it resolves on
90        // the current session's provider with that provider's research-class
91        // default model.
92        let (meta_provider, raw_meta_model) = self
93            .resolve_feature_model_backend("", OpenRouter::default_research_model)
94            .unwrap_or_else(|| (default_provider.clone(), raw_default_model.clone()));
95        let personas = self.db.list_swarm_personas(&session.id).unwrap_or_default();
96        let user_message = self
97            .messages
98            .last()
99            .map(|m| m.content.clone())
100            .unwrap_or_default();
101        let base_history = self.build_history();
102
103        let (tx, rx) = mpsc::unbounded_channel();
104        self.swarm_rx = Some(rx);
105        self.push_status("swarm: starting…".to_string());
106
107        let swarm_session_id = session.id.clone();
108        let task = tokio::spawn(run_swarm_turn(SwarmTurnOptions {
109            backends: self.backends.clone(),
110            meta_provider,
111            raw_meta_model,
112            personas,
113            default_model,
114            default_provider,
115            raw_default_model,
116            base_history,
117            toolbox: self.toolbox.clone(),
118            user_message,
119            session_id: session.id,
120            tx,
121        }));
122        self.swarm_abort = Some(task.abort_handle());
123        self.swarm_session = Some(swarm_session_id);
124    }
125
126    /// Domain half of the roster save: drop blank rows and persist. The
127    /// view clamps its cursor after calling this.
128    pub fn save_swarm_roster(&mut self) -> Result<()> {
129        self.swarm_cache.retain(|p| !p.name.trim().is_empty());
130        if let Some(session) = &self.session {
131            self.db
132                .save_swarm_personas(&session.id, &self.swarm_cache)?;
133        }
134        Ok(())
135    }
136
137    /// Stop the running swarm immediately. Persona model/tool streams are
138    /// children of the aborted orchestration task and are dropped with it.
139    /// The view closes its popup after calling this.
140    pub fn stop_swarm(&mut self) {
141        if let Some(abort) = self.swarm_abort.take() {
142            abort.abort();
143        }
144        if self.swarm_rx.take().is_some() {
145            if let Some(id) = self.swarm_session.take() {
146                let _ = self
147                    .db
148                    .upsert_research_stage_message(&id, "swarm", "stopped by user");
149            }
150            self.push_status("swarm stopped".to_string());
151        } else {
152            self.push_status("no swarm is running".to_string());
153        }
154    }
155
156    /// A swarm turn update: persist it, and mirror it into the live
157    /// transcript if the session it belongs to is the one being viewed.
158    /// `None` = the job's channel closed (fires once, right after the last update).
159    // Long by design (event dispatch).
160    #[allow(clippy::too_many_lines)]
161    pub fn on_swarm_update(&mut self, r: Option<SwarmMsg>) {
162        let Some((session_id, update)) = r else {
163            self.swarm_rx = None;
164            self.swarm_abort = None;
165            self.swarm_session = None;
166            return;
167        };
168        let viewing = self.session.as_ref().is_some_and(|s| s.id == session_id);
169        match update {
170            SwarmUpdate::RosterSuggested(personas) => {
171                let _ = self.db.save_swarm_personas(&session_id, &personas);
172                if viewing {
173                    self.swarm_cache = personas;
174                    self.push_status(
175                        "swarm: roster suggested — starting the conversation…".to_string(),
176                    );
177                }
178            }
179            SwarmUpdate::Progress(s) => {
180                let _ = self
181                    .db
182                    .upsert_research_stage_message(&session_id, "swarm", &s);
183                if viewing {
184                    let text = crate::db::stage_content("swarm", &s);
185                    if let Some(row) = self.messages.iter_mut().rev().find(|m| {
186                        m.role == "research_stage"
187                            && (m.content == "swarm" || m.content.starts_with("swarm:"))
188                    }) {
189                        row.content = text;
190                        self.push_history_invalidated();
191                    } else {
192                        self.messages.push(crate::db::Message {
193                            role: "research_stage".to_string(),
194                            content: text,
195                            model: None,
196                            reasoning: None,
197                            tokens: None,
198                            secs: None,
199                            cost: None,
200                            phrase: None,
201                            persona: None,
202                            created_at: None,
203                        });
204                    }
205                    self.push_status(format!("swarm: {s}"));
206                }
207            }
208            SwarmUpdate::Reply {
209                persona,
210                model,
211                content,
212            } => {
213                if self
214                    .db
215                    .add_persona_message(&session_id, &content, &persona, &model)
216                    .is_ok()
217                    && viewing
218                {
219                    self.messages.push(crate::db::Message {
220                        role: "assistant".to_string(),
221                        content,
222                        model: Some(model),
223                        reasoning: None,
224                        tokens: None,
225                        secs: None,
226                        cost: None,
227                        phrase: None,
228                        persona: Some(persona),
229                        created_at: None,
230                    });
231                }
232            }
233            SwarmUpdate::PersonaJoined(persona) => {
234                let mut roster = self.db.list_swarm_personas(&session_id).unwrap_or_default();
235                if !roster.iter().any(|p| p.name == persona.name) {
236                    roster.push(persona.clone());
237                    let _ = self.db.save_swarm_personas(&session_id, &roster);
238                }
239                if viewing {
240                    if !self.swarm_cache.iter().any(|p| p.name == persona.name) {
241                        self.swarm_cache.push(persona.clone());
242                    }
243                    self.push_status(format!("swarm: {} joined the conversation", persona.name));
244                }
245            }
246            SwarmUpdate::Synthesis(content) => {
247                let _ = self.db.add_assistant_message(
248                    &session_id,
249                    &content,
250                    None,
251                    None,
252                    None,
253                    None,
254                    None,
255                    None,
256                );
257                if viewing {
258                    self.messages.push(crate::db::Message {
259                        role: "assistant".to_string(),
260                        content,
261                        model: None,
262                        reasoning: None,
263                        tokens: None,
264                        secs: None,
265                        cost: None,
266                        phrase: Some("Discussed".to_string()),
267                        persona: None,
268                        created_at: None,
269                    });
270                    self.push_status("swarm turn complete".to_string());
271                    self.maybe_generate_title();
272                    self.maybe_extract_memory();
273                    self.maybe_compact();
274                }
275            }
276            SwarmUpdate::Error(e) => {
277                // Keep failures separate from the upserted progress row so a
278                // later round/tool update cannot overwrite and erase them.
279                let _ = self
280                    .db
281                    .add_error_message(&session_id, &format!("swarm: {e}"));
282                if viewing {
283                    self.messages.push(crate::db::Message {
284                        role: "error".to_string(),
285                        content: format!("swarm: {e}"),
286                        model: None,
287                        reasoning: None,
288                        tokens: None,
289                        secs: None,
290                        cost: None,
291                        phrase: None,
292                        persona: None,
293                        created_at: None,
294                    });
295                    self.push_status(format!("swarm error: {e}"));
296                }
297            }
298        }
299    }
300}
301
302pub fn parse_persona_editor(text: &str) -> Result<Persona, String> {
303    let mut lines = text.lines();
304    let name = lines
305        .next()
306        .and_then(|line| line.strip_prefix("name:"))
307        .map(str::trim)
308        .filter(|value| !value.is_empty())
309        .ok_or_else(|| "first line must be `name: <non-empty name>`".to_string())?;
310    let model = lines
311        .next()
312        .and_then(|line| line.strip_prefix("model:"))
313        .map(str::trim)
314        .filter(|value| !value.is_empty())
315        .ok_or_else(|| "second line must be `model: <model id>`".to_string())?;
316    let remaining: Vec<&str> = lines.collect();
317    let separator = remaining
318        .iter()
319        .position(|line| line.trim() == "---")
320        .ok_or_else(|| "missing `---` before the blurb".to_string())?;
321    let blurb = remaining[separator + 1..].join("\n").trim().to_string();
322    Ok(Persona {
323        name: name.to_string(),
324        model: model.to_string(),
325        blurb,
326    })
327}
328
329#[allow(clippy::too_many_arguments)]
330/// Everything needed to start one swarm conversation turn: the persona
331/// roster, the model/provider config, the conversation history, and the
332/// orchestration plumbing (toolbox, session identity, update channel).
333pub struct SwarmTurnOptions {
334    pub backends: Backends,
335    pub meta_provider: OpenRouter,
336    pub raw_meta_model: String,
337    pub personas: Vec<Persona>,
338    pub default_model: String,
339    pub default_provider: OpenRouter,
340    pub raw_default_model: String,
341    pub base_history: Vec<ChatMessage>,
342    pub toolbox: Arc<dyn ToolExecutor>,
343    pub user_message: String,
344    pub session_id: String,
345    pub tx: mpsc::UnboundedSender<SwarmMsg>,
346}
347
348// Long by design (roundtable orchestration).
349#[allow(clippy::too_many_lines)]
350async fn run_swarm_turn(opts: SwarmTurnOptions) {
351    let SwarmTurnOptions {
352        backends,
353        meta_provider,
354        raw_meta_model,
355        mut personas,
356        default_model,
357        default_provider,
358        raw_default_model,
359        base_history,
360        toolbox,
361        user_message,
362        session_id,
363        tx,
364    } = opts;
365    let send = |u: SwarmUpdate| {
366        let _ = tx.send((session_id.clone(), u));
367    };
368    if personas.is_empty() {
369        send(SwarmUpdate::Progress("suggesting personas…".to_string()));
370        match suggest_personas(
371            &meta_provider,
372            &raw_meta_model,
373            &user_message,
374            &default_model,
375        )
376        .await
377        {
378            Ok(p) => {
379                personas = p;
380                send(SwarmUpdate::RosterSuggested(personas.clone()));
381            }
382            Err(e) => {
383                send(SwarmUpdate::Error(format!(
384                    "couldn't suggest personas: {e}"
385                )));
386                return;
387            }
388        }
389    }
390    if personas.is_empty() {
391        send(SwarmUpdate::Error("no personas to run".to_string()));
392        return;
393    }
394
395    // (persona name, reply content), in speaking order — a running transcript
396    // every persona converses through. A successful reply marks that persona
397    // as having answered; the moderator is gated on every current persona
398    // having answered at least once and on the current round being complete.
399    let mut discussion: Vec<(String, String)> = Vec::new();
400    // Track roster positions, not names: users may configure duplicate names,
401    // but every row still deserves its own response opportunity.
402    let mut answered: HashSet<usize> = HashSet::new();
403    'convo: for round in 1..=MAX_ROUNDS {
404        // Snapshot the roster for this round. Personas can only be added by the
405        // moderator after the round, so everyone in this snapshot gets exactly
406        // one opportunity before moderation.
407        let round_personas = personas.clone();
408        let round_size = round_personas.len();
409        for (idx, p) in round_personas.into_iter().enumerate() {
410            send(SwarmUpdate::Progress(format!(
411                "round {round}/{MAX_ROUNDS} · persona {}/{} — {} is responding",
412                idx + 1,
413                round_size,
414                p.name
415            )));
416
417            let mut messages = base_history.clone();
418            messages.push(ChatMessage::text(
419                "system",
420                format!(
421                    "You are {} in a live group conversation with the other personas below. \
422                     Your persona: {}. This discussion proceeds in rounds and every persona \
423                     gets one response per round. Speak naturally and briefly (a few sentences), \
424                     stay in character, and actually engage with what the prior speakers said — \
425                     agree, push back, build on it, or ask a question. Use the available tools \
426                     when they would improve your answer. Don't just restate your opening \
427                     position every round.",
428                    p.name, p.blurb
429                ),
430            ));
431            match render_discussion(&discussion) {
432                Some(transcript) => messages.push(ChatMessage::text(
433                    "user",
434                    format!(
435                        "Conversation so far:\n\n{transcript}\n\nRound {round}: it's your chance \
436                         to respond, {}. Engage with what was said.",
437                        p.name
438                    ),
439                )),
440                None => messages.push(ChatMessage::text(
441                    "user",
442                    format!(
443                        "Round {round}: you're opening the discussion, {}. Give your first take.",
444                        p.name
445                    ),
446                )),
447            }
448
449            let (persona_provider, raw_persona_model) =
450                backends.resolve(&p.model).unwrap_or_else(|| {
451                    send(SwarmUpdate::Error(format!(
452                        "{} model unavailable ({}); using session model",
453                        p.name, p.model
454                    )));
455                    (default_provider.clone(), raw_default_model.clone())
456                });
457            let persona_toolbox = toolbox.clone();
458            let tools = persona_toolbox.defs();
459            let (mut rx, abort) = persona_provider.stream_chat(
460                raw_persona_model,
461                messages,
462                ChatParams::default(),
463                tools,
464                persona_toolbox,
465                SWARM_PERSONA_MAX_TOOL_ITERS,
466            );
467            let abort = super::AbortOnDrop(abort);
468            let response = timeout(Duration::from_mins(2), async {
469                let mut content = String::new();
470                while let Some(event) = rx.recv().await {
471                    match event {
472                        StreamEvent::Token(token) => content.push_str(&token),
473                        StreamEvent::Status(status) => send(SwarmUpdate::Progress(format!(
474                            "round {round}/{MAX_ROUNDS} — {}: {status}",
475                            p.name
476                        ))),
477                        StreamEvent::ToolCall {
478                            name,
479                            arguments,
480                            result,
481                        } => {
482                            let summary = crate::app::tool_call_summary(&name, &arguments, &result);
483                            send(SwarmUpdate::Progress(format!(
484                                "round {round}/{MAX_ROUNDS} — {}: {summary}",
485                                p.name
486                            )));
487                        }
488                        StreamEvent::Error(error) => return Err(anyhow::anyhow!(error)),
489                        StreamEvent::Done => break,
490                        _ => {}
491                    }
492                }
493                if content.trim().is_empty() {
494                    Err(anyhow::anyhow!("returned an empty response"))
495                } else {
496                    Ok(content)
497                }
498            })
499            .await;
500            let response = if let Ok(result) = response {
501                result
502            } else {
503                abort.0.abort();
504                Err(anyhow::anyhow!("timed out after 120s"))
505            };
506            match response {
507                Ok(content) => {
508                    send(SwarmUpdate::Reply {
509                        persona: p.name.clone(),
510                        model: p.model.clone(),
511                        content: content.clone(),
512                    });
513                    answered.insert(idx);
514                    discussion.push((p.name, content));
515                }
516                Err(e) => send(SwarmUpdate::Error(format!("{} failed: {e}", p.name))),
517            }
518        }
519
520        // The hard cap synthesizes after the final complete round. Never ask
521        // the moderator to add someone who would have no round left to speak.
522        if round == MAX_ROUNDS {
523            break 'convo;
524        }
525        if !all_personas_answered(personas.len(), &answered) {
526            let missing = personas
527                .iter()
528                .enumerate()
529                .filter(|(idx, _)| !answered.contains(idx))
530                .map(|(_, p)| p.name.as_str())
531                .collect::<Vec<_>>()
532                .join(", ");
533            send(SwarmUpdate::Progress(format!(
534                "round {round} complete — waiting for replies from: {missing}"
535            )));
536            continue;
537        }
538
539        send(SwarmUpdate::Progress(format!(
540            "round {round} complete — moderator checking for a conclusion…"
541        )));
542        match moderator_check(&meta_provider, &raw_meta_model, &user_message, &discussion).await {
543            Ok(ModeratorVerdict::Converged) | Err(_) => break 'convo,
544            Ok(ModeratorVerdict::Continue) => {}
545            Ok(ModeratorVerdict::AddPersona { name, blurb }) => {
546                if personas.len() >= MAX_PERSONAS
547                    || personas.iter().any(|existing| existing.name == name)
548                {
549                    continue;
550                }
551                let new_persona = Persona {
552                    name,
553                    model: default_model.clone(),
554                    blurb,
555                };
556                send(SwarmUpdate::PersonaJoined(new_persona.clone()));
557                personas.push(new_persona);
558            }
559        }
560    }
561
562    send(SwarmUpdate::Progress("writing synthesis…".to_string()));
563    match synthesize(&meta_provider, &raw_meta_model, &user_message, &discussion).await {
564        Ok(content) => send(SwarmUpdate::Synthesis(content)),
565        Err(e) => send(SwarmUpdate::Error(format!("synthesis failed: {e}"))),
566    }
567}
568
569fn all_personas_answered(persona_count: usize, answered: &HashSet<usize>) -> bool {
570    persona_count > 0 && (0..persona_count).all(|idx| answered.contains(&idx))
571}
572
573fn render_discussion(discussion: &[(String, String)]) -> Option<String> {
574    if discussion.is_empty() {
575        return None;
576    }
577    Some(
578        discussion
579            .iter()
580            .map(|(name, content)| format!("**{name}**: {content}"))
581            .collect::<Vec<_>>()
582            .join("\n\n"),
583    )
584}
585
586#[derive(serde::Deserialize)]
587struct Suggested {
588    name: String,
589    blurb: String,
590}
591
592async fn suggest_personas(
593    provider: &OpenRouter,
594    meta_model: &str,
595    topic: &str,
596    default_model: &str,
597) -> Result<Vec<Persona>, String> {
598    let prompt = format!(
599        "Suggest exactly 3 distinct personas to discuss the following message from different \
600         points of view. Reply with ONLY a JSON array, no prose, no markdown fences, shaped \
601         like: [{{\"name\": \"short persona name\", \"blurb\": \"one-sentence personality and \
602         perspective\"}}, ...]\n\nMessage: {topic}"
603    );
604    let raw = timeout(
605        Duration::from_mins(1),
606        provider.complete(meta_model, vec![ChatMessage::text("user", prompt)]),
607    )
608    .await
609    .map_err(|_| "timed out after 60s".to_string())?
610    .map_err(|e| e.to_string())?;
611    parse_suggested_personas(&raw, default_model)
612}
613
614/// Extract a `[{"name","blurb"}, ...]` JSON array from a (possibly
615/// prose/fence-wrapped) model reply. Split out from `suggest_personas` so
616/// the parsing logic is testable without a network call.
617fn parse_suggested_personas(raw: &str, default_model: &str) -> Result<Vec<Persona>, String> {
618    let start = raw.find('[').ok_or("no JSON array in response")?;
619    let end = raw.rfind(']').ok_or("no JSON array in response")?;
620    let json = raw.get(start..=end).ok_or("no JSON array in response")?;
621    let parsed: Vec<Suggested> = serde_json::from_str(json).map_err(|e| e.to_string())?;
622    Ok(parsed
623        .into_iter()
624        .filter(|s| !s.name.trim().is_empty())
625        .map(|s| Persona {
626            name: s.name,
627            model: default_model.to_string(),
628            blurb: s.blurb,
629        })
630        .collect())
631}
632
633enum ModeratorVerdict {
634    Converged,
635    Continue,
636    /// The conversation would benefit from a new voice — its name and a
637    /// one-sentence personality/perspective blurb.
638    AddPersona {
639        name: String,
640        blurb: String,
641    },
642}
643
644async fn moderator_check(
645    provider: &OpenRouter,
646    meta_model: &str,
647    topic: &str,
648    discussion: &[(String, String)],
649) -> Result<ModeratorVerdict, String> {
650    let transcript = render_discussion(discussion).unwrap_or_default();
651    let prompt = format!(
652        "A panel of personas is having a live conversation about this message:\n\n{topic}\n\n\
653         Conversation so far:\n\n{transcript}\n\nDecide one of three things and reply with ONLY \
654         one line, no other text:\n\
655         - If they've reached a conclusion (agreement, a clear resolution, or a settled \
656         trade-off), reply exactly: CONVERGED\n\
657         - If the conversation is missing an important point of view that would meaningfully \
658         change it, reply exactly: ADD: <short persona name> | <one-sentence personality and \
659         perspective>\n\
660         - Otherwise, reply exactly: CONTINUE"
661    );
662    let raw = timeout(
663        Duration::from_secs(30),
664        provider.complete(meta_model, vec![ChatMessage::text("user", prompt)]),
665    )
666    .await
667    .map_err(|_| "timed out after 30s".to_string())?
668    .map_err(|e| e.to_string())?;
669    Ok(parse_moderator_verdict(&raw))
670}
671
672fn parse_moderator_verdict(raw: &str) -> ModeratorVerdict {
673    // ASCII-only uppercasing so byte offsets stay valid for slicing `raw`
674    // (full Unicode `to_uppercase` can change a string's byte length).
675    let upper = raw.to_ascii_uppercase();
676    if let Some(rest) = upper.find("ADD:").map(|i| &raw[i + 4..])
677        && let Some((name, blurb)) = rest.split_once('|')
678    {
679        let name = name.trim().trim_matches('*').to_string();
680        let blurb = blurb.lines().next().unwrap_or("").trim().to_string();
681        if !name.is_empty() && !blurb.is_empty() {
682            return ModeratorVerdict::AddPersona { name, blurb };
683        }
684    }
685    if upper.contains("CONVERGED") {
686        return ModeratorVerdict::Converged;
687    }
688    ModeratorVerdict::Continue
689}
690
691async fn synthesize(
692    provider: &OpenRouter,
693    meta_model: &str,
694    topic: &str,
695    discussion: &[(String, String)],
696) -> Result<String> {
697    let transcript = render_discussion(discussion).unwrap_or_default();
698    let prompt = format!(
699        "A panel of personas had this conversation about a message:\n\n{topic}\n\n\
700         Conversation:\n\n{transcript}\n\nWrite one final reply to the original message that \
701         states the conclusion this conversation actually reached — the agreement, resolution, \
702         or trade-off they settled on — as a clear, useful answer. If they genuinely disagreed \
703         to the end, say so and give the best-supported call. Markdown allowed. Don't mention \
704         that this came from a panel — just answer well, informed by the conversation."
705    );
706    timeout(
707        Duration::from_secs(90),
708        provider.complete(meta_model, vec![ChatMessage::text("user", prompt)]),
709    )
710    .await
711    .map_err(|_| anyhow::anyhow!("timed out after 90s"))?
712}
713
714#[cfg(test)]
715mod tests {
716    use super::*;
717
718    #[test]
719    fn persona_editor_round_trips_name_model_and_multiline_blurb() {
720        let persona = parse_persona_editor(
721            "name: Skeptic\nmodel: codex:gpt-5.4-mini\n---\npokes holes\nand checks evidence\n",
722        )
723        .unwrap();
724        assert_eq!(persona.name, "Skeptic");
725        assert_eq!(persona.model, "codex:gpt-5.4-mini");
726        assert_eq!(persona.blurb, "pokes holes\nand checks evidence");
727        assert!(parse_persona_editor("name: \nmodel: m\n---\n").is_err());
728    }
729
730    #[test]
731    fn parse_suggested_personas_tolerates_prose_and_fences() {
732        let raw = "sure, here you go:\n```json\n[{\"name\":\"Skeptic\",\"blurb\":\"pokes holes\"},\
733                   {\"name\":\"Advocate\",\"blurb\":\"user-first\"}]\n```";
734        let personas = parse_suggested_personas(raw, "a/one").unwrap();
735        assert_eq!(personas.len(), 2);
736        assert_eq!(personas[0].name, "Skeptic");
737        assert_eq!(personas[0].model, "a/one");
738        assert_eq!(personas[1].blurb, "user-first");
739    }
740
741    #[test]
742    fn parse_suggested_personas_drops_blank_names() {
743        let raw = r#"[{"name":"","blurb":"nameless"},{"name":"Real","blurb":"ok"}]"#;
744        let personas = parse_suggested_personas(raw, "a/one").unwrap();
745        assert_eq!(personas.len(), 1);
746        assert_eq!(personas[0].name, "Real");
747    }
748
749    #[test]
750    fn parse_suggested_personas_errors_without_an_array() {
751        assert!(parse_suggested_personas("no json here", "a/one").is_err());
752    }
753
754    #[test]
755    fn moderator_gate_requires_every_current_persona_row_to_have_answered() {
756        let mut answered = HashSet::from([0usize]);
757        assert!(!all_personas_answered(2, &answered));
758
759        answered.insert(1);
760        assert!(all_personas_answered(2, &answered));
761        assert!(
762            !all_personas_answered(3, &answered),
763            "a newly joined persona must get a round before moderation"
764        );
765    }
766
767    #[test]
768    fn render_discussion_none_when_empty_some_when_not() {
769        assert!(render_discussion(&[]).is_none());
770        let d = vec![("Skeptic".to_string(), "wait, really?".to_string())];
771        let rendered = render_discussion(&d).unwrap();
772        assert!(rendered.contains("**Skeptic**"));
773        assert!(rendered.contains("wait, really?"));
774    }
775
776    #[test]
777    fn parse_moderator_verdict_recognizes_converged() {
778        assert!(matches!(
779            parse_moderator_verdict("CONVERGED"),
780            ModeratorVerdict::Converged
781        ));
782        assert!(matches!(
783            parse_moderator_verdict("  converged.\n"),
784            ModeratorVerdict::Converged
785        ));
786    }
787
788    #[test]
789    fn parse_moderator_verdict_defaults_to_continue() {
790        assert!(matches!(
791            parse_moderator_verdict("CONTINUE"),
792            ModeratorVerdict::Continue
793        ));
794        assert!(matches!(
795            parse_moderator_verdict("something unexpected"),
796            ModeratorVerdict::Continue
797        ));
798    }
799
800    #[test]
801    fn parse_moderator_verdict_extracts_name_and_blurb_for_add() {
802        match parse_moderator_verdict("ADD: Realist | grounds the discussion in constraints") {
803            ModeratorVerdict::AddPersona { name, blurb } => {
804                assert_eq!(name, "Realist");
805                assert_eq!(blurb, "grounds the discussion in constraints");
806            }
807            _ => panic!("expected AddPersona"),
808        }
809    }
810
811    #[test]
812    fn parse_moderator_verdict_falls_back_to_continue_on_malformed_add() {
813        assert!(matches!(
814            parse_moderator_verdict("ADD: no separator here"),
815            ModeratorVerdict::Continue
816        ));
817    }
818}