Skip to main content

tuwunel_service/
services.rs

1use std::{fmt, sync::Arc};
2
3use futures::{StreamExt, TryStreamExt};
4use tokio::sync::Mutex;
5use tuwunel_core::{
6	Result, Server, debug, debug_info, implement, info, trace, utils::stream::IterStream,
7};
8use tuwunel_database::Database;
9
10pub(crate) use crate::OnceServices;
11use crate::{
12	account_data, admin, appservice, client, config, deactivate, emergency, federation, fetcher,
13	globals, key_backups, login_ratelimit,
14	manager::Manager,
15	media, membership, oauth, presence, profile, pusher, registration_tokens, rendezvous,
16	resolver,
17	rooms::{self, retention},
18	sending, sendmail, server_keys,
19	service::{Args, Service},
20	storage, sync, tasks, threepid, transaction_ids, uiaa, users,
21};
22
23pub struct Services {
24	pub account_data: Arc<account_data::Service>,
25	pub admin: Arc<admin::Service>,
26	pub appservice: Arc<appservice::Service>,
27	pub config: Arc<config::Service>,
28	pub client: Arc<client::Service>,
29	pub emergency: Arc<emergency::Service>,
30	pub fetcher: Arc<fetcher::Service>,
31	pub globals: Arc<globals::Service>,
32	pub key_backups: Arc<key_backups::Service>,
33	pub login_ratelimit: Arc<login_ratelimit::Service>,
34	pub media: Arc<media::Service>,
35	pub presence: Arc<presence::Service>,
36	pub pusher: Arc<pusher::Service>,
37	pub resolver: Arc<resolver::Service>,
38	pub alias: Arc<rooms::alias::Service>,
39	pub auth_chain: Arc<rooms::auth_chain::Service>,
40	pub delete: Arc<rooms::delete::Service>,
41	pub directory: Arc<rooms::directory::Service>,
42	pub event_handler: Arc<rooms::event_handler::Service>,
43	pub lazy_loading: Arc<rooms::lazy_loading::Service>,
44	pub metadata: Arc<rooms::metadata::Service>,
45	pub pdu_metadata: Arc<rooms::pdu_metadata::Service>,
46	pub read_receipt: Arc<rooms::read_receipt::Service>,
47	pub search: Arc<rooms::search::Service>,
48	pub short: Arc<rooms::short::Service>,
49	pub spaces: Arc<rooms::spaces::Service>,
50	pub state: Arc<rooms::state::Service>,
51	pub state_accessor: Arc<rooms::state_accessor::Service>,
52	pub state_cache: Arc<rooms::state_cache::Service>,
53	pub state_compressor: Arc<rooms::state_compressor::Service>,
54	pub storage: Arc<storage::Service>,
55	pub threads: Arc<rooms::threads::Service>,
56	pub timeline: Arc<rooms::timeline::Service>,
57	pub typing: Arc<rooms::typing::Service>,
58	pub federation: Arc<federation::Service>,
59	pub sending: Arc<sending::Service>,
60	pub server_keys: Arc<server_keys::Service>,
61	pub sync: Arc<sync::Service>,
62	pub tasks: Arc<tasks::Service>,
63	pub transaction_ids: Arc<transaction_ids::Service>,
64	pub uiaa: Arc<uiaa::Service>,
65	pub users: Arc<users::Service>,
66	pub membership: Arc<membership::Service>,
67	pub deactivate: Arc<deactivate::Service>,
68	pub oauth: Arc<oauth::Service>,
69	pub retention: Arc<retention::Service>,
70	pub registration_tokens: Arc<registration_tokens::Service>,
71	pub rendezvous: Arc<rendezvous::Service>,
72	pub sendmail: Arc<sendmail::Service>,
73	pub threepid: Arc<threepid::Service>,
74	pub profile: Arc<profile::Service>,
75
76	manager: Mutex<Option<Arc<Manager>>>,
77	pub server: Arc<Server>,
78	pub db: Arc<Database>,
79}
80
81#[implement(Services)]
82pub async fn build(server: Arc<Server>) -> Result<Arc<Self>> {
83	let db = Database::open(&server).await?;
84	let services = Arc::new(OnceServices::default());
85	let args = Args {
86		db: &db,
87		server: &server,
88		services: &services,
89	};
90
91	let res = Arc::new(Self {
92		account_data: account_data::Service::build(&args)?,
93		admin: admin::Service::build(&args)?,
94		appservice: appservice::Service::build(&args)?,
95		resolver: resolver::Service::build(&args)?,
96		client: client::Service::build(&args)?,
97		config: config::Service::build(&args)?,
98		emergency: emergency::Service::build(&args)?,
99		fetcher: fetcher::Service::build(&args)?,
100		globals: globals::Service::build(&args)?,
101		key_backups: key_backups::Service::build(&args)?,
102		media: media::Service::build(&args)?,
103		presence: presence::Service::build(&args)?,
104		pusher: pusher::Service::build(&args)?,
105		alias: rooms::alias::Service::build(&args)?,
106		auth_chain: rooms::auth_chain::Service::build(&args)?,
107		delete: rooms::delete::Service::build(&args)?,
108		directory: rooms::directory::Service::build(&args)?,
109		event_handler: rooms::event_handler::Service::build(&args)?,
110		lazy_loading: rooms::lazy_loading::Service::build(&args)?,
111		metadata: rooms::metadata::Service::build(&args)?,
112		pdu_metadata: rooms::pdu_metadata::Service::build(&args)?,
113		read_receipt: rooms::read_receipt::Service::build(&args)?,
114		search: rooms::search::Service::build(&args)?,
115		short: rooms::short::Service::build(&args)?,
116		spaces: rooms::spaces::Service::build(&args)?,
117		state: rooms::state::Service::build(&args)?,
118		state_accessor: rooms::state_accessor::Service::build(&args)?,
119		state_cache: rooms::state_cache::Service::build(&args)?,
120		state_compressor: rooms::state_compressor::Service::build(&args)?,
121		storage: storage::Service::build(&args)?,
122		threads: rooms::threads::Service::build(&args)?,
123		timeline: rooms::timeline::Service::build(&args)?,
124		typing: rooms::typing::Service::build(&args)?,
125		federation: federation::Service::build(&args)?,
126		sending: sending::Service::build(&args)?,
127		server_keys: server_keys::Service::build(&args)?,
128		sync: sync::Service::build(&args)?,
129		tasks: tasks::Service::build(&args)?,
130		transaction_ids: transaction_ids::Service::build(&args)?,
131		login_ratelimit: login_ratelimit::Service::build(&args)?,
132		uiaa: uiaa::Service::build(&args)?,
133		users: users::Service::build(&args)?,
134		membership: membership::Service::build(&args)?,
135		deactivate: deactivate::Service::build(&args)?,
136		oauth: oauth::Service::build(&args)?,
137		retention: retention::Service::build(&args)?,
138		registration_tokens: registration_tokens::Service::build(&args)?,
139		rendezvous: rendezvous::Service::build(&args)?,
140		sendmail: sendmail::Service::build(&args)?,
141		threepid: threepid::Service::build(&args)?,
142		profile: profile::Service::build(&args)?,
143
144		manager: Mutex::new(None),
145		server,
146		db,
147	});
148
149	Ok(services.set(res))
150}
151
152#[implement(Services)]
153pub(crate) fn services(&self) -> impl Iterator<Item = Arc<dyn Service>> + Send {
154	macro_rules! cast {
155		($s:expr) => {
156			<Arc<dyn Service> as Into<_>>::into($s.clone())
157		};
158	}
159
160	[
161		cast!(self.account_data),
162		cast!(self.admin),
163		cast!(self.appservice),
164		cast!(self.resolver),
165		cast!(self.client),
166		cast!(self.config),
167		cast!(self.emergency),
168		cast!(self.fetcher),
169		cast!(self.globals),
170		cast!(self.key_backups),
171		cast!(self.media),
172		cast!(self.presence),
173		cast!(self.pusher),
174		cast!(self.alias),
175		cast!(self.auth_chain),
176		cast!(self.delete),
177		cast!(self.directory),
178		cast!(self.event_handler),
179		cast!(self.lazy_loading),
180		cast!(self.metadata),
181		cast!(self.pdu_metadata),
182		cast!(self.read_receipt),
183		cast!(self.search),
184		cast!(self.short),
185		cast!(self.spaces),
186		cast!(self.state),
187		cast!(self.state_accessor),
188		cast!(self.state_cache),
189		cast!(self.state_compressor),
190		cast!(self.storage),
191		cast!(self.threads),
192		cast!(self.timeline),
193		cast!(self.typing),
194		cast!(self.federation),
195		cast!(self.sending),
196		cast!(self.server_keys),
197		cast!(self.sync),
198		cast!(self.tasks),
199		cast!(self.transaction_ids),
200		cast!(self.uiaa),
201		cast!(self.users),
202		cast!(self.membership),
203		cast!(self.deactivate),
204		cast!(self.oauth),
205		cast!(self.retention),
206		cast!(self.registration_tokens),
207		cast!(self.rendezvous),
208		cast!(self.profile),
209	]
210	.into_iter()
211}
212
213impl fmt::Debug for Services {
214	fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
215		f.debug_struct("Services").finish()
216	}
217}
218
219#[implement(Services)]
220pub async fn start(self: &Arc<Self>) -> Result<Arc<Self>> {
221	debug_info!("Starting services...");
222
223	super::migrations::migrations(self).await?;
224	self.users.validate_server_user().await?;
225
226	self.manager
227		.lock()
228		.await
229		.insert(Manager::new(self))
230		.clone()
231		.start()
232		.await?;
233
234	debug_info!("Services startup complete.");
235
236	Ok(Arc::clone(self))
237}
238
239#[implement(Services)]
240pub async fn stop(&self) {
241	info!("Shutting down services...");
242
243	self.interrupt().await;
244	if let Some(manager) = self.manager.lock().await.as_ref() {
245		manager.stop().await;
246	}
247
248	debug_info!("Services shutdown complete.");
249}
250
251#[implement(Services)]
252pub(crate) async fn interrupt(&self) {
253	debug!("Interrupting services...");
254	for service in self.services() {
255		trace!(
256			name = ?service.name(),
257			"Interrupting Service"
258		);
259
260		service.interrupt().await;
261	}
262}
263
264#[implement(Services)]
265pub async fn poll(&self) -> Result {
266	if let Some(manager) = self.manager.lock().await.as_ref() {
267		trace!("Polling service manager...");
268		return manager.poll().await;
269	}
270
271	Ok(())
272}
273
274#[implement(Services)]
275pub async fn clear_cache(&self) {
276	// Uncorked, every per-key delete in a database-backed cache flushes the
277	// write-ahead log; the rows are reconstructible, so no fsync is owed.
278	let _cork = self.db.cork_and_flush();
279
280	self.services()
281		.stream()
282		.for_each(async |service| {
283			service.clear_cache().await;
284		})
285		.await;
286}
287
288#[implement(Services)]
289pub async fn memory_usage(&self) -> Result<String> {
290	self.services()
291		.try_stream()
292		.try_fold(String::new(), async |mut out, service| {
293			service.memory_usage(&mut out).await?;
294			Ok(out)
295		})
296		.await
297}