Skip to main content

tuwunel_service/media/thumbnail/
request.rs

1//! Routes thumbnail requests through stored and generated media.
2
3use std::{sync::Arc, time::Duration};
4
5use bytes::Bytes;
6use futures::{StreamExt, pin_mut};
7use ruma::{Mxc, UserId, http_headers::ContentDisposition};
8use tokio::{sync::Notify, time::timeout};
9use tuwunel_core::{
10	Err, Result, async_noinline, err, implement,
11	utils::{result::LogDebugErr, stream::IterStream},
12};
13
14#[cfg(feature = "media_thumbnail")]
15use super::generate::picture_dim;
16#[cfg(all(test, feature = "media_thumbnail"))]
17use super::tests::source_fetched;
18use super::{
19	super::{Media, data::Metadata},
20	Animate, Dim,
21};
22
23impl super::super::Service {
24	/// Uploads or replaces a file thumbnail.
25	///
26	/// Metadata is written first, then the supplied bytes replace the stored
27	/// media body for the requested dimensions.
28	#[tracing::instrument(
29		level = "debug",
30		ret(level = "debug")
31		skip(self, file),
32	)]
33	pub async fn upload_thumbnail(
34		&self,
35		mxc: &Mxc<'_>,
36		content_disposition: Option<&ContentDisposition>,
37		content_type: Option<&str>,
38		dim: &Dim,
39		file: &[u8],
40	) -> Result {
41		let key =
42			self.db
43				.create_file_metadata(mxc, None, dim, content_disposition, content_type)?;
44
45		// TODO: Remove dangling metadata when file creation fails.
46		self.create_media_file(&key, file).await?;
47
48		Ok(())
49	}
50
51	/// Answers a thumbnail request, fetching from the peer when it is remote.
52	///
53	/// The dimension a fetch asks the peer for is the dimension the answer is
54	/// filed under, and every later lookup seeks the normalized one, so the
55	/// three have to agree or nothing found on one request is found on the
56	/// next. Past every bucket there is no dimension to ask at, and the request
57	/// names the original file instead.
58	#[tracing::instrument(
59		level = "debug",
60		err(level = "debug")
61		skip(self),
62	)]
63	pub async fn get_or_fetch_thumbnail(
64		&self,
65		mxc: &Mxc<'_>,
66		dim: &Dim,
67		animate: Animate,
68		timeout_ms: Duration,
69		user: &UserId,
70	) -> Result<Media> {
71		if let Ok(media) = self
72			.get_thumbnail(mxc, dim, animate, Some(timeout_ms))
73			.await
74		{
75			return Ok(media);
76		}
77
78		if self
79			.services
80			.globals
81			.server_is_ours(mxc.server_name)
82		{
83			return Err!(Request(NotFound("Local thumbnail not found.")));
84		}
85
86		let lock = self.federation_mutex.lock(&mxc.to_string()).await;
87		let normalized = dim.normalized();
88
89		if self
90			.db
91			.file_metadata_exists(mxc, &normalized)
92			.await
93		{
94			drop(lock);
95			return self.get_thumbnail(mxc, dim, animate, None).await;
96		}
97
98		let fetched = match normalized.is_original() {
99			| true =>
100				self.fetch_remote_content(mxc, None, timeout_ms)
101					.await?,
102			| false =>
103				self.fetch_remote_thumbnail(mxc, None, timeout_ms, &normalized, animate)
104					.await?,
105		};
106
107		// a peer may ignore the parameter and this answers the caller directly,
108		// so the variant is repaired before it leaves rather than one request on
109		if animate.accepts_fetched(&fetched) {
110			return Ok(fetched.media);
111		}
112
113		self.store_still(mxc, &normalized, fetched.media)
114			.await
115	}
116
117	/// Downloads a thumbnail, waiting for a pending upload when requested.
118	///
119	/// The supplied duration bounds that wait. The future is boxed because the
120	/// still-repair path pulls the thumbnailer
121	/// into it, and inlining that into every caller overflows the layout depth
122	/// limit in the federation handler.
123	#[async_noinline]
124	#[tracing::instrument(
125		level = "debug",
126		err(level = "debug")
127		skip(self),
128	)]
129	pub async fn get_thumbnail<'a>(
130		&'a self,
131		mxc: &'a Mxc<'_>,
132		dim: &'a Dim,
133		animate: Animate,
134		timeout_duration: Option<Duration>,
135	) -> Result<Media> {
136		if let Ok(media) = self.get_stored_thumbnail(mxc, dim, animate).await {
137			return Ok(media);
138		}
139
140		let Some(timeout_duration) = timeout_duration else {
141			return Err!(Request(NotFound("Media thumbnail not found.")));
142		};
143
144		let Ok(_pending) = self.db.search_pending_mxc(mxc).await else {
145			return Err!(Request(NotFound("Media thumbnail not found.")));
146		};
147
148		let notifier = self
149			.mxc_state
150			.notifiers
151			.lock()?
152			.entry(mxc.to_string().into())
153			.or_insert_with(|| Arc::new(Notify::new()))
154			.clone();
155
156		if timeout(timeout_duration, notifier.notified())
157			.await
158			.is_err()
159		{
160			return Err!(Request(NotYetUploaded("Media has not been uploaded yet.")));
161		}
162
163		self.get_stored_thumbnail(mxc, dim, animate).await
164	}
165
166	/// Downloads a stored or generated thumbnail.
167	///
168	/// Requests normalize to a bounded storage bucket. An existing variant is
169	/// returned directly; a missing variant is generated from the original or
170	/// answered from promoted storage when the original row is absent.
171	#[tracing::instrument(
172		name = "thumbnail",
173		level = "debug",
174		err(level = "trace")
175		skip(self),
176	)]
177	pub async fn get_stored_thumbnail(
178		&self,
179		mxc: &Mxc<'_>,
180		dim: &Dim,
181		animate: Animate,
182	) -> Result<Media> {
183		let dim = dim.normalized();
184
185		// the sentinel is the key the original is stored under rather than a
186		// size, so a request reaching it is answered from the original itself
187		if dim.is_original() {
188			return self.answer_original(mxc, animate).await;
189		}
190
191		if let Ok(metadata) = self
192			.db
193			.search_file_metadata(mxc, &dim, animate)
194			.await
195		{
196			return self
197				.answer_stored(mxc, &dim, animate, metadata)
198				.await;
199		}
200
201		let Ok(metadata) = self.original_metadata(mxc).await else {
202			return self.answer_promoted(mxc, animate).await;
203		};
204
205		self.get_thumbnail_generate(mxc, &dim, animate, metadata)
206			.await
207	}
208}
209
210/// Answers a request past every bucket, which names the original file.
211///
212/// The original's own row is keyed at the sentinel, where there is no size to
213/// re-encode at, so a picture the request will not accept stands in at its own
214/// dimensions. Those are the largest a still may carry without upscaling, so
215/// the stand-in covers any request the original itself covers and is the best
216/// available where it does not.
217///
218/// The future is not inlined: this branch reaches the thumbnailer, and folding
219/// that into every caller overflows the layout depth limit.
220#[cfg(feature = "media_thumbnail")]
221#[implement(super::super::Service)]
222#[async_noinline]
223#[tracing::instrument(name = "original", level = "debug", skip(self))]
224async fn answer_original<'a>(&'a self, mxc: &'a Mxc<'_>, animate: Animate) -> Result<Media> {
225	let Ok(data) = self.original_metadata(mxc).await else {
226		return self.answer_promoted(mxc, animate).await;
227	};
228
229	let bytes = self.fetch_bytes(&data.key).await?;
230
231	if animate.accepts_picture(&bytes) {
232		return Ok(into_media(data, bytes.into()));
233	}
234
235	let dim = picture_dim(&bytes);
236
237	if let Ok(metadata) = self
238		.db
239		.search_file_metadata(mxc, &dim, animate)
240		.await
241	{
242		drop((bytes, data));
243
244		return self
245			.answer_stored(mxc, &dim, animate, metadata)
246			.await;
247	}
248
249	// no still can be derived from a picture that will not decode, and every
250	// bucket answers that with the original rather than refusing
251	let Ok(image) = self.decode(&bytes) else {
252		return match animate.accepts_fallback(&bytes) {
253			| true => Ok(into_media(data, bytes.into())),
254			| false => Err!(Request(NotFound("Media thumbnail not found."))),
255		};
256	};
257
258	drop((bytes, data));
259
260	self.store_thumbnail(mxc, &dim, image).await
261}
262
263/// Hands the original back, there being no thumbnailer to stand a still in.
264///
265/// A build without the feature keeps serving what it holds rather than
266/// refusing, so a request that forbade animation goes unhonored here.
267#[cfg(not(feature = "media_thumbnail"))]
268#[implement(super::super::Service)]
269#[tracing::instrument(name = "original", level = "debug", skip_all)]
270async fn answer_original(&self, mxc: &Mxc<'_>, animate: Animate) -> Result<Media> {
271	let Ok(data) = self.original_metadata(mxc).await else {
272		return self.answer_promoted(mxc, animate).await;
273	};
274
275	self.get_thumbnail_saved(data).await
276}
277
278/// Metadata for the original file, which every thumbnail is derived from.
279///
280/// Its row is keyed at the sentinel dimension, and nothing is withheld from
281/// this lookup: the row is the original rather than a variant of it.
282#[implement(super::super::Service)]
283pub(super) async fn original_metadata(&self, mxc: &Mxc<'_>) -> Result<Metadata> {
284	self.db
285		.search_file_metadata(mxc, &Dim::default(), Animate::Allowed)
286		.await
287}
288
289/// Answers from stored media when no metadata row names the original.
290///
291/// The original may be lazy preview media promoted on this very request, which
292/// leaves no row behind, and only a picture is worth serving in a thumbnail's
293/// place. No row having named it, the request's own preference is the only gate
294/// it passes, so the picture is read here as it is anywhere else.
295#[implement(super::super::Service)]
296#[tracing::instrument(name = "promoted", level = "debug", skip(self))]
297async fn answer_promoted(&self, mxc: &Mxc<'_>, animate: Animate) -> Result<Media> {
298	let media = self.get_stored(mxc).await?;
299	let servable = media
300		.content_type
301		.as_deref()
302		.is_some_and(|content_type| content_type.starts_with("image/"))
303		&& animate.accepts_picture(&media.content);
304
305	servable
306		.then_some(media)
307		.ok_or_else(|| err!(Request(NotFound("Media not found."))))
308}
309
310/// The stored bytes for a thumbnail row, from the first provider holding them.
311///
312/// Returning the shared `Bytes` spares a caller that only reads the picture
313/// the `Vec` copy `Media` would force on it.
314#[implement(super::super::Service)]
315#[tracing::instrument(name = "fetch", level = "debug", skip_all)]
316pub(super) async fn fetch_bytes(&self, key: &[u8]) -> Result<Bytes> {
317	#[cfg(all(test, feature = "media_thumbnail"))]
318	source_fetched();
319
320	let path = self.get_media_name_sha256(key);
321	let fetch = self
322		.storage_providers()
323		.stream()
324		.filter_map(async |provider| {
325			provider
326				.get(path.as_str())
327				.await
328				.log_debug_err()
329				.ok()
330		});
331
332	pin_mut!(fetch);
333
334	fetch
335		.next()
336		.await
337		.ok_or_else(|| err!(Request(NotFound("Media not found."))))
338}
339/// Answers a request from a stored row, re-deriving a still if it animates.
340///
341/// The type a row is stored under is whatever produced it claimed, and a peer
342/// is free to claim wrongly, so the picture itself decides once it is in hand.
343/// A row that may not answer is re-encoded and the still left behind for the
344/// next request.
345#[cfg(feature = "media_thumbnail")]
346#[implement(super::super::Service)]
347#[tracing::instrument(name = "stored", level = "debug", skip(self, data))]
348async fn answer_stored(
349	&self,
350	mxc: &Mxc<'_>,
351	dim: &Dim,
352	animate: Animate,
353	data: Metadata,
354) -> Result<Media> {
355	let bytes = self.fetch_bytes(&data.key).await?;
356
357	if animate.accepts_picture(&bytes) {
358		return Ok(into_media(data, bytes.into()));
359	}
360
361	let image = self.decode_still(&bytes)?;
362
363	drop((bytes, data));
364
365	self.store_thumbnail(mxc, dim, image).await
366}
367
368/// Hands the stored row back whatever its picture holds.
369///
370/// Judging it would only be worth doing if a still could then be derived, and
371/// a build without the feature has no thumbnailer to derive one with.
372#[cfg(not(feature = "media_thumbnail"))]
373#[implement(super::super::Service)]
374#[tracing::instrument(name = "stored", level = "debug", skip_all)]
375async fn answer_stored(
376	&self,
377	_mxc: &Mxc<'_>,
378	_dim: &Dim,
379	_animate: Animate,
380	data: Metadata,
381) -> Result<Media> {
382	self.get_thumbnail_saved(data).await
383}
384
385/// Hands the picture back as it stands, there being no thumbnailer.
386///
387/// A build without the feature keeps serving what it holds rather than
388/// refusing, so a request for a still goes unhonored instead of failing.
389#[cfg(not(feature = "media_thumbnail"))]
390#[implement(super::super::Service)]
391#[tracing::instrument(name = "still", level = "debug", skip_all)]
392pub(in super::super) async fn store_still(
393	&self,
394	_mxc: &Mxc<'_>,
395	_dim: &Dim,
396	animated: Media,
397) -> Result<Media> {
398	Ok(animated)
399}
400
401/// Hands the original back in place of the thumbnail it cannot generate.
402///
403/// The row this is given is the original's own, so the caller is answered with
404/// media the server holds rather than being refused outright.
405#[cfg(not(feature = "media_thumbnail"))]
406#[implement(super::super::Service)]
407#[tracing::instrument(name = "fallback", level = "debug", skip_all)]
408pub(super) async fn get_thumbnail_generate(
409	&self,
410	_mxc: &Mxc<'_>,
411	_dim: &Dim,
412	_animate: Animate,
413	data: Metadata,
414) -> Result<Media> {
415	self.get_thumbnail_saved(data).await
416}
417
418/// Answers a thumbnail row from the picture already in storage.
419///
420/// The response hands its body over by value, so the bytes are converted
421/// rather than read in place; the conversion reclaims the storage buffer
422/// whenever this holds the only handle to it.
423#[cfg(not(feature = "media_thumbnail"))]
424#[implement(super::super::Service)]
425#[tracing::instrument(name = "saved", level = "debug", skip_all)]
426pub(super) async fn get_thumbnail_saved(&self, data: Metadata) -> Result<Media> {
427	let bytes = self.fetch_bytes(&data.key).await?;
428
429	Ok(into_media(data, bytes.into()))
430}
431/// Transfers stored metadata and content into a media response.
432///
433/// Content type and disposition move out of the metadata row without cloning.
434pub(super) fn into_media(data: Metadata, content: Vec<u8>) -> Media {
435	Media {
436		content,
437		content_type: data.content_type,
438		content_disposition: data.content_disposition,
439	}
440}