tuwunel_api/server/
backfill.rs1use std::cmp;
2
3use axum::extract::State;
4use futures::{FutureExt, StreamExt, TryStreamExt};
5use ruma::{MilliSecondsSinceUnixEpoch, api::federation::backfill::get_backfill};
6use tuwunel_core::{
7 PduCount, Result,
8 utils::{IterStream, ReadyExt, math::usize_from_ruma_bounded},
9};
10
11use super::AccessCheck;
12use crate::Ruma;
13
14const LIMIT_MAX: usize = 150;
17const LIMIT_DEFAULT: usize = 50;
19
20pub(crate) async fn get_backfill_route(
25 State(services): State<crate::State>,
26 ref body: Ruma<get_backfill::v1::Request>,
27) -> Result<get_backfill::v1::Response> {
28 AccessCheck {
29 services: &services,
30 origin: body.origin(),
31 room_id: &body.room_id,
32 event_id: None,
33 }
34 .check()
35 .await?;
36
37 let limit = usize_from_ruma_bounded(body.limit, LIMIT_DEFAULT, LIMIT_MAX);
38
39 let from = body
40 .v
41 .iter()
42 .stream()
43 .filter_map(|event_id| {
44 services
45 .timeline
46 .get_pdu_count(event_id)
47 .map(Result::ok)
48 })
49 .ready_fold(PduCount::min(), cmp::max)
50 .await;
51
52 let room_version = services
53 .state
54 .get_room_version(&body.room_id)
55 .await
56 .ok();
57
58 Ok(get_backfill::v1::Response {
59 origin_server_ts: MilliSecondsSinceUnixEpoch::now(),
60
61 origin: services.globals.server_name().to_owned(),
62
63 pdus: services
64 .timeline
65 .pdus_rev(None, &body.room_id, Some(from.saturating_add(1)))
66 .try_filter_map(async |(_, pdu)| {
67 Ok(services
68 .state_accessor
69 .server_can_see_event(body.origin(), &pdu.room_id, &pdu.event_id)
70 .await
71 .then_some(pdu))
72 })
73 .try_filter_map(async |pdu| {
74 Ok(services
75 .timeline
76 .get_pdu_json(&pdu.event_id)
77 .await
78 .ok())
79 })
80 .take(limit)
81 .and_then(|pdu| {
82 services
83 .state_accessor
84 .erased_for_server(body.origin(), pdu)
85 .map(Ok)
86 })
87 .and_then(|pdu| {
88 services
89 .federation
90 .format_pdu_into(pdu, room_version.as_ref())
91 .map(Ok)
92 })
93 .try_collect()
94 .boxed()
95 .await?,
96 })
97}