Skip to main content

tuwunel_api/server/
backfill.rs

1use 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
14/// arbitrary number but synapse's is 100 and we can handle lots of these
15/// anyways
16const LIMIT_MAX: usize = 150;
17/// no spec defined number but we can handle a lot of these
18const LIMIT_DEFAULT: usize = 50;
19
20/// # `GET /_matrix/federation/v1/backfill/<room_id>`
21///
22/// Retrieves events from before the sender joined the room, if the room's
23/// history visibility allows.
24pub(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}