magi-cli 0.79.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
//! 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,
}

/// 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 created for an ownerless run; `None` when an existing task
    /// only gained the run.
    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 just gains the run in `runs`; with none, a task is filed and
/// claimed for the duration, and [`Adopted::finish`] settles it as the daemon
/// would. 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()?;
    let queue = Queue::open();
    if let Some(parent) = parent_task
        && queue.link_run(parent, &state.id).is_ok()
    {
        return Some(Adopted {
            queue,
            task: None,
            claim: None,
            quota_before: Vec::new(),
        });
    }
    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");
        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_eq!(runs, ["20260101-000000-abcd"]);
    }
}