Skip to main content

moq_archive/
store.rs

1use std::num::NonZeroUsize;
2use std::ops::RangeInclusive;
3
4use bytes::Bytes;
5use futures::StreamExt;
6use futures::stream::BoxStream;
7use object_store::list::{PaginatedListOptions, PaginatedListResult, PaginatedListStore};
8use object_store::path::Path;
9use object_store::{ListResult, ObjectMeta, ObjectStore, ObjectStoreExt, PutMode, PutPayload};
10
11use crate::info::Info;
12use crate::path::Key;
13use crate::segment::Object;
14use crate::{Error, Result};
15
16/// Recording object listing types.
17///
18/// Listing scopes and offsets are relative to the recording prefix supplied to
19/// [`super::Store::new`]. The same [`Query`] can drive streaming and paginated
20/// listing. A page's [`Page::next`] query retains the opaque backend token and
21/// must be passed to [`super::Store::list_paginated`], not [`super::Store::list`].
22///
23/// ```
24/// use std::num::NonZeroUsize;
25///
26/// use moq_archive::store::list::Query;
27///
28/// let query = Query::segments("timeline.z")?
29///     .page_size(NonZeroUsize::new(100).unwrap());
30/// # Ok::<(), moq_archive::Error>(())
31/// ```
32pub mod list {
33	use std::num::NonZeroUsize;
34
35	use object_store::path::Path;
36
37	use crate::path;
38	use crate::{Key, Result};
39
40	/// A recording-relative object listing query.
41	#[derive(Debug, Clone, Default, PartialEq, Eq)]
42	pub struct Query {
43		pub(crate) prefix: Option<Path>,
44		pub(crate) offset: Option<Path>,
45		pub(crate) max_keys: Option<NonZeroUsize>,
46		pub(crate) page_token: Option<String>,
47	}
48
49	impl Query {
50		/// List every object in the recording.
51		pub fn new() -> Self {
52			Self::default()
53		}
54
55		fn prefix(mut self, prefix: impl Into<Path>) -> Self {
56			self.prefix = Some(prefix.into());
57			self.page_token = None;
58			self
59		}
60
61		fn offset(mut self, offset: impl Into<Path>) -> Self {
62			self.offset = Some(offset.into());
63			self.page_token = None;
64			self
65		}
66
67		/// Request at most this many entries per page.
68		pub fn page_size(mut self, max_keys: NonZeroUsize) -> Self {
69			self.max_keys = Some(max_keys);
70			self.page_token = None;
71			self
72		}
73
74		/// List every object for one track.
75		pub fn track(track: &str) -> Result<Self> {
76			Ok(Self::new().prefix(path::track_prefix(track)?))
77		}
78
79		/// List group objects for one track.
80		pub fn groups(track: &str) -> Result<Self> {
81			Ok(Self::new().prefix(path::groups_prefix(&Path::ROOT, track)?))
82		}
83
84		/// List timeline segment objects for one track.
85		pub fn segments(track: &str) -> Result<Self> {
86			Ok(Self::new().prefix(path::segments_prefix(&Path::ROOT, track)?))
87		}
88
89		/// Start strictly after this recording object.
90		pub fn after(self, key: &Key) -> Result<Self> {
91			Ok(self.offset(key.path(&Path::ROOT)?))
92		}
93
94		/// List group objects beginning at the exclusive lexical offset for `group`.
95		pub fn groups_from(track: &str, group: u64) -> Result<Self> {
96			Ok(Self::groups(track)?.offset(path::groups_offset(&Path::ROOT, track, group)?))
97		}
98
99		pub(crate) fn next(&self, page_token: String) -> Self {
100			let mut next = self.clone();
101			next.page_token = Some(page_token);
102			next
103		}
104	}
105
106	/// One recording object returned by a listing.
107	#[derive(Debug, Clone, PartialEq, Eq)]
108	pub struct Entry {
109		/// Parsed recording key.
110		pub key: Key,
111		/// Object size in bytes.
112		pub size: u64,
113	}
114
115	/// One page of recording objects.
116	#[derive(Debug, Clone)]
117	pub struct Page {
118		/// Entries in the backend's order.
119		pub entries: Vec<Entry>,
120		/// The same query advanced to the next page, if any.
121		/// Pass it to `Store::list_paginated`; `Store::list` rejects it.
122		pub next: Option<Query>,
123	}
124}
125
126use list::{Entry, Page, Query};
127
128/// Versioned recording objects on a generic [`ObjectStore`].
129#[derive(Clone, Debug)]
130pub struct Store<T> {
131	inner: T,
132	prefix: Path,
133}
134
135impl<T: ObjectStore> Store<T> {
136	/// Wrap `inner` with an application-defined recording prefix.
137	pub fn new(inner: T, prefix: impl Into<Path>) -> Self {
138		Self {
139			inner,
140			prefix: prefix.into(),
141		}
142	}
143
144	/// The underlying object store.
145	pub fn inner(&self) -> &T {
146		&self.inner
147	}
148
149	/// The recording prefix.
150	pub fn prefix(&self) -> &Path {
151		&self.prefix
152	}
153
154	/// Encode `key` under this store's prefix.
155	pub fn path(&self, key: &Key) -> Result<Path> {
156		key.path(&self.prefix)
157	}
158
159	/// Create `.info`, or accept an existing object with the same parsed properties.
160	pub async fn put_info(&self, track: &str, info: &Info) -> Result<Key> {
161		let key = Key::info(track)?;
162		let path = self.path(&key)?;
163		let bytes = info.encode()?;
164		match self.create(&path, bytes).await? {
165			Create::Created => Ok(key),
166			Create::Exists => {
167				let existing = self.get_bytes(&path).await?;
168				let parsed = Info::decode(&existing)?;
169				if parsed.priority != info.priority {
170					return Err(Error::Priority {
171						existing: parsed.priority,
172						intended: info.priority,
173					});
174				}
175				if parsed.timescale != info.timescale {
176					return Err(Error::TimescaleMismatch {
177						existing: parsed.timescale,
178						intended: info.timescale,
179					});
180				}
181				Ok(key)
182			}
183		}
184	}
185
186	/// Fetch and validate a track's `.info`.
187	pub async fn get_info(&self, track: &str) -> Result<Info> {
188		let path = self.path(&Key::info(track)?)?;
189		Info::decode(&self.get_bytes(&path).await?)
190	}
191
192	/// Create a range-named groups object. A collision is accepted only when the bytes match.
193	pub async fn put_groups(&self, track: &str, object: &Object) -> Result<Key> {
194		let key = Key::groups(track, object.bounds()?)?;
195		self.put_segment(&key, object.encode()?).await?;
196		Ok(key)
197	}
198
199	/// Fetch a groups object and require its table to match the filename bounds.
200	pub async fn get_groups(&self, track: &str, range: RangeInclusive<u64>) -> Result<Object> {
201		let path = self.path(&Key::groups(track, range.clone())?)?;
202		Object::decode_groups(self.get_bytes(&path).await?, range)
203	}
204
205	/// Create a timeline object at `segments/<segment>`. A collision is accepted only when the bytes match.
206	pub async fn put_segments(&self, track: &str, segment: u64, object: &Object) -> Result<Key> {
207		let key = Key::segments(track, segment)?;
208		self.put_segment(&key, object.encode()?).await?;
209		Ok(key)
210	}
211
212	/// Fetch and validate a timeline object.
213	pub async fn get_segments(&self, track: &str, segment: u64) -> Result<Object> {
214		let path = self.path(&Key::segments(track, segment)?)?;
215		Object::decode(self.get_bytes(&path).await?)
216	}
217
218	/// Delete the object at `key`.
219	pub async fn delete(&self, key: &Key) -> Result<()> {
220		let path = self.path(key)?;
221		self.inner.delete(&path).await?;
222		Ok(())
223	}
224
225	/// Stream matching entries. Order is unspecified and `max_keys` is ignored.
226	/// Continuation queries from `Page::next` fail; pass them to `list_paginated`.
227	pub fn list(&self, query: &Query) -> BoxStream<'static, Result<Entry>> {
228		if query.page_token.is_some() {
229			return futures::stream::once(async { Err(Error::Pagination) }).boxed();
230		}
231		let prefix = self.list_path(query.prefix.as_ref());
232		let offset = query.offset.as_ref().map(|offset| self.recording_path(offset));
233		let store_prefix = self.prefix.clone();
234		let stream = match offset.as_ref() {
235			Some(offset) => self.inner.list_with_offset(prefix.as_ref(), offset),
236			None => self.inner.list(prefix.as_ref()),
237		};
238		stream.map(move |item| Entry::from_meta(&store_prefix, item?)).boxed()
239	}
240
241	async fn put_segment(&self, key: &Key, bytes: Bytes) -> Result<()> {
242		let path = self.path(key)?;
243		match self.create(&path, bytes.clone()).await? {
244			Create::Created => Ok(()),
245			Create::Exists => {
246				let existing = self.get_bytes(&path).await?;
247				if existing == bytes {
248					Ok(())
249				} else {
250					Err(Error::Conflict(path.to_string()))
251				}
252			}
253		}
254	}
255
256	async fn create(&self, path: &Path, bytes: Bytes) -> Result<Create> {
257		match self
258			.inner
259			.put_opts(path, PutPayload::from(bytes), PutMode::Create.into())
260			.await
261		{
262			Ok(_) => Ok(Create::Created),
263			Err(object_store::Error::AlreadyExists { .. }) => Ok(Create::Exists),
264			Err(err) => Err(err.into()),
265		}
266	}
267
268	async fn get_bytes(&self, path: &Path) -> Result<Bytes> {
269		Ok(self.inner.get(path).await?.bytes().await?)
270	}
271
272	fn recording_path(&self, relative: &Path) -> Path {
273		let mut path = self.prefix.clone();
274		path.extend(relative.parts());
275		path
276	}
277
278	fn list_path(&self, relative: Option<&Path>) -> Option<Path> {
279		let path = relative.map_or_else(|| self.prefix.clone(), |relative| self.recording_path(relative));
280		(!path.as_ref().is_empty()).then_some(path)
281	}
282
283	fn paginated_prefix(&self, relative: Option<&Path>) -> Option<String> {
284		self.list_path(relative).map(|path| format!("{path}/"))
285	}
286
287	fn paginated_options(&self, query: &Query) -> PaginatedListOptions {
288		PaginatedListOptions {
289			offset: query
290				.offset
291				.as_ref()
292				.map(|offset| self.recording_path(offset).to_string()),
293			max_keys: query.max_keys.map(NonZeroUsize::get),
294			page_token: query.page_token.clone(),
295			..Default::default()
296		}
297	}
298
299	fn page(&self, query: &Query, result: PaginatedListResult) -> Result<Page> {
300		let PaginatedListResult { result, page_token } = result;
301		let entries = self.entries(result)?;
302		let next = page_token.map(|token| query.next(token));
303		Ok(Page { entries, next })
304	}
305
306	fn entries(&self, result: ListResult) -> Result<Vec<Entry>> {
307		if let Some(prefix) = result.common_prefixes.first() {
308			return Err(Error::Directory(prefix.to_string()));
309		}
310		result
311			.objects
312			.into_iter()
313			.map(|meta| Entry::from_meta(&self.prefix, meta))
314			.collect()
315	}
316}
317
318impl<T: ObjectStore + PaginatedListStore> Store<T> {
319	/// List one page while keeping the backend continuation token inside the returned query.
320	pub async fn list_paginated(&self, query: &Query) -> Result<Page> {
321		let prefix = self.paginated_prefix(query.prefix.as_ref());
322		let opts = self.paginated_options(query);
323		let result = self.inner.list_paginated(prefix.as_deref(), opts).await?;
324		self.page(query, result)
325	}
326}
327
328impl Entry {
329	fn from_meta(prefix: &Path, meta: ObjectMeta) -> Result<Self> {
330		let key = Key::parse(prefix, &meta.location)?;
331		Ok(Self { key, size: meta.size })
332	}
333}
334
335enum Create {
336	Created,
337	Exists,
338}
339
340#[cfg(test)]
341mod tests {
342	use std::collections::BTreeSet;
343	use std::num::NonZeroUsize;
344
345	use futures::TryStreamExt;
346	use object_store::list::{PaginatedListOptions, PaginatedListResult, PaginatedListStore};
347	use object_store::memory::InMemory;
348	use object_store::path::Path;
349	use object_store::{
350		CopyOptions, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta, ObjectStore, PutMultipartOptions,
351		PutOptions, PutPayload, PutResult,
352	};
353
354	use super::*;
355	use crate::ID_MAX;
356	use crate::path::encode_track;
357	use crate::segment::{Frame, Group};
358
359	/// In-memory store with a trivial offset-based paginated listing implementation.
360	/// Page tokens are decimal indexes into the filtered, sorted key list.
361	#[derive(Debug, Clone)]
362	struct PaginatedMemory {
363		inner: InMemory,
364	}
365
366	impl std::fmt::Display for PaginatedMemory {
367		fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
368			write!(f, "PaginatedMemory")
369		}
370	}
371
372	#[async_trait::async_trait]
373	impl ObjectStore for PaginatedMemory {
374		async fn put_opts(
375			&self,
376			location: &Path,
377			payload: PutPayload,
378			opts: PutOptions,
379		) -> object_store::Result<PutResult> {
380			self.inner.put_opts(location, payload, opts).await
381		}
382
383		async fn put_multipart_opts(
384			&self,
385			location: &Path,
386			opts: PutMultipartOptions,
387		) -> object_store::Result<Box<dyn MultipartUpload>> {
388			self.inner.put_multipart_opts(location, opts).await
389		}
390
391		async fn get_opts(&self, location: &Path, options: GetOptions) -> object_store::Result<GetResult> {
392			self.inner.get_opts(location, options).await
393		}
394
395		fn delete_stream(
396			&self,
397			locations: BoxStream<'static, object_store::Result<Path>>,
398		) -> BoxStream<'static, object_store::Result<Path>> {
399			self.inner.delete_stream(locations)
400		}
401
402		fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, object_store::Result<ObjectMeta>> {
403			self.inner.list(prefix)
404		}
405
406		async fn list_with_delimiter(&self, prefix: Option<&Path>) -> object_store::Result<ListResult> {
407			self.inner.list_with_delimiter(prefix).await
408		}
409
410		async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> object_store::Result<()> {
411			self.inner.copy_opts(from, to, options).await
412		}
413	}
414
415	#[async_trait::async_trait]
416	impl PaginatedListStore for PaginatedMemory {
417		async fn list_paginated(
418			&self,
419			prefix: Option<&str>,
420			opts: PaginatedListOptions,
421		) -> object_store::Result<PaginatedListResult> {
422			let mut metas: Vec<ObjectMeta> = self.inner.list(None).try_collect().await?;
423			metas.sort_by(|a, b| a.location.as_ref().cmp(b.location.as_ref()));
424			if let Some(prefix) = prefix {
425				metas.retain(|meta| meta.location.as_ref().starts_with(prefix));
426			}
427			if let Some(offset) = opts.offset.as_deref() {
428				metas.retain(|meta| meta.location.as_ref() > offset);
429			}
430			let start: usize = match opts.page_token {
431				Some(token) => token.parse().map_err(|err| object_store::Error::Generic {
432					store: "PaginatedMemory",
433					source: Box::new(err),
434				})?,
435				None => 0,
436			};
437			let remaining = &metas[start.min(metas.len())..];
438			let take = opts.max_keys.unwrap_or(remaining.len()).min(remaining.len());
439			let objects = remaining[..take].to_vec();
440			let next = start.saturating_add(take);
441			let page_token = (next < metas.len()).then(|| next.to_string());
442			Ok(PaginatedListResult {
443				result: ListResult {
444					objects,
445					common_prefixes: Vec::new(),
446					extensions: Default::default(),
447				},
448				page_token,
449			})
450		}
451	}
452
453	fn memory() -> Store<InMemory> {
454		Store::new(InMemory::new(), "rec")
455	}
456
457	fn frame(timestamp: u64, payload: &'static [u8]) -> Frame {
458		Frame {
459			timestamp,
460			payload: Bytes::from_static(payload),
461		}
462	}
463
464	fn one_group(sequence: u64, payload: &'static [u8]) -> Object {
465		Object {
466			groups: vec![Group {
467				sequence,
468				frames: vec![frame(sequence, payload)],
469			}],
470		}
471	}
472
473	fn two_groups() -> Object {
474		Object {
475			groups: vec![
476				Group {
477					sequence: 5,
478					frames: vec![frame(0, b"a")],
479				},
480				Group {
481					sequence: 7,
482					frames: vec![frame(1, b"bb")],
483				},
484			],
485		}
486	}
487
488	async fn names(store: &Store<InMemory>, query: &Query) -> BTreeSet<String> {
489		store
490			.list(query)
491			.map_ok(|entry| entry.key.path(&Path::ROOT).unwrap().to_string())
492			.try_collect()
493			.await
494			.unwrap()
495	}
496
497	#[tokio::test]
498	async fn info_create_get_and_idempotent_retry() {
499		let store = memory();
500		let info = Info::new(1, 1_000).unwrap();
501		store.put_info("catalog.json", &info).await.unwrap();
502		assert_eq!(store.get_info("catalog.json").await.unwrap(), info);
503
504		let path = store.path(&Key::info("catalog.json").unwrap()).unwrap();
505		store
506			.inner()
507			.put(
508				&path,
509				br#"{ "timescale": 1000, "priority": 1, "version": 1 }"#.as_ref().into(),
510			)
511			.await
512			.unwrap();
513		store.put_info("catalog.json", &info).await.unwrap();
514		let kept = ObjectStoreExt::get(store.inner(), &path)
515			.await
516			.unwrap()
517			.bytes()
518			.await
519			.unwrap();
520		assert_eq!(&kept[..], br#"{ "timescale": 1000, "priority": 1, "version": 1 }"#);
521	}
522
523	#[tokio::test]
524	async fn info_property_mismatch_is_a_hard_error() {
525		let store = memory();
526		store.put_info("video", &Info::new(0, 1_000).unwrap()).await.unwrap();
527		assert!(matches!(
528			store.put_info("video", &Info::new(1, 1_000).unwrap()).await,
529			Err(Error::Priority {
530				existing: 0,
531				intended: 1
532			})
533		));
534		assert!(matches!(
535			store.put_info("video", &Info::new(0, 2_000).unwrap()).await,
536			Err(Error::TimescaleMismatch {
537				existing: 1000,
538				intended: 2000
539			})
540		));
541		assert_eq!(store.get_info("video").await.unwrap(), Info::new(0, 1_000).unwrap());
542	}
543
544	#[tokio::test]
545	async fn groups_put_get_and_identical_collision() {
546		let store = memory();
547		let object = two_groups();
548		let range = object.bounds().unwrap();
549		let key = store.put_groups("video", &object).await.unwrap();
550		assert_eq!(key, Key::groups("video", range.clone()).unwrap());
551		assert_eq!(
552			store.path(&key).unwrap().as_ref(),
553			"rec/video/groups/0000000000000000007.0000000000000000005"
554		);
555		assert_eq!(store.get_groups("video", range).await.unwrap(), object);
556		store.put_groups("video", &object).await.unwrap();
557	}
558
559	#[tokio::test]
560	async fn groups_collision_with_different_bytes_fails() {
561		let store = memory();
562		store.put_groups("video", &two_groups()).await.unwrap();
563		let other = Object {
564			groups: vec![
565				Group {
566					sequence: 5,
567					frames: vec![frame(0, b"X")],
568				},
569				Group {
570					sequence: 7,
571					frames: vec![frame(1, b"bb")],
572				},
573			],
574		};
575		assert!(matches!(
576			store.put_groups("video", &other).await,
577			Err(Error::Conflict(_))
578		));
579	}
580
581	#[tokio::test]
582	async fn segments_put_get_and_delete() {
583		let store = memory();
584		let object = one_group(0, b"tl");
585		store.put_segments("timeline.z", 0, &object).await.unwrap();
586		assert_eq!(store.get_segments("timeline.z", 0).await.unwrap(), object);
587		store.delete(&Key::segments("timeline.z", 0).unwrap()).await.unwrap();
588		assert!(matches!(
589			store.get_segments("timeline.z", 0).await,
590			Err(Error::NotFound(_))
591		));
592	}
593
594	#[tokio::test]
595	async fn listing_names_build_a_range_index() {
596		let store = memory();
597		store.put_info("video", &Info::new(0, 1).unwrap()).await.unwrap();
598		store.put_groups("video", &one_group(1, b"a")).await.unwrap();
599		store.put_groups("video", &one_group(3, b"b")).await.unwrap();
600		store.put_segments("timeline.z", 2, &one_group(0, b"t")).await.unwrap();
601
602		let listed: Vec<_> = store.list(&Query::new()).try_collect().await.unwrap();
603		let mut keys: Vec<_> = listed.into_iter().map(|listed| listed.key).collect();
604		keys.sort_by(|a, b| store.path(a).unwrap().as_ref().cmp(store.path(b).unwrap().as_ref()));
605		assert_eq!(
606			keys,
607			vec![
608				Key::segments("timeline.z", 2).unwrap(),
609				Key::info("video").unwrap(),
610				Key::groups("video", 1..=1).unwrap(),
611				Key::groups("video", 3..=3).unwrap(),
612			]
613		);
614	}
615
616	#[tokio::test]
617	async fn list_with_offset_finishes_before_the_caller_sorts() -> Result<()> {
618		let store = memory();
619		for segment in [0u64, 2, 1] {
620			store
621				.put_segments("timeline.z", segment, &one_group(segment, b"t"))
622				.await
623				.unwrap();
624		}
625		let query = Query::segments("timeline.z")?.after(&Key::segments("timeline.z", 0)?)?;
626		let mut rest: Vec<_> = store
627			.list(&query)
628			.map_ok(|entry| match entry.key {
629				Key::Segments { segment, .. } => segment,
630				_ => panic!("expected a segment key"),
631			})
632			.try_collect()
633			.await?;
634		rest.sort();
635		assert_eq!(rest, vec![1, 2]);
636		Ok(())
637	}
638
639	#[tokio::test]
640	async fn percent_encoded_track_listing() {
641		let store = memory();
642		store.put_info("catalog.json", &Info::new(0, 1).unwrap()).await.unwrap();
643		let locations = names(&store, &Query::new()).await;
644		assert!(locations.contains(&format!("{}/.info", encode_track("catalog.json").unwrap())));
645	}
646
647	#[tokio::test]
648	async fn listing_query_is_recording_relative() {
649		let store = memory();
650		store.put_info("video", &Info::new(0, 1).unwrap()).await.unwrap();
651		store.put_groups("video", &one_group(1, b"a")).await.unwrap();
652		store.put_info("video-alt", &Info::new(0, 1).unwrap()).await.unwrap();
653
654		let query = Query::track("video").unwrap();
655		let paths = names(&store, &query).await;
656		assert_eq!(
657			paths,
658			BTreeSet::from([
659				"video/.info".to_string(),
660				"video/groups/0000000000000000001.0000000000000000001".to_string(),
661			])
662		);
663		assert_eq!(
664			store.paginated_prefix(query.prefix.as_ref()).as_deref(),
665			Some("rec/video/")
666		);
667	}
668
669	#[tokio::test]
670	async fn paginated_recording_prefix_excludes_siblings() {
671		let store = memory();
672		store.put_info("video", &Info::new(0, 1).unwrap()).await.unwrap();
673		store
674			.inner()
675			.put(
676				&Path::from("rec-other/video/.info"),
677				Bytes::from_static(b"sibling").into(),
678			)
679			.await
680			.unwrap();
681
682		let prefix = store.paginated_prefix(None).unwrap();
683		assert_eq!(prefix, "rec/");
684		let objects = store
685			.inner()
686			.list(None)
687			.try_filter(|meta| futures::future::ready(meta.location.as_ref().starts_with(&prefix)))
688			.try_collect()
689			.await
690			.unwrap();
691		let entries = store
692			.entries(ListResult {
693				objects,
694				common_prefixes: Vec::new(),
695				extensions: Default::default(),
696			})
697			.unwrap();
698		assert_eq!(entries.len(), 1);
699		assert_eq!(entries[0].key, Key::info("video").unwrap());
700	}
701
702	#[tokio::test]
703	async fn paginated_query_preserves_every_supported_option() {
704		let store = memory();
705		let groups = Query::groups_from("video", 5).unwrap();
706		assert_eq!(
707			store.paginated_prefix(groups.prefix.as_ref()).as_deref(),
708			Some("rec/video/groups/")
709		);
710		assert_eq!(
711			store.paginated_options(&groups).offset.as_deref(),
712			Some("rec/video/groups/0000000000000000005")
713		);
714
715		let query = Query::segments("timeline.z")
716			.unwrap()
717			.after(&Key::segments("timeline.z", 1).unwrap())
718			.unwrap()
719			.page_size(NonZeroUsize::new(2).unwrap());
720		let options = store.paginated_options(&query);
721		assert_eq!(
722			store.paginated_prefix(query.prefix.as_ref()).as_deref(),
723			Some("rec/timeline%2Ez/segments/")
724		);
725		assert_eq!(
726			options.offset.as_deref(),
727			Some("rec/timeline%2Ez/segments/0000000000000000001")
728		);
729		assert_eq!(options.max_keys, Some(2));
730		assert_eq!(options.page_token, None);
731
732		let page = store
733			.page(
734				&query,
735				PaginatedListResult {
736					result: ListResult {
737						objects: Vec::new(),
738						common_prefixes: Vec::new(),
739						extensions: Default::default(),
740					},
741					page_token: Some("opaque".to_string()),
742				},
743			)
744			.unwrap();
745		let next = page.next.unwrap();
746		assert_eq!(next.prefix, query.prefix);
747		assert_eq!(next.offset, query.offset);
748		assert_eq!(next.max_keys, query.max_keys);
749		assert_eq!(store.paginated_options(&next).page_token.as_deref(), Some("opaque"));
750
751		let changed = next.page_size(NonZeroUsize::new(3).unwrap());
752		assert_eq!(store.paginated_options(&changed).page_token, None);
753	}
754
755	#[tokio::test]
756	async fn paginated_pages_match_streaming_results() {
757		let store = memory();
758		for segment in 0..3 {
759			store
760				.put_segments("timeline.z", segment, &one_group(segment, b"t"))
761				.await
762				.unwrap();
763		}
764		let query = Query::segments("timeline.z")
765			.unwrap()
766			.after(&Key::segments("timeline.z", 0).unwrap())
767			.unwrap()
768			.page_size(NonZeroUsize::new(1).unwrap());
769		let streamed = names(&store, &query).await;
770		let mut metas: Vec<_> = store
771			.inner()
772			.list(store.list_path(query.prefix.as_ref()).as_ref())
773			.try_collect()
774			.await
775			.unwrap();
776		metas.sort_by(|a, b| a.location.cmp(&b.location));
777		metas.retain(|meta| meta.location > store.recording_path(query.offset.as_ref().unwrap()));
778
779		let mut paged = BTreeSet::new();
780		let mut current = query;
781		for (index, meta) in metas.iter().cloned().enumerate() {
782			let page = store
783				.page(
784					&current,
785					PaginatedListResult {
786						result: ListResult {
787							objects: vec![meta],
788							common_prefixes: Vec::new(),
789							extensions: Default::default(),
790						},
791						page_token: (index + 1 < metas.len()).then(|| format!("page-{index}")),
792					},
793				)
794				.unwrap();
795			paged.extend(
796				page.entries
797					.into_iter()
798					.map(|entry| entry.key.path(&Path::ROOT).unwrap().to_string()),
799			);
800			if let Some(next) = page.next {
801				current = next;
802			}
803		}
804		assert_eq!(paged, streamed);
805	}
806
807	#[tokio::test]
808	async fn streaming_rejects_continuation_queries() {
809		let store = memory();
810		store.put_segments("timeline.z", 0, &one_group(0, b"t")).await.unwrap();
811		let query = Query::segments("timeline.z")
812			.unwrap()
813			.page_size(NonZeroUsize::new(1).unwrap());
814		let next = query.next("0".to_string());
815		let err = store.list(&next).try_collect::<Vec<_>>().await.unwrap_err();
816		assert_eq!(err, Error::Pagination);
817	}
818
819	#[tokio::test]
820	async fn paginated_listing_walks_pages_and_excludes_siblings() {
821		let store = Store::new(PaginatedMemory { inner: InMemory::new() }, "rec");
822		for segment in 0..3 {
823			store
824				.put_segments("timeline.z", segment, &one_group(segment, b"t"))
825				.await
826				.unwrap();
827		}
828		store
829			.inner()
830			.inner
831			.put(
832				&Path::from("rec-other/video/.info"),
833				Bytes::from_static(b"sibling").into(),
834			)
835			.await
836			.unwrap();
837
838		let query = Query::segments("timeline.z")
839			.unwrap()
840			.after(&Key::segments("timeline.z", 0).unwrap())
841			.unwrap()
842			.page_size(NonZeroUsize::new(1).unwrap());
843		let streamed: BTreeSet<String> = store
844			.list(&query)
845			.map_ok(|entry| entry.key.path(&Path::ROOT).unwrap().to_string())
846			.try_collect()
847			.await
848			.unwrap();
849		assert_eq!(streamed.len(), 2);
850
851		let mut paged = BTreeSet::new();
852		let mut current = Some(query);
853		let mut pages = 0;
854		while let Some(query) = current {
855			let page = store.list_paginated(&query).await.unwrap();
856			assert_eq!(page.entries.len(), 1);
857			paged.extend(
858				page.entries
859					.into_iter()
860					.map(|entry| entry.key.path(&Path::ROOT).unwrap().to_string()),
861			);
862			current = page.next;
863			pages += 1;
864			assert!(pages <= 2, "pagination did not terminate");
865		}
866		assert_eq!(pages, 2);
867		assert_eq!(paged, streamed);
868
869		// A recording-wide walk uses the trailing-slash scope, so the
870		// `rec-other` sibling must not appear in any page.
871		let query = Query::new().page_size(NonZeroUsize::new(1).unwrap());
872		let streamed: BTreeSet<String> = store
873			.list(&query)
874			.map_ok(|entry| entry.key.path(&Path::ROOT).unwrap().to_string())
875			.try_collect()
876			.await
877			.unwrap();
878		assert_eq!(streamed.len(), 3);
879
880		let mut paged = BTreeSet::new();
881		let mut current = Some(query);
882		let mut pages = 0;
883		while let Some(query) = current {
884			let page = store.list_paginated(&query).await.unwrap();
885			assert_eq!(page.entries.len(), 1);
886			paged.extend(
887				page.entries
888					.into_iter()
889					.map(|entry| entry.key.path(&Path::ROOT).unwrap().to_string()),
890			);
891			current = page.next;
892			pages += 1;
893			assert!(pages <= 3, "pagination did not terminate");
894		}
895		assert_eq!(pages, 3);
896		assert_eq!(paged, streamed);
897	}
898
899	#[test]
900	fn directory_results_are_not_silently_dropped() {
901		let store = memory();
902		let err = store
903			.entries(ListResult {
904				objects: Vec::new(),
905				common_prefixes: vec![Path::from("rec/video")],
906				extensions: Default::default(),
907			})
908			.unwrap_err();
909		assert_eq!(err, Error::Directory("rec/video".to_string()));
910	}
911
912	#[tokio::test]
913	async fn id_endpoints_are_valid_keys() {
914		let store = memory();
915		let object = one_group(ID_MAX, b"z");
916		store.put_groups("v", &object).await.unwrap();
917		store.put_segments("t", ID_MAX, &object).await.unwrap();
918		assert_eq!(store.get_groups("v", ID_MAX..=ID_MAX).await.unwrap(), object);
919		assert_eq!(store.get_segments("t", ID_MAX).await.unwrap(), object);
920	}
921
922	#[tokio::test]
923	async fn local_disk_roundtrip() {
924		let dir = tempfile::tempdir().unwrap();
925		let inner = object_store::local::LocalFileSystem::new_with_prefix(dir.path()).unwrap();
926		let store = Store::new(inner, Path::from("rec"));
927		let info = Info::new(3, 90_000).unwrap();
928		store.put_info("audio", &info).await.unwrap();
929		let object = two_groups();
930		store.put_groups("audio", &object).await.unwrap();
931		assert_eq!(store.get_info("audio").await.unwrap(), info);
932		assert_eq!(store.get_groups("audio", 5..=7).await.unwrap(), object);
933	}
934}