tuwunel_service/admin/
mod.rs1mod attach;
2pub mod console;
3pub mod context;
4pub mod create;
5mod execute;
6mod grant;
7mod notices;
8mod processor;
9mod register;
10mod respond;
11
12use std::{
13 collections::BTreeMap,
14 sync::{Arc, Mutex as StdMutex, RwLock as StdRwLock},
15 time::Instant,
16};
17
18use async_trait::async_trait;
19pub use context::Context;
20pub use create::create_admin_room;
21use futures::TryFutureExt;
22use ruma::{
23 OwnedEventId, OwnedRoomAliasId, OwnedRoomId, OwnedUserId, RoomId, RoomOrAliasId, UserId,
24};
25use tokio::sync::mpsc;
26use tuwunel_core::{
27 Err, Event, Result, debug, err, error::default_log, implement, matrix::event::MsgType,
28 utils::ReadyExt, warn,
29};
30
31use crate::rooms::state::RoomMutexGuard;
32
33pub struct Service {
34 services: Arc<crate::services::OnceServices>,
35 channel: StdRwLock<Option<mpsc::Sender<CommandInput>>>,
36 pub command: StdRwLock<Option<Arc<dyn Command>>>,
37 pub admin_alias: OwnedRoomAliasId,
38 register_nonces: StdMutex<BTreeMap<String, Instant>>,
39 #[cfg(feature = "console")]
40 pub console: Arc<console::Console>,
41}
42
43#[derive(Clone, Debug, Default)]
50pub struct CommandInput {
51 pub command: String,
53
54 pub reply_id: Option<OwnedEventId>,
56
57 pub sender: Option<OwnedUserId>,
59}
60
61#[async_trait]
63pub trait Command: Send + Sync + 'static {
64 fn clap(&self) -> clap::Command;
67
68 async fn dispatch(&self, matches: clap::ArgMatches, context: &Context<'_>) -> Result;
70}
71
72pub type ProcessorResult = Result<Option<CommandOutput>, CommandOutput>;
77
78pub enum CommandOutput {
81 Markdown(String),
82 Plain(String),
83}
84
85impl From<String> for CommandInput {
86 fn from(command: String) -> Self { Self { command, ..Default::default() } }
87}
88
89impl From<&str> for CommandInput {
90 fn from(command: &str) -> Self { command.to_owned().into() }
91}
92
93impl CommandOutput {
94 #[inline]
95 #[must_use]
96 pub fn as_str(&self) -> &str {
97 match self {
98 | Self::Markdown(text) | Self::Plain(text) => text,
99 }
100 }
101}
102
103const COMMAND_QUEUE_LIMIT: usize = 512;
105
106#[async_trait]
107impl crate::Service for Service {
108 fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
109 Ok(Arc::new(Self {
110 services: args.services.clone(),
111 channel: StdRwLock::new(None),
112 command: StdRwLock::new(None),
113 admin_alias: OwnedRoomAliasId::try_from(format!("#admins:{}", args.server.name))
114 .expect("#admins:server_name is valid alias name"),
115 register_nonces: StdMutex::new(BTreeMap::new()),
116 #[cfg(feature = "console")]
117 console: console::Console::new(args),
118 }))
119 }
120
121 async fn worker(self: Arc<Self>) -> Result {
122 let mut signals = self.services.server.signal.subscribe();
123 let (sender, mut receiver) = mpsc::channel(COMMAND_QUEUE_LIMIT);
124 _ = self
125 .channel
126 .write()
127 .expect("locked for writing")
128 .insert(sender);
129
130 self.console_auto_start().await;
131
132 loop {
133 tokio::select! {
134 command = receiver.recv() => match command {
135 Some(command) => self.handle_command(command).await,
136 None => break,
137 },
138 sig = signals.recv() => if let Ok(sig) = sig {
139 self.handle_signal(sig).await;
140 },
141 }
142 }
143
144 self.interrupt().await;
146 self.console_auto_stop().await;
147
148 Ok(())
149 }
150
151 async fn interrupt(&self) {
152 #[cfg(feature = "console")]
153 self.console.interrupt();
154
155 _ = self
156 .channel
157 .write()
158 .expect("locked for writing")
159 .take();
160 }
161
162 fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
163}
164
165impl Service {
166 pub async fn command(&self, input: CommandInput) -> Result {
172 let Some(queue) = self
173 .channel
174 .read()
175 .expect("locked for reading")
176 .clone()
177 else {
178 return Err!("Admin command queue unavailable.");
179 };
180
181 queue
182 .send(input)
183 .map_err(|e| err!("Failed to enqueue admin command: {e:?}"))
184 .await
185 }
186
187 pub async fn command_in_place(&self, input: CommandInput) -> ProcessorResult {
193 self.process_command(&input).await
194 }
195
196 pub fn complete_command(&self, command: &str) -> Option<String> {
199 self.command
200 .read()
201 .expect("locked for reading")
202 .as_ref()
203 .map(|root| processor::complete(root.clap(), command))
204 }
205
206 async fn handle_signal(&self, sig: &'static str) {
207 if sig == execute::SIGNAL {
208 self.signal_execute().await.ok();
209 }
210
211 #[cfg(feature = "console")]
212 self.console.handle_signal(sig);
213 }
214
215 async fn handle_command(&self, command: CommandInput) {
216 match self.process_command(&command).await {
217 | Ok(None) => debug!("Command successful with no response"),
218 | Err(output) | Ok(Some(output)) => self
219 .handle_response(output, command.reply_id.as_deref())
220 .await
221 .unwrap_or_else(default_log),
222 }
223 }
224
225 async fn process_command(&self, command: &CommandInput) -> ProcessorResult {
226 let root = self
227 .command
228 .read()
229 .expect("locked for reading")
230 .clone()
231 .expect("Admin module is not loaded");
232
233 processor::handle_command(root, Arc::clone(self.services.get()), command).await
234 }
235
236 pub async fn user_is_admin(&self, user_id: &UserId) -> bool {
238 if user_id == self.services.globals.server_user {
239 return true;
240 }
241
242 let Ok(admin_room) = self.get_admin_room().await else {
243 return false;
244 };
245
246 self.services
247 .state_cache
248 .is_joined(user_id, &admin_room)
249 .await
250 }
251
252 pub async fn user_is_last_admin(&self, user_id: &UserId) -> bool {
260 let server_user: &UserId = &self.services.globals.server_user;
261 if user_id == server_user {
262 return false;
263 }
264
265 let Ok(admin_room) = self.get_admin_room().await else {
266 return false;
267 };
268
269 if !self
270 .services
271 .state_cache
272 .is_joined(user_id, &admin_room)
273 .await
274 {
275 return false;
276 }
277
278 !self
279 .services
280 .state_cache
281 .active_local_users_in_room(&admin_room)
282 .ready_any(|member| member != user_id && member != server_user)
283 .await
284 }
285
286 pub async fn get_admin_room(&self) -> Result<OwnedRoomId> {
291 let room_id = self
292 .services
293 .alias
294 .resolve_local_alias(&self.admin_alias)
295 .await?;
296
297 self.services
298 .state_cache
299 .is_joined(&self.services.globals.server_user, &room_id)
300 .await
301 .then_some(room_id)
302 .ok_or_else(|| err!(Request(NotFound("Admin user not joined to admin room"))))
303 }
304
305 pub async fn get_report_room(&self) -> Result<OwnedRoomId> {
308 let Some(report_room) = self.services.server.config.report_room.as_ref() else {
309 return self.get_admin_room().await;
310 };
311
312 match self.resolve_report_room(report_room).await {
313 | Ok(room_id) => Ok(room_id),
314 | Err(e) => {
315 warn!(%report_room, error = %e, "Falling back to the admin room for reports");
316 self.get_admin_room().await
317 },
318 }
319 }
320
321 async fn resolve_report_room(&self, report_room: &RoomOrAliasId) -> Result<OwnedRoomId> {
322 let room_id = self
323 .services
324 .alias
325 .maybe_resolve(report_room)
326 .await?;
327
328 self.services
329 .state_cache
330 .is_joined(&self.services.globals.server_user, &room_id)
331 .await
332 .then_some(room_id)
333 .ok_or_else(|| err!("server user is not joined to the configured report room"))
334 }
335
336 pub async fn is_admin_command<Pdu>(&self, event: &Pdu, body: &str) -> bool
344 where
345 Pdu: Event,
346 {
347 let body = body.trim_start();
348
349 let is_escape = body.starts_with('\\');
351 let is_public_escape = is_escape
352 && body
353 .trim_start_matches('\\')
354 .starts_with("!admin");
355
356 let server_user = &self.services.globals.server_user;
358 let is_public_prefix =
359 body.starts_with("!admin") || body.starts_with(server_user.as_str());
360
361 if !is_public_escape && !is_public_prefix {
363 return false;
364 }
365
366 let user_is_local = self
367 .services
368 .globals
369 .user_is_local(event.sender());
370
371 if is_public_escape && !user_is_local {
373 return false;
374 }
375
376 if is_public_escape && !self.services.server.config.admin_escape_commands {
378 return false;
379 }
380
381 if event.msgtype() != Some(MsgType::Text) {
384 return false;
385 }
386
387 if is_public_prefix && !self.is_admin_room(event.room_id()).await {
389 return false;
390 }
391
392 if !self.user_is_admin(event.sender()).await {
394 return false;
395 }
396
397 let emergency_password_set = self
400 .services
401 .server
402 .config
403 .emergency_password
404 .is_some();
405 let from_server = event.sender() == server_user && !emergency_password_set;
406 if from_server && self.is_admin_room(event.room_id()).await {
407 return false;
408 }
409
410 true
412 }
413
414 #[must_use]
415 pub async fn is_admin_room(&self, room_id_: &RoomId) -> bool {
416 self.get_admin_room()
417 .map_ok(|room_id| room_id == room_id_)
418 .await
419 .unwrap_or(false)
420 }
421}
422
423#[implement(Service)]
430pub async fn lock_admin_room(&self) -> Option<RoomMutexGuard> {
431 let admin_room = self.get_admin_room().await.ok()?;
432
433 self.services
434 .state
435 .mutex
436 .lock(&admin_room)
437 .await
438 .into()
439}