Skip to main content

tuwunel_service/admin/
mod.rs

1mod 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/// Inputs to a command: its multi-line text, the event to reply to, and who
44/// sent it.
45///
46/// An input without a sender is the operator's, from the console or the
47/// `admin_execute` and `admin_signal_execute` lists; converting a bare command
48/// string builds one.
49#[derive(Clone, Debug, Default)]
50pub struct CommandInput {
51	/// The command line, followed by any body lines.
52	pub command: String,
53
54	/// The event the command's response replies to.
55	pub reply_id: Option<OwnedEventId>,
56
57	/// The user who sent the command, or `None` for the operator.
58	pub sender: Option<OwnedUserId>,
59}
60
61/// Root of a clap command tree installed by a downstream crate.
62#[async_trait]
63pub trait Command: Send + Sync + 'static {
64	/// The clap command tree; equivalent to
65	/// `<C as clap::CommandFactory>::command()`.
66	fn clap(&self) -> clap::Command;
67
68	/// Dispatch already-parsed argument matches to the matching handler.
69	async fn dispatch(&self, matches: clap::ArgMatches, context: &Context<'_>) -> Result;
70}
71
72/// Carries a rendered command outcome while preserving its status.
73///
74/// `Ok(Some(output))` reports success, `Err(output)` reports failure, and
75/// `Ok(None)` suppresses the response. Callers do not infer status from text.
76pub type ProcessorResult = Result<Option<CommandOutput>, CommandOutput>;
77
78/// Textual output of a completed command. Markdown is the norm; Plain carries
79/// clap usage and error text, which must never be markdown-rendered.
80pub 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
103/// Maximum number of commands which can be queued for dispatch.
104const 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		//TODO: not unwind safe
145		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	/// Queues a command for the service worker and returns once it is queued.
167	///
168	/// The worker processes it later and replies to the `reply_id` event in its
169	/// room, posting nothing without one. Queueing waits while the queue is
170	/// full, and errors when the queue is unavailable or closed.
171	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	/// Dispatches a command to the processor on the current task and waits for
188	/// completion.
189	///
190	/// The queue is bypassed, so the outcome returns to the caller rather than
191	/// being posted as a reply.
192	pub async fn command_in_place(&self, input: CommandInput) -> ProcessorResult {
193		self.process_command(&input).await
194	}
195
196	/// Invokes the tab-completer to complete the command. When unavailable,
197	/// None is returned.
198	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	/// Checks whether a given user is an admin of this server
237	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	/// Checks whether a given user is the only active admin left on this server.
253	///
254	/// The server user is never counted: it can sign in only while an emergency
255	/// password is configured. Deactivated accounts still joined to the admin
256	/// room are not counted either, since none of them can sign in to act. Nor
257	/// is a passwordless account, such as an appservice's user, since it stores
258	/// the same empty password as a deactivated one.
259	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	/// Gets the room ID of the admin room
287	///
288	/// Errors are propagated from the database, and will have None if there is
289	/// no admin room
290	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	/// Gets the room reports are posted to: the configured report room when set
306	/// and usable, otherwise the admin room.
307	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	/// Returns whether a message event is an admin command to run.
337	///
338	/// Only an `m.text` from an admin qualifies: prefixed with `!admin` or the
339	/// server user's ID in the admin room, or escaped as `\!admin` by a local
340	/// admin anywhere when escape commands are enabled. The server user's own
341	/// messages in the admin room are refused unless the emergency password is
342	/// set.
343	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		// Server-side command-escape with public echo
350		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		// Admin command with public echo (in admin room)
357		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		// Expected backward branch
362		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		// only allow public escaped commands by local admins
372		if is_public_escape && !user_is_local {
373			return false;
374		}
375
376		// Check if server-side command-escape is disabled by configuration
377		if is_public_escape && !self.services.server.config.admin_escape_commands {
378			return false;
379		}
380
381		// Spec: an m.notice must never be answered automatically. Other msgtypes'
382		// bodies are captions, filenames or emotes a forward or repost can carry.
383		if event.msgtype() != Some(MsgType::Text) {
384			return false;
385		}
386
387		// Prevent unescaped !admin from being used outside of the admin room
388		if is_public_prefix && !self.is_admin_room(event.room_id()).await {
389			return false;
390		}
391
392		// Only senders who are admin can proceed
393		if !self.user_is_admin(event.sender()).await {
394			return false;
395		}
396
397		// This will evaluate to false if the emergency password is set up so that
398		// the administrator can execute commands as the server user
399		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		// Authentic admin command
411		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/// Locks the admins room's state, when there is an admins room.
424///
425/// Hold the guard from the [`Service::user_is_last_admin`] check until the
426/// change it permits is made. The admins room's leave and ban guard runs under
427/// the same lock, so two concurrent removals cannot each see the other as the
428/// admin who remains.
429#[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}