use crate::serve::error::ServeError;
use crate::serve::logs::{LogEvent, log_events};
use crate::serve::state::ServerState;
use axum::extract::{Path, State};
use axum::response::Sse;
use axum::response::sse::{Event, KeepAlive};
use futures::{Stream, StreamExt};
use std::convert::Infallible;
use std::time::Duration;
use tokio::sync::broadcast;
const KEEP_ALIVE_SECS: u64 = 15;
pub async fn stream_logs(
State(state): State<ServerState>,
Path(id): Path<String>,
) -> Result<Sse<impl Stream<Item = Result<Event, Infallible>>>, ServeError> {
let (snapshot, rx, ended) = match state.log_hub().reader(&id) {
Some(reader) => reader,
None => {
let known = state
.history()
.get(&id)
.await
.map_err(|e| ServeError::Internal(e.to_string()))?
.is_some();
if !known {
return Err(ServeError::NotFound);
}
(Vec::new(), broadcast::channel(1).1, true)
}
};
let stream =
log_events(snapshot, rx, ended).map(|ev| Ok::<Event, Infallible>(to_sse_event(ev)));
Ok(
Sse::new(stream)
.keep_alive(KeepAlive::new().interval(Duration::from_secs(KEEP_ALIVE_SECS))),
)
}
fn to_sse_event(ev: LogEvent) -> Event {
match ev {
LogEvent::Log(line) => Event::default().event("log").data(line),
LogEvent::Truncated(n) => Event::default().event("truncated").data(format!(
"{n} log line(s) dropped; rely on the centralized log sink"
)),
LogEvent::End => Event::default().event("end").data("done"),
}
}