Skip to main content

tuwunel_service/media/
mod.rs

1mod data;
2pub(super) mod migrations;
3mod preview;
4mod remote;
5mod tests;
6mod thumbnail;
7#[cfg(feature = "media_thumbnail")]
8mod video;
9use std::{
10	collections::{HashMap, HashSet},
11	path::PathBuf,
12	sync::{Arc, Mutex},
13	time::{Duration, Instant, SystemTime},
14};
15
16use async_trait::async_trait;
17use base64::{Engine as _, engine::general_purpose};
18use futures::{FutureExt, Stream, StreamExt, TryFutureExt, TryStreamExt, pin_mut};
19use http::StatusCode;
20use object_store::ObjectMeta;
21use ruma::{
22	Mxc, OwnedMxcUri, OwnedUserId, UserId,
23	api::error::{ErrorKind, RetryAfter},
24	http_headers::ContentDisposition,
25};
26#[cfg(feature = "media_thumbnail")]
27use tokio::sync::Semaphore;
28use tokio::{fs, sync::Notify};
29use tuwunel_core::{
30	Err, Error, Result, debug, debug_error, debug_info, debug_warn, err, trace,
31	utils::{
32		self, BoolExt, MutexMap,
33		result::LogDebugErr,
34		stream::{BroadbandExt, IterStream, ReadyExt, TryReadyExt},
35		time::now_millis,
36	},
37	warn,
38};
39use url::Url;
40
41#[cfg(feature = "media_thumbnail")]
42use self::video::{FAILURES, Failures, sweep_staging_dir};
43use self::{data::Data, preview::Agent, remote::Fetch, thumbnail::sequence};
44pub use self::{
45	data::Metadata,
46	preview::UrlPreviewData,
47	thumbnail::{Animate, Dim},
48};
49use crate::storage::Provider;
50
51#[derive(Debug)]
52pub struct Media {
53	pub content: Vec<u8>,
54	pub content_type: Option<String>,
55	pub content_disposition: Option<ContentDisposition>,
56}
57
58/// A picture fetched from a peer, and what a walk of its container settled.
59///
60/// The walk picking the type a fetched row is filed under settles the same
61/// question a caller asks before serving the picture, so the answer travels
62/// with the bytes rather than being read out of them a second time. A redirect
63/// files no row and so walks nothing, which is what an absent answer means.
64#[derive(Debug)]
65pub struct Fetched {
66	pub media: Media,
67	/// Whether the bytes carry a frame sequence, where storing them settled it.
68	///
69	/// Private to the crate so nobody outside it can pair an answer with bytes
70	/// it was not walked from, which would defeat the gate reading it.
71	pub(crate) animates: Option<bool>,
72}
73
74/// One row of a user's uploaded media, holding only the fields tuwunel can
75/// derive from the uploader index and storage-provider object metadata. The
76/// uploader is carried when the index holds a row for it.
77#[derive(Clone, Debug)]
78pub struct UserMediaEntry {
79	pub mxc: OwnedMxcUri,
80	pub media_type: Option<String>,
81	pub upload_name: Option<String>,
82	pub media_length: Option<u64>,
83	pub created_ts: u64,
84	pub user_id: Option<OwnedUserId>,
85}
86
87/// One locally-uploaded media item's uploader and storage-object byte length
88/// and modification time, the row shape of the media-statistics scan.
89#[derive(Clone, Debug)]
90pub struct UploadStat {
91	pub user_id: OwnedUserId,
92	pub media_length: u64,
93	pub created_ts: u64,
94}
95
96/// For MSC2246
97struct MXCState {
98	/// Save the notifier for each pending media upload
99	notifiers: Mutex<HashMap<OwnedMxcUri, Arc<Notify>>>,
100	/// Save the ratelimiter for each user
101	ratelimiter: Mutex<HashMap<OwnedUserId, (Instant, f64)>>,
102}
103
104pub struct Service {
105	pub(super) db: Data,
106	services: Arc<crate::services::OnceServices>,
107	url_preview_mutex: MutexMap<String, ()>,
108	federation_mutex: MutexMap<String, ()>,
109	mxc_state: MXCState,
110	#[cfg(feature = "media_thumbnail")]
111	animated_thumbnail_slots: Arc<Semaphore>,
112	#[cfg(feature = "media_thumbnail")]
113	video_thumbnail_slots: Semaphore,
114	#[cfg(feature = "media_thumbnail")]
115	video_thumbnail_failures: Mutex<Failures>,
116}
117
118/// generated MXC ID (`media-id`) length
119pub const MXC_LENGTH: usize = 32;
120
121/// Cache control for immutable objects.
122pub const CACHE_CONTROL_IMMUTABLE: &str = "private,max-age=31536000,immutable";
123
124/// Default cross-origin resource policy.
125pub const CORP_CROSS_ORIGIN: &str = "cross-origin";
126
127/// Validity window for a presigned media download redirect (MSC3860).
128const REDIRECT_TTL: Duration = Duration::from_mins(5);
129
130#[async_trait]
131impl crate::Service for Service {
132	fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
133		#[cfg(feature = "media_thumbnail")]
134		let animated_thumbnail_slots = Arc::new(Semaphore::new(
135			args.server
136				.config
137				.media_thumbnail_animated_concurrency
138				.max(1),
139		));
140
141		let service = Arc::new(Self {
142			db: Data::new(args.db),
143			services: args.services.clone(),
144			url_preview_mutex: MutexMap::new(),
145			federation_mutex: MutexMap::new(),
146			mxc_state: MXCState {
147				notifiers: Mutex::new(HashMap::new()),
148				ratelimiter: Mutex::new(HashMap::new()),
149			},
150			#[cfg(feature = "media_thumbnail")]
151			video_thumbnail_failures: Failures::new(FAILURES).into(),
152			#[cfg(feature = "media_thumbnail")]
153			animated_thumbnail_slots,
154			#[cfg(feature = "media_thumbnail")]
155			video_thumbnail_slots: Semaphore::new(
156				args.server
157					.config
158					.media_video_thumbnail_concurrency
159					.max(1),
160			),
161		});
162
163		#[cfg(feature = "media_thumbnail")]
164		sweep_staging_dir(&args.server.config);
165
166		Ok(service)
167	}
168
169	fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
170}
171
172impl Service {
173	/// Create a pending media upload ID.
174	#[tracing::instrument(level = "debug", skip(self))]
175	pub async fn create_pending(
176		&self,
177		mxc: &Mxc<'_>,
178		user: &UserId,
179		unused_expires_at: u64,
180	) -> Result {
181		let config = &self.services.server.config;
182
183		// Rate limiting (rc_media_create)
184		let rate = f64::from(config.media_rc_create_per_second);
185		let burst = f64::from(config.media_rc_create_burst_count);
186
187		// Check rate limiting
188		if rate > 0.0 && burst > 0.0 {
189			let now = Instant::now();
190			let mut ratelimiter = self.mxc_state.ratelimiter.lock()?;
191
192			let (last_time, tokens) = ratelimiter
193				.entry(user.to_owned())
194				.or_insert_with(|| (now, burst));
195
196			let elapsed = now.duration_since(*last_time).as_secs_f64();
197			let new_tokens = elapsed.mul_add(rate, *tokens).min(burst);
198
199			if new_tokens >= 1.0 {
200				*last_time = now;
201				*tokens = new_tokens - 1.0;
202			} else {
203				return Err(Error::Request(
204					ErrorKind::LimitExceeded(ruma::api::error::LimitExceededErrorData {
205						retry_after: None,
206					}),
207					"Too many pending media creation requests.".into(),
208					StatusCode::TOO_MANY_REQUESTS,
209				));
210			}
211		}
212
213		let max_uploads = config.max_pending_media_uploads;
214		let (current_uploads, earliest_expiration) =
215			self.db.count_pending_mxc_for_user(user).await;
216
217		// Check if the user has reached the maximum number of pending media uploads
218		if current_uploads >= max_uploads {
219			let retry_after = earliest_expiration.saturating_sub(now_millis());
220			return Err(Error::Request(
221				ErrorKind::LimitExceeded(ruma::api::error::LimitExceededErrorData {
222					retry_after: Some(RetryAfter::Delay(Duration::from_millis(retry_after))),
223				}),
224				"Maximum number of pending media uploads reached.".into(),
225				StatusCode::TOO_MANY_REQUESTS,
226			));
227		}
228
229		self.db
230			.insert_pending_mxc(mxc, user, unused_expires_at);
231
232		Ok(())
233	}
234
235	/// Uploads content to a pending media ID.
236	#[tracing::instrument(level = "debug", skip(self))]
237	pub async fn upload_pending(
238		&self,
239		mxc: &Mxc<'_>,
240		user: &UserId,
241		content_disposition: Option<&ContentDisposition>,
242		content_type: Option<&str>,
243		file: &[u8],
244	) -> Result {
245		let Ok((owner_id, expires_at)) = self.db.search_pending_mxc(mxc).await else {
246			if self.get_metadata(mxc).await.is_some() {
247				return Err!(Request(CannotOverwriteMedia("Media ID already has content")));
248			}
249
250			return Err!(Request(NotFound("Media not found")));
251		};
252
253		if owner_id != user {
254			return Err!(Request(Forbidden("You did not create this media ID")));
255		}
256
257		let current_time = now_millis();
258		if expires_at < current_time {
259			return Err!(Request(NotFound("Pending media ID expired")));
260		}
261
262		self.create(mxc, Some(user), content_disposition, content_type, file)
263			.await?;
264
265		self.db.remove_pending_mxc(mxc);
266
267		let mxc_uri: OwnedMxcUri = mxc.to_string().into();
268		let notifier = self.mxc_state.notifiers.lock()?.remove(&mxc_uri);
269
270		if let Some(notifier) = notifier {
271			notifier.notify_waiters();
272		}
273
274		Ok(())
275	}
276
277	/// Uploads a file and reports whether its own bytes carry a sequence.
278	///
279	/// The declared type is whoever uploaded it saying so, and a picture that
280	/// animates is stored under the type its own container names instead. That
281	/// is the only record of whether the media animates that a later lookup can
282	/// read without fetching the whole of it. The disposition is left as the
283	/// caller computed it, so a file already bound for download stays there.
284	pub async fn create(
285		&self,
286		mxc: &Mxc<'_>,
287		user: Option<&UserId>,
288		content_disposition: Option<&ContentDisposition>,
289		content_type: Option<&str>,
290		file: &[u8],
291	) -> Result<bool> {
292		let walk = sequence(file);
293
294		// Width, Height = 0 if it's not a thumbnail
295		let key = self.db.create_file_metadata(
296			mxc,
297			user,
298			&Dim::default(),
299			content_disposition,
300			walk.stored_type(content_type),
301		)?;
302
303		//TODO: Dangling metadata in database if creation fails
304		self.create_media_file(&key, file)
305			.map_ok(|()| walk.animates())
306			.await
307	}
308
309	/// Deletes a file in the database and from the media directory via an MXC
310	#[tracing::instrument(level = "trace", skip(self))]
311	pub async fn delete(&self, mxc: &Mxc<'_>) -> Result {
312		// lazy URL-preview media has no file keys of its own; drop its reference
313		// and staged bytes whenever present so a delete can't re-mint the media
314		let had_lazy = self.db.search_lazy_media(mxc).await.is_ok();
315		if had_lazy {
316			let key = mxc.to_string();
317			let mut txn = self.services.db.txn();
318
319			self.db.remove_lazy_media(&mut txn, &key);
320			self.db.remove_lazy_content(&mut txn, &key);
321			txn.execute();
322		}
323
324		match self.db.search_mxc_metadata_prefix(mxc).await {
325			| Ok(keys) => {
326				for key in keys {
327					trace!(?mxc, "MXC Key: {key:?}");
328					debug_info!(?mxc, "Deleting from storage provider");
329
330					if let Err(e) = self.remove_media_file(&key).await {
331						debug_error!(?mxc, "Failed to remove media file: {e}");
332					}
333
334					debug_info!(?mxc, "Deleting from database");
335					self.db.delete_file_mxc(mxc).await;
336				}
337
338				Ok(())
339			},
340			| _ if had_lazy => Ok(()),
341			| _ => Err!(Database(error!(
342				"Failed to find any media keys for MXC {mxc} in our database."
343			))),
344		}
345	}
346
347	/// Deletes all media by the specified user
348	///
349	/// currently, this is only practical for local users
350	#[tracing::instrument(level = "trace", skip(self))]
351	pub async fn delete_from_user(&self, user: &UserId) -> Result<usize> {
352		let mxcs = self.db.get_all_user_mxcs(user).await;
353		let mut deletion_count: usize = 0;
354
355		for mxc in mxcs {
356			let Ok(mxc) = mxc.as_str().try_into().inspect_err(|e| {
357				debug_error!(?mxc, "Failed to parse MXC URI from database: {e}");
358			}) else {
359				continue;
360			};
361
362			debug_info!(
363				%deletion_count,
364				"Deleting MXC {mxc} by user {user} from database and filesystem",
365			);
366			match self.delete(&mxc).await {
367				| Ok(()) => {
368					deletion_count = deletion_count.saturating_add(1);
369				},
370				| Err(e) => {
371					debug_error!(
372						%deletion_count,
373						"Failed to delete {mxc} from user {user}, ignoring error: {e}"
374					);
375				},
376			}
377		}
378
379		Ok(deletion_count)
380	}
381
382	/// Get file from local storage or make a federation request if it
383	/// originates remotely.
384	#[tracing::instrument(
385		level = "debug",
386		err(level = "debug")
387		skip(self),
388	)]
389	pub async fn get_or_fetch(&self, mxc: &Mxc<'_>, timeout_ms: Duration) -> Result<Media> {
390		if let Ok(media) = self.get(mxc, Some(timeout_ms)).await {
391			return Ok(media);
392		}
393
394		if self
395			.services
396			.globals
397			.server_is_ours(mxc.server_name)
398		{
399			return Err!(Request(NotFound("Local media not found.")));
400		}
401
402		let lock = self.federation_mutex.lock(&mxc.to_string()).await;
403
404		if self
405			.db
406			.file_metadata_exists(mxc, &Dim::default())
407			.await
408		{
409			drop(lock);
410			return self.get(mxc, None).await;
411		}
412
413		self.fetch_remote_content(mxc, None, timeout_ms)
414			.map_ok(|fetched| fetched.media)
415			.await
416	}
417
418	/// Get file from local storage while waiting up to a timeout_ms if it is
419	/// pending.
420	#[tracing::instrument(
421		level = "debug",
422		err(level = "trace")
423		skip(self),
424	)]
425	pub async fn get(&self, mxc: &Mxc<'_>, timeout: Option<Duration>) -> Result<Media> {
426		if let Ok(meta) = self.get_stored(mxc).await {
427			return Ok(meta);
428		}
429
430		let Some(timeout) = timeout else {
431			return Err!(Request(NotFound("Media not found.")));
432		};
433
434		let Ok(_pending) = self.db.search_pending_mxc(mxc).await else {
435			return Err!(Request(NotFound("Media not found.")));
436		};
437
438		let notifier = self
439			.mxc_state
440			.notifiers
441			.lock()?
442			.entry(mxc.to_string().into())
443			.or_insert_with(|| Arc::new(Notify::new()))
444			.clone();
445
446		if tokio::time::timeout(timeout, notifier.notified())
447			.await
448			.is_err()
449		{
450			return Err!(Request(NotYetUploaded("Media has not been uploaded yet")));
451		}
452
453		self.get_stored(mxc).await
454	}
455
456	/// Get file from local storage.
457	///
458	/// A metadata miss falls through to the lazy-media path, which serializes
459	/// on a per-mxc mutex and fetches the origin once.
460	#[tracing::instrument(level = "debug", skip(self))]
461	pub async fn get_stored(&self, mxc: &Mxc<'_>) -> Result<Media> {
462		let meta = self
463			.db
464			.search_file_metadata(mxc, &Dim::default(), Animate::Allowed)
465			.await;
466
467		let Ok(Metadata { content_type, content_disposition, key }) = meta else {
468			// query-depth firewall
469			return Box::pin(self.fetch_lazy_media(mxc)).await;
470		};
471
472		let path = self.get_media_name_sha256(&key);
473		let fetch = self
474			.storage_providers()
475			.stream()
476			.filter_map(async |provider| {
477				provider
478					.get(path.as_str())
479					.await
480					.log_debug_err()
481					.ok()
482			});
483
484		pin_mut!(fetch);
485		let Some(bytes) = fetch.next().await else {
486			return Err!(Request(NotFound("Media not found.")));
487		};
488
489		Ok(Media {
490			content: bytes.to_vec(),
491			content_type,
492			content_disposition,
493		})
494	}
495
496	/// Resolve lazy URL-preview media on first download, promoting the staged
497	/// or freshly-fetched bytes into the media store so the origin is fetched
498	/// at most once per item.
499	#[tracing::instrument(level = "debug", skip(self))]
500	async fn fetch_lazy_media(&self, mxc: &Mxc<'_>) -> Result<Media> {
501		let key = mxc.to_string();
502
503		// bound outbound amplification to one in-flight fetch per mxc; requests
504		// for the same mxc queue rather than fan out
505		let _lock = self.federation_mutex.lock(&key).await;
506
507		// a queued waiter for the same mxc may have promoted it already
508		if self
509			.db
510			.search_file_metadata(mxc, &Dim::default(), Animate::Allowed)
511			.await
512			.is_ok()
513		{
514			// box the cold re-entry to break the get_stored/fetch_lazy_media cycle
515			return Box::pin(self.get_stored(mxc)).await;
516		}
517
518		let media = match self.db.get_lazy_content(&key).await {
519			| Ok(media) => media,
520			| Err(_) => {
521				let Ok(url) = self.db.search_lazy_media(mxc).await else {
522					return Err!(Request(NotFound("Media not found.")));
523				};
524
525				let limit = self.services.config.url_preview_max_media_size;
526
527				self.location_request(Fetch::Preview(Agent::Media), &url, limit)
528					.await?
529			},
530		};
531
532		// promote through the ordinary upload path so real file metadata makes
533		// every later download and thumbnail a normal-media hit
534		if let Err(e) = self
535			.create(
536				mxc,
537				None,
538				media.content_disposition.as_ref(),
539				media.content_type.as_deref(),
540				&media.content,
541			)
542			.await
543		{
544			// a failed promotion leaves metadata without bytes, masking the lazy
545			// fallback on every later read; drop it so the next read retries
546			self.db.delete_file_mxc(mxc).await;
547
548			return Err(e);
549		}
550
551		let mut txn = self.services.db.txn();
552
553		self.db.remove_lazy_media(&mut txn, &key);
554		self.db.remove_lazy_content(&mut txn, &key);
555		txn.execute();
556
557		Ok(media)
558	}
559
560	/// Presigned redirect URL for locally-stored media (MSC3860).
561	///
562	/// Returns the first configured provider's signed URL for the object, or
563	/// `None` when redirects are disabled, the media is unknown, or no
564	/// provider can presign it (filesystem-only media). A stored file that
565	/// animates also answers `None` to a request forbidding animation, since
566	/// handing over the object as it stands cannot re-encode it.
567	#[tracing::instrument(level = "debug", skip(self))]
568	pub async fn redirect_url(
569		&self,
570		mxc: &Mxc<'_>,
571		dim: &Dim,
572		animate: Animate,
573	) -> Result<Option<Url>> {
574		if !self.services.config.media_allow_redirect {
575			return Ok(None);
576		}
577
578		let Ok(Metadata { key, content_type, .. }) = self
579			.db
580			.search_file_metadata(mxc, dim, animate)
581			.await
582		else {
583			return Ok(None);
584		};
585
586		// a redirect hands over the stored object as it is, so a variant the
587		// request forbids has to fall through to the re-encoding path
588		if !animate.accepts_type(content_type.as_deref()) {
589			return Ok(None);
590		}
591
592		let path = self.get_media_name_sha256(&key);
593		let urls = self
594			.storage_providers()
595			.stream()
596			.filter_map(async |provider| {
597				provider
598					.signed_get_url(path.as_str(), REDIRECT_TTL)
599					.await
600					.log_debug_err()
601					.ok()
602					.flatten()
603			});
604
605		pin_mut!(urls);
606
607		Ok(urls.next().await)
608	}
609
610	/// Gets all the MXC URIs in our media database
611	pub async fn get_all_mxcs(&self) -> Result<Vec<OwnedMxcUri>> {
612		let all_keys = self.db.get_all_media_keys().await;
613
614		let mut mxcs = Vec::with_capacity(all_keys.len());
615
616		for key in all_keys {
617			trace!("Full MXC key from database: {key:?}");
618
619			let mut parts = key.split(|&b| b == 0xFF);
620			let mxc = parts
621				.next()
622				.map(|bytes| {
623					utils::string_from_bytes(bytes).map_err(|e| {
624						err!(Database(error!(
625							"Failed to parse MXC unicode bytes from our database: {e}"
626						)))
627					})
628				})
629				.transpose()?;
630
631			let Some(mxc_s) = mxc else {
632				debug_warn!(
633					?mxc,
634					"Parsed MXC URL unicode bytes from database but is still invalid"
635				);
636				continue;
637			};
638
639			trace!("Parsed MXC key to URL: {mxc_s}");
640			let mxc = OwnedMxcUri::from(mxc_s);
641
642			if mxc.is_valid() {
643				mxcs.push(mxc);
644			} else {
645				debug_warn!("{mxc:?} from database was found to not be valid");
646			}
647		}
648
649		Ok(mxcs)
650	}
651
652	/// Every media item uploaded by a local user, carrying the fields tuwunel
653	/// can derive: content type, upload name, byte length and modification time
654	/// from storage-provider object metadata. Untracked columns (last-access,
655	/// quarantine, url-cache) are not represented.
656	#[tracing::instrument(level = "debug", skip(self))]
657	pub async fn user_media(&self, user: &UserId) -> Result<Vec<UserMediaEntry>> {
658		let entries = self
659			.db
660			.get_all_user_mxcs(user)
661			.await
662			.into_iter()
663			.stream()
664			.broad_filter_map(async |mxc| self.user_media_entry(Some(user), mxc).await)
665			.collect()
666			.await;
667
668		Ok(entries)
669	}
670
671	/// Derivable metadata for the single media item at the given MXC, with the
672	/// uploading local user resolved from the uploader index when present.
673	#[tracing::instrument(level = "debug", skip(self))]
674	pub async fn media_entry(&self, mxc: &Mxc<'_>) -> Option<UserMediaEntry> {
675		let user = self.db.mxc_user(mxc).await;
676
677		self.user_media_entry(user.as_deref(), mxc.to_string().into())
678			.await
679	}
680
681	async fn user_media_entry(
682		&self,
683		user: Option<&UserId>,
684		mxc: OwnedMxcUri,
685	) -> Option<UserMediaEntry> {
686		let parts = mxc.parts().ok()?;
687		let Metadata { content_type, content_disposition, key } =
688			self.get_metadata(&parts).await?;
689
690		let object = self.head_meta(&key).await;
691		let upload_name = content_disposition.and_then(|disposition| disposition.filename);
692
693		Some(UserMediaEntry {
694			media_type: content_type,
695			upload_name,
696			media_length: object.as_ref().map(|object| object.size),
697			created_ts: object.as_ref().map(mtime_millis).unwrap_or(0),
698			user_id: user.map(ToOwned::to_owned),
699			mxc,
700		})
701	}
702
703	/// Uploader, byte length and storage modification time of every media item
704	/// uploaded by a local user, one row per upload; media missing from every
705	/// storage provider are skipped.
706	pub fn upload_stats(&self) -> impl Stream<Item = UploadStat> + Send + '_ {
707		self.db
708			.all_uploads()
709			.broad_filter_map(async |(mxc, user_id)| {
710				let parts = mxc.parts().ok()?;
711				let Metadata { key, .. } = self.get_metadata(&parts).await?;
712				let object = self.head_meta(&key).await?;
713
714				Some(UploadStat {
715					user_id,
716					media_length: object.size,
717					created_ts: mtime_millis(&object),
718				})
719			})
720	}
721
722	/// Deletes local media older than `before_ts` (by storage-provider mtime)
723	/// and strictly larger than `size_gt` bytes, sparing profile and
724	/// room-avatar media when `keep_profiles`. Returns the deleted MXCs, empty
725	/// when none match. Quarantine and protection flags are not tracked, so no
726	/// media is spared on those grounds.
727	#[tracing::instrument(level = "debug", skip(self))]
728	pub async fn delete_by_date_size(
729		&self,
730		before_ts: u64,
731		size_gt: u64,
732		keep_profiles: bool,
733	) -> Result<Vec<OwnedMxcUri>> {
734		let spared = keep_profiles
735			.then_async(|| self.avatar_mxcs())
736			.await
737			.unwrap_or_default();
738
739		let candidates = self.get_all_mxcs().await?;
740
741		let deleted: Vec<OwnedMxcUri> = candidates
742			.into_iter()
743			.stream()
744			.ready_filter(|mxc| self.is_local(mxc) && !spared.contains(mxc))
745			.broad_filter_map(async |mxc| {
746				let parts = mxc.parts().ok()?;
747				let Metadata { key, .. } = self.get_metadata(&parts).await?;
748				let object = self.head_meta(&key).await?;
749
750				let eligible = mtime_millis(&object) < before_ts && object.size > size_gt;
751
752				eligible
753					.then_async(|| self.delete(&parts))
754					.await
755					.and_then(Result::ok)
756					.map(|()| mxc)
757			})
758			.collect()
759			.await;
760
761		Ok(deleted)
762	}
763
764	/// The MXCs of every local user's profile avatar and every room's avatar,
765	/// the spare-set honoured by `keep_profiles`.
766	async fn avatar_mxcs(&self) -> HashSet<OwnedMxcUri> {
767		let user_avatars = self
768			.services
769			.users
770			.list_local_users()
771			.map(ToOwned::to_owned)
772			.broad_filter_map(async |user| self.services.profile.avatar_url(&user).await.ok());
773
774		let room_avatars = self
775			.services
776			.metadata
777			.iter_ids()
778			.map(ToOwned::to_owned)
779			.broad_filter_map(async |room_id| {
780				self.services
781					.state_accessor
782					.get_avatar(&room_id)
783					.await
784					.ok()
785					.and_then(|avatar| avatar.url)
786			});
787
788		user_avatars.chain(room_avatars).collect().await
789	}
790
791	fn is_local(&self, mxc: &OwnedMxcUri) -> bool {
792		mxc.server_name()
793			.is_ok_and(|server| self.services.globals.server_is_ours(server))
794	}
795
796	/// First storage provider's object metadata for the media stored under
797	/// `key` (byte length and modification time), or `None` when no provider
798	/// holds it.
799	async fn head_meta(&self, key: &[u8]) -> Option<ObjectMeta> {
800		let path = self.get_media_name_sha256(key);
801
802		let stream = self
803			.storage_providers()
804			.stream()
805			.filter_map(async |provider| provider.head(&path).await.log_debug_err().ok());
806
807		pin_mut!(stream);
808		stream.next().await
809	}
810
811	/// Deletes all media files before or after the given time. Returns a usize
812	/// with the number of media files deleted.
813	pub async fn delete_range(
814		&self,
815		time: SystemTime,
816		older_than: bool,
817		newer_than: bool,
818		yes_i_want_to_delete_local_media: bool,
819	) -> Result<usize> {
820		let all_keys = self.db.get_all_media_keys().await;
821		let mut remote_mxcs = Vec::with_capacity(all_keys.len());
822
823		for key in all_keys {
824			trace!("Full MXC key from database: {key:?}");
825			let mut parts = key.split(|&b| b == 0xFF);
826			let mxc = parts
827				.next()
828				.map(|bytes| {
829					utils::string_from_bytes(bytes).map_err(|e| {
830						err!(Database(error!(
831							"Failed to parse MXC unicode bytes from our database: {e}"
832						)))
833					})
834				})
835				.transpose()?;
836
837			let Some(mxc_s) = mxc else {
838				debug_warn!(
839					?mxc,
840					"Parsed MXC URL unicode bytes from database but is still invalid"
841				);
842				continue;
843			};
844
845			trace!("Parsed MXC key to URL: {mxc_s}");
846			let mxc = OwnedMxcUri::from(mxc_s);
847			if (mxc.server_name() == Ok(self.services.globals.server_name())
848				&& !yes_i_want_to_delete_local_media)
849				|| !mxc.is_valid()
850			{
851				debug!("Ignoring local or broken media MXC: {mxc}");
852				continue;
853			}
854
855			let file_created_at = if let Some(file_metadata) = self
856				.storage_providers()
857				.stream()
858				.filter_map(async |provider| {
859					let path = self.get_media_name_sha256(&key);
860					match provider.head(&path).await {
861						| Ok(file_metadata) => {
862							trace!(%mxc, ?path, "Provider file metadata: {file_metadata:?}");
863							Some(file_metadata)
864						},
865						| Err(e) => {
866							debug_warn!(
867								"Failed to obtain {:?} file metadata for MXC {mxc} at file path \
868								 {path:?}\", skipping: {e}",
869								provider.name,
870							);
871							None
872						},
873					}
874				})
875				.boxed()
876				.next()
877				.await
878			{
879				SystemTime::from(file_metadata.last_modified)
880			} else {
881				continue;
882			};
883
884			debug!("File created at: {file_created_at:?}");
885
886			if file_created_at <= time && older_than {
887				debug!(
888					"File is older than user duration, pushing to list of file paths and keys \
889					 to delete."
890				);
891				remote_mxcs.push(mxc.to_string());
892			} else if file_created_at >= time && newer_than {
893				debug!(
894					"File is newer than user duration, pushing to list of file paths and keys \
895					 to delete."
896				);
897				remote_mxcs.push(mxc.to_string());
898			}
899		}
900
901		debug_info!("Deleting media now in the past {time:?}");
902
903		let mut deletion_count: usize = 0;
904
905		for mxc in remote_mxcs {
906			let Ok(mxc) = mxc.as_str().try_into() else {
907				debug_warn!("Invalid MXC in database, skipping");
908				continue;
909			};
910
911			debug_info!("Deleting MXC {mxc} from database and filesystem");
912
913			match self.delete(&mxc).await {
914				| Ok(()) => {
915					deletion_count = deletion_count.saturating_add(1);
916				},
917				| Err(e) => {
918					warn!("Failed to delete {mxc}, ignoring error and skipping: {e}");
919				},
920			}
921		}
922
923		Ok(deletion_count)
924	}
925
926	pub async fn create_media_dir(&self) -> Result {
927		let dir = self.get_media_dir();
928		Ok(fs::create_dir_all(dir).await?)
929	}
930
931	async fn remove_media_file(&self, key: &[u8]) -> Result {
932		let path = self.get_media_name_sha256(key);
933		self.storage_providers()
934			.stream()
935			.filter_map(async |provider| {
936				debug!(
937					?key, ?path, provider = ?provider.name,
938					"Deleting media file from provider",
939				);
940
941				provider
942					.delete_one(&path)
943					.await
944					.log_debug_err()
945					.ok()
946			})
947			.count()
948			.map(|count| {
949				count
950					.ge(&0)
951					.into_option()
952					.ok_or_else(|| err!(Request(NotFound("Failed to remove on any provider."))))
953			})
954			.await
955	}
956
957	async fn create_media_file(&self, key: &[u8], file: &[u8]) -> Result {
958		self.storage_providers()
959			.try_stream()
960			.ready_try_filter(|provider| {
961				let store_media_on_providers = &self.services.config.store_media_on_providers;
962
963				store_media_on_providers.is_empty()
964					|| store_media_on_providers.contains(&provider.name)
965			})
966			.and_then(async |provider| {
967				let path = self.get_media_name_sha256(key);
968				debug!(
969					?key, ?path,
970					len = ?file.len(),
971					provider = ?provider.name,
972					"Creating media file on storage provider."
973				);
974
975				if let Err(e) = provider
976					.put_one(path.as_str(), file.to_vec())
977					.await
978				{
979					return Err!(Database(error!(
980						?path,
981						?provider,
982						"Failed to store media on provider: {e:?}"
983					)));
984				}
985
986				Ok(1)
987			})
988			.ready_try_fold(0_usize, |a, c| Ok(a.saturating_add(c)))
989			.inspect_ok(|&uploads| assert!(uploads > 0, "Successfully saved to nowhere."))
990			.map_ok(|_| ())
991			.await
992	}
993
994	fn storage_providers(&self) -> impl Iterator<Item = &Arc<Provider>> + Send + '_ {
995		let explicit_providers = &self.services.config.media_storage_providers;
996
997		let or_all_providers = explicit_providers
998			.is_empty()
999			.then(|| self.services.storage.providers())
1000			.into_iter()
1001			.flatten();
1002
1003		explicit_providers
1004			.iter()
1005			.filter_map(|id| self.services.storage.provider(id).ok())
1006			.chain(or_all_providers)
1007	}
1008
1009	#[inline]
1010	pub async fn get_metadata(&self, mxc: &Mxc<'_>) -> Option<Metadata> {
1011		self.db
1012			.search_file_metadata(mxc, &Dim::default(), Animate::Allowed)
1013			.await
1014			.ok()
1015	}
1016
1017	#[inline]
1018	#[must_use]
1019	pub fn get_media_path_sha256(&self, key: &[u8]) -> PathBuf {
1020		self.get_media_dir()
1021			.join(self.get_media_name_sha256(key))
1022	}
1023
1024	/// new SHA256 file name media function. requires database migrated. uses
1025	/// SHA256 hash of the base64 key as the file name
1026	#[inline]
1027	#[must_use]
1028	pub fn get_media_name_sha256(&self, key: &[u8]) -> String {
1029		// Using the hash of the base64 key as the filename prevents the total
1030		// length of the path from exceeding the maximum length in most
1031		// filesystems
1032		let digest = <sha2::Sha256 as sha2::Digest>::digest(key);
1033		encode_key(&digest)
1034	}
1035
1036	/// old base64 file name media function
1037	/// This is the old version of `get_media_path_sha256` that uses the full
1038	/// base64 key as the filename.
1039	#[must_use]
1040	pub fn get_media_path_b64(&self, key: &[u8]) -> PathBuf {
1041		self.get_media_dir().join(encode_key(key))
1042	}
1043
1044	#[must_use]
1045	pub fn get_media_dir(&self) -> PathBuf {
1046		self.services
1047			.server
1048			.config
1049			.database_path
1050			.join("media")
1051	}
1052}
1053
1054#[inline]
1055#[must_use]
1056pub fn encode_key(key: &[u8]) -> String { general_purpose::URL_SAFE_NO_PAD.encode(key) }
1057
1058fn mtime_millis(object: &ObjectMeta) -> u64 {
1059	u64::try_from(object.last_modified.timestamp_millis()).unwrap_or(0)
1060}