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
136/// Watch `id` until the loop that owns it is done with it, calling `seen` with
137/// each run id the first time it appears on the task, and `progress` with a
138/// line whenever the newest run's status changes.
139///
140/// Finished is `done`, `held` or `blocked`: a `failed` or `queued` task is the
141/// loop's to retry. A task that stays unfinished while no loop is alive, or
142/// past `max_wait`, is an error - never taken over here, because the loop may
143/// merely be slow to heartbeat and a second driver would race it. Recovery is
144/// the loop's own claim reclaim and `Runner::resume`. Running out of
145/// `max_wait` is not an error: it returns with `finished: false`.
146pub async fn follow(
147    queue: &Queue,
148    home: &Path,
149    id: &str,
150    poll: Duration,
151    max_wait: Option<Duration>,
152    mut seen: impl FnMut(&str),
153    mut progress: impl FnMut(&str),
154) -> Result<Followed> {
155    let began = Instant::now();
156    let mut announced = 0usize;
157    let mut last_status = String::new();
158    loop {
159        let task = queue.get(id)?;
160        for run in task.runs.iter().skip(announced) {
161            seen(run);
162        }
163        announced = task.runs.len();
164        if let Some(run) = task.runs.last()
165            && crate::run::try_home().is_some()
166            && let Ok(state) = crate::run::RunState::load(run)
167        {
168            let line = format!("run {}: {}", state.short(), state.status.as_str());
169            if line != last_status {
170                progress(&line);
171                last_status = line;
172            }
173        }
174        if matches!(
175            task.status,
176            TaskStatus::Done | TaskStatus::Held | TaskStatus::Blocked
177        ) {
178            return Ok(Followed {
179                task,
180                finished: true,
181            });
182        }
183        let reading = crate::daemon::read_status(home);
184        if crate::daemon::foreign_loop(reading.as_ref(), Timestamp::now(), std::process::id())
185            .is_none()
186        {
187            bail!(
188                "no magi loop is serving the queue any more, and task {} is still {}; it was \
189                 not taken over here (the loop may only be slow to heartbeat). Start `magi \
190                 serve` to carry on, or follow it with `magi task show {}`",
191                task.short(),
192                task.status.as_str(),
193                task.short()
194            );
195        }
196        if max_wait.is_some_and(|m| began.elapsed() >= m) {
197            return Ok(Followed {
198                task,
199                finished: false,
200            });
201        }
202        tokio::time::sleep(poll).await;
203    }
204}
205
206/// A run opened inside this process (the follow-up review of `magi fix`) and
207/// the task that owns it. See [`adopt`].
208#[derive(Debug)]
209pub struct Adopted {
210    queue: Queue,
211    /// The task this process claimed and started for the run; `None` when the
212    /// claim could not be taken and the run was only linked.
213    task: Option<Task>,
214    claim: Option<Claim>,
215    quota_before: Vec<crate::run::QuotaLoss>,
216}
217
218/// Give a run this process is about to execute an owning task. A parent task
219/// that exists and is queued, failed or held is claimed, started with the
220/// run and settled by [`Adopted::finish`] exactly like an ownerless run's; if
221/// 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
222/// parent, a task is filed and claimed for the duration. A no-op (`None`) when
223/// no magi home is pinned.
224pub fn adopt(state: &crate::run::RunState, parent_task: Option<&str>) -> Option<Adopted> {
225    crate::run::try_home()?;
226    adopt_in(Queue::open(), state, parent_task)
227}
228
229fn adopt_in(
230    queue: Queue,
231    state: &crate::run::RunState,
232    parent_task: Option<&str>,
233) -> Option<Adopted> {
234    if let Some(parent) = parent_task
235        && let Ok(id) = queue.resolve_id(parent)
236    {
237        // The parent exists: never file a second owner for this run.
238        return match queue.claim(&id) {
239            Ok(claim) => {
240                let started = queue.get(&id).and_then(|mut task| {
241                    // A task the daemon could pick up, or one a failed
242                    // hand-started run left held (that is what a requested
243                    // follow-up review re-verifies), is this run's to settle.
244                    // A Done, blocked or running one is settled already or
245                    // owned by somebody else: restarting it would let a
246                    // follow-up review overwrite its outcome, so it only
247                    // gains the run.
248                    if !(task.status.runnable() || task.status == TaskStatus::Held) {
249                        return Ok(None);
250                    }
251                    task.start(state.id.clone());
252                    queue.put(&mut task)?;
253                    Ok(Some(task))
254                });
255                match started {
256                    Ok(None) => {
257                        drop(claim);
258                        let _ = queue.link_run(&id, &state.id);
259                        None
260                    }
261                    Ok(Some(task)) => Some(Adopted {
262                        queue,
263                        task: Some(task),
264                        claim: Some(claim),
265                        quota_before: state.quota.clone(),
266                    }),
267                    Err(e) => {
268                        tracing::warn!("could not start task {id} for run {}: {e:#}", state.id);
269                        drop(claim);
270                        let _ = queue.link_run(&id, &state.id);
271                        None
272                    }
273                }
274            }
275            Err(e) => {
276                // Another driver owns the task; do not settle over its result.
277                tracing::warn!("task {id} is claimed elsewhere ({e:#}); linking run only");
278                let _ = queue.link_run(&id, &state.id);
279                None
280            }
281        };
282    }
283    let mut task = Task::new(
284        crate::queue::title_from(&state.instruction, 72),
285        state.instruction.clone(),
286        state.repo.clone(),
287        Source::Human,
288    );
289    let claim = queue.claim(&task.id).ok()?;
290    task.start(state.id.clone());
291    queue.put(&mut task).ok()?;
292    Some(Adopted {
293        queue,
294        task: Some(task),
295        claim: Some(claim),
296        quota_before: state.quota.clone(),
297    })
298}
299
300impl Adopted {
301    /// Settle the task of an adopted ownerless run after it executed.
302    pub fn finish(mut self, state: &crate::run::RunState, result: Result<()>) {
303        if let Some(task) = &mut self.task {
304            crate::daemon::finish_attempt(
305                crate::daemon::Opts::default().max_attempts,
306                &self.queue,
307                task,
308                state,
309                &self.quota_before,
310                result,
311            );
312            crate::daemon::hold_if_runnable(&self.queue, task);
313        }
314        drop(self.claim.take());
315    }
316}
317
318#[cfg(test)]
319mod tests {
320    use super::*;
321    use crate::daemon::Status;
322
323    fn filing(repo: &Path) -> Filing {
324        Filing {
325            instruction: "review the branch".to_owned(),
326            title: "review x".to_owned(),
327            repo: repo.to_path_buf(),
328            source: Source::Human,
329            solo: false,
330            overrides: RunOverrides {
331                merge: Some("none".to_owned()),
332                ..RunOverrides::default()
333            },
334            review_of: Some("feat/x".to_owned()),
335        }
336    }
337
338    fn heartbeat(home: &Path, pid: u32, age_secs: i64) {
339        let mut status = Status::new();
340        status.pid = pid;
341        status.updated_at = Timestamp::now()
342            .checked_sub(jiff::SignedDuration::from_secs(age_secs))
343            .unwrap();
344        crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
345    }
346
347    #[test]
348    fn without_a_loop_the_task_is_claimed_before_anyone_else_can_take_it() {
349        let dir = tempfile::tempdir().unwrap();
350        let q = Queue::at(dir.path().join("queue"));
351        let filed = file(
352            &q,
353            dir.path(),
354            Timestamp::now(),
355            1,
356            filing(dir.path()),
357            false,
358        )
359        .unwrap();
360        let Filed::Standalone { task, claim } = filed else {
361            panic!("no loop is alive");
362        };
363        assert!(!task.urgent);
364        assert_eq!(task.review_of.as_deref(), Some("feat/x"));
365        assert_eq!(
366            task.overrides.as_ref().unwrap().merge.as_deref(),
367            Some("none")
368        );
369        assert!(
370            q.claim(&task.id).is_err(),
371            "a second process must not be able to claim a standalone run's task"
372        );
373        drop(claim);
374        assert!(q.claim(&task.id).is_ok(), "released with the claim");
375    }
376
377    #[test]
378    fn a_live_foreign_loop_gets_an_urgent_unclaimed_task() {
379        let dir = tempfile::tempdir().unwrap();
380        let q = Queue::at(dir.path().join("queue"));
381        heartbeat(dir.path(), 4242, 0);
382        let filed = file(
383            &q,
384            dir.path(),
385            Timestamp::now(),
386            1,
387            filing(dir.path()),
388            false,
389        )
390        .unwrap();
391        let Filed::Daemon { task, pid } = filed else {
392            panic!("a loop is alive");
393        };
394        assert_eq!(pid, Some(4242));
395        assert!(task.urgent);
396        let stored = q.get(&task.id).unwrap();
397        assert_eq!(stored.status, TaskStatus::Queued);
398        assert!(stored.runs.is_empty(), "nothing ran in this process");
399        assert!(q.claim(&task.id).is_ok(), "and this process holds no claim");
400    }
401
402    #[test]
403    fn a_stale_heartbeat_our_own_pid_and_a_dry_run_do_not_count_as_a_loop() {
404        let dir = tempfile::tempdir().unwrap();
405        let q = Queue::at(dir.path().join("queue"));
406        heartbeat(dir.path(), 4242, 3600);
407        assert!(matches!(
408            file(
409                &q,
410                dir.path(),
411                Timestamp::now(),
412                1,
413                filing(dir.path()),
414                false
415            )
416            .unwrap(),
417            Filed::Standalone { .. }
418        ));
419        heartbeat(dir.path(), 7, 0);
420        assert!(matches!(
421            file(
422                &q,
423                dir.path(),
424                Timestamp::now(),
425                7,
426                filing(dir.path()),
427                false
428            )
429            .unwrap(),
430            Filed::Standalone { .. }
431        ));
432        heartbeat(dir.path(), 4242, 0);
433        assert!(matches!(
434            file(
435                &q,
436                dir.path(),
437                Timestamp::now(),
438                1,
439                filing(dir.path()),
440                true
441            )
442            .unwrap(),
443            Filed::Standalone { .. }
444        ));
445    }
446
447    #[tokio::test]
448    async fn following_gives_up_on_a_task_nobody_is_serving() {
449        let dir = tempfile::tempdir().unwrap();
450        let q = Queue::at(dir.path().join("queue"));
451        let mut t = filing(dir.path()).into_task();
452        q.put(&mut t).unwrap();
453        let err = follow(
454            &q,
455            dir.path(),
456            &t.id,
457            Duration::from_millis(5),
458            None,
459            |_| {},
460            |_| {},
461        )
462        .await
463        .unwrap_err();
464        assert!(format!("{err}").contains("no magi loop"), "{err}");
465        assert!(q.claim(&t.id).is_ok(), "never taken over");
466    }
467
468    #[tokio::test]
469    async fn following_returns_when_the_loop_finishes_the_task_and_times_out_otherwise() {
470        let dir = tempfile::tempdir().unwrap();
471        let q = Queue::at(dir.path().join("queue"));
472        heartbeat(dir.path(), 4242, 0);
473        let mut t = filing(dir.path()).into_task();
474        q.put(&mut t).unwrap();
475        let waited = follow(
476            &q,
477            dir.path(),
478            &t.id,
479            Duration::from_millis(5),
480            Some(Duration::from_millis(30)),
481            |_| {},
482            |_| {},
483        )
484        .await
485        .unwrap();
486        assert!(!waited.finished, "the wait ran out, the loop still owns it");
487        t.link_run("20260101-000000-abcd");
488        t.succeed();
489        q.put(&mut t).unwrap();
490        let mut runs = Vec::new();
491        let done = follow(
492            &q,
493            dir.path(),
494            &t.id,
495            Duration::from_millis(5),
496            Some(Duration::from_secs(5)),
497            |r| runs.push(r.to_owned()),
498            |_| {},
499        )
500        .await
501        .unwrap();
502        assert_eq!(done.task.status, TaskStatus::Done);
503        assert_eq!(runs, ["20260101-000000-abcd"]);
504    }
505
506    fn parent_in(q: &Queue, dir: &Path) -> Task {
507        let mut t = filing(dir).into_task();
508        q.put(&mut t).unwrap();
509        t
510    }
511
512    fn follow_up_state(dir: &Path) -> crate::run::RunState {
513        crate::run::RunState::new(
514            dir.to_path_buf(),
515            "main".to_owned(),
516            "abc1234".to_owned(),
517            "review the branch".to_owned(),
518            Config::default(),
519        )
520    }
521
522    #[test]
523    fn an_adopted_review_claims_starts_and_settles_a_runnable_task() {
524        let dir = tempfile::tempdir().unwrap();
525        let q = Queue::at(dir.path().join("queue"));
526        let parent = parent_in(&q, dir.path());
527        let state = follow_up_state(dir.path());
528        let adopted = adopt_in(q.clone(), &state, Some(&parent.id)).expect("adopted");
529        let running = q.get(&parent.id).unwrap();
530        assert_eq!(running.status, TaskStatus::Running);
531        assert_eq!(running.attempts, parent.attempts + 1);
532        assert_eq!(running.runs, std::slice::from_ref(&state.id));
533        assert!(q.claim(&parent.id).is_err(), "exclusive while it runs");
534        adopted.finish(&state, Err(anyhow::anyhow!("boom")));
535        let settled = q.get(&parent.id).unwrap();
536        assert_eq!(settled.status, TaskStatus::Held, "{settled:?}");
537        assert_eq!(
538            settled.runs,
539            std::slice::from_ref(&state.id),
540            "no duplicate run"
541        );
542        assert!(q.claim(&parent.id).is_ok(), "claim released");
543        assert_eq!(q.list().len(), 1, "no second owner was filed");
544    }
545
546    #[test]
547    fn a_task_claimed_elsewhere_only_gains_the_run() {
548        let dir = tempfile::tempdir().unwrap();
549        let q = Queue::at(dir.path().join("queue"));
550        let parent = parent_in(&q, dir.path());
551        let _theirs = q.claim(&parent.id).unwrap();
552        let state = follow_up_state(dir.path());
553        assert!(adopt_in(q.clone(), &state, Some(&parent.id)).is_none());
554        let after = q.get(&parent.id).unwrap();
555        assert_eq!(after.status, parent.status);
556        assert_eq!(after.attempts, parent.attempts);
557        assert_eq!(after.runs, std::slice::from_ref(&state.id));
558        assert_eq!(q.list().len(), 1, "no second owner was filed");
559    }
560
561    #[test]
562    fn a_done_task_only_gains_the_run() {
563        let dir = tempfile::tempdir().unwrap();
564        let q = Queue::at(dir.path().join("queue"));
565        let mut parent = parent_in(&q, dir.path());
566        parent.succeed();
567        q.put(&mut parent).unwrap();
568        let state = follow_up_state(dir.path());
569        assert!(adopt_in(q.clone(), &state, Some(&parent.id)).is_none());
570        let after = q.get(&parent.id).unwrap();
571        assert_eq!(after.status, TaskStatus::Done);
572        assert_eq!(after.attempts, parent.attempts);
573        assert_eq!(after.runs, std::slice::from_ref(&state.id));
574        assert!(q.claim(&parent.id).is_ok(), "claim released");
575    }
576
577    #[test]
578    fn a_held_task_is_claimed_and_settled_by_its_follow_up_review() {
579        let dir = tempfile::tempdir().unwrap();
580        let q = Queue::at(dir.path().join("queue"));
581        let mut parent = parent_in(&q, dir.path());
582        parent.hold_manual(Some("the run did not finish: stale".to_owned()));
583        q.put(&mut parent).unwrap();
584        let state = follow_up_state(dir.path());
585        let adopted = adopt_in(q.clone(), &state, Some(&parent.id)).expect("adopted");
586        assert_eq!(q.get(&parent.id).unwrap().status, TaskStatus::Running);
587        assert!(q.claim(&parent.id).is_err(), "claimed while it runs");
588        adopted.finish(&state, Err(anyhow::anyhow!("fresh failure")));
589        let after = q.get(&parent.id).unwrap();
590        assert_eq!(after.status, TaskStatus::Held);
591        assert!(
592            after
593                .hold_reason
594                .as_deref()
595                .unwrap()
596                .contains("fresh failure"),
597            "{after:?}"
598        );
599        assert_eq!(after.runs, std::slice::from_ref(&state.id));
600        assert!(q.claim(&parent.id).is_ok(), "claim released");
601    }
602}