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#[derive(Debug)]
65pub struct Fetched {
66 pub media: Media,
67 pub(crate) animates: Option<bool>,
72}
73
74#[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#[derive(Clone, Debug)]
90pub struct UploadStat {
91 pub user_id: OwnedUserId,
92 pub media_length: u64,
93 pub created_ts: u64,
94}
95
96struct MXCState {
98 notifiers: Mutex<HashMap<OwnedMxcUri, Arc<Notify>>>,
100 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
118pub const MXC_LENGTH: usize = 32;
120
121pub const CACHE_CONTROL_IMMUTABLE: &str = "private,max-age=31536000,immutable";
123
124pub const CORP_CROSS_ORIGIN: &str = "cross-origin";
126
127const 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 #[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 let rate = f64::from(config.media_rc_create_per_second);
185 let burst = f64::from(config.media_rc_create_burst_count);
186
187 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 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 #[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 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 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 self.create_media_file(&key, file)
305 .map_ok(|()| walk.animates())
306 .await
307 }
308
309 #[tracing::instrument(level = "trace", skip(self))]
311 pub async fn delete(&self, mxc: &Mxc<'_>) -> Result {
312 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 #[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 #[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 #[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 #[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 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 #[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 let _lock = self.federation_mutex.lock(&key).await;
506
507 if self
509 .db
510 .search_file_metadata(mxc, &Dim::default(), Animate::Allowed)
511 .await
512 .is_ok()
513 {
514 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 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 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 #[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 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 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 #[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 #[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 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 #[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 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 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 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 #[inline]
1027 #[must_use]
1028 pub fn get_media_name_sha256(&self, key: &[u8]) -> String {
1029 let digest = <sha2::Sha256 as sha2::Digest>::digest(key);
1033 encode_key(&digest)
1034 }
1035
1036 #[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}