use std::sync::Arc;
use bytes::Bytes;
use taquba::Queue;
use crate::error::Result;
#[derive(Clone)]
pub struct KvReadHandle {
queue: Option<Arc<Queue>>,
}
impl std::fmt::Debug for KvReadHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match &self.queue {
Some(_) => f.debug_struct("KvReadHandle").finish_non_exhaustive(),
None => f.write_str("KvReadHandle::detached"),
}
}
}
impl KvReadHandle {
pub fn detached() -> Self {
Self { queue: None }
}
pub(crate) fn for_delivery(queue: Arc<Queue>) -> Self {
Self { queue: Some(queue) }
}
pub async fn get(&self, key: &[u8]) -> Result<Option<Bytes>> {
match &self.queue {
Some(queue) => Ok(queue.view().kv_get(key).await?),
None => Ok(None),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn a_detached_handle_reads_no_value() {
let handle = KvReadHandle::detached();
assert_eq!(handle.get(b"any").await.unwrap(), None);
}
}