magi-cli 0.80.0

Blind multi-agent implementation competition: N agents implement, M judges rank blind, deliberate, vote privately, winner survives double review + E2E gate
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
//! Runs started by hand are tasks.
//!
//! `magi run` and `magi review` used to mint a run no task knew about: no
//! detail page, no attempt history, invisible to duplicate detection, and an
//! orphan the moment anything cleaned up after it. They now file a task first,
//! through the queue like `magi task add`, and then either hand it to a live
//! loop or execute it here through the daemon's own attempt
//! ([`crate::daemon::run_claimed`]). There is no second dispatcher in this
//! module: it decides *who* runs the task and watches it, nothing more.
//!
//! * A live loop in another process (a fresh `daemon.json` heartbeat, see
//!   [`crate::daemon::foreign_loop`]) gets the task as `urgent`, runs it in
//!   its own urgent slot, and the caller [`follow`]s it. The caller never
//!   claims or executes it: two processes must not drive one task.
//! * Otherwise the caller claims the task **before it is written**, so a loop
//!   that starts mid-run finds the claim and leaves it alone, and executes it.

use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};

use anyhow::{Context, Result, bail};
use jiff::Timestamp;

use crate::config::Config;
use crate::queue::{Claim, Queue, RunOverrides, Source, Task, TaskStatus};

/// What to file. Everything `magi run` / `magi review` take that changes what
/// the run does rides on the task, so a loop that executes it behaves the same
/// as this process would have.
#[derive(Debug, Clone)]
pub struct Filing {
    /// The task text; for a review, a line saying which branch.
    pub instruction: String,
    /// The task's title.
    pub title: String,
    /// Absolute repository path.
    pub repo: PathBuf,
    /// Who asked: the operator, or the agent seat named by `MAGI_RUN` / `MAGI_NODE`.
    pub source: Source,
    /// `--solo`.
    pub solo: bool,
    /// Command-line choices carried onto the task.
    pub overrides: RunOverrides,
    /// `magi review <branch>`: the branch to review.
    pub review_of: Option<String>,
}

impl Filing {
    fn into_task(self) -> Task {
        let mut task = Task::new(self.title, self.instruction, self.repo, self.source);
        task.solo = self.solo;
        task.review_of = self.review_of;
        task.overrides = Some(self.overrides);
        task
    }
}

/// Who will execute a filed task.
#[derive(Debug)]
pub enum Filed {
    /// A loop in another process owns the queue; the task is `urgent` and
    /// unclaimed. `pid` is the loop's, when it published one.
    Daemon {
        /// The filed task.
        task: Task,
        /// The loop's pid.
        pid: Option<u32>,
    },
    /// Nobody else is serving: the task is claimed for this process, and the
    /// claim lives as long as this value.
    Standalone {
        /// The filed task.
        task: Task,
        /// Held until dropped, which is what keeps a later loop away.
        claim: Claim,
    },
}

/// File `filing` and decide who runs it. `home` is the magi home whose
/// `daemon.json` is judged; `own_pid` is this process (a loop of our own is
/// not "another" one). `force_local` (`--dry-run`) never hands over: a dry run
/// spends no agent call and a loop would not honour that.
pub fn file(
    queue: &Queue,
    home: &Path,
    now: Timestamp,
    own_pid: u32,
    filing: Filing,
    force_local: bool,
) -> Result<Filed> {
    let mut task = filing.into_task();
    let owner = if force_local {
        None
    } else {
        crate::daemon::foreign_loop(crate::daemon::read_status(home).as_ref(), now, own_pid)
    };
    match owner {
        Some(pid) => {
            task.urgent = true;
            queue.put(&mut task).context("file the task")?;
            Ok(Filed::Daemon { task, pid })
        }
        None => {
            // Claimed first: between `put` and the claim a loop could take it.
            let claim = queue.claim(&task.id)?;
            queue.put(&mut task).context("file the task")?;
            Ok(Filed::Standalone { task, claim })
        }
    }
}

/// The config a task's run is built with: the task's `--config` layer, then
/// its overrides and `solo`. What the loop's `attempt` does, for the dry run.
pub fn config_for(task: &Task, repo: &Path) -> Result<Config> {
    let path = task.overrides.as_ref().and_then(|o| o.config.as_deref());
    let (mut cfg, _) = Config::discover(repo, path)?;
    if let Some(o) = &task.overrides {
        o.apply(&mut cfg);
    }
    if task.solo {
        cfg.graph.candidates = 1;
    }
    Ok(cfg)
}

/// How a followed task ended up.
#[derive(Debug)]
pub struct Followed {
    /// The task as last read.
    pub task: Task,
    /// False when `max_wait` ran out first: the loop still owns the task and
    /// nothing is known of its outcome yet.
    pub finished: bool,
}

impl Followed {
    /// The error a caller must exit with when the wait ran out before the
    /// task finished: that is not a success, and the outcome is not known.
    /// `None` once the task finished (whatever its outcome).
    pub fn unfinished(&self) -> Option<anyhow::Error> {
        if self.finished {
            return None;
        }
        Some(anyhow::anyhow!(
            "task {id} has not finished: it is still running under the loop. \
             This is not a success; the outcome is not known yet. Read it with \
             `magi task show {id}`. The loop keeps the task, so do not run it again.",
            id = self.task.short()
        ))
    }
}

/// Watch `id` until the loop that owns it is done with it, calling `seen` with
/// each run id the first time it appears on the task, and `progress` with a
/// line whenever the newest run's status changes.
///
/// Finished is `done`, `held` or `blocked`: a `failed` or `queued` task is the
/// loop's to retry. A task that stays unfinished while no loop is alive, or
/// past `max_wait`, is an error - never taken over here, because the loop may
/// merely be slow to heartbeat and a second driver would race it. Recovery is
/// the loop's own claim reclaim and `Runner::resume`. Running out of
/// `max_wait` is not an error: it returns with `finished: false`.
pub async fn follow(
    queue: &Queue,
    home: &Path,
    id: &str,
    poll: Duration,
    max_wait: Option<Duration>,
    mut seen: impl FnMut(&str),
    mut progress: impl FnMut(&str),
) -> Result<Followed> {
    let began = Instant::now();
    let mut announced = 0usize;
    let mut last_status = String::new();
    loop {
        let task = queue.get(id)?;
        for run in task.runs.iter().skip(announced) {
            seen(run);
        }
        announced = task.runs.len();
        if let Some(run) = task.runs.last()
            && crate::run::try_home().is_some()
            && let Ok(state) = crate::run::RunState::load(run)
        {
            let line = format!("run {}: {}", state.short(), state.status.as_str());
            if line != last_status {
                progress(&line);
                last_status = line;
            }
        }
        if matches!(
            task.status,
            TaskStatus::Done | TaskStatus::Held | TaskStatus::Blocked
        ) {
            return Ok(Followed {
                task,
                finished: true,
            });
        }
        let reading = crate::daemon::read_status(home);
        if crate::daemon::foreign_loop(reading.as_ref(), Timestamp::now(), std::process::id())
            .is_none()
        {
            bail!(
                "no magi loop is serving the queue any more, and task {} is still {}; it was \
                 not taken over here (the loop may only be slow to heartbeat). Start `magi \
                 serve` to carry on, or follow it with `magi task show {}`",
                task.short(),
                task.status.as_str(),
                task.short()
            );
        }
        if max_wait.is_some_and(|m| began.elapsed() >= m) {
            return Ok(Followed {
                task,
                finished: false,
            });
        }
        tokio::time::sleep(poll).await;
    }
}

/// A run opened inside this process (the follow-up review of `magi fix`) and
/// the task that owns it. See [`adopt`].
#[derive(Debug)]
pub struct Adopted {
    queue: Queue,
    /// The task this process claimed and started for the run; `None` when the
    /// claim could not be taken and the run was only linked.
    task: Option<Task>,
    claim: Option<Claim>,
    quota_before: Vec<crate::run::QuotaLoss>,
}

/// Give a run this process is about to execute an owning task. A parent task
/// that exists and is queued, failed or held is claimed, started with the
/// run and settled by [`Adopted::finish`] exactly like an ownerless run's; if
/// 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
/// parent, a task is filed and claimed for the duration. A no-op (`None`) when
/// no magi home is pinned.
pub fn adopt(state: &crate::run::RunState, parent_task: Option<&str>) -> Option<Adopted> {
    crate::run::try_home()?;
    adopt_in(Queue::open(), state, parent_task)
}

fn adopt_in(
    queue: Queue,
    state: &crate::run::RunState,
    parent_task: Option<&str>,
) -> Option<Adopted> {
    if let Some(parent) = parent_task
        && let Ok(id) = queue.resolve_id(parent)
    {
        // The parent exists: never file a second owner for this run.
        return match queue.claim(&id) {
            Ok(claim) => {
                let started = queue.get(&id).and_then(|mut task| {
                    // A task the daemon could pick up, or one a failed
                    // hand-started run left held (that is what a requested
                    // follow-up review re-verifies), is this run's to settle.
                    // A Done, blocked or running one is settled already or
                    // owned by somebody else: restarting it would let a
                    // follow-up review overwrite its outcome, so it only
                    // gains the run.
                    if !(task.status.runnable() || task.status == TaskStatus::Held) {
                        return Ok(None);
                    }
                    task.start(state.id.clone());
                    queue.put(&mut task)?;
                    Ok(Some(task))
                });
                match started {
                    Ok(None) => {
                        drop(claim);
                        let _ = queue.link_run(&id, &state.id);
                        None
                    }
                    Ok(Some(task)) => Some(Adopted {
                        queue,
                        task: Some(task),
                        claim: Some(claim),
                        quota_before: state.quota.clone(),
                    }),
                    Err(e) => {
                        tracing::warn!("could not start task {id} for run {}: {e:#}", state.id);
                        drop(claim);
                        let _ = queue.link_run(&id, &state.id);
                        None
                    }
                }
            }
            Err(e) => {
                // Another driver owns the task; do not settle over its result.
                tracing::warn!("task {id} is claimed elsewhere ({e:#}); linking run only");
                let _ = queue.link_run(&id, &state.id);
                None
            }
        };
    }
    let mut task = Task::new(
        crate::queue::title_from(&state.instruction, 72),
        state.instruction.clone(),
        state.repo.clone(),
        Source::Human,
    );
    let claim = queue.claim(&task.id).ok()?;
    task.start(state.id.clone());
    queue.put(&mut task).ok()?;
    Some(Adopted {
        queue,
        task: Some(task),
        claim: Some(claim),
        quota_before: state.quota.clone(),
    })
}

impl Adopted {
    /// Settle the task of an adopted ownerless run after it executed.
    pub fn finish(mut self, state: &crate::run::RunState, result: Result<()>) {
        if let Some(task) = &mut self.task {
            crate::daemon::finish_attempt(
                crate::daemon::Opts::default().max_attempts,
                &self.queue,
                task,
                state,
                &self.quota_before,
                result,
            );
            crate::daemon::hold_if_runnable(&self.queue, task);
        }
        drop(self.claim.take());
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::daemon::Status;

    fn filing(repo: &Path) -> Filing {
        Filing {
            instruction: "review the branch".to_owned(),
            title: "review x".to_owned(),
            repo: repo.to_path_buf(),
            source: Source::Human,
            solo: false,
            overrides: RunOverrides {
                merge: Some("none".to_owned()),
                ..RunOverrides::default()
            },
            review_of: Some("feat/x".to_owned()),
        }
    }

    fn heartbeat(home: &Path, pid: u32, age_secs: i64) {
        let mut status = Status::new();
        status.pid = pid;
        status.updated_at = Timestamp::now()
            .checked_sub(jiff::SignedDuration::from_secs(age_secs))
            .unwrap();
        crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
    }

    #[test]
    fn without_a_loop_the_task_is_claimed_before_anyone_else_can_take_it() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        let filed = file(
            &q,
            dir.path(),
            Timestamp::now(),
            1,
            filing(dir.path()),
            false,
        )
        .unwrap();
        let Filed::Standalone { task, claim } = filed else {
            panic!("no loop is alive");
        };
        assert!(!task.urgent);
        assert_eq!(task.review_of.as_deref(), Some("feat/x"));
        assert_eq!(
            task.overrides.as_ref().unwrap().merge.as_deref(),
            Some("none")
        );
        assert!(
            q.claim(&task.id).is_err(),
            "a second process must not be able to claim a standalone run's task"
        );
        drop(claim);
        assert!(q.claim(&task.id).is_ok(), "released with the claim");
    }

    #[test]
    fn a_live_foreign_loop_gets_an_urgent_unclaimed_task() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        heartbeat(dir.path(), 4242, 0);
        let filed = file(
            &q,
            dir.path(),
            Timestamp::now(),
            1,
            filing(dir.path()),
            false,
        )
        .unwrap();
        let Filed::Daemon { task, pid } = filed else {
            panic!("a loop is alive");
        };
        assert_eq!(pid, Some(4242));
        assert!(task.urgent);
        let stored = q.get(&task.id).unwrap();
        assert_eq!(stored.status, TaskStatus::Queued);
        assert!(stored.runs.is_empty(), "nothing ran in this process");
        assert!(q.claim(&task.id).is_ok(), "and this process holds no claim");
    }

    #[test]
    fn a_stale_heartbeat_our_own_pid_and_a_dry_run_do_not_count_as_a_loop() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        heartbeat(dir.path(), 4242, 3600);
        assert!(matches!(
            file(
                &q,
                dir.path(),
                Timestamp::now(),
                1,
                filing(dir.path()),
                false
            )
            .unwrap(),
            Filed::Standalone { .. }
        ));
        heartbeat(dir.path(), 7, 0);
        assert!(matches!(
            file(
                &q,
                dir.path(),
                Timestamp::now(),
                7,
                filing(dir.path()),
                false
            )
            .unwrap(),
            Filed::Standalone { .. }
        ));
        heartbeat(dir.path(), 4242, 0);
        assert!(matches!(
            file(
                &q,
                dir.path(),
                Timestamp::now(),
                1,
                filing(dir.path()),
                true
            )
            .unwrap(),
            Filed::Standalone { .. }
        ));
    }

    #[tokio::test]
    async fn following_gives_up_on_a_task_nobody_is_serving() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        let mut t = filing(dir.path()).into_task();
        q.put(&mut t).unwrap();
        let err = follow(
            &q,
            dir.path(),
            &t.id,
            Duration::from_millis(5),
            None,
            |_| {},
            |_| {},
        )
        .await
        .unwrap_err();
        assert!(format!("{err}").contains("no magi loop"), "{err}");
        assert!(q.claim(&t.id).is_ok(), "never taken over");
    }

    #[tokio::test]
    async fn following_returns_when_the_loop_finishes_the_task_and_times_out_otherwise() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        heartbeat(dir.path(), 4242, 0);
        let mut t = filing(dir.path()).into_task();
        q.put(&mut t).unwrap();
        let waited = follow(
            &q,
            dir.path(),
            &t.id,
            Duration::from_millis(5),
            Some(Duration::from_millis(30)),
            |_| {},
            |_| {},
        )
        .await
        .unwrap();
        assert!(!waited.finished, "the wait ran out, the loop still owns it");
        let err = waited
            .unfinished()
            .expect("an unfinished wait is an error")
            .to_string();
        assert!(err.contains(t.short()), "{err}");
        assert!(
            err.contains("not finished") || err.contains("not a success"),
            "{err}"
        );
        assert!(err.contains("magi task show"), "{err}");
        assert_eq!(
            q.get(&t.id).unwrap().status,
            TaskStatus::Queued,
            "the queue is untouched"
        );
        t.link_run("20260101-000000-abcd");
        t.succeed();
        q.put(&mut t).unwrap();
        let mut runs = Vec::new();
        let done = follow(
            &q,
            dir.path(),
            &t.id,
            Duration::from_millis(5),
            Some(Duration::from_secs(5)),
            |r| runs.push(r.to_owned()),
            |_| {},
        )
        .await
        .unwrap();
        assert_eq!(done.task.status, TaskStatus::Done);
        assert!(done.unfinished().is_none());
        assert_eq!(runs, ["20260101-000000-abcd"]);
    }

    fn parent_in(q: &Queue, dir: &Path) -> Task {
        let mut t = filing(dir).into_task();
        q.put(&mut t).unwrap();
        t
    }

    fn follow_up_state(dir: &Path) -> crate::run::RunState {
        crate::run::RunState::new(
            dir.to_path_buf(),
            "main".to_owned(),
            "abc1234".to_owned(),
            "review the branch".to_owned(),
            Config::default(),
        )
    }

    #[test]
    fn an_adopted_review_claims_starts_and_settles_a_runnable_task() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        let parent = parent_in(&q, dir.path());
        let state = follow_up_state(dir.path());
        let adopted = adopt_in(q.clone(), &state, Some(&parent.id)).expect("adopted");
        let running = q.get(&parent.id).unwrap();
        assert_eq!(running.status, TaskStatus::Running);
        assert_eq!(running.attempts, parent.attempts + 1);
        assert_eq!(running.runs, std::slice::from_ref(&state.id));
        assert!(q.claim(&parent.id).is_err(), "exclusive while it runs");
        adopted.finish(&state, Err(anyhow::anyhow!("boom")));
        let settled = q.get(&parent.id).unwrap();
        assert_eq!(settled.status, TaskStatus::Held, "{settled:?}");
        assert_eq!(
            settled.runs,
            std::slice::from_ref(&state.id),
            "no duplicate run"
        );
        assert!(q.claim(&parent.id).is_ok(), "claim released");
        assert_eq!(q.list().len(), 1, "no second owner was filed");
    }

    #[test]
    fn a_task_claimed_elsewhere_only_gains_the_run() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        let parent = parent_in(&q, dir.path());
        let _theirs = q.claim(&parent.id).unwrap();
        let state = follow_up_state(dir.path());
        assert!(adopt_in(q.clone(), &state, Some(&parent.id)).is_none());
        let after = q.get(&parent.id).unwrap();
        assert_eq!(after.status, parent.status);
        assert_eq!(after.attempts, parent.attempts);
        assert_eq!(after.runs, std::slice::from_ref(&state.id));
        assert_eq!(q.list().len(), 1, "no second owner was filed");
    }

    #[test]
    fn a_done_task_only_gains_the_run() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        let mut parent = parent_in(&q, dir.path());
        parent.succeed();
        q.put(&mut parent).unwrap();
        let state = follow_up_state(dir.path());
        assert!(adopt_in(q.clone(), &state, Some(&parent.id)).is_none());
        let after = q.get(&parent.id).unwrap();
        assert_eq!(after.status, TaskStatus::Done);
        assert_eq!(after.attempts, parent.attempts);
        assert_eq!(after.runs, std::slice::from_ref(&state.id));
        assert!(q.claim(&parent.id).is_ok(), "claim released");
    }

    #[test]
    fn a_held_task_is_claimed_and_settled_by_its_follow_up_review() {
        let dir = tempfile::tempdir().unwrap();
        let q = Queue::at(dir.path().join("queue"));
        let mut parent = parent_in(&q, dir.path());
        parent.hold_manual(Some("the run did not finish: stale".to_owned()));
        q.put(&mut parent).unwrap();
        let state = follow_up_state(dir.path());
        let adopted = adopt_in(q.clone(), &state, Some(&parent.id)).expect("adopted");
        assert_eq!(q.get(&parent.id).unwrap().status, TaskStatus::Running);
        assert!(q.claim(&parent.id).is_err(), "claimed while it runs");
        adopted.finish(&state, Err(anyhow::anyhow!("fresh failure")));
        let after = q.get(&parent.id).unwrap();
        assert_eq!(after.status, TaskStatus::Held);
        assert!(
            after
                .hold_reason
                .as_deref()
                .unwrap()
                .contains("fresh failure"),
            "{after:?}"
        );
        assert_eq!(after.runs, std::slice::from_ref(&state.id));
        assert!(q.claim(&parent.id).is_ok(), "claim released");
    }
}