mirror of
https://forgejo.ellis.link/continuwuation/continuwuity/
synced 2026-10-06 02:37:25 +00:00
feat: Continuwuity notarizes the notary
This commit is contained in:
@@ -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.
|
||||
@@ -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.
|
||||
#
|
||||
|
||||
@@ -203,6 +203,8 @@ pub fn build(router: Router<State>, state: State) -> Router<State> {
|
||||
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)
|
||||
|
||||
+277
-5
@@ -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<crate::State>,
|
||||
body: Ruma<get_remote_server_keys_batch::v2::Request>,
|
||||
) -> Result<get_remote_server_keys_batch::v2::Response> {
|
||||
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<Raw<ServerSigningKeys>> {
|
||||
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::<ServerSigningKeys>::from_json)
|
||||
.map_err(Into::into)
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip(server_keys, my_name, queries, start))]
|
||||
async fn acquire_keys_as_notary(
|
||||
server_keys: Arc<server_keys::Service>,
|
||||
my_name: OwnedServerName,
|
||||
remote: OwnedServerName,
|
||||
mut queries: BTreeMap<OwnedServerSigningKeyId, QueryCriteria>,
|
||||
start: MilliSecondsSinceUnixEpoch,
|
||||
) -> Vec<Raw<ServerSigningKeys>> {
|
||||
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::<Vec<_>>();
|
||||
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::<HashSet<_>>()
|
||||
.into_iter()
|
||||
.map(|idx| results.index(idx).to_owned())
|
||||
.collect()
|
||||
}
|
||||
|
||||
pub(crate) async fn get_remote_server_keys_route(
|
||||
State(services): State<crate::State>,
|
||||
body: Ruma<get_remote_server_keys::v2::Request>,
|
||||
) -> Result<get_remote_server_keys::v2::Response> {
|
||||
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::<Vec<_>>()
|
||||
.await;
|
||||
Ok(get_remote_server_keys::v2::Response::new(response))
|
||||
}
|
||||
|
||||
+7
-11
@@ -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<TrustedServer>,
|
||||
|
||||
@@ -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<TrustedServer> {
|
||||
// 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)
|
||||
|
||||
@@ -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<HashMap<OwnedServerName, Instant>>,
|
||||
}
|
||||
|
||||
struct Services {
|
||||
@@ -61,6 +67,7 @@ fn build(args: crate::Args<'_>) -> Result<Arc<Self>> {
|
||||
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))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<ServerSigningKeys> {
|
||||
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
|
||||
|
||||
@@ -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)]
|
||||
|
||||
Reference in New Issue
Block a user