1use crate::events::EventBroadcaster;
19use mold_core::ServerEvent;
20use serde::Serialize;
21use std::sync::{Arc, RwLock};
22use std::time::{SystemTime, UNIX_EPOCH};
23use tokio::sync::Notify;
24
25#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, utoipa::ToSchema)]
32#[serde(rename_all = "snake_case")]
33pub enum JobLifecycle {
34 Queued,
35 Running,
36}
37
38#[derive(Debug, Clone, Serialize, utoipa::ToSchema)]
48pub struct JobEntry {
49 pub id: String,
50 pub model: String,
51 pub state: JobLifecycle,
52 pub started_at_unix_ms: u64,
53 pub position: usize,
54 #[serde(skip_serializing_if = "Option::is_none")]
56 pub gpu: Option<usize>,
57 #[serde(skip_serializing_if = "Option::is_none")]
59 pub target_gpu: Option<usize>,
60 #[serde(skip_serializing_if = "Option::is_none")]
63 pub seed_pinned: Option<bool>,
64 #[serde(skip_serializing_if = "Option::is_none")]
68 #[schema(value_type = Object)]
69 pub metadata: Option<Box<mold_core::OutputMetadata>>,
70}
71
72#[derive(Debug, Clone, Serialize, utoipa::ToSchema)]
76pub struct QueueListing {
77 pub entries: Vec<JobEntry>,
78}
79
80#[derive(Debug, Clone)]
81struct EntryInternal {
82 id: String,
83 model: String,
84 state: JobLifecycle,
85 started_at_unix_ms: u64,
86 gpu: Option<usize>,
87 target_gpu: Option<usize>,
88 seed_pinned: Option<bool>,
89 metadata: Option<Box<mold_core::OutputMetadata>>,
90 cancel: Arc<Notify>,
96}
97
98#[derive(Debug, Clone, Copy, PartialEq, Eq)]
99pub enum TargetGpuUpdateError {
100 NotFound,
101 AlreadyRunning,
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq)]
106pub enum QueuedJobCancelError {
107 NotFound,
108 AlreadyRunning,
109}
110
111#[derive(Debug, Clone, Copy, PartialEq, Eq)]
114pub enum QueueReorderError {
115 NotFound,
116 AlreadyRunning,
117}
118
119pub struct JobRegistry {
124 inner: RwLock<Vec<EntryInternal>>,
125 events: Option<Arc<EventBroadcaster>>,
130}
131
132pub type SharedJobRegistry = Arc<JobRegistry>;
134
135impl JobRegistry {
136 pub fn new() -> SharedJobRegistry {
137 Arc::new(Self {
138 inner: RwLock::new(Vec::new()),
139 events: None,
140 })
141 }
142
143 pub fn with_events(events: Arc<EventBroadcaster>) -> SharedJobRegistry {
146 Arc::new(Self {
147 inner: RwLock::new(Vec::new()),
148 events: Some(events),
149 })
150 }
151
152 fn emit(&self, event: ServerEvent) {
156 if let Some(events) = &self.events {
157 events.publish(event);
158 }
159 }
160
161 pub fn register(&self, id: impl Into<String>, model: impl Into<String>) -> Arc<Notify> {
164 self.register_with_target_gpu(id, model, None)
165 }
166
167 pub fn register_with_target_gpu(
175 &self,
176 id: impl Into<String>,
177 model: impl Into<String>,
178 target_gpu: Option<usize>,
179 ) -> Arc<Notify> {
180 self.register_job(id, model, target_gpu, None, None)
181 }
182
183 pub fn register_job(
186 &self,
187 id: impl Into<String>,
188 model: impl Into<String>,
189 target_gpu: Option<usize>,
190 seed_pinned: Option<bool>,
191 metadata: Option<Box<mold_core::OutputMetadata>>,
192 ) -> Arc<Notify> {
193 let started_at_unix_ms = SystemTime::now()
194 .duration_since(UNIX_EPOCH)
195 .unwrap_or_default()
196 .as_millis() as u64;
197 let id = id.into();
198 let model = model.into();
199 let cancel = Arc::new(Notify::new());
200 {
201 let mut entries = self.inner.write().unwrap_or_else(|e| e.into_inner());
202 entries.push(EntryInternal {
203 id: id.clone(),
204 model: model.clone(),
205 state: JobLifecycle::Queued,
206 started_at_unix_ms,
207 gpu: None,
208 target_gpu,
209 seed_pinned,
210 metadata,
211 cancel: cancel.clone(),
212 });
213 }
214 self.emit(ServerEvent::JobQueued { id, model });
215 cancel
216 }
217
218 pub fn cancel_queued(&self, id: &str) -> Result<(), QueuedJobCancelError> {
228 {
229 let mut entries = self.inner.write().unwrap_or_else(|e| e.into_inner());
230 let Some(pos) = entries.iter().position(|e| e.id == id) else {
231 return Err(QueuedJobCancelError::NotFound);
232 };
233 if entries[pos].state == JobLifecycle::Running {
234 return Err(QueuedJobCancelError::AlreadyRunning);
235 }
236 let entry = entries.remove(pos);
237 entry.cancel.notify_one();
238 }
239 self.emit(ServerEvent::JobEnded { id: id.to_string() });
240 Ok(())
241 }
242
243 pub fn cancel_all_queued(&self) -> usize {
250 let cancelled_ids = {
251 let mut entries = self.inner.write().unwrap_or_else(|e| e.into_inner());
252 let mut ids = Vec::new();
253 entries.retain(|e| {
254 if e.state == JobLifecycle::Queued {
255 e.cancel.notify_one();
256 ids.push(e.id.clone());
257 false
258 } else {
259 true
260 }
261 });
262 ids
263 };
264 for id in &cancelled_ids {
265 self.emit(ServerEvent::JobEnded { id: id.clone() });
266 }
267 cancelled_ids.len()
268 }
269
270 pub fn mark_running(&self, id: &str, gpu: Option<usize>) {
273 let model = {
274 let mut entries = self.inner.write().unwrap_or_else(|e| e.into_inner());
275 entries.iter_mut().find(|e| e.id == id).map(|e| {
276 e.state = JobLifecycle::Running;
277 e.gpu = gpu;
278 e.target_gpu = None;
279 e.model.clone()
280 })
281 };
282 if let Some(model) = model {
283 self.emit(ServerEvent::JobStarted {
284 id: id.to_string(),
285 model,
286 gpu,
287 });
288 }
289 }
290
291 pub fn set_target_gpu(
292 &self,
293 id: &str,
294 target_gpu: Option<usize>,
295 ) -> Result<(), TargetGpuUpdateError> {
296 let mut entries = self.inner.write().unwrap_or_else(|e| e.into_inner());
297 let Some(e) = entries.iter_mut().find(|e| e.id == id) else {
298 return Err(TargetGpuUpdateError::NotFound);
299 };
300 if e.state == JobLifecycle::Running {
301 return Err(TargetGpuUpdateError::AlreadyRunning);
302 }
303 e.target_gpu = target_gpu;
304 Ok(())
305 }
306
307 pub fn reorder_queued(&self, id: &str, position: usize) -> Result<(), QueueReorderError> {
321 let mut entries = self.inner.write().unwrap_or_else(|e| e.into_inner());
322 let Some(from) = entries.iter().position(|e| e.id == id) else {
323 return Err(QueueReorderError::NotFound);
324 };
325 if entries[from].state == JobLifecycle::Running {
326 return Err(QueueReorderError::AlreadyRunning);
327 }
328 let queued_slots: Vec<usize> = entries
331 .iter()
332 .enumerate()
333 .filter(|(_, e)| e.state == JobLifecycle::Queued)
334 .map(|(i, _)| i)
335 .collect();
336 let cur = queued_slots
338 .iter()
339 .position(|&i| i == from)
340 .expect("queued job must appear in its own queued-slot list");
341 let target = position.min(queued_slots.len() - 1);
342 if cur == target {
343 return Ok(());
344 }
345 let entry = entries.remove(from);
348 let remaining_slots: Vec<usize> = entries
349 .iter()
350 .enumerate()
351 .filter(|(_, e)| e.state == JobLifecycle::Queued)
352 .map(|(i, _)| i)
353 .collect();
354 let dest = remaining_slots
355 .get(target)
356 .copied()
357 .unwrap_or(entries.len());
358 entries.insert(dest, entry);
359 Ok(())
360 }
361
362 pub fn target_gpu(&self, id: &str) -> Option<Option<usize>> {
363 let entries = self.inner.read().unwrap_or_else(|e| e.into_inner());
364 entries.iter().find(|e| e.id == id).map(|e| e.target_gpu)
365 }
366
367 pub fn entry(&self, id: &str) -> Option<JobEntry> {
368 let entries = self.inner.read().unwrap_or_else(|e| e.into_inner());
369 entries.iter().enumerate().find_map(|(i, e)| {
370 (e.id == id).then(|| JobEntry {
371 id: e.id.clone(),
372 model: e.model.clone(),
373 state: e.state,
374 started_at_unix_ms: e.started_at_unix_ms,
375 position: i,
376 gpu: e.gpu,
377 target_gpu: e.target_gpu,
378 seed_pinned: e.seed_pinned,
379 metadata: e.metadata.clone(),
380 })
381 })
382 }
383
384 pub fn remove(&self, id: &str) {
387 if id.is_empty() {
388 return;
389 }
390 let removed = {
391 let mut entries = self.inner.write().unwrap_or_else(|e| e.into_inner());
392 let before = entries.len();
393 entries.retain(|e| e.id != id);
394 entries.len() != before
395 };
396 if removed {
400 self.emit(ServerEvent::JobEnded { id: id.to_string() });
401 }
402 }
403
404 pub fn queued_ids_in_order(&self) -> Vec<String> {
414 let entries = self.inner.read().unwrap_or_else(|e| e.into_inner());
415 entries
416 .iter()
417 .filter(|e| e.state == JobLifecycle::Queued)
418 .map(|e| e.id.clone())
419 .collect()
420 }
421
422 pub fn snapshot(&self) -> QueueListing {
426 let entries = self.inner.read().unwrap_or_else(|e| e.into_inner());
427 let out = entries
428 .iter()
429 .enumerate()
430 .map(|(i, e)| JobEntry {
431 id: e.id.clone(),
432 model: e.model.clone(),
433 state: e.state,
434 started_at_unix_ms: e.started_at_unix_ms,
435 position: i,
436 gpu: e.gpu,
437 target_gpu: e.target_gpu,
438 seed_pinned: e.seed_pinned,
439 metadata: e.metadata.clone(),
440 })
441 .collect();
442 QueueListing { entries: out }
443 }
444
445 pub fn len(&self) -> usize {
447 self.inner.read().unwrap_or_else(|e| e.into_inner()).len()
448 }
449
450 pub fn is_empty(&self) -> bool {
454 self.len() == 0
455 }
456}
457
458#[cfg(test)]
459mod tests {
460 use super::*;
461
462 #[test]
463 fn register_appends_in_fifo_order_with_queued_state() {
464 let reg = JobRegistry::new();
465 reg.register("a", "flux-dev:fp16");
466 reg.register("b", "sdxl:q8");
467 let snap = reg.snapshot();
468 assert_eq!(snap.entries.len(), 2);
469 assert_eq!(snap.entries[0].id, "a");
470 assert_eq!(snap.entries[0].position, 0);
471 assert_eq!(snap.entries[0].state, JobLifecycle::Queued);
472 assert_eq!(snap.entries[1].id, "b");
473 assert_eq!(snap.entries[1].position, 1);
474 }
475
476 #[test]
477 fn mark_running_flips_state_and_records_gpu_ordinal() {
478 let reg = JobRegistry::new();
479 reg.register("a", "flux-dev:fp16");
480 reg.mark_running("a", Some(1));
481 let snap = reg.snapshot();
482 assert_eq!(snap.entries[0].state, JobLifecycle::Running);
483 assert_eq!(snap.entries[0].gpu, Some(1));
484 }
485
486 #[test]
487 fn queued_entries_can_carry_target_gpu_metadata() {
488 let reg = JobRegistry::new();
489 reg.register_with_target_gpu("a", "flux-dev:fp16", Some(1));
490 let snap = reg.snapshot();
491 assert_eq!(snap.entries[0].state, JobLifecycle::Queued);
492 assert_eq!(snap.entries[0].target_gpu, Some(1));
493 assert_eq!(snap.entries[0].gpu, None);
494 }
495
496 #[test]
497 fn target_gpu_updates_only_apply_to_queued_entries() {
498 let reg = JobRegistry::new();
499 reg.register("a", "flux-dev:fp16");
500 reg.set_target_gpu("a", Some(1)).unwrap();
501 assert_eq!(reg.target_gpu("a"), Some(Some(1)));
502
503 reg.mark_running("a", Some(1));
504 let err = reg.set_target_gpu("a", None).unwrap_err();
505 assert_eq!(err, TargetGpuUpdateError::AlreadyRunning);
506 assert_eq!(reg.target_gpu("a"), Some(None));
507 }
508
509 fn queued_order(reg: &JobRegistry) -> Vec<String> {
510 reg.snapshot().entries.into_iter().map(|e| e.id).collect()
511 }
512
513 #[test]
514 fn reorder_moves_queued_job_to_the_front() {
515 let reg = JobRegistry::new();
516 reg.register("a", "flux-dev:fp16");
517 reg.register("b", "sdxl:q8");
518 reg.register("c", "ltx-video:q8");
519 reg.reorder_queued("c", 0).unwrap();
520 assert_eq!(queued_order(®), ["c", "a", "b"]);
521 let snap = reg.snapshot();
523 assert_eq!(snap.entries[0].position, 0);
524 assert_eq!(snap.entries[2].position, 2);
525 }
526
527 #[test]
528 fn reorder_moves_queued_job_to_the_back() {
529 let reg = JobRegistry::new();
530 reg.register("a", "flux-dev:fp16");
531 reg.register("b", "sdxl:q8");
532 reg.register("c", "ltx-video:q8");
533 reg.reorder_queued("a", 2).unwrap();
534 assert_eq!(queued_order(®), ["b", "c", "a"]);
535 }
536
537 #[test]
538 fn reorder_moves_a_middle_job_up_and_preserves_the_others_order() {
539 let reg = JobRegistry::new();
540 for id in ["a", "b", "c", "d", "e"] {
541 reg.register(id, "flux-dev:fp16");
542 }
543 reg.reorder_queued("e", 1).unwrap();
545 assert_eq!(queued_order(®), ["a", "e", "b", "c", "d"]);
546 }
547
548 #[test]
549 fn reorder_clamps_an_out_of_range_position_to_the_back() {
550 let reg = JobRegistry::new();
551 reg.register("a", "flux-dev:fp16");
552 reg.register("b", "sdxl:q8");
553 reg.register("c", "ltx-video:q8");
554 reg.reorder_queued("a", 999).unwrap();
555 assert_eq!(queued_order(®), ["b", "c", "a"]);
556 }
557
558 #[test]
559 fn reorder_to_the_current_position_is_a_noop() {
560 let reg = JobRegistry::new();
561 reg.register("a", "flux-dev:fp16");
562 reg.register("b", "sdxl:q8");
563 reg.register("c", "ltx-video:q8");
564 reg.reorder_queued("b", 1).unwrap();
565 assert_eq!(queued_order(®), ["a", "b", "c"]);
566 }
567
568 #[test]
569 fn reorder_indexes_among_queued_only_and_keeps_running_jobs_in_place() {
570 let reg = JobRegistry::new();
574 reg.register("a", "flux-dev:fp16");
575 reg.register("b", "sdxl:q8");
576 reg.register("c", "ltx-video:q8");
577 reg.mark_running("a", Some(0));
578 reg.reorder_queued("c", 0).unwrap();
580 assert_eq!(queued_order(®), ["a", "c", "b"]);
581 let snap = reg.snapshot();
582 assert_eq!(snap.entries[0].state, JobLifecycle::Running);
583 assert_eq!(snap.entries[1].id, "c");
584 assert_eq!(snap.entries[1].state, JobLifecycle::Queued);
585 }
586
587 #[test]
588 fn reorder_rejects_a_running_job() {
589 let reg = JobRegistry::new();
590 reg.register("a", "flux-dev:fp16");
591 reg.mark_running("a", Some(0));
592 let err = reg.reorder_queued("a", 0).unwrap_err();
593 assert_eq!(err, QueueReorderError::AlreadyRunning);
594 }
595
596 #[test]
597 fn reorder_unknown_id_is_not_found() {
598 let reg = JobRegistry::new();
599 reg.register("a", "flux-dev:fp16");
600 let err = reg.reorder_queued("ghost", 0).unwrap_err();
601 assert_eq!(err, QueueReorderError::NotFound);
602 }
603
604 #[test]
605 fn mark_running_is_a_noop_for_unknown_ids() {
606 let reg = JobRegistry::new();
607 reg.register("a", "flux-dev:fp16");
608 reg.mark_running("not-here", Some(0));
610 let snap = reg.snapshot();
611 assert_eq!(snap.entries.len(), 1);
612 assert_eq!(snap.entries[0].state, JobLifecycle::Queued);
613 }
614
615 #[test]
616 fn remove_compacts_positions_for_the_survivors() {
617 let reg = JobRegistry::new();
618 reg.register("a", "flux-dev:fp16");
619 reg.register("b", "sdxl:q8");
620 reg.register("c", "ltx-video:q8");
621 reg.remove("b");
622 let snap = reg.snapshot();
623 assert_eq!(snap.entries.len(), 2);
624 assert_eq!(snap.entries[0].id, "a");
625 assert_eq!(snap.entries[0].position, 0);
626 assert_eq!(snap.entries[1].id, "c");
627 assert_eq!(snap.entries[1].position, 1);
628 }
629
630 #[test]
631 fn remove_is_idempotent_and_ignores_empty_ids() {
632 let reg = JobRegistry::new();
637 reg.register("a", "flux-dev:fp16");
638 reg.remove("a");
639 reg.remove("a"); reg.remove("");
641 reg.remove("never-existed");
642 assert!(reg.is_empty());
643 }
644
645 #[test]
646 fn cancel_queued_removes_the_entry() {
647 let reg = JobRegistry::new();
648 reg.register("a", "flux-dev:fp16");
649 reg.register("b", "sdxl:q8");
650 reg.cancel_queued("a").unwrap();
651 let snap = reg.snapshot();
652 assert_eq!(snap.entries.len(), 1);
653 assert_eq!(snap.entries[0].id, "b");
654 assert_eq!(snap.entries[0].position, 0);
655 }
656
657 #[test]
658 fn cancel_queued_rejects_running_jobs_and_keeps_the_entry() {
659 let reg = JobRegistry::new();
660 reg.register("a", "flux-dev:fp16");
661 reg.mark_running("a", Some(0));
662 let err = reg.cancel_queued("a").unwrap_err();
663 assert_eq!(err, QueuedJobCancelError::AlreadyRunning);
664 assert_eq!(reg.len(), 1, "running entry must survive a cancel attempt");
665 }
666
667 #[test]
668 fn cancel_queued_unknown_id_is_not_found() {
669 let reg = JobRegistry::new();
670 let err = reg.cancel_queued("never-existed").unwrap_err();
671 assert_eq!(err, QueuedJobCancelError::NotFound);
672 }
673
674 #[test]
675 fn cancel_all_queued_removes_only_queued_and_returns_count() {
676 let reg = JobRegistry::new();
677 reg.register("a", "flux-dev:fp16");
678 reg.register("b", "sdxl:q8");
679 reg.register("c", "ltx-video:q8");
680 reg.mark_running("b", Some(0));
682
683 let cancelled = reg.cancel_all_queued();
684 assert_eq!(cancelled, 2, "both queued jobs cancelled, running one kept");
685 let snap = reg.snapshot();
686 assert_eq!(snap.entries.len(), 1);
687 assert_eq!(snap.entries[0].id, "b");
688 assert_eq!(snap.entries[0].state, JobLifecycle::Running);
689 }
690
691 #[test]
692 fn cancel_all_queued_on_empty_registry_returns_zero() {
693 let reg = JobRegistry::new();
694 assert_eq!(reg.cancel_all_queued(), 0);
695 }
696
697 #[tokio::test]
698 async fn cancel_all_queued_signals_every_registered_waiter() {
699 let reg = JobRegistry::new();
703 let cancel_a = reg.register("a", "flux-dev:fp16");
704 let cancel_b = reg.register("b", "sdxl:q8");
705 assert_eq!(reg.cancel_all_queued(), 2);
706 tokio::time::timeout(std::time::Duration::from_secs(1), cancel_a.notified())
707 .await
708 .expect("cancel signal for a must resolve");
709 tokio::time::timeout(std::time::Duration::from_secs(1), cancel_b.notified())
710 .await
711 .expect("cancel signal for b must resolve");
712 }
713
714 #[tokio::test]
715 async fn cancel_queued_signals_the_registered_waiter() {
716 let reg = JobRegistry::new();
720 let cancel = reg.register("a", "flux-dev:fp16");
721 reg.cancel_queued("a").unwrap();
722 tokio::time::timeout(std::time::Duration::from_secs(1), cancel.notified())
723 .await
724 .expect("cancel signal must resolve the waiter");
725 }
726
727 #[test]
728 fn snapshot_serializes_with_snake_case_state_and_omits_gpu_when_queued() {
729 let reg = JobRegistry::new();
733 reg.register("a", "flux-dev:fp16");
734 let snap = reg.snapshot();
735 let json = serde_json::to_string(&snap.entries[0]).unwrap();
736 assert!(json.contains(r#""state":"queued""#), "got: {json}");
737 assert!(
738 !json.contains("gpu"),
739 "queued row leaked a gpu field: {json}"
740 );
741
742 reg.mark_running("a", Some(0));
743 let snap2 = reg.snapshot();
744 let json2 = serde_json::to_string(&snap2.entries[0]).unwrap();
745 assert!(json2.contains(r#""state":"running""#));
746 assert!(json2.contains(r#""gpu":0"#));
747 }
748
749 #[test]
750 fn snapshot_carries_request_metadata_only_when_registered_with_it() {
751 let reg = JobRegistry::new();
752 reg.register("plain", "flux-dev:fp16");
753 let req: mold_core::GenerateRequest = serde_json::from_str(
754 r#"{"prompt":"a cat","model":"flux-dev:fp16","width":512,"height":512,"steps":4,"guidance":3.5}"#,
755 )
756 .expect("minimal request parses");
757 let meta = Box::new(mold_core::OutputMetadata::from_generate_request(
758 &req,
759 0,
760 None,
761 "test-version",
762 ));
763 reg.register_job("rich", "flux-dev:fp16", None, Some(true), Some(meta));
764
765 let snap = reg.snapshot();
766 let plain_json = serde_json::to_string(&snap.entries[0]).unwrap();
768 assert!(!plain_json.contains("metadata"), "got: {plain_json}");
769 assert_eq!(snap.entries[1].seed_pinned, Some(true));
770 let rich = snap.entries[1].metadata.as_ref().expect("metadata rides");
771 assert_eq!(rich.prompt, "a cat");
772 assert_eq!(rich.width, 512);
773 }
774
775 mod event_emission {
776 use super::*;
777 use crate::events::EventBroadcaster;
778 use mold_core::ServerEvent;
779 use tokio::sync::broadcast::error::TryRecvError;
780
781 fn wired() -> (
782 SharedJobRegistry,
783 tokio::sync::broadcast::Receiver<ServerEvent>,
784 ) {
785 let events = EventBroadcaster::new();
786 let rx = events.subscribe();
787 (JobRegistry::with_events(events), rx)
788 }
789
790 #[test]
791 fn register_emits_job_queued() {
792 let (reg, mut rx) = wired();
793 reg.register("a", "flux-dev:fp16");
794 match rx.try_recv().unwrap() {
795 ServerEvent::JobQueued { id, model } => {
796 assert_eq!(id, "a");
797 assert_eq!(model, "flux-dev:fp16");
798 }
799 other => panic!("expected job_queued, got {other:?}"),
800 }
801 }
802
803 #[test]
804 fn mark_running_emits_job_started_with_model_and_gpu() {
805 let (reg, mut rx) = wired();
806 reg.register("a", "flux-dev:fp16");
807 let _ = rx.try_recv(); reg.mark_running("a", Some(1));
809 match rx.try_recv().unwrap() {
810 ServerEvent::JobStarted { id, model, gpu } => {
811 assert_eq!(id, "a");
812 assert_eq!(model, "flux-dev:fp16");
813 assert_eq!(gpu, Some(1));
814 }
815 other => panic!("expected job_started, got {other:?}"),
816 }
817 }
818
819 #[test]
820 fn mark_running_unknown_id_emits_nothing() {
821 let (reg, mut rx) = wired();
822 reg.mark_running("ghost", None);
823 assert!(matches!(rx.try_recv(), Err(TryRecvError::Empty)));
824 }
825
826 #[test]
827 fn remove_emits_job_ended_exactly_once_across_double_call() {
828 let (reg, mut rx) = wired();
829 reg.register("a", "flux-dev:fp16");
830 let _ = rx.try_recv(); reg.remove("a");
832 reg.remove("a"); match rx.try_recv().unwrap() {
834 ServerEvent::JobEnded { id } => assert_eq!(id, "a"),
835 other => panic!("expected job_ended, got {other:?}"),
836 }
837 assert!(
838 matches!(rx.try_recv(), Err(TryRecvError::Empty)),
839 "second remove must not emit a duplicate job_ended"
840 );
841 }
842
843 #[test]
844 fn cancel_queued_emits_job_ended() {
845 let (reg, mut rx) = wired();
846 reg.register("a", "flux-dev:fp16");
847 let _ = rx.try_recv(); reg.cancel_queued("a").unwrap();
849 match rx.try_recv().unwrap() {
850 ServerEvent::JobEnded { id } => assert_eq!(id, "a"),
851 other => panic!("expected job_ended, got {other:?}"),
852 }
853 }
854
855 #[test]
856 fn cancel_all_queued_emits_exactly_one_job_ended_per_cancelled_job() {
857 let (reg, mut rx) = wired();
858 reg.register("a", "flux-dev:fp16");
859 reg.register("b", "sdxl:q8");
860 reg.register("c", "ltx-video:q8");
861 reg.mark_running("c", Some(0));
862 while rx.try_recv().is_ok() {}
864
865 assert_eq!(reg.cancel_all_queued(), 2);
866 let mut ended = Vec::new();
867 while let Ok(ev) = rx.try_recv() {
868 match ev {
869 ServerEvent::JobEnded { id } => ended.push(id),
870 other => panic!("expected only job_ended, got {other:?}"),
871 }
872 }
873 ended.sort();
874 assert_eq!(ended, vec!["a".to_string(), "b".to_string()]);
875 }
876 }
877}