Skip to main content

magi/
waiter.rs

1//! The waiter: something that keeps waiting on a question after the process
2//! that asked it is gone.
3//!
4//! `magi ask` blocks in the agent's own shell tool, and that tool kills a wait
5//! that runs long (see `ask::WAIT_SLICE`). An agent can also background the
6//! call and simply finish. Either way the question outlives the only process
7//! that would have read the owner's answer, and the owner's reply lands in a
8//! file nobody is polling - the phone then says "waiting for the agent" while
9//! there is no agent to wait for.
10//!
11//! This module is the guarantee that something is always on the other end. It
12//! runs inside `magi serve`, on its own cadence rather than the queue's, and
13//! for every question a `magi ask` filed it decides one of three things:
14//!
15//! - **the asker is still there** (its [`Lease`] is fresh): do nothing. A
16//!   second agent must never be started while the first is blocked;
17//! - **the asker is gone and the owner has said or answered something the agent
18//!   has not read**: resume the *same seat's* CLI session with the question,
19//!   the thread so far and the owner's words, so the agent that holds the
20//!   context carries on;
21//! - **nobody answered before `answer_timeout`**: abandon it, in the same words
22//!   the asker would have used.
23//!
24//! # Shape
25//!
26//! Same split as [`crate::queue`] and [`crate::ask`]: [`decide`] is pure - a
27//! question, a lease, a clock and a bool in, an [`Action`] out - and every
28//! filesystem or process effect is in [`Waiter::tick`]. State lives on disk (the
29//! question record and its lease), so a restarted daemon picks up exactly where
30//! the last one stopped, including an answer nobody has read yet.
31//!
32//! # What it will not do
33//!
34//! It never starts a fresh consultant. If the seat's session cannot be resumed
35//! (`agent::has_session` says no, the working directory is gone, the run is
36//! over, the agent left the roster) the question stays as it is, the owner
37//! gets a notification saying why, and nothing else runs. A new agent would
38//! have to be caught up on everything the first one knew, and would act on the
39//! owner's words without the context they were written against.
40//!
41//! It never assumes a resume worked either: a turn that fails or times out
42//! leaves the word undelivered, and a later tick tries again.
43
44use std::collections::HashMap;
45use std::path::PathBuf;
46use std::time::{Duration, Instant};
47
48use anyhow::Result;
49use jiff::Timestamp;
50
51use crate::agent::{self, Invocation, SeatState};
52use crate::ask::{Lease, Question, QuestionStatus, Questions, Waiter as Note, WaiterKind, Who};
53use crate::config::{AgentSpec, Config};
54use crate::notices::{self, Notice};
55use crate::prompt;
56use crate::run::RunState;
57
58/// How often the waiter looks at the store.
59pub const TICK: Duration = Duration::from_secs(5);
60
61/// Wall-clock limit for one resumed turn. A resumed agent may go on to do real
62/// work once it has the owner's word, but it is one conversation turn, not a
63/// whole node.
64const DELIVERY_TIMEOUT: Duration = Duration::from_secs(15 * 60);
65
66/// How long a failed or refused delivery is left alone before the same input
67/// is tried again. Long enough not to hammer a seat that is out of quota,
68/// short enough that a fixed problem does not wait a day.
69const RETRY_AFTER: Duration = Duration::from_secs(10 * 60);
70
71/// What the agent has not read yet.
72#[derive(Debug, Clone, PartialEq, Eq)]
73pub enum Word {
74    /// The owner spoke back without deciding.
75    Said(String),
76    /// The owner decided.
77    Answered(String),
78}
79
80/// What to do about one question, this tick.
81#[derive(Debug, Clone, PartialEq, Eq)]
82pub enum Action {
83    /// Nothing: not ours, someone is waiting, or nothing new to hand over.
84    Idle,
85    /// The deadline passed with no answer and nobody waiting.
86    Expire,
87    /// Resume the asking seat with this.
88    Deliver(Word),
89}
90
91/// The whole decision, with no I/O.
92///
93/// `seat_busy` is whether the asking seat's CLI is still mid-turn - the agent
94/// that exited its `magi ask` may not have exited itself, and resuming a
95/// conversation that is running would fork it.
96///
97/// Order matters: a word the owner gave before the deadline is delivered even
98/// if the deadline has passed by the time the waiter looks, and only then does
99/// an unanswered question expire.
100pub fn decide(
101    q: &Question,
102    lease: Option<&Lease>,
103    seat_busy: bool,
104    default_timeout: u64,
105    now: Timestamp,
106) -> Action {
107    if q.cwd.is_none() || lease.is_some_and(|l| l.fresh(now)) {
108        return Action::Idle;
109    }
110    match q.status {
111        QuestionStatus::Abandoned => Action::Idle,
112        QuestionStatus::Answered => match q.resolution() {
113            Some(a) if !q.answer_delivered && !seat_busy => Action::Deliver(Word::Answered(a)),
114            _ => Action::Idle,
115        },
116        QuestionStatus::Open => {
117            if let Some(said) = q.unread_from_owner() {
118                if seat_busy {
119                    return Action::Idle;
120                }
121                return Action::Deliver(Word::Said(said));
122            }
123            let secs = if q.answer_timeout > 0 {
124                q.answer_timeout
125            } else {
126                default_timeout
127            };
128            let deadline = q.last_activity().saturating_add(secs as i64);
129            if now.as_second() > deadline {
130                Action::Expire
131            } else {
132                Action::Idle
133            }
134        }
135    }
136}
137
138/// Everything needed to resume one seat.
139struct Target {
140    spec: AgentSpec,
141    seat: SeatState,
142    cwd: PathBuf,
143    sessions: bool,
144    allow_write: bool,
145    language: String,
146}
147
148/// The waiter's state: the store it watches and what it must remember between
149/// ticks (only which inputs it has already failed on).
150pub struct Waiter {
151    store: Questions,
152    home: PathBuf,
153    /// The config the conductor runs under, for a question the conductor asked
154    /// (it has no run whose snapshot could say).
155    conduct_cfg: Option<Config>,
156    /// `answer_timeout` for a question that did not record its own.
157    default_timeout: u64,
158    /// Inputs already tried and failed, by question id: `(key, until)`.
159    memo: HashMap<String, (String, Instant)>,
160}
161
162impl Waiter {
163    /// A waiter over `store`, with `home` as the magi home its runs live under.
164    pub fn new(store: Questions, home: PathBuf, conduct_cfg: Option<Config>) -> Self {
165        let default_timeout = conduct_cfg
166            .as_ref()
167            .map_or(Config::default().graph.answer_timeout, |c| {
168                c.graph.answer_timeout
169            });
170        Self {
171            store,
172            home,
173            conduct_cfg,
174            default_timeout,
175            memo: HashMap::new(),
176        }
177    }
178
179    /// Is this question's asking seat still mid-turn?
180    fn seat_busy(&self, q: &Question, now: Timestamp) -> bool {
181        if q.run == crate::conduct::NODE {
182            return crate::conduct::busy(&self.home);
183        }
184        RunState::load_under(&q.run, &self.home).is_ok_and(|s| {
185            s.seats_active()
186                .any(|(k, a)| *k == q.seat && a.remaining_secs(now) > 0)
187        })
188    }
189
190    /// Look at every question once and act on what needs acting on. `halt` is
191    /// asked between questions and during a delivery: true means the daemon is
192    /// parking, and whatever is in flight is dropped, undelivered and
193    /// resumable.
194    pub async fn tick(&mut self, now: Timestamp, halt: &(dyn Fn() -> bool + Sync)) {
195        for q in self.store.list() {
196            if halt() {
197                return;
198            }
199            let lease = self.store.read_lease(&q.id);
200            match decide(&q, lease.as_ref(), false, self.default_timeout, now) {
201                Action::Idle => {}
202                Action::Expire => self.expire(&q),
203                Action::Deliver(word) => {
204                    if self.seat_busy(&q, now) {
205                        continue;
206                    }
207                    let key = format!("{}:{:?}", q.thread.len(), word);
208                    if matches!(self.memo.get(&q.id), Some((k, until)) if *k == key && Instant::now() < *until)
209                    {
210                        continue;
211                    }
212                    if let Err(e) = self.deliver(&q, &word, halt).await {
213                        tracing::warn!("question {}: {e:#}", q.short());
214                    }
215                }
216            }
217        }
218    }
219
220    /// Abandon an unanswered question whose deadline passed with nobody
221    /// waiting - the words the asker's own wait would have used.
222    fn expire(&self, q: &Question) {
223        let secs = if q.answer_timeout > 0 {
224            q.answer_timeout
225        } else {
226            self.default_timeout
227        };
228        let why = format!("no answer within {}s of asking", secs.max(1));
229        let done = self.store.update(&q.id, |r| {
230            // Decided again on the record as it is now: the owner may have
231            // spoken since the tick read it, and a word given in time is
232            // delivered, not abandoned.
233            let now = Timestamp::now();
234            let lease = self.store.read_lease(&r.id);
235            if decide(r, lease.as_ref(), false, self.default_timeout, now) != Action::Expire {
236                return Ok(false);
237            }
238            r.abandon(&why);
239            r.waiter = None;
240            Ok(true)
241        });
242        match done {
243            Ok((_, false)) => {}
244            Ok((_, true)) => {
245                self.store.drop_lease(&q.id);
246                tracing::warn!(
247                    "question {} went unanswered for {secs}s with nobody waiting; \
248                     it stays as the record of it",
249                    q.short()
250                );
251            }
252            Err(e) => tracing::warn!("could not expire question {}: {e:#}", q.short()),
253        }
254    }
255
256    /// Work out how to resume the seat, or why it cannot be.
257    fn target(&self, q: &Question) -> std::result::Result<Target, String> {
258        let cwd = q
259            .cwd
260            .as_deref()
261            .map(PathBuf::from)
262            .filter(|p| p.is_dir())
263            .ok_or("the directory the agent was working in is gone")?;
264        let (cfg, seat, allow_write) = if q.run == crate::conduct::NODE {
265            let cfg = self
266                .conduct_cfg
267                .clone()
268                .ok_or("the conductor's configuration is not available")?;
269            let seat = crate::conduct::load_seat(&self.home)
270                .ok_or("the conductor has no recorded session")?;
271            (cfg, seat, false)
272        } else {
273            let state = RunState::load_under(&q.run, &self.home)
274                .map_err(|_| "the run that asked is gone".to_owned())?;
275            if !state.status.resumable() {
276                return Err(format!(
277                    "the run that asked has already {}",
278                    state.status.as_str()
279                ));
280            }
281            let seat = state
282                .seats
283                .get(&q.seat)
284                .cloned()
285                .ok_or("the run has no record of the asking seat")?;
286            let write = matches!(q.node.as_str(), "implement" | "fix");
287            (state.config, seat, write)
288        };
289        let spec = cfg
290            .agent(&seat.agent)
291            .map_err(|_| format!("agent `{}` is no longer in the roster", seat.agent))?
292            .clone();
293        if !agent::has_session(spec.kind, &seat, cfg.graph.sessions) {
294            return Err("the seat has no session to resume".to_owned());
295        }
296        Ok(Target {
297            spec,
298            seat,
299            cwd,
300            sessions: cfg.graph.sessions,
301            allow_write,
302            language: cfg.graph.language.clone(),
303        })
304    }
305
306    /// Say why the owner's word cannot reach anybody, once per input.
307    fn refuse(&mut self, q: &Question, key: String, why: &str) {
308        tracing::warn!("question {} cannot be delivered: {why}", q.short());
309        notices::raise_in(
310            &self.home,
311            Notice::warn(
312                &format!("question:{}", q.id),
313                format!(
314                    "Question {} \"{}\": what you said cannot reach the agent \
315                     ({why}). It is recorded, but nothing will read it.",
316                    q.short(),
317                    q.summary
318                ),
319            ),
320        );
321        let _ = self.store.update(&q.id, |r| {
322            r.waiter = None;
323            Ok(())
324        });
325        self.memo
326            .insert(q.id.clone(), (key, Instant::now() + RETRY_AFTER));
327    }
328
329    /// Resume the asking seat with the owner's word.
330    async fn deliver(
331        &mut self,
332        q: &Question,
333        word: &Word,
334        halt: &(dyn Fn() -> bool + Sync),
335    ) -> Result<()> {
336        let key = format!("{}:{:?}", q.thread.len(), word);
337        let mut target = match self.target(q) {
338            Ok(t) => t,
339            Err(why) => {
340                self.refuse(q, key, &why);
341                return Ok(());
342            }
343        };
344
345        let claim = self.store.root().join(format!("{}.claim", q.id));
346        if !take_claim(&claim) {
347            return Ok(());
348        }
349        let _release = Release(claim);
350
351        // Re-read under the claim: the asker may have come back, or the owner
352        // may have said more, since the tick's copy was taken.
353        let q = self.store.get(&q.id)?;
354        let now = Timestamp::now();
355        let lease = self.store.read_lease(&q.id);
356        let Action::Deliver(word) = decide(&q, lease.as_ref(), false, self.default_timeout, now)
357        else {
358            return Ok(());
359        };
360        let snapshot = q.thread.len();
361
362        self.store.beat(&q.id, WaiterKind::Daemon);
363        self.store.update(&q.id, |r| {
364            r.waiter = Some(Note {
365                kind: WaiterKind::Daemon,
366                since: now,
367            });
368            Ok(())
369        })?;
370        tracing::info!(
371            "question {}: the asker is gone, resuming seat {} with the owner's word",
372            q.short(),
373            q.seat
374        );
375
376        let thread: Vec<(&str, &str)> = q
377            .thread
378            .iter()
379            .map(|t| {
380                (
381                    if t.who == Who::Operator {
382                        "operator"
383                    } else {
384                        "agent"
385                    },
386                    t.body.as_str(),
387                )
388            })
389            .collect();
390        // What the owner said and the agent has not read has its own section in
391        // the prompt, so the recap stops at what it already read.
392        let shown = match word {
393            Word::Said(_) => &thread[..q.delivered_turns.min(thread.len())],
394            Word::Answered(_) => &thread[..],
395        };
396        let owner_word = match &word {
397            Word::Said(s) => prompt::OwnerWord::Said(s),
398            Word::Answered(a) => prompt::OwnerWord::Answered(a),
399        };
400        let body = prompt::question_resumed(
401            &q.id,
402            &q.summary,
403            &q.detail,
404            shown,
405            &owner_word,
406            &target.language,
407        );
408
409        let artifacts = self.store.root().join(format!("{}.delivery", q.id));
410        let stem = format!("deliver-{}", now.as_second());
411        let cache_dir = None;
412        let inv = Invocation {
413            cwd: &target.cwd,
414            prompt: &body,
415            timeout: DELIVERY_TIMEOUT,
416            allow_write: target.allow_write,
417            sessions: target.sessions,
418            artifacts: &artifacts,
419            stem: &stem,
420            run: &q.run,
421            node: &q.node,
422            cache_dir,
423            attachments: &[],
424        };
425
426        let store = self.store.clone();
427        let id = q.id.clone();
428        let out = {
429            let fut = agent::invoke(&target.spec, &mut target.seat, &inv);
430            tokio::pin!(fut);
431            let mut beat = tokio::time::interval(Duration::from_secs(1));
432            let mut beats = 0u32;
433            loop {
434                tokio::select! {
435                    r = &mut fut => break Some(r),
436                    _ = beat.tick() => {
437                        if halt() {
438                            break None;
439                        }
440                        beats += 1;
441                        if beats % 20 == 0 {
442                            store.beat(&id, WaiterKind::Daemon);
443                        }
444                    }
445                }
446            }
447        };
448
449        let Some(out) = out else {
450            // Parking: leave everything as it was, resumable.
451            let _ = self.store.update(&q.id, |r| {
452                r.waiter = None;
453                Ok(())
454            });
455            return Ok(());
456        };
457        let out = match out {
458            Ok(o) if o.usable() => o,
459            other => {
460                let why = match other {
461                    Ok(o) if o.timed_out => "the resumed turn timed out".to_owned(),
462                    Ok(o) if o.quota_exhausted() => "the agent is out of quota".to_owned(),
463                    Ok(o) => format!("the resumed turn failed (exit {:?})", o.exit_code),
464                    Err(e) => format!("the agent could not be started: {e:#}"),
465                };
466                self.refuse(&q, key, &why);
467                return Ok(());
468            }
469        };
470
471        let text = out.text.trim().to_owned();
472        self.store.update(&q.id, |r| {
473            r.waiter = None;
474            match &word {
475                Word::Answered(_) => r.answer_delivered = true,
476                Word::Said(_) => {
477                    r.delivered_turns = r.delivered_turns.max(snapshot);
478                    // An agent that answered in prose instead of calling
479                    // `magi ask --thread` still spoke; keep it on the record so
480                    // the owner reads it. If it did use --thread, or the owner
481                    // said more meanwhile, the record already moved on.
482                    let agent_spoke = r.thread[snapshot.min(r.thread.len())..]
483                        .iter()
484                        .any(|t| t.who == Who::Agent);
485                    if !agent_spoke && r.thread.len() == snapshot && r.status.open() {
486                        let choices = r.choices.clone();
487                        r.reply(text, choices)?;
488                    }
489                }
490            }
491            Ok(())
492        })?;
493        self.memo.remove(&q.id);
494        Ok(())
495    }
496}
497
498/// Take the delivery claim, so two waiters (`magi serve` twice, or a stray
499/// second daemon) cannot resume the same seat at once. A claim older than a
500/// whole delivery plus slack was left by a process that died.
501fn take_claim(path: &std::path::Path) -> bool {
502    let attempt = || {
503        std::fs::OpenOptions::new()
504            .write(true)
505            .create_new(true)
506            .open(path)
507            .is_ok()
508    };
509    if attempt() {
510        return true;
511    }
512    let stale = std::fs::metadata(path)
513        .and_then(|m| m.modified())
514        .ok()
515        .and_then(|t| t.elapsed().ok())
516        .is_some_and(|age| age > DELIVERY_TIMEOUT + Duration::from_secs(60));
517    if stale {
518        let _ = std::fs::remove_file(path);
519        return attempt();
520    }
521    false
522}
523
524/// Removes the delivery claim when the delivery ends, however it ends.
525struct Release(PathBuf);
526
527impl Drop for Release {
528    fn drop(&mut self) {
529        let _ = std::fs::remove_file(&self.0);
530    }
531}
532
533/// The waiter's loop, until `stop` is asked for.
534///
535/// A stop ends the loop between ticks; a park also drops a delivery in flight.
536/// Neither writes anything to a question: the lease just ages out and the next
537/// waiter - this daemon restarted, or another - reads the same state off disk.
538pub async fn run(mut waiter: Waiter, stop: crate::daemon::Stop) {
539    let halt = {
540        let stop = stop.clone();
541        move || stop.parking()
542    };
543    while !stop.stopped() {
544        waiter.tick(Timestamp::now(), &halt).await;
545        tokio::time::sleep(TICK).await;
546    }
547}
548
549#[cfg(test)]
550mod tests {
551    use super::*;
552
553    fn ts(secs: i64) -> Timestamp {
554        Timestamp::from_second(secs).unwrap()
555    }
556
557    fn asked(at: i64) -> Question {
558        let mut q = Question::new(
559            "run".into(),
560            "implement".into(),
561            "impl-A".into(),
562            "which?".into(),
563            String::new(),
564            vec![],
565        );
566        q.cwd = Some("/tmp".into());
567        q.asked_at = ts(at);
568        q.answer_timeout = 1000;
569        q
570    }
571
572    fn lease(at: i64) -> Lease {
573        Lease {
574            kind: WaiterKind::Asker,
575            pid: 1,
576            beat_at: ts(at),
577        }
578    }
579
580    #[test]
581    fn a_fresh_lease_means_nobody_else_acts() {
582        let mut q = asked(0);
583        q.say("why?").unwrap();
584        assert_eq!(
585            decide(&q, Some(&lease(95)), false, 86_400, ts(100)),
586            Action::Idle
587        );
588    }
589
590    #[test]
591    fn a_stale_lease_delivers_what_the_owner_said() {
592        let mut q = asked(0);
593        q.say("why?").unwrap();
594        assert_eq!(
595            decide(&q, Some(&lease(0)), false, 86_400, ts(500)),
596            Action::Deliver(Word::Said("why?".into()))
597        );
598        assert_eq!(
599            decide(&q, None, true, 86_400, ts(500)),
600            Action::Idle,
601            "a seat still mid-turn is not resumed"
602        );
603    }
604
605    #[test]
606    fn an_answer_nobody_read_is_delivered_once() {
607        let mut q = asked(0);
608        q.choices = vec!["A".into()];
609        q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
610        assert_eq!(
611            decide(&q, None, false, 86_400, ts(10)),
612            Action::Deliver(Word::Answered("A".into()))
613        );
614        q.answer_delivered = true;
615        assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
616    }
617
618    #[test]
619    fn only_asked_questions_and_only_before_the_deadline() {
620        let mut q = asked(0);
621        assert_eq!(decide(&q, None, false, 86_400, ts(999)), Action::Idle);
622        assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Expire);
623        q.cwd = None;
624        assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Idle);
625    }
626
627    #[test]
628    fn the_deadline_runs_from_the_last_turn_not_from_asking() {
629        let mut q = asked(0);
630        q.thread.push(crate::ask::Turn {
631            who: Who::Agent,
632            body: "context".into(),
633            at: ts(4000),
634        });
635        q.delivered_turns = 1;
636        assert_eq!(decide(&q, None, false, 86_400, ts(4500)), Action::Idle);
637        assert_eq!(decide(&q, None, false, 86_400, ts(5001)), Action::Expire);
638    }
639
640    #[test]
641    fn a_word_given_in_time_is_delivered_after_the_deadline() {
642        let mut q = asked(0);
643        q.say("wait").unwrap();
644        assert!(matches!(
645            decide(&q, None, false, 86_400, ts(5000)),
646            Action::Deliver(_)
647        ));
648    }
649
650    #[test]
651    fn an_action_answer_is_delivered_until_the_daemon_marks_it_handled() {
652        let mut q = asked(0);
653        q.choices = vec!["A".into()];
654        q.actions
655            .insert("A".into(), crate::ask::ChoiceAction::Requeue);
656        q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
657        assert_eq!(
658            decide(&q, None, false, 86_400, ts(10)),
659            Action::Deliver(Word::Answered("A".into())),
660            "an action nobody applied must not swallow the answer"
661        );
662        q.answer_delivered = true;
663        assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
664    }
665}