Skip to content
//! 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"))
}