Skip to main content

magi/
direct.rs

1//! Runs started by hand are tasks.
2//!
3//! `magi run` and `magi review` used to mint a run no task knew about: no
4//! detail page, no attempt history, invisible to duplicate detection, and an
5//! orphan the moment anything cleaned up after it. They now file a task first,
6//! through the queue like `magi task add`, and then either hand it to a live
7//! loop or execute it here through the daemon's own attempt
8//! ([`crate::daemon::run_claimed`]). There is no second dispatcher in this
9//! module: it decides *who* runs the task and watches it, nothing more.
10//!
11//! * A live loop in another process (a fresh `daemon.json` heartbeat, see
12//!   [`crate::daemon::foreign_loop`]) gets the task as `urgent`, runs it in
13//!   its own urgent slot, and the caller [`follow`]s it. The caller never
14//!   claims or executes it: two processes must not drive one task.
15//! * Otherwise the caller claims the task **before it is written**, so a loop
16//!   that starts mid-run finds the claim and leaves it alone, and executes it.
17
18use std::path::{Path, PathBuf};
19use std::time::{Duration, Instant};
20
21use anyhow::{Context, Result, bail};
22use jiff::Timestamp;
23
24use crate::config::Config;
25use crate::queue::{Claim, Queue, RunOverrides, Source, Task, TaskStatus};
26
27/// What to file. Everything `magi run` / `magi review` take that changes what
28/// the run does rides on the task, so a loop that executes it behaves the same
29/// as this process would have.
30#[derive(Debug, Clone)]
31pub struct Filing {
32    /// The task text; for a review, a line saying which branch.
33    pub instruction: String,
34    /// The task's title.
35    pub title: String,
36    /// Absolute repository path.
37    pub repo: PathBuf,
38    /// Who asked: the operator, or the agent seat named by `MAGI_RUN` / `MAGI_NODE`.
39    pub source: Source,
40    /// `--solo`.
41    pub solo: bool,
42    /// Command-line choices carried onto the task.
43    pub overrides: RunOverrides,
44    /// `magi review <branch>`: the branch to review.
45    pub review_of: Option<String>,
46}
47
48impl Filing {
49    fn into_task(self) -> Task {
50        let mut task = Task::new(self.title, self.instruction, self.repo, self.source);
51        task.solo = self.solo;
52        task.review_of = self.review_of;
53        task.overrides = Some(self.overrides);
54        task
55    }
56}
57
58/// Who will execute a filed task.
59#[derive(Debug)]
60pub enum Filed {
61    /// A loop in another process owns the queue; the task is `urgent` and
62    /// unclaimed. `pid` is the loop's, when it published one.
63    Daemon {
64        /// The filed task.
65        task: Task,
66        /// The loop's pid.
67        pid: Option<u32>,
68    },
69    /// Nobody else is serving: the task is claimed for this process, and the
70    /// claim lives as long as this value.
71    Standalone {
72        /// The filed task.
73        task: Task,
74        /// Held until dropped, which is what keeps a later loop away.
75        claim: Claim,
76    },
77}
78
79/// File `filing` and decide who runs it. `home` is the magi home whose
80/// `daemon.json` is judged; `own_pid` is this process (a loop of our own is
81/// not "another" one). `force_local` (`--dry-run`) never hands over: a dry run
82/// spends no agent call and a loop would not honour that.
83pub fn file(
84    queue: &Queue,
85    home: &Path,
86    now: Timestamp,
87    own_pid: u32,
88    filing: Filing,
89    force_local: bool,
90) -> Result<Filed> {
91    let mut task = filing.into_task();
92    let owner = if force_local {
93        None
94    } else {
95        crate::daemon::foreign_loop(crate::daemon::read_status(home).as_ref(), now, own_pid)
96    };
97    match owner {
98        Some(pid) => {
99            task.urgent = true;
100            queue.put(&mut task).context("file the task")?;
101            Ok(Filed::Daemon { task, pid })
102        }
103        None => {
104            // Claimed first: between `put` and the claim a loop could take it.
105            let claim = queue.claim(&task.id)?;
106            queue.put(&mut task).context("file the task")?;
107            Ok(Filed::Standalone { task, claim })
108        }
109    }
110}
111
112/// The config a task's run is built with: the task's `--config` layer, then
113/// its overrides and `solo`. What the loop's `attempt` does, for the dry run.
114pub fn config_for(task: &Task, repo: &Path) -> Result<Config> {
115    let path = task.overrides.as_ref().and_then(|o| o.config.as_deref());
116    let (mut cfg, _) = Config::discover(repo, path)?;
117    if let Some(o) = &task.overrides {
118        o.apply(&mut cfg);
119    }
120    if task.solo {
121        cfg.graph.candidates = 1;
122    }
123    Ok(cfg)
124}
125
126/// How a followed task ended up.
127#[derive(Debug)]
128pub struct Followed {
129    /// The task as last read.
130    pub task: Task,
131    /// False when `max_wait` ran out first: the loop still owns the task and
132    /// nothing is known of its outcome yet.
133    pub finished: bool,
134}
135
136impl Followed {
137    /// The error a caller must exit with when the wait ran out before the
138    /// task finished: that is not a success, and the outcome is not known.
139    /// `None` once the task finished (whatever its outcome).
140    pub fn unfinished(&self) -> Option<anyhow::Error> {
141        if self.finished {
142            return None;
143        }
144        Some(anyhow::anyhow!(
145            "task {id} has not finished: it is still running under the loop. \
146             This is not a success; the outcome is not known yet. Read it with \
147             `magi task show {id}`. The loop keeps the task, so do not run it again.",
148            id = self.task.short()
149        ))
150    }
151}
152
153/// Watch `id` until the loop that owns it is done with it, calling `seen` with
154/// each run id the first time it appears on the task, and `progress` with a
155/// line whenever the newest run's status changes.
156///
157/// Finished is `done`, `held` or `blocked`: a `failed` or `queued` task is the
158/// loop's to retry. A task that stays unfinished while no loop is alive, or
159/// past `max_wait`, is an error - never taken over here, because the loop may
160/// merely be slow to heartbeat and a second driver would race it. Recovery is
161/// the loop's own claim reclaim and `Runner::resume`. Running out of
162/// `max_wait` is not an error: it returns with `finished: false`.
163pub async fn follow(
164    queue: &Queue,
165    home: &Path,
166    id: &str,
167    poll: Duration,
168    max_wait: Option<Duration>,
169    mut seen: impl FnMut(&str),
170    mut progress: impl FnMut(&str),
171) -> Result<Followed> {
172    let began = Instant::now();
173    let mut announced = 0usize;
174    let mut last_status = String::new();
175    loop {
176        let task = queue.get(id)?;
177        for run in task.runs.iter().skip(announced) {
178            seen(run);
179        }
180        announced = task.runs.len();
181        if let Some(run) = task.runs.last()
182            && crate::run::try_home().is_some()
183            && let Ok(state) = crate::run::RunState::load(run)
184        {
185            let line = format!("run {}: {}", state.short(), state.status.as_str());
186            if line != last_status {
187                progress(&line);
188                last_status = line;
189            }
190        }
191        if matches!(
192            task.status,
193            TaskStatus::Done | TaskStatus::Held | TaskStatus::Blocked
194        ) {
195            return Ok(Followed {
196                task,
197                finished: true,
198            });
199        }
200        let reading = crate::daemon::read_status(home);
201        if crate::daemon::foreign_loop(reading.as_ref(), Timestamp::now(), std::process::id())
202            .is_none()
203        {
204            bail!(
205                "no magi loop is serving the queue any more, and task {} is still {}; it was \
206                 not taken over here (the loop may only be slow to heartbeat). Start `magi \
207                 serve` to carry on, or follow it with `magi task show {}`",
208                task.short(),
209                task.status.as_str(),
210                task.short()
211            );
212        }
213        if max_wait.is_some_and(|m| began.elapsed() >= m) {
214            return Ok(Followed {
215                task,
216                finished: false,
217            });
218        }
219        tokio::time::sleep(poll).await;
220    }
221}
222
223/// A run opened inside this process (the follow-up review of `magi fix`) and
224/// the task that owns it. See [`adopt`].
225#[derive(Debug)]
226pub struct Adopted {
227    queue: Queue,
228    /// The task this process claimed and started for the run; `None` when the
229    /// claim could not be taken and the run was only linked.
230    task: Option<Task>,
231    claim: Option<Claim>,
232    quota_before: Vec<crate::run::QuotaLoss>,
233}
234
235/// Give a run this process is about to execute an owning task. A parent task
236/// that exists and is queued, failed or held is claimed, started with the
237/// run and settled by [`Adopted::finish`] exactly like an ownerless run's; if
238/// somebody else holds its claim, or the task is Done, blocked or running, the run is only linked and nothing is settled. With no such
239/// parent, a task is filed and claimed for the duration. A no-op (`None`) when
240/// no magi home is pinned.
241pub fn adopt(state: &crate::run::RunState, parent_task: Option<&str>) -> Option<Adopted> {
242    crate::run::try_home()?;
243    adopt_in(Queue::open(), state, parent_task)
244}
245
246fn adopt_in(
247    queue: Queue,
248    state: &crate::run::RunState,
249    parent_task: Option<&str>,
250) -> Option<Adopted> {
251    if let Some(parent) = parent_task
252        && let Ok(id) = queue.resolve_id(parent)
253    {
254        // The parent exists: never file a second owner for this run.
255        return match queue.claim(&id) {
256            Ok(claim) => {
257                let started = queue.get(&id).and_then(|mut task| {
258                    // A task the daemon could pick up, or one a failed
259                    // hand-started run left held (that is what a requested
260                    // follow-up review re-verifies), is this run's to settle.
261                    // A Done, blocked or running one is settled already or
262                    // owned by somebody else: restarting it would let a
263                    // follow-up review overwrite its outcome, so it only
264                    // gains the run.
265                    if !(task.status.runnable() || task.status == TaskStatus::Held) {
266                        return Ok(None);
267                    }
268                    task.start(state.id.clone());
269                    queue.put(&mut task)?;
270                    Ok(Some(task))
271                });
272                match started {
273                    Ok(None) => {
274                        drop(claim);
275                        let _ = queue.link_run(&id, &state.id);
276                        None
277                    }
278                    Ok(Some(task)) => Some(Adopted {
279                        queue,
280                        task: Some(task),
281                        claim: Some(claim),
282                        quota_before: state.quota.clone(),
283                    }),
284                    Err(e) => {
285                        tracing::warn!("could not start task {id} for run {}: {e:#}", state.id);
286                        drop(claim);
287                        let _ = queue.link_run(&id, &state.id);
288                        None
289                    }
290                }
291            }
292            Err(e) => {
293                // Another driver owns the task; do not settle over its result.
294                tracing::warn!("task {id} is claimed elsewhere ({e:#}); linking run only");
295                let _ = queue.link_run(&id, &state.id);
296                None
297            }
298        };
299    }
300    let mut task = Task::new(
301        crate::queue::title_from(&state.instruction, 72),
302        state.instruction.clone(),
303        state.repo.clone(),
304        Source::Human,
305    );
306    let claim = queue.claim(&task.id).ok()?;
307    task.start(state.id.clone());
308    queue.put(&mut task).ok()?;
309    Some(Adopted {
310        queue,
311        task: Some(task),
312        claim: Some(claim),
313        quota_before: state.quota.clone(),
314    })
315}
316
317impl Adopted {
318    /// Settle the task of an adopted ownerless run after it executed.
319    pub fn finish(mut self, state: &crate::run::RunState, result: Result<()>) {
320        if let Some(task) = &mut self.task {
321            crate::daemon::finish_attempt(
322                crate::daemon::Opts::default().max_attempts,
323                &self.queue,
324                task,
325                state,
326                &self.quota_before,
327                result,
328            );
329            crate::daemon::hold_if_runnable(&self.queue, task);
330        }
331        drop(self.claim.take());
332    }
333}
334
335#[cfg(test)]
336mod tests {
337    use super::*;
338    use crate::daemon::Status;
339
340    fn filing(repo: &Path) -> Filing {
341        Filing {
342            instruction: "review the branch".to_owned(),
343            title: "review x".to_owned(),
344            repo: repo.to_path_buf(),
345            source: Source::Human,
346            solo: false,
347            overrides: RunOverrides {
348                merge: Some("none".to_owned()),
349                ..RunOverrides::default()
350            },
351            review_of: Some("feat/x".to_owned()),
352        }
353    }
354
355    fn heartbeat(home: &Path, pid: u32, age_secs: i64) {
356        let mut status = Status::new();
357        status.pid = pid;
358        status.updated_at = Timestamp::now()
359            .checked_sub(jiff::SignedDuration::from_secs(age_secs))
360            .unwrap();
361        crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
362    }
363
364    #[test]
365    fn without_a_loop_the_task_is_claimed_before_anyone_else_can_take_it() {
366        let dir = tempfile::tempdir().unwrap();
367        let q = Queue::at(dir.path().join("queue"));
368        let filed = file(
369            &q,
370            dir.path(),
371            Timestamp::now(),
372            1,
373            filing(dir.path()),
374            false,
375        )
376        .unwrap();
377        let Filed::Standalone { task, claim } = filed else {
378            panic!("no loop is alive");
379        };
380        assert!(!task.urgent);
381        assert_eq!(task.review_of.as_deref(), Some("feat/x"));
382        assert_eq!(
383            task.overrides.as_ref().unwrap().merge.as_deref(),
384            Some("none")
385        );
386        assert!(
387            q.claim(&task.id).is_err(),
388            "a second process must not be able to claim a standalone run's task"
389        );
390        drop(claim);
391        assert!(q.claim(&task.id).is_ok(), "released with the claim");
392    }
393
394    #[test]
395    fn a_live_foreign_loop_gets_an_urgent_unclaimed_task() {
396        let dir = tempfile::tempdir().unwrap();
397        let q = Queue::at(dir.path().join("queue"));
398        heartbeat(dir.path(), 4242, 0);
399        let filed = file(
400            &q,
401            dir.path(),
402            Timestamp::now(),
403            1,
404            filing(dir.path()),
405            false,
406        )
407        .unwrap();
408        let Filed::Daemon { task, pid } = filed else {
409            panic!("a loop is alive");
410        };
411        assert_eq!(pid, Some(4242));
412        assert!(task.urgent);
413        let stored = q.get(&task.id).unwrap();
414        assert_eq!(stored.status, TaskStatus::Queued);
415        assert!(stored.runs.is_empty(), "nothing ran in this process");
416        assert!(q.claim(&task.id).is_ok(), "and this process holds no claim");
417    }
418
419    #[test]
420    fn a_stale_heartbeat_our_own_pid_and_a_dry_run_do_not_count_as_a_loop() {
421        let dir = tempfile::tempdir().unwrap();
422        let q = Queue::at(dir.path().join("queue"));
423        heartbeat(dir.path(), 4242, 3600);
424        assert!(matches!(
425            file(
426                &q,
427                dir.path(),
428                Timestamp::now(),
429                1,
430                filing(dir.path()),
431                false
432            )
433            .unwrap(),
434            Filed::Standalone { .. }
435        ));
436        heartbeat(dir.path(), 7, 0);
437        assert!(matches!(
438            file(
439                &q,
440                dir.path(),
441                Timestamp::now(),
442                7,
443                filing(dir.path()),
444                false
445            )
446            .unwrap(),
447            Filed::Standalone { .. }
448        ));
449        heartbeat(dir.path(), 4242, 0);
450        assert!(matches!(
451            file(
452                &q,
453                dir.path(),
454                Timestamp::now(),
455                1,
456                filing(dir.path()),
457                true
458            )
459            .unwrap(),
460            Filed::Standalone { .. }
461        ));
462    }
463
464    #[tokio::test]
465    async fn following_gives_up_on_a_task_nobody_is_serving() {
466        let dir = tempfile::tempdir().unwrap();
467        let q = Queue::at(dir.path().join("queue"));
468        let mut t = filing(dir.path()).into_task();
469        q.put(&mut t).unwrap();
470        let err = follow(
471            &q,
472            dir.path(),
473            &t.id,
474            Duration::from_millis(5),
475            None,
476            |_| {},
477            |_| {},
478        )
479        .await
480        .unwrap_err();
481        assert!(format!("{err}").contains("no magi loop"), "{err}");
482        assert!(q.claim(&t.id).is_ok(), "never taken over");
483    }
484
485    #[tokio::test]
486    async fn following_returns_when_the_loop_finishes_the_task_and_times_out_otherwise() {
487        let dir = tempfile::tempdir().unwrap();
488        let q = Queue::at(dir.path().join("queue"));
489        heartbeat(dir.path(), 4242, 0);
490        let mut t = filing(dir.path()).into_task();
491        q.put(&mut t).unwrap();
492        let waited = follow(
493            &q,
494            dir.path(),
495            &t.id,
496            Duration::from_millis(5),
497            Some(Duration::from_millis(30)),
498            |_| {},
499            |_| {},
500        )
501        .await
502        .unwrap();
503        assert!(!waited.finished, "the wait ran out, the loop still owns it");
504        let err = waited
505            .unfinished()
506            .expect("an unfinished wait is an error")
507            .to_string();
508        assert!(err.contains(t.short()), "{err}");
509        assert!(
510            err.contains("not finished") || err.contains("not a success"),
511            "{err}"
512        );
513        assert!(err.contains("magi task show"), "{err}");
514        assert_eq!(
515            q.get(&t.id).unwrap().status,
516            TaskStatus::Queued,
517            "the queue is untouched"
518        );
519        t.link_run("20260101-000000-abcd");
520        t.succeed();
521        q.put(&mut t).unwrap();
522        let mut runs = Vec::new();
523        let done = follow(
524            &q,
525            dir.path(),
526            &t.id,
527            Duration::from_millis(5),
528            Some(Duration::from_secs(5)),
529            |r| runs.push(r.to_owned()),
530            |_| {},
531        )
532        .await
533        .unwrap();
534        assert_eq!(done.task.status, TaskStatus::Done);
535        assert!(done.unfinished().is_none());
536        assert_eq!(runs, ["20260101-000000-abcd"]);
537    }
538
539    fn parent_in(q: &Queue, dir: &Path) -> Task {
540        let mut t = filing(dir).into_task();
541        q.put(&mut t).unwrap();
542        t
543    }
544
545    fn follow_up_state(dir: &Path) -> crate::run::RunState {
546        crate::run::RunState::new(
547            dir.to_path_buf(),
548            "main".to_owned(),
549            "abc1234".to_owned(),
550            "review the branch".to_owned(),
551            Config::default(),
552        )
553    }
554
555    #[test]
556    fn an_adopted_review_claims_starts_and_settles_a_runnable_task() {
557        let dir = tempfile::tempdir().unwrap();
558        let q = Queue::at(dir.path().join("queue"));
559        let parent = parent_in(&q, dir.path());
560        let state = follow_up_state(dir.path());
561        let adopted = adopt_in(q.clone(), &state, Some(&parent.id)).expect("adopted");
562        let running = q.get(&parent.id).unwrap();
563        assert_eq!(running.status, TaskStatus::Running);
564        assert_eq!(running.attempts, parent.attempts + 1);
565        assert_eq!(running.runs, std::slice::from_ref(&state.id));
566        assert!(q.claim(&parent.id).is_err(), "exclusive while it runs");
567        adopted.finish(&state, Err(anyhow::anyhow!("boom")));
568        let settled = q.get(&parent.id).unwrap();
569        assert_eq!(settled.status, TaskStatus::Held, "{settled:?}");
570        assert_eq!(
571            settled.runs,
572            std::slice::from_ref(&state.id),
573            "no duplicate run"
574        );
575        assert!(q.claim(&parent.id).is_ok(), "claim released");
576        assert_eq!(q.list().len(), 1, "no second owner was filed");
577    }
578
579    #[test]
580    fn a_task_claimed_elsewhere_only_gains_the_run() {
581        let dir = tempfile::tempdir().unwrap();
582        let q = Queue::at(dir.path().join("queue"));
583        let parent = parent_in(&q, dir.path());
584        let _theirs = q.claim(&parent.id).unwrap();
585        let state = follow_up_state(dir.path());
586        assert!(adopt_in(q.clone(), &state, Some(&parent.id)).is_none());
587        let after = q.get(&parent.id).unwrap();
588        assert_eq!(after.status, parent.status);
589        assert_eq!(after.attempts, parent.attempts);
590        assert_eq!(after.runs, std::slice::from_ref(&state.id));
591        assert_eq!(q.list().len(), 1, "no second owner was filed");
592    }
593
594    #[test]
595    fn a_done_task_only_gains_the_run() {
596        let dir = tempfile::tempdir().unwrap();
597        let q = Queue::at(dir.path().join("queue"));
598        let mut parent = parent_in(&q, dir.path());
599        parent.succeed();
600        q.put(&mut parent).unwrap();
601        let state = follow_up_state(dir.path());
602        assert!(adopt_in(q.clone(), &state, Some(&parent.id)).is_none());
603        let after = q.get(&parent.id).unwrap();
604        assert_eq!(after.status, TaskStatus::Done);
605        assert_eq!(after.attempts, parent.attempts);
606        assert_eq!(after.runs, std::slice::from_ref(&state.id));
607        assert!(q.claim(&parent.id).is_ok(), "claim released");
608    }
609
610    #[test]
611    fn a_held_task_is_claimed_and_settled_by_its_follow_up_review() {
612        let dir = tempfile::tempdir().unwrap();
613        let q = Queue::at(dir.path().join("queue"));
614        let mut parent = parent_in(&q, dir.path());
615        parent.hold_manual(Some("the run did not finish: stale".to_owned()));
616        q.put(&mut parent).unwrap();
617        let state = follow_up_state(dir.path());
618        let adopted = adopt_in(q.clone(), &state, Some(&parent.id)).expect("adopted");
619        assert_eq!(q.get(&parent.id).unwrap().status, TaskStatus::Running);
620        assert!(q.claim(&parent.id).is_err(), "claimed while it runs");
621        adopted.finish(&state, Err(anyhow::anyhow!("fresh failure")));
622        let after = q.get(&parent.id).unwrap();
623        assert_eq!(after.status, TaskStatus::Held);
624        assert!(
625            after
626                .hold_reason
627                .as_deref()
628                .unwrap()
629                .contains("fresh failure"),
630            "{after:?}"
631        );
632        assert_eq!(after.runs, std::slice::from_ref(&state.id));
633        assert!(q.claim(&parent.id).is_ok(), "claim released");
634    }
635}