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/// Length of a consult stored before quoted markers were defused, which may
208/// quote the end phrase or a whole earlier consult in its detail. The closer is
209/// the generated tail's full final sentence (an owner reply does not repeat
210/// it), and a real heading met before it opens a nested block that needs its
211/// own closer. A block that never closes runs to the end of the text, which can
212/// only exclude more.
213fn legacy_len(rest: &str) -> usize {
214    use crate::prompt::CHAT_CONSULT_HEADING;
215    const CLOSERS: [&str; 2] = [
216        "cannot be revived. Do not edit the repository.",
217        "make: do not edit the repository.",
218    ];
219    let mut events: Vec<(usize, usize)> = Vec::new(); // (at, 0 = open, else closer length)
220    let first = rest.find(CHAT_CONSULT_HEADING).unwrap_or(0);
221    for (i, _) in rest.match_indices(CHAT_CONSULT_HEADING) {
222        if i > first
223            && rest[..i].ends_with("# ")
224            && rest[..i - 2].ends_with('\n')
225            && rest[i + CHAT_CONSULT_HEADING.len()..].starts_with("\n\nThe operator passed you")
226        {
227            events.push((i, 0));
228        }
229    }
230    for c in CLOSERS {
231        events.extend(rest.match_indices(c).map(|(i, _)| (i, c.len())));
232    }
233    events.sort_unstable();
234    let mut depth = 1usize;
235    for (at, len) in events {
236        if len == 0 {
237            depth += 1;
238        } else {
239            depth -= 1;
240            if depth == 0 {
241                return at + len;
242            }
243        }
244    }
245    rest.len()
246}
247
248/// The owner's own words in an operator turn, with magi's generated hand-over
249/// text cut out. `drain` stores a queued consult as an operator turn, and
250/// `talk::queue` joins an owner reply sent meanwhile onto the same draft, so
251/// the turn's origin cannot be told from the turn as a whole. Each generated
252/// block runs from its [`CHAT_CONSULT_HEADING`](crate::prompt::CHAT_CONSULT_HEADING)
253/// to the first [`CHAT_CONSULT_END`](crate::prompt::CHAT_CONSULT_END) after it
254/// (quoted markers are defused by `prompt::chat_consult`; a consult stored
255/// before that, recognised by its opening sentence, ends at the last one before
256/// the next heading); with no end the
257/// block runs to the next heading. Only the blocks are cut: owner replies
258/// between them are kept, in order, even when they contain the phrase. With `after_block_of` (the turn
259/// carrying the hand-over itself) everything up to the end of the last block
260/// naming that question id predates the hand-over and is dropped.
261fn owner_words(body: &str, after_block_of: Option<&str>) -> String {
262    use crate::prompt::{CHAT_CONSULT_END, CHAT_CONSULT_HEADING};
263    // A real heading is the generated one: `# ` at a line start, followed by
264    // the fixed opening sentence. A heading phrase quoted inside a question's
265    // detail is not a boundary; it stays inside the block it is part of.
266    let heads: Vec<usize> = body
267        .match_indices(CHAT_CONSULT_HEADING)
268        .map(|(i, _)| i)
269        .filter(|&i| {
270            let line_start = body[..i]
271                .strip_suffix("# ")
272                .is_some_and(|b| b.is_empty() || b.ends_with('\n'));
273            line_start
274                && body[i + CHAT_CONSULT_HEADING.len()..].starts_with("\n\nThe operator passed you")
275        })
276        .collect();
277    if heads.is_empty() {
278        return body.trim().to_owned();
279    }
280    // (start, end) of each generated block, the heading's own "# " included.
281    // `chat_consult` defuses every marker quoted in a question's text, so the
282    // first end phrase after a heading is the block's own. A block with no end
283    // before the next heading runs up to that heading (or the text's end).
284    // Turns stored before the markers were defused may still quote an end
285    // phrase; their tail is not recognised as part of the block.
286    let mut blocks: Vec<(usize, usize)> = Vec::new();
287    for &h in &heads {
288        let start = if body[..h].ends_with("# ") { h - 2 } else { h };
289        if blocks.last().is_some_and(|&(_, e)| start < e) {
290            continue; // quoted inside the block above
291        }
292        let rest = &body[h..];
293        let head = rest.split("\n\n## ").next().unwrap_or(rest);
294        let end = if head.contains(crate::prompt::CHAT_CONSULT_DEFUSED) {
295            // Current format: quoted markers are defused, so the first end
296            // phrase is the block's own.
297            rest.find(CHAT_CONSULT_END)
298                .map_or(body.len(), |p| h + p + CHAT_CONSULT_END.len())
299        } else {
300            h + legacy_len(rest)
301        };
302        blocks.push((start, end));
303    }
304    let mut from = 0;
305    if let Some(id) = after_block_of {
306        if let Some(&(_, end)) = blocks.iter().rfind(|&&(s, e)| body[s..e].contains(id)) {
307            from = end;
308        } else if let Some(&(_, end)) = blocks.last() {
309            from = end;
310        }
311    }
312    let mut parts = Vec::new();
313    let mut at = from;
314    for &(s, e) in &blocks {
315        if s >= at {
316            parts.push(body[at..s].trim());
317        }
318        at = at.max(e);
319    }
320    parts.push(body[at..].trim());
321    parts
322        .into_iter()
323        .filter(|p| !p.is_empty())
324        .collect::<Vec<_>>()
325        .join("\n\n")
326}
327
328/// What [`start_turn`] did.
329#[derive(Debug, PartialEq, Eq)]
330pub enum Started {
331    /// The question was handed over and the chat's turn ran to its end.
332    Answered,
333    /// Another process holds the talk's turn lease; nothing was changed.
334    Busy,
335    /// The question had already been handed to the chat; nothing was changed.
336    Nothing,
337}
338
339/// [`begin`], then run the chat's turn here, under the talk's own lease.
340///
341/// The lease is taken *before* anything is written, so a talk somebody else is
342/// running leaves no consult record and no draft behind (`Started::Busy`) and
343/// the caller can simply try again later. A turn that fails after `begin`
344/// leaves the record and the draft, so the chat can resume it.
345pub async fn start_turn(
346    questions: &Questions,
347    talks: &Talks,
348    q: &Question,
349    talk: &Talk,
350    cfg: &Config,
351) -> Result<Started> {
352    if q.consult.is_some() {
353        return Ok(Started::Nothing);
354    }
355    let Some(mut lease) = talks.claim_turn(&talk.id)? else {
356        return Ok(Started::Busy);
357    };
358    if !begin(questions, talks, q, talk)? {
359        return Ok(Started::Nothing);
360    }
361    let mut talk = talks.get(&talk.id)?;
362    // Other starters that found the lease held queued drafts and left them to
363    // us, so drain until nothing is left (as the web's drain loop does).
364    // A failed turn is kept and reported at the end, not returned at once:
365    // drafts accepted meanwhile are still owed an answer.
366    let mut failed = None;
367    loop {
368        while let Some(text) = talk::drain(&mut talk, talks)? {
369            if let Err(e) = talk::respond(&lease, &mut talk, talks, cfg, &text).await {
370                failed.get_or_insert(e);
371            }
372            if !lease.beat()? {
373                bail!("the turn lease for chat {} was lost", talk.short());
374            }
375        }
376        // A draft queued after the last drain, before the lease is gone, was
377        // left to us by a starter that found the lease held. Release first,
378        // then look again: whoever queues later either sees no lease (and
379        // starts its own turn) or is seen by this re-check.
380        drop(lease);
381        talk = talks.get(&talk.id)?;
382        let owed = talk.status.open()
383            && (!talk.pending.is_empty() || !talk.pending_attachments.is_empty());
384        if !owed {
385            break;
386        }
387        // Somebody else took the lease meanwhile: they drain it.
388        let Some(again) = talks.claim_turn(&talk.id)? else {
389            break;
390        };
391        lease = again;
392    }
393    if let Some(e) = failed {
394        return Err(e);
395    }
396    Ok(Started::Answered)
397}
398
399#[cfg(test)]
400mod tests {
401    use std::collections::BTreeMap;
402    use std::path::PathBuf;
403
404    use super::*;
405    use crate::config::{AgentKind, AgentSpec, Config};
406    use crate::queue::Source;
407    use crate::talk::TalkStatus;
408
409    const RUN: &str = "20260902-000000-beef";
410
411    fn talks() -> (tempfile::TempDir, Talks, Talk) {
412        let tmp = tempfile::tempdir().expect("tempdir");
413        let store = Talks::at(tmp.path().join("talks"));
414        let cfg = Config {
415            agents: vec![AgentSpec {
416                id: "mock".to_owned(),
417                kind: AgentKind::Command,
418                model: None,
419                command: vec!["true".to_owned()],
420                extra_args: Vec::new(),
421                env: BTreeMap::new(),
422                prompt_delivery: None,
423            }],
424            ..Config::default()
425        };
426        let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
427        (tmp, store, talk)
428    }
429
430    fn task(source: Source) -> Task {
431        let mut t = Task::new(
432            "t".to_owned(),
433            "Do it".to_owned(),
434            PathBuf::from("/repo"),
435            source,
436        );
437        t.start(RUN.to_owned());
438        t
439    }
440
441    fn from_chat(talk: &Talk) -> Task {
442        task(Source::Agent {
443            run: talk.id.clone(),
444            node: CHAT_NODE.to_owned(),
445        })
446    }
447
448    fn question(node: &str) -> Question {
449        Question::new(
450            RUN.to_owned(),
451            node.to_owned(),
452            "impl-A".to_owned(),
453            "Which backend?".to_owned(),
454            "SQLite is simpler.".to_owned(),
455            vec!["SQLite".to_owned(), "Redis".to_owned()],
456        )
457    }
458
459    #[tokio::test]
460    async fn start_turn_is_busy_and_writes_nothing_while_the_lease_is_held() {
461        let (tmp, store, talk) = talks();
462        let questions = Questions::at(tmp.path().join("questions"));
463        let mut q = question("implement");
464        questions.put(&mut q).expect("put");
465        let held = store.claim_turn(&talk.id).expect("claim").expect("first");
466
467        let got = start_turn(&questions, &store, &q, &talk, &Config::default())
468            .await
469            .expect("start");
470        assert_eq!(got, Started::Busy);
471        assert!(questions.get(&q.id).unwrap().consult.is_none());
472        assert!(store.get(&talk.id).unwrap().pending.is_empty());
473        drop(held);
474    }
475
476    /// A `kind = "command"` agent running `script`; `LEASE` names the turn
477    /// lease file so the script can look at it while the turn is in flight.
478    fn scripted(tmp: &std::path::Path, script: &str, lease: &std::path::Path) -> Config {
479        let path = tmp.join("mock-consult-agent.sh");
480        std::fs::write(&path, script).expect("write mock");
481        Config {
482            agents: vec![AgentSpec {
483                id: "mock".to_owned(),
484                kind: AgentKind::Command,
485                model: None,
486                command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
487                extra_args: Vec::new(),
488                env: BTreeMap::from([("LEASE".to_owned(), lease.to_string_lossy().into_owned())]),
489                prompt_delivery: None,
490            }],
491            ..Config::default()
492        }
493    }
494
495    /// Run `start_turn` with `script` as the agent and hand back its result,
496    /// the store and the talk, for the caller to claim the lease afterwards.
497    async fn run_turn(
498        script: &str,
499    ) -> (
500        tempfile::TempDir,
501        Result<Started>,
502        Questions,
503        Question,
504        Talk,
505    ) {
506        crate::run::set_home(crate::run::test_home());
507        let (tmp, store, talk) = talks();
508        let questions = Questions::at(tmp.path().join("questions"));
509        let mut q = question("implement");
510        questions.put(&mut q).expect("put");
511        let cfg = scripted(tmp.path(), script, &store.turn_path(&talk.id));
512        let got = start_turn(&questions, &store, &q, &talk, &cfg).await;
513        (tmp, got, questions, q, talk)
514    }
515
516    #[tokio::test]
517    async fn the_lease_is_held_during_the_turn_and_free_once_it_ends() {
518        // exit 7 if the lease file is missing while the agent runs.
519        let (tmp, got, _questions, _q, talk) =
520            run_turn("#!/bin/sh\ncat >/dev/null\n[ -f \"$LEASE\" ] || exit 7\nprintf 'ok\\n'\n")
521                .await;
522        assert_eq!(got.expect("turn"), Started::Answered);
523        let other = Talks::at(tmp.path().join("talks"));
524        assert!(
525            other.claim_turn(&talk.id).expect("claim").is_some(),
526            "free once the turn ended"
527        );
528    }
529
530    #[tokio::test]
531    async fn the_lease_is_free_again_after_a_failed_turn() {
532        let (tmp, got, questions, q, talk) = run_turn("#!/bin/sh\ncat >/dev/null\nexit 3\n").await;
533        assert!(got.is_err(), "the agent failed");
534        assert!(
535            questions.get(&q.id).unwrap().consult.is_some(),
536            "the turn got as far as running"
537        );
538        let other = Talks::at(tmp.path().join("talks"));
539        assert!(
540            other.claim_turn(&talk.id).expect("claim").is_some(),
541            "free after the failed turn"
542        );
543    }
544
545    #[test]
546    fn origin_talk_is_decided_in_one_table() {
547        let (_tmp, _store, talk) = talks();
548        let mut closed = talk.clone();
549        closed.status = TalkStatus::Closed;
550        let chat = from_chat(&talk);
551        let q = question("implement");
552
553        let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
554        assert_eq!(
555            hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
556            Some(talk.id.clone())
557        );
558        // Not filed from a chat.
559        assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
560        let other = task(Source::Agent {
561            run: talk.id.clone(),
562            node: "implement".to_owned(),
563        });
564        assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
565        // The conversation is gone, or no longer takes turns.
566        assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
567        assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
568        // No task owns the question.
569        assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
570        // Merge approvals can be consulted; release notices cannot.
571        assert!(
572            hit(
573                std::slice::from_ref(&chat),
574                std::slice::from_ref(&talk),
575                &question(crate::land::APPROVAL_NODE)
576            )
577            .is_some()
578        );
579        assert!(
580            hit(
581                std::slice::from_ref(&chat),
582                std::slice::from_ref(&talk),
583                &question(crate::bump::NOTICE_NODE)
584            )
585            .is_none()
586        );
587        // A settled question has nothing left to ask.
588        let mut answered = question("implement");
589        answered
590            .answer(crate::ask::Answer::Choice("Redis".to_owned()))
591            .unwrap();
592        assert!(hit(&[chat], &[talk], &answered).is_none());
593    }
594
595    #[test]
596    fn approval_origin_requires_an_open_chat_task_and_question() {
597        let (_tmp, _store, talk) = talks();
598        let chat = from_chat(&talk);
599        let mut q = question(crate::land::APPROVAL_NODE);
600        assert!(origin_talk(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
601        assert!(
602            origin_talk(
603                &[task(Source::Agent {
604                    run: talk.id.clone(),
605                    node: "implement".into(),
606                })],
607                std::slice::from_ref(&talk),
608                &q
609            )
610            .is_none()
611        );
612        assert!(origin_talk(std::slice::from_ref(&chat), &[], &q).is_none());
613        let mut closed = talk.clone();
614        closed.status = TalkStatus::Closed;
615        assert!(origin_talk(std::slice::from_ref(&chat), &[closed], &q).is_none());
616        q.abandon("expired");
617        assert!(
618            origin_talk(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).is_none()
619        );
620        let mut q = question(crate::land::APPROVAL_NODE);
621        q.answer(crate::ask::Answer::Choice("Redis".into()))
622            .unwrap();
623        assert!(origin_talk(&[chat], &[talk], &q).is_none());
624    }
625
626    #[test]
627    fn a_conductor_question_is_found_through_the_task_id() {
628        let (_tmp, _store, talk) = talks();
629        let chat = from_chat(&talk);
630        let mut q = question(crate::conduct::NODE);
631        q.run = chat.id.clone();
632        assert!(origin_talk(&[chat], &[talk], &q).is_some());
633    }
634
635    fn followup_of(parent: &Task, n: u32) -> Task {
636        let mut t = Task::new(
637            "f".to_owned(),
638            "Fix".to_owned(),
639            PathBuf::from("/repo"),
640            Source::Agent {
641                run: format!("merged-{n}"),
642                node: "followup".to_owned(),
643            },
644        );
645        t.start(format!("run-f{n}"));
646        t.followup = Some(crate::queue::FollowUp {
647            run: format!("merged-{n}"),
648            origin_task: Some(parent.id.clone()),
649            pr: "https://example.invalid/pr/1".to_owned(),
650            findings: Vec::new(),
651            generation: n,
652        });
653        t
654    }
655
656    fn q_for(t: &Task) -> Question {
657        let mut q = question("implement");
658        q.run = t.runs[0].clone();
659        q
660    }
661
662    #[test]
663    fn a_followup_traces_back_to_the_chat() {
664        let (_tmp, _store, talk) = talks();
665        let mut chat = from_chat(&talk);
666        chat.runs = vec!["chat-run".to_owned()];
667        let f1 = followup_of(&chat, 1);
668        let f2 = followup_of(&f1, 2);
669        let ts = [chat, f1.clone(), f2.clone()];
670        let id = Some(talk.id.clone());
671        let tk = std::slice::from_ref(&talk);
672        assert_eq!(origin_talk(&ts, tk, &q_for(&f1)).map(|t| t.id), id);
673        assert_eq!(origin_talk(&ts, tk, &q_for(&f2)).map(|t| t.id), id);
674    }
675
676    #[test]
677    fn a_followup_without_origin_task_is_found_through_runs() {
678        let (_tmp, _store, talk) = talks();
679        let mut chat = from_chat(&talk);
680        chat.runs = vec!["merged-1".to_owned()];
681        let mut f1 = followup_of(&chat, 1);
682        f1.followup.as_mut().unwrap().origin_task = None;
683        let q = q_for(&f1);
684        assert!(origin_talk(&[chat, f1], &[talk], &q).is_some());
685    }
686
687    #[test]
688    fn a_dangling_origin_task_has_no_chat() {
689        let (_tmp, _store, talk) = talks();
690        let mut chat = from_chat(&talk);
691        chat.runs = vec!["chat-run".to_owned()];
692        let mut f1 = followup_of(&chat, 1);
693        f1.followup.as_mut().unwrap().origin_task = Some("gone".to_owned());
694        let q = q_for(&f1);
695        assert!(origin_talk(&[chat, f1], &[talk], &q).is_none());
696    }
697
698    #[test]
699    fn a_followup_cycle_ends_without_a_chat() {
700        let (_tmp, _store, talk) = talks();
701        let mut a = followup_of(&from_chat(&talk), 1);
702        let mut b = followup_of(&a, 2);
703        a.followup.as_mut().unwrap().origin_task = Some(b.id.clone());
704        b.followup.as_mut().unwrap().origin_task = Some(a.id.clone());
705        let q = q_for(&a);
706        assert!(origin_talk(&[a, b], &[talk], &q).is_none());
707    }
708
709    #[test]
710    fn pending_consults_follow_the_store_not_the_turn_text() {
711        let (tmp, store, talk) = talks();
712        let questions = Questions::at(tmp.path().join("questions"));
713        assert!(!pending_consults(&questions, &talk.id), "empty store");
714
715        let mut plain = question("implement");
716        questions.put(&mut plain).unwrap();
717        assert!(!pending_consults(&questions, &talk.id), "no consult record");
718
719        let mut q = question("implement");
720        questions.put(&mut q).unwrap();
721        begin(&questions, &store, &q, &talk).unwrap();
722        assert!(pending_consults(&questions, &talk.id));
723        assert!(!pending_consults(&questions, "other-talk"));
724
725        questions
726            .update(&q.id, |q| {
727                q.abandon("test");
728                Ok(())
729            })
730            .unwrap();
731        assert!(!pending_consults(&questions, &talk.id), "closed question");
732    }
733
734    #[test]
735    fn begin_queues_once_and_leaves_the_question_open() {
736        let (tmp, store, talk) = talks();
737        let questions = Questions::at(tmp.path().join("questions"));
738        let mut q = question("implement");
739        questions.put(&mut q).unwrap();
740
741        assert!(begin(&questions, &store, &q, &talk).unwrap());
742        assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
743
744        let after = questions.get(&q.id).unwrap();
745        assert!(after.status.open());
746        assert!(after.thread.is_empty());
747        assert_eq!(after.choices, q.choices, "choices are not touched");
748        assert_eq!(
749            after.consult.as_ref().map(|c| c.talk.as_str()),
750            Some(talk.id.as_str())
751        );
752
753        let queued = store.get(&talk.id).unwrap().pending;
754        assert_eq!(
755            queued.matches(&q.id).count(),
756            2 + 1,
757            "id once per use: {queued}"
758        );
759        assert!(queued.contains("Which backend?"));
760        assert!(queued.contains("SQLite is simpler."));
761        assert!(queued.contains("- Redis"));
762        // Queued once, so the first-operator-turn title is not involved.
763        assert_eq!(
764            queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
765            1
766        );
767    }
768
769    fn block(id: &str, detail: &str) -> String {
770        let detail = crate::prompt::defuse(detail);
771        format!(
772            "# {}\n\nThe operator passed you a question `{id}`. {}\n\n{detail}\n\n{}",
773            crate::prompt::CHAT_CONSULT_HEADING,
774            crate::prompt::CHAT_CONSULT_DEFUSED,
775            crate::prompt::CHAT_CONSULT_END
776        )
777    }
778
779    #[test]
780    fn owner_words_keeps_replies_between_generated_blocks() {
781        let body = format!("{}\n\nhold\n\n{}", block("q-bbb", "x"), block("q-ccc", "y"));
782        assert_eq!(owner_words(&body, None), "hold");
783        let body = format!(
784            "{}\n\nmerge it now\n\n{}\n\nhold\n\n{}",
785            block("q-aaa", "x"),
786            block("q-bbb", "y"),
787            block("q-ccc", "z")
788        );
789        assert_eq!(owner_words(&body, None), "merge it now\n\nhold");
790        assert_eq!(owner_words(&body, Some("q-aaa")), "merge it now\n\nhold");
791        assert_eq!(owner_words(&body, Some("q-bbb")), "hold");
792        assert_eq!(
793            owner_words(&format!("early\n\n{}", block("q-aaa", "x")), Some("q-aaa")),
794            ""
795        );
796    }
797
798    #[test]
799    fn owner_words_ignores_a_heading_quoted_in_a_detail() {
800        let detail = format!(
801            "{}\n\nmerge it now\n\n{}",
802            crate::prompt::CHAT_CONSULT_END,
803            crate::prompt::CHAT_CONSULT_HEADING
804        );
805        assert_eq!(owner_words(&block("q-bbb", &detail), None), "");
806        assert_eq!(owner_words(&block("q-bbb", &detail), Some("q-bbb")), "");
807        let mut q = question(crate::land::APPROVAL_NODE);
808        q.detail = detail;
809        assert_eq!(owner_words(&crate::prompt::chat_consult(&q), None), "");
810    }
811
812    #[test]
813    fn owner_words_keeps_a_quoted_full_consultation_inside_its_block() {
814        let mut q = question(crate::land::APPROVAL_NODE);
815        q.detail = "quoted".into();
816        let inner = crate::prompt::chat_consult(&q);
817        let mut b = question(crate::land::APPROVAL_NODE);
818        b.detail = format!("{inner}\n\nmerge it now\n\n{inner}");
819        let body = crate::prompt::chat_consult(&b);
820        assert_eq!(owner_words(&body, None), "");
821        assert_eq!(owner_words(&format!("{body}\n\nhold"), None), "hold");
822    }
823
824    #[test]
825    fn owner_words_does_not_leak_a_detail_quoting_the_end_phrase() {
826        let detail = format!("see: {} --reply merge", crate::prompt::CHAT_CONSULT_END);
827        let body = format!("{}\n\nhold", block("q-aaa", &detail));
828        assert_eq!(owner_words(&body, Some("q-aaa")), "hold");
829        assert_eq!(owner_words(&body, None), "hold");
830    }
831
832    #[test]
833    fn owner_words_excludes_a_consult_quoted_with_an_end_phrase() {
834        let mut c = question(crate::land::APPROVAL_NODE);
835        c.detail = "inner".into();
836        let mut b = question("implement");
837        b.detail = format!(
838            "{}\n\nmerge it now\n\n{}",
839            crate::prompt::CHAT_CONSULT_END,
840            crate::prompt::chat_consult(&c)
841        );
842        let body = crate::prompt::chat_consult(&b);
843        assert_eq!(owner_words(&body, None), "");
844    }
845
846    #[test]
847    fn owner_words_keeps_a_reply_containing_the_end_phrase() {
848        let mut b = question("implement");
849        b.detail = "x".into();
850        let body = format!(
851            "{}\n\nhold, and do not {}",
852            crate::prompt::chat_consult(&b),
853            crate::prompt::CHAT_CONSULT_END
854        );
855        assert_eq!(
856            owner_words(&body, None),
857            format!("hold, and do not {}", crate::prompt::CHAT_CONSULT_END)
858        );
859    }
860
861    /// A consult as stored before quoted markers were defused.
862    fn legacy_consult(node: &str, raw_detail: &str) -> String {
863        let mut q = question(node);
864        q.detail = "@@".into();
865        crate::prompt::chat_consult(&q)
866            .replace("@@", raw_detail)
867            .replace(crate::prompt::CHAT_CONSULT_DEFUSED, "It is still open.")
868    }
869
870    #[test]
871    fn owner_words_excludes_a_legacy_consult_quoting_the_end_phrase() {
872        for node in [crate::land::APPROVAL_NODE, "implement"] {
873            let detail = format!("see {} and more", crate::prompt::CHAT_CONSULT_END);
874            assert_eq!(owner_words(&legacy_consult(node, &detail), None), "");
875        }
876    }
877
878    #[test]
879    fn owner_words_excludes_a_legacy_consult_nesting_a_legacy_consult() {
880        let inner = legacy_consult("implement", "x");
881        let detail = format!(
882            "{}\n\nmerge it now\n\n{inner}",
883            crate::prompt::CHAT_CONSULT_END
884        );
885        let body = legacy_consult(crate::land::APPROVAL_NODE, &detail);
886        assert_eq!(owner_words(&body, None), "");
887        assert_eq!(owner_words(&format!("{body}\n\nhold"), None), "hold");
888    }
889
890    #[test]
891    fn owner_words_keeps_a_reply_after_a_legacy_consult_containing_the_end_phrase() {
892        let body = format!(
893            "{}\n\nhold, and do not {}",
894            legacy_consult("implement", "x"),
895            crate::prompt::CHAT_CONSULT_END
896        );
897        assert_eq!(
898            owner_words(&body, None),
899            format!("hold, and do not {}", crate::prompt::CHAT_CONSULT_END)
900        );
901    }
902
903    #[test]
904    fn chat_consult_has_each_marker_once() {
905        use crate::prompt::{CHAT_CONSULT_END as E, CHAT_CONSULT_HEADING as H};
906        for node in ["implement", crate::land::APPROVAL_NODE] {
907            let mut q = question(node);
908            let quoted = format!("{H} {E}");
909            q.summary = quoted.clone();
910            q.detail = quoted.clone();
911            q.choices = vec![quoted.clone(), "b".into()];
912            let s = crate::prompt::chat_consult(&q);
913            assert_eq!(s.matches(H).count(), 1, "{s}");
914            assert_eq!(s.matches(E).count(), 1, "{s}");
915        }
916    }
917
918    #[test]
919    fn approval_consult_waits_for_latest_owner_confirmation() {
920        let (tmp, store, mut talk) = talks();
921        let questions = Questions::at(tmp.path().join("questions"));
922        let mut q = question(crate::land::APPROVAL_NODE);
923        q.choices = vec![crate::land::APPROVE.into(), crate::land::HOLD.into()];
924        questions.put(&mut q).unwrap();
925        assert!(begin(&questions, &store, &q, &talk).unwrap());
926        assert!(!begin(&questions, &store, &q, &talk).unwrap());
927        q = questions.get(&q.id).unwrap();
928        assert!(q.status.open());
929        assert!(q.answer.is_none());
930        assert!(q.thread.is_empty());
931        let queued = store.get(&talk.id).unwrap().pending;
932        assert!(queued.contains("Never answer it yourself"));
933        assert!(queued.contains("Silence holds"));
934        assert!(queued.contains("--reply merge --quote"));
935        assert!(!queued.contains("answer it yourself with"));
936
937        let owner_turn = |body: &str| talk::Turn {
938            who: talk::Who::Operator,
939            body: body.into(),
940            at: Timestamp::now(),
941            attachments: Vec::new(),
942            usage: None,
943        };
944        // An approval the owner gave before the hand-over is not reusable.
945        let mut old = talk.clone();
946        old.turns
947            .insert(0, owner_turn("Please merge PR 12 after review"));
948        store.put(&mut old).unwrap();
949        talk.turns = old.turns.clone();
950        // The drained consult draft is stored as an operator turn; it must not
951        // count as the owner's words even though it contains "merge".
952        talk.turns.push(owner_turn(&queued));
953        store.put(&mut talk).unwrap();
954        assert!(validate_answer(&q, &store, &talk.id, CHAT_NODE, "merge", Some("merge")).is_err());
955        assert!(
956            validate_answer(
957                &q,
958                &store,
959                &talk.id,
960                CHAT_NODE,
961                "merge",
962                Some("Please merge PR 12")
963            )
964            .is_err()
965        );
966        // An owner reply coalesced onto the same draft still counts.
967        let n = talk.turns.len();
968        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now");
969        store.put(&mut talk).unwrap();
970        assert!(
971            validate_answer(
972                &q,
973                &store,
974                &talk.id,
975                CHAT_NODE,
976                "merge",
977                Some("merge it now")
978            )
979            .is_ok()
980        );
981        // A later retraction in the same coalesced turn wins.
982        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now\n\nhold");
983        store.put(&mut talk).unwrap();
984        assert!(
985            validate_answer(
986                &q,
987                &store,
988                &talk.id,
989                CHAT_NODE,
990                "merge",
991                Some("merge it now")
992            )
993            .is_err()
994        );
995        assert!(validate_answer(&q, &store, &talk.id, CHAT_NODE, "hold", Some("hold")).is_ok());
996        // So does a reply still waiting in the draft.
997        talk.turns[n - 1].body = format!("{queued}\n\nmerge it now");
998        talk.pending = "hold".into();
999        store.put(&mut talk).unwrap();
1000        assert!(
1001            validate_answer(
1002                &q,
1003                &store,
1004                &talk.id,
1005                CHAT_NODE,
1006                "merge",
1007                Some("merge it now")
1008            )
1009            .is_err()
1010        );
1011        // Another question's hand-over queued after the retraction hides nothing.
1012        talk.pending = format!("hold\n\n{queued}");
1013        store.put(&mut talk).unwrap();
1014        assert!(
1015            validate_answer(
1016                &q,
1017                &store,
1018                &talk.id,
1019                CHAT_NODE,
1020                "merge",
1021                Some("merge it now")
1022            )
1023            .is_err()
1024        );
1025        talk.pending.clear();
1026        talk.turns[n - 1].body = queued.clone();
1027        store.put(&mut talk).unwrap();
1028        talk.turns
1029            .push(owner_turn("Merge this pull request please"));
1030        store.put(&mut talk).unwrap();
1031        let check = |q: &Question, run: &str, reply: &str, quote: Option<&str>| {
1032            validate_answer(q, &store, run, CHAT_NODE, reply, quote)
1033        };
1034        assert!(check(&q, "other-talk", "merge", Some("Merge this")).is_err());
1035        assert!(check(&q, &talk.id, "merge", None).is_err());
1036        assert!(check(&q, &talk.id, "merge", Some("never said")).is_err());
1037        assert!(
1038            check(
1039                &q,
1040                &talk.id,
1041                "merge",
1042                Some("Merge this pull request please")
1043            )
1044            .is_ok()
1045        );
1046        assert!(check(&q, &talk.id, "hold", Some("hold")).is_err());
1047        talk.turns
1048            .push(owner_turn("Wait, explain the checks first"));
1049        store.put(&mut talk).unwrap();
1050        assert!(
1051            check(
1052                &q,
1053                &talk.id,
1054                "merge",
1055                Some("Merge this pull request please")
1056            )
1057            .is_err()
1058        );
1059        talk.turns.push(owner_turn("Please hold"));
1060        store.put(&mut talk).unwrap();
1061        assert!(check(&q, &talk.id, "hold", Some("hold")).is_err());
1062        talk.turns.push(owner_turn("hold"));
1063        store.put(&mut talk).unwrap();
1064        assert!(check(&q, &talk.id, "hold", Some("hold")).is_ok());
1065        talk.turns
1066            .push(owner_turn("Merge this pull request please"));
1067        store.put(&mut talk).unwrap();
1068        q.abandon("approval expired or head changed");
1069        assert!(
1070            check(
1071                &q,
1072                &talk.id,
1073                "merge",
1074                Some("Merge this pull request please")
1075            )
1076            .is_err()
1077        );
1078        assert!(
1079            q.answer(crate::ask::Answer::Choice("merge".into()))
1080                .is_err()
1081        );
1082    }
1083
1084    #[test]
1085    fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
1086        let (tmp, store, talk) = talks();
1087        let questions = Questions::at(tmp.path().join("questions"));
1088        let mut q = question("implement");
1089        questions.put(&mut q).unwrap();
1090        std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
1091        assert!(begin(&questions, &store, &q, &talk).is_err());
1092        assert!(questions.get(&q.id).unwrap().consult.is_none());
1093    }
1094
1095    #[test]
1096    fn an_older_question_file_reads_without_a_consult() {
1097        let mut q = question("implement");
1098        q.schema = 5;
1099        let mut v = serde_json::to_value(&q).unwrap();
1100        v.as_object_mut().unwrap().remove("consult");
1101        let back: Question = serde_json::from_value(v).unwrap();
1102        assert!(back.consult.is_none());
1103    }
1104}