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