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