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. `Ok(None)` when no
240/// magi home is pinned, or when the run was only linked to its parent. An
241/// error means the run has no owner and must not be executed: filing or
242/// claiming the new task failed.
243pub fn adopt(
244    state: &crate::run::RunState,
245    parent_task: Option<&str>,
246) -> anyhow::Result<Option<Adopted>> {
247    if crate::run::try_home().is_none() {
248        return Ok(None);
249    }
250    adopt_in(Queue::open(), state, parent_task)
251}
252
253fn adopt_in(
254    queue: Queue,
255    state: &crate::run::RunState,
256    parent_task: Option<&str>,
257) -> anyhow::Result<Option<Adopted>> {
258    use anyhow::Context as _;
259    if let Some(parent) = parent_task
260        && let Ok(id) = queue.resolve_id(parent)
261    {
262        // The parent exists: never file a second owner for this run.
263        return Ok(match queue.claim(&id) {
264            Ok(claim) => {
265                let started = queue.get(&id).and_then(|mut task| {
266                    // A task the daemon could pick up, or one a failed
267                    // hand-started run left held (that is what a requested
268                    // follow-up review re-verifies), is this run's to settle.
269                    // A Done, blocked or running one is settled already or
270                    // owned by somebody else: restarting it would let a
271                    // follow-up review overwrite its outcome, so it only
272                    // gains the run.
273                    if !(task.status.runnable() || task.status == TaskStatus::Held) {
274                        return Ok(None);
275                    }
276                    task.start(state.id.clone());
277                    queue.put(&mut task)?;
278                    Ok(Some(task))
279                });
280                match started {
281                    Ok(None) => {
282                        drop(claim);
283                        let _ = queue.link_run(&id, &state.id);
284                        None
285                    }
286                    Ok(Some(task)) => Some(Adopted {
287                        queue,
288                        task: Some(task),
289                        claim: Some(claim),
290                        quota_before: state.quota.clone(),
291                    }),
292                    Err(e) => {
293                        tracing::warn!("could not start task {id} for run {}: {e:#}", state.id);
294                        drop(claim);
295                        let _ = queue.link_run(&id, &state.id);
296                        None
297                    }
298                }
299            }
300            Err(e) => {
301                // Another driver owns the task; do not settle over its result.
302                tracing::warn!("task {id} is claimed elsewhere ({e:#}); linking run only");
303                let _ = queue.link_run(&id, &state.id);
304                None
305            }
306        });
307    }
308    let mut task = Task::new(
309        crate::queue::title_from(&state.instruction, 72),
310        state.instruction.clone(),
311        state.repo.clone(),
312        Source::Human,
313    );
314    let claim = queue
315        .claim(&task.id)
316        .with_context(|| format!("could not claim a new task for run {}", state.id))?;
317    task.start(state.id.clone());
318    queue
319        .put(&mut task)
320        .with_context(|| format!("could not file a new task for run {}", state.id))?;
321    Ok(Some(Adopted {
322        queue,
323        task: Some(task),
324        claim: Some(claim),
325        quota_before: state.quota.clone(),
326    }))
327}
328
329impl Adopted {
330    /// Settle the task of an adopted ownerless run after it executed.
331    pub fn finish(mut self, state: &crate::run::RunState, result: Result<()>) {
332        if let Some(task) = &mut self.task {
333            crate::daemon::finish_attempt(
334                crate::daemon::Opts::default().max_attempts,
335                &self.queue,
336                task,
337                state,
338                &self.quota_before,
339                result,
340            );
341            crate::daemon::hold_if_runnable(&self.queue, task);
342        }
343        drop(self.claim.take());
344    }
345}
346
347#[cfg(test)]
348mod tests {
349    use super::*;
350    use crate::daemon::Status;
351
352    fn filing(repo: &Path) -> Filing {
353        Filing {
354            instruction: "review the branch".to_owned(),
355            title: "review x".to_owned(),
356            repo: repo.to_path_buf(),
357            source: Source::Human,
358            solo: false,
359            overrides: RunOverrides {
360                merge: Some("none".to_owned()),
361                ..RunOverrides::default()
362            },
363            review_of: Some("feat/x".to_owned()),
364        }
365    }
366
367    fn heartbeat(home: &Path, pid: u32, age_secs: i64) {
368        let mut status = Status::new();
369        status.pid = pid;
370        status.updated_at = Timestamp::now()
371            .checked_sub(jiff::SignedDuration::from_secs(age_secs))
372            .unwrap();
373        crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
374    }
375
376    #[test]
377    fn without_a_loop_the_task_is_claimed_before_anyone_else_can_take_it() {
378        let dir = tempfile::tempdir().unwrap();
379        let q = Queue::at(dir.path().join("queue"));
380        let filed = file(
381            &q,
382            dir.path(),
383            Timestamp::now(),
384            1,
385            filing(dir.path()),
386            false,
387        )
388        .unwrap();
389        let Filed::Standalone { task, claim } = filed else {
390            panic!("no loop is alive");
391        };
392        assert!(!task.urgent);
393        assert_eq!(task.review_of.as_deref(), Some("feat/x"));
394        assert_eq!(
395            task.overrides.as_ref().unwrap().merge.as_deref(),
396            Some("none")
397        );
398        assert!(
399            q.claim(&task.id).is_err(),
400            "a second process must not be able to claim a standalone run's task"
401        );
402        drop(claim);
403        assert!(q.claim(&task.id).is_ok(), "released with the claim");
404    }
405
406    #[test]
407    fn a_live_foreign_loop_gets_an_urgent_unclaimed_task() {
408        let dir = tempfile::tempdir().unwrap();
409        let q = Queue::at(dir.path().join("queue"));
410        heartbeat(dir.path(), 4242, 0);
411        let filed = file(
412            &q,
413            dir.path(),
414            Timestamp::now(),
415            1,
416            filing(dir.path()),
417            false,
418        )
419        .unwrap();
420        let Filed::Daemon { task, pid } = filed else {
421            panic!("a loop is alive");
422        };
423        assert_eq!(pid, Some(4242));
424        assert!(task.urgent);
425        let stored = q.get(&task.id).unwrap();
426        assert_eq!(stored.status, TaskStatus::Queued);
427        assert!(stored.runs.is_empty(), "nothing ran in this process");
428        assert!(q.claim(&task.id).is_ok(), "and this process holds no claim");
429    }
430
431    #[test]
432    fn a_stale_heartbeat_our_own_pid_and_a_dry_run_do_not_count_as_a_loop() {
433        let dir = tempfile::tempdir().unwrap();
434        let q = Queue::at(dir.path().join("queue"));
435        heartbeat(dir.path(), 4242, 3600);
436        assert!(matches!(
437            file(
438                &q,
439                dir.path(),
440                Timestamp::now(),
441                1,
442                filing(dir.path()),
443                false
444            )
445            .unwrap(),
446            Filed::Standalone { .. }
447        ));
448        heartbeat(dir.path(), 7, 0);
449        assert!(matches!(
450            file(
451                &q,
452                dir.path(),
453                Timestamp::now(),
454                7,
455                filing(dir.path()),
456                false
457            )
458            .unwrap(),
459            Filed::Standalone { .. }
460        ));
461        heartbeat(dir.path(), 4242, 0);
462        assert!(matches!(
463            file(
464                &q,
465                dir.path(),
466                Timestamp::now(),
467                1,
468                filing(dir.path()),
469                true
470            )
471            .unwrap(),
472            Filed::Standalone { .. }
473        ));
474    }
475
476    #[tokio::test]
477    async fn following_gives_up_on_a_task_nobody_is_serving() {
478        let dir = tempfile::tempdir().unwrap();
479        let q = Queue::at(dir.path().join("queue"));
480        let mut t = filing(dir.path()).into_task();
481        q.put(&mut t).unwrap();
482        let err = follow(
483            &q,
484            dir.path(),
485            &t.id,
486            Duration::from_millis(5),
487            None,
488            |_| {},
489            |_| {},
490        )
491        .await
492        .unwrap_err();
493        assert!(format!("{err}").contains("no magi loop"), "{err}");
494        assert!(q.claim(&t.id).is_ok(), "never taken over");
495    }
496
497    #[tokio::test]
498    async fn following_returns_when_the_loop_finishes_the_task_and_times_out_otherwise() {
499        let dir = tempfile::tempdir().unwrap();
500        let q = Queue::at(dir.path().join("queue"));
501        heartbeat(dir.path(), 4242, 0);
502        let mut t = filing(dir.path()).into_task();
503        q.put(&mut t).unwrap();
504        let waited = follow(
505            &q,
506            dir.path(),
507            &t.id,
508            Duration::from_millis(5),
509            Some(Duration::from_millis(30)),
510            |_| {},
511            |_| {},
512        )
513        .await
514        .unwrap();
515        assert!(!waited.finished, "the wait ran out, the loop still owns it");
516        let err = waited
517            .unfinished()
518            .expect("an unfinished wait is an error")
519            .to_string();
520        assert!(err.contains(t.short()), "{err}");
521        assert!(
522            err.contains("not finished") || err.contains("not a success"),
523            "{err}"
524        );
525        assert!(err.contains("magi task show"), "{err}");
526        assert_eq!(
527            q.get(&t.id).unwrap().status,
528            TaskStatus::Queued,
529            "the queue is untouched"
530        );
531        t.link_run("20260101-000000-abcd");
532        t.succeed();
533        q.put(&mut t).unwrap();
534        let mut runs = Vec::new();
535        let done = follow(
536            &q,
537            dir.path(),
538            &t.id,
539            Duration::from_millis(5),
540            Some(Duration::from_secs(5)),
541            |r| runs.push(r.to_owned()),
542            |_| {},
543        )
544        .await
545        .unwrap();
546        assert_eq!(done.task.status, TaskStatus::Done);
547        assert!(done.unfinished().is_none());
548        assert_eq!(runs, ["20260101-000000-abcd"]);
549    }
550
551    fn parent_in(q: &Queue, dir: &Path) -> Task {
552        let mut t = filing(dir).into_task();
553        q.put(&mut t).unwrap();
554        t
555    }
556
557    fn follow_up_state(dir: &Path) -> crate::run::RunState {
558        crate::run::RunState::new(
559            dir.to_path_buf(),
560            "main".to_owned(),
561            "abc1234".to_owned(),
562            "review the branch".to_owned(),
563            Config::default(),
564        )
565    }
566
567    #[test]
568    fn an_adopted_review_claims_starts_and_settles_a_runnable_task() {
569        let dir = tempfile::tempdir().unwrap();
570        let q = Queue::at(dir.path().join("queue"));
571        let parent = parent_in(&q, dir.path());
572        let state = follow_up_state(dir.path());
573        let adopted = adopt_in(q.clone(), &state, Some(&parent.id))
574            .unwrap()
575            .expect("adopted");
576        let running = q.get(&parent.id).unwrap();
577        assert_eq!(running.status, TaskStatus::Running);
578        assert_eq!(running.attempts, parent.attempts + 1);
579        assert_eq!(running.runs, std::slice::from_ref(&state.id));
580        assert!(q.claim(&parent.id).is_err(), "exclusive while it runs");
581        adopted.finish(&state, Err(anyhow::anyhow!("boom")));
582        let settled = q.get(&parent.id).unwrap();
583        assert_eq!(settled.status, TaskStatus::Held, "{settled:?}");
584        assert_eq!(
585            settled.runs,
586            std::slice::from_ref(&state.id),
587            "no duplicate run"
588        );
589        assert!(q.claim(&parent.id).is_ok(), "claim released");
590        assert_eq!(q.list().len(), 1, "no second owner was filed");
591    }
592
593    #[test]
594    fn an_ownerless_run_whose_task_cannot_be_filed_is_an_error() {
595        let dir = tempfile::tempdir().unwrap();
596        // The queue root sits under a regular file, so claiming cannot create it.
597        let blocker = dir.path().join("blocker");
598        std::fs::write(&blocker, "x").unwrap();
599        let q = Queue::at(blocker.join("q"));
600        let state = follow_up_state(dir.path());
601        let err = adopt_in(q, &state, None).expect_err("must not be swallowed");
602        assert!(format!("{err:#}").contains(&state.id), "{err:#}");
603    }
604
605    #[test]
606    fn a_task_claimed_elsewhere_only_gains_the_run() {
607        let dir = tempfile::tempdir().unwrap();
608        let q = Queue::at(dir.path().join("queue"));
609        let parent = parent_in(&q, dir.path());
610        let _theirs = q.claim(&parent.id).unwrap();
611        let state = follow_up_state(dir.path());
612        assert!(
613            adopt_in(q.clone(), &state, Some(&parent.id))
614                .unwrap()
615                .is_none()
616        );
617        let after = q.get(&parent.id).unwrap();
618        assert_eq!(after.status, parent.status);
619        assert_eq!(after.attempts, parent.attempts);
620        assert_eq!(after.runs, std::slice::from_ref(&state.id));
621        assert_eq!(q.list().len(), 1, "no second owner was filed");
622    }
623
624    #[test]
625    fn a_done_task_only_gains_the_run() {
626        let dir = tempfile::tempdir().unwrap();
627        let q = Queue::at(dir.path().join("queue"));
628        let mut parent = parent_in(&q, dir.path());
629        parent.succeed();
630        q.put(&mut parent).unwrap();
631        let state = follow_up_state(dir.path());
632        assert!(
633            adopt_in(q.clone(), &state, Some(&parent.id))
634                .unwrap()
635                .is_none()
636        );
637        let after = q.get(&parent.id).unwrap();
638        assert_eq!(after.status, TaskStatus::Done);
639        assert_eq!(after.attempts, parent.attempts);
640        assert_eq!(after.runs, std::slice::from_ref(&state.id));
641        assert!(q.claim(&parent.id).is_ok(), "claim released");
642    }
643
644    #[test]
645    fn a_held_task_is_claimed_and_settled_by_its_follow_up_review() {
646        let dir = tempfile::tempdir().unwrap();
647        let q = Queue::at(dir.path().join("queue"));
648        let mut parent = parent_in(&q, dir.path());
649        parent.hold_manual(Some("the run did not finish: stale".to_owned()));
650        q.put(&mut parent).unwrap();
651        let state = follow_up_state(dir.path());
652        let adopted = adopt_in(q.clone(), &state, Some(&parent.id))
653            .unwrap()
654            .expect("adopted");
655        assert_eq!(q.get(&parent.id).unwrap().status, TaskStatus::Running);
656        assert!(q.claim(&parent.id).is_err(), "claimed while it runs");
657        adopted.finish(&state, Err(anyhow::anyhow!("fresh failure")));
658        let after = q.get(&parent.id).unwrap();
659        assert_eq!(after.status, TaskStatus::Held);
660        assert!(
661            after
662                .hold_reason
663                .as_deref()
664                .unwrap()
665                .contains("fresh failure"),
666            "{after:?}"
667        );
668        assert_eq!(after.runs, std::slice::from_ref(&state.id));
669        assert!(q.claim(&parent.id).is_ok(), "claim released");
670    }
671}