Skip to main content

magi/
consult.rs

1//! Handing an open question to the chat conversation its task came from.
2//!
3//! A task filed by the standing chat (`Source::Agent` with
4//! [`crate::queue::CHAT_NODE`], whose `run` is the conversation's id) can raise
5//! a question the owner would rather talk through where the task started. The
6//! owner passes it on; the chat agent either answers it with `magi answer` or
7//! puts the decision points to the owner in the conversation.
8//!
9//! Three rules, each with a reason:
10//!
11//! - **One place decides.** [`origin_talk`] is the only judge of "does this
12//!   question have a chat to ask"; the web view and both entry points use it.
13//! - **Not a choice.** The hand-over never goes into `Question::choices` and
14//!   never touches `Question::answer`: the question stays `Open`, its
15//!   validation and the choice actions stay as they were.
16//! - **No new seat, no new waiter.** The text is queued as a draft of the
17//!   existing talk, so the chat's own turn machinery (its session, its turn
18//!   gate) runs it exactly as if the owner had typed it.
19
20use anyhow::{Context, Result, bail};
21use jiff::Timestamp;
22
23use crate::ask::{ChatConsult, Question, Questions};
24use crate::config::Config;
25use crate::queue::{CHAT_NODE, Source, Task};
26use crate::talk::{self, Talk, Talks};
27
28/// The conversation `q` may be handed to, if any.
29///
30/// Requires an open question whose task was filed from a chat that is still
31/// open. A merge approval and a release notice are left out: their answer is
32/// gated by the deputy's merge-intent rules and must not be settled by a side
33/// door.
34pub fn origin_talk(tasks: &[Task], talks: &[Talk], q: &Question) -> Option<Talk> {
35    if !q.status.open()
36        || q.node == crate::land::APPROVAL_NODE
37        || q.node == crate::bump::NOTICE_NODE
38    {
39        return None;
40    }
41    let task = crate::daemon::task_of_question(tasks, q)?;
42    let Source::Agent { run, .. } = &chat_origin(tasks, task)?.source else {
43        return None;
44    };
45    talks
46        .iter()
47        .find(|t| &t.id == run && t.status.open())
48        .cloned()
49}
50
51/// Walk the provenance of `start` to the task a chat filed.
52///
53/// A follow-up goes to the task its merged run served: `FollowUp::origin_task`
54/// when recorded (a missing one ends the walk - guessing another task could
55/// hand the question to an unrelated chat), else the task whose `runs` hold
56/// `FollowUp::run`. At most `MAX_FOLLOWUP_GENERATION + 1` tasks are looked at,
57/// each once, so a cycle ends in `None`.
58fn chat_origin<'a>(tasks: &'a [Task], start: &'a Task) -> Option<&'a Task> {
59    let mut seen = std::collections::HashSet::new();
60    let mut cur = start;
61    for _ in 0..=crate::followup::MAX_FOLLOWUP_GENERATION {
62        if !seen.insert(cur.id.as_str()) {
63            return None;
64        }
65        if matches!(&cur.source, Source::Agent { node, .. } if node == CHAT_NODE) {
66            return Some(cur);
67        }
68        let f = cur.followup.as_ref()?;
69        cur = match &f.origin_task {
70            Some(id) => tasks.iter().find(|t| &t.id == id)?,
71            None => tasks.iter().find(|t| t.runs.contains(&f.run))?,
72        };
73    }
74    None
75}
76
77/// Is any question handed to talk `talk_id` still open?
78///
79/// Decided from the question store, not from the newest turn's wording: the
80/// owner's decision usually arrives in a later turn that carries no hand-over
81/// heading, and `magi answer` still has to be able to write then. Answered or
82/// abandoned questions stop counting, so the turn goes back to read-only.
83pub fn pending_consults(questions: &Questions, talk_id: &str) -> bool {
84    questions
85        .list()
86        .iter()
87        .any(|q| q.status.open() && q.consult.as_ref().is_some_and(|c| c.talk == talk_id))
88}
89
90/// Record the hand-over and queue the question as the chat's next message.
91///
92/// Returns `false` (and changes nothing) when the question was already handed
93/// over. The question is not otherwise touched: still `Open`, no thread turn.
94/// When queueing fails the record is withdrawn, so the owner can try again.
95pub fn begin(questions: &Questions, talks: &Talks, q: &Question, talk: &Talk) -> Result<bool> {
96    let (q, fresh) = questions.update(&q.id, |r| {
97        if !r.status.open() {
98            bail!("question {} is already {}", r.short(), r.status.as_str());
99        }
100        if r.consult.is_some() {
101            return Ok(false);
102        }
103        r.consult = Some(ChatConsult {
104            talk: talk.id.clone(),
105            at: Timestamp::now(),
106        });
107        Ok(true)
108    })?;
109    if !fresh {
110        return Ok(false);
111    }
112    let mut talk = talk.clone();
113    if let Err(e) = talk::queue(
114        &mut talk,
115        talks,
116        &crate::prompt::chat_consult(&q),
117        Vec::new(),
118    ) {
119        let _ = questions.update(&q.id, |r| {
120            r.consult = None;
121            Ok(())
122        });
123        return Err(e).context("queue the question into the chat");
124    }
125    Ok(true)
126}
127
128/// What [`start_turn`] did.
129#[derive(Debug, PartialEq, Eq)]
130pub enum Started {
131    /// The question was handed over and the chat's turn ran to its end.
132    Answered,
133    /// Another process holds the talk's turn lease; nothing was changed.
134    Busy,
135    /// The question had already been handed to the chat; nothing was changed.
136    Nothing,
137}
138
139/// [`begin`], then run the chat's turn here, under the talk's own lease.
140///
141/// The lease is taken *before* anything is written, so a talk somebody else is
142/// running leaves no consult record and no draft behind (`Started::Busy`) and
143/// the caller can simply try again later. A turn that fails after `begin`
144/// leaves the record and the draft, so the chat can resume it.
145pub async fn start_turn(
146    questions: &Questions,
147    talks: &Talks,
148    q: &Question,
149    talk: &Talk,
150    cfg: &Config,
151) -> Result<Started> {
152    if q.consult.is_some() {
153        return Ok(Started::Nothing);
154    }
155    let Some(mut lease) = talks.claim_turn(&talk.id)? else {
156        return Ok(Started::Busy);
157    };
158    if !begin(questions, talks, q, talk)? {
159        return Ok(Started::Nothing);
160    }
161    let mut talk = talks.get(&talk.id)?;
162    // Other starters that found the lease held queued drafts and left them to
163    // us, so drain until nothing is left (as the web's drain loop does).
164    // A failed turn is kept and reported at the end, not returned at once:
165    // drafts accepted meanwhile are still owed an answer.
166    let mut failed = None;
167    loop {
168        while let Some(text) = talk::drain(&mut talk, talks)? {
169            if let Err(e) = talk::respond(&lease, &mut talk, talks, cfg, &text).await {
170                failed.get_or_insert(e);
171            }
172            if !lease.beat()? {
173                bail!("the turn lease for chat {} was lost", talk.short());
174            }
175        }
176        // A draft queued after the last drain, before the lease is gone, was
177        // left to us by a starter that found the lease held. Release first,
178        // then look again: whoever queues later either sees no lease (and
179        // starts its own turn) or is seen by this re-check.
180        drop(lease);
181        talk = talks.get(&talk.id)?;
182        let owed = talk.status.open()
183            && (!talk.pending.is_empty() || !talk.pending_attachments.is_empty());
184        if !owed {
185            break;
186        }
187        // Somebody else took the lease meanwhile: they drain it.
188        let Some(again) = talks.claim_turn(&talk.id)? else {
189            break;
190        };
191        lease = again;
192    }
193    if let Some(e) = failed {
194        return Err(e);
195    }
196    Ok(Started::Answered)
197}
198
199#[cfg(test)]
200mod tests {
201    use std::collections::BTreeMap;
202    use std::path::PathBuf;
203
204    use super::*;
205    use crate::config::{AgentKind, AgentSpec, Config};
206    use crate::queue::Source;
207    use crate::talk::TalkStatus;
208
209    const RUN: &str = "20260902-000000-beef";
210
211    fn talks() -> (tempfile::TempDir, Talks, Talk) {
212        let tmp = tempfile::tempdir().expect("tempdir");
213        let store = Talks::at(tmp.path().join("talks"));
214        let cfg = Config {
215            agents: vec![AgentSpec {
216                id: "mock".to_owned(),
217                kind: AgentKind::Command,
218                model: None,
219                command: vec!["true".to_owned()],
220                extra_args: Vec::new(),
221                env: BTreeMap::new(),
222                prompt_delivery: None,
223            }],
224            ..Config::default()
225        };
226        let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
227        (tmp, store, talk)
228    }
229
230    fn task(source: Source) -> Task {
231        let mut t = Task::new(
232            "t".to_owned(),
233            "Do it".to_owned(),
234            PathBuf::from("/repo"),
235            source,
236        );
237        t.start(RUN.to_owned());
238        t
239    }
240
241    fn from_chat(talk: &Talk) -> Task {
242        task(Source::Agent {
243            run: talk.id.clone(),
244            node: CHAT_NODE.to_owned(),
245        })
246    }
247
248    fn question(node: &str) -> Question {
249        Question::new(
250            RUN.to_owned(),
251            node.to_owned(),
252            "impl-A".to_owned(),
253            "Which backend?".to_owned(),
254            "SQLite is simpler.".to_owned(),
255            vec!["SQLite".to_owned(), "Redis".to_owned()],
256        )
257    }
258
259    #[tokio::test]
260    async fn start_turn_is_busy_and_writes_nothing_while_the_lease_is_held() {
261        let (tmp, store, talk) = talks();
262        let questions = Questions::at(tmp.path().join("questions"));
263        let mut q = question("implement");
264        questions.put(&mut q).expect("put");
265        let held = store.claim_turn(&talk.id).expect("claim").expect("first");
266
267        let got = start_turn(&questions, &store, &q, &talk, &Config::default())
268            .await
269            .expect("start");
270        assert_eq!(got, Started::Busy);
271        assert!(questions.get(&q.id).unwrap().consult.is_none());
272        assert!(store.get(&talk.id).unwrap().pending.is_empty());
273        drop(held);
274    }
275
276    #[test]
277    fn origin_talk_is_decided_in_one_table() {
278        let (_tmp, _store, talk) = talks();
279        let mut closed = talk.clone();
280        closed.status = TalkStatus::Closed;
281        let chat = from_chat(&talk);
282        let q = question("implement");
283
284        let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
285        assert_eq!(
286            hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
287            Some(talk.id.clone())
288        );
289        // Not filed from a chat.
290        assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
291        let other = task(Source::Agent {
292            run: talk.id.clone(),
293            node: "implement".to_owned(),
294        });
295        assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
296        // The conversation is gone, or no longer takes turns.
297        assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
298        assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
299        // No task owns the question.
300        assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
301        // Approvals and release notices are never handed over.
302        for node in [crate::land::APPROVAL_NODE, crate::bump::NOTICE_NODE] {
303            assert!(
304                hit(
305                    std::slice::from_ref(&chat),
306                    std::slice::from_ref(&talk),
307                    &question(node)
308                )
309                .is_none()
310            );
311        }
312        // A settled question has nothing left to ask.
313        let mut answered = question("implement");
314        answered
315            .answer(crate::ask::Answer::Choice("Redis".to_owned()))
316            .unwrap();
317        assert!(hit(&[chat], &[talk], &answered).is_none());
318    }
319
320    #[test]
321    fn a_conductor_question_is_found_through_the_task_id() {
322        let (_tmp, _store, talk) = talks();
323        let chat = from_chat(&talk);
324        let mut q = question(crate::conduct::NODE);
325        q.run = chat.id.clone();
326        assert!(origin_talk(&[chat], &[talk], &q).is_some());
327    }
328
329    fn followup_of(parent: &Task, n: u32) -> Task {
330        let mut t = Task::new(
331            "f".to_owned(),
332            "Fix".to_owned(),
333            PathBuf::from("/repo"),
334            Source::Agent {
335                run: format!("merged-{n}"),
336                node: "followup".to_owned(),
337            },
338        );
339        t.start(format!("run-f{n}"));
340        t.followup = Some(crate::queue::FollowUp {
341            run: format!("merged-{n}"),
342            origin_task: Some(parent.id.clone()),
343            pr: "https://example.invalid/pr/1".to_owned(),
344            findings: Vec::new(),
345            generation: n,
346        });
347        t
348    }
349
350    fn q_for(t: &Task) -> Question {
351        let mut q = question("implement");
352        q.run = t.runs[0].clone();
353        q
354    }
355
356    #[test]
357    fn a_followup_traces_back_to_the_chat() {
358        let (_tmp, _store, talk) = talks();
359        let mut chat = from_chat(&talk);
360        chat.runs = vec!["chat-run".to_owned()];
361        let f1 = followup_of(&chat, 1);
362        let f2 = followup_of(&f1, 2);
363        let ts = [chat, f1.clone(), f2.clone()];
364        let id = Some(talk.id.clone());
365        let tk = std::slice::from_ref(&talk);
366        assert_eq!(origin_talk(&ts, tk, &q_for(&f1)).map(|t| t.id), id);
367        assert_eq!(origin_talk(&ts, tk, &q_for(&f2)).map(|t| t.id), id);
368    }
369
370    #[test]
371    fn a_followup_without_origin_task_is_found_through_runs() {
372        let (_tmp, _store, talk) = talks();
373        let mut chat = from_chat(&talk);
374        chat.runs = vec!["merged-1".to_owned()];
375        let mut f1 = followup_of(&chat, 1);
376        f1.followup.as_mut().unwrap().origin_task = None;
377        let q = q_for(&f1);
378        assert!(origin_talk(&[chat, f1], &[talk], &q).is_some());
379    }
380
381    #[test]
382    fn a_dangling_origin_task_has_no_chat() {
383        let (_tmp, _store, talk) = talks();
384        let mut chat = from_chat(&talk);
385        chat.runs = vec!["chat-run".to_owned()];
386        let mut f1 = followup_of(&chat, 1);
387        f1.followup.as_mut().unwrap().origin_task = Some("gone".to_owned());
388        let q = q_for(&f1);
389        assert!(origin_talk(&[chat, f1], &[talk], &q).is_none());
390    }
391
392    #[test]
393    fn a_followup_cycle_ends_without_a_chat() {
394        let (_tmp, _store, talk) = talks();
395        let mut a = followup_of(&from_chat(&talk), 1);
396        let mut b = followup_of(&a, 2);
397        a.followup.as_mut().unwrap().origin_task = Some(b.id.clone());
398        b.followup.as_mut().unwrap().origin_task = Some(a.id.clone());
399        let q = q_for(&a);
400        assert!(origin_talk(&[a, b], &[talk], &q).is_none());
401    }
402
403    #[test]
404    fn pending_consults_follow_the_store_not_the_turn_text() {
405        let (tmp, store, talk) = talks();
406        let questions = Questions::at(tmp.path().join("questions"));
407        assert!(!pending_consults(&questions, &talk.id), "empty store");
408
409        let mut plain = question("implement");
410        questions.put(&mut plain).unwrap();
411        assert!(!pending_consults(&questions, &talk.id), "no consult record");
412
413        let mut q = question("implement");
414        questions.put(&mut q).unwrap();
415        begin(&questions, &store, &q, &talk).unwrap();
416        assert!(pending_consults(&questions, &talk.id));
417        assert!(!pending_consults(&questions, "other-talk"));
418
419        questions
420            .update(&q.id, |q| {
421                q.abandon("test");
422                Ok(())
423            })
424            .unwrap();
425        assert!(!pending_consults(&questions, &talk.id), "closed question");
426    }
427
428    #[test]
429    fn begin_queues_once_and_leaves_the_question_open() {
430        let (tmp, store, talk) = talks();
431        let questions = Questions::at(tmp.path().join("questions"));
432        let mut q = question("implement");
433        questions.put(&mut q).unwrap();
434
435        assert!(begin(&questions, &store, &q, &talk).unwrap());
436        assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
437
438        let after = questions.get(&q.id).unwrap();
439        assert!(after.status.open());
440        assert!(after.thread.is_empty());
441        assert_eq!(after.choices, q.choices, "choices are not touched");
442        assert_eq!(
443            after.consult.as_ref().map(|c| c.talk.as_str()),
444            Some(talk.id.as_str())
445        );
446
447        let queued = store.get(&talk.id).unwrap().pending;
448        assert_eq!(
449            queued.matches(&q.id).count(),
450            2 + 1,
451            "id once per use: {queued}"
452        );
453        assert!(queued.contains("Which backend?"));
454        assert!(queued.contains("SQLite is simpler."));
455        assert!(queued.contains("- Redis"));
456        // Queued once, so the first-operator-turn title is not involved.
457        assert_eq!(
458            queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
459            1
460        );
461    }
462
463    #[test]
464    fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
465        let (tmp, store, talk) = talks();
466        let questions = Questions::at(tmp.path().join("questions"));
467        let mut q = question("implement");
468        questions.put(&mut q).unwrap();
469        std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
470        assert!(begin(&questions, &store, &q, &talk).is_err());
471        assert!(questions.get(&q.id).unwrap().consult.is_none());
472    }
473
474    #[test]
475    fn an_older_question_file_reads_without_a_consult() {
476        let mut q = question("implement");
477        q.schema = 5;
478        let mut v = serde_json::to_value(&q).unwrap();
479        v.as_object_mut().unwrap().remove("consult");
480        let back: Question = serde_json::from_value(v).unwrap();
481        assert!(back.consult.is_none());
482    }
483}