1use 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
20pub struct JobRun {
22 observers: WorkerObservers,
23 job: CurrentJob,
24 span: tracing::Span,
25 started: Instant,
26}
27
28impl JobRun {
29 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 pub fn job_id(&self) -> &str {
53 &self.job.job_id
54 }
55
56 pub fn span(&self) -> &tracing::Span {
58 &self.span
59 }
60
61 pub fn set_prompt(&mut self, prompt: &str) {
63 self.job.prompt = crate::runtime::truncate_prompt(prompt);
64 }
65
66 pub fn keep_thumbnail(&self, result: &TaskResult) {
69 self.thumbnail_keeper().keep(result);
70 }
71
72 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 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#[derive(Clone)]
123pub struct ThumbnailKeeper {
124 thumbnails: crate::thumbnail::Thumbnails,
125 job_id: String,
126 span: tracing::Span,
127}
128
129impl ThumbnailKeeper {
130 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}