1use std::collections::HashMap;
6use std::io;
7use std::path::PathBuf;
8use std::sync::atomic::{AtomicU64, Ordering};
9use std::future::Future;
10use std::sync::Arc;
11
12use tokio::sync::{oneshot, Mutex};
13use tokio::task::JoinHandle;
14
15use super::stream::BoundedStream;
16use crate::interpreter::ExecResult;
17
18pub use kaish_types::{JobId, JobInfo, JobStatus};
20
21pub struct Job {
23 pub id: JobId,
25 session_id: u64,
29 pub command: String,
31 handle: Option<JoinHandle<ExecResult>>,
33 result_rx: Option<oneshot::Receiver<ExecResult>>,
35 result: Option<ExecResult>,
37 output_file: Option<PathBuf>,
39 persist_output: bool,
47 stdout_stream: Option<Arc<BoundedStream>>,
49 stderr_stream: Option<Arc<BoundedStream>>,
51 pid: Option<u32>,
53 pgid: Option<u32>,
55 stopped: bool,
57 cancel: Option<tokio_util::sync::CancellationToken>,
63 pgids: Vec<u32>,
68}
69
70impl Job {
71 pub fn new(id: JobId, session_id: u64, command: String, handle: JoinHandle<ExecResult>) -> Self {
73 Self {
74 id,
75 session_id,
76 command,
77 handle: Some(handle),
78 result_rx: None,
79 result: None,
80 output_file: None,
81 persist_output: true,
82 stdout_stream: None,
83 stderr_stream: None,
84 pid: None,
85 pgid: None,
86 stopped: false,
87 cancel: None,
88 pgids: Vec::new(),
89 }
90 }
91
92 pub fn from_channel(id: JobId, session_id: u64, command: String, rx: oneshot::Receiver<ExecResult>) -> Self {
94 Self {
95 id,
96 session_id,
97 command,
98 handle: None,
99 result_rx: Some(rx),
100 result: None,
101 output_file: None,
102 persist_output: true,
103 stdout_stream: None,
104 stderr_stream: None,
105 pid: None,
106 pgid: None,
107 stopped: false,
108 cancel: None,
109 pgids: Vec::new(),
110 }
111 }
112
113 pub fn with_streams(
117 id: JobId,
118 session_id: u64,
119 command: String,
120 rx: oneshot::Receiver<ExecResult>,
121 stdout: Arc<BoundedStream>,
122 stderr: Arc<BoundedStream>,
123 ) -> Self {
124 Self {
125 id,
126 session_id,
127 command,
128 handle: None,
129 result_rx: Some(rx),
130 result: None,
131 output_file: None,
132 persist_output: true,
133 stdout_stream: Some(stdout),
134 stderr_stream: Some(stderr),
135 pid: None,
136 pgid: None,
137 stopped: false,
138 cancel: None,
139 pgids: Vec::new(),
140 }
141 }
142
143 pub fn stopped(id: JobId, session_id: u64, command: String, pid: u32, pgid: u32) -> Self {
145 Self {
146 id,
147 session_id,
148 command,
149 handle: None,
150 result_rx: None,
151 result: None,
152 output_file: None,
153 persist_output: true,
154 stdout_stream: None,
155 stderr_stream: None,
156 pid: Some(pid),
157 pgid: Some(pgid),
158 stopped: true,
159 cancel: None,
160 pgids: Vec::new(),
161 }
162 }
163
164 pub fn output_file(&self) -> Option<&PathBuf> {
166 self.output_file.as_ref()
167 }
168
169 pub fn is_done(&mut self) -> bool {
173 if self.stopped {
174 return false;
175 }
176 self.try_poll();
177 self.result.is_some()
178 }
179
180 pub fn status(&mut self) -> JobStatus {
182 if self.stopped {
183 return JobStatus::Stopped;
184 }
185 self.try_poll();
186 match &self.result {
187 Some(r) if r.ok() => JobStatus::Done,
188 Some(_) => JobStatus::Failed,
189 None => JobStatus::Running,
190 }
191 }
192
193 pub fn status_string(&mut self) -> String {
200 self.try_poll();
201 match &self.result {
202 Some(r) if r.ok() => "done:0".to_string(),
203 Some(r) => format!("failed:{}", r.code),
204 None => "running".to_string(),
205 }
206 }
207
208 pub fn stdout_stream(&self) -> Option<&Arc<BoundedStream>> {
210 self.stdout_stream.as_ref()
211 }
212
213 pub fn stderr_stream(&self) -> Option<&Arc<BoundedStream>> {
215 self.stderr_stream.as_ref()
216 }
217
218 fn write_output_file(&self, result: &ExecResult) -> Option<PathBuf> {
220 let is_bytes = result.is_bytes();
224 let text = if is_bytes {
225 std::borrow::Cow::Borrowed("")
226 } else {
227 result.text_out()
228 };
229 if !is_bytes && text.is_empty() && result.err.is_empty() {
230 return None;
231 }
232
233 let tmp_dir = std::env::temp_dir().join("kaish").join("jobs");
234 if std::fs::create_dir_all(&tmp_dir).is_err() {
235 tracing::warn!("Failed to create job output directory");
236 return None;
237 }
238
239 let filename = format!(
248 "session_{}_job_{}.{}.txt",
249 self.session_id,
250 self.id.0,
251 std::process::id()
252 );
253 let path = tmp_dir.join(filename);
254
255 let mut content = String::new();
256 content.push_str(&format!("# Job {}: {}\n", self.id, self.command));
257 content.push_str(&format!("# Status: {}\n\n", if result.ok() { "Done" } else { "Failed" }));
258
259 if is_bytes {
260 let n = result.out_bytes().map(|b| b.len()).unwrap_or(0);
261 content.push_str(&format!(
262 "## STDOUT\n[binary output: {n} bytes — omitted from this text log]\n"
263 ));
264 } else if !text.is_empty() {
265 content.push_str("## STDOUT\n");
266 content.push_str(&text);
267 if !text.ends_with('\n') {
268 content.push('\n');
269 }
270 }
271
272 if !result.err.is_empty() {
273 content.push_str("\n## STDERR\n");
274 content.push_str(&result.err);
275 if !result.err.ends_with('\n') {
276 content.push('\n');
277 }
278 }
279
280 match std::fs::write(&path, content) {
281 Ok(()) => Some(path),
282 Err(e) => {
283 tracing::warn!("Failed to write job output file: {}", e);
284 None
285 }
286 }
287 }
288
289 pub fn cleanup_files(&mut self) {
291 if let Some(path) = self.output_file.take() {
292 if let Err(e) = std::fs::remove_file(&path) {
293 if e.kind() != io::ErrorKind::NotFound {
295 tracing::warn!("Failed to clean up job output file {}: {}", path.display(), e);
296 }
297 }
298 }
299 }
300
301 pub fn try_result(&self) -> Option<&ExecResult> {
303 self.result.as_ref()
304 }
305
306 pub fn try_poll(&mut self) -> bool {
311 if self.result.is_some() {
312 return true;
313 }
314
315 if let Some(rx) = self.result_rx.as_mut() {
317 match rx.try_recv() {
318 Ok(result) => {
319 self.result = Some(result);
320 self.result_rx = None;
321 return true;
322 }
323 Err(tokio::sync::oneshot::error::TryRecvError::Empty) => {
324 return false;
326 }
327 Err(tokio::sync::oneshot::error::TryRecvError::Closed) => {
328 self.result = Some(ExecResult::failure(1, "job channel closed"));
330 self.result_rx = None;
331 return true;
332 }
333 }
334 }
335
336 if let Some(handle) = self.handle.as_mut()
338 && handle.is_finished() {
339 let Some(mut handle) = self.handle.take() else {
341 return false;
342 };
343 let waker = std::task::Waker::noop();
345 let mut cx = std::task::Context::from_waker(waker);
346 let result = match std::pin::Pin::new(&mut handle).poll(&mut cx) {
347 std::task::Poll::Ready(Ok(r)) => r,
348 std::task::Poll::Ready(Err(e)) => {
349 ExecResult::failure(1, format!("job panicked: {}", e))
350 }
351 std::task::Poll::Pending => return false, };
353 self.result = Some(result);
354 return true;
355 }
356
357 false
358 }
359}
360
361static NEXT_SESSION_ID: AtomicU64 = AtomicU64::new(0);
367
368fn prune_orphaned_job_files() {
380 #[cfg(target_os = "linux")]
382 {
383 let jobs_dir = std::env::temp_dir().join("kaish").join("jobs");
384 let Ok(entries) = std::fs::read_dir(&jobs_dir) else {
385 return; };
387 let current_pid = std::process::id();
388 for entry in entries.flatten() {
389 let name = entry.file_name();
390 let name_str = name.to_string_lossy();
391 let file_pid: Option<u32> = name_str
394 .strip_suffix(".txt")
395 .and_then(|s| s.rsplit_once('.'))
396 .and_then(|(_, pid_str)| pid_str.parse().ok());
397 let Some(pid) = file_pid else {
398 continue; };
400 if pid == current_pid {
401 continue; }
403 if std::path::Path::new(&format!("/proc/{}", pid)).exists() {
405 continue; }
407 let _ = std::fs::remove_file(entry.path());
409 }
410 }
411}
412
413pub struct JobManager {
415 session_id: u64,
417 next_id: AtomicU64,
419 jobs: Arc<Mutex<HashMap<JobId, Job>>>,
421 persist_output_files: std::sync::atomic::AtomicBool,
427}
428
429impl JobManager {
430 pub fn new() -> Self {
445 static PRUNE_ONCE: std::sync::Once = std::sync::Once::new();
450 PRUNE_ONCE.call_once(prune_orphaned_job_files);
451 Self {
452 session_id: NEXT_SESSION_ID.fetch_add(1, Ordering::SeqCst),
453 next_id: AtomicU64::new(1),
454 jobs: Arc::new(Mutex::new(HashMap::new())),
455 persist_output_files: std::sync::atomic::AtomicBool::new(true),
456 }
457 }
458
459 pub fn set_persist_output_files(&self, on: bool) {
469 self.persist_output_files.store(on, Ordering::Relaxed);
470 }
471
472 pub fn persist_output_files(&self) -> bool {
474 self.persist_output_files.load(Ordering::Relaxed)
475 }
476
477 pub async fn spawn<F>(&self, command: String, future: F) -> JobId
482 where
483 F: std::future::Future<Output = ExecResult> + Send + 'static,
484 {
485 let id = JobId(self.next_id.fetch_add(1, Ordering::SeqCst));
486 let handle = tokio::spawn(crate::telemetry::bind_current_context(future));
489 let mut job = Job::new(id, self.session_id, command, handle);
490 job.persist_output = self.persist_output_files();
491
492 self.jobs.lock().await.insert(id, job);
499
500 id
501 }
502
503 pub async fn register(&self, command: String, rx: oneshot::Receiver<ExecResult>) -> JobId {
505 let id = JobId(self.next_id.fetch_add(1, Ordering::SeqCst));
506 let mut job = Job::from_channel(id, self.session_id, command, rx);
507 job.persist_output = self.persist_output_files();
508
509 let mut jobs = self.jobs.lock().await;
510 jobs.insert(id, job);
511
512 id
513 }
514
515 pub async fn register_with_streams(
519 &self,
520 command: String,
521 rx: oneshot::Receiver<ExecResult>,
522 stdout: Arc<BoundedStream>,
523 stderr: Arc<BoundedStream>,
524 ) -> JobId {
525 let id = JobId(self.next_id.fetch_add(1, Ordering::SeqCst));
526 let mut job = Job::with_streams(id, self.session_id, command, rx, stdout, stderr);
527 job.persist_output = self.persist_output_files();
528
529 let mut jobs = self.jobs.lock().await;
530 jobs.insert(id, job);
531
532 id
533 }
534
535 pub async fn wait(&self, id: JobId) -> Option<ExecResult> {
545 loop {
557 {
558 let mut jobs = self.jobs.lock().await;
559 let job = jobs.get_mut(&id)?;
560 if job.is_done() {
561 let result = job
562 .result
563 .clone()
564 .unwrap_or_else(|| ExecResult::failure(1, "no result"));
565 if job.persist_output
567 && job.output_file.is_none()
568 && let Some(path) = job.write_output_file(&result)
569 {
570 job.output_file = Some(path);
571 }
572 return Some(result);
573 }
574 }
575 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
577 }
578 }
579
580 pub async fn wait_all(&self) -> Vec<(JobId, ExecResult)> {
582 let mut results = Vec::new();
583
584 let ids: Vec<JobId> = {
586 let jobs = self.jobs.lock().await;
587 jobs.keys().copied().collect()
588 };
589
590 for id in ids {
591 if let Some(result) = self.wait(id).await {
592 results.push((id, result));
593 }
594 }
595
596 results
597 }
598
599 pub async fn list(&self) -> Vec<JobInfo> {
601 let mut jobs = self.jobs.lock().await;
602 jobs.values_mut()
603 .map(|job| JobInfo {
604 id: job.id,
605 command: job.command.clone(),
606 status: job.status(),
607 output_file: job.output_file.clone(),
608 pid: job.pid,
609 })
610 .collect()
611 }
612
613 pub async fn running_count(&self) -> usize {
615 let mut jobs = self.jobs.lock().await;
616 let mut count = 0;
617 for job in jobs.values_mut() {
618 if !job.is_done() {
619 count += 1;
620 }
621 }
622 count
623 }
624
625 pub async fn cleanup(&self) {
627 let mut jobs = self.jobs.lock().await;
628 jobs.retain(|_, job| {
629 if job.is_done() {
630 job.cleanup_files();
631 false
632 } else {
633 true
634 }
635 });
636 }
637
638 pub async fn exists(&self, id: JobId) -> bool {
640 let jobs = self.jobs.lock().await;
641 jobs.contains_key(&id)
642 }
643
644 pub async fn get(&self, id: JobId) -> Option<JobInfo> {
646 let mut jobs = self.jobs.lock().await;
647 jobs.get_mut(&id).map(|job| JobInfo {
648 id: job.id,
649 command: job.command.clone(),
650 status: job.status(),
651 output_file: job.output_file.clone(),
652 pid: job.pid,
653 })
654 }
655
656 pub async fn get_command(&self, id: JobId) -> Option<String> {
658 let jobs = self.jobs.lock().await;
659 jobs.get(&id).map(|job| job.command.clone())
660 }
661
662 pub async fn get_status_string(&self, id: JobId) -> Option<String> {
664 let mut jobs = self.jobs.lock().await;
665 jobs.get_mut(&id).map(|job| job.status_string())
666 }
667
668 pub async fn read_stdout(&self, id: JobId) -> Option<Vec<u8>> {
672 let jobs = self.jobs.lock().await;
673 if let Some(job) = jobs.get(&id)
674 && let Some(stream) = job.stdout_stream() {
675 return Some(stream.read().await);
676 }
677 None
678 }
679
680 pub async fn read_stderr(&self, id: JobId) -> Option<Vec<u8>> {
684 let jobs = self.jobs.lock().await;
685 if let Some(job) = jobs.get(&id)
686 && let Some(stream) = job.stderr_stream() {
687 return Some(stream.read().await);
688 }
689 None
690 }
691
692 pub async fn list_ids(&self) -> Vec<JobId> {
694 let jobs = self.jobs.lock().await;
695 jobs.keys().copied().collect()
696 }
697
698 pub async fn register_stopped(&self, command: String, pid: u32, pgid: u32) -> JobId {
700 let id = JobId(self.next_id.fetch_add(1, Ordering::SeqCst));
701 let job = Job::stopped(id, self.session_id, command, pid, pgid);
702 let mut jobs = self.jobs.lock().await;
703 jobs.insert(id, job);
704 id
705 }
706
707 pub async fn stop_job(&self, id: JobId, pid: u32, pgid: u32) {
709 let mut jobs = self.jobs.lock().await;
710 if let Some(job) = jobs.get_mut(&id) {
711 job.stopped = true;
712 job.pid = Some(pid);
713 job.pgid = Some(pgid);
714 }
715 }
716
717 pub async fn resume_job(&self, id: JobId) {
719 let mut jobs = self.jobs.lock().await;
720 if let Some(job) = jobs.get_mut(&id) {
721 job.stopped = false;
722 }
723 }
724
725 pub async fn last_stopped(&self) -> Option<JobId> {
727 let mut jobs = self.jobs.lock().await;
728 let mut best: Option<JobId> = None;
730 for job in jobs.values_mut() {
731 if job.stopped {
732 match best {
733 None => best = Some(job.id),
734 Some(b) if job.id.0 > b.0 => best = Some(job.id),
735 _ => {}
736 }
737 }
738 }
739 best
740 }
741
742 pub async fn get_process_info(&self, id: JobId) -> Option<(u32, u32)> {
744 let jobs = self.jobs.lock().await;
745 jobs.get(&id).and_then(|job| {
746 match (job.pid, job.pgid) {
747 (Some(pid), Some(pgid)) => Some((pid, pgid)),
748 _ => None,
749 }
750 })
751 }
752
753 pub async fn set_cancel_token(&self, id: JobId, token: tokio_util::sync::CancellationToken) {
757 let mut jobs = self.jobs.lock().await;
758 if let Some(job) = jobs.get_mut(&id) {
759 job.cancel = Some(token);
760 }
761 }
762
763 pub async fn cancel(&self, id: JobId) -> bool {
767 let jobs = self.jobs.lock().await;
768 match jobs.get(&id).and_then(|job| job.cancel.clone()) {
769 Some(token) => {
770 token.cancel();
771 true
772 }
773 None => false,
774 }
775 }
776
777 pub async fn add_pgid(&self, id: JobId, pgid: u32) {
781 let mut jobs = self.jobs.lock().await;
782 if let Some(job) = jobs.get_mut(&id) {
783 if !job.pgids.contains(&pgid) {
784 job.pgids.push(pgid);
785 }
786 }
787 }
788
789 pub async fn job_pgids(&self, id: JobId) -> Vec<u32> {
793 let jobs = self.jobs.lock().await;
794 jobs.get(&id)
795 .map(|job| {
796 let mut v = job.pgids.clone();
797 if let Some(pg) = job.pgid {
798 if !v.contains(&pg) {
799 v.push(pg);
800 }
801 }
802 v
803 })
804 .unwrap_or_default()
805 }
806
807 pub async fn remove(&self, id: JobId) {
809 let mut jobs = self.jobs.lock().await;
810 if let Some(mut job) = jobs.remove(&id) {
811 job.cleanup_files();
812 }
813 }
814}
815
816impl Default for JobManager {
817 fn default() -> Self {
818 Self::new()
819 }
820}
821
822#[cfg(test)]
823mod tests {
824 use super::*;
825 use std::time::Duration;
826
827 #[tokio::test]
828 async fn test_no_host_output_file_when_persistence_disabled() {
829 let manager = JobManager::new();
833 assert!(manager.persist_output_files(), "default is to persist");
834 manager.set_persist_output_files(false);
835 assert!(!manager.persist_output_files());
836
837 let id = manager.spawn("leaky".to_string(), async {
838 ExecResult::success("output that must not hit host disk")
839 }).await;
840 tokio::time::sleep(Duration::from_millis(10)).await;
841 let result = manager.wait(id).await;
842 assert!(result.is_some());
843
844 let output_file = {
846 let jobs = manager.jobs.lock().await;
847 jobs.get(&id).and_then(|j| j.output_file().cloned())
848 };
849 assert!(
850 output_file.is_none(),
851 "no host output file should be written when persistence is disabled, got {output_file:?}"
852 );
853 }
854
855 #[tokio::test]
856 async fn test_spawn_and_wait() {
857 let manager = JobManager::new();
858
859 let id = manager.spawn("test".to_string(), async {
860 tokio::time::sleep(Duration::from_millis(10)).await;
861 ExecResult::success("done")
862 }).await;
863
864 tokio::time::sleep(Duration::from_millis(5)).await;
866
867 let result = manager.wait(id).await;
868 assert!(result.is_some());
869 let result = result.unwrap();
870 assert!(result.ok());
871 assert_eq!(&*result.text_out(), "done");
872 }
873
874 #[tokio::test]
875 async fn test_wait_all() {
876 let manager = JobManager::new();
877
878 manager.spawn("job1".to_string(), async {
879 tokio::time::sleep(Duration::from_millis(10)).await;
880 ExecResult::success("one")
881 }).await;
882
883 manager.spawn("job2".to_string(), async {
884 tokio::time::sleep(Duration::from_millis(5)).await;
885 ExecResult::success("two")
886 }).await;
887
888 tokio::time::sleep(Duration::from_millis(5)).await;
890
891 let results = manager.wait_all().await;
892 assert_eq!(results.len(), 2);
893 }
894
895 #[tokio::test]
896 async fn test_list_jobs() {
897 let manager = JobManager::new();
898
899 manager.spawn("test job".to_string(), async {
900 tokio::time::sleep(Duration::from_millis(50)).await;
901 ExecResult::success("")
902 }).await;
903
904 tokio::time::sleep(Duration::from_millis(5)).await;
906
907 let jobs = manager.list().await;
908 assert_eq!(jobs.len(), 1);
909 assert_eq!(jobs[0].command, "test job");
910 assert_eq!(jobs[0].status, JobStatus::Running);
911 }
912
913 #[tokio::test]
914 async fn test_job_status_after_completion() {
915 let manager = JobManager::new();
916
917 let id = manager.spawn("quick".to_string(), async {
918 ExecResult::success("")
919 }).await;
920
921 tokio::time::sleep(Duration::from_millis(10)).await;
923 let _ = manager.wait(id).await;
924
925 let info = manager.get(id).await;
926 assert!(info.is_some());
927 assert_eq!(info.unwrap().status, JobStatus::Done);
928 }
929
930 #[tokio::test]
931 async fn test_cleanup() {
932 let manager = JobManager::new();
933
934 let id = manager.spawn("done".to_string(), async {
935 ExecResult::success("")
936 }).await;
937
938 tokio::time::sleep(Duration::from_millis(10)).await;
940 let _ = manager.wait(id).await;
941
942 assert_eq!(manager.list().await.len(), 1);
944
945 manager.cleanup().await;
947
948 assert_eq!(manager.list().await.len(), 0);
950 }
951
952 #[tokio::test]
953 async fn test_cleanup_removes_temp_files() {
954 let manager = JobManager::new();
956
957 let id = manager.spawn("output job".to_string(), async {
958 ExecResult::success("some output that gets written to a temp file")
959 }).await;
960
961 tokio::time::sleep(Duration::from_millis(10)).await;
963 let result = manager.wait(id).await;
964 assert!(result.is_some());
965
966 let output_file = {
970 let jobs = manager.jobs.lock().await;
971 jobs.get(&id).and_then(|j| j.output_file().cloned())
972 };
973 let path = output_file.expect("job with output should have written a temp file");
974 assert!(path.exists(), "temp file should exist before cleanup: {}", path.display());
975
976 manager.cleanup().await;
978
979 assert!(
980 !path.exists(),
981 "temp file should be removed after cleanup: {}",
982 path.display()
983 );
984 }
985
986 #[tokio::test]
987 async fn test_register_with_channel() {
988 let manager = JobManager::new();
989 let (tx, rx) = oneshot::channel();
990
991 let id = manager.register("channel job".to_string(), rx).await;
992
993 tx.send(ExecResult::success("from channel")).unwrap();
995
996 let result = manager.wait(id).await;
997 assert!(result.is_some());
998 assert_eq!(&*result.unwrap().text_out(), "from channel");
999 }
1000
1001 #[tokio::test]
1002 async fn test_spawn_immediately_available() {
1003 let manager = JobManager::new();
1005
1006 let id = manager.spawn("instant".to_string(), async {
1007 tokio::time::sleep(Duration::from_millis(100)).await;
1008 ExecResult::success("done")
1009 }).await;
1010
1011 let exists = manager.exists(id).await;
1013 assert!(exists, "job should be immediately available after spawn()");
1014
1015 let info = manager.get(id).await;
1016 assert!(info.is_some(), "job info should be available immediately");
1017 }
1018
1019 #[tokio::test]
1020 async fn test_nonexistent_job() {
1021 let manager = JobManager::new();
1022 let result = manager.wait(JobId(999)).await;
1023 assert!(result.is_none());
1024 }
1025
1026 #[tokio::test]
1027 async fn test_cancel_token_fires() {
1028 let manager = JobManager::new();
1031 let token = tokio_util::sync::CancellationToken::new();
1032 let id = manager.spawn("bg".to_string(), async { ExecResult::success("") }).await;
1033 manager.set_cancel_token(id, token.clone()).await;
1034
1035 assert!(!token.is_cancelled());
1036 assert!(manager.cancel(id).await, "cancel should report success");
1037 assert!(token.is_cancelled(), "the job's token must be tripped");
1038 }
1039
1040 #[tokio::test]
1041 async fn test_cancel_without_token_returns_false() {
1042 let manager = JobManager::new();
1043 let id = manager.spawn("bg".to_string(), async { ExecResult::success("") }).await;
1044 assert!(!manager.cancel(id).await);
1046 assert!(!manager.cancel(JobId(999)).await);
1048 }
1049
1050 #[tokio::test]
1051 async fn test_pgids_recorded_and_deduped() {
1052 let manager = JobManager::new();
1053 let id = manager.spawn("bg".to_string(), async { ExecResult::success("") }).await;
1054 assert!(manager.job_pgids(id).await.is_empty());
1055
1056 manager.add_pgid(id, 4242).await;
1057 manager.add_pgid(id, 4243).await;
1058 manager.add_pgid(id, 4242).await; assert_eq!(manager.job_pgids(id).await, vec![4242, 4243]);
1060
1061 assert!(manager.job_pgids(JobId(999)).await.is_empty());
1063 }
1064
1065 #[tokio::test]
1066 async fn wait_does_not_block_other_job_ops() {
1067 let manager = Arc::new(JobManager::new());
1074 manager.set_persist_output_files(false);
1075
1076 let (tx, rx) = oneshot::channel::<()>();
1078 let id = manager
1079 .spawn("blocker".to_string(), async move {
1080 let _ = rx.await;
1081 ExecResult::success("done")
1082 })
1083 .await;
1084
1085 let waiter = {
1088 let m = manager.clone();
1089 tokio::spawn(async move { m.wait(id).await })
1090 };
1091 tokio::time::sleep(Duration::from_millis(50)).await;
1093
1094 let listed = tokio::time::timeout(Duration::from_secs(2), manager.list()).await;
1096 assert!(
1097 listed.is_ok(),
1098 "list() blocked while wait() was parked — jobs lock held across await"
1099 );
1100 let second = tokio::time::timeout(
1101 Duration::from_secs(2),
1102 manager.spawn("second".to_string(), async { ExecResult::success("2") }),
1103 )
1104 .await;
1105 assert!(
1106 second.is_ok(),
1107 "spawn() blocked/spun while wait() was parked"
1108 );
1109
1110 let _ = tx.send(());
1112 let result = tokio::time::timeout(Duration::from_secs(2), waiter)
1113 .await
1114 .expect("waiter join timed out")
1115 .expect("waiter task panicked");
1116 assert_eq!(result.map(|r| r.code), Some(0), "waiter should see exit 0");
1117 }
1118
1119 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1120 async fn wait_survives_a_dropped_waiter() {
1121 let manager = Arc::new(JobManager::new());
1128 manager.set_persist_output_files(false);
1129
1130 let (tx, rx) = oneshot::channel::<()>();
1131 let id = manager
1132 .spawn("blocker".to_string(), async move {
1133 let _ = rx.await;
1134 ExecResult::success("done")
1135 })
1136 .await;
1137
1138 {
1140 let m = manager.clone();
1141 let a = tokio::spawn(async move { m.wait(id).await });
1142 tokio::time::sleep(Duration::from_millis(20)).await;
1143 a.abort();
1144 let _ = a.await;
1145 }
1146
1147 let _ = tx.send(());
1149
1150 let res = tokio::time::timeout(Duration::from_secs(2), manager.wait(id))
1152 .await
1153 .expect("wait must not hang after a prior waiter was dropped");
1154 assert_eq!(res.map(|r| r.code), Some(0), "B should see the completed job");
1155 }
1156}