//! Job routes: snapshots and a live server-sent event stream.
use std::convert::Infallible;
use axum::Json;
use axum::extract::{Path, State};
use axum::response::Sse;
use axum::response::sse::{Event, KeepAlive};
use futures_util::stream::{self, Stream, StreamExt as _};
use tokio::sync::broadcast::error::RecvError;
use twice_core::{JobId, JobView, JobsResponse};
use crate::error::ApiError;
use crate::jobs::JobEvent;
use crate::state::AppState;
/// `GET /jobs` โ every retained job, newest first.
pub(crate) async fn list(State(state): State<AppState>) -> Json<JobsResponse> {
Json(JobsResponse {
jobs: state.jobs.list().await,
})
}
/// `GET /jobs/{id}` โ one job including its buffered log.
pub(crate) async fn get_one(
State(state): State<AppState>,
Path(id): Path<String>,
) -> Result<Json<JobView>, ApiError> {
let id: JobId = id.parse().map_err(|_ignored| ApiError::UnknownJob)?;
state
.jobs
.get(id)
.await
.map(Json)
.ok_or(ApiError::UnknownJob)
}
/// `GET /jobs/{id}/stream` โ buffered log, then live output, as SSE.
///
/// Events are named `line` and `finished`. The stream closes right after
/// `finished`, so a client can treat end-of-stream as end-of-job.
pub(crate) async fn stream(
State(state): State<AppState>,
Path(id): Path<String>,
) -> Result<Sse<impl Stream<Item = Result<Event, Infallible>> + Send>, ApiError> {
let id: JobId = id.parse().map_err(|_ignored| ApiError::UnknownJob)?;
let (view, receiver) = state.jobs.subscribe(id).await.ok_or(ApiError::UnknownJob)?;
let replay: Vec<Result<Event, Infallible>> = view
.log
.iter()
.map(|line| Ok(encode("line", line)))
.collect();
let backlog = stream::iter(replay);
// A job that already finished has nothing live to wait for; subscribing to
// its broadcast channel would hang until the server shuts down.
let tail = if view.outcome.is_terminal() {
stream::once(std::future::ready(Ok(encode("finished", &view.outcome)))).boxed()
} else {
live_events(receiver).boxed()
};
Ok(Sse::new(backlog.chain(tail)).keep_alive(KeepAlive::default()))
}
/// Turns the job's broadcast channel into a stream that ends at `finished`.
fn live_events(
receiver: tokio::sync::broadcast::Receiver<JobEvent>,
) -> impl Stream<Item = Result<Event, Infallible>> + Send {
stream::unfold((receiver, false), |(mut receiver, finished)| async move {
if finished {
return None;
}
loop {
match receiver.recv().await {
Ok(JobEvent::Line(line)) => {
return Some((Ok(encode("line", &line)), (receiver, false)));
}
Ok(JobEvent::Finished(outcome)) => {
return Some((Ok(encode("finished", &outcome)), (receiver, true)));
}
// A slow client missed lines; the buffered log still has them.
Err(RecvError::Lagged(_)) => {}
Err(RecvError::Closed) => return None,
}
}
})
}
fn encode<T: serde::Serialize>(name: &str, payload: &T) -> Event {
Event::default()
.event(name)
.json_data(payload)
.unwrap_or_else(|_ignored| Event::default().event("error").data("encoding failed"))
}