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, false);
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))
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. The generated
211/// block runs from [`CHAT_CONSULT_HEADING`](crate::prompt::CHAT_CONSULT_HEADING)
212/// to the last [`CHAT_CONSULT_END`](crate::prompt::CHAT_CONSULT_END); an owner
213/// reply that happens to contain that phrase is cut too, which only ever
214/// refuses. With `after_block_only` (the turn carrying the hand-over itself)
215/// words queued before the block predate it and are dropped.
216fn owner_words(body: &str, after_block_only: bool) -> String {
217    use crate::prompt::{CHAT_CONSULT_END, CHAT_CONSULT_HEADING};
218    let Some(start) = body.find(CHAT_CONSULT_HEADING) else {
219        return body.trim().to_owned();
220    };
221    let after = match body.rfind(CHAT_CONSULT_END) {
222        Some(e) if e >= start => &body[e + CHAT_CONSULT_END.len()..],
223        _ => "",
224    };
225    if after_block_only {
226        return after.trim().to_owned();
227    }
228    let before = &body[..start];
229    let before = before.strip_suffix("# ").unwrap_or(before);
230    format!("{before}\n\n{after}").trim().to_owned()
231}
232
233/// What [`start_turn`] did.
234#[derive(Debug, PartialEq, Eq)]
235pub enum Started {
236    /// The question was handed over and the chat's turn ran to its end.
237    Answered,
238    /// Another process holds the talk's turn lease; nothing was changed.
239    Busy,
240    /// The question had already been handed to the chat; nothing was changed.
241    Nothing,
242}
243
244/// [`begin`], then run the chat's turn here, under the talk's own lease.
245///
246/// The lease is taken *before* anything is written, so a talk somebody else is
247/// running leaves no consult record and no draft behind (`Started::Busy`) and
248/// the caller can simply try again later. A turn that fails after `begin`
249/// leaves the record and the draft, so the chat can resume it.
250pub async fn start_turn(
251    questions: &Questions,
252    talks: &Talks,
253    q: &Question,
254    talk: &Talk,
255    cfg: &Config,
256) -> Result<Started> {
257    if q.consult.is_some() {
258        return Ok(Started::Nothing);
259    }
260    let Some(mut lease) = talks.claim_turn(&talk.id)? else {
261        return Ok(Started::Busy);
262    };
263    if !begin(questions, talks, q, talk)? {
264        return Ok(Started::Nothing);
265    }
266    let mut talk = talks.get(&talk.id)?;
267    // Other starters that found the lease held queued drafts and left them to
268    // us, so drain until nothing is left (as the web's drain loop does).
269    // A failed turn is kept and reported at the end, not returned at once:
270    // drafts accepted meanwhile are still owed an answer.
271    let mut failed = None;
272    loop {
273        while let Some(text) = talk::drain(&mut talk, talks)? {
274            if let Err(e) = talk::respond(&lease, &mut talk, talks, cfg, &text).await {
275                failed.get_or_insert(e);
276            }
277            if !lease.beat()? {
278                bail!("the turn lease for chat {} was lost", talk.short());
279            }
280        }
281        // A draft queued after the last drain, before the lease is gone, was
282        // left to us by a starter that found the lease held. Release first,
283        // then look again: whoever queues later either sees no lease (and
284        // starts its own turn) or is seen by this re-check.
285        drop(lease);
286        talk = talks.get(&talk.id)?;
287        let owed = talk.status.open()
288            && (!talk.pending.is_empty() || !talk.pending_attachments.is_empty());
289        if !owed {
290            break;
291        }
292        // Somebody else took the lease meanwhile: they drain it.
293        let Some(again) = talks.claim_turn(&talk.id)? else {
294            break;
295        };
296        lease = again;
297    }
298    if let Some(e) = failed {
299        return Err(e);
300    }
301    Ok(Started::Answered)
302}
303
304#[cfg(test)]
305mod tests {
306    use std::collections::BTreeMap;
307    use std::path::PathBuf;
308
309    use super::*;
310    use crate::config::{AgentKind, AgentSpec, Config};
311    use crate::queue::Source;
312    use crate::talk::TalkStatus;
313
314    const RUN: &str = "20260902-000000-beef";
315
316    fn talks() -> (tempfile::TempDir, Talks, Talk) {
317        let tmp = tempfile::tempdir().expect("tempdir");
318        let store = Talks::at(tmp.path().join("talks"));
319        let cfg = Config {
320            agents: vec![AgentSpec {
321                id: "mock".to_owned(),
322                kind: AgentKind::Command,
323                model: None,
324                command: vec!["true".to_owned()],
325                extra_args: Vec::new(),
326                env: BTreeMap::new(),
327                prompt_delivery: None,
328            }],
329            ..Config::default()
330        };
331        let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
332        (tmp, store, talk)
333    }
334
335    fn task(source: Source) -> Task {
336        let mut t = Task::new(
337            "t".to_owned(),
338            "Do it".to_owned(),
339            PathBuf::from("/repo"),
340            source,
341        );
342        t.start(RUN.to_owned());
343        t
344    }
345
346    fn from_chat(talk: &Talk) -> Task {
347        task(Source::Agent {
348            run: talk.id.clone(),
349            node: CHAT_NODE.to_owned(),
350        })
351    }
352
353    fn question(node: &str) -> Question {
354        Question::new(
355            RUN.to_owned(),
356            node.to_owned(),
357            "impl-A".to_owned(),
358            "Which backend?".to_owned(),
359            "SQLite is simpler.".to_owned(),
360            vec!["SQLite".to_owned(), "Redis".to_owned()],
361        )
362    }
363
364    #[tokio::test]
365    async fn start_turn_is_busy_and_writes_nothing_while_the_lease_is_held() {
366        let (tmp, store, talk) = talks();
367        let questions = Questions::at(tmp.path().join("questions"));
368        let mut q = question("implement");
369        questions.put(&mut q).expect("put");
370        let held = store.claim_turn(&talk.id).expect("claim").expect("first");
371
372        let got = start_turn(&questions, &store, &q, &talk, &Config::default())
373            .await
374            .expect("start");
375        assert_eq!(got, Started::Busy);
376        assert!(questions.get(&q.id).unwrap().consult.is_none());
377        assert!(store.get(&talk.id).unwrap().pending.is_empty());
378        drop(held);
379    }
380
381    /// A `kind = "command"` agent running `script`; `LEASE` names the turn
382    /// lease file so the script can look at it while the turn is in flight.
383    fn scripted(tmp: &std::path::Path, script: &str, lease: &std::path::Path) -> Config {
384        let path = tmp.join("mock-consult-agent.sh");
385        std::fs::write(&path, script).expect("write mock");
386        Config {
387            agents: vec![AgentSpec {
388                id: "mock".to_owned(),
389                kind: AgentKind::Command,
390                model: None,
391                command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
392                extra_args: Vec::new(),
393                env: BTreeMap::from([("LEASE".to_owned(), lease.to_string_lossy().into_owned())]),
394                prompt_delivery: None,
395            }],
396            ..Config::default()
397        }
398    }
399
400    /// Run `start_turn` with `script` as the agent and hand back its result,
401    /// the store and the talk, for the caller to claim the lease afterwards.
402    async fn run_turn(
403        script: &str,
404    ) -> (
405        tempfile::TempDir,
406        Result<Started>,
407        Questions,
408        Question,
409        Talk,
410    ) {
411        crate::run::set_home(crate::run::test_home());
412        let (tmp, store, talk) = talks();
413        let questions = Questions::at(tmp.path().join("questions"));
414        let mut q = question("implement");
415        questions.put(&mut q).expect("put");
416        let cfg = scripted(tmp.path(), script, &store.turn_path(&talk.id));
417        let got = start_turn(&questions, &store, &q, &talk, &cfg).await;
418        (tmp, got, questions, q, talk)
419    }
420
421    #[tokio::test]
422    async fn the_lease_is_held_during_the_turn_and_free_once_it_ends() {
423        // exit 7 if the lease file is missing while the agent runs.
424        let (tmp, got, _questions, _q, talk) =
425            run_turn("#!/bin/sh\ncat >/dev/null\n[ -f \"$LEASE\" ] || exit 7\nprintf 'ok\\n'\n")
426                .await;
427        assert_eq!(got.expect("turn"), Started::Answered);
428        let other = Talks::at(tmp.path().join("talks"));
429        assert!(
430            other.claim_turn(&talk.id).expect("claim").is_some(),
431            "free once the turn ended"
432        );
433    }
434
435    #[tokio::test]
436    async fn the_lease_is_free_again_after_a_failed_turn() {
437        let (tmp, got, questions, q, talk) = run_turn("#!/bin/sh\ncat >/dev/null\nexit 3\n").await;
438        assert!(got.is_err(), "the agent failed");
439        assert!(
440            questions.get(&q.id).unwrap().consult.is_some(),
441            "the turn got as far as running"
442        );
443        let other = Talks::at(tmp.path().join("talks"));
444        assert!(
445            other.claim_turn(&talk.id).expect("claim").is_some(),
446            "free after the failed turn"
447        );
448    }
449
450    #[test]
451    fn origin_talk_is_decided_in_one_table() {
452        let (_tmp, _store, talk) = talks();
453        let mut closed = talk.clone();
454        closed.status = TalkStatus::Closed;
455        let chat = from_chat(&talk);
456        let q = question("implement");
457
458        let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
459        assert_eq!(
460            hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
461            Some(talk.id.clone())
462        );
463        // Not filed from a chat.
464        assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
465        let other = task(Source::Agent {
466            run: talk.id.clone(),
467            node: "implement".to_owned(),
468        });
469        assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
470        // The conversation is gone, or no longer takes turns.
471        assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
472        assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
473        // No task owns the question.
474        assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
475        // Merge approvals can be consulted; release notices cannot.
476        assert!(
477            hit(
478                std::slice::from_ref(&chat),
479                std::slice::from_ref(&talk),
480                &question(crate::land::APPROVAL_NODE)
481            )
482            .is_some()
483        );
484        assert!(
485            hit(
486                std::slice::from_ref(&chat),
487                std::slice::from_ref(&talk),
488                &question(crate::bump::NOTICE_NODE)
489            )
490            .is_none()
491        );
492        // A settled question has nothing left to ask.
493        let mut answered = question("implement");
494        answered
495            .answer(crate::ask::Answer::Choice("Redis".to_owned()))
496            .unwrap();
497        assert!(hit(&[chat], &[talk], &answered).is_none());
498    }
499
500    #[test]
501    fn approval_origin_requires_an_open_chat_task_and_question() {
502        let (_tmp, _store, talk) = talks();
503        let chat = from_chat(&talk);
504        let mut q = question(crate::land::APPROVAL_NODE);
505        assert!(origin_talk(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
506        assert!(
507            origin_talk(
508                &[task(Source::Agent {
509                    run: talk.id.clone(),
510                    node: "implement".into(),
511                })],
512                std::slice::from_ref(&talk),
513                &q
514            )
515            .is_none()
516        );
517        assert!(origin_talk(std::slice::from_ref(&chat), &[], &q).is_none());
518        let mut closed = talk.clone();
519        closed.status = TalkStatus::Closed;
520        assert!(origin_talk(std::slice::from_ref(&chat), &[closed], &q).is_none());
521        q.abandon("expired");
522        assert!(
523            origin_talk(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).is_none()
524        );
525        let mut q = question(crate::land::APPROVAL_NODE);
526        q.answer(crate::ask::Answer::Choice("Redis".into()))
527            .unwrap();
528        assert!(origin_talk(&[chat], &[talk], &q).is_none());
529    }
530
531    #[test]
532    fn a_conductor_question_is_found_through_the_task_id() {
533        let (_tmp, _store, talk) = talks();
534        let chat = from_chat(&talk);
535        let mut q = question(crate::conduct::NODE);
536        q.run = chat.id.clone();
537        assert!(origin_talk(&[chat], &[talk], &q).is_some());
538    }
539
540    fn followup_of(parent: &Task, n: u32) -> Task {
541        let mut t = Task::new(
542            "f".to_owned(),
543            "Fix".to_owned(),
544            PathBuf::from("/repo"),
545            Source::Agent {
546                run: format!("merged-{n}"),
547                node: "followup".to_owned(),
548            },
549        );
550        t.start(format!("run-f{n}"));
551        t.followup = Some(crate::queue::FollowUp {
552            run: format!("merged-{n}"),
553            origin_task: Some(parent.id.clone()),
554            pr: "https://example.invalid/pr/1".to_owned(),
555            findings: Vec::new(),
556            generation: n,
557        });
558        t
559    }
560
561    fn q_for(t: &Task) -> Question {
562        let mut q = question("implement");
563        q.run = t.runs[0].clone();
564        q
565    }
566
567    #[test]
568    fn a_followup_traces_back_to_the_chat() {
569        let (_tmp, _store, talk) = talks();
570        let mut chat = from_chat(&talk);
571        chat.runs = vec!["chat-run".to_owned()];
572        let f1 = followup_of(&chat, 1);
573        let f2 = followup_of(&f1, 2);
574        let ts = [chat, f1.clone(), f2.clone()];
575        let id = Some(talk.id.clone());
576        let tk = std::slice::from_ref(&talk);
577        assert_eq!(origin_talk(&ts, tk, &q_for(&f1)).map(|t| t.id), id);
578        assert_eq!(origin_talk(&ts, tk, &q_for(&f2)).map(|t| t.id), id);
579    }
580
581    #[test]
582    fn a_followup_without_origin_task_is_found_through_runs() {
583        let (_tmp, _store, talk) = talks();
584        let mut chat = from_chat(&talk);
585        chat.runs = vec!["merged-1".to_owned()];
586        let mut f1 = followup_of(&chat, 1);
587        f1.followup.as_mut().unwrap().origin_task = None;
588        let q = q_for(&f1);
589        assert!(origin_talk(&[chat, f1], &[talk], &q).is_some());
590    }
591
592    #[test]
593    fn a_dangling_origin_task_has_no_chat() {
594        let (_tmp, _store, talk) = talks();
595        let mut chat = from_chat(&talk);
596        chat.runs = vec!["chat-run".to_owned()];
597        let mut f1 = followup_of(&chat, 1);
598        f1.followup.as_mut().unwrap().origin_task = Some("gone".to_owned());
599        let q = q_for(&f1);
600        assert!(origin_talk(&[chat, f1], &[talk], &q).is_none());
601    }
602
603    #[test]
604    fn a_followup_cycle_ends_without_a_chat() {
605        let (_tmp, _store, talk) = talks();
606        let mut a = followup_of(&from_chat(&talk), 1);
607        let mut b = followup_of(&a, 2);
608        a.followup.as_mut().unwrap().origin_task = Some(b.id.clone());
609        b.followup.as_mut().unwrap().origin_task = Some(a.id.clone());
610        let q = q_for(&a);
611        assert!(origin_talk(&[a, b], &[talk], &q).is_none());
612    }
613
614    #[test]
615    fn pending_consults_follow_the_store_not_the_turn_text() {
616        let (tmp, store, talk) = talks();
617        let questions = Questions::at(tmp.path().join("questions"));
618        assert!(!pending_consults(&questions, &talk.id), "empty store");
619
620        let mut plain = question("implement");
621        questions.put(&mut plain).unwrap();
622        assert!(!pending_consults(&questions, &talk.id), "no consult record");
623
624        let mut q = question("implement");
625        questions.put(&mut q).unwrap();
626        begin(&questions, &store, &q, &talk).unwrap();
627        assert!(pending_consults(&questions, &talk.id));
628        assert!(!pending_consults(&questions, "other-talk"));
629
630        questions
631            .update(&q.id, |q| {
632                q.abandon("test");
633                Ok(())
634            })
635            .unwrap();
636        assert!(!pending_consults(&questions, &talk.id), "closed question");
637    }
638
639    #[test]
640    fn begin_queues_once_and_leaves_the_question_open() {
641        let (tmp, store, talk) = talks();
642        let questions = Questions::at(tmp.path().join("questions"));
643        let mut q = question("implement");
644        questions.put(&mut q).unwrap();
645
646        assert!(begin(&questions, &store, &q, &talk).unwrap());
647        assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
648
649        let after = questions.get(&q.id).unwrap();
650        assert!(after.status.open());
651        assert!(after.thread.is_empty());
652        assert_eq!(after.choices, q.choices, "choices are not touched");
653        assert_eq!(
654            after.consult.as_ref().map(|c| c.talk.as_str()),
655            Some(talk.id.as_str())
656        );
657
658        let queued = store.get(&talk.id).unwrap().pending;
659        assert_eq!(
660            queued.matches(&q.id).count(),
661            2 + 1,
662            "id once per use: {queued}"
663        );
664        assert!(queued.contains("Which backend?"));
665        assert!(queued.contains("SQLite is simpler."));
666        assert!(queued.contains("- Redis"));
667        // Queued once, so the first-operator-turn title is not involved.
668        assert_eq!(
669            queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
670            1
671        );
672    }
673
674    #[test]
675    fn approval_consult_waits_for_latest_owner_confirmation() {
676        let (tmp, store, mut talk) = talks();
677        let questions = Questions::at(tmp.path().join("questions"));
678        let mut q = question(crate::land::APPROVAL_NODE);
679        q.choices = vec![crate::land::APPROVE.into(), crate::land::HOLD.into()];
680        questions.put(&mut q).unwrap();
681        assert!(begin(&questions, &store, &q, &talk).unwrap());
682        assert!(!begin(&questions, &store, &q, &talk).unwrap());
683        q = questions.get(&q.id).unwrap();
684        assert!(q.status.open());
685        assert!(q.answer.is_none());
686        assert!(q.thread.is_empty());
687        let queued = store.get(&talk.id).unwrap().pending;
688        assert!(queued.contains("Never answer it yourself"));
689        assert!(queued.contains("Silence holds"));
690        assert!(queued.contains("--reply merge --quote"));
691        assert!(!queued.contains("answer it yourself with"));
692
693        let owner_turn = |body: &str| talk::Turn {
694            who: talk::Who::Operator,
695            body: body.into(),
696            at: Timestamp::now(),
697            attachments: Vec::new(),
698            usage: None,
699        };
700        // An approval the owner gave before the hand-over is not reusable.
701        let mut old = talk.clone();
702        old.turns
703            .insert(0, owner_turn("Please merge PR 12 after review"));
704        store.put(&mut old).unwrap();
705        talk.turns = old.turns.clone();
706        // The drained consult draft is stored as an operator turn; it must not
707        // count as the owner's words even though it contains "merge".
708        talk.turns.push(owner_turn(&queued));
709        store.put(&mut talk).unwrap();
710        assert!(validate_answer(&q, &store, &talk.id, CHAT_NODE, "merge", Some("merge")).is_err());
711        assert!(
712            validate_answer(
713                &q,
714                &store,
715                &talk.id,
716                CHAT_NODE,
717                "merge",
718                Some("Please merge PR 12")
719            )
720            .is_err()
721        );
722        // An owner reply coalesced onto the same draft still counts.
723        let n = talk.turns.len();
724        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now");
725        store.put(&mut talk).unwrap();
726        assert!(
727            validate_answer(
728                &q,
729                &store,
730                &talk.id,
731                CHAT_NODE,
732                "merge",
733                Some("merge it now")
734            )
735            .is_ok()
736        );
737        // A later retraction in the same coalesced turn wins.
738        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now\n\nhold");
739        store.put(&mut talk).unwrap();
740        assert!(
741            validate_answer(
742                &q,
743                &store,
744                &talk.id,
745                CHAT_NODE,
746                "merge",
747                Some("merge it now")
748            )
749            .is_err()
750        );
751        assert!(validate_answer(&q, &store, &talk.id, CHAT_NODE, "hold", Some("hold")).is_ok());
752        // So does a reply still waiting in the draft.
753        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now");
754        talk.pending = "hold".into();
755        store.put(&mut talk).unwrap();
756        assert!(
757            validate_answer(
758                &q,
759                &store,
760                &talk.id,
761                CHAT_NODE,
762                "merge",
763                Some("merge it now")
764            )
765            .is_err()
766        );
767        // Another question's hand-over queued after the retraction hides nothing.
768        talk.pending = format!("hold\n\n{queued}");
769        store.put(&mut talk).unwrap();
770        assert!(
771            validate_answer(
772                &q,
773                &store,
774                &talk.id,
775                CHAT_NODE,
776                "merge",
777                Some("merge it now")
778            )
779            .is_err()
780        );
781        talk.pending.clear();
782        talk.turns[n - 1].body = queued.clone();
783        store.put(&mut talk).unwrap();
784        talk.turns
785            .push(owner_turn("Merge this pull request please"));
786        store.put(&mut talk).unwrap();
787        let check = |q: &Question, run: &str, reply: &str, quote: Option<&str>| {
788            validate_answer(q, &store, run, CHAT_NODE, reply, quote)
789        };
790        assert!(check(&q, "other-talk", "merge", Some("Merge this")).is_err());
791        assert!(check(&q, &talk.id, "merge", None).is_err());
792        assert!(check(&q, &talk.id, "merge", Some("never said")).is_err());
793        assert!(
794            check(
795                &q,
796                &talk.id,
797                "merge",
798                Some("Merge this pull request please")
799            )
800            .is_ok()
801        );
802        assert!(check(&q, &talk.id, "hold", Some("hold")).is_err());
803        talk.turns
804            .push(owner_turn("Wait, explain the checks first"));
805        store.put(&mut talk).unwrap();
806        assert!(
807            check(
808                &q,
809                &talk.id,
810                "merge",
811                Some("Merge this pull request please")
812            )
813            .is_err()
814        );
815        talk.turns.push(owner_turn("Please hold"));
816        store.put(&mut talk).unwrap();
817        assert!(check(&q, &talk.id, "hold", Some("hold")).is_err());
818        talk.turns.push(owner_turn("hold"));
819        store.put(&mut talk).unwrap();
820        assert!(check(&q, &talk.id, "hold", Some("hold")).is_ok());
821        talk.turns
822            .push(owner_turn("Merge this pull request please"));
823        store.put(&mut talk).unwrap();
824        q.abandon("approval expired or head changed");
825        assert!(
826            check(
827                &q,
828                &talk.id,
829                "merge",
830                Some("Merge this pull request please")
831            )
832            .is_err()
833        );
834        assert!(
835            q.answer(crate::ask::Answer::Choice("merge".into()))
836                .is_err()
837        );
838    }
839
840    #[test]
841    fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
842        let (tmp, store, talk) = talks();
843        let questions = Questions::at(tmp.path().join("questions"));
844        let mut q = question("implement");
845        questions.put(&mut q).unwrap();
846        std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
847        assert!(begin(&questions, &store, &q, &talk).is_err());
848        assert!(questions.get(&q.id).unwrap().consult.is_none());
849    }
850
851    #[test]
852    fn an_older_question_file_reads_without_a_consult() {
853        let mut q = question("implement");
854        q.schema = 5;
855        let mut v = serde_json::to_value(&q).unwrap();
856        v.as_object_mut().unwrap().remove("consult");
857        let back: Question = serde_json::from_value(v).unwrap();
858        assert!(back.consult.is_none());
859    }
860}