1use 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#[derive(Debug, Clone)]
31pub struct Filing {
32 pub instruction: String,
34 pub title: String,
36 pub repo: PathBuf,
38 pub source: Source,
40 pub solo: bool,
42 pub overrides: RunOverrides,
44 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#[derive(Debug)]
60pub enum Filed {
61 Daemon {
64 task: Task,
66 pid: Option<u32>,
68 },
69 Standalone {
72 task: Task,
74 claim: Claim,
76 },
77}
78
79pub 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 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
112pub 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#[derive(Debug)]
128pub struct Followed {
129 pub task: Task,
131 pub finished: bool,
134}
135
136pub 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#[derive(Debug)]
209pub struct Adopted {
210 queue: Queue,
211 task: Option<Task>,
214 claim: Option<Claim>,
215 quota_before: Vec<crate::run::QuotaLoss>,
216}
217
218pub 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 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}