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
|
use axum::{
extract::State,
response::{
sse::{self, Sse},
IntoResponse, Response,
},
routing::get,
Router,
};
use futures::stream::{Stream, StreamExt as _};
use super::{
extract::LastEventId,
types::{self, ResumePoint},
};
use crate::{app::App, error::Internal, repo::login::Login};
#[cfg(test)]
mod test;
pub fn router() -> Router<App> {
Router::new().route("/api/events", get(events))
}
async fn events(
State(app): State<App>,
_: Login, // requires auth, but doesn't actually care who you are
last_event_id: Option<LastEventId<ResumePoint>>,
) -> Result<Events<impl Stream<Item = types::ResumableEvent> + std::fmt::Debug>, Internal> {
let resume_at = last_event_id
.map(LastEventId::into_inner)
.unwrap_or_default();
let stream = app.events().subscribe(resume_at).await?;
Ok(Events(stream))
}
#[derive(Debug)]
struct Events<S>(S);
impl<S> IntoResponse for Events<S>
where
S: Stream<Item = types::ResumableEvent> + 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<types::ResumableEvent> for sse::Event {
type Error = serde_json::Error;
fn try_from(value: types::ResumableEvent) -> Result<Self, Self::Error> {
let types::ResumableEvent(resume_at, data) = value;
let id = serde_json::to_string(&resume_at)?;
let data = serde_json::to_string_pretty(&data)?;
let event = Self::default().id(id).data(data);
Ok(event)
}
}
|