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