1use std::sync::Arc;
2
3use futures::{Stream, StreamExt, pin_mut};
4use ruma::{Mxc, OwnedMxcUri, OwnedUserId, UserId, http_headers::ContentDisposition};
5use serde::Deserialize;
6#[cfg(feature = "url_preview")]
7use serde::Serialize;
8use tuwunel_core::{
9 Err, Result, at, debug, debug_info, err,
10 utils::{
11 ReadyExt, str_from_bytes,
12 stream::{TryExpect, TryIgnore},
13 string_from_bytes,
14 },
15};
16use tuwunel_database::{Cbor, Database, Deserialized, Ignore, Interfix, Map, Txn, serialize_key};
17
18use super::{
19 Media,
20 preview::CachedPreview,
21 thumbnail::{Animate, Dim},
22};
23
24pub(crate) struct Data {
25 db: Arc<Database>,
26 mediaid_file: Arc<Map>,
27 mediaid_lazy: Arc<Map>,
28 mediaid_lazycontent: Arc<Map>,
29 mediaid_pending: Arc<Map>,
30 mediaid_user: Arc<Map>,
31 url_preview: Arc<Map>,
32}
33
34#[derive(Debug)]
35pub struct Metadata {
36 pub content_disposition: Option<ContentDisposition>,
37 pub content_type: Option<String>,
38 pub(super) key: Vec<u8>,
39}
40
41#[cfg(feature = "url_preview")]
43#[derive(Serialize)]
44struct LazyContentRef<'a> {
45 content_type: Option<&'a str>,
46 content_disposition: Option<&'a str>,
47 #[serde(with = "serde_bytes")]
48 content: &'a [u8],
49}
50
51#[derive(Deserialize)]
54struct LazyContent {
55 content_type: Option<String>,
56 content_disposition: Option<String>,
57 #[serde(with = "serde_bytes")]
58 content: Vec<u8>,
59}
60
61impl From<LazyContent> for Media {
62 fn from(lazy: LazyContent) -> Self {
63 let content_disposition = lazy
64 .content_disposition
65 .and_then(|disposition| disposition.parse().ok());
66
67 Self {
68 content: lazy.content,
69 content_type: lazy.content_type,
70 content_disposition,
71 }
72 }
73}
74
75impl Data {
76 pub(super) fn new(db: &Arc<Database>) -> Self {
77 Self {
78 db: db.clone(),
79 mediaid_file: db["mediaid_file"].clone(),
80 mediaid_lazy: db["mediaid_lazy"].clone(),
81 mediaid_lazycontent: db["mediaid_lazycontent"].clone(),
82 mediaid_pending: db["mediaid_pending"].clone(),
83 mediaid_user: db["mediaid_user"].clone(),
84 url_preview: db["url_preview"].clone(),
85 }
86 }
87
88 pub(super) fn create_file_metadata(
89 &self,
90 mxc: &Mxc<'_>,
91 user: Option<&UserId>,
92 dim: &Dim,
93 content_disposition: Option<&ContentDisposition>,
94 content_type: Option<&str>,
95 ) -> Result<Vec<u8>> {
96 let dim: &[u32] = &[dim.width, dim.height];
97 let key = (mxc, dim, content_disposition, content_type);
98 let key = serialize_key(key)?;
99 let mut txn = self.db.txn();
100
101 txn.insert_raw(&self.mediaid_file, &key, []);
102 if let Some(user) = user {
103 let key = (mxc, user);
104
105 txn.put_raw(&self.mediaid_user, key, user);
106 }
107
108 txn.execute();
109
110 Ok(key.to_vec())
111 }
112
113 pub(super) fn insert_pending_mxc(
115 &self,
116 mxc: &Mxc<'_>,
117 user: &UserId,
118 unused_expires_at: u64,
119 ) {
120 let value = (unused_expires_at, user);
121 debug!(?mxc, ?user, ?unused_expires_at, "Inserting pending");
122
123 self.mediaid_pending
124 .raw_put(mxc.to_string(), value);
125 }
126
127 pub(super) fn remove_pending_mxc(&self, mxc: &Mxc<'_>) {
129 self.mediaid_pending.remove(&mxc.to_string());
130 }
131
132 pub(super) async fn count_pending_mxc_for_user(&self, user_id: &UserId) -> (usize, u64) {
134 type KeyVal<'a> = (Ignore, (u64, &'a UserId));
135
136 self.mediaid_pending
137 .stream()
138 .expect_ok()
139 .ready_filter(|(_, (_, pending_user_id)): &KeyVal<'_>| user_id == *pending_user_id)
140 .ready_fold(
141 (0_usize, u64::MAX),
142 |(count, earliest_expiration), (_, (expires_at, _))| {
143 (count.saturating_add(1), earliest_expiration.min(expires_at))
144 },
145 )
146 .await
147 }
148
149 pub(super) async fn search_pending_mxc(&self, mxc: &Mxc<'_>) -> Result<(OwnedUserId, u64)> {
151 type Value<'a> = (u64, OwnedUserId);
152
153 self.mediaid_pending
154 .get(&mxc.to_string())
155 .await
156 .deserialized()
157 .map(|(expires_at, user_id): Value<'_>| (user_id, expires_at))
158 .inspect(|(user_id, expires_at)| debug!(?mxc, ?user_id, ?expires_at, "Found pending"))
159 .map_err(|e| err!(Request(NotFound("Pending not found or error: {e}"))))
160 }
161
162 #[cfg(feature = "url_preview")]
165 pub(super) fn insert_lazy_media(&self, mxc: &str, url: &str) {
166 debug!(?mxc, ?url, "Registering lazy media");
167
168 self.mediaid_lazy.insert(mxc, url.as_bytes());
169 }
170
171 #[cfg(feature = "url_preview")]
172 pub(super) fn queue_lazy_media(&self, txn: &mut Txn, mxc: &str, url: &str) {
173 debug!(?mxc, ?url, "Registering lazy media");
174
175 txn.insert_raw(&self.mediaid_lazy, mxc, url.as_bytes());
176 }
177
178 pub(super) fn remove_lazy_media(&self, txn: &mut Txn, mxc: &str) {
181 txn.del_raw(&self.mediaid_lazy, mxc);
182 }
183
184 pub(super) async fn search_lazy_media(&self, mxc: &Mxc<'_>) -> Result<String> {
186 let handle = self.mediaid_lazy.get(&mxc.to_string()).await?;
187
188 string_from_bytes(&handle)
189 .map_err(|e| err!(Database(error!(?mxc, "Lazy media URL is invalid: {e}"))))
190 }
191
192 #[cfg(feature = "url_preview")]
195 pub(super) fn set_lazy_content(
196 &self,
197 txn: &mut Txn,
198 mxc: &str,
199 content_type: Option<&str>,
200 content_disposition: Option<&str>,
201 content: &[u8],
202 ) {
203 let value = LazyContentRef {
204 content_type,
205 content_disposition,
206 content,
207 };
208
209 txn.raw_put(&self.mediaid_lazycontent, mxc, Cbor(&value));
210 }
211
212 pub(super) async fn get_lazy_content(&self, mxc: &str) -> Result<Media> {
214 self.mediaid_lazycontent
215 .get(mxc)
216 .await
217 .deserialized::<Cbor<LazyContent>>()
218 .map(at!(0))
219 .map(Into::into)
220 }
221
222 pub(super) fn remove_lazy_content(&self, txn: &mut Txn, mxc: &str) {
223 txn.del_raw(&self.mediaid_lazycontent, mxc);
224 }
225
226 pub(super) async fn delete_file_mxc(&self, mxc: &Mxc<'_>) {
227 debug!("MXC URI: {mxc}");
228
229 let prefix = (mxc, Interfix);
230 let txn = self
231 .mediaid_file
232 .keys_prefix_raw(&prefix)
233 .ignore_err()
234 .ready_fold(self.db.txn(), |mut txn, key| {
235 txn.del_raw(&self.mediaid_file, key);
236
237 txn
238 })
239 .await;
240
241 let txn = self
242 .mediaid_user
243 .stream_prefix_raw(&prefix)
244 .ignore_err()
245 .ready_fold(txn, |mut txn, (key, val)| {
246 debug_assert!(
247 key.starts_with(mxc.to_string().as_bytes()),
248 "key should start with the mxc"
249 );
250
251 let user = str_from_bytes(val).unwrap_or_default();
252 debug_info!("Deleting key {key:?} which was uploaded by user {user}");
253
254 txn.del_raw(&self.mediaid_user, key);
255
256 txn
257 })
258 .await;
259
260 txn.execute();
261 }
262
263 pub(super) async fn search_mxc_metadata_prefix(&self, mxc: &Mxc<'_>) -> Result<Vec<Vec<u8>>> {
265 debug!("MXC URI: {mxc}");
266
267 let prefix = (mxc, Interfix);
268 let keys: Vec<Vec<u8>> = self
269 .mediaid_file
270 .keys_prefix_raw(&prefix)
271 .ignore_err()
272 .map(<[u8]>::to_vec)
273 .collect()
274 .await;
275
276 if keys.is_empty() {
277 return Err!(Database("Failed to find any keys in database for `{mxc}`",));
278 }
279
280 debug!("Got the following keys: {keys:?}");
281
282 Ok(keys)
283 }
284
285 pub(super) async fn file_metadata_exists(&self, mxc: &Mxc<'_>, dim: &Dim) -> bool {
286 let dim: &[u32] = &[dim.width, dim.height];
287 let prefix = (mxc, dim, Interfix);
288 let keys = self
289 .mediaid_file
290 .keys_prefix_raw(&prefix)
291 .ignore_err();
292
293 pin_mut!(keys);
294 keys.next().await.is_some()
295 }
296
297 pub(super) async fn search_file_metadata(
304 &self,
305 mxc: &Mxc<'_>,
306 dim: &Dim,
307 animate: Animate,
308 ) -> Result<Metadata> {
309 let dim: &[u32] = &[dim.width, dim.height];
310 let prefix = (mxc, dim, Interfix);
311
312 let (wanted, other) = self
313 .mediaid_file
314 .keys_prefix_raw(&prefix)
315 .ignore_err()
316 .ready_fold((None, None), |(wanted, other), key: &[u8]| {
317 match animate.prefers_type(key_content_type(key)) {
320 | true if wanted.is_none() => (Some(key.to_owned()), other),
321 | false if other.is_none() => (wanted, Some(key.to_owned())),
322 | _ => (wanted, other),
323 }
324 })
325 .await;
326
327 let key = wanted
328 .or(other)
329 .ok_or_else(|| err!(Request(NotFound("Media not found"))))?;
330
331 let (content_type, content_disposition) = {
333 let mut parts = key_parts(&key);
334
335 let content_type = parts
336 .next()
337 .map(string_from_bytes)
338 .transpose()
339 .map_err(|e| err!(Database(error!(?mxc, "Content-type is invalid: {e}"))))?;
340
341 let content_disposition = parts
342 .next()
343 .map(Some)
344 .ok_or_else(|| err!(Database(error!(?mxc, "Media ID in db is invalid."))))?
345 .filter(|bytes| !bytes.is_empty())
346 .map(string_from_bytes)
347 .transpose()
348 .map_err(|e| err!(Database(error!(?mxc, "Content-disposition is invalid: {e}"))))?
349 .as_deref()
350 .map(str::parse)
351 .transpose()
352 .map_err(|e| {
353 err!(Database(error!(?mxc, "Content-disposition is invalid: {e}")))
354 })?;
355
356 (content_type, content_disposition)
357 };
358
359 Ok(Metadata { content_disposition, content_type, key })
360 }
361
362 pub(super) async fn mxc_user(&self, mxc: &Mxc<'_>) -> Option<OwnedUserId> {
365 let prefix = (mxc, Interfix);
366 let users = self
367 .mediaid_user
368 .stream_prefix(&prefix)
369 .ignore_err()
370 .map(|(_, user): (Ignore, &UserId)| user.to_owned());
371
372 pin_mut!(users);
373 users.next().await
374 }
375
376 pub(super) async fn get_all_user_mxcs(&self, user_id: &UserId) -> Vec<OwnedMxcUri> {
378 self.mediaid_user
379 .stream()
380 .ignore_err()
381 .ready_filter_map(|((key, _), user): ((&str, Ignore), &UserId)| {
382 (user == user_id).then(|| key.into())
383 })
384 .collect()
385 .await
386 }
387
388 pub(crate) async fn get_all_media_keys(&self) -> Vec<Vec<u8>> {
391 self.mediaid_file
392 .raw_keys()
393 .ignore_err()
394 .map(<[u8]>::to_vec)
395 .collect()
396 .await
397 }
398
399 pub(super) fn set_url_preview(&self, url: &str, cached: &CachedPreview) -> Result {
400 self.url_preview.raw_put(url, Cbor(cached));
401
402 Ok(())
403 }
404
405 pub(super) async fn get_url_preview(&self, url: &str) -> Result<CachedPreview> {
406 self.url_preview
407 .get(url)
408 .await
409 .deserialized::<Cbor<_>>()
410 .map(at!(0))
411 .ok()
412 .filter(CachedPreview::valid)
413 .ok_or(err!(Request(NotFound("Expired from cache"))))
414 }
415
416 pub(super) fn all_uploads(
418 &self,
419 ) -> impl Stream<Item = (OwnedMxcUri, OwnedUserId)> + Send + '_ {
420 self.mediaid_user
421 .keys()
422 .ignore_err()
423 .map(|(mxc, user): (&str, &UserId)| (mxc.into(), user.to_owned()))
424 }
425}
426
427fn key_parts(key: &[u8]) -> impl Iterator<Item = &[u8]> { key.rsplit(|&b| b == 0xFF) }
433
434fn key_content_type(key: &[u8]) -> Option<&str> {
435 key_parts(key)
436 .next()
437 .map(str_from_bytes)
438 .and_then(Result::ok)
439}
440
441#[cfg(feature = "url_preview")]
442#[cfg(test)]
443mod tests {
444 use minicbor_serde::{from_slice, to_vec};
445
446 use super::{LazyContent, LazyContentRef, Media};
447
448 #[test]
449 fn lazy_content_roundtrip() {
450 let content: &[u8] = b"\x00\x01\xFF\xFE arbitrary staged bytes";
451 let value = LazyContentRef {
452 content_type: Some("image/png"),
453 content_disposition: Some("inline; filename=\"cat.png\""),
454 content,
455 };
456
457 let bytes = to_vec(&value).expect("encodes");
458 let decoded: LazyContent = from_slice(&bytes).expect("decodes");
459
460 assert_eq!(decoded.content_type.as_deref(), Some("image/png"));
461 assert_eq!(decoded.content.as_slice(), content);
462
463 let media = Media::from(decoded);
464 assert_eq!(media.content.as_slice(), content);
465 assert!(media.content_disposition.is_some(), "disposition re-parses to the ruma type");
466 }
467
468 #[test]
469 fn lazy_content_bytes_compact() {
470 let content = vec![0xAB_u8; 4096];
471 let value = LazyContentRef {
472 content_type: None,
473 content_disposition: None,
474 content: content.as_slice(),
475 };
476
477 let bytes = to_vec(&value).expect("encodes");
478
479 assert!(bytes.len() <= content.len() + 64);
482 }
483}