#[cfg(all(feature = "sled", feature = "serde"))]
#[test]
fn typed_bincode_survives_a_crash_on_sled() {
use persistent_queue::{Bincode, Builder, SledStore};
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize, Debug, PartialEq)]
struct Job {
id: u64,
name: String,
}
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("queue");
{
let store = SledStore::open(&path).unwrap();
let (tx, rx) = Builder::new(store).open_typed(Bincode).unwrap();
for id in 0..10u64 {
tx.push(&Job {
id,
name: format!("job-{id}"),
})
.unwrap();
}
let mut held = Vec::new();
while let Some(item) = rx.reserve().unwrap() {
if item.id % 2 == 0 {
item.ack().unwrap();
} else {
held.push(item);
}
}
}
let store = SledStore::open(&path).unwrap();
let (_tx, rx) = Builder::new(store).open_typed::<Job, _>(Bincode).unwrap();
let mut recovered = Vec::new();
while let Some(item) = rx.reserve().unwrap() {
assert_eq!(item.name, format!("job-{}", item.id));
recovered.push(item.id);
item.ack().unwrap();
}
recovered.sort_unstable();
assert_eq!(recovered, vec![1u64, 3, 5, 7, 9]);
}
#[cfg(all(feature = "redb", feature = "rkyv"))]
#[test]
fn typed_rkyv_survives_a_crash_on_redb() {
use persistent_queue::{Builder, RedbStore, Rkyv};
#[derive(rkyv::Archive, rkyv::Serialize, rkyv::Deserialize, Debug, PartialEq)]
struct Job {
id: u64,
name: String,
}
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("queue.redb");
{
let store = RedbStore::open(&path).unwrap();
let (tx, rx) = Builder::new(store).open_typed(Rkyv).unwrap();
for id in 0..10u64 {
tx.push(&Job {
id,
name: format!("job-{id}"),
})
.unwrap();
}
let mut held = Vec::new();
while let Some(item) = rx.reserve().unwrap() {
if item.id % 2 == 0 {
item.ack().unwrap();
} else {
held.push(item);
}
}
}
let store = RedbStore::open(&path).unwrap();
let (_tx, rx) = Builder::new(store).open_typed::<Job, _>(Rkyv).unwrap();
let mut recovered = Vec::new();
while let Some(item) = rx.reserve().unwrap() {
assert_eq!(item.name, format!("job-{}", item.id));
recovered.push(item.id);
item.ack().unwrap();
}
recovered.sort_unstable();
assert_eq!(recovered, vec![1u64, 3, 5, 7, 9]);
}
#[cfg(all(feature = "tokio", feature = "sled"))]
#[tokio::test]
async fn async_survives_a_crash_on_sled() {
use persistent_queue::{Builder, SledStore};
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("queue");
{
let store = SledStore::open(&path).unwrap();
let (tx, rx) = Builder::new(store).open_async().await.unwrap();
for i in 0..6u64 {
tx.push(i.to_le_bytes().to_vec()).await.unwrap();
}
for _ in 0..3 {
rx.reserve().await.unwrap().unwrap().ack().await.unwrap();
}
}
let store = SledStore::open(&path).unwrap();
let (tx, rx) = Builder::new(store).open_async().await.unwrap();
tx.close(); let mut recovered = 0;
while let Some(item) = rx.reserve().await.unwrap() {
recovered += 1;
item.ack().await.unwrap();
}
assert_eq!(recovered, 3, "the 3 unacked items survived the crash");
}