use std::time::Instant;
use chrono::Utc;
use crate::job_log::{job_span, TRACE_TARGET};
use crate::runtime::{
record_local_job, record_recent_job, CurrentJob, JobOutcome, JobSource, RecentJob,
WorkerObservers,
};
use crate::types::TaskResult;
pub struct JobRun {
observers: WorkerObservers,
job: CurrentJob,
span: tracing::Span,
started: Instant,
}
impl JobRun {
pub fn begin(observers: &WorkerObservers, job: CurrentJob) -> Self {
let span = job_span(&job.job_id);
span.in_scope(|| {
tracing::info!(
target: TRACE_TARGET,
op = "job",
source = job.source.as_str(),
kind = job.kind.as_str(),
model = %job.model,
"job started"
);
});
observers.active_jobs.lock().push(job.clone());
Self {
observers: observers.clone(),
job,
span,
started: Instant::now(),
}
}
pub fn job_id(&self) -> &str {
&self.job.job_id
}
pub fn span(&self) -> &tracing::Span {
&self.span
}
pub fn set_prompt(&mut self, prompt: &str) {
self.job.prompt = crate::runtime::truncate_prompt(prompt);
}
pub fn keep_thumbnail(&self, result: &TaskResult) {
self.thumbnail_keeper().keep(result);
}
pub fn thumbnail_keeper(&self) -> ThumbnailKeeper {
ThumbnailKeeper {
thumbnails: self.observers.thumbnails.clone(),
job_id: self.job.job_id.clone(),
span: self.span.clone(),
}
}
pub fn finish(self, outcome: JobOutcome) -> RecentJob {
let elapsed_ms = self.started.elapsed().as_millis() as u64;
self.span.in_scope(|| match &outcome {
JobOutcome::Completed => tracing::info!(
target: TRACE_TARGET,
op = "job",
outcome = "completed",
elapsed_ms,
"job finished"
),
JobOutcome::Failed { reason } => tracing::warn!(
target: TRACE_TARGET,
op = "job",
outcome = "failed",
elapsed_ms,
reason = %reason,
"job finished"
),
});
let recent = RecentJob {
job_id: self.job.job_id.clone(),
kind: self.job.kind,
model: self.job.model.clone(),
prompt: self.job.prompt.clone(),
outcome,
started_at: self.job.started_at,
finished_at: Utc::now(),
source: self.job.source,
};
match self.job.source {
JobSource::Studio => record_recent_job(&self.observers, recent.clone()),
_ => record_local_job(&self.observers, recent.clone()),
}
recent
}
}
#[derive(Clone)]
pub struct ThumbnailKeeper {
thumbnails: crate::thumbnail::Thumbnails,
job_id: String,
span: tracing::Span,
}
impl ThumbnailKeeper {
pub fn keep(&self, result: &TaskResult) {
let TaskResult::Image { bytes, .. } = result else {
return;
};
match crate::thumbnail::make_thumbnail(bytes) {
Ok(png) => self.thumbnails.insert(&self.job_id, png),
Err(err) => self.span.in_scope(|| {
tracing::warn!(
target: TRACE_TARGET,
op = "thumbnail",
error = %format!("{err:#}"),
"could not make a thumbnail of the job's image"
);
}),
}
}
}
impl Drop for JobRun {
fn drop(&mut self) {
self.observers
.active_jobs
.lock()
.retain(|j| j.job_id != self.job.job_id);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::types::TaskKind;
fn job(id: &str, source: JobSource) -> CurrentJob {
CurrentJob {
job_id: id.into(),
kind: TaskKind::Image,
model: "m".into(),
prompt: "a fox".into(),
started_at: Utc::now(),
source,
}
}
fn webp() -> Vec<u8> {
let image = image::RgbImage::from_pixel(8, 8, image::Rgb([1, 2, 3]));
let mut out = Vec::new();
image::DynamicImage::ImageRgb8(image)
.write_to(
&mut std::io::Cursor::new(&mut out),
image::ImageFormat::WebP,
)
.unwrap();
out
}
#[test]
fn a_running_job_is_listed_until_it_finishes() {
let observers = WorkerObservers::default();
let run = JobRun::begin(&observers, job("run-1", JobSource::Local));
assert_eq!(observers.active_jobs.lock()[0].job_id, "run-1");
let recent = run.finish(JobOutcome::Completed);
assert!(observers.active_jobs.lock().is_empty());
assert_eq!(recent.source, JobSource::Local);
assert_eq!(observers.local_jobs.lock()[0], recent);
}
#[test]
fn a_studio_job_is_recorded_in_the_recent_ring() {
let observers = WorkerObservers::default();
JobRun::begin(&observers, job("studio-1", JobSource::Studio)).finish(JobOutcome::Failed {
reason: "boom".into(),
});
assert_eq!(observers.recent_jobs.lock()[0].job_id, "studio-1");
assert!(observers.local_jobs.lock().is_empty());
}
#[test]
fn lane_and_stream_jobs_are_recorded_in_the_local_ring() {
let observers = WorkerObservers::default();
JobRun::begin(&observers, job("lane-1", JobSource::Lane)).finish(JobOutcome::Completed);
JobRun::begin(&observers, job("stream-1", JobSource::Stream)).finish(JobOutcome::Completed);
let ids: Vec<_> = observers
.local_jobs
.lock()
.iter()
.map(|j| j.job_id.clone())
.collect();
assert_eq!(ids, ["stream-1", "lane-1"]);
}
#[test]
fn the_prompt_can_be_set_once_known() {
let observers = WorkerObservers::default();
let mut run = JobRun::begin(&observers, job("stream-2", JobSource::Stream));
run.set_prompt("hello there");
assert_eq!(run.finish(JobOutcome::Completed).prompt, "hello there");
}
#[test]
fn dropping_an_unfinished_job_takes_it_off_the_running_list() {
let observers = WorkerObservers::default();
{
let _run = JobRun::begin(&observers, job("dropped", JobSource::Local));
assert_eq!(observers.active_jobs.lock().len(), 1);
}
assert!(observers.active_jobs.lock().is_empty());
assert!(observers.local_jobs.lock().is_empty());
}
#[test]
fn an_image_result_keeps_a_thumbnail() {
let observers = WorkerObservers::default();
let run = JobRun::begin(&observers, job("thumb-1", JobSource::Local));
run.keep_thumbnail(&TaskResult::Image {
bytes: webp(),
ext: "webp".into(),
});
assert!(observers.thumbnails.contains("thumb-1"));
}
#[test]
fn a_non_image_result_keeps_no_thumbnail() {
let observers = WorkerObservers::default();
let run = JobRun::begin(&observers, job("thumb-2", JobSource::Local));
run.keep_thumbnail(&TaskResult::Llm {
json: serde_json::json!({}),
});
assert!(observers.thumbnails.is_empty());
}
#[test]
fn the_job_log_records_start_finish_and_a_failed_thumbnail() {
crate::test_support::install_job_log_capture();
let observers = WorkerObservers::default();
let run = JobRun::begin(&observers, job("logged-1", JobSource::Local));
run.keep_thumbnail(&TaskResult::Image {
bytes: b"garbage".to_vec(),
ext: "webp".into(),
});
run.finish(JobOutcome::Failed {
reason: "engine exploded".into(),
});
let log = crate::job_log::global().get("logged-1").expect("job log");
let messages: Vec<_> = log.lines.iter().map(|l| l.message.as_str()).collect();
assert!(messages[0].starts_with("job started"), "{messages:?}");
assert!(messages[0].contains("source=\"local\""), "{messages:?}");
assert!(
messages[1].starts_with("could not make a thumbnail"),
"{messages:?}"
);
assert!(messages[1].contains("op=\"thumbnail\""), "{messages:?}");
assert!(messages[2].starts_with("job finished"), "{messages:?}");
assert!(
messages[2].contains("reason=engine exploded"),
"{messages:?}"
);
assert_eq!(log.lines[2].level, "warn");
assert!(observers.thumbnails.is_empty());
}
}