Compare commits

...

29 Commits

Author SHA1 Message Date
strawberry b3986b3350 sdkfjsadklfsdklfdjsakfljas
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-21 22:05:56 -04:00
strawberry d5db11eb45 significantly drop URL preview timeouts
theres no reason for us to spend so long trying to get
a preview

Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 22:18:23 -04:00
strawberry ba80fbd2a4 raise connection pooling idle timeout to 50 seconds
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 22:17:17 -04:00
strawberry 3a1941a972 raise get_keys_helper timeout even more
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 22:16:39 -04:00
strawberry e6ab3ac2ad update book.toml for conduwuit
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 21:33:08 -04:00
strawberry fd428e9512 slight request logging improvements
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 21:20:04 -04:00
strawberry 32188ba1f9 auto join rooms from admin room created users too
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 21:16:03 -04:00
strawberry a51bc163f5 fix wrong error message about presence
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 18:28:34 -04:00
strawberry f27d98cfb5 skip rooms we have not joined before for auto-join
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 18:09:07 -04:00
strawberry d62c01b3c0 default to None if "name" in m.room.name is empty
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 17:43:48 -04:00
strawberry debc8b6164 simplify heroes get_avatar
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 17:41:16 -04:00
strawberry fee442c5c5 feat: automatically join rooms on registration
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 12:11:25 -04:00
strawberry 924adfb4e0 use unwrap_or_default if timestamp conversion fails
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 00:53:39 -04:00
strawberry f25764a158 check+clarify online backups are RocksDB only
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 00:50:22 -04:00
strawberry 59cdec4932 return helpful message instead of empty message if no backups
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 00:48:40 -04:00
strawberry b036a4fa75 make database_backup_path a PathBuf
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 00:44:02 -04:00
strawberry 00baef9c00 make database_path a PathBuf
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 00:27:32 -04:00
renovate[bot] 0e95b2d3cb chore(deps): update docker docker tag to v25.0.5 2024-03-20 00:27:32 -04:00
strawberry ea3834b19b fix lints
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 00:03:07 -04:00
Jason Volk b9d7185290 add database backup with admin commands
Signed-off-by: Jason Volk <jason@zemos.net>
2024-03-20 00:00:50 -04:00
strawberry e5f00926d7 db_cache_capacity_mb defaults to 256.0 now
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 00:00:43 -04:00
Jason Volk ede82e7b90 reconfigure and optimize rocksdb options.
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-20 00:00:43 -04:00
Jason Volk 8d1b597b6b add sync() to db abstraction for fsync(2). 2024-03-20 00:00:43 -04:00
Jason Volk f43b8d449b add rocksdb env to options. keep options in engine state.
Signed-off-by: Jason Volk <jason@zemos.net>
2024-03-19 19:14:00 -04:00
Jason Volk 15c3e03908 add abstract fallbacks for kv batch methods.
Signed-off-by: Jason Volk <jason@zemos.net>
2024-03-19 19:13:46 -04:00
strawberry fdc6e05443 bump rocksdb, deps, switch to hickory dns/resolver
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-19 19:12:49 -04:00
renovate[bot] 05002267a7 fix(deps): update rust crate serde_yaml to 0.9.33 2024-03-19 19:01:27 -04:00
strawberry 2aec32c007 fix docs
Signed-off-by: strawberry <strawberry@puppygock.gay>
2024-03-19 00:51:24 -04:00
Jason Volk 148628dbc8 fix zealous client connection close (regression 809c9b4481)
Signed-off-by: Jason Volk <jason@zemos.net>
2024-03-19 00:50:12 -04:00
24 changed files with 583 additions and 173 deletions
+2 -2
View File
@@ -129,9 +129,9 @@ artifacts:
.push-oci-image: .push-oci-image:
stage: publish stage: publish
image: docker:25.0.4 image: docker:25.0.5
services: services:
- docker:25.0.4-dind - docker:25.0.5-dind
variables: variables:
IMAGE_SUFFIX_AMD64: amd64 IMAGE_SUFFIX_AMD64: amd64
IMAGE_SUFFIX_ARM64V8: arm64v8 IMAGE_SUFFIX_ARM64V8: arm64v8
Generated
+63 -54
View File
@@ -355,6 +355,15 @@ version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fd16c4719339c4530435d38e511904438d07cce7950afa3718a84ac36c10e89e" checksum = "fd16c4719339c4530435d38e511904438d07cce7950afa3718a84ac36c10e89e"
[[package]]
name = "chrono"
version = "0.4.35"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8eaf5903dcbc0a39312feb77df2ff4c76387d591b9fc7b04a238dcf8bb62639a"
dependencies = [
"num-traits",
]
[[package]] [[package]]
name = "clang-sys" name = "clang-sys"
version = "1.7.0" version = "1.7.0"
@@ -421,11 +430,13 @@ dependencies = [
"axum-server-dual-protocol", "axum-server-dual-protocol",
"base64 0.22.0", "base64 0.22.0",
"bytes", "bytes",
"chrono",
"clap", "clap",
"cyborgtime", "cyborgtime",
"either", "either",
"figment", "figment",
"futures-util", "futures-util",
"hickory-resolver",
"hmac", "hmac",
"http", "http",
"hyper", "hyper",
@@ -467,7 +478,6 @@ dependencies = [
"tracing-flame", "tracing-flame",
"tracing-opentelemetry", "tracing-opentelemetry",
"tracing-subscriber", "tracing-subscriber",
"trust-dns-resolver",
"webpage", "webpage",
] ]
@@ -950,6 +960,51 @@ version = "0.4.3"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70"
[[package]]
name = "hickory-proto"
version = "0.24.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "091a6fbccf4860009355e3efc52ff4acf37a63489aad7435372d44ceeb6fbbcf"
dependencies = [
"async-trait",
"cfg-if",
"data-encoding",
"enum-as-inner",
"futures-channel",
"futures-io",
"futures-util",
"idna 0.4.0",
"ipnet",
"once_cell",
"rand",
"thiserror",
"tinyvec",
"tokio",
"tracing",
"url",
]
[[package]]
name = "hickory-resolver"
version = "0.24.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "35b8f021164e6a984c9030023544c57789c51760065cd510572fedcfb04164e8"
dependencies = [
"cfg-if",
"futures-util",
"hickory-proto",
"ipconfig",
"lru-cache",
"once_cell",
"parking_lot",
"rand",
"resolv-conf",
"smallvec",
"thiserror",
"tokio",
"tracing",
]
[[package]] [[package]]
name = "hmac" name = "hmac"
version = "0.12.1" version = "0.12.1"
@@ -2029,9 +2084,9 @@ checksum = "c08c74e62047bb2de4ff487b251e4a92e24f48745648451635cec7d591162d9f"
[[package]] [[package]]
name = "reqwest" name = "reqwest"
version = "0.11.26" version = "0.11.27"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "78bf93c4af7a8bb7d879d51cebe797356ff10ae8516ace542b5182d9dcac10b2" checksum = "dd67538700a17451e7cba03ac727fb961abb7607553461627b97de0b89cf4a62"
dependencies = [ dependencies = [
"base64 0.21.7", "base64 0.21.7",
"bytes", "bytes",
@@ -2039,6 +2094,7 @@ dependencies = [
"futures-core", "futures-core",
"futures-util", "futures-util",
"h2", "h2",
"hickory-resolver",
"http", "http",
"http-body", "http-body",
"hyper", "hyper",
@@ -2062,7 +2118,6 @@ dependencies = [
"tokio-rustls", "tokio-rustls",
"tokio-socks", "tokio-socks",
"tower-service", "tower-service",
"trust-dns-resolver",
"url", "url",
"wasm-bindgen", "wasm-bindgen",
"wasm-bindgen-futures", "wasm-bindgen-futures",
@@ -2300,8 +2355,8 @@ dependencies = [
[[package]] [[package]]
name = "rust-librocksdb-sys" name = "rust-librocksdb-sys"
version = "0.18.1+8.11.3" version = "0.19.0+9.0.0"
source = "git+https://github.com/zaidoon1/rust-rocksdb?rev=3e4a0f632a8c0c2839c7d183725c53895110d907#3e4a0f632a8c0c2839c7d183725c53895110d907" source = "git+https://github.com/zaidoon1/rust-rocksdb?rev=4b8491b9db45066115435881184902bbb23d3811#4b8491b9db45066115435881184902bbb23d3811"
dependencies = [ dependencies = [
"bindgen", "bindgen",
"bzip2-sys", "bzip2-sys",
@@ -2316,8 +2371,8 @@ dependencies = [
[[package]] [[package]]
name = "rust-rocksdb" name = "rust-rocksdb"
version = "0.22.7" version = "0.22.8"
source = "git+https://github.com/zaidoon1/rust-rocksdb?rev=3e4a0f632a8c0c2839c7d183725c53895110d907#3e4a0f632a8c0c2839c7d183725c53895110d907" source = "git+https://github.com/zaidoon1/rust-rocksdb?rev=4b8491b9db45066115435881184902bbb23d3811#4b8491b9db45066115435881184902bbb23d3811"
dependencies = [ dependencies = [
"libc", "libc",
"rust-librocksdb-sys", "rust-librocksdb-sys",
@@ -3169,52 +3224,6 @@ dependencies = [
"tracing-log", "tracing-log",
] ]
[[package]]
name = "trust-dns-proto"
version = "0.23.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3119112651c157f4488931a01e586aa459736e9d6046d3bd9105ffb69352d374"
dependencies = [
"async-trait",
"cfg-if",
"data-encoding",
"enum-as-inner",
"futures-channel",
"futures-io",
"futures-util",
"idna 0.4.0",
"ipnet",
"once_cell",
"rand",
"smallvec",
"thiserror",
"tinyvec",
"tokio",
"tracing",
"url",
]
[[package]]
name = "trust-dns-resolver"
version = "0.23.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "10a3e6c3aff1718b3c73e395d1f35202ba2ffa847c6a62eea0db8fb4cfe30be6"
dependencies = [
"cfg-if",
"futures-util",
"ipconfig",
"lru-cache",
"once_cell",
"parking_lot",
"rand",
"resolv-conf",
"smallvec",
"thiserror",
"tokio",
"tracing",
"trust-dns-proto",
]
[[package]] [[package]]
name = "try-lock" name = "try-lock"
version = "0.2.5" version = "0.2.5"
+10 -4
View File
@@ -27,7 +27,7 @@ base64 = "0.22.0"
ring = "0.17.8" ring = "0.17.8"
# Used when querying the SRV record of other servers # Used when querying the SRV record of other servers
trust-dns-resolver = "0.23.2" hickory-resolver = "0.24.0"
# Used to find matching events for appservices # Used to find matching events for appservices
regex = "1.10.3" regex = "1.10.3"
@@ -67,6 +67,12 @@ cyborgtime = "2.1.1"
bytes = "1.5.0" bytes = "1.5.0"
http = "0.2.12" http = "0.2.12"
# standard date and time tools
[dependencies.chrono]
version = "0.4.35"
features = ["alloc"]
default-features = false
# Web framework # Web framework
[dependencies.axum] [dependencies.axum]
version = "0.6.20" version = "0.6.20"
@@ -101,7 +107,7 @@ features = [
] ]
[dependencies.reqwest] [dependencies.reqwest]
version = "0.11.26" version = "0.11.27"
default-features = false default-features = false
features = [ features = [
"rustls-tls-native-roots", "rustls-tls-native-roots",
@@ -116,7 +122,7 @@ version = "1.0.197"
features = ["rc"] features = ["rc"]
# Used for appservice registration files # Used for appservice registration files
[dependencies.serde_yaml] [dependencies.serde_yaml]
version = "0.9.32" version = "0.9.33"
# Used for ruma wrapper # Used for ruma wrapper
[dependencies.serde_json] [dependencies.serde_json]
version = "1.0.114" version = "1.0.114"
@@ -255,7 +261,7 @@ features = [
[dependencies.rust-rocksdb] [dependencies.rust-rocksdb]
git = "https://github.com/zaidoon1/rust-rocksdb" git = "https://github.com/zaidoon1/rust-rocksdb"
#branch = "master" #branch = "master"
rev = "3e4a0f632a8c0c2839c7d183725c53895110d907" rev = "4b8491b9db45066115435881184902bbb23d3811"
optional = true optional = true
default-features = true default-features = true
features = [ features = [
+1 -5
View File
@@ -35,11 +35,7 @@ from time to time.
There are still a few nice to have features missing that some users may notice: There are still a few nice to have features missing that some users may notice:
- Outgoing read receipts and typing indicators (receiving works) - Outgoing typing indicators (receiving works)
#### What's different about your fork than upstream Conduit?
See [docs/differences.md](docs/differences.md)
#### Why does this fork exist? Why don't you contribute back upstream? #### Why does this fork exist? Why don't you contribute back upstream?
+5 -5
View File
@@ -1,6 +1,6 @@
[book] [book]
title = "Conduit" title = "conduwuit"
description = "Conduit is a simple, fast and reliable chat server for the Matrix protocol" description = "conduwuit, which is a fork of Conduit, is a simple, fast and reliable chat server for the Matrix protocol"
language = "en" language = "en"
multilingual = false multilingual = false
src = "docs" src = "docs"
@@ -10,9 +10,9 @@ build-dir = "public"
create-missing = true create-missing = true
[output.html] [output.html]
git-repository-url = "https://gitlab.com/famedly/conduit" git-repository-url = "https://github.com/girlbossceo/conduwuit"
edit-url-template = "https://gitlab.com/famedly/conduit/-/edit/next/{path}" edit-url-template = "https://github.com/girlbossceo/conduwuit/edit/main/{path}"
git-repository-icon = "fa-git-square" git-repository-icon = "fa-github-square"
[output.html.search] [output.html.search]
limit-results = 15 limit-results = 15
+9 -3
View File
@@ -263,6 +263,12 @@ url_preview_check_root_domain = false
# Defaults to true as this is the fastest option for federation. # Defaults to true as this is the fastest option for federation.
#query_trusted_key_servers_first = true #query_trusted_key_servers_first = true
# List/vector of room **IDs** that conduwuit will make newly registered users join.
# The room IDs specified must be rooms that you have joined at least once on the server, and must be public.
#
# No default.
#auto_join_rooms = []
### Generic database options ### Generic database options
@@ -274,8 +280,8 @@ url_preview_check_root_domain = false
# Set this to any float value in megabytes for conduwuit to tell the database engine that this much memory is available for database-related caches. # Set this to any float value in megabytes for conduwuit to tell the database engine that this much memory is available for database-related caches.
# May be useful if you have significant memory to spare to increase performance. # May be useful if you have significant memory to spare to increase performance.
# Defaults to 300.0 # Defaults to 256.0
#db_cache_capacity_mb = 300.0 #db_cache_capacity_mb = 256.0
# Interval in seconds when conduwuit will run database cleanup operations. # Interval in seconds when conduwuit will run database cleanup operations.
# #
@@ -314,7 +320,7 @@ url_preview_check_root_domain = false
# Amount of threads that RocksDB will use for parallelism. Set to 0 to use all your physical cores. # Amount of threads that RocksDB will use for parallelism. Set to 0 to use all your physical cores.
# Conduit eagerly spawns threads mainly for federation, so it may not be desirable to use all your cores / logical threads. # Conduit eagerly spawns threads mainly for federation, so it may not be desirable to use all your cores / logical threads.
# #
# Defaults to your CPU physical core count (not logical threads) count divided by 2 (half) # Defaults to your CPU physical core count (not logical threads).
#rocksdb_parallelism_threads = 0 #rocksdb_parallelism_threads = 0
# Maximum number of LOG files RocksDB will keep. This must *not* be set to 0. It must be at least 1. # Maximum number of LOG files RocksDB will keep. This must *not* be set to 0. It must be at least 1.
+4
View File
@@ -4,6 +4,10 @@
{{#include ../README.md:body}} {{#include ../README.md:body}}
#### What's different about your fork than upstream Conduit?
See [differences.md](differences.md)
#### How can I deploy my own? #### How can I deploy my own?
- [Deployment options](deploying.md) - [Deployment options](deploying.md)
+34 -2
View File
@@ -12,10 +12,13 @@ use ruma::{
events::{room::message::RoomMessageEventContent, GlobalAccountDataEventType}, events::{room::message::RoomMessageEventContent, GlobalAccountDataEventType},
push, UserId, push, UserId,
}; };
use tracing::{info, warn}; use tracing::{error, info, warn};
use super::{DEVICE_ID_LENGTH, SESSION_ID_LENGTH, TOKEN_LENGTH}; use super::{DEVICE_ID_LENGTH, SESSION_ID_LENGTH, TOKEN_LENGTH};
use crate::{api::client_server, services, utils, Error, Result, Ruma}; use crate::{
api::client_server::{self, join_room_by_id_helper},
services, utils, Error, Result, Ruma,
};
const RANDOM_USER_ID_LENGTH: usize = 10; const RANDOM_USER_ID_LENGTH: usize = 10;
@@ -285,6 +288,35 @@ pub async fn register_route(body: Ruma<register::v3::Request>) -> Result<registe
} }
} }
if !services().globals.config.auto_join_rooms.is_empty() {
for room in &services().globals.config.auto_join_rooms {
if !services().rooms.state_cache.server_in_room(services().globals.server_name(), room)? {
warn!("Skipping room {room} to automatically join as we have never joined before.");
continue;
}
if let Some(room_id_server_name) = room.server_name() {
match join_room_by_id_helper(
Some(&user_id),
room,
Some("Automatically joining this room upon registration".to_owned()),
&[room_id_server_name.to_owned(), services().globals.server_name().to_owned()],
None,
)
.await
{
Ok(_) => {
info!("Automatically joined room {room} for user {user_id}");
},
Err(e) => {
// don't return this error so we don't fail registrations
error!("Failed to automatically join room {room} for user {user_id}: {e}");
},
};
}
}
}
Ok(register::v3::Response { Ok(register::v3::Response {
access_token: Some(token), access_token: Some(token),
user_id, user_id,
+2 -2
View File
@@ -313,7 +313,7 @@ pub(crate) async fn get_keys_helper<F: Fn(&UserId) -> bool>(
( (
server, server,
tokio::time::timeout( tokio::time::timeout(
Duration::from_secs(50), Duration::from_secs(90),
services().sending.send_federation_request( services().sending.send_federation_request(
server, server,
federation::keys::get_keys::v1::Request { federation::keys::get_keys::v1::Request {
@@ -323,7 +323,7 @@ pub(crate) async fn get_keys_helper<F: Fn(&UserId) -> bool>(
) )
.await .await
.map_err(|e| { .map_err(|e| {
error!("get_keys_helper query took too long: {}", e); error!("get_keys_helper query took too long: {e}");
Error::BadServerResponse("get_keys_helper query took too long") Error::BadServerResponse("get_keys_helper query took too long")
}), }),
) )
+1 -1
View File
@@ -474,7 +474,7 @@ pub async fn joined_members_route(body: Ruma<joined_members::v3::Request>) -> Re
}) })
} }
async fn join_room_by_id_helper( pub(crate) async fn join_room_by_id_helper(
sender_user: Option<&UserId>, room_id: &RoomId, reason: Option<String>, servers: &[OwnedServerName], sender_user: Option<&UserId>, room_id: &RoomId, reason: Option<String>, servers: &[OwnedServerName],
_third_party_signed: Option<&ThirdPartySigned>, _third_party_signed: Option<&ThirdPartySigned>,
) -> Result<join_room_by_id::v3::Response> { ) -> Result<join_room_by_id::v3::Response> {
+1 -3
View File
@@ -1333,9 +1333,7 @@ pub async fn sync_events_v4_route(
ruma::JsOption::Some(heroes_avatar) ruma::JsOption::Some(heroes_avatar)
} else { } else {
match services().rooms.state_accessor.get_avatar(room_id)? { match services().rooms.state_accessor.get_avatar(room_id)? {
ruma::JsOption::Some(avatar) => { ruma::JsOption::Some(avatar) => ruma::JsOption::from_option(avatar.url),
avatar.url.map_or(ruma::JsOption::Undefined, ruma::JsOption::Some)
},
ruma::JsOption::Null => ruma::JsOption::Null, ruma::JsOption::Null => ruma::JsOption::Null,
ruma::JsOption::Undefined => ruma::JsOption::Undefined, ruma::JsOption::Undefined => ruma::JsOption::Undefined,
} }
+1 -1
View File
@@ -13,6 +13,7 @@ use std::{
use axum::{response::IntoResponse, Json}; use axum::{response::IntoResponse, Json};
use futures_util::future::TryFutureExt; use futures_util::future::TryFutureExt;
use get_profile_information::v1::ProfileField; use get_profile_information::v1::ProfileField;
use hickory_resolver::{error::ResolveError, lookup::SrvLookup};
use http::header::{HeaderValue, AUTHORIZATION}; use http::header::{HeaderValue, AUTHORIZATION};
use ipaddress::IPAddress; use ipaddress::IPAddress;
use ruma::{ use ruma::{
@@ -53,7 +54,6 @@ use ruma::{
use serde_json::value::{to_raw_value, RawValue as RawJsonValue}; use serde_json::value::{to_raw_value, RawValue as RawJsonValue};
use tokio::sync::RwLock; use tokio::sync::RwLock;
use tracing::{debug, error, info, warn}; use tracing::{debug, error, info, warn};
use trust_dns_resolver::{error::ResolveError, lookup::SrvLookup};
use crate::{ use crate::{
api::client_server::{self, claim_keys_helper, get_keys_helper}, api::client_server::{self, claim_keys_helper, get_keys_helper},
+28 -5
View File
@@ -10,7 +10,7 @@ use either::Either;
use figment::Figment; use figment::Figment;
use itertools::Itertools; use itertools::Itertools;
use regex::RegexSet; use regex::RegexSet;
use ruma::{OwnedServerName, RoomVersionId}; use ruma::{OwnedRoomId, OwnedServerName, RoomVersionId};
use serde::{de::IgnoredAny, Deserialize}; use serde::{de::IgnoredAny, Deserialize};
use tracing::{debug, error, warn}; use tracing::{debug, error, warn};
@@ -41,7 +41,10 @@ pub struct Config {
pub server_name: OwnedServerName, pub server_name: OwnedServerName,
#[serde(default = "default_database_backend")] #[serde(default = "default_database_backend")]
pub database_backend: String, pub database_backend: String,
pub database_path: String, pub database_path: PathBuf,
pub database_backup_path: Option<PathBuf>,
#[serde(default = "default_database_backups_to_keep")]
pub database_backups_to_keep: i16,
#[serde(default = "default_db_cache_capacity_mb")] #[serde(default = "default_db_cache_capacity_mb")]
pub db_cache_capacity_mb: f64, pub db_cache_capacity_mb: f64,
#[serde(default = "default_new_user_displayname_suffix")] #[serde(default = "default_new_user_displayname_suffix")]
@@ -107,6 +110,9 @@ pub struct Config {
#[serde(default = "default_turn_ttl")] #[serde(default = "default_turn_ttl")]
pub turn_ttl: u64, pub turn_ttl: u64,
#[serde(default = "Vec::new")]
pub auto_join_rooms: Vec<OwnedRoomId>,
#[serde(default = "default_rocksdb_log_level")] #[serde(default = "default_rocksdb_log_level")]
pub rocksdb_log_level: String, pub rocksdb_log_level: String,
#[serde(default = "default_rocksdb_max_log_file_size")] #[serde(default = "default_rocksdb_max_log_file_size")]
@@ -251,7 +257,15 @@ impl fmt::Display for Config {
let lines = [ let lines = [
("Server name", self.server_name.host()), ("Server name", self.server_name.host()),
("Database backend", &self.database_backend), ("Database backend", &self.database_backend),
("Database path", &self.database_path), ("Database path", &self.database_path.to_string_lossy()),
(
"Database backup path",
match &self.database_backup_path {
Some(path) => path.to_str().unwrap(),
None => "",
},
),
("Database backups to keep", &self.database_backups_to_keep.to_string()),
("Database cache capacity (MB)", &self.db_cache_capacity_mb.to_string()), ("Database cache capacity (MB)", &self.db_cache_capacity_mb.to_string()),
("Cache capacity modifier", &self.conduit_cache_capacity_modifier.to_string()), ("Cache capacity modifier", &self.conduit_cache_capacity_modifier.to_string()),
("PDU cache capacity", &self.pdu_cache_capacity.to_string()), ("PDU cache capacity", &self.pdu_cache_capacity.to_string()),
@@ -353,6 +367,13 @@ impl fmt::Display for Config {
} }
&lst.join(", ") &lst.join(", ")
}), }),
("Auto Join Rooms", {
let mut lst = vec![];
for room in &self.auto_join_rooms {
lst.push(room);
}
&lst.into_iter().join(", ")
}),
#[cfg(feature = "compression-zstd")] #[cfg(feature = "compression-zstd")]
("zstd Response Body Compression", &self.zstd_compression.to_string()), ("zstd Response Body Compression", &self.zstd_compression.to_string()),
#[cfg(feature = "rocksdb")] #[cfg(feature = "rocksdb")]
@@ -446,9 +467,11 @@ fn default_port() -> ListeningPort {
fn default_unix_socket_perms() -> u32 { 660 } fn default_unix_socket_perms() -> u32 { 660 }
fn default_database_backups_to_keep() -> i16 { 1 }
fn default_database_backend() -> String { "rocksdb".to_owned() } fn default_database_backend() -> String { "rocksdb".to_owned() }
fn default_db_cache_capacity_mb() -> f64 { 300.0 } fn default_db_cache_capacity_mb() -> f64 { 256.0 }
fn default_conduit_cache_capacity_modifier() -> f64 { 1.0 } fn default_conduit_cache_capacity_modifier() -> f64 { 1.0 }
@@ -489,7 +512,7 @@ fn default_rocksdb_max_log_file_size() -> usize {
4 * 1024 * 1024 4 * 1024 * 1024
} }
fn default_rocksdb_parallelism_threads() -> usize { num_cpus::get_physical() / 2 } fn default_rocksdb_parallelism_threads() -> usize { 0 }
fn default_rocksdb_compression_algo() -> String { "zstd".to_owned() } fn default_rocksdb_compression_algo() -> String { "zstd".to_owned() }
+28 -5
View File
@@ -1,4 +1,4 @@
use std::{future::Future, pin::Pin, sync::Arc}; use std::{error::Error, future::Future, pin::Pin, sync::Arc};
use super::Config; use super::Config;
use crate::Result; use crate::Result;
@@ -18,6 +18,8 @@ pub(crate) trait KeyValueDatabaseEngine: Send + Sync {
Self: Sized; Self: Sized;
fn open_tree(&self, name: &'static str) -> Result<Arc<dyn KvTree>>; fn open_tree(&self, name: &'static str) -> Result<Arc<dyn KvTree>>;
fn flush(&self) -> Result<()>; fn flush(&self) -> Result<()>;
#[allow(dead_code)]
fn sync(&self) -> Result<()> { Ok(()) }
fn cleanup(&self) -> Result<()> { Ok(()) } fn cleanup(&self) -> Result<()> { Ok(()) }
fn memory_usage(&self) -> Result<String> { fn memory_usage(&self) -> Result<String> {
Ok("Current database engine does not support memory usage reporting.".to_owned()) Ok("Current database engine does not support memory usage reporting.".to_owned())
@@ -25,6 +27,10 @@ pub(crate) trait KeyValueDatabaseEngine: Send + Sync {
#[allow(dead_code)] #[allow(dead_code)]
fn clear_caches(&self) {} fn clear_caches(&self) {}
fn backup(&self) -> Result<(), Box<dyn Error>> { unimplemented!() }
fn backup_list(&self) -> Result<String> { Ok(String::new()) }
} }
pub(crate) trait KvTree: Send + Sync { pub(crate) trait KvTree: Send + Sync {
@@ -39,20 +45,37 @@ pub(crate) trait KvTree: Send + Sync {
} }
fn insert(&self, key: &[u8], value: &[u8]) -> Result<()>; fn insert(&self, key: &[u8], value: &[u8]) -> Result<()>;
fn insert_batch(&self, iter: &mut dyn Iterator<Item = (Vec<u8>, Vec<u8>)>) -> Result<()>; fn insert_batch(&self, iter: &mut dyn Iterator<Item = (Vec<u8>, Vec<u8>)>) -> Result<()> {
for (key, value) in iter {
self.insert(&key, &value)?;
}
Ok(())
}
fn remove(&self, key: &[u8]) -> Result<()>; fn remove(&self, key: &[u8]) -> Result<()>;
#[allow(dead_code)] #[allow(dead_code)]
#[cfg(feature = "rocksdb")] fn remove_batch(&self, iter: &mut dyn Iterator<Item = Vec<u8>>) -> Result<()> {
fn remove_batch(&self, _iter: &mut dyn Iterator<Item = Vec<u8>>) -> Result<()> { unimplemented!() } for key in iter {
self.remove(&key)?;
}
Ok(())
}
fn iter<'a>(&'a self) -> Box<dyn Iterator<Item = (Vec<u8>, Vec<u8>)> + 'a>; fn iter<'a>(&'a self) -> Box<dyn Iterator<Item = (Vec<u8>, Vec<u8>)> + 'a>;
fn iter_from<'a>(&'a self, from: &[u8], backwards: bool) -> Box<dyn Iterator<Item = (Vec<u8>, Vec<u8>)> + 'a>; fn iter_from<'a>(&'a self, from: &[u8], backwards: bool) -> Box<dyn Iterator<Item = (Vec<u8>, Vec<u8>)> + 'a>;
fn increment(&self, key: &[u8]) -> Result<Vec<u8>>; fn increment(&self, key: &[u8]) -> Result<Vec<u8>>;
fn increment_batch(&self, iter: &mut dyn Iterator<Item = Vec<u8>>) -> Result<()>; fn increment_batch(&self, iter: &mut dyn Iterator<Item = Vec<u8>>) -> Result<()> {
for key in iter {
self.increment(&key)?;
}
Ok(())
}
fn scan_prefix<'a>(&'a self, prefix: Vec<u8>) -> Box<dyn Iterator<Item = (Vec<u8>, Vec<u8>)> + 'a>; fn scan_prefix<'a>(&'a self, prefix: Vec<u8>) -> Box<dyn Iterator<Item = (Vec<u8>, Vec<u8>)> + 'a>;
+152 -60
View File
@@ -4,19 +4,24 @@ use std::{
sync::{Arc, RwLock}, sync::{Arc, RwLock},
}; };
use chrono::{DateTime, Utc};
use rust_rocksdb::{ use rust_rocksdb::{
backup::{BackupEngine, BackupEngineOptions},
LogLevel::{Debug, Error, Fatal, Info, Warn}, LogLevel::{Debug, Error, Fatal, Info, Warn},
WriteBatchWithTransaction, WriteBatchWithTransaction,
}; };
use tracing::{debug, info}; use tracing::{debug, error, info};
use super::{super::Config, watchers::Watchers, KeyValueDatabaseEngine, KvTree}; use super::{super::Config, watchers::Watchers, KeyValueDatabaseEngine, KvTree};
use crate::{utils, Result}; use crate::{utils, Result};
pub(crate) struct Engine { pub(crate) struct Engine {
rocks: rust_rocksdb::DBWithThreadMode<rust_rocksdb::MultiThreaded>, rocks: rust_rocksdb::DBWithThreadMode<rust_rocksdb::MultiThreaded>,
cache: rust_rocksdb::Cache, row_cache: rust_rocksdb::Cache,
col_cache: rust_rocksdb::Cache,
old_cfs: Vec<String>, old_cfs: Vec<String>,
opts: rust_rocksdb::Options,
env: rust_rocksdb::Env,
config: Config, config: Config,
} }
@@ -27,24 +32,13 @@ struct RocksDbEngineTree<'a> {
write_lock: RwLock<()>, write_lock: RwLock<()>,
} }
fn db_options(rocksdb_cache: &rust_rocksdb::Cache, config: &Config) -> rust_rocksdb::Options { fn db_options(
// block-based options: https://docs.rs/rocksdb/latest/rocksdb/struct.BlockBasedOptions.html# config: &Config, env: &rust_rocksdb::Env, row_cache: &rust_rocksdb::Cache, col_cache: &rust_rocksdb::Cache,
let mut block_based_options = rust_rocksdb::BlockBasedOptions::default(); ) -> rust_rocksdb::Options {
block_based_options.set_block_cache(rocksdb_cache);
// "Difference of spinning disk"
// https://zhangyuchi.gitbooks.io/rocksdbbook/content/RocksDB-Tuning-Guide.html
block_based_options.set_block_size(64 * 1024);
block_based_options.set_cache_index_and_filter_blocks(true);
block_based_options.set_bloom_filter(10.0, false);
block_based_options.set_pin_l0_filter_and_index_blocks_in_cache(true);
block_based_options.set_optimize_filters_for_memory(true);
// database options: https://docs.rs/rocksdb/latest/rocksdb/struct.Options.html# // database options: https://docs.rs/rocksdb/latest/rocksdb/struct.Options.html#
let mut db_opts = rust_rocksdb::Options::default(); let mut db_opts = rust_rocksdb::Options::default();
// Logging
let rocksdb_log_level = match config.rocksdb_log_level.as_ref() { let rocksdb_log_level = match config.rocksdb_log_level.as_ref() {
"debug" => Debug, "debug" => Debug,
"info" => Info, "info" => Info,
@@ -52,7 +46,52 @@ fn db_options(rocksdb_cache: &rust_rocksdb::Cache, config: &Config) -> rust_rock
"fatal" => Fatal, "fatal" => Fatal,
_ => Error, _ => Error,
}; };
db_opts.set_log_level(rocksdb_log_level);
db_opts.set_max_log_file_size(config.rocksdb_max_log_file_size);
db_opts.set_log_file_time_to_roll(config.rocksdb_log_time_to_roll);
db_opts.set_keep_log_file_num(config.rocksdb_max_log_files);
// Processing
let threads = if config.rocksdb_parallelism_threads == 0 {
num_cpus::get_physical() // max cores if user specified 0
} else {
config.rocksdb_parallelism_threads
};
db_opts.set_max_background_jobs(threads.try_into().unwrap());
db_opts.set_max_subcompactions(threads.try_into().unwrap());
// IO
db_opts.set_use_direct_reads(true);
db_opts.set_use_direct_io_for_flush_and_compaction(true);
if config.rocksdb_optimize_for_spinning_disks {
db_opts.set_skip_stats_update_on_db_open(true); // speeds up opening DB on hard
// drives
}
// Blocks
let mut block_based_options = rust_rocksdb::BlockBasedOptions::default();
block_based_options.set_block_size(4 * 1024);
block_based_options.set_metadata_block_size(4 * 1024);
block_based_options.set_bloom_filter(9.6, true);
block_based_options.set_optimize_filters_for_memory(true);
block_based_options.set_cache_index_and_filter_blocks(true);
block_based_options.set_pin_top_level_index_and_filter(true);
block_based_options.set_block_cache(col_cache);
db_opts.set_row_cache(row_cache);
// Buffers
db_opts.set_write_buffer_size(2 * 1024 * 1024);
db_opts.set_max_write_buffer_number(2);
db_opts.set_min_write_buffer_number(1);
// Files
db_opts.set_level_zero_file_num_compaction_trigger(1);
db_opts.set_target_file_size_base(64 * 1024 * 1024);
db_opts.set_max_bytes_for_level_base(128 * 1024 * 1024);
db_opts.set_ttl(14 * 24 * 60 * 60);
// Compression
let rocksdb_compression_algo = match config.rocksdb_compression_algo.as_ref() { let rocksdb_compression_algo = match config.rocksdb_compression_algo.as_ref() {
"zstd" => rust_rocksdb::DBCompressionType::Zstd, "zstd" => rust_rocksdb::DBCompressionType::Zstd,
"zlib" => rust_rocksdb::DBCompressionType::Zlib, "zlib" => rust_rocksdb::DBCompressionType::Zlib,
@@ -61,29 +100,6 @@ fn db_options(rocksdb_cache: &rust_rocksdb::Cache, config: &Config) -> rust_rock
_ => rust_rocksdb::DBCompressionType::Zstd, _ => rust_rocksdb::DBCompressionType::Zstd,
}; };
let threads = if config.rocksdb_parallelism_threads == 0 {
num_cpus::get_physical() // max cores if user specified 0
} else {
config.rocksdb_parallelism_threads
};
db_opts.set_log_level(rocksdb_log_level);
db_opts.set_max_log_file_size(config.rocksdb_max_log_file_size);
db_opts.set_log_file_time_to_roll(config.rocksdb_log_time_to_roll);
db_opts.set_keep_log_file_num(config.rocksdb_max_log_files);
if config.rocksdb_optimize_for_spinning_disks {
db_opts.set_skip_stats_update_on_db_open(true); // speeds up opening DB on hard drives
db_opts.set_compaction_readahead_size(4 * 1024 * 1024); // "If youre running RocksDB on spinning disks, you should set this to at least
// 2MB. That way RocksDBs compaction is doing sequential instead of random
// reads."
db_opts.set_target_file_size_base(256 * 1024 * 1024);
} else {
db_opts.set_max_bytes_for_level_base(512 * 1024 * 1024);
db_opts.set_use_direct_reads(true);
db_opts.set_use_direct_io_for_flush_and_compaction(true);
}
if config.rocksdb_bottommost_compression { if config.rocksdb_bottommost_compression {
db_opts.set_bottommost_compression_type(rocksdb_compression_algo); db_opts.set_bottommost_compression_type(rocksdb_compression_algo);
db_opts.set_bottommost_zstd_max_train_bytes(0, true); db_opts.set_bottommost_zstd_max_train_bytes(0, true);
@@ -94,18 +110,10 @@ fn db_options(rocksdb_cache: &rust_rocksdb::Cache, config: &Config) -> rust_rock
// -14 w_bits is only read by zlib. // -14 w_bits is only read by zlib.
db_opts.set_compression_options(-14, config.rocksdb_compression_level, 0, 0); db_opts.set_compression_options(-14, config.rocksdb_compression_level, 0, 0);
db_opts.set_block_based_table_factory(&block_based_options);
db_opts.create_if_missing(true);
db_opts.increase_parallelism(
threads.try_into().expect("Failed to convert \"rocksdb_parallelism_threads\" usize into i32"),
);
db_opts.set_compression_type(rocksdb_compression_algo); db_opts.set_compression_type(rocksdb_compression_algo);
// https://github.com/facebook/rocksdb/wiki/Setup-Options-and-Basic-Tuning // Misc
db_opts.set_level_compaction_dynamic_level_bytes(true); db_opts.create_if_missing(true);
db_opts.set_max_background_jobs(6);
db_opts.set_bytes_per_sync(1_048_576);
// https://github.com/facebook/rocksdb/wiki/WAL-Recovery-Modes#ktoleratecorruptedtailrecords // https://github.com/facebook/rocksdb/wiki/WAL-Recovery-Modes#ktoleratecorruptedtailrecords
// //
@@ -118,15 +126,21 @@ fn db_options(rocksdb_cache: &rust_rocksdb::Cache, config: &Config) -> rust_rock
let prefix_extractor = rust_rocksdb::SliceTransform::create_fixed_prefix(1); let prefix_extractor = rust_rocksdb::SliceTransform::create_fixed_prefix(1);
db_opts.set_prefix_extractor(prefix_extractor); db_opts.set_prefix_extractor(prefix_extractor);
db_opts.set_block_based_table_factory(&block_based_options);
db_opts.set_env(env);
db_opts db_opts
} }
impl KeyValueDatabaseEngine for Arc<Engine> { impl KeyValueDatabaseEngine for Arc<Engine> {
fn open(config: &Config) -> Result<Self> { fn open(config: &Config) -> Result<Self> {
let cache_capacity_bytes = (config.db_cache_capacity_mb * 1024.0 * 1024.0) as usize; let cache_capacity_bytes = config.db_cache_capacity_mb * 1024.0 * 1024.0;
let rocksdb_cache = rust_rocksdb::Cache::new_lru_cache(cache_capacity_bytes); let row_cache_capacity_bytes = (cache_capacity_bytes * 0.25) as usize;
let col_cache_capacity_bytes = (cache_capacity_bytes * 0.75) as usize;
let db_opts = db_options(&rocksdb_cache, config); let db_env = rust_rocksdb::Env::new()?;
let row_cache = rust_rocksdb::Cache::new_lru_cache(row_cache_capacity_bytes);
let col_cache = rust_rocksdb::Cache::new_lru_cache(col_cache_capacity_bytes);
let db_opts = db_options(config, &db_env, &row_cache, &col_cache);
debug!("Listing column families in database"); debug!("Listing column families in database");
let cfs = let cfs =
@@ -138,13 +152,16 @@ impl KeyValueDatabaseEngine for Arc<Engine> {
let db = rust_rocksdb::DBWithThreadMode::<rust_rocksdb::MultiThreaded>::open_cf_descriptors( let db = rust_rocksdb::DBWithThreadMode::<rust_rocksdb::MultiThreaded>::open_cf_descriptors(
&db_opts, &db_opts,
&config.database_path, &config.database_path,
cfs.iter().map(|name| rust_rocksdb::ColumnFamilyDescriptor::new(name, db_options(&rocksdb_cache, config))), cfs.iter().map(|name| rust_rocksdb::ColumnFamilyDescriptor::new(name, db_opts.clone())),
)?; )?;
Ok(Arc::new(Engine { Ok(Arc::new(Engine {
rocks: db, rocks: db,
cache: rocksdb_cache, row_cache,
col_cache,
old_cfs: cfs, old_cfs: cfs,
opts: db_opts,
env: db_env,
config: config.clone(), config: config.clone(),
})) }))
} }
@@ -153,7 +170,7 @@ impl KeyValueDatabaseEngine for Arc<Engine> {
if !self.old_cfs.contains(&name.to_owned()) { if !self.old_cfs.contains(&name.to_owned()) {
// Create if it didn't exist // Create if it didn't exist
debug!("Creating new column family in database: {}", name); debug!("Creating new column family in database: {}", name);
let _ = self.rocks.create_cf(name, &db_options(&self.cache, &self.config)); let _ = self.rocks.create_cf(name, &self.opts);
} }
Ok(Arc::new(RocksDbEngineTree { Ok(Arc::new(RocksDbEngineTree {
@@ -165,23 +182,35 @@ impl KeyValueDatabaseEngine for Arc<Engine> {
} }
fn flush(&self) -> Result<()> { fn flush(&self) -> Result<()> {
debug!("Running flush_wal (no sync)");
rust_rocksdb::DBCommon::flush_wal(&self.rocks, false)?; rust_rocksdb::DBCommon::flush_wal(&self.rocks, false)?;
Ok(()) Ok(())
} }
fn sync(&self) -> Result<()> {
rust_rocksdb::DBCommon::flush_wal(&self.rocks, true)?;
Ok(())
}
fn memory_usage(&self) -> Result<String> { fn memory_usage(&self) -> Result<String> {
let stats = rust_rocksdb::perf::get_memory_usage_stats(Some(&[&self.rocks]), Some(&[&self.cache]))?; let stats = rust_rocksdb::perf::get_memory_usage_stats(
Some(&[&self.rocks]),
Some(&[&self.row_cache, &self.col_cache]),
)?;
Ok(format!( Ok(format!(
"Approximate memory usage of all the mem-tables: {:.3} MB\nApproximate memory usage of un-flushed \ "Approximate memory usage of all the mem-tables: {:.3} MB\nApproximate memory usage of un-flushed \
mem-tables: {:.3} MB\nApproximate memory usage of all the table readers: {:.3} MB\nApproximate memory \ mem-tables: {:.3} MB\nApproximate memory usage of all the table readers: {:.3} MB\nApproximate memory \
usage by cache: {:.3} MB\nApproximate memory usage by cache pinned: {:.3} MB\n", usage by cache: {:.3} MB\nApproximate memory usage by row cache: {:.3} MB pinned: {:.3} MB\nApproximate \
memory usage by column cache: {:.3} MB pinned: {:.3} MB\n",
stats.mem_table_total as f64 / 1024.0 / 1024.0, stats.mem_table_total as f64 / 1024.0 / 1024.0,
stats.mem_table_unflushed as f64 / 1024.0 / 1024.0, stats.mem_table_unflushed as f64 / 1024.0 / 1024.0,
stats.mem_table_readers_total as f64 / 1024.0 / 1024.0, stats.mem_table_readers_total as f64 / 1024.0 / 1024.0,
stats.cache_total as f64 / 1024.0 / 1024.0, stats.cache_total as f64 / 1024.0 / 1024.0,
self.cache.get_pinned_usage() as f64 / 1024.0 / 1024.0, self.row_cache.get_usage() as f64 / 1024.0 / 1024.0,
self.row_cache.get_pinned_usage() as f64 / 1024.0 / 1024.0,
self.col_cache.get_usage() as f64 / 1024.0 / 1024.0,
self.col_cache.get_pinned_usage() as f64 / 1024.0 / 1024.0,
)) ))
} }
@@ -194,6 +223,69 @@ impl KeyValueDatabaseEngine for Arc<Engine> {
Ok(()) Ok(())
} }
fn backup(&self) -> Result<(), Box<dyn std::error::Error>> {
let path = self.config.database_backup_path.as_ref();
if path.is_none() || path.is_some_and(|path| path.as_os_str().is_empty()) {
return Ok(());
}
let options = BackupEngineOptions::new(path.unwrap())?;
let mut engine = BackupEngine::open(&options, &self.env)?;
let ret = if self.config.database_backups_to_keep > 0 {
match engine.create_new_backup_flush(&self.rocks, true) {
Err(e) => return Err(Box::new(e)),
Ok(_) => {
let _info = engine.get_backup_info();
let info = &_info.last().unwrap();
info!(
"Created database backup #{} using {} bytes in {} files",
info.backup_id, info.size, info.num_files,
);
Ok(())
},
}
} else {
Ok(())
};
if self.config.database_backups_to_keep >= 0 {
let keep = u32::try_from(self.config.database_backups_to_keep)?;
if let Err(e) = engine.purge_old_backups(keep.try_into()?) {
error!("Failed to purge old backup: {:?}", e.to_string())
}
}
ret
}
fn backup_list(&self) -> Result<String> {
let path = self.config.database_backup_path.as_ref();
if path.is_none() || path.is_some_and(|path| path.as_os_str().is_empty()) {
return Ok(
"Configure database_backup_path to enable backups, or the path specified is not valid".to_owned(),
);
}
let mut res = String::new();
let options = BackupEngineOptions::new(path.unwrap())?;
let engine = BackupEngine::open(&options, &self.env)?;
for info in engine.get_backup_info() {
std::fmt::write(
&mut res,
format_args!(
"#{} {}: {} bytes, {} files\n",
info.backup_id,
DateTime::<Utc>::from_timestamp(info.timestamp, 0).unwrap_or_default().to_rfc2822(),
info.size,
info.num_files,
),
)
.unwrap();
}
Ok(res)
}
// TODO: figure out if this is needed for rocksdb // TODO: figure out if this is needed for rocksdb
#[allow(dead_code)] #[allow(dead_code)]
fn clear_caches(&self) {} fn clear_caches(&self) {}
+4
View File
@@ -272,4 +272,8 @@ lasttimelinecount_cache: {lasttimelinecount_cache}\n"
self.global.insert(b"version", &new_version.to_be_bytes())?; self.global.insert(b"version", &new_version.to_be_bytes())?;
Ok(()) Ok(())
} }
fn backup(&self) -> Result<(), Box<dyn std::error::Error>> { self.db.backup() }
fn backup_list(&self) -> Result<String> { self.db.backup_list() }
} }
+26 -3
View File
@@ -1,4 +1,4 @@
use ruma::api::client::error::ErrorKind; use ruma::{api::client::error::ErrorKind, OwnedUserId};
use tracing::debug; use tracing::debug;
use crate::{ use crate::{
@@ -25,8 +25,12 @@ impl service::media::Data for KeyValueDatabase {
self.mediaid_file.insert(&key, &[])?; self.mediaid_file.insert(&key, &[])?;
if let Some(user) = sender_user { if let Some(user) = sender_user {
let key = mxc.as_bytes().to_vec(); let mut key = mxc.as_bytes().to_vec();
let user = user.as_bytes().to_vec(); key.push(0xFF);
let mut user = user.as_bytes().to_vec();
user.push(0xFF);
self.mediaid_user.insert(&key, &user)?; self.mediaid_user.insert(&key, &user)?;
} }
@@ -65,6 +69,8 @@ impl service::media::Data for KeyValueDatabase {
let mut prefix = mxc.as_bytes().to_vec(); let mut prefix = mxc.as_bytes().to_vec();
prefix.push(0xFF); prefix.push(0xFF);
debug!("AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA prefix: {:?}", prefix);
let mut keys: Vec<Vec<u8>> = vec![]; let mut keys: Vec<Vec<u8>> = vec![];
for (key, _) in self.mediaid_file.scan_prefix(prefix) { for (key, _) in self.mediaid_file.scan_prefix(prefix) {
@@ -82,6 +88,23 @@ impl service::media::Data for KeyValueDatabase {
Ok(keys) Ok(keys)
} }
fn get_all_media_keys_by_user(&self, user_id: OwnedUserId) -> Result<Vec<Vec<u8>>> {
debug!("User ID: {user_id:?}");
let mut user = user_id.as_bytes().to_vec();
user.push(0xFF);
let mut keys: Vec<Vec<u8>> = vec![];
for (key, value) in self.mediaid_user.iter() {
if value == user {
keys.push(key);
}
}
Ok(keys)
}
fn search_file_metadata( fn search_file_metadata(
&self, mxc: String, width: u32, height: u32, &self, mxc: String, width: u32, height: u32,
) -> Result<(Option<String>, Option<String>, Vec<u8>)> { ) -> Result<(Option<String>, Option<String>, Vec<u8>)> {
+13 -4
View File
@@ -280,7 +280,7 @@ async fn main() {
} }
if config.allow_outgoing_presence && !config.allow_local_presence { if config.allow_outgoing_presence && !config.allow_local_presence {
error!("Outgoing presence requires allowing local presence. Please enable \"allow_outgoing_presence\"."); error!("Outgoing presence requires allowing local presence. Please enable \"allow_local_presence\".");
return; return;
} }
@@ -535,11 +535,15 @@ async fn unrecognized_method<B: Send + 'static>(
let uri = req.uri().clone(); let uri = req.uri().clone();
let inner = next.run(req).await; let inner = next.run(req).await;
if inner.status() == StatusCode::METHOD_NOT_ALLOWED { if inner.status() == StatusCode::METHOD_NOT_ALLOWED {
warn!("Method not allowed: {method} {uri}"); if uri.path().contains("_matrix/") {
warn!("Method not allowed: {method} {uri}");
} else {
info!("Method not allowed: {method} {uri}");
}
return Ok(RumaResponse(UiaaResponse::MatrixError(RumaError { return Ok(RumaResponse(UiaaResponse::MatrixError(RumaError {
body: ErrorBody::Standard { body: ErrorBody::Standard {
kind: ErrorKind::Unrecognized, kind: ErrorKind::Unrecognized,
message: "M_UNRECOGNIZED: Unrecognized request".to_owned(), message: "M_UNRECOGNIZED: Method not allowed for endpoint".to_owned(),
}, },
status_code: StatusCode::METHOD_NOT_ALLOWED, status_code: StatusCode::METHOD_NOT_ALLOWED,
})) }))
@@ -806,7 +810,12 @@ async fn shutdown_signal(handle: ServerHandle, tx: Sender<()>) -> Result<()> {
} }
async fn not_found(uri: Uri) -> impl IntoResponse { async fn not_found(uri: Uri) -> impl IntoResponse {
warn!("Not found: {uri}"); if uri.path().contains("_matrix/") {
warn!("Not found: {uri}");
} else {
info!("Not found: {uri}");
}
Error::BadRequest(ErrorKind::Unrecognized, "Unrecognized request") Error::BadRequest(ErrorKind::Unrecognized, "Unrecognized request")
} }
+111 -1
View File
@@ -30,7 +30,9 @@ use tracing::{debug, error, info, warn};
use super::pdu::PduBuilder; use super::pdu::PduBuilder;
use crate::{ use crate::{
api::{ api::{
client_server::{get_alias_helper, leave_all_rooms, leave_room, AUTO_GEN_PASSWORD_LENGTH}, client_server::{
get_alias_helper, join_room_by_id_helper, leave_all_rooms, leave_room, AUTO_GEN_PASSWORD_LENGTH,
},
server_server::parse_incoming_pdu, server_server::parse_incoming_pdu,
}, },
services, services,
@@ -92,6 +94,12 @@ enum MediaCommand {
event_id: Option<Box<EventId>>, event_id: Option<Box<EventId>>,
}, },
/// Deletes all **uploaded** local media from the specified user.
DeleteFromUser {
/// User ID to delete all media from
user_id: Box<UserId>,
},
/// - Deletes a codeblock list of MXC URLs from our database and on the /// - Deletes a codeblock list of MXC URLs from our database and on the
/// filesystem /// filesystem
DeleteList, DeleteList,
@@ -103,6 +111,13 @@ enum MediaCommand {
/// past 5 minutes /// past 5 minutes
duration: String, duration: String,
}, },
/// - Lists all **uploaded** local media from the specified user with their
/// MXC URLs.
ListFromUser {
/// User ID to list all media from
user_id: Box<UserId>,
},
} }
#[cfg_attr(test, derive(Debug))] #[cfg_attr(test, derive(Debug))]
@@ -436,6 +451,13 @@ enum ServerCommand {
ClearServiceCaches { ClearServiceCaches {
amount: u32, amount: u32,
}, },
/// - Performs an online backup of the database (only available for RocksDB
/// at the moment)
BackupDatabase,
/// - List database backups
ListBackups,
} }
#[derive(Debug)] #[derive(Debug)]
@@ -797,6 +819,17 @@ impl Service {
)); ));
} }
}, },
MediaCommand::DeleteFromUser {
user_id,
} => {
let deleted_count = services().media.delete_from_user(user_id.into()).await?;
debug!("Deleted {deleted_count} total media files.");
return Ok(RoomMessageEventContent::text_plain(format!(
"Deleted {deleted_count} total media files.",
)));
},
MediaCommand::DeleteList => { MediaCommand::DeleteList => {
if body.len() > 2 && body[0].trim().starts_with("```") && body.last().unwrap().trim() == "```" { if body.len() > 2 && body[0].trim().starts_with("```") && body.last().unwrap().trim() == "```" {
let mxc_list = body.clone().drain(1..body.len() - 1).collect::<Vec<_>>(); let mxc_list = body.clone().drain(1..body.len() - 1).collect::<Vec<_>>();
@@ -830,6 +863,25 @@ impl Service {
deleted_count deleted_count
))); )));
}, },
MediaCommand::ListFromUser {
user_id,
} => {
let mxc_list = services().media.list_all_media_by_user(user_id.clone().into()).await?;
let output_plain = format!(
"All MXC URLs {user_id} Uploaded:\n{}",
mxc_list.iter().map(ToOwned::to_owned).collect::<Vec<_>>().join("\n")
);
let output_html = format!(
"<table><caption>All MXC URLs {user_id} \
Uploaded</caption>\n<tr><th>URL</th></tr>\n{}</table>",
mxc_list.iter().fold(String::new(), |mut output, mxc| {
writeln!(output, "<tr><td>{}</td></tr>", mxc).unwrap();
output
})
);
RoomMessageEventContent::text_html(output_plain, output_html)
},
} }
}, },
AdminCommand::Users(command) => match command { AdminCommand::Users(command) => match command {
@@ -893,6 +945,35 @@ impl Service {
.expect("to json value always works"), .expect("to json value always works"),
)?; )?;
if !services().globals.config.auto_join_rooms.is_empty() {
for room in &services().globals.config.auto_join_rooms {
if !services().rooms.state_cache.server_in_room(services().globals.server_name(), room)? {
warn!("Skipping room {room} to automatically join as we have never joined before.");
continue;
}
if let Some(room_id_server_name) = room.server_name() {
match join_room_by_id_helper(
Some(&user_id),
room,
Some("Automatically joining this room upon registration".to_owned()),
&[room_id_server_name.to_owned(), services().globals.server_name().to_owned()],
None,
)
.await
{
Ok(_) => {
info!("Automatically joined room {room} for user {user_id}");
},
Err(e) => {
// don't return this error so we don't fail registrations
error!("Failed to automatically join room {room} for user {user_id}: {e}");
},
};
}
}
}
// we dont add a device since we're not the user, just the creator // we dont add a device since we're not the user, just the creator
// Inhibit login does not work for guests // Inhibit login does not work for guests
@@ -1866,6 +1947,35 @@ impl Service {
RoomMessageEventContent::text_plain("Done.") RoomMessageEventContent::text_plain("Done.")
}, },
ServerCommand::ListBackups => {
let result = services().globals.db.backup_list()?;
if result.is_empty() {
return Ok(RoomMessageEventContent::text_plain("No backups found."));
}
RoomMessageEventContent::text_plain(result)
},
ServerCommand::BackupDatabase => {
if !cfg!(feature = "rocksdb") {
return Ok(RoomMessageEventContent::text_plain(
"Only RocksDB supports online backups in conduwuit.",
));
}
let mut result = tokio::task::spawn_blocking(move || match services().globals.db.backup() {
Ok(_) => String::new(),
Err(e) => (*e).to_string(),
})
.await
.unwrap();
if result.is_empty() {
result = services().globals.db.backup_list()?;
}
RoomMessageEventContent::text_plain(&result)
},
}, },
AdminCommand::Debug(command) => match command { AdminCommand::Debug(command) => match command {
DebugCommand::GetAuthChain { DebugCommand::GetAuthChain {
+3 -1
View File
@@ -1,4 +1,4 @@
use std::collections::BTreeMap; use std::{collections::BTreeMap, error::Error};
use async_trait::async_trait; use async_trait::async_trait;
use ruma::{ use ruma::{
@@ -32,4 +32,6 @@ pub trait Data: Send + Sync {
fn signing_keys_for(&self, origin: &ServerName) -> Result<BTreeMap<OwnedServerSigningKeyId, VerifyKey>>; fn signing_keys_for(&self, origin: &ServerName) -> Result<BTreeMap<OwnedServerSigningKeyId, VerifyKey>>;
fn database_version(&self) -> Result<u64>; fn database_version(&self) -> Result<u64>;
fn bump_database_version(&self, new_version: u64) -> Result<()>; fn bump_database_version(&self, new_version: u64) -> Result<()>;
fn backup(&self) -> Result<(), Box<dyn Error>> { unimplemented!() }
fn backup_list(&self) -> Result<String> { Ok(String::new()) }
} }
+9 -6
View File
@@ -17,6 +17,7 @@ use argon2::Argon2;
use base64::{engine::general_purpose, Engine as _}; use base64::{engine::general_purpose, Engine as _};
pub use data::Data; pub use data::Data;
use futures_util::FutureExt; use futures_util::FutureExt;
use hickory_resolver::TokioAsyncResolver;
use hyper::{ use hyper::{
client::connect::dns::{GaiResolver, Name}, client::connect::dns::{GaiResolver, Name},
service::Service as HyperService, service::Service as HyperService,
@@ -34,7 +35,6 @@ use ruma::{
}; };
use tokio::sync::{broadcast, watch::Receiver, Mutex, RwLock, Semaphore}; use tokio::sync::{broadcast, watch::Receiver, Mutex, RwLock, Semaphore};
use tracing::{error, info}; use tracing::{error, info};
use trust_dns_resolver::TokioAsyncResolver;
use crate::{api::server_server::FedDest, services, Config, Error, Result}; use crate::{api::server_server::FedDest, services, Config, Error, Result};
@@ -327,6 +327,8 @@ impl Service<'_> {
pub fn turn_secret(&self) -> &String { &self.config.turn_secret } pub fn turn_secret(&self) -> &String { &self.config.turn_secret }
pub fn auto_join_rooms(&self) -> &[OwnedRoomId] { &self.config.auto_join_rooms }
pub fn notification_push_path(&self) -> &String { &self.config.notification_push_path } pub fn notification_push_path(&self) -> &String { &self.config.notification_push_path }
pub fn emergency_password(&self) -> &Option<String> { &self.config.emergency_password } pub fn emergency_password(&self) -> &Option<String> { &self.config.emergency_password }
@@ -497,8 +499,9 @@ fn reqwest_client_builder(config: &Config) -> Result<reqwest::ClientBuilder> {
}); });
let mut reqwest_client_builder = reqwest::Client::builder() let mut reqwest_client_builder = reqwest::Client::builder()
.trust_dns(true) .hickory_dns(true)
.pool_max_idle_per_host(0) .pool_max_idle_per_host(1)
.pool_idle_timeout(Duration::from_secs(50))
.connect_timeout(Duration::from_secs(60)) .connect_timeout(Duration::from_secs(60))
.timeout(Duration::from_secs(60 * 5)) .timeout(Duration::from_secs(60 * 5))
.redirect(redirect_policy) .redirect(redirect_policy)
@@ -525,10 +528,10 @@ fn url_preview_reqwest_client_builder(config: &Config) -> Result<reqwest::Client
}); });
let mut reqwest_client_builder = reqwest::Client::builder() let mut reqwest_client_builder = reqwest::Client::builder()
.trust_dns(true) .hickory_dns(true)
.pool_max_idle_per_host(0) .pool_max_idle_per_host(0)
.connect_timeout(Duration::from_secs(60)) .connect_timeout(Duration::from_secs(20))
.timeout(Duration::from_secs(60 * 5)) .timeout(Duration::from_secs(30))
.redirect(redirect_policy) .redirect(redirect_policy)
.user_agent("Conduwuit".to_owned() + "/" + env!("CARGO_PKG_VERSION")); .user_agent("Conduwuit".to_owned() + "/" + env!("CARGO_PKG_VERSION"));
+4
View File
@@ -1,3 +1,5 @@
use ruma::OwnedUserId;
use crate::Result; use crate::Result;
pub trait Data: Send + Sync { pub trait Data: Send + Sync {
@@ -15,6 +17,8 @@ pub trait Data: Send + Sync {
fn search_mxc_metadata_prefix(&self, mxc: String) -> Result<Vec<Vec<u8>>>; fn search_mxc_metadata_prefix(&self, mxc: String) -> Result<Vec<Vec<u8>>>;
fn get_all_media_keys_by_user(&self, user_id: OwnedUserId) -> Result<Vec<Vec<u8>>>;
fn get_all_media_keys(&self) -> Result<Vec<Vec<u8>>>; fn get_all_media_keys(&self) -> Result<Vec<Vec<u8>>>;
fn remove_url_preview(&self, url: &str) -> Result<()>; fn remove_url_preview(&self, url: &str) -> Result<()>;
+69 -2
View File
@@ -12,7 +12,7 @@ use tokio::{
}; };
use tracing::{debug, error}; use tracing::{debug, error};
use crate::{services, utils, Error, Result}; use crate::{services, utils::string_from_bytes, Error, Result};
#[derive(Debug)] #[derive(Debug)]
pub struct FileMeta { pub struct FileMeta {
@@ -110,6 +110,71 @@ impl Service {
} }
} }
/// Deletes all media in the database and from the media directory by a
/// specific user
pub async fn delete_from_user(&self, user_id: OwnedUserId) -> Result<u32> {
if let Ok(user_keys) = self.db.get_all_media_keys_by_user(user_id.clone()) {
if user_keys.is_empty() {
error!("User \"{user_id}\" has not uploaded any media.");
return Err(Error::bad_database("User has not uploaded any media."));
}
let mut mxc_deletion_count = 0;
for key in user_keys {
let mxc = String::from_utf8_lossy(&key);
let keys = self.db.search_mxc_metadata_prefix(mxc.to_string())?; // the MXC alone does not determine the file path, it is only the prefix
for key in keys {
let mxc = String::from_utf8_lossy(&key);
debug!("Deleting MXC {mxc} from database");
self.delete(mxc.to_string()).await?;
mxc_deletion_count += 1;
}
}
Ok(mxc_deletion_count)
} else {
error!("Failed to find any media keys for user \"{user_id}\" in our database");
Err(Error::bad_database(
"Failed to find any media keys for the provided user in our database",
))
}
}
pub async fn list_all_media_by_user(&self, user_id: OwnedUserId) -> Result<Vec<String>> {
if let Ok(keys) = self.db.get_all_media_keys_by_user(user_id.clone()) {
if keys.is_empty() {
error!("User \"{user_id}\" has not uploaded any media.");
return Err(Error::bad_database("User has not uploaded any media."));
}
let mut mxc_list: Vec<String> = vec![];
for key in keys {
let mxc_string = string_from_bytes(&key).map_err(|e| {
error!("Failed to convert MXC key to string in database for user \"{user_id}\": {e}");
Error::bad_database("Failed to convert MXC key to string in database")
})?;
// TODO: add the file name
mxc_list.push(mxc_string);
}
Ok(mxc_list)
} else {
error!("Failed to find any media keys for user \"{user_id}\" in our database");
Err(Error::bad_database(
"Failed to find any media keys for the provided user in our database",
))
}
}
/// Uploads or replaces a file thumbnail. /// Uploads or replaces a file thumbnail.
#[allow(clippy::too_many_arguments)] #[allow(clippy::too_many_arguments)]
pub async fn upload_thumbnail( pub async fn upload_thumbnail(
@@ -200,7 +265,7 @@ impl Service {
let mxc = parts let mxc = parts
.next() .next()
.map(|bytes| { .map(|bytes| {
utils::string_from_bytes(bytes).map_err(|e| { string_from_bytes(bytes).map_err(|e| {
error!("Failed to parse MXC unicode bytes from our database: {}", e); error!("Failed to parse MXC unicode bytes from our database: {}", e);
Error::bad_database("Failed to parse MXC unicode bytes from our database") Error::bad_database("Failed to parse MXC unicode bytes from our database")
}) })
@@ -502,6 +567,8 @@ mod tests {
fn search_mxc_metadata_prefix(&self, _mxc: String) -> Result<Vec<Vec<u8>>> { todo!() } fn search_mxc_metadata_prefix(&self, _mxc: String) -> Result<Vec<Vec<u8>>> { todo!() }
fn get_all_media_keys_by_user(&self, _user_id: OwnedUserId) -> Result<Vec<Vec<u8>>> { todo!() }
fn get_all_media_keys(&self) -> Result<Vec<Vec<u8>>> { todo!() } fn get_all_media_keys(&self) -> Result<Vec<Vec<u8>>> { todo!() }
fn search_file_metadata( fn search_file_metadata(
+3 -4
View File
@@ -230,10 +230,9 @@ impl Service {
pub fn get_name(&self, room_id: &RoomId) -> Result<Option<String>> { pub fn get_name(&self, room_id: &RoomId) -> Result<Option<String>> {
services().rooms.state_accessor.room_state_get(room_id, &StateEventType::RoomName, "")?.map_or(Ok(None), |s| { services().rooms.state_accessor.room_state_get(room_id, &StateEventType::RoomName, "")?.map_or(Ok(None), |s| {
serde_json::from_str(s.content.get()).map(|c: RoomNameEventContent| Some(c.name)).map_err(|e| { Ok(serde_json::from_str(s.content.get())
error!("Invalid room name event in database for room {}. {}", room_id, e); .map(|c: RoomNameEventContent| Some(c.name))
Error::bad_database("Invalid room name event in database.") .unwrap_or_else(|_| None))
})
}) })
} }