From 83bf1ee1df97e787e70b22ef1609d8aafbde5a19 Mon Sep 17 00:00:00 2001 From: Jacob Taylor Date: Mon, 5 Oct 2026 07:47:01 -0700 Subject: [PATCH] feat: Continuwuity notarizes the notary --- changelog.d/2314.feature.md | 1 + conduwuit-example.toml | 7 +- src/api/router.rs | 2 + src/api/server/key.rs | 282 ++++++++++++++++++++++++++++- src/core/config/mod.rs | 18 +- src/service/server_keys/mod.rs | 28 ++- src/service/server_keys/request.rs | 12 +- src/service/server_keys/verify.rs | 30 +-- 8 files changed, 339 insertions(+), 41 deletions(-) create mode 100644 changelog.d/2314.feature.md diff --git a/changelog.d/2314.feature.md b/changelog.d/2314.feature.md new file mode 100644 index 000000000..8fdd14ea1 --- /dev/null +++ b/changelog.d/2314.feature.md @@ -0,0 +1 @@ +Added support for acting as a trusted server, meaning the `trusted_servers` list can consist of servers that are either Synapse OR continuwuity. diff --git a/conduwuit-example.toml b/conduwuit-example.toml index d0bf4e51f..db8dffb25 100644 --- a/conduwuit-example.toml +++ b/conduwuit-example.toml @@ -894,10 +894,7 @@ # (trusted key notary servers), as well as serve as trusted servers for # other operations like backfill and event fetching. # -# Currently, continuwuity doesn't support inbound key requests, so -# this list should only contain other Synapse servers. -# -# example: ["matrix.org", "starstruck.systems"] +# example: ["continuwuity.org", "matrix.org"] # # It is possible to restrict which signing keys trusted servers are # allowed to sign responses with. Without configuring this, all responses @@ -934,7 +931,7 @@ # empty list (`[]`) to operate without trusted server assistance, but this # is discouraged for performance and reliability reasons. # -#trusted_servers = ["matrix.org"] +#trusted_servers = ["continuwuity.org", "continuwuity.rocks", "matrix.org"] # Whether to prioritise trusted server lookups over origin server lookups. # diff --git a/src/api/router.rs b/src/api/router.rs index de2736e2b..2496be696 100644 --- a/src/api/router.rs +++ b/src/api/router.rs @@ -203,6 +203,8 @@ pub fn build(router: Router, state: State) -> Router { router = router .ruma_route(&server::get_server_version_route) .route("/_matrix/key/v2/server", get(server::get_server_keys_route)) + .ruma_route(&server::get_remote_server_keys_batch_route) + .ruma_route(&server::get_remote_server_keys_route) .ruma_route(&server::get_public_rooms_route) .ruma_route(&server::get_public_rooms_filtered_route) .ruma_route(&server::send_transaction_message_route) diff --git a/src/api/server/key.rs b/src/api/server/key.rs index b547ea15c..42dcc26d3 100644 --- a/src/api/server/key.rs +++ b/src/api/server/key.rs @@ -1,19 +1,37 @@ -use std::{collections::BTreeMap, mem::take, time::Duration}; +use std::{ + collections::{BTreeMap, HashMap, HashSet}, + mem::take, + ops::Index, + sync::Arc, + time::{Duration, Instant}, +}; use axum::{Json, extract::State, response::IntoResponse}; -use conduwuit::{Result, utils::timepoint_from_now}; -use futures::StreamExt; +use conduwuit::{ + Err, Result, debug, debug_info, error, + utils::{ReadyExt, stream::BroadbandExt, timepoint_from_now, to_canonical_object}, + warn, +}; +use futures::{StreamExt, stream::FuturesUnordered}; use ruma::{ - MilliSecondsSinceUnixEpoch, + MilliSecondsSinceUnixEpoch, OwnedServerName, OwnedServerSigningKeyId, ServerName, api::{ OutgoingResponseExt, - federation::discovery::{OldVerifyKey, ServerSigningKeys, get_server_keys}, + federation::discovery::{ + OldVerifyKey, ServerSigningKeys, get_remote_server_keys, + get_remote_server_keys_batch, get_remote_server_keys_batch::v2::QueryCriteria, + get_server_keys, + }, }, assign, serde::Raw, + uint, }; +use serde_json::value::to_raw_value; use service::{server_keys, server_keys::in_one_week}; +use crate::router::Ruma; + /// # `GET /_matrix/key/v2/server` /// /// Gets the public signing keys of this server. @@ -87,3 +105,257 @@ fn valid_until_ts() -> MilliSecondsSinceUnixEpoch { let timepoint = timepoint_from_now(dur).expect("SystemTime should not overflow"); MilliSecondsSinceUnixEpoch::from_system_time(timepoint).expect("UInt should not overflow") } + +const MAX_KEYS_PER_QUERY: usize = 16384; +const MAX_SERVERS_PER_QUERY: usize = 4096; + +pub(crate) async fn get_remote_server_keys_batch_route( + State(services): State, + body: Ruma, +) -> Result { + let start = (Instant::now(), MilliSecondsSinceUnixEpoch::now()); + if body.server_keys.is_empty() { + return Ok(get_remote_server_keys_batch::v2::Response::new(Vec::new())); + } + + let total_queried_servers = body.server_keys.len(); + let total_queried_keys = body + .server_keys + .values() + .fold(body.server_keys.len(), |acc, q| acc.saturating_add(q.len())); + + if body.server_keys.len() > MAX_SERVERS_PER_QUERY { + // TODO(nex): enforce once MSC4556 is merged + warn!( + %total_queried_servers, + %total_queried_keys, + "Received a large notary request (too many servers)" + ); + // return Err!(Request(TooLarge( + // "Too many server keys requested ({} > {MAX_SERVERS_PER_QUERY}", + // body.server_keys.len() + // ))); + } + if total_queried_keys > MAX_KEYS_PER_QUERY { + // We shouldn't really enforce this before MSC4456 either, but not doing + // so may cause performance degradation. + warn!( + %total_queried_servers, + %total_queried_keys, + "Received a huge notary request (too many keys), rejecting", + ); + return Err!(Request(TooLarge( + "Too many keys requested ({total_queried_keys} > {MAX_KEYS_PER_QUERY})" + ))); + } + + debug_info!("Fetching {total_queried_keys} keys across {} servers", body.server_keys.len()); + let mut futs: FuturesUnordered<_> = FuturesUnordered::new(); + for (server_name, queries) in body.server_keys.clone() { + futs.push(acquire_keys_as_notary( + services.server_keys.clone(), + services.globals.server_name().to_owned(), + server_name, + queries, + start.1, + )); + } + let mut response = Vec::with_capacity(total_queried_keys); + while let Some(v) = futs.next().await { + response.extend(v); + } + + debug_info!( + elapsed=?start.0.elapsed(), + %total_queried_servers, + %total_queried_keys, + "Fetched {} key responses", + response.len() + ); + Ok(get_remote_server_keys_batch::v2::Response::new(response)) +} + +async fn sign_ssk( + server_keys: &server_keys::Service, + ssk: ServerSigningKeys, + server_name: &ServerName, + our_name: &ServerName, +) -> Result> { + let mut canonical = to_canonical_object(&ssk)?; + server_keys.sign_json(&mut canonical)?; + server_keys::strip_extraneous_signatures(&mut canonical, server_name, &[our_name]); + to_raw_value(&canonical) + .map(Raw::::from_json) + .map_err(Into::into) +} + +#[tracing::instrument(skip(server_keys, my_name, queries, start))] +async fn acquire_keys_as_notary( + server_keys: Arc, + my_name: OwnedServerName, + remote: OwnedServerName, + mut queries: BTreeMap, + start: MilliSecondsSinceUnixEpoch, +) -> Vec> { + let mut results = Vec::with_capacity(queries.len().max(1)); + let mut keymap = HashMap::with_capacity(queries.len().max(1)); + + // First contact the origin (if we're allowed to) + if server_keys.notary_may_contact_origin(&remote) { + debug_info!("Asking remote directly for verify keys"); + if let Ok(res) = server_keys.origin_request(remote.clone(), start).await + && let Ok(ssk) = sign_ssk(&server_keys, res.clone(), &remote, &my_name).await + { + let index = results.len(); + results.push(ssk); + for key_id in res.verify_keys.keys().chain(res.old_verify_keys.keys()) { + if queries.is_empty() + || queries + .get(key_id) + .and_then(|c| c.minimum_valid_until_ts) + .is_none_or(|m| res.valid_until_ts >= m) + { + queries.remove(key_id); + keymap.insert(key_id.to_owned(), index); + } + } + } + } else { + debug!("Not asking remote for keys (already asked recently)"); + } + debug!(keys=?keymap.keys(), "Live verify keys"); + if queries.is_empty() && keymap.is_empty() { + // If the server asked for all keys, AND we didn't get anything from the + // origin, just fetch any fresh responses we have. + server_keys + .signing_keys_for(&remote) + .broad_filter_map(|ssk| { + let server_name = &remote; + let my_name = &my_name; + let server_keys = &server_keys; + async move { + if ssk.valid_until_ts > in_one_week() || ssk.valid_until_ts < start { + return None; + } + + sign_ssk(server_keys, ssk, server_name, my_name).await.ok() + } + }) + .ready_for_each(|signed_ssk| { + let ssk = signed_ssk.deserialize().unwrap(); + let idx = results.len(); + results.push(signed_ssk); + for key_id in ssk.verify_keys.keys().chain(ssk.old_verify_keys.keys()) { + keymap.insert(key_id.to_owned(), idx); + } + }) + .await; + } + + // If we're still missing some keys, fetch them from the local cache + for (key_id, criteria) in queries { + debug!(%key_id, "Fetching verify key from local repository"); + + if keymap.contains_key(&key_id) { + debug!(%key_id, "Already found key"); + continue; + } + + let minimum_valid_until_ts = criteria + .clone() + .minimum_valid_until_ts + .unwrap_or_else(|| MilliSecondsSinceUnixEpoch(uint!(0))); + + let Some(ssk) = server_keys.get_signing_key(&remote, &key_id).await else { + debug!(%remote, %key_id, "Could not find a matching signing key locally."); + continue; + }; + + if ssk.valid_until_ts < minimum_valid_until_ts || ssk.valid_until_ts > in_one_week() { + debug!( + %key_id, + valid_until_ts=?ssk.valid_until_ts, + ?minimum_valid_until_ts, + "Stored verify key does not satisfy query criteria" + ); + continue; + } + + debug!(%key_id, ?ssk, "Found key locally"); + let rep_key_ids = ssk + .verify_keys + .keys() + .chain(ssk.old_verify_keys.keys()) + .cloned() + .collect::>(); + let Ok(signed_ssk) = sign_ssk(&server_keys, ssk.clone(), &remote, &my_name) + .await + .inspect_err( + |e| error!(%key_id, %remote, "Failed to sign signing keys chunk: {e:?}"), + ) + else { + continue; + }; + + let idx = results.len(); + results.push(signed_ssk); + for key_id in rep_key_ids { + keymap.insert(key_id, idx); + } + } + + keymap + .into_values() + .collect::>() + .into_iter() + .map(|idx| results.index(idx).to_owned()) + .collect() +} + +pub(crate) async fn get_remote_server_keys_route( + State(services): State, + body: Ruma, +) -> Result { + let min_valid_ts = body.minimum_valid_until_ts; + if services + .server_keys + .notary_may_contact_origin(&body.server_name) + && let Ok(response) = services + .server_keys + .origin_request(body.server_name.clone(), min_valid_ts) + .await + { + return sign_ssk( + &services.server_keys, + response, + &body.server_name, + services.globals.server_name(), + ) + .await + .map(|r| Ok(get_remote_server_keys::v2::Response::new(vec![r])))?; + } + + let response = services + .server_keys + .signing_keys_for(&body.server_name) + .broad_filter_map(|ssk| { + let server_name = body.server_name.clone(); + async move { + if ssk.valid_until_ts > in_one_week() || ssk.valid_until_ts < min_valid_ts { + return None; + } + + sign_ssk( + &services.server_keys, + ssk, + server_name.as_ref(), + services.globals.server_name(), + ) + .await + .ok() + } + }) + .collect::>() + .await; + Ok(get_remote_server_keys::v2::Response::new(response)) +} diff --git a/src/core/config/mod.rs b/src/core/config/mod.rs index ecdba258c..ad88fcb62 100644 --- a/src/core/config/mod.rs +++ b/src/core/config/mod.rs @@ -1092,10 +1092,7 @@ pub struct Config { /// (trusted key notary servers), as well as serve as trusted servers for /// other operations like backfill and event fetching. /// - /// Currently, continuwuity doesn't support inbound key requests, so - /// this list should only contain other Synapse servers. - /// - /// example: ["matrix.org", "starstruck.systems"] + /// example: ["continuwuity.org", "matrix.org"] /// /// It is possible to restrict which signing keys trusted servers are /// allowed to sign responses with. Without configuring this, all responses @@ -1132,7 +1129,7 @@ pub struct Config { /// empty list (`[]`) to operate without trusted server assistance, but this /// is discouraged for performance and reliability reasons. /// - /// default: ["matrix.org"] + /// default: ["continuwuity.org", "continuwuity.rocks", "matrix.org"] #[serde(default = "default_trusted_servers")] pub trusted_servers: Vec, @@ -3193,12 +3190,11 @@ fn default_otlp_protocol() -> String { "http".to_owned() } fn default_tracing_flame_output_path() -> String { "./tracing.folded".to_owned() } fn default_trusted_servers() -> Vec { - // TODO(nex): Once we can be a notary, add maintainer(?) homeservers here. - // Rationale: Users are already running our code, arguably that's a higher - // level of trust than is assigned to notaries in the first place. - // We should still remind everyone that notaries are evil and out to get you - // and to replace this list with servers they actually trust. - let trusted_servers = vec![ruma::owned_server_name!("matrix.org")]; + let trusted_servers = vec![ + ruma::owned_server_name!("continuwuity.org"), + ruma::owned_server_name!("continuwuity.rocks"), + ruma::owned_server_name!("matrix.org"), + ]; trusted_servers .into_iter() .map(TrustedServer::Name) diff --git a/src/service/server_keys/mod.rs b/src/service/server_keys/mod.rs index ae1352d2b..07f1dd360 100644 --- a/src/service/server_keys/mod.rs +++ b/src/service/server_keys/mod.rs @@ -5,22 +5,27 @@ mod util; mod verify; -use std::{collections::BTreeMap, sync::Arc}; +use std::{ + collections::{BTreeMap, HashMap}, + sync::Arc, + time::{Duration, Instant}, +}; use conduwuit::{ - Result, Server, + Result, Server, SyncRwLock, utils::{IterStream, ReadyExt, stream::TryIgnore}, }; use database::{Deserialized, Ignore, Interfix, Json, Map}; use futures::{Stream, StreamExt}; pub use request::in_one_week; use ruma::{ - CanonicalJsonObject, MilliSecondsSinceUnixEpoch, OwnedServerSigningKeyId, ServerName, - ServerSigningKeyId, + CanonicalJsonObject, MilliSecondsSinceUnixEpoch, OwnedServerName, OwnedServerSigningKeyId, + ServerName, ServerSigningKeyId, api::federation::discovery::{ServerSigningKeys, VerifyKey}, room_version_rules::RoomVersionRules, signatures::{Ed25519KeyPair, PublicKeyMap, PublicKeySet}, }; +pub use verify::strip_extraneous_signatures; use crate::{Dep, globals, sending, server_keys::util::required_keys}; @@ -29,6 +34,7 @@ pub struct Service { verify_keys: VerifyKeys, services: Services, db: Data, + last_lookup: SyncRwLock>, } struct Services { @@ -61,6 +67,7 @@ fn build(args: crate::Args<'_>) -> Result> { db: Data { servernamekeyid_response: args.db["servernamekeyid_response"].clone(), }, + last_lookup: SyncRwLock::new(HashMap::new()), })) } @@ -183,6 +190,8 @@ pub async fn verify_keys_for(&self, origin: &ServerName) -> VerifyKeys { } /// Returns the stored server signing keys responses for the origin. + /// + /// Does not imply that the stored responses are still in-date. pub fn signing_keys_for<'a>( &'a self, origin: &'a ServerName, @@ -193,4 +202,15 @@ pub fn signing_keys_for<'a>( .ignore_err() .map(|(_, v): (Ignore, ServerSigningKeys)| v) } + + /// Determines if the server may contact the origin server to fetch keys + /// when acting as a notary server. This limits origin lookups to once per + /// minute, which prevents amplification attacks. + #[must_use] + pub fn notary_may_contact_origin(&self, server_name: &ServerName) -> bool { + self.last_lookup + .read() + .get(server_name) + .is_none_or(|last| last.elapsed() >= Duration::from_mins(1)) + } } diff --git a/src/service/server_keys/request.rs b/src/service/server_keys/request.rs index 5a4ce6a10..bda5540e3 100644 --- a/src/service/server_keys/request.rs +++ b/src/service/server_keys/request.rs @@ -1,4 +1,4 @@ -use std::collections::BTreeMap; +use std::{collections::BTreeMap, time::Instant}; use conduwuit::{Err, Result, trace, utils::millis_since_unix_epoch, warn}; use ruma::{ @@ -30,6 +30,16 @@ pub async fn origin_request( ) -> Result { use get_server_keys::v2::Request; + // N.B. The "last lookup" is written before the request is actually made + // to prevent concurrent notary requests from spawning... concurrent + // origin requests, especially if the origin is slow. + // This has the downside that, if the origin is temporarily unreachable + // (including if we're backing off from it), the notary might take an + // additional minute to recover compared to the rest of the server. This + // is deemed acceptable. + self.last_lookup + .write() + .insert(target.clone(), Instant::now()); let server_signing_key = self .services .sending diff --git a/src/service/server_keys/verify.rs b/src/service/server_keys/verify.rs index a9c0c24e4..74e2a1a85 100644 --- a/src/service/server_keys/verify.rs +++ b/src/service/server_keys/verify.rs @@ -78,7 +78,7 @@ pub fn verify_server_keys_response( pubkey_map.insert(server_keys.server_name.to_string(), set); let mut canonical = to_canonical_object(server_keys)?; - Self::strip_extraneous_signatures(&mut canonical, &server_keys.server_name, &[]); + strip_extraneous_signatures(&mut canonical, &server_keys.server_name, &[]); ruma::signatures::verify_json(&pubkey_map, &canonical).map_err(Into::into) } @@ -100,7 +100,7 @@ fn verify_notary_signature( ))); }; let mut canonical_object = to_canonical_object(notary_signatures)?; - Self::strip_extraneous_signatures(&mut canonical_object, &server_keys.server_name, &[ + strip_extraneous_signatures(&mut canonical_object, &server_keys.server_name, &[ notary_name, ]); let for_verify = ruma::signatures::to_canonical_json_string_for_signing( @@ -137,22 +137,22 @@ fn verify_notary_signature( "No valid signature from {notary_name} present on signing keys response" ))) } +} - fn strip_extraneous_signatures( - canonical: &mut CanonicalJsonObject, - origin: &ServerName, - notaries: &[&ServerName], - ) { - canonical.entry("signatures".to_owned()).and_modify(|sigs| { - sigs.as_object_mut().map(|s| { - s.retain(|server_name, _| { - let Ok(server_name) = ServerName::parse(server_name) else { return false }; - origin == server_name || notaries.iter().any(|ns| *ns == server_name) - }); - Some(s) +pub fn strip_extraneous_signatures( + canonical: &mut CanonicalJsonObject, + origin: &ServerName, + notaries: &[&ServerName], +) { + canonical.entry("signatures".to_owned()).and_modify(|sigs| { + sigs.as_object_mut().map(|s| { + s.retain(|server_name, _| { + let Ok(server_name) = ServerName::parse(server_name) else { return false }; + origin == server_name || notaries.iter().any(|ns| *ns == server_name) }); + Some(s) }); - } + }); } #[cfg(test)]