Skip to main content

mold_server/
job_registry.rs

1//! Server-side ledger of in-flight generation jobs.
2//!
3//! The web UI tracks each "Generate" click as a card in `useGenerateStream`
4//! and relies on the SSE stream as the only signal that the work is still
5//! happening. When that stream silently drops (network blip, proxy idle
6//! timeout, server restart, browser tab suspended past keepalive), the card
7//! gets stuck `running` forever because no terminal event arrives.
8//!
9//! `JobRegistry` is the server's authoritative list of "things still owed an
10//! output" — every job between `submit()` and worker completion. The new
11//! `GET /api/queue` endpoint exposes this list so the SPA can poll it and
12//! dead-letter cards whose server-assigned `id` is no longer present.
13//!
14//! The registry deliberately doesn't track *completed* jobs — the gallery DB
15//! is the source of truth for those. Anything in here is currently queued or
16//! actively running on some worker.
17
18use 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/// Wire-facing job state. Mirrors the actual lifecycle:
26///
27/// - `queued` — accepted by `submit()`, sitting in the channel awaiting a
28///   dispatcher decision OR the dispatcher is mid-retry across workers.
29/// - `running` — handed off to a GPU worker thread; flipping happens when
30///   the worker pulls the job off its channel and starts loading / inferring.
31#[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/// One row in the `GET /api/queue` response.
39///
40/// `position` is the 0-based index in the registry's dispatch-priority order
41/// at the time of the snapshot — 0 is at the head (about to be dispatched, or
42/// already running on a worker), N-1 is last in line. It starts as insertion
43/// order but a `PATCH /api/queue/:id {position}` reorder moves a job within it,
44/// and it shifts as earlier jobs finish and drop out. The dispatch loops align
45/// their lookahead buffer to this order, so it reflects real execution order,
46/// not just the reported snapshot.
47#[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    /// GPU ordinal currently running this job (`null` for queued rows).
55    #[serde(skip_serializing_if = "Option::is_none")]
56    pub gpu: Option<usize>,
57    /// Preferred GPU ordinal for queued jobs (`None` means Auto).
58    #[serde(skip_serializing_if = "Option::is_none")]
59    pub target_gpu: Option<usize>,
60    /// Whether the submitted request pinned a seed. Additive — `metadata`'s
61    /// required seed field can't distinguish an explicit 0 from unpinned.
62    #[serde(skip_serializing_if = "Option::is_none")]
63    pub seed_pinned: Option<bool>,
64    /// The submitted request's parameters, metadata-shaped so any client can
65    /// inspect a queued job and reuse its settings. Additive — absent on
66    /// older servers; never carries image payloads.
67    #[serde(skip_serializing_if = "Option::is_none")]
68    #[schema(value_type = Object)]
69    pub metadata: Option<Box<mold_core::OutputMetadata>>,
70}
71
72/// Whole-queue listing returned by `GET /api/queue`. Wrapped in a struct so
73/// the response can grow extra fields (totals, byte counts, etc.) without a
74/// breaking change.
75#[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    /// Cancellation signal for `DELETE /api/queue/:id`. The submitting
91    /// handler holds the clone returned by `register*()` and selects on
92    /// `notified()` alongside the job's result channel; `cancel_queued`
93    /// fires `notify_one()` so the permit survives even when the cancel
94    /// lands before the waiter starts awaiting.
95    cancel: Arc<Notify>,
96}
97
98#[derive(Debug, Clone, Copy, PartialEq, Eq)]
99pub enum TargetGpuUpdateError {
100    NotFound,
101    AlreadyRunning,
102}
103
104/// Why a `DELETE /api/queue/:id` cancel attempt was rejected.
105#[derive(Debug, Clone, Copy, PartialEq, Eq)]
106pub enum QueuedJobCancelError {
107    NotFound,
108    AlreadyRunning,
109}
110
111/// Why a `PATCH /api/queue/:id` `position` reorder was rejected. Same shape as
112/// the target-gpu/cancel errors — only queued jobs are movable.
113#[derive(Debug, Clone, Copy, PartialEq, Eq)]
114pub enum QueueReorderError {
115    NotFound,
116    AlreadyRunning,
117}
118
119/// The registry itself. Construct via `JobRegistry::new` and share through
120/// `AppState`. All mutation is fire-and-forget — if the inner lock is
121/// poisoned (extremely unlikely in practice) we recover from the inner
122/// state rather than propagating the panic into the dispatcher hot path.
123pub struct JobRegistry {
124    inner: RwLock<Vec<EntryInternal>>,
125    /// Optional lifecycle broadcast (`GET /api/events`). Emitting from the
126    /// registry — rather than each call site — guarantees every submit /
127    /// promote / terminal path produces exactly one event. `None` keeps
128    /// event-less construction (tests) cheap.
129    events: Option<Arc<EventBroadcaster>>,
130}
131
132/// Cheap-cloneable handle. Workers and routes pass this around by value.
133pub 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    /// Like [`JobRegistry::new`] but mirrors every lifecycle change onto the
144    /// server-wide event broadcast.
145    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    /// Publish outside the registry lock — callers must have dropped the
153    /// write guard first so a slow broadcast can never extend the critical
154    /// section.
155    fn emit(&self, event: ServerEvent) {
156        if let Some(events) = &self.events {
157            events.publish(event);
158        }
159    }
160
161    /// Insert a freshly-submitted job at the tail in `Queued` state.
162    /// Returns the job's cancellation signal — see `register_with_target_gpu`.
163    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    /// Insert a freshly-submitted job with an optional queued lane target.
168    ///
169    /// Returns the job's cancellation signal. The submitting handler must
170    /// hold it and select on `notified()` alongside the result channel —
171    /// `cancel_queued` resolves it when `DELETE /api/queue/:id` removes the
172    /// entry. Callers that never wait (tests poking the registry directly)
173    /// can drop the handle.
174    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    /// Full-form insert: `register_with_target_gpu` plus the request's
184    /// metadata so `GET /api/queue` can expose the job's settings.
185    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    /// Cancel a still-queued job: remove its entry and fire the cancel
219    /// signal returned by `register*()` so the waiting request future
220    /// resolves with a cancellation error. Running jobs are not cancelable
221    /// — the GPU worker owns them and there is no safe preemption point.
222    ///
223    /// The state check and removal happen under the same write lock that
224    /// `mark_running` takes, so a job can never be both cancelled and
225    /// promoted. (A worker that already dequeued the job before the cancel
226    /// landed will observe the closed result channel and skip it.)
227    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    /// Cancel every still-queued job in one pass, backing `DELETE /api/queue`.
244    /// Under a single write lock this removes each `Queued` entry and fires its
245    /// cancel signal; running jobs are left untouched (same rule as
246    /// [`cancel_queued`](Self::cancel_queued) — a GPU worker owns them). After
247    /// dropping the lock it emits one `JobEnded` per cancelled job (the
248    /// emit-outside-lock discipline). Returns the number of jobs cancelled.
249    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    /// Promote a registry entry from `Queued` to `Running`. No-op if `id`
271    /// isn't present (the entry may have been removed concurrently).
272    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    /// Move a still-queued job to be the `position`-th queued entry (0-based)
308    /// among the queued jobs. This order is authoritative for dispatch: the
309    /// single- and multi-GPU loops align their lookahead buffer to
310    /// [`queued_ids_in_order`](Self::queued_ids_in_order) before picking, so a
311    /// reorder changes real execution order, not just the `GET /api/queue`
312    /// snapshot.
313    ///
314    /// `position` is clamped into the queued range, so a caller can pass a
315    /// large number to mean "send to the back". Running jobs are ignored for
316    /// indexing (they've left the reorderable set) and keep their absolute
317    /// slot, so a queued job can never jump ahead of one already executing.
318    /// The relative order of the other queued jobs is preserved. Moving a job
319    /// that is already at `position` is a successful no-op.
320    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        // Vec indices occupied by queued entries, in order — running entries
329        // are skipped so `position` indexes purely among movable jobs.
330        let queued_slots: Vec<usize> = entries
331            .iter()
332            .enumerate()
333            .filter(|(_, e)| e.state == JobLifecycle::Queued)
334            .map(|(i, _)| i)
335            .collect();
336        // `from` is queued (checked above), so it is present here.
337        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        // Pull the job out, then reinsert at the Vec index that makes it the
346        // `target`-th queued entry among the survivors.
347        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    /// Drop the entry — call once on every terminal path (success, error,
385    /// client-disconnect skip, dispatch failure). Idempotent.
386    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        // `remove` is called on every terminal path and is deliberately
397        // idempotent — only the call that actually dropped the entry emits,
398        // so subscribers see exactly one `job_ended` per job.
399        if removed {
400            self.emit(ServerEvent::JobEnded { id: id.to_string() });
401        }
402    }
403
404    /// The ids of currently-queued jobs, in dispatch-priority order — the same
405    /// order [`reorder_queued`](Self::reorder_queued) maintains and `GET
406    /// /api/queue` reports. Running jobs are excluded (they've left the
407    /// reorderable set).
408    ///
409    /// The dispatch loops call this each iteration to align their lookahead
410    /// buffer to the registry, which is what makes the registry the single
411    /// source of truth for order: a `PATCH /api/queue/:id {position}` reorder
412    /// changes real dispatch order, not just the snapshot.
413    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    /// Snapshot the registry as a wire-shaped listing. Positions are assigned
423    /// in the registry's dispatch-priority order at snapshot time (insertion
424    /// order until a `reorder_queued` moves a job).
425    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    /// Currently-tracked job count. Exposed for tests and metrics.
446    pub fn len(&self) -> usize {
447        self.inner.read().unwrap_or_else(|e| e.into_inner()).len()
448    }
449
450    /// Returns true when nothing is queued or running. Public so other
451    /// callers (metrics, integration tests) can check emptiness without
452    /// allocating a full snapshot.
453    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(&reg), ["c", "a", "b"]);
521        // Positions recompute from the new order.
522        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(&reg), ["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        // Move the last job to slot 1 — everything else keeps relative order.
544        reg.reorder_queued("e", 1).unwrap();
545        assert_eq!(queued_order(&reg), ["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(&reg), ["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(&reg), ["a", "b", "c"]);
566    }
567
568    #[test]
569    fn reorder_indexes_among_queued_only_and_keeps_running_jobs_in_place() {
570        // `a` is running: it must stay at the head and be skipped when
571        // `position` indexes the queued jobs, so a reorder can never promote a
572        // queued job ahead of the job already executing.
573        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        // Among the queued jobs [b, c], send c to slot 0.
579        reg.reorder_queued("c", 0).unwrap();
580        assert_eq!(queued_order(&reg), ["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        // No panic, no insertion — bogus id is ignored entirely.
609        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        // The worker's QueueSlot drop guard removes unconditionally — if the
633        // dispatcher already removed the entry on an error path, the worker's
634        // remove must not panic. Same for jobs that bypassed the registry
635        // entirely (id == "").
636        let reg = JobRegistry::new();
637        reg.register("a", "flux-dev:fp16");
638        reg.remove("a");
639        reg.remove("a"); // second remove is a no-op
640        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        // `b` is running — it must survive the bulk cancel.
681        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        // Each queued job's cancel handle must resolve — cancel_all_queued
700        // fires notify_one() per entry, so the permit survives even when the
701        // cancel lands before the waiter awaits.
702        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        // The handle returned by register() must resolve `notified()` even
717        // when the cancel fires before the waiter starts awaiting — Notify
718        // stores the permit from notify_one().
719        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        // Wire contract: queued rows must NOT carry a `gpu` field at all
730        // (clients shouldn't see `"gpu": null` and infer GPU 0). The state
731        // tag is lowercase to match the rest of the SSE/JSON style.
732        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        // Wire contract: rows without settings omit the key entirely.
767        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(); // drain job_queued
808            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(); // drain job_queued
831            reg.remove("a");
832            reg.remove("a"); // idempotent second call on another terminal path
833            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(); // drain job_queued
848            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            // Drain the three job_queued + one job_started emissions.
863            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}