summaryrefslogtreecommitdiff
path: root/src/event/routes.rs
blob: 5b9c7e3634c9ded7b3689bb620215252d1b374d3 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
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::{Sequence, Sequenced as _},
    token::{app::ValidateError, extract::Identity},
};

#[cfg(test)]
mod test;

pub fn router() -> Router<App> {
    Router::new().route("/api/events", get(events))
}

#[derive(Default, serde::Deserialize)]
struct EventsQuery {
    resume_point: Option<Sequence>,
}

async fn events(
    State(app): State<App>,
    identity: Identity,
    last_event_id: Option<LastEventId<Sequence>>,
    Query(query): Query<EventsQuery>,
) -> Result<Events<impl Stream<Item = Event> + 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>(S);

impl<S> IntoResponse for Events<S>
where
    S: Stream<Item = Event> + 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<Event> for sse::Event {
    type Error = serde_json::Error;

    fn try_from(event: Event) -> Result<Self, Self::Error> {
        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(),
        }
    }
}