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 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}