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
16pub 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 #[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 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 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 pub fn track(track: &str) -> Result<Self> {
76 Ok(Self::new().prefix(path::track_prefix(track)?))
77 }
78
79 pub fn groups(track: &str) -> Result<Self> {
81 Ok(Self::new().prefix(path::groups_prefix(&Path::ROOT, track)?))
82 }
83
84 pub fn segments(track: &str) -> Result<Self> {
86 Ok(Self::new().prefix(path::segments_prefix(&Path::ROOT, track)?))
87 }
88
89 pub fn after(self, key: &Key) -> Result<Self> {
91 Ok(self.offset(key.path(&Path::ROOT)?))
92 }
93
94 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 #[derive(Debug, Clone, PartialEq, Eq)]
108 pub struct Entry {
109 pub key: Key,
111 pub size: u64,
113 }
114
115 #[derive(Debug, Clone)]
117 pub struct Page {
118 pub entries: Vec<Entry>,
120 pub next: Option<Query>,
123 }
124}
125
126use list::{Entry, Page, Query};
127
128#[derive(Clone, Debug)]
130pub struct Store<T> {
131 inner: T,
132 prefix: Path,
133}
134
135impl<T: ObjectStore> Store<T> {
136 pub fn new(inner: T, prefix: impl Into<Path>) -> Self {
138 Self {
139 inner,
140 prefix: prefix.into(),
141 }
142 }
143
144 pub fn inner(&self) -> &T {
146 &self.inner
147 }
148
149 pub fn prefix(&self) -> &Path {
151 &self.prefix
152 }
153
154 pub fn path(&self, key: &Key) -> Result<Path> {
156 key.path(&self.prefix)
157 }
158
159 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 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 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 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 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 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 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 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 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 #[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 ¤t,
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 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}