#![cfg(feature = "write")]
use std::collections::BTreeMap;
use dendro::archive::{Archive, SourceMeta, WalRow};
use dendro::keys;
use dendro::segment::{EncodeResult, SegmentEncoder};
use dendro::writer::Writer;
struct Never;
impl SegmentEncoder for Never {
fn encode(&self, _stream: &str, _rows: &[WalRow]) -> EncodeResult {
Ok(None)
}
}
fn source() -> SourceMeta {
SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::from([("kept".to_string(), "yes".to_string())]),
clock_anchor_wall_ns: 0,
}
}
fn row(ts: i64) -> WalRow {
WalRow {
stream: "s".to_string(),
ts,
wall_offset: 0,
row: vec![1],
}
}
fn patch(k: &str, v: &str) -> BTreeMap<String, String> {
BTreeMap::from([(k.to_string(), v.to_string())])
}
#[test]
fn a_patch_lands_in_order_with_the_ticks_and_without_a_finalize() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("m.dendro");
let mut archive = Writer::create(&path, Box::new(Never)).unwrap();
let mut w = archive.add_source(source()).unwrap();
let id = w.source_id();
w.wal(vec![row(1_000)]).unwrap();
w.update_metadata(patch(keys::PRODUCER_EPOCH, "epoch-a"))
.unwrap();
w.wal(vec![row(2_000)]).unwrap();
w.update_metadata(patch(keys::PRODUCER_EPOCHS, "[]"))
.unwrap();
w.sync().unwrap();
let seen = Archive::open(&path).unwrap().source_metadata(id).unwrap();
assert_eq!(
seen.get(keys::PRODUCER_EPOCH).map(String::as_str),
Some("epoch-a")
);
assert_eq!(
seen.get(keys::PRODUCER_EPOCHS).map(String::as_str),
Some("[]")
);
assert_eq!(
seen.get("kept").map(String::as_str),
Some("yes"),
"seed keys survive a patch"
);
w.update_metadata(patch(keys::PRODUCER_EPOCH, "epoch-b"))
.unwrap();
archive.finalize_single(w, (2_000, 0)).unwrap();
let after = Archive::open(&path).unwrap().source_metadata(id).unwrap();
assert_eq!(
after.get(keys::PRODUCER_EPOCH).map(String::as_str),
Some("epoch-b")
);
assert_eq!(
after.get(keys::PRODUCER_EPOCHS).map(String::as_str),
Some("[]")
);
assert!(after.contains_key(keys::WRITER_SESSIONS));
assert_eq!(after.len(), 4);
}
#[test]
fn a_patch_that_cannot_land_does_not_stop_the_writer() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("m.dendro");
let mut archive = Writer::create(&path, Box::new(Never)).unwrap();
let mut w = archive.add_source(source()).unwrap();
let id = w.source_id();
w.sync().unwrap();
rusqlite::Connection::open(&path)
.unwrap()
.execute("DELETE FROM sources WHERE id = ?1", [id])
.unwrap();
w.update_metadata(patch("k", "v")).unwrap();
w.sync().unwrap();
drop(w);
archive.join().unwrap();
}