Skip to main content

tuwunel_service/sending/sender/dispatch/
federation.rs

1use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD};
2use futures::StreamExt;
3use ruma::{
4	MilliSecondsSinceUnixEpoch, OwnedServerName,
5	api::federation::transactions::{edu::Edu, send_transaction_message::v1::Request},
6	serde::Raw,
7};
8use tuwunel_core::{
9	extract_variant, implement,
10	utils::{IterStream, calculate_hash, future::TryExtExt, stream::WidebandExt},
11	warn,
12};
13
14use super::SendingResult;
15use crate::sending::{Destination, EduBuf, SendingEvent, Service};
16
17/// Send a federation transaction, reporting whether one went out at all.
18///
19/// Rows that all fail to load leave nothing to send; they still succeed, so
20/// their keys are acknowledged.
21#[implement(Service)]
22#[tracing::instrument(
23	name = "federation",
24	level = "debug",
25	skip(self, events),
26	fields(
27		events = %events.len(),
28	),
29)]
30pub(super) async fn send_events_dest_federation(
31	&self,
32	server: OwnedServerName,
33	events: Vec<SendingEvent>,
34) -> (SendingResult, bool) {
35	let pdus: Vec<_> = events
36		.iter()
37		.filter_map(|event| extract_variant!(event, SendingEvent::Pdu))
38		.stream()
39		.wide_filter_map(|pdu_id| {
40			self.services
41				.timeline
42				.get_pdu_json_from_id(pdu_id)
43				.ok()
44		})
45		.wide_then(|pdu| {
46			self.services
47				.state_accessor
48				.erased_for_server(&server, pdu)
49		})
50		.wide_then(|pdu| {
51			self.services
52				.federation
53				.format_pdu_into(pdu, None)
54		})
55		.collect()
56		.await;
57
58	let edus: Vec<Raw<Edu>> = events
59		.iter()
60		.filter_map(|event| extract_variant!(event, SendingEvent::Edu))
61		.map(EduBuf::as_slice)
62		.map(serde_json::from_slice)
63		.filter_map(Result::ok)
64		.collect();
65
66	if pdus.is_empty() && edus.is_empty() {
67		return (Ok(Destination::Federation(server)), false);
68	}
69
70	let preimage = pdus
71		.iter()
72		.map(|raw| raw.get().as_bytes())
73		.chain(edus.iter().map(|raw| raw.json().get().as_bytes()));
74
75	let txn_hash = calculate_hash(preimage);
76	let txn_id = &*URL_SAFE_NO_PAD.encode(txn_hash);
77	let request = Request {
78		transaction_id: txn_id.into(),
79		origin: self.server.name.clone(),
80		origin_server_ts: MilliSecondsSinceUnixEpoch::now(),
81		pdus,
82		edus,
83	};
84
85	let result = self
86		.services
87		.federation
88		.execute_on(&self.services.client.sender, &server, request)
89		.await;
90
91	result
92		.iter()
93		.flat_map(|resp| resp.pdus.iter())
94		.filter_map(|(event_id, result)| {
95			result
96				.as_ref()
97				.err()
98				.map(|error| (event_id, error))
99		})
100		.for_each(|(event_id, error)| {
101			warn!(%txn_id, %server, %event_id, %error, "error sending PDU to remote server");
102		});
103
104	let result = match result {
105		| Ok(_) => Ok(Destination::Federation(server)),
106		| Err(error) => Err((Destination::Federation(server), error)),
107	};
108
109	(result, true)
110}