tuwunel_service/sending/sender/dispatch/
federation.rs1use 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#[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}