use axum::{ extract::State, response::{ sse::{self, Sse}, IntoResponse, Response, }, routing::get, Router, }; use axum_extra::extract::Query; use futures::stream::{Stream, StreamExt as _}; use super::{extract::LastEventId, Event}; use crate::{ app::App, error::{Internal, Unauthorized}, event::{ResumePoint, Sequence, Sequenced as _}, token::{app::ValidateError, extract::Identity}, }; #[cfg(test)] mod test; pub fn router() -> Router { Router::new().route("/api/events", get(events)) } #[derive(Default, serde::Deserialize)] struct EventsQuery { resume_point: ResumePoint, } async fn events( State(app): State, identity: Identity, last_event_id: Option>, Query(query): Query, ) -> Result + std::fmt::Debug>, EventsError> { let resume_at = last_event_id .map(LastEventId::into_inner) .or(query.resume_point); let stream = app.events().subscribe(resume_at).await?; let stream = app.tokens().limit_stream(identity.token, stream).await?; Ok(Events(stream)) } #[derive(Debug)] struct Events(S); impl IntoResponse for Events where S: Stream + Send + 'static, { fn into_response(self) -> Response { let Self(stream) = self; let stream = stream.map(sse::Event::try_from); Sse::new(stream) .keep_alive(sse::KeepAlive::default()) .into_response() } } impl TryFrom for sse::Event { type Error = serde_json::Error; fn try_from(event: Event) -> Result { let id = serde_json::to_string(&event.sequence())?; let data = serde_json::to_string_pretty(&event)?; let event = Self::default().id(id).data(data); Ok(event) } } #[derive(Debug, thiserror::Error)] #[error(transparent)] pub enum EventsError { DatabaseError(#[from] sqlx::Error), ValidateError(#[from] ValidateError), } impl IntoResponse for EventsError { fn into_response(self) -> Response { match self { Self::ValidateError(ValidateError::InvalidToken) => Unauthorized.into_response(), other => Internal::from(other).into_response(), } } }