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