Skip to main content

mail4agent_server/http/
fed_net.rs

1//! Outgoing federation: signed requests, outbox delivery, PDU verification,
2//! remote join, and the remote halves of key query/claim and to-device.
3
4use 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
22/// Send one signed request to `dest` (path without the `/_matrix` prefix, query included).
23pub(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
32/// A signed GET whose answer is binary (media); `path` is without the `/_matrix` prefix.
33pub(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
42/// Verify a PDU's signature with its signer's keys. `Ok(false)`: signature
43/// good, content hash wrong (keep only the redacted form).
44pub(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/// Public keys of everyone who signed `pdus` (our own from the database, peers' fetched).
56#[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/// A peer's live event for a DAG room. When it cites events this server does not hold, they are
74/// fetched from `origin` (get_missing_events), verified and stored first, then the event is
75/// processed again; the fork this makes is merged by state resolution.
76#[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/// Ask `origin` for the events between our forward extremities and `pdu`, check and store them.
92#[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/// Verify (keys fetched as needed) and store events of the past.
111#[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
131/// Export new local events and deliver everything due. Returns the number of delivered items.
132pub 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                // flush the batch first to keep order
173                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
265/// Background worker: drains whenever poked or every few seconds.
266pub 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
275// -------------------------------------------------------------- remote join
276
277/// Whether `room_id` belongs to another server's namespace.
278pub(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
290/// Join a room that lives on another server (make_join, send_join, ingest).
291pub(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; // our own join, echoed back; stored below
334        }
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/// Join a DAG room: the peer's template (with `prev_events`/`auth_events`/`depth`) is filled in and
355/// signed here; the answer is the room as it stood right after the join, from which the replica starts.
356#[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    // A late joiner catches up the history, parallel branches included, when the room shares it.
385    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/// Join a room of a spec-following server (Synapse and kin): their make_join template is filled in
401/// and signed here, their send_join answer (state before the join, plus auth chain) is verified in
402/// full, and the replica starts from that state plus our own join. Room version 11 only.
403#[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
470// ------------------------------------------------------------ remote keys
471
472/// Split `{mxid: v}` into (local part, per-remote-domain part).
473pub(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
487/// Query remote servers for device keys and merge into `out` (a client keys/query response).
488pub(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
508/// Claim one-time keys from remote servers and merge into `out`.
509pub(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
527/// Queue a `m.direct_to_device` EDU for a remote server.
528pub(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}