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