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                    // A question with a deputy is the deputy's to read
205                    // (`crate::deputy`); resuming the conductor's own seat
206                    // here would fork a conversation that is not its own.
207                    if q.deputy.is_some() {
208                        // Unless that deputy is spent: then nobody will read
209                        // the word, and the deadline still has to retire the
210                        // question.
211                        if crate::deputy::exhausted_past_deadline(
212                            &q,
213                            crate::deputy::can_start(
214                                self.conduct_cfg.as_ref(),
215                                crate::deputy::agent_of(&q),
216                            ),
217                            self.default_timeout,
218                            now,
219                        ) {
220                            self.expire(&q);
221                        }
222                        continue;
223                    }
224                    if self.seat_busy(&q, now) {
225                        continue;
226                    }
227                    let key = format!("{}:{:?}", q.thread.len(), word);
228                    if matches!(self.memo.get(&q.id), Some((k, until)) if *k == key && Instant::now() < *until)
229                    {
230                        continue;
231                    }
232                    if let Err(e) = self.deliver(&q, &word, halt).await {
233                        tracing::warn!("question {}: {e:#}", q.short());
234                    }
235                }
236            }
237        }
238    }
239
240    /// Abandon an unanswered question whose deadline passed with nobody
241    /// waiting - the words the asker's own wait would have used.
242    fn expire(&self, q: &Question) {
243        let secs = if q.answer_timeout > 0 {
244            q.answer_timeout
245        } else {
246            self.default_timeout
247        };
248        let startable =
249            crate::deputy::can_start(self.conduct_cfg.as_ref(), crate::deputy::agent_of(q));
250        let why = format!("no answer within {}s of asking", secs.max(1));
251        let done = self.store.update(&q.id, |r| {
252            // Decided again on the record as it is now: the owner may have
253            // spoken since the tick read it, and a word given in time is
254            // delivered, not abandoned.
255            let now = Timestamp::now();
256            let lease = self.store.read_lease(&r.id);
257            if decide(r, lease.as_ref(), false, self.default_timeout, now) != Action::Expire
258                && !(lease.as_ref().is_none_or(|l| !l.fresh(now))
259                    && crate::deputy::exhausted_past_deadline(
260                        r,
261                        startable,
262                        self.default_timeout,
263                        now,
264                    ))
265            {
266                return Ok(false);
267            }
268            r.abandon(&why);
269            r.waiter = None;
270            Ok(true)
271        });
272        match done {
273            Ok((_, false)) => {}
274            Ok((_, true)) => {
275                self.store.drop_lease(&q.id);
276                tracing::warn!(
277                    "question {} went unanswered for {secs}s with nobody waiting; \
278                     it stays as the record of it",
279                    q.short()
280                );
281            }
282            Err(e) => tracing::warn!("could not expire question {}: {e:#}", q.short()),
283        }
284    }
285
286    /// Work out how to resume the seat, or why it cannot be.
287    fn target(&self, q: &Question) -> std::result::Result<Target, String> {
288        let cwd = q
289            .cwd
290            .as_deref()
291            .map(PathBuf::from)
292            .filter(|p| p.is_dir())
293            .ok_or("the directory the agent was working in is gone")?;
294        let (cfg, seat, allow_write) = if q.run == crate::conduct::NODE {
295            let cfg = self
296                .conduct_cfg
297                .clone()
298                .ok_or("the conductor's configuration is not available")?;
299            let seat = crate::conduct::load_seat(&self.home)
300                .ok_or("the conductor has no recorded session")?;
301            (cfg, seat, false)
302        } else {
303            let state = RunState::load_under(&q.run, &self.home)
304                .map_err(|_| "the run that asked is gone".to_owned())?;
305            if !state.status.resumable() {
306                return Err(format!(
307                    "the run that asked has already {}",
308                    state.status.as_str()
309                ));
310            }
311            let seat = state
312                .seats
313                .get(&q.seat)
314                .cloned()
315                .ok_or("the run has no record of the asking seat")?;
316            let write = matches!(q.node.as_str(), "implement" | "fix");
317            (state.config, seat, write)
318        };
319        let spec = cfg
320            .agent(&seat.agent)
321            .map_err(|_| format!("agent `{}` is no longer in the roster", seat.agent))?
322            .clone();
323        if !agent::has_session(spec.kind, &seat, cfg.graph.sessions) {
324            return Err("the seat has no session to resume".to_owned());
325        }
326        Ok(Target {
327            spec,
328            seat,
329            cwd,
330            sessions: cfg.graph.sessions,
331            allow_write,
332            language: cfg.graph.language.clone(),
333        })
334    }
335
336    /// Say why the owner's word cannot reach anybody, once per input.
337    fn refuse(&mut self, q: &Question, key: String, why: &str) {
338        tracing::warn!("question {} cannot be delivered: {why}", q.short());
339        notices::raise_in(
340            &self.home,
341            Notice::warn(
342                &format!("question:{}", q.id),
343                format!(
344                    "Question {} \"{}\": what you said cannot reach the agent \
345                     ({why}). It is recorded, but nothing will read it.",
346                    q.short(),
347                    q.summary
348                ),
349            ),
350        );
351        let _ = self.store.update(&q.id, |r| {
352            r.waiter = None;
353            Ok(())
354        });
355        self.memo
356            .insert(q.id.clone(), (key, Instant::now() + RETRY_AFTER));
357    }
358
359    /// Resume the asking seat with the owner's word.
360    async fn deliver(
361        &mut self,
362        q: &Question,
363        word: &Word,
364        halt: &(dyn Fn() -> bool + Sync),
365    ) -> Result<()> {
366        let key = format!("{}:{:?}", q.thread.len(), word);
367        let mut target = match self.target(q) {
368            Ok(t) => t,
369            Err(why) => {
370                self.refuse(q, key, &why);
371                return Ok(());
372            }
373        };
374
375        let claim = self.store.root().join(format!("{}.claim", q.id));
376        if !take_claim(&claim, DELIVERY_TIMEOUT + Duration::from_secs(60)) {
377            return Ok(());
378        }
379        let _release = Release(claim);
380
381        // Re-read under the claim: the asker may have come back, or the owner
382        // may have said more, since the tick's copy was taken.
383        let q = self.store.get(&q.id)?;
384        let now = Timestamp::now();
385        let lease = self.store.read_lease(&q.id);
386        let Action::Deliver(word) = decide(&q, lease.as_ref(), false, self.default_timeout, now)
387        else {
388            return Ok(());
389        };
390        let snapshot = q.thread.len();
391
392        self.store.beat(&q.id, WaiterKind::Daemon);
393        self.store.update(&q.id, |r| {
394            r.waiter = Some(Note {
395                kind: WaiterKind::Daemon,
396                since: now,
397            });
398            Ok(())
399        })?;
400        tracing::info!(
401            "question {}: the asker is gone, resuming seat {} with the owner's word",
402            q.short(),
403            q.seat
404        );
405
406        let thread: Vec<(&str, &str)> = q
407            .thread
408            .iter()
409            .map(|t| {
410                (
411                    if t.who == Who::Operator {
412                        "operator"
413                    } else {
414                        "agent"
415                    },
416                    t.body.as_str(),
417                )
418            })
419            .collect();
420        // What the owner said and the agent has not read has its own section in
421        // the prompt, so the recap stops at what it already read.
422        let shown = match word {
423            Word::Said(_) => &thread[..q.delivered_turns.min(thread.len())],
424            Word::Answered(_) => &thread[..],
425        };
426        let owner_word = match &word {
427            Word::Said(s) => prompt::OwnerWord::Said(s),
428            Word::Answered(a) => prompt::OwnerWord::Answered(a),
429        };
430        let body = prompt::question_resumed(
431            &q.id,
432            &q.summary,
433            &q.detail,
434            shown,
435            &owner_word,
436            &target.language,
437        );
438
439        let artifacts = self.store.root().join(format!("{}.delivery", q.id));
440        let stem = format!("deliver-{}", now.as_second());
441        let cache_dir = None;
442        let inv = Invocation {
443            cwd: &target.cwd,
444            prompt: &body,
445            timeout: DELIVERY_TIMEOUT,
446            allow_write: target.allow_write,
447            sessions: target.sessions,
448            artifacts: &artifacts,
449            stem: &stem,
450            run: &q.run,
451            node: &q.node,
452            cache_dir,
453            attachments: &[],
454            writable: &[],
455        };
456
457        let store = self.store.clone();
458        let id = q.id.clone();
459        let out = {
460            let fut = agent::invoke(&target.spec, &mut target.seat, &inv);
461            tokio::pin!(fut);
462            let mut beat = tokio::time::interval(Duration::from_secs(1));
463            let mut beats = 0u32;
464            loop {
465                tokio::select! {
466                    r = &mut fut => break Some(r),
467                    _ = beat.tick() => {
468                        if halt() {
469                            break None;
470                        }
471                        beats += 1;
472                        if beats % 20 == 0 {
473                            store.beat(&id, WaiterKind::Daemon);
474                        }
475                    }
476                }
477            }
478        };
479
480        let Some(out) = out else {
481            // Parking: leave everything as it was, resumable.
482            let _ = self.store.update(&q.id, |r| {
483                r.waiter = None;
484                Ok(())
485            });
486            return Ok(());
487        };
488        let out = match out {
489            Ok(o) if o.usable() => o,
490            other => {
491                let why = match other {
492                    Ok(o) if o.timed_out => "the resumed turn timed out".to_owned(),
493                    Ok(o) if o.quota_exhausted() => "the agent is out of quota".to_owned(),
494                    Ok(o) => format!("the resumed turn failed (exit {:?})", o.exit_code),
495                    Err(e) => format!("the agent could not be started: {e:#}"),
496                };
497                self.refuse(&q, key, &why);
498                return Ok(());
499            }
500        };
501
502        let text = out.text.trim().to_owned();
503        self.store.update(&q.id, |r| {
504            r.waiter = None;
505            match &word {
506                Word::Answered(_) => r.answer_delivered = true,
507                Word::Said(_) => {
508                    r.delivered_turns = r.delivered_turns.max(snapshot);
509                    // An agent that answered in prose instead of calling
510                    // `magi ask --thread` still spoke; keep it on the record so
511                    // the owner reads it. If it did use --thread, or the owner
512                    // said more meanwhile, the record already moved on.
513                    let agent_spoke = r.thread[snapshot.min(r.thread.len())..]
514                        .iter()
515                        .any(|t| t.who == Who::Agent);
516                    if !agent_spoke && r.thread.len() == snapshot && r.status.open() {
517                        let choices = r.choices.clone();
518                        r.reply(text, choices)?;
519                    }
520                }
521            }
522            Ok(())
523        })?;
524        self.memo.remove(&q.id);
525        Ok(())
526    }
527}
528
529/// Take the delivery claim, so two waiters (`magi serve` twice, or a stray
530/// second daemon) cannot resume the same seat at once. A claim older than a
531/// whole delivery plus slack was left by a process that died.
532pub(crate) fn take_claim(path: &std::path::Path, stale_after: Duration) -> bool {
533    let attempt = || {
534        std::fs::OpenOptions::new()
535            .write(true)
536            .create_new(true)
537            .open(path)
538            .is_ok()
539    };
540    if attempt() {
541        return true;
542    }
543    let stale = std::fs::metadata(path)
544        .and_then(|m| m.modified())
545        .ok()
546        .and_then(|t| t.elapsed().ok())
547        .is_some_and(|age| age > stale_after);
548    if stale {
549        let _ = std::fs::remove_file(path);
550        return attempt();
551    }
552    false
553}
554
555/// Removes the delivery claim when the delivery ends, however it ends.
556pub(crate) struct Release(pub(crate) PathBuf);
557
558impl Drop for Release {
559    fn drop(&mut self) {
560        let _ = std::fs::remove_file(&self.0);
561    }
562}
563
564/// The waiter's loop, until `stop` is asked for.
565///
566/// A stop ends the loop between ticks; a park also drops a delivery in flight.
567/// Neither writes anything to a question: the lease just ages out and the next
568/// waiter - this daemon restarted, or another - reads the same state off disk.
569pub async fn run(mut waiter: Waiter, stop: crate::daemon::Stop) {
570    let halt = {
571        let stop = stop.clone();
572        move || stop.parking()
573    };
574    while !stop.stopped() {
575        waiter.tick(Timestamp::now(), &halt).await;
576        tokio::time::sleep(TICK).await;
577    }
578}
579
580#[cfg(test)]
581mod tests {
582    use super::*;
583
584    fn ts(secs: i64) -> Timestamp {
585        Timestamp::from_second(secs).unwrap()
586    }
587
588    fn asked(at: i64) -> Question {
589        let mut q = Question::new(
590            "run".into(),
591            "implement".into(),
592            "impl-A".into(),
593            "which?".into(),
594            String::new(),
595            vec![],
596        );
597        q.cwd = Some("/tmp".into());
598        q.asked_at = ts(at);
599        q.answer_timeout = 1000;
600        q
601    }
602
603    fn lease(at: i64) -> Lease {
604        Lease {
605            kind: WaiterKind::Asker,
606            pid: 1,
607            beat_at: ts(at),
608        }
609    }
610
611    #[test]
612    fn a_fresh_lease_means_nobody_else_acts() {
613        let mut q = asked(0);
614        q.say("why?").unwrap();
615        assert_eq!(
616            decide(&q, Some(&lease(95)), false, 86_400, ts(100)),
617            Action::Idle
618        );
619    }
620
621    #[test]
622    fn a_stale_lease_delivers_what_the_owner_said() {
623        let mut q = asked(0);
624        q.say("why?").unwrap();
625        assert_eq!(
626            decide(&q, Some(&lease(0)), false, 86_400, ts(500)),
627            Action::Deliver(Word::Said("why?".into()))
628        );
629        assert_eq!(
630            decide(&q, None, true, 86_400, ts(500)),
631            Action::Idle,
632            "a seat still mid-turn is not resumed"
633        );
634    }
635
636    #[test]
637    fn an_answer_nobody_read_is_delivered_once() {
638        let mut q = asked(0);
639        q.choices = vec!["A".into()];
640        q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
641        assert_eq!(
642            decide(&q, None, false, 86_400, ts(10)),
643            Action::Deliver(Word::Answered("A".into()))
644        );
645        q.answer_delivered = true;
646        assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
647    }
648
649    #[test]
650    fn only_asked_questions_and_only_before_the_deadline() {
651        let mut q = asked(0);
652        assert_eq!(decide(&q, None, false, 86_400, ts(999)), Action::Idle);
653        assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Expire);
654        q.cwd = None;
655        assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Idle);
656    }
657
658    #[test]
659    fn the_deadline_runs_from_the_last_turn_not_from_asking() {
660        let mut q = asked(0);
661        q.thread.push(crate::ask::Turn {
662            who: Who::Agent,
663            body: "context".into(),
664            at: ts(4000),
665        });
666        q.delivered_turns = 1;
667        assert_eq!(decide(&q, None, false, 86_400, ts(4500)), Action::Idle);
668        assert_eq!(decide(&q, None, false, 86_400, ts(5001)), Action::Expire);
669    }
670
671    #[test]
672    fn a_word_given_in_time_is_delivered_after_the_deadline() {
673        let mut q = asked(0);
674        q.say("wait").unwrap();
675        assert!(matches!(
676            decide(&q, None, false, 86_400, ts(5000)),
677            Action::Deliver(_)
678        ));
679    }
680
681    #[test]
682    fn an_action_answer_is_delivered_until_the_daemon_marks_it_handled() {
683        let mut q = asked(0);
684        q.choices = vec!["A".into()];
685        q.actions
686            .insert("A".into(), crate::ask::ChoiceAction::Requeue);
687        q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
688        assert_eq!(
689            decide(&q, None, false, 86_400, ts(10)),
690            Action::Deliver(Word::Answered("A".into())),
691            "an action nobody applied must not swallow the answer"
692        );
693        q.answer_delivered = true;
694        assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
695    }
696}