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
136impl Followed {
137 pub fn unfinished(&self) -> Option<anyhow::Error> {
141 if self.finished {
142 return None;
143 }
144 Some(anyhow::anyhow!(
145 "task {id} has not finished: it is still running under the loop. \
146 This is not a success; the outcome is not known yet. Read it with \
147 `magi task show {id}`. The loop keeps the task, so do not run it again.",
148 id = self.task.short()
149 ))
150 }
151}
152
153pub async fn follow(
164 queue: &Queue,
165 home: &Path,
166 id: &str,
167 poll: Duration,
168 max_wait: Option<Duration>,
169 mut seen: impl FnMut(&str),
170 mut progress: impl FnMut(&str),
171) -> Result<Followed> {
172 let began = Instant::now();
173 let mut announced = 0usize;
174 let mut last_status = String::new();
175 loop {
176 let task = queue.get(id)?;
177 for run in task.runs.iter().skip(announced) {
178 seen(run);
179 }
180 announced = task.runs.len();
181 if let Some(run) = task.runs.last()
182 && crate::run::try_home().is_some()
183 && let Ok(state) = crate::run::RunState::load(run)
184 {
185 let line = format!("run {}: {}", state.short(), state.status.as_str());
186 if line != last_status {
187 progress(&line);
188 last_status = line;
189 }
190 }
191 if matches!(
192 task.status,
193 TaskStatus::Done | TaskStatus::Held | TaskStatus::Blocked
194 ) {
195 return Ok(Followed {
196 task,
197 finished: true,
198 });
199 }
200 let reading = crate::daemon::read_status(home);
201 if crate::daemon::foreign_loop(reading.as_ref(), Timestamp::now(), std::process::id())
202 .is_none()
203 {
204 bail!(
205 "no magi loop is serving the queue any more, and task {} is still {}; it was \
206 not taken over here (the loop may only be slow to heartbeat). Start `magi \
207 serve` to carry on, or follow it with `magi task show {}`",
208 task.short(),
209 task.status.as_str(),
210 task.short()
211 );
212 }
213 if max_wait.is_some_and(|m| began.elapsed() >= m) {
214 return Ok(Followed {
215 task,
216 finished: false,
217 });
218 }
219 tokio::time::sleep(poll).await;
220 }
221}
222
223#[derive(Debug)]
226pub struct Adopted {
227 queue: Queue,
228 task: Option<Task>,
231 claim: Option<Claim>,
232 quota_before: Vec<crate::run::QuotaLoss>,
233}
234
235pub fn adopt(
244 state: &crate::run::RunState,
245 parent_task: Option<&str>,
246) -> anyhow::Result<Option<Adopted>> {
247 if crate::run::try_home().is_none() {
248 return Ok(None);
249 }
250 adopt_in(Queue::open(), state, parent_task)
251}
252
253fn adopt_in(
254 queue: Queue,
255 state: &crate::run::RunState,
256 parent_task: Option<&str>,
257) -> anyhow::Result<Option<Adopted>> {
258 use anyhow::Context as _;
259 if let Some(parent) = parent_task
260 && let Ok(id) = queue.resolve_id(parent)
261 {
262 return Ok(match queue.claim(&id) {
264 Ok(claim) => {
265 let started = queue.get(&id).and_then(|mut task| {
266 if !(task.status.runnable() || task.status == TaskStatus::Held) {
274 return Ok(None);
275 }
276 task.start(state.id.clone());
277 queue.put(&mut task)?;
278 Ok(Some(task))
279 });
280 match started {
281 Ok(None) => {
282 drop(claim);
283 let _ = queue.link_run(&id, &state.id);
284 None
285 }
286 Ok(Some(task)) => Some(Adopted {
287 queue,
288 task: Some(task),
289 claim: Some(claim),
290 quota_before: state.quota.clone(),
291 }),
292 Err(e) => {
293 tracing::warn!("could not start task {id} for run {}: {e:#}", state.id);
294 drop(claim);
295 let _ = queue.link_run(&id, &state.id);
296 None
297 }
298 }
299 }
300 Err(e) => {
301 tracing::warn!("task {id} is claimed elsewhere ({e:#}); linking run only");
303 let _ = queue.link_run(&id, &state.id);
304 None
305 }
306 });
307 }
308 let mut task = Task::new(
309 crate::queue::title_from(&state.instruction, 72),
310 state.instruction.clone(),
311 state.repo.clone(),
312 Source::Human,
313 );
314 let claim = queue
315 .claim(&task.id)
316 .with_context(|| format!("could not claim a new task for run {}", state.id))?;
317 task.start(state.id.clone());
318 queue
319 .put(&mut task)
320 .with_context(|| format!("could not file a new task for run {}", state.id))?;
321 Ok(Some(Adopted {
322 queue,
323 task: Some(task),
324 claim: Some(claim),
325 quota_before: state.quota.clone(),
326 }))
327}
328
329impl Adopted {
330 pub fn finish(mut self, state: &crate::run::RunState, result: Result<()>) {
332 if let Some(task) = &mut self.task {
333 crate::daemon::finish_attempt(
334 crate::daemon::Opts::default().max_attempts,
335 &self.queue,
336 task,
337 state,
338 &self.quota_before,
339 result,
340 );
341 crate::daemon::hold_if_runnable(&self.queue, task);
342 }
343 drop(self.claim.take());
344 }
345}
346
347#[cfg(test)]
348mod tests {
349 use super::*;
350 use crate::daemon::Status;
351
352 fn filing(repo: &Path) -> Filing {
353 Filing {
354 instruction: "review the branch".to_owned(),
355 title: "review x".to_owned(),
356 repo: repo.to_path_buf(),
357 source: Source::Human,
358 solo: false,
359 overrides: RunOverrides {
360 merge: Some("none".to_owned()),
361 ..RunOverrides::default()
362 },
363 review_of: Some("feat/x".to_owned()),
364 }
365 }
366
367 fn heartbeat(home: &Path, pid: u32, age_secs: i64) {
368 let mut status = Status::new();
369 status.pid = pid;
370 status.updated_at = Timestamp::now()
371 .checked_sub(jiff::SignedDuration::from_secs(age_secs))
372 .unwrap();
373 crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
374 }
375
376 #[test]
377 fn without_a_loop_the_task_is_claimed_before_anyone_else_can_take_it() {
378 let dir = tempfile::tempdir().unwrap();
379 let q = Queue::at(dir.path().join("queue"));
380 let filed = file(
381 &q,
382 dir.path(),
383 Timestamp::now(),
384 1,
385 filing(dir.path()),
386 false,
387 )
388 .unwrap();
389 let Filed::Standalone { task, claim } = filed else {
390 panic!("no loop is alive");
391 };
392 assert!(!task.urgent);
393 assert_eq!(task.review_of.as_deref(), Some("feat/x"));
394 assert_eq!(
395 task.overrides.as_ref().unwrap().merge.as_deref(),
396 Some("none")
397 );
398 assert!(
399 q.claim(&task.id).is_err(),
400 "a second process must not be able to claim a standalone run's task"
401 );
402 drop(claim);
403 assert!(q.claim(&task.id).is_ok(), "released with the claim");
404 }
405
406 #[test]
407 fn a_live_foreign_loop_gets_an_urgent_unclaimed_task() {
408 let dir = tempfile::tempdir().unwrap();
409 let q = Queue::at(dir.path().join("queue"));
410 heartbeat(dir.path(), 4242, 0);
411 let filed = file(
412 &q,
413 dir.path(),
414 Timestamp::now(),
415 1,
416 filing(dir.path()),
417 false,
418 )
419 .unwrap();
420 let Filed::Daemon { task, pid } = filed else {
421 panic!("a loop is alive");
422 };
423 assert_eq!(pid, Some(4242));
424 assert!(task.urgent);
425 let stored = q.get(&task.id).unwrap();
426 assert_eq!(stored.status, TaskStatus::Queued);
427 assert!(stored.runs.is_empty(), "nothing ran in this process");
428 assert!(q.claim(&task.id).is_ok(), "and this process holds no claim");
429 }
430
431 #[test]
432 fn a_stale_heartbeat_our_own_pid_and_a_dry_run_do_not_count_as_a_loop() {
433 let dir = tempfile::tempdir().unwrap();
434 let q = Queue::at(dir.path().join("queue"));
435 heartbeat(dir.path(), 4242, 3600);
436 assert!(matches!(
437 file(
438 &q,
439 dir.path(),
440 Timestamp::now(),
441 1,
442 filing(dir.path()),
443 false
444 )
445 .unwrap(),
446 Filed::Standalone { .. }
447 ));
448 heartbeat(dir.path(), 7, 0);
449 assert!(matches!(
450 file(
451 &q,
452 dir.path(),
453 Timestamp::now(),
454 7,
455 filing(dir.path()),
456 false
457 )
458 .unwrap(),
459 Filed::Standalone { .. }
460 ));
461 heartbeat(dir.path(), 4242, 0);
462 assert!(matches!(
463 file(
464 &q,
465 dir.path(),
466 Timestamp::now(),
467 1,
468 filing(dir.path()),
469 true
470 )
471 .unwrap(),
472 Filed::Standalone { .. }
473 ));
474 }
475
476 #[tokio::test]
477 async fn following_gives_up_on_a_task_nobody_is_serving() {
478 let dir = tempfile::tempdir().unwrap();
479 let q = Queue::at(dir.path().join("queue"));
480 let mut t = filing(dir.path()).into_task();
481 q.put(&mut t).unwrap();
482 let err = follow(
483 &q,
484 dir.path(),
485 &t.id,
486 Duration::from_millis(5),
487 None,
488 |_| {},
489 |_| {},
490 )
491 .await
492 .unwrap_err();
493 assert!(format!("{err}").contains("no magi loop"), "{err}");
494 assert!(q.claim(&t.id).is_ok(), "never taken over");
495 }
496
497 #[tokio::test]
498 async fn following_returns_when_the_loop_finishes_the_task_and_times_out_otherwise() {
499 let dir = tempfile::tempdir().unwrap();
500 let q = Queue::at(dir.path().join("queue"));
501 heartbeat(dir.path(), 4242, 0);
502 let mut t = filing(dir.path()).into_task();
503 q.put(&mut t).unwrap();
504 let waited = follow(
505 &q,
506 dir.path(),
507 &t.id,
508 Duration::from_millis(5),
509 Some(Duration::from_millis(30)),
510 |_| {},
511 |_| {},
512 )
513 .await
514 .unwrap();
515 assert!(!waited.finished, "the wait ran out, the loop still owns it");
516 let err = waited
517 .unfinished()
518 .expect("an unfinished wait is an error")
519 .to_string();
520 assert!(err.contains(t.short()), "{err}");
521 assert!(
522 err.contains("not finished") || err.contains("not a success"),
523 "{err}"
524 );
525 assert!(err.contains("magi task show"), "{err}");
526 assert_eq!(
527 q.get(&t.id).unwrap().status,
528 TaskStatus::Queued,
529 "the queue is untouched"
530 );
531 t.link_run("20260101-000000-abcd");
532 t.succeed();
533 q.put(&mut t).unwrap();
534 let mut runs = Vec::new();
535 let done = follow(
536 &q,
537 dir.path(),
538 &t.id,
539 Duration::from_millis(5),
540 Some(Duration::from_secs(5)),
541 |r| runs.push(r.to_owned()),
542 |_| {},
543 )
544 .await
545 .unwrap();
546 assert_eq!(done.task.status, TaskStatus::Done);
547 assert!(done.unfinished().is_none());
548 assert_eq!(runs, ["20260101-000000-abcd"]);
549 }
550
551 fn parent_in(q: &Queue, dir: &Path) -> Task {
552 let mut t = filing(dir).into_task();
553 q.put(&mut t).unwrap();
554 t
555 }
556
557 fn follow_up_state(dir: &Path) -> crate::run::RunState {
558 crate::run::RunState::new(
559 dir.to_path_buf(),
560 "main".to_owned(),
561 "abc1234".to_owned(),
562 "review the branch".to_owned(),
563 Config::default(),
564 )
565 }
566
567 #[test]
568 fn an_adopted_review_claims_starts_and_settles_a_runnable_task() {
569 let dir = tempfile::tempdir().unwrap();
570 let q = Queue::at(dir.path().join("queue"));
571 let parent = parent_in(&q, dir.path());
572 let state = follow_up_state(dir.path());
573 let adopted = adopt_in(q.clone(), &state, Some(&parent.id))
574 .unwrap()
575 .expect("adopted");
576 let running = q.get(&parent.id).unwrap();
577 assert_eq!(running.status, TaskStatus::Running);
578 assert_eq!(running.attempts, parent.attempts + 1);
579 assert_eq!(running.runs, std::slice::from_ref(&state.id));
580 assert!(q.claim(&parent.id).is_err(), "exclusive while it runs");
581 adopted.finish(&state, Err(anyhow::anyhow!("boom")));
582 let settled = q.get(&parent.id).unwrap();
583 assert_eq!(settled.status, TaskStatus::Held, "{settled:?}");
584 assert_eq!(
585 settled.runs,
586 std::slice::from_ref(&state.id),
587 "no duplicate run"
588 );
589 assert!(q.claim(&parent.id).is_ok(), "claim released");
590 assert_eq!(q.list().len(), 1, "no second owner was filed");
591 }
592
593 #[test]
594 fn an_ownerless_run_whose_task_cannot_be_filed_is_an_error() {
595 let dir = tempfile::tempdir().unwrap();
596 let blocker = dir.path().join("blocker");
598 std::fs::write(&blocker, "x").unwrap();
599 let q = Queue::at(blocker.join("q"));
600 let state = follow_up_state(dir.path());
601 let err = adopt_in(q, &state, None).expect_err("must not be swallowed");
602 assert!(format!("{err:#}").contains(&state.id), "{err:#}");
603 }
604
605 #[test]
606 fn a_task_claimed_elsewhere_only_gains_the_run() {
607 let dir = tempfile::tempdir().unwrap();
608 let q = Queue::at(dir.path().join("queue"));
609 let parent = parent_in(&q, dir.path());
610 let _theirs = q.claim(&parent.id).unwrap();
611 let state = follow_up_state(dir.path());
612 assert!(
613 adopt_in(q.clone(), &state, Some(&parent.id))
614 .unwrap()
615 .is_none()
616 );
617 let after = q.get(&parent.id).unwrap();
618 assert_eq!(after.status, parent.status);
619 assert_eq!(after.attempts, parent.attempts);
620 assert_eq!(after.runs, std::slice::from_ref(&state.id));
621 assert_eq!(q.list().len(), 1, "no second owner was filed");
622 }
623
624 #[test]
625 fn a_done_task_only_gains_the_run() {
626 let dir = tempfile::tempdir().unwrap();
627 let q = Queue::at(dir.path().join("queue"));
628 let mut parent = parent_in(&q, dir.path());
629 parent.succeed();
630 q.put(&mut parent).unwrap();
631 let state = follow_up_state(dir.path());
632 assert!(
633 adopt_in(q.clone(), &state, Some(&parent.id))
634 .unwrap()
635 .is_none()
636 );
637 let after = q.get(&parent.id).unwrap();
638 assert_eq!(after.status, TaskStatus::Done);
639 assert_eq!(after.attempts, parent.attempts);
640 assert_eq!(after.runs, std::slice::from_ref(&state.id));
641 assert!(q.claim(&parent.id).is_ok(), "claim released");
642 }
643
644 #[test]
645 fn a_held_task_is_claimed_and_settled_by_its_follow_up_review() {
646 let dir = tempfile::tempdir().unwrap();
647 let q = Queue::at(dir.path().join("queue"));
648 let mut parent = parent_in(&q, dir.path());
649 parent.hold_manual(Some("the run did not finish: stale".to_owned()));
650 q.put(&mut parent).unwrap();
651 let state = follow_up_state(dir.path());
652 let adopted = adopt_in(q.clone(), &state, Some(&parent.id))
653 .unwrap()
654 .expect("adopted");
655 assert_eq!(q.get(&parent.id).unwrap().status, TaskStatus::Running);
656 assert!(q.claim(&parent.id).is_err(), "claimed while it runs");
657 adopted.finish(&state, Err(anyhow::anyhow!("fresh failure")));
658 let after = q.get(&parent.id).unwrap();
659 assert_eq!(after.status, TaskStatus::Held);
660 assert!(
661 after
662 .hold_reason
663 .as_deref()
664 .unwrap()
665 .contains("fresh failure"),
666 "{after:?}"
667 );
668 assert_eq!(after.runs, std::slice::from_ref(&state.id));
669 assert!(q.claim(&parent.id).is_ok(), "claim released");
670 }
671}