1use std::collections::BTreeMap;
5use std::sync::Arc;
6
7use serde_json::{json, Value};
8
9use super::{with_conn_pub, Homeserver};
10use crate::error::MatrixError;
11use crate::fed_rooms as fr;
12use crate::federation::{self as fed, enc};
13
14fn internal<E>(_: E) -> MatrixError {
15 MatrixError::internal()
16}
17
18fn now_rfc3339() -> String {
19 chrono::Utc::now().to_rfc3339()
20}
21
22pub(crate) async fn fed_request(state: &Arc<Homeserver>, dest: &str, method: &str, path: &str, body: Option<Value>) -> Result<(u16, Value), MatrixError> {
24 let transport = state.fed_transport.get().cloned().ok_or_else(|| MatrixError::unknown("federation transport is not configured"))?;
25 let local = crate::store::matrix_server_name();
26 let uri = format!("/_matrix{path}");
27 let (key_id, key) = with_conn_pub(state, |c| fed::active_signing_key(c, fed::now_ms()).map_err(internal)).await?;
28 let auth = fed::build_x_matrix_header(local, dest, &key_id, &key, method, &uri, body.as_ref());
29 transport.request(dest, method, &uri, &auth, body.as_ref()).await.map_err(|e| MatrixError::unknown(format!("federation to {dest}: {e}")))
30}
31
32pub(crate) async fn fed_request_raw(state: &Arc<Homeserver>, dest: &str, path: &str, max_bytes: usize) -> Result<crate::federation::RawResponse, MatrixError> {
34 let transport = state.fed_transport.get().cloned().ok_or_else(|| MatrixError::unknown("federation transport is not configured"))?;
35 let local = crate::store::matrix_server_name();
36 let uri = format!("/_matrix{path}");
37 let (key_id, key) = with_conn_pub(state, |c| fed::active_signing_key(c, fed::now_ms()).map_err(internal)).await?;
38 let auth = fed::build_x_matrix_header(local, dest, &key_id, &key, "GET", &uri, None);
39 transport.request_raw(dest, "GET", &uri, &auth, max_bytes).await.map_err(|e| MatrixError::unknown(format!("federation to {dest}: {e}")))
40}
41
42pub(crate) async fn verify_pdu(state: &Arc<Homeserver>, pdu: &Value) -> Result<bool, MatrixError> {
45 let (signer, key_id) = fr::pdu_signer(pdu).ok_or_else(|| MatrixError::bad_json("pdu is not signed by its sender's server"))?;
46 if crate::store::is_local_server_name(&signer) {
47 let id = pdu.get("event_id").and_then(Value::as_str).unwrap_or_default().to_string();
48 let known = with_conn_pub(state, move |c| Ok(crate::store::get_event(c, &id)?.is_some())).await?;
49 return if known { Ok(true) } else { Err(MatrixError::forbidden("pdu claims to be ours but is unknown")) };
50 }
51 let pk = super::federation::remote_key(state, &signer, &key_id).await?;
52 fr::verify_pdu_with_key(pdu, &signer, &key_id, &pk).map_err(|_| MatrixError::forbidden("bad pdu signature"))
53}
54
55#[cfg(feature = "f3-hash-ids")]
57pub(crate) async fn f3_keys(state: &Arc<Homeserver>, pdus: &[Value]) -> Result<m4a_matrix_core::PublicKeyMap, MatrixError> {
58 let owned = pdus.to_vec();
59 let (mut keys, need) = with_conn_pub(state, move |c| {
60 let mut keys = m4a_matrix_core::PublicKeyMap::new();
61 let refs: Vec<&Value> = owned.iter().collect();
62 let need = crate::f3::missing_keys(c, &refs, &mut keys);
63 Ok((keys, need))
64 })
65 .await?;
66 for (server, key_id) in need {
67 let pk = super::federation::remote_key(state, &server, &key_id).await?;
68 crate::f3::add_key(&mut keys, &server, &key_id, &pk)?;
69 }
70 Ok(keys)
71}
72
73#[cfg(feature = "f3-hash-ids")]
77pub(crate) async fn f3_receive_live(state: &Arc<Homeserver>, origin: &str, pdu: &Value) -> Result<crate::f3::Received, MatrixError> {
78 for _ in 0..3 {
79 let keys = f3_keys(state, std::slice::from_ref(pdu)).await?;
80 let (p, o) = (pdu.clone(), origin.to_string());
81 let res = with_conn_pub(state, move |c| Ok(crate::f3::receive_detail(c, Some(&o), &p, &keys, &now_rfc3339()))).await?;
82 match res {
83 Ok(r) => return Ok(r),
84 Err(crate::f3::RecvErr::Other(e)) => return Err(e),
85 Err(crate::f3::RecvErr::Missing(_)) => f3_fetch_missing(state, origin, pdu).await?,
86 }
87 }
88 Err(MatrixError::forbidden("the event still cites unknown events after fetching"))
89}
90
91#[cfg(feature = "f3-hash-ids")]
93async fn f3_fetch_missing(state: &Arc<Homeserver>, origin: &str, pdu: &Value) -> Result<(), MatrixError> {
94 let room = pdu.get("room_id").and_then(Value::as_str).ok_or_else(|| MatrixError::bad_json("room_id"))?.to_string();
95 let latest = crate::f3::wire_id(pdu)?;
96 let r2 = room.clone();
97 let earliest = with_conn_pub(state, move |c| Ok(crate::f3::extremity_ids(c, &r2))).await?;
98 let body = json!({ "earliest_events": earliest, "latest_events": [latest], "limit": 50, "min_depth": 0 });
99 let (status, resp) = fed_request(state, origin, "POST", &format!("/federation/v1/get_missing_events/{}", enc(&room)), Some(body)).await?;
100 if status != 200 {
101 return Err(map_status(status, "get_missing_events"));
102 }
103 let events = resp.get("events").and_then(Value::as_array).cloned().unwrap_or_default();
104 if events.is_empty() {
105 return Err(MatrixError::forbidden("the origin returned no missing events"));
106 }
107 f3_store_history(state, &room, events).await.map(|_| ())
108}
109
110#[cfg(feature = "f3-hash-ids")]
112pub(crate) async fn f3_store_history(state: &Arc<Homeserver>, room: &str, events: Vec<Value>) -> Result<usize, MatrixError> {
113 let keys = f3_keys(state, &events).await?;
114 let room = room.to_string();
115 with_conn_pub(state, move |c| crate::f3::process_historic(c, &room, &events, &keys, &now_rfc3339())).await
116}
117
118fn backoff_ms(attempts: i64) -> i64 {
119 (1000i64 << attempts.clamp(0, 10)).min(600_000)
120}
121
122struct Row {
123 id: i64,
124 kind: String,
125 room_id: String,
126 event_id: String,
127 payload: Value,
128 attempts: i64,
129}
130
131pub async fn drain_outbox(state: &Arc<Homeserver>) -> usize {
133 if state.federation_enabled.get().is_none() || state.fed_transport.get().is_none() {
134 return 0;
135 }
136 let local = crate::store::matrix_server_name().to_string();
137 let now = fed::now_ms();
138 let l2 = local.clone();
139 let rows = with_conn_pub(state, move |c| {
140 if let Err(e) = fr::export_local_events(c, &l2, now) {
141 tracing::warn!("federation export: {e}");
142 }
143 let mut stmt = c
144 .prepare("SELECT id, destination, kind, room_id, event_id, payload, attempts FROM fed_outbox WHERE next_try_ms <= ?1 ORDER BY id LIMIT 300")
145 .map_err(internal)?;
146 let it = stmt
147 .query_map([now], |r| {
148 Ok((
149 r.get::<_, String>(1)?,
150 Row { id: r.get(0)?, kind: r.get(2)?, room_id: r.get(3)?, event_id: r.get(4)?, payload: serde_json::from_str(&r.get::<_, String>(5)?).unwrap_or(Value::Null), attempts: r.get(6)? },
151 ))
152 })
153 .map_err(internal)?;
154 it.collect::<Result<Vec<_>, _>>().map_err(internal)
155 })
156 .await
157 .unwrap_or_default();
158 let mut by_dest: BTreeMap<String, Vec<Row>> = BTreeMap::new();
159 for (d, r) in rows {
160 by_dest.entry(d).or_default().push(r);
161 }
162 let mut delivered = 0;
163 for (dest, items) in by_dest {
164 let mut batch: Vec<Row> = Vec::new();
165 let mut failed_from: Option<usize> = None;
166 let mut idx = 0;
167 let total = items.len();
168 let mut items = items.into_iter().peekable();
169 while let Some(item) = items.next() {
170 let is_invite = item.kind == "invite";
171 if is_invite {
172 if !batch.is_empty() {
174 let b = std::mem::take(&mut batch);
175 match send_batch(state, &local, &dest, &b).await {
176 Ok(n) => delivered += n,
177 Err(_) => {
178 mark_failed(state, &b).await;
179 mark_failed(state, std::slice::from_ref(&item)).await;
180 failed_from = Some(idx);
181 break;
182 }
183 }
184 }
185 #[allow(unused_mut)]
186 let mut body = json!({ "room_version": "11", "event": item.payload["event"], "room_info": item.payload["room_info"], "invite_room_state": item.payload["state"] });
187 if let Some(snap) = item.payload.get("f3") {
188 body["f3"] = snap.clone();
189 }
190 let path = format!("/federation/v2/invite/{}/{}", enc(&item.room_id), enc(&item.event_id));
191 match fed_request(state, &dest, "PUT", &path, Some(body)).await {
192 Ok((200, _)) | Ok((403, _)) => {
193 delete_rows(state, &[item.id]).await;
194 delivered += 1;
195 }
196 _ => {
197 mark_failed(state, std::slice::from_ref(&item)).await;
198 failed_from = Some(idx);
199 break;
200 }
201 }
202 } else {
203 batch.push(item);
204 if batch.len() >= 20 || items.peek().is_none_or(|n| n.kind == "invite") {
205 let b = std::mem::take(&mut batch);
206 match send_batch(state, &local, &dest, &b).await {
207 Ok(n) => delivered += n,
208 Err(_) => {
209 mark_failed(state, &b).await;
210 failed_from = Some(idx);
211 break;
212 }
213 }
214 }
215 }
216 idx += 1;
217 }
218 let _ = (failed_from, total);
219 }
220 delivered
221}
222
223async fn delete_rows(state: &Arc<Homeserver>, ids: &[i64]) {
224 let ids = ids.to_vec();
225 let _ = with_conn_pub(state, move |c| {
226 for id in ids {
227 c.execute("DELETE FROM fed_outbox WHERE id = ?1", [id]).map_err(internal)?;
228 }
229 Ok(())
230 })
231 .await;
232}
233
234async fn mark_failed(state: &Arc<Homeserver>, rows: &[Row]) {
235 let now = fed::now_ms();
236 let data: Vec<(i64, i64)> = rows.iter().map(|r| (r.id, r.attempts)).collect();
237 let _ = with_conn_pub(state, move |c| {
238 for (id, attempts) in data {
239 c.execute("UPDATE fed_outbox SET attempts = attempts + 1, next_try_ms = ?2 WHERE id = ?1", rusqlite::params![id, now + backoff_ms(attempts)]).map_err(internal)?;
240 }
241 Ok(())
242 })
243 .await;
244}
245
246async fn send_batch(state: &Arc<Homeserver>, local: &str, dest: &str, rows: &[Row]) -> Result<usize, MatrixError> {
247 let (mut pdus, mut edus) = (Vec::new(), Vec::new());
248 for r in rows {
249 if r.kind == "edu" {
250 edus.push(r.payload.clone());
251 } else {
252 pdus.push(r.payload.clone());
253 }
254 }
255 let txn = format!("m4a{}", rows[0].id);
256 let body = json!({ "origin": local, "origin_server_ts": fed::now_ms(), "pdus": pdus, "edus": edus });
257 let (status, _) = fed_request(state, dest, "PUT", &format!("/federation/v1/send/{}", enc(&txn)), Some(body)).await?;
258 if status != 200 {
259 return Err(MatrixError::unknown(format!("destination answered {status}")));
260 }
261 delete_rows(state, &rows.iter().map(|r| r.id).collect::<Vec<_>>()).await;
262 Ok(rows.len())
263}
264
265pub fn spawn_outbox_worker(state: Arc<Homeserver>) {
267 tokio::spawn(async move {
268 loop {
269 drain_outbox(&state).await;
270 let _ = tokio::time::timeout(std::time::Duration::from_secs(3), state.fed_notify.notified()).await;
271 }
272 });
273}
274
275pub(crate) fn room_domain_is_remote(room_id: &str) -> bool {
279 room_id.split_once(':').is_some_and(|(_, d)| !crate::store::is_local_server_name(d))
280}
281
282fn map_status(status: u16, what: &str) -> MatrixError {
283 match status {
284 403 => MatrixError::forbidden(format!("{what}: refused by the remote server")),
285 404 => MatrixError::not_found(format!("{what}: not found on the remote server")),
286 _ => MatrixError::unknown(format!("{what}: remote server answered {status}")),
287 }
288}
289
290pub(crate) async fn federated_join(state: &Arc<Homeserver>, caller: &super::Caller, room_id: &str) -> Result<Vec<i64>, MatrixError> {
292 let dest = room_id.split_once(':').map(|(_, d)| d.to_string()).ok_or_else(|| MatrixError::invalid_param("room id"))?;
293 let local = crate::store::matrix_server_name().to_string();
294 let (status, tmpl) = fed_request(state, &dest, "GET", &format!("/federation/v1/make_join/{}/{}{}", enc(room_id), enc(&caller.mxid), if cfg!(feature = "f3-hash-ids") { "?ver=11" } else { "" }), None).await?;
295 if status != 200 {
296 return Err(map_status(status, "make_join"));
297 }
298 #[cfg(feature = "f3-hash-ids")]
299 if tmpl.get("m4a_f3").and_then(Value::as_bool) == Some(true) {
300 return federated_join_f3(state, caller, room_id, &dest, &tmpl).await;
301 }
302 #[cfg(feature = "f3-hash-ids")]
303 if tmpl.get("room_version").and_then(Value::as_str).is_some() && tmpl.pointer("/event/prev_events").is_some() && tmpl.get("m4a_f3").is_none() {
304 return federated_join_spec(state, caller, room_id, &dest, &tmpl).await;
305 }
306 let mut ev = tmpl.get("event").and_then(Value::as_object).cloned().ok_or_else(|| MatrixError::unknown("make_join: no event"))?;
307 let uid = caller.user_id;
308 let label = with_conn_pub(state, move |c| crate::nick::effective_label(c, uid).map_err(internal)).await?;
309 let event_id = crate::store::new_event_id();
310 ev.insert("event_id".into(), json!(event_id));
311 ev.insert("origin".into(), json!(local));
312 let mut content = ev.get("content").cloned().unwrap_or_else(|| json!({}));
313 if !label.is_empty() {
314 content["displayname"] = json!(label);
315 }
316 content["membership"] = json!("join");
317 ev.insert("content".into(), content);
318 if ev.get("sender").and_then(Value::as_str) != Some(caller.mxid.as_str()) || ev.get("room_id").and_then(Value::as_str) != Some(room_id) {
319 return Err(MatrixError::unknown("make_join: template does not match the request"));
320 }
321 let l2 = local.clone();
322 let pdu = with_conn_pub(state, move |c| fr::finalize_pdu(c, ev, &l2, fed::now_ms()).map_err(internal)).await?;
323 let (status, resp) = fed_request(state, &dest, "PUT", &format!("/federation/v2/send_join/{}/{}", enc(room_id), enc(&event_id)), Some(pdu.clone())).await?;
324 if status != 200 {
325 return Err(map_status(status, "send_join"));
326 }
327 let list = |k: &str| resp.get(k).and_then(Value::as_array).cloned().unwrap_or_default();
328 let (state_pdus, timeline) = (list("state"), list("m4a_timeline"));
329 let info = resp.get("m4a_room").cloned().ok_or_else(|| MatrixError::unknown("send_join: no room info"))?;
330 let mut verified: Vec<(Value, bool)> = Vec::new();
331 for p in state_pdus.iter().chain(timeline.iter()) {
332 if p.get("event_id").and_then(Value::as_str) == Some(event_id.as_str()) {
333 continue; }
335 let ok = verify_pdu(state, p).await?;
336 verified.push((p.clone(), ok));
337 }
338 let room = room_id.to_string();
339 let (mxid, user_id) = (caller.mxid.clone(), caller.user_id);
340 let ids = with_conn_pub(state, move |c| {
341 let now = now_rfc3339();
342 fr::create_replica_room(c, &room, &info, &now)?;
343 for (p, ok) in &verified {
344 fr::ingest_pdu(c, p, *ok, fr::Mode::Trusted, &now)?;
345 }
346 fr::store_own_pdu(c, &pdu, user_id, &mxid, &now)?;
347 let ids = crate::rooms::member_and_invited_ids(c, &room).map_err(internal)?;
348 Ok(ids.into_iter().collect::<Vec<_>>())
349 })
350 .await?;
351 Ok(ids)
352}
353
354#[cfg(feature = "f3-hash-ids")]
357async fn federated_join_f3(state: &Arc<Homeserver>, caller: &super::Caller, room_id: &str, dest: &str, tmpl: &Value) -> Result<Vec<i64>, MatrixError> {
358 let mut ev = tmpl.get("event").cloned().ok_or_else(|| MatrixError::unknown("make_join: no event"))?;
359 if ev.get("sender").and_then(Value::as_str) != Some(caller.mxid.as_str()) || ev.get("room_id").and_then(Value::as_str) != Some(room_id) {
360 return Err(MatrixError::unknown("make_join: template does not match the request"));
361 }
362 let uid = caller.user_id;
363 let label = with_conn_pub(state, move |c| crate::nick::effective_label(c, uid).map_err(internal)).await?;
364 if !label.is_empty() {
365 ev["content"]["displayname"] = json!(label);
366 }
367 let (event_id, pdu) = with_conn_pub(state, move |c| crate::f3::sign_own(c, &ev, fed::now_ms())).await?;
368 let (status, resp) = fed_request(state, dest, "PUT", &format!("/federation/v2/send_join/{}/{}", enc(room_id), enc(&event_id)), Some(pdu)).await?;
369 if status != 200 {
370 return Err(map_status(status, "send_join"));
371 }
372 let snap = resp.get("m4a_f3").cloned().ok_or_else(|| MatrixError::unknown("send_join: no room snapshot"))?;
373 let info = resp.get("m4a_room").cloned().ok_or_else(|| MatrixError::unknown("send_join: no room info"))?;
374 let all = crate::f3::snapshot_events(&snap);
375 let keys = f3_keys(state, &all).await?;
376 let room = room_id.to_string();
377 let shared = matches!(info.get("history_visibility").and_then(Value::as_str), Some("shared") | Some("world_readable"));
378 let ids = with_conn_pub(state, move |c| {
379 crate::f3::import_snapshot(c, &room, &info, &snap, &keys, &now_rfc3339())?;
380 let ids = crate::rooms::member_and_invited_ids(c, &room).map_err(internal)?;
381 Ok(ids.into_iter().collect::<Vec<_>>())
382 })
383 .await?;
384 if shared {
386 let q = format!("/federation/v1/backfill/{}?v={}&limit=100", enc(room_id), enc(&event_id));
387 match fed_request(state, dest, "GET", &q, None).await {
388 Ok((200, resp)) => {
389 let events = resp.get("pdus").and_then(Value::as_array).cloned().unwrap_or_default();
390 if let Err(e) = f3_store_history(state, room_id, events).await {
391 tracing::warn!("backfill after join: {}", e.error);
392 }
393 }
394 other => tracing::warn!("backfill after join: {:?}", other.map(|x| x.0)),
395 }
396 }
397 Ok(ids)
398}
399
400#[cfg(feature = "f3-hash-ids")]
404async fn federated_join_spec(state: &Arc<Homeserver>, caller: &super::Caller, room_id: &str, dest: &str, tmpl: &Value) -> Result<Vec<i64>, MatrixError> {
405 if tmpl.get("room_version").and_then(Value::as_str) != Some("11") {
406 return Err(MatrixError::new(400, "M_INCOMPATIBLE_ROOM_VERSION", "only room version 11 can be joined over federation"));
407 }
408 let t = tmpl.get("event").cloned().unwrap_or_default();
409 if t.get("sender").and_then(Value::as_str) != Some(caller.mxid.as_str()) || t.get("room_id").and_then(Value::as_str) != Some(room_id) {
410 return Err(MatrixError::unknown("make_join: template does not match the request"));
411 }
412 let uid = caller.user_id;
413 let label = with_conn_pub(state, move |c| crate::nick::effective_label(c, uid).map_err(internal)).await?;
414 let mut content = t.get("content").cloned().unwrap_or_else(|| json!({}));
415 content["membership"] = json!("join");
416 if !label.is_empty() {
417 content["displayname"] = json!(label);
418 }
419 let mut ev = json!({ "content": content, "origin_server_ts": fed::now_ms() });
420 for k in ["room_id", "sender", "type", "state_key", "prev_events", "auth_events", "depth"] {
421 if let Some(v) = t.get(k) {
422 ev[k] = v.clone();
423 }
424 }
425 let (event_id, pdu) = with_conn_pub(state, move |c| crate::f3::sign_own(c, &ev, fed::now_ms())).await?;
426 let (status, resp) = fed_request(state, dest, "PUT", &format!("/federation/v2/send_join/{}/{}", enc(room_id), enc(&event_id)), Some(pdu.clone())).await?;
427 if status != 200 {
428 return Err(map_status(status, "send_join"));
429 }
430 let list = |k: &str| resp.get(k).and_then(Value::as_array).cloned().unwrap_or_default();
431 let st = list("state");
432 if resp.get("members_omitted").and_then(Value::as_bool) == Some(true) {
433 return Err(MatrixError::unknown("send_join: partial state is not supported"));
434 }
435 let find = |ty: &str| st.iter().find(|e| e.get("type").and_then(Value::as_str) == Some(ty) && e.get("state_key").and_then(Value::as_str) == Some("")).cloned();
436 let content_of = |ty: &str, k: &str| find(ty).and_then(|e| e.pointer(&format!("/content/{k}")).and_then(Value::as_str).map(str::to_string));
437 let info = json!({
438 "kind": "group",
439 "is_encrypted": find("m.room.encryption").is_some(),
440 "join_rule": content_of("m.room.join_rules", "join_rule").unwrap_or_else(|| "invite".into()),
441 "history_visibility": content_of("m.room.history_visibility", "history_visibility").unwrap_or_else(|| "shared".into()),
442 "room_version": "11",
443 "creator": find("m.room.create").and_then(|e| e.get("sender").and_then(Value::as_str).map(str::to_string)).unwrap_or_default(),
444 });
445 let mut state_v = st.clone();
446 state_v.push(pdu.clone());
447 let snap = json!({ "state": state_v, "extremities": [pdu], "auth_chain": list("auth_chain") });
448 let all = crate::f3::snapshot_events(&snap);
449 let keys = f3_keys(state, &all).await?;
450 let room = room_id.to_string();
451 let shared = matches!(info["history_visibility"].as_str(), Some("shared") | Some("world_readable"));
452 let ids = with_conn_pub(state, move |c| {
453 crate::f3::import_snapshot(c, &room, &info, &snap, &keys, &now_rfc3339())?;
454 let ids = crate::rooms::member_and_invited_ids(c, &room).map_err(internal)?;
455 Ok(ids.into_iter().collect::<Vec<_>>())
456 })
457 .await?;
458 if shared {
459 let q = format!("/federation/v1/backfill/{}?v={}&limit=100", enc(room_id), enc(&event_id));
460 if let Ok((200, resp)) = fed_request(state, dest, "GET", &q, None).await {
461 let events = resp.get("pdus").and_then(Value::as_array).cloned().unwrap_or_default();
462 if let Err(e) = f3_store_history(state, room_id, events).await {
463 tracing::warn!("backfill after join: {}", e.error);
464 }
465 }
466 }
467 Ok(ids)
468}
469
470pub(crate) fn split_by_domain<V: Clone>(map: &BTreeMap<String, V>) -> (BTreeMap<String, V>, BTreeMap<String, BTreeMap<String, V>>) {
474 let (mut local, mut remote): (BTreeMap<String, V>, BTreeMap<String, BTreeMap<String, V>>) = (BTreeMap::new(), BTreeMap::new());
475 for (k, v) in map {
476 if fr::is_remote_mxid(k) {
477 if let Some(d) = fr::domain_of(k) {
478 remote.entry(d.to_string()).or_default().insert(k.clone(), v.clone());
479 }
480 } else {
481 local.insert(k.clone(), v.clone());
482 }
483 }
484 (local, remote)
485}
486
487pub(crate) async fn merge_remote_keys_query(state: &Arc<Homeserver>, remote: BTreeMap<String, BTreeMap<String, Vec<String>>>, out: &mut Value) {
489 for (domain, users) in remote {
490 let body = json!({ "device_keys": users });
491 match fed_request(state, &domain, "POST", "/federation/v1/user/keys/query", Some(body)).await {
492 Ok((200, resp)) => {
493 for k in ["device_keys", "master_keys", "self_signing_keys"] {
494 if let (Some(dst), Some(src)) = (out[k].as_object_mut(), resp.get(k).and_then(Value::as_object)) {
495 for (u, v) in src {
496 dst.insert(u.clone(), v.clone());
497 }
498 }
499 }
500 }
501 _ => {
502 out["failures"][&domain] = json!({ "status": 502, "errcode": "M_UNREACHABLE" });
503 }
504 }
505 }
506}
507
508pub(crate) async fn merge_remote_keys_claim(state: &Arc<Homeserver>, remote: BTreeMap<String, BTreeMap<String, BTreeMap<String, String>>>, out: &mut Value) {
510 for (domain, users) in remote {
511 let body = json!({ "one_time_keys": users });
512 match fed_request(state, &domain, "POST", "/federation/v1/user/keys/claim", Some(body)).await {
513 Ok((200, resp)) => {
514 if let (Some(dst), Some(src)) = (out["one_time_keys"].as_object_mut(), resp.get("one_time_keys").and_then(Value::as_object)) {
515 for (u, v) in src {
516 dst.insert(u.clone(), v.clone());
517 }
518 }
519 }
520 _ => {
521 out["failures"][&domain] = json!({ "status": 502, "errcode": "M_UNREACHABLE" });
522 }
523 }
524 }
525}
526
527pub(crate) fn enqueue_to_device_edu(conn: &rusqlite::Connection, domain: &str, sender: &str, event_type: &str, message_id: &str, messages: &Value) -> Result<(), MatrixError> {
529 let edu = json!({ "edu_type": "m.direct_to_device", "content": { "sender": sender, "type": event_type, "message_id": message_id, "messages": messages } });
530 fr::enqueue(conn, domain, "edu", "", "", &edu, fed::now_ms()).map_err(internal)
531}