use std::num::NonZeroUsize;
use std::ops::RangeInclusive;
use bytes::Bytes;
use futures::StreamExt;
use futures::stream::BoxStream;
use object_store::list::{PaginatedListOptions, PaginatedListResult, PaginatedListStore};
use object_store::path::Path;
use object_store::{ListResult, ObjectMeta, ObjectStore, ObjectStoreExt, PutMode, PutPayload};
use crate::info::Info;
use crate::path::Key;
use crate::segment::Object;
use crate::{Error, Result};
pub mod list {
use std::num::NonZeroUsize;
use object_store::path::Path;
use crate::path;
use crate::{Key, Result};
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Query {
pub(crate) prefix: Option<Path>,
pub(crate) offset: Option<Path>,
pub(crate) max_keys: Option<NonZeroUsize>,
pub(crate) page_token: Option<String>,
}
impl Query {
pub fn new() -> Self {
Self::default()
}
fn prefix(mut self, prefix: impl Into<Path>) -> Self {
self.prefix = Some(prefix.into());
self.page_token = None;
self
}
fn offset(mut self, offset: impl Into<Path>) -> Self {
self.offset = Some(offset.into());
self.page_token = None;
self
}
pub fn page_size(mut self, max_keys: NonZeroUsize) -> Self {
self.max_keys = Some(max_keys);
self.page_token = None;
self
}
pub fn track(track: &str) -> Result<Self> {
Ok(Self::new().prefix(path::track_prefix(track)?))
}
pub fn groups(track: &str) -> Result<Self> {
Ok(Self::new().prefix(path::groups_prefix(&Path::ROOT, track)?))
}
pub fn segments(track: &str) -> Result<Self> {
Ok(Self::new().prefix(path::segments_prefix(&Path::ROOT, track)?))
}
pub fn after(self, key: &Key) -> Result<Self> {
Ok(self.offset(key.path(&Path::ROOT)?))
}
pub fn groups_from(track: &str, group: u64) -> Result<Self> {
Ok(Self::groups(track)?.offset(path::groups_offset(&Path::ROOT, track, group)?))
}
pub(crate) fn next(&self, page_token: String) -> Self {
let mut next = self.clone();
next.page_token = Some(page_token);
next
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Entry {
pub key: Key,
pub size: u64,
}
#[derive(Debug, Clone)]
pub struct Page {
pub entries: Vec<Entry>,
pub next: Option<Query>,
}
}
use list::{Entry, Page, Query};
#[derive(Clone, Debug)]
pub struct Store<T> {
inner: T,
prefix: Path,
}
impl<T: ObjectStore> Store<T> {
pub fn new(inner: T, prefix: impl Into<Path>) -> Self {
Self {
inner,
prefix: prefix.into(),
}
}
pub fn inner(&self) -> &T {
&self.inner
}
pub fn prefix(&self) -> &Path {
&self.prefix
}
pub fn path(&self, key: &Key) -> Result<Path> {
key.path(&self.prefix)
}
pub async fn put_info(&self, track: &str, info: &Info) -> Result<Key> {
let key = Key::info(track)?;
let path = self.path(&key)?;
let bytes = info.encode()?;
match self.create(&path, bytes).await? {
Create::Created => Ok(key),
Create::Exists => {
let existing = self.get_bytes(&path).await?;
let parsed = Info::decode(&existing)?;
if parsed.priority != info.priority {
return Err(Error::Priority {
existing: parsed.priority,
intended: info.priority,
});
}
if parsed.timescale != info.timescale {
return Err(Error::TimescaleMismatch {
existing: parsed.timescale,
intended: info.timescale,
});
}
Ok(key)
}
}
}
pub async fn get_info(&self, track: &str) -> Result<Info> {
let path = self.path(&Key::info(track)?)?;
Info::decode(&self.get_bytes(&path).await?)
}
pub async fn put_groups(&self, track: &str, object: &Object) -> Result<Key> {
let key = Key::groups(track, object.bounds()?)?;
self.put_segment(&key, object.encode()?).await?;
Ok(key)
}
pub async fn get_groups(&self, track: &str, range: RangeInclusive<u64>) -> Result<Object> {
let path = self.path(&Key::groups(track, range.clone())?)?;
Object::decode_groups(self.get_bytes(&path).await?, range)
}
pub async fn put_segments(&self, track: &str, segment: u64, object: &Object) -> Result<Key> {
let key = Key::segments(track, segment)?;
self.put_segment(&key, object.encode()?).await?;
Ok(key)
}
pub async fn get_segments(&self, track: &str, segment: u64) -> Result<Object> {
let path = self.path(&Key::segments(track, segment)?)?;
Object::decode(self.get_bytes(&path).await?)
}
pub async fn delete(&self, key: &Key) -> Result<()> {
let path = self.path(key)?;
self.inner.delete(&path).await?;
Ok(())
}
pub fn list(&self, query: &Query) -> BoxStream<'static, Result<Entry>> {
if query.page_token.is_some() {
return futures::stream::once(async { Err(Error::Pagination) }).boxed();
}
let prefix = self.list_path(query.prefix.as_ref());
let offset = query.offset.as_ref().map(|offset| self.recording_path(offset));
let store_prefix = self.prefix.clone();
let stream = match offset.as_ref() {
Some(offset) => self.inner.list_with_offset(prefix.as_ref(), offset),
None => self.inner.list(prefix.as_ref()),
};
stream.map(move |item| Entry::from_meta(&store_prefix, item?)).boxed()
}
async fn put_segment(&self, key: &Key, bytes: Bytes) -> Result<()> {
let path = self.path(key)?;
match self.create(&path, bytes.clone()).await? {
Create::Created => Ok(()),
Create::Exists => {
let existing = self.get_bytes(&path).await?;
if existing == bytes {
Ok(())
} else {
Err(Error::Conflict(path.to_string()))
}
}
}
}
async fn create(&self, path: &Path, bytes: Bytes) -> Result<Create> {
match self
.inner
.put_opts(path, PutPayload::from(bytes), PutMode::Create.into())
.await
{
Ok(_) => Ok(Create::Created),
Err(object_store::Error::AlreadyExists { .. }) => Ok(Create::Exists),
Err(err) => Err(err.into()),
}
}
async fn get_bytes(&self, path: &Path) -> Result<Bytes> {
Ok(self.inner.get(path).await?.bytes().await?)
}
fn recording_path(&self, relative: &Path) -> Path {
let mut path = self.prefix.clone();
path.extend(relative.parts());
path
}
fn list_path(&self, relative: Option<&Path>) -> Option<Path> {
let path = relative.map_or_else(|| self.prefix.clone(), |relative| self.recording_path(relative));
(!path.as_ref().is_empty()).then_some(path)
}
fn paginated_prefix(&self, relative: Option<&Path>) -> Option<String> {
self.list_path(relative).map(|path| format!("{path}/"))
}
fn paginated_options(&self, query: &Query) -> PaginatedListOptions {
PaginatedListOptions {
offset: query
.offset
.as_ref()
.map(|offset| self.recording_path(offset).to_string()),
max_keys: query.max_keys.map(NonZeroUsize::get),
page_token: query.page_token.clone(),
..Default::default()
}
}
fn page(&self, query: &Query, result: PaginatedListResult) -> Result<Page> {
let PaginatedListResult { result, page_token } = result;
let entries = self.entries(result)?;
let next = page_token.map(|token| query.next(token));
Ok(Page { entries, next })
}
fn entries(&self, result: ListResult) -> Result<Vec<Entry>> {
if let Some(prefix) = result.common_prefixes.first() {
return Err(Error::Directory(prefix.to_string()));
}
result
.objects
.into_iter()
.map(|meta| Entry::from_meta(&self.prefix, meta))
.collect()
}
}
impl<T: ObjectStore + PaginatedListStore> Store<T> {
pub async fn list_paginated(&self, query: &Query) -> Result<Page> {
let prefix = self.paginated_prefix(query.prefix.as_ref());
let opts = self.paginated_options(query);
let result = self.inner.list_paginated(prefix.as_deref(), opts).await?;
self.page(query, result)
}
}
impl Entry {
fn from_meta(prefix: &Path, meta: ObjectMeta) -> Result<Self> {
let key = Key::parse(prefix, &meta.location)?;
Ok(Self { key, size: meta.size })
}
}
enum Create {
Created,
Exists,
}
#[cfg(test)]
mod tests {
use std::collections::BTreeSet;
use std::num::NonZeroUsize;
use futures::TryStreamExt;
use object_store::list::{PaginatedListOptions, PaginatedListResult, PaginatedListStore};
use object_store::memory::InMemory;
use object_store::path::Path;
use object_store::{
CopyOptions, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta, ObjectStore, PutMultipartOptions,
PutOptions, PutPayload, PutResult,
};
use super::*;
use crate::ID_MAX;
use crate::path::encode_track;
use crate::segment::{Frame, Group};
#[derive(Debug, Clone)]
struct PaginatedMemory {
inner: InMemory,
}
impl std::fmt::Display for PaginatedMemory {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "PaginatedMemory")
}
}
#[async_trait::async_trait]
impl ObjectStore for PaginatedMemory {
async fn put_opts(
&self,
location: &Path,
payload: PutPayload,
opts: PutOptions,
) -> object_store::Result<PutResult> {
self.inner.put_opts(location, payload, opts).await
}
async fn put_multipart_opts(
&self,
location: &Path,
opts: PutMultipartOptions,
) -> object_store::Result<Box<dyn MultipartUpload>> {
self.inner.put_multipart_opts(location, opts).await
}
async fn get_opts(&self, location: &Path, options: GetOptions) -> object_store::Result<GetResult> {
self.inner.get_opts(location, options).await
}
fn delete_stream(
&self,
locations: BoxStream<'static, object_store::Result<Path>>,
) -> BoxStream<'static, object_store::Result<Path>> {
self.inner.delete_stream(locations)
}
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, object_store::Result<ObjectMeta>> {
self.inner.list(prefix)
}
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> object_store::Result<ListResult> {
self.inner.list_with_delimiter(prefix).await
}
async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> object_store::Result<()> {
self.inner.copy_opts(from, to, options).await
}
}
#[async_trait::async_trait]
impl PaginatedListStore for PaginatedMemory {
async fn list_paginated(
&self,
prefix: Option<&str>,
opts: PaginatedListOptions,
) -> object_store::Result<PaginatedListResult> {
let mut metas: Vec<ObjectMeta> = self.inner.list(None).try_collect().await?;
metas.sort_by(|a, b| a.location.as_ref().cmp(b.location.as_ref()));
if let Some(prefix) = prefix {
metas.retain(|meta| meta.location.as_ref().starts_with(prefix));
}
if let Some(offset) = opts.offset.as_deref() {
metas.retain(|meta| meta.location.as_ref() > offset);
}
let start: usize = match opts.page_token {
Some(token) => token.parse().map_err(|err| object_store::Error::Generic {
store: "PaginatedMemory",
source: Box::new(err),
})?,
None => 0,
};
let remaining = &metas[start.min(metas.len())..];
let take = opts.max_keys.unwrap_or(remaining.len()).min(remaining.len());
let objects = remaining[..take].to_vec();
let next = start.saturating_add(take);
let page_token = (next < metas.len()).then(|| next.to_string());
Ok(PaginatedListResult {
result: ListResult {
objects,
common_prefixes: Vec::new(),
extensions: Default::default(),
},
page_token,
})
}
}
fn memory() -> Store<InMemory> {
Store::new(InMemory::new(), "rec")
}
fn frame(timestamp: u64, payload: &'static [u8]) -> Frame {
Frame {
timestamp,
payload: Bytes::from_static(payload),
}
}
fn one_group(sequence: u64, payload: &'static [u8]) -> Object {
Object {
groups: vec![Group {
sequence,
frames: vec![frame(sequence, payload)],
}],
}
}
fn two_groups() -> Object {
Object {
groups: vec![
Group {
sequence: 5,
frames: vec![frame(0, b"a")],
},
Group {
sequence: 7,
frames: vec![frame(1, b"bb")],
},
],
}
}
async fn names(store: &Store<InMemory>, query: &Query) -> BTreeSet<String> {
store
.list(query)
.map_ok(|entry| entry.key.path(&Path::ROOT).unwrap().to_string())
.try_collect()
.await
.unwrap()
}
#[tokio::test]
async fn info_create_get_and_idempotent_retry() {
let store = memory();
let info = Info::new(1, 1_000).unwrap();
store.put_info("catalog.json", &info).await.unwrap();
assert_eq!(store.get_info("catalog.json").await.unwrap(), info);
let path = store.path(&Key::info("catalog.json").unwrap()).unwrap();
store
.inner()
.put(
&path,
br#"{ "timescale": 1000, "priority": 1, "version": 1 }"#.as_ref().into(),
)
.await
.unwrap();
store.put_info("catalog.json", &info).await.unwrap();
let kept = ObjectStoreExt::get(store.inner(), &path)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(&kept[..], br#"{ "timescale": 1000, "priority": 1, "version": 1 }"#);
}
#[tokio::test]
async fn info_property_mismatch_is_a_hard_error() {
let store = memory();
store.put_info("video", &Info::new(0, 1_000).unwrap()).await.unwrap();
assert!(matches!(
store.put_info("video", &Info::new(1, 1_000).unwrap()).await,
Err(Error::Priority {
existing: 0,
intended: 1
})
));
assert!(matches!(
store.put_info("video", &Info::new(0, 2_000).unwrap()).await,
Err(Error::TimescaleMismatch {
existing: 1000,
intended: 2000
})
));
assert_eq!(store.get_info("video").await.unwrap(), Info::new(0, 1_000).unwrap());
}
#[tokio::test]
async fn groups_put_get_and_identical_collision() {
let store = memory();
let object = two_groups();
let range = object.bounds().unwrap();
let key = store.put_groups("video", &object).await.unwrap();
assert_eq!(key, Key::groups("video", range.clone()).unwrap());
assert_eq!(
store.path(&key).unwrap().as_ref(),
"rec/video/groups/0000000000000000007.0000000000000000005"
);
assert_eq!(store.get_groups("video", range).await.unwrap(), object);
store.put_groups("video", &object).await.unwrap();
}
#[tokio::test]
async fn groups_collision_with_different_bytes_fails() {
let store = memory();
store.put_groups("video", &two_groups()).await.unwrap();
let other = Object {
groups: vec![
Group {
sequence: 5,
frames: vec![frame(0, b"X")],
},
Group {
sequence: 7,
frames: vec![frame(1, b"bb")],
},
],
};
assert!(matches!(
store.put_groups("video", &other).await,
Err(Error::Conflict(_))
));
}
#[tokio::test]
async fn segments_put_get_and_delete() {
let store = memory();
let object = one_group(0, b"tl");
store.put_segments("timeline.z", 0, &object).await.unwrap();
assert_eq!(store.get_segments("timeline.z", 0).await.unwrap(), object);
store.delete(&Key::segments("timeline.z", 0).unwrap()).await.unwrap();
assert!(matches!(
store.get_segments("timeline.z", 0).await,
Err(Error::NotFound(_))
));
}
#[tokio::test]
async fn listing_names_build_a_range_index() {
let store = memory();
store.put_info("video", &Info::new(0, 1).unwrap()).await.unwrap();
store.put_groups("video", &one_group(1, b"a")).await.unwrap();
store.put_groups("video", &one_group(3, b"b")).await.unwrap();
store.put_segments("timeline.z", 2, &one_group(0, b"t")).await.unwrap();
let listed: Vec<_> = store.list(&Query::new()).try_collect().await.unwrap();
let mut keys: Vec<_> = listed.into_iter().map(|listed| listed.key).collect();
keys.sort_by(|a, b| store.path(a).unwrap().as_ref().cmp(store.path(b).unwrap().as_ref()));
assert_eq!(
keys,
vec![
Key::segments("timeline.z", 2).unwrap(),
Key::info("video").unwrap(),
Key::groups("video", 1..=1).unwrap(),
Key::groups("video", 3..=3).unwrap(),
]
);
}
#[tokio::test]
async fn list_with_offset_finishes_before_the_caller_sorts() -> Result<()> {
let store = memory();
for segment in [0u64, 2, 1] {
store
.put_segments("timeline.z", segment, &one_group(segment, b"t"))
.await
.unwrap();
}
let query = Query::segments("timeline.z")?.after(&Key::segments("timeline.z", 0)?)?;
let mut rest: Vec<_> = store
.list(&query)
.map_ok(|entry| match entry.key {
Key::Segments { segment, .. } => segment,
_ => panic!("expected a segment key"),
})
.try_collect()
.await?;
rest.sort();
assert_eq!(rest, vec![1, 2]);
Ok(())
}
#[tokio::test]
async fn percent_encoded_track_listing() {
let store = memory();
store.put_info("catalog.json", &Info::new(0, 1).unwrap()).await.unwrap();
let locations = names(&store, &Query::new()).await;
assert!(locations.contains(&format!("{}/.info", encode_track("catalog.json").unwrap())));
}
#[tokio::test]
async fn listing_query_is_recording_relative() {
let store = memory();
store.put_info("video", &Info::new(0, 1).unwrap()).await.unwrap();
store.put_groups("video", &one_group(1, b"a")).await.unwrap();
store.put_info("video-alt", &Info::new(0, 1).unwrap()).await.unwrap();
let query = Query::track("video").unwrap();
let paths = names(&store, &query).await;
assert_eq!(
paths,
BTreeSet::from([
"video/.info".to_string(),
"video/groups/0000000000000000001.0000000000000000001".to_string(),
])
);
assert_eq!(
store.paginated_prefix(query.prefix.as_ref()).as_deref(),
Some("rec/video/")
);
}
#[tokio::test]
async fn paginated_recording_prefix_excludes_siblings() {
let store = memory();
store.put_info("video", &Info::new(0, 1).unwrap()).await.unwrap();
store
.inner()
.put(
&Path::from("rec-other/video/.info"),
Bytes::from_static(b"sibling").into(),
)
.await
.unwrap();
let prefix = store.paginated_prefix(None).unwrap();
assert_eq!(prefix, "rec/");
let objects = store
.inner()
.list(None)
.try_filter(|meta| futures::future::ready(meta.location.as_ref().starts_with(&prefix)))
.try_collect()
.await
.unwrap();
let entries = store
.entries(ListResult {
objects,
common_prefixes: Vec::new(),
extensions: Default::default(),
})
.unwrap();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].key, Key::info("video").unwrap());
}
#[tokio::test]
async fn paginated_query_preserves_every_supported_option() {
let store = memory();
let groups = Query::groups_from("video", 5).unwrap();
assert_eq!(
store.paginated_prefix(groups.prefix.as_ref()).as_deref(),
Some("rec/video/groups/")
);
assert_eq!(
store.paginated_options(&groups).offset.as_deref(),
Some("rec/video/groups/0000000000000000005")
);
let query = Query::segments("timeline.z")
.unwrap()
.after(&Key::segments("timeline.z", 1).unwrap())
.unwrap()
.page_size(NonZeroUsize::new(2).unwrap());
let options = store.paginated_options(&query);
assert_eq!(
store.paginated_prefix(query.prefix.as_ref()).as_deref(),
Some("rec/timeline%2Ez/segments/")
);
assert_eq!(
options.offset.as_deref(),
Some("rec/timeline%2Ez/segments/0000000000000000001")
);
assert_eq!(options.max_keys, Some(2));
assert_eq!(options.page_token, None);
let page = store
.page(
&query,
PaginatedListResult {
result: ListResult {
objects: Vec::new(),
common_prefixes: Vec::new(),
extensions: Default::default(),
},
page_token: Some("opaque".to_string()),
},
)
.unwrap();
let next = page.next.unwrap();
assert_eq!(next.prefix, query.prefix);
assert_eq!(next.offset, query.offset);
assert_eq!(next.max_keys, query.max_keys);
assert_eq!(store.paginated_options(&next).page_token.as_deref(), Some("opaque"));
let changed = next.page_size(NonZeroUsize::new(3).unwrap());
assert_eq!(store.paginated_options(&changed).page_token, None);
}
#[tokio::test]
async fn paginated_pages_match_streaming_results() {
let store = memory();
for segment in 0..3 {
store
.put_segments("timeline.z", segment, &one_group(segment, b"t"))
.await
.unwrap();
}
let query = Query::segments("timeline.z")
.unwrap()
.after(&Key::segments("timeline.z", 0).unwrap())
.unwrap()
.page_size(NonZeroUsize::new(1).unwrap());
let streamed = names(&store, &query).await;
let mut metas: Vec<_> = store
.inner()
.list(store.list_path(query.prefix.as_ref()).as_ref())
.try_collect()
.await
.unwrap();
metas.sort_by(|a, b| a.location.cmp(&b.location));
metas.retain(|meta| meta.location > store.recording_path(query.offset.as_ref().unwrap()));
let mut paged = BTreeSet::new();
let mut current = query;
for (index, meta) in metas.iter().cloned().enumerate() {
let page = store
.page(
¤t,
PaginatedListResult {
result: ListResult {
objects: vec![meta],
common_prefixes: Vec::new(),
extensions: Default::default(),
},
page_token: (index + 1 < metas.len()).then(|| format!("page-{index}")),
},
)
.unwrap();
paged.extend(
page.entries
.into_iter()
.map(|entry| entry.key.path(&Path::ROOT).unwrap().to_string()),
);
if let Some(next) = page.next {
current = next;
}
}
assert_eq!(paged, streamed);
}
#[tokio::test]
async fn streaming_rejects_continuation_queries() {
let store = memory();
store.put_segments("timeline.z", 0, &one_group(0, b"t")).await.unwrap();
let query = Query::segments("timeline.z")
.unwrap()
.page_size(NonZeroUsize::new(1).unwrap());
let next = query.next("0".to_string());
let err = store.list(&next).try_collect::<Vec<_>>().await.unwrap_err();
assert_eq!(err, Error::Pagination);
}
#[tokio::test]
async fn paginated_listing_walks_pages_and_excludes_siblings() {
let store = Store::new(PaginatedMemory { inner: InMemory::new() }, "rec");
for segment in 0..3 {
store
.put_segments("timeline.z", segment, &one_group(segment, b"t"))
.await
.unwrap();
}
store
.inner()
.inner
.put(
&Path::from("rec-other/video/.info"),
Bytes::from_static(b"sibling").into(),
)
.await
.unwrap();
let query = Query::segments("timeline.z")
.unwrap()
.after(&Key::segments("timeline.z", 0).unwrap())
.unwrap()
.page_size(NonZeroUsize::new(1).unwrap());
let streamed: BTreeSet<String> = store
.list(&query)
.map_ok(|entry| entry.key.path(&Path::ROOT).unwrap().to_string())
.try_collect()
.await
.unwrap();
assert_eq!(streamed.len(), 2);
let mut paged = BTreeSet::new();
let mut current = Some(query);
let mut pages = 0;
while let Some(query) = current {
let page = store.list_paginated(&query).await.unwrap();
assert_eq!(page.entries.len(), 1);
paged.extend(
page.entries
.into_iter()
.map(|entry| entry.key.path(&Path::ROOT).unwrap().to_string()),
);
current = page.next;
pages += 1;
assert!(pages <= 2, "pagination did not terminate");
}
assert_eq!(pages, 2);
assert_eq!(paged, streamed);
let query = Query::new().page_size(NonZeroUsize::new(1).unwrap());
let streamed: BTreeSet<String> = store
.list(&query)
.map_ok(|entry| entry.key.path(&Path::ROOT).unwrap().to_string())
.try_collect()
.await
.unwrap();
assert_eq!(streamed.len(), 3);
let mut paged = BTreeSet::new();
let mut current = Some(query);
let mut pages = 0;
while let Some(query) = current {
let page = store.list_paginated(&query).await.unwrap();
assert_eq!(page.entries.len(), 1);
paged.extend(
page.entries
.into_iter()
.map(|entry| entry.key.path(&Path::ROOT).unwrap().to_string()),
);
current = page.next;
pages += 1;
assert!(pages <= 3, "pagination did not terminate");
}
assert_eq!(pages, 3);
assert_eq!(paged, streamed);
}
#[test]
fn directory_results_are_not_silently_dropped() {
let store = memory();
let err = store
.entries(ListResult {
objects: Vec::new(),
common_prefixes: vec![Path::from("rec/video")],
extensions: Default::default(),
})
.unwrap_err();
assert_eq!(err, Error::Directory("rec/video".to_string()));
}
#[tokio::test]
async fn id_endpoints_are_valid_keys() {
let store = memory();
let object = one_group(ID_MAX, b"z");
store.put_groups("v", &object).await.unwrap();
store.put_segments("t", ID_MAX, &object).await.unwrap();
assert_eq!(store.get_groups("v", ID_MAX..=ID_MAX).await.unwrap(), object);
assert_eq!(store.get_segments("t", ID_MAX).await.unwrap(), object);
}
#[tokio::test]
async fn local_disk_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let inner = object_store::local::LocalFileSystem::new_with_prefix(dir.path()).unwrap();
let store = Store::new(inner, Path::from("rec"));
let info = Info::new(3, 90_000).unwrap();
store.put_info("audio", &info).await.unwrap();
let object = two_groups();
store.put_groups("audio", &object).await.unwrap();
assert_eq!(store.get_info("audio").await.unwrap(), info);
assert_eq!(store.get_groups("audio", 5..=7).await.unwrap(), object);
}
}