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(state: &crate::run::RunState, parent_task: Option<&str>) -> Option<Adopted> {
242 crate::run::try_home()?;
243 adopt_in(Queue::open(), state, parent_task)
244}
245
246fn adopt_in(
247 queue: Queue,
248 state: &crate::run::RunState,
249 parent_task: Option<&str>,
250) -> Option<Adopted> {
251 if let Some(parent) = parent_task
252 && let Ok(id) = queue.resolve_id(parent)
253 {
254 return match queue.claim(&id) {
256 Ok(claim) => {
257 let started = queue.get(&id).and_then(|mut task| {
258 if !(task.status.runnable() || task.status == TaskStatus::Held) {
266 return Ok(None);
267 }
268 task.start(state.id.clone());
269 queue.put(&mut task)?;
270 Ok(Some(task))
271 });
272 match started {
273 Ok(None) => {
274 drop(claim);
275 let _ = queue.link_run(&id, &state.id);
276 None
277 }
278 Ok(Some(task)) => Some(Adopted {
279 queue,
280 task: Some(task),
281 claim: Some(claim),
282 quota_before: state.quota.clone(),
283 }),
284 Err(e) => {
285 tracing::warn!("could not start task {id} for run {}: {e:#}", state.id);
286 drop(claim);
287 let _ = queue.link_run(&id, &state.id);
288 None
289 }
290 }
291 }
292 Err(e) => {
293 tracing::warn!("task {id} is claimed elsewhere ({e:#}); linking run only");
295 let _ = queue.link_run(&id, &state.id);
296 None
297 }
298 };
299 }
300 let mut task = Task::new(
301 crate::queue::title_from(&state.instruction, 72),
302 state.instruction.clone(),
303 state.repo.clone(),
304 Source::Human,
305 );
306 let claim = queue.claim(&task.id).ok()?;
307 task.start(state.id.clone());
308 queue.put(&mut task).ok()?;
309 Some(Adopted {
310 queue,
311 task: Some(task),
312 claim: Some(claim),
313 quota_before: state.quota.clone(),
314 })
315}
316
317impl Adopted {
318 pub fn finish(mut self, state: &crate::run::RunState, result: Result<()>) {
320 if let Some(task) = &mut self.task {
321 crate::daemon::finish_attempt(
322 crate::daemon::Opts::default().max_attempts,
323 &self.queue,
324 task,
325 state,
326 &self.quota_before,
327 result,
328 );
329 crate::daemon::hold_if_runnable(&self.queue, task);
330 }
331 drop(self.claim.take());
332 }
333}
334
335#[cfg(test)]
336mod tests {
337 use super::*;
338 use crate::daemon::Status;
339
340 fn filing(repo: &Path) -> Filing {
341 Filing {
342 instruction: "review the branch".to_owned(),
343 title: "review x".to_owned(),
344 repo: repo.to_path_buf(),
345 source: Source::Human,
346 solo: false,
347 overrides: RunOverrides {
348 merge: Some("none".to_owned()),
349 ..RunOverrides::default()
350 },
351 review_of: Some("feat/x".to_owned()),
352 }
353 }
354
355 fn heartbeat(home: &Path, pid: u32, age_secs: i64) {
356 let mut status = Status::new();
357 status.pid = pid;
358 status.updated_at = Timestamp::now()
359 .checked_sub(jiff::SignedDuration::from_secs(age_secs))
360 .unwrap();
361 crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
362 }
363
364 #[test]
365 fn without_a_loop_the_task_is_claimed_before_anyone_else_can_take_it() {
366 let dir = tempfile::tempdir().unwrap();
367 let q = Queue::at(dir.path().join("queue"));
368 let filed = file(
369 &q,
370 dir.path(),
371 Timestamp::now(),
372 1,
373 filing(dir.path()),
374 false,
375 )
376 .unwrap();
377 let Filed::Standalone { task, claim } = filed else {
378 panic!("no loop is alive");
379 };
380 assert!(!task.urgent);
381 assert_eq!(task.review_of.as_deref(), Some("feat/x"));
382 assert_eq!(
383 task.overrides.as_ref().unwrap().merge.as_deref(),
384 Some("none")
385 );
386 assert!(
387 q.claim(&task.id).is_err(),
388 "a second process must not be able to claim a standalone run's task"
389 );
390 drop(claim);
391 assert!(q.claim(&task.id).is_ok(), "released with the claim");
392 }
393
394 #[test]
395 fn a_live_foreign_loop_gets_an_urgent_unclaimed_task() {
396 let dir = tempfile::tempdir().unwrap();
397 let q = Queue::at(dir.path().join("queue"));
398 heartbeat(dir.path(), 4242, 0);
399 let filed = file(
400 &q,
401 dir.path(),
402 Timestamp::now(),
403 1,
404 filing(dir.path()),
405 false,
406 )
407 .unwrap();
408 let Filed::Daemon { task, pid } = filed else {
409 panic!("a loop is alive");
410 };
411 assert_eq!(pid, Some(4242));
412 assert!(task.urgent);
413 let stored = q.get(&task.id).unwrap();
414 assert_eq!(stored.status, TaskStatus::Queued);
415 assert!(stored.runs.is_empty(), "nothing ran in this process");
416 assert!(q.claim(&task.id).is_ok(), "and this process holds no claim");
417 }
418
419 #[test]
420 fn a_stale_heartbeat_our_own_pid_and_a_dry_run_do_not_count_as_a_loop() {
421 let dir = tempfile::tempdir().unwrap();
422 let q = Queue::at(dir.path().join("queue"));
423 heartbeat(dir.path(), 4242, 3600);
424 assert!(matches!(
425 file(
426 &q,
427 dir.path(),
428 Timestamp::now(),
429 1,
430 filing(dir.path()),
431 false
432 )
433 .unwrap(),
434 Filed::Standalone { .. }
435 ));
436 heartbeat(dir.path(), 7, 0);
437 assert!(matches!(
438 file(
439 &q,
440 dir.path(),
441 Timestamp::now(),
442 7,
443 filing(dir.path()),
444 false
445 )
446 .unwrap(),
447 Filed::Standalone { .. }
448 ));
449 heartbeat(dir.path(), 4242, 0);
450 assert!(matches!(
451 file(
452 &q,
453 dir.path(),
454 Timestamp::now(),
455 1,
456 filing(dir.path()),
457 true
458 )
459 .unwrap(),
460 Filed::Standalone { .. }
461 ));
462 }
463
464 #[tokio::test]
465 async fn following_gives_up_on_a_task_nobody_is_serving() {
466 let dir = tempfile::tempdir().unwrap();
467 let q = Queue::at(dir.path().join("queue"));
468 let mut t = filing(dir.path()).into_task();
469 q.put(&mut t).unwrap();
470 let err = follow(
471 &q,
472 dir.path(),
473 &t.id,
474 Duration::from_millis(5),
475 None,
476 |_| {},
477 |_| {},
478 )
479 .await
480 .unwrap_err();
481 assert!(format!("{err}").contains("no magi loop"), "{err}");
482 assert!(q.claim(&t.id).is_ok(), "never taken over");
483 }
484
485 #[tokio::test]
486 async fn following_returns_when_the_loop_finishes_the_task_and_times_out_otherwise() {
487 let dir = tempfile::tempdir().unwrap();
488 let q = Queue::at(dir.path().join("queue"));
489 heartbeat(dir.path(), 4242, 0);
490 let mut t = filing(dir.path()).into_task();
491 q.put(&mut t).unwrap();
492 let waited = follow(
493 &q,
494 dir.path(),
495 &t.id,
496 Duration::from_millis(5),
497 Some(Duration::from_millis(30)),
498 |_| {},
499 |_| {},
500 )
501 .await
502 .unwrap();
503 assert!(!waited.finished, "the wait ran out, the loop still owns it");
504 let err = waited
505 .unfinished()
506 .expect("an unfinished wait is an error")
507 .to_string();
508 assert!(err.contains(t.short()), "{err}");
509 assert!(
510 err.contains("not finished") || err.contains("not a success"),
511 "{err}"
512 );
513 assert!(err.contains("magi task show"), "{err}");
514 assert_eq!(
515 q.get(&t.id).unwrap().status,
516 TaskStatus::Queued,
517 "the queue is untouched"
518 );
519 t.link_run("20260101-000000-abcd");
520 t.succeed();
521 q.put(&mut t).unwrap();
522 let mut runs = Vec::new();
523 let done = follow(
524 &q,
525 dir.path(),
526 &t.id,
527 Duration::from_millis(5),
528 Some(Duration::from_secs(5)),
529 |r| runs.push(r.to_owned()),
530 |_| {},
531 )
532 .await
533 .unwrap();
534 assert_eq!(done.task.status, TaskStatus::Done);
535 assert!(done.unfinished().is_none());
536 assert_eq!(runs, ["20260101-000000-abcd"]);
537 }
538
539 fn parent_in(q: &Queue, dir: &Path) -> Task {
540 let mut t = filing(dir).into_task();
541 q.put(&mut t).unwrap();
542 t
543 }
544
545 fn follow_up_state(dir: &Path) -> crate::run::RunState {
546 crate::run::RunState::new(
547 dir.to_path_buf(),
548 "main".to_owned(),
549 "abc1234".to_owned(),
550 "review the branch".to_owned(),
551 Config::default(),
552 )
553 }
554
555 #[test]
556 fn an_adopted_review_claims_starts_and_settles_a_runnable_task() {
557 let dir = tempfile::tempdir().unwrap();
558 let q = Queue::at(dir.path().join("queue"));
559 let parent = parent_in(&q, dir.path());
560 let state = follow_up_state(dir.path());
561 let adopted = adopt_in(q.clone(), &state, Some(&parent.id)).expect("adopted");
562 let running = q.get(&parent.id).unwrap();
563 assert_eq!(running.status, TaskStatus::Running);
564 assert_eq!(running.attempts, parent.attempts + 1);
565 assert_eq!(running.runs, std::slice::from_ref(&state.id));
566 assert!(q.claim(&parent.id).is_err(), "exclusive while it runs");
567 adopted.finish(&state, Err(anyhow::anyhow!("boom")));
568 let settled = q.get(&parent.id).unwrap();
569 assert_eq!(settled.status, TaskStatus::Held, "{settled:?}");
570 assert_eq!(
571 settled.runs,
572 std::slice::from_ref(&state.id),
573 "no duplicate run"
574 );
575 assert!(q.claim(&parent.id).is_ok(), "claim released");
576 assert_eq!(q.list().len(), 1, "no second owner was filed");
577 }
578
579 #[test]
580 fn a_task_claimed_elsewhere_only_gains_the_run() {
581 let dir = tempfile::tempdir().unwrap();
582 let q = Queue::at(dir.path().join("queue"));
583 let parent = parent_in(&q, dir.path());
584 let _theirs = q.claim(&parent.id).unwrap();
585 let state = follow_up_state(dir.path());
586 assert!(adopt_in(q.clone(), &state, Some(&parent.id)).is_none());
587 let after = q.get(&parent.id).unwrap();
588 assert_eq!(after.status, parent.status);
589 assert_eq!(after.attempts, parent.attempts);
590 assert_eq!(after.runs, std::slice::from_ref(&state.id));
591 assert_eq!(q.list().len(), 1, "no second owner was filed");
592 }
593
594 #[test]
595 fn a_done_task_only_gains_the_run() {
596 let dir = tempfile::tempdir().unwrap();
597 let q = Queue::at(dir.path().join("queue"));
598 let mut parent = parent_in(&q, dir.path());
599 parent.succeed();
600 q.put(&mut parent).unwrap();
601 let state = follow_up_state(dir.path());
602 assert!(adopt_in(q.clone(), &state, Some(&parent.id)).is_none());
603 let after = q.get(&parent.id).unwrap();
604 assert_eq!(after.status, TaskStatus::Done);
605 assert_eq!(after.attempts, parent.attempts);
606 assert_eq!(after.runs, std::slice::from_ref(&state.id));
607 assert!(q.claim(&parent.id).is_ok(), "claim released");
608 }
609
610 #[test]
611 fn a_held_task_is_claimed_and_settled_by_its_follow_up_review() {
612 let dir = tempfile::tempdir().unwrap();
613 let q = Queue::at(dir.path().join("queue"));
614 let mut parent = parent_in(&q, dir.path());
615 parent.hold_manual(Some("the run did not finish: stale".to_owned()));
616 q.put(&mut parent).unwrap();
617 let state = follow_up_state(dir.path());
618 let adopted = adopt_in(q.clone(), &state, Some(&parent.id)).expect("adopted");
619 assert_eq!(q.get(&parent.id).unwrap().status, TaskStatus::Running);
620 assert!(q.claim(&parent.id).is_err(), "claimed while it runs");
621 adopted.finish(&state, Err(anyhow::anyhow!("fresh failure")));
622 let after = q.get(&parent.id).unwrap();
623 assert_eq!(after.status, TaskStatus::Held);
624 assert!(
625 after
626 .hold_reason
627 .as_deref()
628 .unwrap()
629 .contains("fresh failure"),
630 "{after:?}"
631 );
632 assert_eq!(after.runs, std::slice::from_ref(&state.id));
633 assert!(q.claim(&parent.id).is_ok(), "claim released");
634 }
635}