Skip to main content

tuwunel_service/media/
data.rs

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/// Borrowed staging-cache value: written zero-copy from the measured bytes.
42#[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/// Owned staging-cache value read back at promotion. `ContentDisposition` is
52/// Serialize-only, so the disposition rides as its header string.
53#[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	/// Insert a pending MXC URI into the database
114	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	/// Remove a pending MXC URI from the database
128	pub(super) fn remove_pending_mxc(&self, mxc: &Mxc<'_>) {
129		self.mediaid_pending.remove(&mxc.to_string());
130	}
131
132	/// Count the number of pending MXC URIs for a specific user
133	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	/// Search for a pending MXC URI in the database
150	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	/// Map a minted mxc:// URI to the external URL it resolves to on first
163	/// download (see `Service::fetch_lazy_media`).
164	#[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	/// Remove a lazy media reference by its mxc:// URI string, unregistering
179	/// the mxc.
180	pub(super) fn remove_lazy_media(&self, txn: &mut Txn, mxc: &str) {
181		txn.del_raw(&self.mediaid_lazy, mxc);
182	}
183
184	/// Look up the external URL a lazy media MXC URI refers to.
185	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	/// Stage the measured preview media bytes under its minted mxc so the
193	/// first client download can promote without touching the origin.
194	#[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	/// Take the staged bytes a preview seeded for a lazy media mxc, if any.
213	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	/// Searches for all files with the given MXC
264	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	/// Metadata for the row at these dimensions that best answers the request.
298	///
299	/// A row of the variant the request asked for wins when the prefix holds
300	/// one, and otherwise any row does, so the caller can repair a variant it
301	/// cannot serve rather than refuse media the server holds. One cursor walk
302	/// harvests both, since a second seek would cost every fallback a scan.
303	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				// owning only what is retained; the rest of the prefix is walked
318				// past and never copied
319				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		// the borrow the parse takes ends before the key is moved into the result
332		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	/// Uploading local user of the media at the given MXC, from the uploader
363	/// index.
364	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	/// Gets all the MXCs associated with a user
377	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	/// Gets all the media keys in our database (this includes all the metadata
389	/// associated with it such as width, height, content-type, etc)
390	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	/// Streams every (mxc, uploader) pair in the user-media index.
417	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
427/// Components of a `mediaid_file` key, innermost first.
428///
429/// The key is joined on `0xFF` with the content type last and the content
430/// disposition before it, so splitting from the right yields the two the
431/// caller wants without walking the mxc and dimensions ahead of them.
432fn 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		// serde_bytes must encode a CBOR byte string, not an array-of-uints
480		// (~1.9x); only a small fixed header of overhead is permitted
481		assert!(bytes.len() <= content.len() + 64);
482	}
483}