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 created for an ownerless run; `None` when an existing task
212    /// only gained the run.
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 just gains the run in `runs`; with none, a task is filed and
220/// claimed for the duration, and [`Adopted::finish`] settles it as the daemon
221/// would. A no-op (`None`) when no magi home is pinned.
222pub fn adopt(state: &crate::run::RunState, parent_task: Option<&str>) -> Option<Adopted> {
223    crate::run::try_home()?;
224    let queue = Queue::open();
225    if let Some(parent) = parent_task
226        && queue.link_run(parent, &state.id).is_ok()
227    {
228        return Some(Adopted {
229            queue,
230            task: None,
231            claim: None,
232            quota_before: Vec::new(),
233        });
234    }
235    let mut task = Task::new(
236        crate::queue::title_from(&state.instruction, 72),
237        state.instruction.clone(),
238        state.repo.clone(),
239        Source::Human,
240    );
241    let claim = queue.claim(&task.id).ok()?;
242    task.start(state.id.clone());
243    queue.put(&mut task).ok()?;
244    Some(Adopted {
245        queue,
246        task: Some(task),
247        claim: Some(claim),
248        quota_before: state.quota.clone(),
249    })
250}
251
252impl Adopted {
253    /// Settle the task of an adopted ownerless run after it executed.
254    pub fn finish(mut self, state: &crate::run::RunState, result: Result<()>) {
255        if let Some(task) = &mut self.task {
256            crate::daemon::finish_attempt(
257                crate::daemon::Opts::default().max_attempts,
258                &self.queue,
259                task,
260                state,
261                &self.quota_before,
262                result,
263            );
264            crate::daemon::hold_if_runnable(&self.queue, task);
265        }
266        drop(self.claim.take());
267    }
268}
269
270#[cfg(test)]
271mod tests {
272    use super::*;
273    use crate::daemon::Status;
274
275    fn filing(repo: &Path) -> Filing {
276        Filing {
277            instruction: "review the branch".to_owned(),
278            title: "review x".to_owned(),
279            repo: repo.to_path_buf(),
280            source: Source::Human,
281            solo: false,
282            overrides: RunOverrides {
283                merge: Some("none".to_owned()),
284                ..RunOverrides::default()
285            },
286            review_of: Some("feat/x".to_owned()),
287        }
288    }
289
290    fn heartbeat(home: &Path, pid: u32, age_secs: i64) {
291        let mut status = Status::new();
292        status.pid = pid;
293        status.updated_at = Timestamp::now()
294            .checked_sub(jiff::SignedDuration::from_secs(age_secs))
295            .unwrap();
296        crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
297    }
298
299    #[test]
300    fn without_a_loop_the_task_is_claimed_before_anyone_else_can_take_it() {
301        let dir = tempfile::tempdir().unwrap();
302        let q = Queue::at(dir.path().join("queue"));
303        let filed = file(
304            &q,
305            dir.path(),
306            Timestamp::now(),
307            1,
308            filing(dir.path()),
309            false,
310        )
311        .unwrap();
312        let Filed::Standalone { task, claim } = filed else {
313            panic!("no loop is alive");
314        };
315        assert!(!task.urgent);
316        assert_eq!(task.review_of.as_deref(), Some("feat/x"));
317        assert_eq!(
318            task.overrides.as_ref().unwrap().merge.as_deref(),
319            Some("none")
320        );
321        assert!(
322            q.claim(&task.id).is_err(),
323            "a second process must not be able to claim a standalone run's task"
324        );
325        drop(claim);
326        assert!(q.claim(&task.id).is_ok(), "released with the claim");
327    }
328
329    #[test]
330    fn a_live_foreign_loop_gets_an_urgent_unclaimed_task() {
331        let dir = tempfile::tempdir().unwrap();
332        let q = Queue::at(dir.path().join("queue"));
333        heartbeat(dir.path(), 4242, 0);
334        let filed = file(
335            &q,
336            dir.path(),
337            Timestamp::now(),
338            1,
339            filing(dir.path()),
340            false,
341        )
342        .unwrap();
343        let Filed::Daemon { task, pid } = filed else {
344            panic!("a loop is alive");
345        };
346        assert_eq!(pid, Some(4242));
347        assert!(task.urgent);
348        let stored = q.get(&task.id).unwrap();
349        assert_eq!(stored.status, TaskStatus::Queued);
350        assert!(stored.runs.is_empty(), "nothing ran in this process");
351        assert!(q.claim(&task.id).is_ok(), "and this process holds no claim");
352    }
353
354    #[test]
355    fn a_stale_heartbeat_our_own_pid_and_a_dry_run_do_not_count_as_a_loop() {
356        let dir = tempfile::tempdir().unwrap();
357        let q = Queue::at(dir.path().join("queue"));
358        heartbeat(dir.path(), 4242, 3600);
359        assert!(matches!(
360            file(
361                &q,
362                dir.path(),
363                Timestamp::now(),
364                1,
365                filing(dir.path()),
366                false
367            )
368            .unwrap(),
369            Filed::Standalone { .. }
370        ));
371        heartbeat(dir.path(), 7, 0);
372        assert!(matches!(
373            file(
374                &q,
375                dir.path(),
376                Timestamp::now(),
377                7,
378                filing(dir.path()),
379                false
380            )
381            .unwrap(),
382            Filed::Standalone { .. }
383        ));
384        heartbeat(dir.path(), 4242, 0);
385        assert!(matches!(
386            file(
387                &q,
388                dir.path(),
389                Timestamp::now(),
390                1,
391                filing(dir.path()),
392                true
393            )
394            .unwrap(),
395            Filed::Standalone { .. }
396        ));
397    }
398
399    #[tokio::test]
400    async fn following_gives_up_on_a_task_nobody_is_serving() {
401        let dir = tempfile::tempdir().unwrap();
402        let q = Queue::at(dir.path().join("queue"));
403        let mut t = filing(dir.path()).into_task();
404        q.put(&mut t).unwrap();
405        let err = follow(
406            &q,
407            dir.path(),
408            &t.id,
409            Duration::from_millis(5),
410            None,
411            |_| {},
412            |_| {},
413        )
414        .await
415        .unwrap_err();
416        assert!(format!("{err}").contains("no magi loop"), "{err}");
417        assert!(q.claim(&t.id).is_ok(), "never taken over");
418    }
419
420    #[tokio::test]
421    async fn following_returns_when_the_loop_finishes_the_task_and_times_out_otherwise() {
422        let dir = tempfile::tempdir().unwrap();
423        let q = Queue::at(dir.path().join("queue"));
424        heartbeat(dir.path(), 4242, 0);
425        let mut t = filing(dir.path()).into_task();
426        q.put(&mut t).unwrap();
427        let waited = follow(
428            &q,
429            dir.path(),
430            &t.id,
431            Duration::from_millis(5),
432            Some(Duration::from_millis(30)),
433            |_| {},
434            |_| {},
435        )
436        .await
437        .unwrap();
438        assert!(!waited.finished, "the wait ran out, the loop still owns it");
439        t.link_run("20260101-000000-abcd");
440        t.succeed();
441        q.put(&mut t).unwrap();
442        let mut runs = Vec::new();
443        let done = follow(
444            &q,
445            dir.path(),
446            &t.id,
447            Duration::from_millis(5),
448            Some(Duration::from_secs(5)),
449            |r| runs.push(r.to_owned()),
450            |_| {},
451        )
452        .await
453        .unwrap();
454        assert_eq!(done.task.status, TaskStatus::Done);
455        assert_eq!(runs, ["20260101-000000-abcd"]);
456    }
457}