Skip to main content

studio_worker/
job_run.rs

1//! One job's bookkeeping, shared by every job path (studio offers, local
2//! transient jobs, lane requests, streaming sessions).
3//!
4//! [`JobRun::begin`] lists the job as running, opens its log span and logs
5//! "job started"; [`JobRun::finish`] logs "job finished" and records the job
6//! in its ring.  Dropping a run without finishing it (an early return or a
7//! panic) still takes it off the running list.
8
9use std::time::Instant;
10
11use chrono::Utc;
12
13use crate::job_log::{job_span, TRACE_TARGET};
14use crate::runtime::{
15    record_local_job, record_recent_job, CurrentJob, JobOutcome, JobSource, RecentJob,
16    WorkerObservers,
17};
18use crate::types::TaskResult;
19
20/// A running job.  See the module docs.
21pub struct JobRun {
22    observers: WorkerObservers,
23    job: CurrentJob,
24    span: tracing::Span,
25    started: Instant,
26}
27
28impl JobRun {
29    /// List `job` as running and log its start inside its span.
30    pub fn begin(observers: &WorkerObservers, job: CurrentJob) -> Self {
31        let span = job_span(&job.job_id);
32        span.in_scope(|| {
33            tracing::info!(
34                target: TRACE_TARGET,
35                op = "job",
36                source = job.source.as_str(),
37                kind = job.kind.as_str(),
38                model = %job.model,
39                "job started"
40            );
41        });
42        observers.active_jobs.lock().push(job.clone());
43        Self {
44            observers: observers.clone(),
45            job,
46            span,
47            started: Instant::now(),
48        }
49    }
50
51    /// The job's id.
52    pub fn job_id(&self) -> &str {
53        &self.job.job_id
54    }
55
56    /// The span to run the job's work in; its events land in the job log.
57    pub fn span(&self) -> &tracing::Span {
58        &self.span
59    }
60
61    /// Replace the job's prompt preview, e.g. with a stream's final text.
62    pub fn set_prompt(&mut self, prompt: &str) {
63        self.job.prompt = crate::runtime::truncate_prompt(prompt);
64    }
65
66    /// Keep a thumbnail when `result` is an image.  A thumbnail that cannot
67    /// be made is logged in the job's log; the job is unaffected.
68    pub fn keep_thumbnail(&self, result: &TaskResult) {
69        self.thumbnail_keeper().keep(result);
70    }
71
72    /// A `Send` handle that keeps this job's thumbnail, for a blocking
73    /// thread that holds the result.
74    pub fn thumbnail_keeper(&self) -> ThumbnailKeeper {
75        ThumbnailKeeper {
76            thumbnails: self.observers.thumbnails.clone(),
77            job_id: self.job.job_id.clone(),
78            span: self.span.clone(),
79        }
80    }
81
82    /// Log the job's end, take it off the running list and record it in its
83    /// ring (studio jobs in the recent ring, the rest in the local ring).
84    pub fn finish(self, outcome: JobOutcome) -> RecentJob {
85        let elapsed_ms = self.started.elapsed().as_millis() as u64;
86        self.span.in_scope(|| match &outcome {
87            JobOutcome::Completed => tracing::info!(
88                target: TRACE_TARGET,
89                op = "job",
90                outcome = "completed",
91                elapsed_ms,
92                "job finished"
93            ),
94            JobOutcome::Failed { reason } => tracing::warn!(
95                target: TRACE_TARGET,
96                op = "job",
97                outcome = "failed",
98                elapsed_ms,
99                reason = %reason,
100                "job finished"
101            ),
102        });
103        let recent = RecentJob {
104            job_id: self.job.job_id.clone(),
105            kind: self.job.kind,
106            model: self.job.model.clone(),
107            prompt: self.job.prompt.clone(),
108            outcome,
109            started_at: self.job.started_at,
110            finished_at: Utc::now(),
111            source: self.job.source,
112        };
113        match self.job.source {
114            JobSource::Studio => record_recent_job(&self.observers, recent.clone()),
115            _ => record_local_job(&self.observers, recent.clone()),
116        }
117        recent
118    }
119}
120
121/// Keeps one job's thumbnail; see [`JobRun::thumbnail_keeper`].
122#[derive(Clone)]
123pub struct ThumbnailKeeper {
124    thumbnails: crate::thumbnail::Thumbnails,
125    job_id: String,
126    span: tracing::Span,
127}
128
129impl ThumbnailKeeper {
130    /// Keep a thumbnail when `result` is an image; log in the job's log
131    /// when one cannot be made.
132    pub fn keep(&self, result: &TaskResult) {
133        let TaskResult::Image { bytes, .. } = result else {
134            return;
135        };
136        match crate::thumbnail::make_thumbnail(bytes) {
137            Ok(png) => self.thumbnails.insert(&self.job_id, png),
138            Err(err) => self.span.in_scope(|| {
139                tracing::warn!(
140                    target: TRACE_TARGET,
141                    op = "thumbnail",
142                    error = %format!("{err:#}"),
143                    "could not make a thumbnail of the job's image"
144                );
145            }),
146        }
147    }
148}
149
150impl Drop for JobRun {
151    fn drop(&mut self) {
152        self.observers
153            .active_jobs
154            .lock()
155            .retain(|j| j.job_id != self.job.job_id);
156    }
157}
158
159#[cfg(test)]
160mod tests {
161    use super::*;
162    use crate::types::TaskKind;
163
164    fn job(id: &str, source: JobSource) -> CurrentJob {
165        CurrentJob {
166            job_id: id.into(),
167            kind: TaskKind::Image,
168            model: "m".into(),
169            prompt: "a fox".into(),
170            started_at: Utc::now(),
171            source,
172        }
173    }
174
175    fn webp() -> Vec<u8> {
176        let image = image::RgbImage::from_pixel(8, 8, image::Rgb([1, 2, 3]));
177        let mut out = Vec::new();
178        image::DynamicImage::ImageRgb8(image)
179            .write_to(
180                &mut std::io::Cursor::new(&mut out),
181                image::ImageFormat::WebP,
182            )
183            .unwrap();
184        out
185    }
186
187    #[test]
188    fn a_running_job_is_listed_until_it_finishes() {
189        let observers = WorkerObservers::default();
190        let run = JobRun::begin(&observers, job("run-1", JobSource::Local));
191        assert_eq!(observers.active_jobs.lock()[0].job_id, "run-1");
192        let recent = run.finish(JobOutcome::Completed);
193        assert!(observers.active_jobs.lock().is_empty());
194        assert_eq!(recent.source, JobSource::Local);
195        assert_eq!(observers.local_jobs.lock()[0], recent);
196    }
197
198    #[test]
199    fn a_studio_job_is_recorded_in_the_recent_ring() {
200        let observers = WorkerObservers::default();
201        JobRun::begin(&observers, job("studio-1", JobSource::Studio)).finish(JobOutcome::Failed {
202            reason: "boom".into(),
203        });
204        assert_eq!(observers.recent_jobs.lock()[0].job_id, "studio-1");
205        assert!(observers.local_jobs.lock().is_empty());
206    }
207
208    #[test]
209    fn lane_and_stream_jobs_are_recorded_in_the_local_ring() {
210        let observers = WorkerObservers::default();
211        JobRun::begin(&observers, job("lane-1", JobSource::Lane)).finish(JobOutcome::Completed);
212        JobRun::begin(&observers, job("stream-1", JobSource::Stream)).finish(JobOutcome::Completed);
213        let ids: Vec<_> = observers
214            .local_jobs
215            .lock()
216            .iter()
217            .map(|j| j.job_id.clone())
218            .collect();
219        assert_eq!(ids, ["stream-1", "lane-1"]);
220    }
221
222    #[test]
223    fn the_prompt_can_be_set_once_known() {
224        let observers = WorkerObservers::default();
225        let mut run = JobRun::begin(&observers, job("stream-2", JobSource::Stream));
226        run.set_prompt("hello there");
227        assert_eq!(run.finish(JobOutcome::Completed).prompt, "hello there");
228    }
229
230    #[test]
231    fn dropping_an_unfinished_job_takes_it_off_the_running_list() {
232        let observers = WorkerObservers::default();
233        {
234            let _run = JobRun::begin(&observers, job("dropped", JobSource::Local));
235            assert_eq!(observers.active_jobs.lock().len(), 1);
236        }
237        assert!(observers.active_jobs.lock().is_empty());
238        assert!(observers.local_jobs.lock().is_empty());
239    }
240
241    #[test]
242    fn an_image_result_keeps_a_thumbnail() {
243        let observers = WorkerObservers::default();
244        let run = JobRun::begin(&observers, job("thumb-1", JobSource::Local));
245        run.keep_thumbnail(&TaskResult::Image {
246            bytes: webp(),
247            ext: "webp".into(),
248        });
249        assert!(observers.thumbnails.contains("thumb-1"));
250    }
251
252    #[test]
253    fn a_non_image_result_keeps_no_thumbnail() {
254        let observers = WorkerObservers::default();
255        let run = JobRun::begin(&observers, job("thumb-2", JobSource::Local));
256        run.keep_thumbnail(&TaskResult::Llm {
257            json: serde_json::json!({}),
258        });
259        assert!(observers.thumbnails.is_empty());
260    }
261
262    #[test]
263    fn the_job_log_records_start_finish_and_a_failed_thumbnail() {
264        crate::test_support::install_job_log_capture();
265        let observers = WorkerObservers::default();
266        let run = JobRun::begin(&observers, job("logged-1", JobSource::Local));
267        run.keep_thumbnail(&TaskResult::Image {
268            bytes: b"garbage".to_vec(),
269            ext: "webp".into(),
270        });
271        run.finish(JobOutcome::Failed {
272            reason: "engine exploded".into(),
273        });
274        let log = crate::job_log::global().get("logged-1").expect("job log");
275        let messages: Vec<_> = log.lines.iter().map(|l| l.message.as_str()).collect();
276        assert!(messages[0].starts_with("job started"), "{messages:?}");
277        assert!(messages[0].contains("source=\"local\""), "{messages:?}");
278        assert!(
279            messages[1].starts_with("could not make a thumbnail"),
280            "{messages:?}"
281        );
282        assert!(messages[1].contains("op=\"thumbnail\""), "{messages:?}");
283        assert!(messages[2].starts_with("job finished"), "{messages:?}");
284        assert!(
285            messages[2].contains("reason=engine exploded"),
286            "{messages:?}"
287        );
288        assert_eq!(log.lines[2].level, "warn");
289        assert!(observers.thumbnails.is_empty());
290    }
291}