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