Skip to main content

magi/
consult.rs

1//! Handing an open question to the chat conversation its task came from.
2//!
3//! A task filed by the standing chat (`Source::Agent` with
4//! [`crate::queue::CHAT_NODE`], whose `run` is the conversation's id) can raise
5//! a question the owner would rather talk through where the task started. The
6//! owner passes it on; the chat agent either answers it with `magi answer` or
7//! puts the decision points to the owner in the conversation.
8//!
9//! Three rules, each with a reason:
10//!
11//! - **One place decides.** [`origin_talk`] is the only judge of "does this
12//!   question have a chat to ask"; the web view and both entry points use it.
13//! - **Not a choice.** The hand-over never goes into `Question::choices` and
14//!   never touches `Question::answer`: the question stays `Open`, its
15//!   validation and the choice actions stay as they were.
16//! - **No new seat, no new waiter.** The text is queued as a draft of the
17//!   existing talk, so the chat's own turn machinery (its session, its turn
18//!   gate) runs it exactly as if the owner had typed it.
19
20use anyhow::{Context, Result, bail};
21use jiff::Timestamp;
22
23use crate::ask::{ChatConsult, Question, Questions};
24use crate::config::Config;
25use crate::queue::{CHAT_NODE, Source, Task};
26use crate::talk::{self, Talk, Talks};
27
28/// The conversation `q` may be handed to, if any.
29///
30/// Requires an open question whose task was filed from a chat that is still
31/// open. Release notices are left out. Merge approvals may be discussed, but
32/// answering from chat requires the owner's latest words.
33pub fn origin_talk(tasks: &[Task], talks: &[Talk], q: &Question) -> Option<Talk> {
34    if !q.status.open() || q.node == crate::bump::NOTICE_NODE {
35        return None;
36    }
37    let task = crate::daemon::task_of_question(tasks, q)?;
38    let Source::Agent { run, .. } = &chat_origin(tasks, task)?.source else {
39        return None;
40    };
41    talks
42        .iter()
43        .find(|t| &t.id == run && t.status.open())
44        .cloned()
45}
46
47/// Walk the provenance of `start` to the task a chat filed.
48///
49/// A follow-up goes to the task its merged run served: `FollowUp::origin_task`
50/// when recorded (a missing one ends the walk - guessing another task could
51/// hand the question to an unrelated chat), else the task whose `runs` hold
52/// `FollowUp::run`. At most `MAX_FOLLOWUP_GENERATION + 1` tasks are looked at,
53/// each once, so a cycle ends in `None`.
54fn chat_origin<'a>(tasks: &'a [Task], start: &'a Task) -> Option<&'a Task> {
55    let mut seen = std::collections::HashSet::new();
56    let mut cur = start;
57    for _ in 0..=crate::followup::MAX_FOLLOWUP_GENERATION {
58        if !seen.insert(cur.id.as_str()) {
59            return None;
60        }
61        if matches!(&cur.source, Source::Agent { node, .. } if node == CHAT_NODE) {
62            return Some(cur);
63        }
64        let f = cur.followup.as_ref()?;
65        cur = match &f.origin_task {
66            Some(id) => tasks.iter().find(|t| &t.id == id)?,
67            None => tasks.iter().find(|t| t.runs.contains(&f.run))?,
68        };
69    }
70    None
71}
72
73/// Is any question handed to talk `talk_id` still open?
74///
75/// Decided from the question store, not from the newest turn's wording: the
76/// owner's decision usually arrives in a later turn that carries no hand-over
77/// heading, and `magi answer` still has to be able to write then. Answered or
78/// abandoned questions stop counting, so the turn goes back to read-only.
79pub fn pending_consults(questions: &Questions, talk_id: &str) -> bool {
80    questions
81        .list()
82        .iter()
83        .any(|q| q.status.open() && q.consult.as_ref().is_some_and(|c| c.talk == talk_id))
84}
85
86/// Record the hand-over and queue the question as the chat's next message.
87///
88/// Returns `false` (and changes nothing) when the question was already handed
89/// over. The question is not otherwise touched: still `Open`, no thread turn.
90/// When queueing fails the record is withdrawn, so the owner can try again.
91pub fn begin(questions: &Questions, talks: &Talks, q: &Question, talk: &Talk) -> Result<bool> {
92    let (q, fresh) = questions.update(&q.id, |r| {
93        if !r.status.open() {
94            bail!("question {} is already {}", r.short(), r.status.as_str());
95        }
96        if r.consult.is_some() {
97            return Ok(false);
98        }
99        r.consult = Some(ChatConsult {
100            talk: talk.id.clone(),
101            at: Timestamp::now(),
102        });
103        Ok(true)
104    })?;
105    if !fresh {
106        return Ok(false);
107    }
108    let mut talk = talk.clone();
109    if let Err(e) = talk::queue(
110        &mut talk,
111        talks,
112        &crate::prompt::chat_consult(&q),
113        Vec::new(),
114    ) {
115        let _ = questions.update(&q.id, |r| {
116            r.consult = None;
117            Ok(())
118        });
119        return Err(e).context("queue the question into the chat");
120    }
121    Ok(true)
122}
123
124/// Validate a merge approval answered by a running chat. Terminal answers and
125/// other question types retain their existing behavior. Intent is judged by
126/// the chat; this gate requires evidence from the correct conversation.
127pub fn validate_answer(
128    q: &Question,
129    talks: &Talks,
130    run: &str,
131    node: &str,
132    reply: &str,
133    quote: Option<&str>,
134) -> Result<()> {
135    if q.node != crate::land::APPROVAL_NODE || node != CHAT_NODE {
136        return Ok(());
137    }
138    let consult = q
139        .consult
140        .as_ref()
141        .context("merge approval was not handed to a chat")?;
142    if consult.talk != run {
143        bail!("merge approval belongs to a different chat");
144    }
145    if !q.status.open() {
146        bail!("merge approval is no longer open");
147    }
148    let talk = talks.get(run)?;
149    if !talk.status.open() {
150        bail!("the consulted chat is closed");
151    }
152    // Only words the owner wrote after this question was handed over count.
153    // The generated hand-over text is stored as an operator turn (possibly
154    // coalesced with owner replies), so it is cut out by its own markers.
155    let at = talk
156        .turns
157        .iter()
158        .rposition(|t| {
159            t.who == talk::Who::Operator
160                && t.body.contains(crate::prompt::CHAT_CONSULT_HEADING)
161                && t.body.contains(&q.id)
162        })
163        .context("the question was not delivered to the chat yet")?;
164    // A reply still waiting in `pending` is newer than every stored turn.
165    let queued = owner_words(&talk.pending, None);
166    let stored = talk.turns[at..]
167        .iter()
168        .enumerate()
169        .rev()
170        .filter(|(_, t)| t.who == talk::Who::Operator)
171        .map(|(i, t)| owner_words(&t.body, (i == 0).then_some(q.id.as_str())))
172        .find(|w| !w.is_empty());
173    let latest = if queued.is_empty() {
174        stored
175    } else {
176        Some(queued)
177    }
178    .context("no owner message after the question was handed to the chat")?;
179    // Replies sent while a turn runs are joined with a blank line into one
180    // turn, so the last paragraph is the only text certain to be the newest.
181    let latest = latest
182        .rsplit("\n\n")
183        .map(str::trim)
184        .find(|p| !p.is_empty())
185        .unwrap_or_default()
186        .to_owned();
187    let quote = quote
188        .map(str::trim)
189        .filter(|s| !s.is_empty())
190        .context("chat merge approval requires --quote from the owner's latest message")?;
191    let valid = match reply {
192        crate::land::APPROVE => latest.contains(quote),
193        crate::land::HOLD => {
194            latest.trim().eq_ignore_ascii_case(crate::land::HOLD)
195                && quote.eq_ignore_ascii_case(crate::land::HOLD)
196        }
197        _ => false,
198    };
199    if !valid {
200        bail!(
201            "merge requires a verbatim quote of the latest owner message; hold requires the whole message to be hold"
202        );
203    }
204    Ok(())
205}
206
207/// The owner's own words in an operator turn, with magi's generated hand-over
208/// text cut out. `drain` stores a queued consult as an operator turn, and
209/// `talk::queue` joins an owner reply sent meanwhile onto the same draft, so
210/// the turn's origin cannot be told from the turn as a whole. Each generated
211/// block runs from its [`CHAT_CONSULT_HEADING`](crate::prompt::CHAT_CONSULT_HEADING)
212/// to the last [`CHAT_CONSULT_END`](crate::prompt::CHAT_CONSULT_END) before the
213/// next heading (a detail may quote the phrase itself); with no end the block
214/// runs to the next heading. Only the blocks are cut: owner replies between
215/// them are kept, in order. An owner reply that happens to contain the phrase
216/// is cut too, which only ever refuses. With `after_block_of` (the turn
217/// carrying the hand-over itself) everything up to the end of the last block
218/// naming that question id predates the hand-over and is dropped.
219fn owner_words(body: &str, after_block_of: Option<&str>) -> String {
220    use crate::prompt::{CHAT_CONSULT_END, CHAT_CONSULT_HEADING};
221    // A real heading is the generated one: `# ` at a line start, followed by
222    // the fixed opening sentence. A heading phrase quoted inside a question's
223    // detail is not a boundary; it stays inside the block it is part of.
224    let heads: Vec<usize> = body
225        .match_indices(CHAT_CONSULT_HEADING)
226        .map(|(i, _)| i)
227        .filter(|&i| {
228            let line_start = body[..i]
229                .strip_suffix("# ")
230                .is_some_and(|b| b.is_empty() || b.ends_with('\n'));
231            line_start
232                && body[i + CHAT_CONSULT_HEADING.len()..].starts_with("\n\nThe operator passed you")
233        })
234        .collect();
235    if heads.is_empty() {
236        return body.trim().to_owned();
237    }
238    // (start, end) of each generated block, the heading's own "# " included.
239    // A full consultation quoted in a detail nests: headings open and end
240    // phrases close, and only text outside every open block can be the owner's.
241    // Once a block has closed, a further end phrase before the next top-level
242    // heading extends it (a detail may quote the phrase alone). A block that
243    // never closes runs to the end of the text.
244    let mut events: Vec<(usize, bool)> = heads.iter().map(|&h| (h, true)).collect();
245    events.extend(
246        body.match_indices(CHAT_CONSULT_END)
247            .map(|(e, _)| (e, false)),
248    );
249    events.sort_unstable();
250    let mut blocks: Vec<(usize, usize)> = Vec::new();
251    let mut cur: Option<(usize, usize, usize)> = None; // start, end, depth
252    for (at, is_head) in events {
253        match (&mut cur, is_head) {
254            (Some((_, _, depth)), true) if *depth > 0 => *depth += 1,
255            (_, true) => {
256                blocks.extend(cur.take().map(|(s, e, _)| (s, e)));
257                let start = if body[..at].ends_with("# ") {
258                    at - 2
259                } else {
260                    at
261                };
262                cur = Some((start, body.len(), 1));
263            }
264            (Some((_, end, depth)), false) => {
265                *depth = depth.saturating_sub(1);
266                if *depth == 0 {
267                    *end = at + CHAT_CONSULT_END.len();
268                }
269            }
270            (None, false) => {}
271        }
272    }
273    blocks.extend(cur.map(|(s, e, _)| (s, e)));
274    let mut from = 0;
275    if let Some(id) = after_block_of {
276        if let Some(&(_, end)) = blocks.iter().rfind(|&&(s, e)| body[s..e].contains(id)) {
277            from = end;
278        } else if let Some(&(_, end)) = blocks.last() {
279            from = end;
280        }
281    }
282    let mut parts = Vec::new();
283    let mut at = from;
284    for &(s, e) in &blocks {
285        if s >= at {
286            parts.push(body[at..s].trim());
287        }
288        at = at.max(e);
289    }
290    parts.push(body[at..].trim());
291    parts
292        .into_iter()
293        .filter(|p| !p.is_empty())
294        .collect::<Vec<_>>()
295        .join("\n\n")
296}
297
298/// What [`start_turn`] did.
299#[derive(Debug, PartialEq, Eq)]
300pub enum Started {
301    /// The question was handed over and the chat's turn ran to its end.
302    Answered,
303    /// Another process holds the talk's turn lease; nothing was changed.
304    Busy,
305    /// The question had already been handed to the chat; nothing was changed.
306    Nothing,
307}
308
309/// [`begin`], then run the chat's turn here, under the talk's own lease.
310///
311/// The lease is taken *before* anything is written, so a talk somebody else is
312/// running leaves no consult record and no draft behind (`Started::Busy`) and
313/// the caller can simply try again later. A turn that fails after `begin`
314/// leaves the record and the draft, so the chat can resume it.
315pub async fn start_turn(
316    questions: &Questions,
317    talks: &Talks,
318    q: &Question,
319    talk: &Talk,
320    cfg: &Config,
321) -> Result<Started> {
322    if q.consult.is_some() {
323        return Ok(Started::Nothing);
324    }
325    let Some(mut lease) = talks.claim_turn(&talk.id)? else {
326        return Ok(Started::Busy);
327    };
328    if !begin(questions, talks, q, talk)? {
329        return Ok(Started::Nothing);
330    }
331    let mut talk = talks.get(&talk.id)?;
332    // Other starters that found the lease held queued drafts and left them to
333    // us, so drain until nothing is left (as the web's drain loop does).
334    // A failed turn is kept and reported at the end, not returned at once:
335    // drafts accepted meanwhile are still owed an answer.
336    let mut failed = None;
337    loop {
338        while let Some(text) = talk::drain(&mut talk, talks)? {
339            if let Err(e) = talk::respond(&lease, &mut talk, talks, cfg, &text).await {
340                failed.get_or_insert(e);
341            }
342            if !lease.beat()? {
343                bail!("the turn lease for chat {} was lost", talk.short());
344            }
345        }
346        // A draft queued after the last drain, before the lease is gone, was
347        // left to us by a starter that found the lease held. Release first,
348        // then look again: whoever queues later either sees no lease (and
349        // starts its own turn) or is seen by this re-check.
350        drop(lease);
351        talk = talks.get(&talk.id)?;
352        let owed = talk.status.open()
353            && (!talk.pending.is_empty() || !talk.pending_attachments.is_empty());
354        if !owed {
355            break;
356        }
357        // Somebody else took the lease meanwhile: they drain it.
358        let Some(again) = talks.claim_turn(&talk.id)? else {
359            break;
360        };
361        lease = again;
362    }
363    if let Some(e) = failed {
364        return Err(e);
365    }
366    Ok(Started::Answered)
367}
368
369#[cfg(test)]
370mod tests {
371    use std::collections::BTreeMap;
372    use std::path::PathBuf;
373
374    use super::*;
375    use crate::config::{AgentKind, AgentSpec, Config};
376    use crate::queue::Source;
377    use crate::talk::TalkStatus;
378
379    const RUN: &str = "20260902-000000-beef";
380
381    fn talks() -> (tempfile::TempDir, Talks, Talk) {
382        let tmp = tempfile::tempdir().expect("tempdir");
383        let store = Talks::at(tmp.path().join("talks"));
384        let cfg = Config {
385            agents: vec![AgentSpec {
386                id: "mock".to_owned(),
387                kind: AgentKind::Command,
388                model: None,
389                command: vec!["true".to_owned()],
390                extra_args: Vec::new(),
391                env: BTreeMap::new(),
392                prompt_delivery: None,
393            }],
394            ..Config::default()
395        };
396        let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
397        (tmp, store, talk)
398    }
399
400    fn task(source: Source) -> Task {
401        let mut t = Task::new(
402            "t".to_owned(),
403            "Do it".to_owned(),
404            PathBuf::from("/repo"),
405            source,
406        );
407        t.start(RUN.to_owned());
408        t
409    }
410
411    fn from_chat(talk: &Talk) -> Task {
412        task(Source::Agent {
413            run: talk.id.clone(),
414            node: CHAT_NODE.to_owned(),
415        })
416    }
417
418    fn question(node: &str) -> Question {
419        Question::new(
420            RUN.to_owned(),
421            node.to_owned(),
422            "impl-A".to_owned(),
423            "Which backend?".to_owned(),
424            "SQLite is simpler.".to_owned(),
425            vec!["SQLite".to_owned(), "Redis".to_owned()],
426        )
427    }
428
429    #[tokio::test]
430    async fn start_turn_is_busy_and_writes_nothing_while_the_lease_is_held() {
431        let (tmp, store, talk) = talks();
432        let questions = Questions::at(tmp.path().join("questions"));
433        let mut q = question("implement");
434        questions.put(&mut q).expect("put");
435        let held = store.claim_turn(&talk.id).expect("claim").expect("first");
436
437        let got = start_turn(&questions, &store, &q, &talk, &Config::default())
438            .await
439            .expect("start");
440        assert_eq!(got, Started::Busy);
441        assert!(questions.get(&q.id).unwrap().consult.is_none());
442        assert!(store.get(&talk.id).unwrap().pending.is_empty());
443        drop(held);
444    }
445
446    /// A `kind = "command"` agent running `script`; `LEASE` names the turn
447    /// lease file so the script can look at it while the turn is in flight.
448    fn scripted(tmp: &std::path::Path, script: &str, lease: &std::path::Path) -> Config {
449        let path = tmp.join("mock-consult-agent.sh");
450        std::fs::write(&path, script).expect("write mock");
451        Config {
452            agents: vec![AgentSpec {
453                id: "mock".to_owned(),
454                kind: AgentKind::Command,
455                model: None,
456                command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
457                extra_args: Vec::new(),
458                env: BTreeMap::from([("LEASE".to_owned(), lease.to_string_lossy().into_owned())]),
459                prompt_delivery: None,
460            }],
461            ..Config::default()
462        }
463    }
464
465    /// Run `start_turn` with `script` as the agent and hand back its result,
466    /// the store and the talk, for the caller to claim the lease afterwards.
467    async fn run_turn(
468        script: &str,
469    ) -> (
470        tempfile::TempDir,
471        Result<Started>,
472        Questions,
473        Question,
474        Talk,
475    ) {
476        crate::run::set_home(crate::run::test_home());
477        let (tmp, store, talk) = talks();
478        let questions = Questions::at(tmp.path().join("questions"));
479        let mut q = question("implement");
480        questions.put(&mut q).expect("put");
481        let cfg = scripted(tmp.path(), script, &store.turn_path(&talk.id));
482        let got = start_turn(&questions, &store, &q, &talk, &cfg).await;
483        (tmp, got, questions, q, talk)
484    }
485
486    #[tokio::test]
487    async fn the_lease_is_held_during_the_turn_and_free_once_it_ends() {
488        // exit 7 if the lease file is missing while the agent runs.
489        let (tmp, got, _questions, _q, talk) =
490            run_turn("#!/bin/sh\ncat >/dev/null\n[ -f \"$LEASE\" ] || exit 7\nprintf 'ok\\n'\n")
491                .await;
492        assert_eq!(got.expect("turn"), Started::Answered);
493        let other = Talks::at(tmp.path().join("talks"));
494        assert!(
495            other.claim_turn(&talk.id).expect("claim").is_some(),
496            "free once the turn ended"
497        );
498    }
499
500    #[tokio::test]
501    async fn the_lease_is_free_again_after_a_failed_turn() {
502        let (tmp, got, questions, q, talk) = run_turn("#!/bin/sh\ncat >/dev/null\nexit 3\n").await;
503        assert!(got.is_err(), "the agent failed");
504        assert!(
505            questions.get(&q.id).unwrap().consult.is_some(),
506            "the turn got as far as running"
507        );
508        let other = Talks::at(tmp.path().join("talks"));
509        assert!(
510            other.claim_turn(&talk.id).expect("claim").is_some(),
511            "free after the failed turn"
512        );
513    }
514
515    #[test]
516    fn origin_talk_is_decided_in_one_table() {
517        let (_tmp, _store, talk) = talks();
518        let mut closed = talk.clone();
519        closed.status = TalkStatus::Closed;
520        let chat = from_chat(&talk);
521        let q = question("implement");
522
523        let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
524        assert_eq!(
525            hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
526            Some(talk.id.clone())
527        );
528        // Not filed from a chat.
529        assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
530        let other = task(Source::Agent {
531            run: talk.id.clone(),
532            node: "implement".to_owned(),
533        });
534        assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
535        // The conversation is gone, or no longer takes turns.
536        assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
537        assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
538        // No task owns the question.
539        assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
540        // Merge approvals can be consulted; release notices cannot.
541        assert!(
542            hit(
543                std::slice::from_ref(&chat),
544                std::slice::from_ref(&talk),
545                &question(crate::land::APPROVAL_NODE)
546            )
547            .is_some()
548        );
549        assert!(
550            hit(
551                std::slice::from_ref(&chat),
552                std::slice::from_ref(&talk),
553                &question(crate::bump::NOTICE_NODE)
554            )
555            .is_none()
556        );
557        // A settled question has nothing left to ask.
558        let mut answered = question("implement");
559        answered
560            .answer(crate::ask::Answer::Choice("Redis".to_owned()))
561            .unwrap();
562        assert!(hit(&[chat], &[talk], &answered).is_none());
563    }
564
565    #[test]
566    fn approval_origin_requires_an_open_chat_task_and_question() {
567        let (_tmp, _store, talk) = talks();
568        let chat = from_chat(&talk);
569        let mut q = question(crate::land::APPROVAL_NODE);
570        assert!(origin_talk(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
571        assert!(
572            origin_talk(
573                &[task(Source::Agent {
574                    run: talk.id.clone(),
575                    node: "implement".into(),
576                })],
577                std::slice::from_ref(&talk),
578                &q
579            )
580            .is_none()
581        );
582        assert!(origin_talk(std::slice::from_ref(&chat), &[], &q).is_none());
583        let mut closed = talk.clone();
584        closed.status = TalkStatus::Closed;
585        assert!(origin_talk(std::slice::from_ref(&chat), &[closed], &q).is_none());
586        q.abandon("expired");
587        assert!(
588            origin_talk(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).is_none()
589        );
590        let mut q = question(crate::land::APPROVAL_NODE);
591        q.answer(crate::ask::Answer::Choice("Redis".into()))
592            .unwrap();
593        assert!(origin_talk(&[chat], &[talk], &q).is_none());
594    }
595
596    #[test]
597    fn a_conductor_question_is_found_through_the_task_id() {
598        let (_tmp, _store, talk) = talks();
599        let chat = from_chat(&talk);
600        let mut q = question(crate::conduct::NODE);
601        q.run = chat.id.clone();
602        assert!(origin_talk(&[chat], &[talk], &q).is_some());
603    }
604
605    fn followup_of(parent: &Task, n: u32) -> Task {
606        let mut t = Task::new(
607            "f".to_owned(),
608            "Fix".to_owned(),
609            PathBuf::from("/repo"),
610            Source::Agent {
611                run: format!("merged-{n}"),
612                node: "followup".to_owned(),
613            },
614        );
615        t.start(format!("run-f{n}"));
616        t.followup = Some(crate::queue::FollowUp {
617            run: format!("merged-{n}"),
618            origin_task: Some(parent.id.clone()),
619            pr: "https://example.invalid/pr/1".to_owned(),
620            findings: Vec::new(),
621            generation: n,
622        });
623        t
624    }
625
626    fn q_for(t: &Task) -> Question {
627        let mut q = question("implement");
628        q.run = t.runs[0].clone();
629        q
630    }
631
632    #[test]
633    fn a_followup_traces_back_to_the_chat() {
634        let (_tmp, _store, talk) = talks();
635        let mut chat = from_chat(&talk);
636        chat.runs = vec!["chat-run".to_owned()];
637        let f1 = followup_of(&chat, 1);
638        let f2 = followup_of(&f1, 2);
639        let ts = [chat, f1.clone(), f2.clone()];
640        let id = Some(talk.id.clone());
641        let tk = std::slice::from_ref(&talk);
642        assert_eq!(origin_talk(&ts, tk, &q_for(&f1)).map(|t| t.id), id);
643        assert_eq!(origin_talk(&ts, tk, &q_for(&f2)).map(|t| t.id), id);
644    }
645
646    #[test]
647    fn a_followup_without_origin_task_is_found_through_runs() {
648        let (_tmp, _store, talk) = talks();
649        let mut chat = from_chat(&talk);
650        chat.runs = vec!["merged-1".to_owned()];
651        let mut f1 = followup_of(&chat, 1);
652        f1.followup.as_mut().unwrap().origin_task = None;
653        let q = q_for(&f1);
654        assert!(origin_talk(&[chat, f1], &[talk], &q).is_some());
655    }
656
657    #[test]
658    fn a_dangling_origin_task_has_no_chat() {
659        let (_tmp, _store, talk) = talks();
660        let mut chat = from_chat(&talk);
661        chat.runs = vec!["chat-run".to_owned()];
662        let mut f1 = followup_of(&chat, 1);
663        f1.followup.as_mut().unwrap().origin_task = Some("gone".to_owned());
664        let q = q_for(&f1);
665        assert!(origin_talk(&[chat, f1], &[talk], &q).is_none());
666    }
667
668    #[test]
669    fn a_followup_cycle_ends_without_a_chat() {
670        let (_tmp, _store, talk) = talks();
671        let mut a = followup_of(&from_chat(&talk), 1);
672        let mut b = followup_of(&a, 2);
673        a.followup.as_mut().unwrap().origin_task = Some(b.id.clone());
674        b.followup.as_mut().unwrap().origin_task = Some(a.id.clone());
675        let q = q_for(&a);
676        assert!(origin_talk(&[a, b], &[talk], &q).is_none());
677    }
678
679    #[test]
680    fn pending_consults_follow_the_store_not_the_turn_text() {
681        let (tmp, store, talk) = talks();
682        let questions = Questions::at(tmp.path().join("questions"));
683        assert!(!pending_consults(&questions, &talk.id), "empty store");
684
685        let mut plain = question("implement");
686        questions.put(&mut plain).unwrap();
687        assert!(!pending_consults(&questions, &talk.id), "no consult record");
688
689        let mut q = question("implement");
690        questions.put(&mut q).unwrap();
691        begin(&questions, &store, &q, &talk).unwrap();
692        assert!(pending_consults(&questions, &talk.id));
693        assert!(!pending_consults(&questions, "other-talk"));
694
695        questions
696            .update(&q.id, |q| {
697                q.abandon("test");
698                Ok(())
699            })
700            .unwrap();
701        assert!(!pending_consults(&questions, &talk.id), "closed question");
702    }
703
704    #[test]
705    fn begin_queues_once_and_leaves_the_question_open() {
706        let (tmp, store, talk) = talks();
707        let questions = Questions::at(tmp.path().join("questions"));
708        let mut q = question("implement");
709        questions.put(&mut q).unwrap();
710
711        assert!(begin(&questions, &store, &q, &talk).unwrap());
712        assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
713
714        let after = questions.get(&q.id).unwrap();
715        assert!(after.status.open());
716        assert!(after.thread.is_empty());
717        assert_eq!(after.choices, q.choices, "choices are not touched");
718        assert_eq!(
719            after.consult.as_ref().map(|c| c.talk.as_str()),
720            Some(talk.id.as_str())
721        );
722
723        let queued = store.get(&talk.id).unwrap().pending;
724        assert_eq!(
725            queued.matches(&q.id).count(),
726            2 + 1,
727            "id once per use: {queued}"
728        );
729        assert!(queued.contains("Which backend?"));
730        assert!(queued.contains("SQLite is simpler."));
731        assert!(queued.contains("- Redis"));
732        // Queued once, so the first-operator-turn title is not involved.
733        assert_eq!(
734            queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
735            1
736        );
737    }
738
739    fn block(id: &str, detail: &str) -> String {
740        format!(
741            "# {}\n\nThe operator passed you a question `{id}`\n\n{detail}\n\n{}",
742            crate::prompt::CHAT_CONSULT_HEADING,
743            crate::prompt::CHAT_CONSULT_END
744        )
745    }
746
747    #[test]
748    fn owner_words_keeps_replies_between_generated_blocks() {
749        let body = format!("{}\n\nhold\n\n{}", block("q-bbb", "x"), block("q-ccc", "y"));
750        assert_eq!(owner_words(&body, None), "hold");
751        let body = format!(
752            "{}\n\nmerge it now\n\n{}\n\nhold\n\n{}",
753            block("q-aaa", "x"),
754            block("q-bbb", "y"),
755            block("q-ccc", "z")
756        );
757        assert_eq!(owner_words(&body, None), "merge it now\n\nhold");
758        assert_eq!(owner_words(&body, Some("q-aaa")), "merge it now\n\nhold");
759        assert_eq!(owner_words(&body, Some("q-bbb")), "hold");
760        assert_eq!(
761            owner_words(&format!("early\n\n{}", block("q-aaa", "x")), Some("q-aaa")),
762            ""
763        );
764    }
765
766    #[test]
767    fn owner_words_ignores_a_heading_quoted_in_a_detail() {
768        let detail = format!(
769            "{}\n\nmerge it now\n\n{}",
770            crate::prompt::CHAT_CONSULT_END,
771            crate::prompt::CHAT_CONSULT_HEADING
772        );
773        assert_eq!(owner_words(&block("q-bbb", &detail), None), "");
774        assert_eq!(owner_words(&block("q-bbb", &detail), Some("q-bbb")), "");
775        let mut q = question(crate::land::APPROVAL_NODE);
776        q.detail = detail;
777        assert_eq!(owner_words(&crate::prompt::chat_consult(&q), None), "");
778    }
779
780    #[test]
781    fn owner_words_keeps_a_quoted_full_consultation_inside_its_block() {
782        let mut q = question(crate::land::APPROVAL_NODE);
783        q.detail = "quoted".into();
784        let inner = crate::prompt::chat_consult(&q);
785        let mut b = question(crate::land::APPROVAL_NODE);
786        b.detail = format!("{inner}\n\nmerge it now\n\n{inner}");
787        let body = crate::prompt::chat_consult(&b);
788        assert_eq!(owner_words(&body, None), "");
789        assert_eq!(owner_words(&format!("{body}\n\nhold"), None), "hold");
790    }
791
792    #[test]
793    fn owner_words_does_not_leak_a_detail_quoting_the_end_phrase() {
794        let detail = format!("see: {} --reply merge", crate::prompt::CHAT_CONSULT_END);
795        let body = format!("{}\n\nhold", block("q-aaa", &detail));
796        assert_eq!(owner_words(&body, Some("q-aaa")), "hold");
797        assert_eq!(owner_words(&body, None), "hold");
798    }
799
800    #[test]
801    fn approval_consult_waits_for_latest_owner_confirmation() {
802        let (tmp, store, mut talk) = talks();
803        let questions = Questions::at(tmp.path().join("questions"));
804        let mut q = question(crate::land::APPROVAL_NODE);
805        q.choices = vec![crate::land::APPROVE.into(), crate::land::HOLD.into()];
806        questions.put(&mut q).unwrap();
807        assert!(begin(&questions, &store, &q, &talk).unwrap());
808        assert!(!begin(&questions, &store, &q, &talk).unwrap());
809        q = questions.get(&q.id).unwrap();
810        assert!(q.status.open());
811        assert!(q.answer.is_none());
812        assert!(q.thread.is_empty());
813        let queued = store.get(&talk.id).unwrap().pending;
814        assert!(queued.contains("Never answer it yourself"));
815        assert!(queued.contains("Silence holds"));
816        assert!(queued.contains("--reply merge --quote"));
817        assert!(!queued.contains("answer it yourself with"));
818
819        let owner_turn = |body: &str| talk::Turn {
820            who: talk::Who::Operator,
821            body: body.into(),
822            at: Timestamp::now(),
823            attachments: Vec::new(),
824            usage: None,
825        };
826        // An approval the owner gave before the hand-over is not reusable.
827        let mut old = talk.clone();
828        old.turns
829            .insert(0, owner_turn("Please merge PR 12 after review"));
830        store.put(&mut old).unwrap();
831        talk.turns = old.turns.clone();
832        // The drained consult draft is stored as an operator turn; it must not
833        // count as the owner's words even though it contains "merge".
834        talk.turns.push(owner_turn(&queued));
835        store.put(&mut talk).unwrap();
836        assert!(validate_answer(&q, &store, &talk.id, CHAT_NODE, "merge", Some("merge")).is_err());
837        assert!(
838            validate_answer(
839                &q,
840                &store,
841                &talk.id,
842                CHAT_NODE,
843                "merge",
844                Some("Please merge PR 12")
845            )
846            .is_err()
847        );
848        // An owner reply coalesced onto the same draft still counts.
849        let n = talk.turns.len();
850        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now");
851        store.put(&mut talk).unwrap();
852        assert!(
853            validate_answer(
854                &q,
855                &store,
856                &talk.id,
857                CHAT_NODE,
858                "merge",
859                Some("merge it now")
860            )
861            .is_ok()
862        );
863        // A later retraction in the same coalesced turn wins.
864        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now\n\nhold");
865        store.put(&mut talk).unwrap();
866        assert!(
867            validate_answer(
868                &q,
869                &store,
870                &talk.id,
871                CHAT_NODE,
872                "merge",
873                Some("merge it now")
874            )
875            .is_err()
876        );
877        assert!(validate_answer(&q, &store, &talk.id, CHAT_NODE, "hold", Some("hold")).is_ok());
878        // So does a reply still waiting in the draft.
879        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now");
880        talk.pending = "hold".into();
881        store.put(&mut talk).unwrap();
882        assert!(
883            validate_answer(
884                &q,
885                &store,
886                &talk.id,
887                CHAT_NODE,
888                "merge",
889                Some("merge it now")
890            )
891            .is_err()
892        );
893        // Another question's hand-over queued after the retraction hides nothing.
894        talk.pending = format!("hold\n\n{queued}");
895        store.put(&mut talk).unwrap();
896        assert!(
897            validate_answer(
898                &q,
899                &store,
900                &talk.id,
901                CHAT_NODE,
902                "merge",
903                Some("merge it now")
904            )
905            .is_err()
906        );
907        talk.pending.clear();
908        talk.turns[n - 1].body = queued.clone();
909        store.put(&mut talk).unwrap();
910        talk.turns
911            .push(owner_turn("Merge this pull request please"));
912        store.put(&mut talk).unwrap();
913        let check = |q: &Question, run: &str, reply: &str, quote: Option<&str>| {
914            validate_answer(q, &store, run, CHAT_NODE, reply, quote)
915        };
916        assert!(check(&q, "other-talk", "merge", Some("Merge this")).is_err());
917        assert!(check(&q, &talk.id, "merge", None).is_err());
918        assert!(check(&q, &talk.id, "merge", Some("never said")).is_err());
919        assert!(
920            check(
921                &q,
922                &talk.id,
923                "merge",
924                Some("Merge this pull request please")
925            )
926            .is_ok()
927        );
928        assert!(check(&q, &talk.id, "hold", Some("hold")).is_err());
929        talk.turns
930            .push(owner_turn("Wait, explain the checks first"));
931        store.put(&mut talk).unwrap();
932        assert!(
933            check(
934                &q,
935                &talk.id,
936                "merge",
937                Some("Merge this pull request please")
938            )
939            .is_err()
940        );
941        talk.turns.push(owner_turn("Please hold"));
942        store.put(&mut talk).unwrap();
943        assert!(check(&q, &talk.id, "hold", Some("hold")).is_err());
944        talk.turns.push(owner_turn("hold"));
945        store.put(&mut talk).unwrap();
946        assert!(check(&q, &talk.id, "hold", Some("hold")).is_ok());
947        talk.turns
948            .push(owner_turn("Merge this pull request please"));
949        store.put(&mut talk).unwrap();
950        q.abandon("approval expired or head changed");
951        assert!(
952            check(
953                &q,
954                &talk.id,
955                "merge",
956                Some("Merge this pull request please")
957            )
958            .is_err()
959        );
960        assert!(
961            q.answer(crate::ask::Answer::Choice("merge".into()))
962                .is_err()
963        );
964    }
965
966    #[test]
967    fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
968        let (tmp, store, talk) = talks();
969        let questions = Questions::at(tmp.path().join("questions"));
970        let mut q = question("implement");
971        questions.put(&mut q).unwrap();
972        std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
973        assert!(begin(&questions, &store, &q, &talk).is_err());
974        assert!(questions.get(&q.id).unwrap().consult.is_none());
975    }
976
977    #[test]
978    fn an_older_question_file_reads_without_a_consult() {
979        let mut q = question("implement");
980        q.schema = 5;
981        let mut v = serde_json::to_value(&q).unwrap();
982        v.as_object_mut().unwrap().remove("consult");
983        let back: Question = serde_json::from_value(v).unwrap();
984        assert!(back.consult.is_none());
985    }
986}