use crate::error::{Error, ReadOnly, Result};
use std::collections::BTreeMap;
use std::fs::OpenOptions;
use std::path::Path;
use rusqlite::{Connection, OpenFlags};
pub const PAGE_SIZE: u32 = 4096;
const WAL_AUTOCHECKPOINT_BYTES: u32 = 4 * 1024 * 1024;
const READER_CACHE_SIZE_KIB: i32 = -262_144;
const WRITER_CACHE_SIZE_KIB: i32 = -16_384;
const _: () = assert!(
WRITER_CACHE_SIZE_KIB > READER_CACHE_SIZE_KIB,
"the writer's page cache must be smaller than the reader's"
);
const SCHEMA_VERSION: i64 = 4;
pub const APPLICATION_ID: u32 = 0x6465_6e64;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum Sniff {
Stamped {
version: i64,
},
NotAnArchive,
}
pub fn sniff_bytes(bytes: &[u8]) -> Sniff {
const MAGIC: &[u8; 16] = b"SQLite format 3\0";
const HEADER_LEN: usize = 100;
const USER_VERSION_OFFSET: usize = 60;
const APPLICATION_ID_OFFSET: usize = 68;
if bytes.len() < HEADER_LEN || &bytes[..MAGIC.len()] != MAGIC {
return Sniff::NotAnArchive;
}
let be = |at: usize| {
let mut b = [0u8; 4];
b.copy_from_slice(&bytes[at..at + 4]);
u32::from_be_bytes(b)
};
match be(APPLICATION_ID_OFFSET) {
APPLICATION_ID => Sniff::Stamped {
version: i64::from(be(USER_VERSION_OFFSET)),
},
_ => Sniff::NotAnArchive,
}
}
pub fn sniff(path: &Path) -> Result<Sniff> {
use std::io::Read as _;
let mut file = std::fs::File::open(path)
.map_err(|e| Error::Message(format!("failed to open {}: {e}", path.display())))?;
let mut header = [0u8; 100];
match file.read_exact(&mut header) {
Ok(()) => Ok(sniff_bytes(&header)),
Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => Ok(Sniff::NotAnArchive),
Err(e) => Err(Error::Message(format!(
"failed to read {}: {e}",
path.display()
))),
}
}
fn is_readonly_media(e: &Error) -> bool {
matches!(
e.sqlite_code(),
Some(rusqlite::ErrorCode::ReadOnly | rusqlite::ErrorCode::CannotOpen)
)
}
fn sidecar_exists(path: &Path, suffix: &str) -> bool {
let mut p = path.as_os_str().to_os_string();
p.push(suffix);
Path::new(&p).exists()
}
fn immutable_uri(path: &Path) -> Option<String> {
let path = path.to_str()?;
let mut uri = String::with_capacity(path.len() + 24);
uri.push_str("file:");
for b in path.bytes() {
match b {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' | b'/' => {
uri.push(b as char)
}
_ => uri.push_str(&format!("%{b:02X}")),
}
}
uri.push_str("?immutable=1");
Some(uri)
}
#[derive(Clone, Debug, PartialEq)]
pub struct SourceMeta {
pub labels: BTreeMap<String, String>,
pub metadata: BTreeMap<String, String>,
pub clock_anchor_wall_ns: i64,
}
#[derive(Clone, Debug, PartialEq)]
#[non_exhaustive]
pub struct SourceRow {
pub id: i64,
pub meta: SourceMeta,
pub uuid: Option<String>,
pub complete: bool,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct SegmentMeta {
pub rows: u64,
pub first_ts: i64,
pub last_ts: i64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct SegmentRow {
pub seq: u64,
pub meta: SegmentMeta,
pub bytes: Vec<u8>,
pub caller_index: Option<Vec<u8>>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct WalRow {
pub stream: String,
pub ts: i64,
pub wall_offset: i64,
pub row: Vec<u8>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct CallerRow {
pub ts: i64,
pub blob: Vec<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[non_exhaustive]
pub struct Evicted {
pub segments: usize,
pub wal_rows: usize,
pub live_rows: usize,
pub caller_rows: usize,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct Span {
pub rows: u64,
pub first_ts: Option<i64>,
pub last_ts: Option<i64>,
}
const LIVE_WAL_PREDICATE: &str = "source_id = ?1 AND stream = ?2 \
AND ( \
ts > (SELECT MAX(last_ts) FROM segments \
WHERE source_id = ?1 AND stream = ?2) \
OR NOT EXISTS (SELECT 1 FROM segments \
WHERE source_id = ?1 AND stream = ?2) \
)";
const LIVE_WAL_PREDICATE_FOR_ROW: &str = "\
ts > (SELECT MAX(last_ts) FROM segments s \
WHERE s.source_id = wal.source_id AND s.stream = wal.stream) \
OR NOT EXISTS (SELECT 1 FROM segments s \
WHERE s.source_id = wal.source_id AND s.stream = wal.stream)";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Depth {
Quick,
Full,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct Report {
pub sources: usize,
pub streams: usize,
pub segments: usize,
pub wal_rows: usize,
pub problems: Vec<Problem>,
}
impl Report {
pub fn is_sound(&self) -> bool {
self.problems.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum Problem {
Corrupt(String),
ForeignKey(String),
Segment {
source_id: i64,
stream: String,
seq: u64,
detail: String,
},
UnreadableWalRows {
source_id: i64,
stream: String,
rows: usize,
},
}
impl std::fmt::Display for Problem {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Problem::Corrupt(m) => write!(f, "corrupt: {m}"),
Problem::ForeignKey(m) => write!(f, "dangling reference: {m}"),
Problem::Segment {
source_id,
stream,
seq,
detail,
} => write!(
f,
"source {source_id}, stream {stream}, segment {seq}: {detail}"
),
Problem::UnreadableWalRows {
source_id,
stream,
rows,
} => write!(
f,
"source {source_id}, stream {stream}: {rows} WAL row(s) at or below the \
sealed watermark, which no read path can reach"
),
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct PageStats {
pub pages: u32,
pub free: u32,
pub page_size: u32,
}
pub struct Archive {
conn: Connection,
}
impl std::fmt::Debug for Archive {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Archive")
.field("path", &self.conn.path())
.finish_non_exhaustive()
}
}
impl Archive {
pub fn serialize(&self) -> Result<Vec<u8>> {
let data = self
.conn
.serialize(rusqlite::MAIN_DB)
.map_err(Error::sqlite("failed to serialize the archive"))?;
Ok(data.to_vec())
}
fn sidecar(path: &Path, suffix: &str) -> std::path::PathBuf {
let mut name = path.as_os_str().to_os_string();
name.push(suffix);
std::path::PathBuf::from(name)
}
pub(crate) fn remove_archive(path: &Path) {
for p in [
path.to_path_buf(),
Self::sidecar(path, "-wal"),
Self::sidecar(path, "-shm"),
] {
match std::fs::remove_file(&p) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => tracing::warn!("failed to remove {}: {e}", p.display()),
}
}
}
fn create_with_page_size(path: &Path, page_size: u32) -> Result<ArchiveMut> {
OpenOptions::new()
.write(true)
.create_new(true)
.open(path)
.map_err(|e| Error::Message(format!("failed to create {}: {e}", path.display())))?;
match Self::init_created(path, page_size) {
Ok(db) => Ok(db),
Err(e) => {
Self::remove_archive(path);
Err(e)
}
}
}
fn init_created(path: &Path, page_size: u32) -> Result<ArchiveMut> {
let conn = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)
.map_err(|e| Error::Message(format!("failed to open {}: {e}", path.display())))?;
let db = Archive { conn };
db.set_pragma("page_size", page_size)?;
db.set_pragma("auto_vacuum", "INCREMENTAL")?;
db.set_journal_mode_wal()?;
db.apply_connection_pragmas(WRITER_CACHE_SIZE_KIB)?;
db.conn
.execute_batch(SCHEMA_SQL)
.map_err(Error::sqlite("failed to create archive schema"))?;
db.stamp_header()?;
let mut db = ArchiveMut::wrap(db);
db.checkpoint_passive()?;
Ok(db)
}
fn open_with_cache(path: &Path, cache_size_kib: i32, exclusive: bool) -> Result<Self> {
let conn = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)
.map_err(|e| Error::Message(format!("failed to open {}: {e}", path.display())))?;
if conn
.is_readonly(rusqlite::MAIN_DB)
.map_err(Error::sqlite(format!("failed to open {}", path.display())))?
{
return Err(Error::ReadOnly(ReadOnly::Media));
}
let mut db = Archive { conn };
if exclusive {
db.take_exclusive_lock()?;
}
db.adopt_schema(&path.display().to_string())?;
db.apply_connection_pragmas(cache_size_kib)?;
Ok(db)
}
pub fn open(path: &Path) -> Result<Self> {
match Self::open_at(path, None) {
Err(e) if is_readonly_media(&e) => {
if sidecar_exists(path, "-wal") {
return Err(e);
}
match immutable_uri(path) {
Some(uri) => Self::open_at(path, Some(&uri)),
None => Err(e),
}
}
other => other,
}
}
fn open_at(path: &Path, uri: Option<&str>) -> Result<Self> {
let conn = match uri {
None => Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY),
Some(uri) => Connection::open_with_flags(
uri,
OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI,
),
}
.map_err(Error::sqlite(format!(
"failed to open {} read-only",
path.display()
)))?;
let mut db = Archive { conn };
let what = path.display().to_string();
db.adopt_schema(&what)?;
db.set_pragma("cache_size", READER_CACHE_SIZE_KIB)?;
db.set_pragma("query_only", "1")?;
Ok(db)
}
pub fn open_bytes(bytes: Vec<u8>) -> Result<Self> {
const HEADER: &[u8] = b"SQLite format 3\0";
const JOURNAL_MODE_ROLLBACK: u8 = 1;
const FILE_FORMAT_WAL: u8 = 2;
let mut bytes = bytes;
if bytes.len() < 20 || !bytes.starts_with(HEADER) {
return Err(Error::NotAnArchive {
what: "<bytes>".to_string(),
reason: "not a SQLite database".to_string(),
});
}
if bytes[18] == FILE_FORMAT_WAL && bytes[19] == FILE_FORMAT_WAL {
bytes[18] = JOURNAL_MODE_ROLLBACK;
bytes[19] = JOURNAL_MODE_ROLLBACK;
}
let mut conn = Connection::open_in_memory()
.map_err(Error::sqlite("failed to open an in-memory database"))?;
let len = bytes.len();
conn.deserialize_read_exact(rusqlite::MAIN_DB, &mut bytes.as_slice(), len, true)
.map_err(Error::sqlite("failed to read the archive"))?;
let mut db = Archive { conn };
db.adopt_schema("<bytes>")?;
db.apply_connection_pragmas(READER_CACHE_SIZE_KIB)?;
Ok(db)
}
fn stamp_header(&self) -> Result<()> {
self.set_pragma("application_id", APPLICATION_ID)?;
self.set_pragma("user_version", SCHEMA_VERSION)
}
fn adopt_schema(&mut self, what: &str) -> Result<()> {
let not_an_archive = |reason: &str| Error::NotAnArchive {
what: what.to_string(),
reason: reason.to_string(),
};
let app = match self
.conn
.pragma_query_value(None, "application_id", |row| row.get::<_, i64>(0))
{
Ok(v) => v,
Err(rusqlite::Error::SqliteFailure(e, _))
if e.code == rusqlite::ErrorCode::NotADatabase =>
{
return Err(not_an_archive("not a SQLite database"));
}
Err(rusqlite::Error::SqliteFailure(e, _))
if e.code == rusqlite::ErrorCode::DatabaseBusy =>
{
return Err(Error::InUse {
what: what.to_string(),
});
}
Err(e) => return Err(Error::sqlite("failed to read application_id")(e)),
};
if app == 0 {
return Err(not_an_archive(
"a SQLite database without dendro's header stamp. dendro reads only \
archives it wrote; a `.rez` recording from before dendro is upgraded \
by rezolus",
));
}
if app != i64::from(APPLICATION_ID) {
return Err(not_an_archive(&format!(
"a SQLite database of another application (id {app:#x})"
)));
}
match self.pragma_i64("user_version")? {
SCHEMA_VERSION => Ok(()),
other => Err(Error::UnsupportedSchema {
found: other,
writes: SCHEMA_VERSION,
reads: SCHEMA_VERSION,
}),
}
}
fn apply_connection_pragmas(&self, cache_size_kib: i32) -> Result<()> {
self.set_pragma("synchronous", "FULL")?;
self.set_pragma("foreign_keys", "ON")?;
let pages = WAL_AUTOCHECKPOINT_BYTES / self.pragma_u32("page_size")?.max(1);
self.set_pragma("wal_autocheckpoint", pages)?;
self.set_pragma("cache_size", cache_size_kib)?;
Ok(())
}
fn take_exclusive_lock(&self) -> Result<()> {
let mode: String = self
.conn
.pragma_update_and_check(None, "locking_mode", "EXCLUSIVE", |row| row.get(0))
.map_err(Error::sqlite("failed to set locking_mode=EXCLUSIVE"))?;
if !mode.eq_ignore_ascii_case("exclusive") {
return Err(Error::Message(format!(
"locking_mode is {mode}, expected exclusive"
)));
}
self.conn
.busy_timeout(std::time::Duration::from_secs(1))
.map_err(Error::sqlite("failed to set busy_timeout"))
}
fn set_journal_mode_wal(&self) -> Result<()> {
let mode: String = self
.conn
.pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
.map_err(Error::sqlite("failed to set journal_mode=WAL"))?;
if !mode.eq_ignore_ascii_case("wal") {
return Err(Error::Message(format!(
"journal_mode is {mode}, expected wal"
)));
}
Ok(())
}
fn set_pragma<V: rusqlite::ToSql>(&self, name: &str, value: V) -> Result<()> {
self.conn
.pragma_update(None, name, value)
.map_err(Error::sqlite(format!("failed to set pragma {name}")))
}
pub fn mint_uuid(&self) -> Result<String> {
mint_uuid(&self.conn)
}
pub fn read_sources(&self) -> Result<Vec<SourceRow>> {
let uuid_col = if has_column(&self.conn, "sources", "uuid")? {
"uuid"
} else {
"NULL"
};
let mut stmt = self
.conn
.prepare(&format!(
"SELECT id, labels, metadata, complete, clock_anchor_wall_ns, {uuid_col} \
FROM sources ORDER BY id"
))
.map_err(Error::sqlite("failed to query sources"))?;
let rows = stmt
.query_map([], |row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, i64>(3)?,
row.get::<_, i64>(4)?,
row.get::<_, Option<String>>(5)?,
))
})
.map_err(Error::sqlite("failed to query sources"))?;
let mut out = Vec::new();
for row in rows {
let (id, labels, metadata, complete, anchor, uuid) =
row.map_err(Error::sqlite("failed to read source"))?;
out.push(SourceRow {
id,
uuid,
meta: SourceMeta {
labels: serde_json::from_str(&labels).map_err(|e| {
Error::Message(format!("source {id} has invalid labels: {e}"))
})?,
metadata: serde_json::from_str(&metadata).map_err(|e| {
Error::Message(format!("source {id} has invalid metadata: {e}"))
})?,
clock_anchor_wall_ns: anchor,
},
complete: complete != 0,
});
}
Ok(out)
}
pub fn read_segment_indexes(
&self,
source_id: i64,
stream: &str,
) -> Result<Vec<(u64, Option<Vec<u8>>)>> {
read_segment_indexes_sql(&self.conn, source_id, stream)
}
pub fn read_segment_meta(
&self,
source_id: i64,
stream: &str,
) -> Result<Vec<(u64, SegmentMeta)>> {
let mut stmt = self
.conn
.prepare(
"SELECT seq, rows, first_ts, last_ts FROM segments \
WHERE source_id = ?1 AND stream = ?2 ORDER BY seq",
)
.map_err(Error::sqlite(format!(
"failed to query segment meta for {stream}"
)))?;
let rows = stmt
.query_map(rusqlite::params![source_id, stream], |r| {
Ok((
r.get::<_, i64>(0)? as u64,
SegmentMeta {
rows: r.get::<_, i64>(1)? as u64,
first_ts: r.get::<_, i64>(2)?,
last_ts: r.get::<_, i64>(3)?,
},
))
})
.map_err(Error::sqlite(format!(
"failed to read segment meta for {stream}"
)))?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::sqlite(format!(
"failed to read segment meta for {stream}"
)))
}
pub fn read_segment_bytes(
&self,
source_id: i64,
stream: &str,
seq: u64,
) -> Result<Option<Vec<u8>>> {
let mut stmt = self
.conn
.prepare(
"SELECT bytes FROM segments \
WHERE source_id = ?1 AND stream = ?2 AND seq = ?3",
)
.map_err(Error::sqlite(format!(
"failed to query segment bytes for {stream}"
)))?;
let mut rows = stmt
.query(rusqlite::params![source_id, stream, seq as i64])
.map_err(Error::sqlite(format!(
"failed to read segment bytes for {stream}"
)))?;
match rows.next() {
Ok(Some(r)) => {
Ok(Some(r.get(0).map_err(|e| {
format!("failed to read segment bytes for {stream}: {e}")
})?))
}
Ok(None) => Ok(None),
Err(e) => Err(Error::sqlite(format!(
"failed to read segment bytes for {stream}"
))(e)),
}
}
pub fn read_segments(&self, source_id: i64, stream: &str) -> Result<Vec<SegmentRow>> {
let mut stmt = self
.conn
.prepare(
"SELECT seq, rows, first_ts, last_ts, bytes, caller_index FROM segments \
WHERE source_id = ?1 AND stream = ?2 ORDER BY seq",
)
.map_err(Error::sqlite(format!(
"failed to query segments for {stream}"
)))?;
Self::collect_segments(&mut stmt, rusqlite::params![source_id, stream], stream)
}
fn collect_segments(
stmt: &mut rusqlite::Statement<'_>,
params: &[&dyn rusqlite::ToSql],
stream: &str,
) -> Result<Vec<SegmentRow>> {
let rows = stmt
.query_map(params, |row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, i64>(3)?,
row.get::<_, Vec<u8>>(4)?,
row.get::<_, Option<Vec<u8>>>(5)?,
))
})
.map_err(Error::sqlite(format!(
"failed to query segments for {stream}"
)))?;
let mut out = Vec::new();
for row in rows {
let (seq, n_rows, first_ts, last_ts, bytes, caller_index) = row.map_err(
Error::sqlite(format!("failed to read segment row for {stream}")),
)?;
out.push(SegmentRow {
seq: seq as u64,
meta: SegmentMeta {
rows: n_rows as u64,
first_ts,
last_ts,
},
bytes,
caller_index,
});
}
Ok(out)
}
pub fn segments_overlapping(
&self,
source_id: i64,
stream: &str,
start: i64,
end: i64,
) -> Result<Vec<SegmentRow>> {
let mut stmt = self
.conn
.prepare(
"SELECT seq, rows, first_ts, last_ts, bytes, caller_index FROM segments \
WHERE source_id = ?1 AND stream = ?2 \
AND last_ts >= ?3 AND first_ts <= ?4 ORDER BY seq",
)
.map_err(Error::sqlite(format!(
"failed to query segments for {stream}"
)))?;
let params = rusqlite::params![source_id, stream, start, end,];
Self::collect_segments(&mut stmt, params, stream)
}
pub fn read_snapshot<T>(&self, f: impl FnOnce(&Self) -> Result<T>) -> Result<T> {
if !self.conn.is_autocommit() {
return f(self);
}
self.conn
.execute_batch("BEGIN DEFERRED")
.map_err(Error::sqlite("failed to open a read snapshot"))?;
struct EndSnapshot<'a>(&'a Connection);
impl Drop for EndSnapshot<'_> {
fn drop(&mut self) {
let _ = self.0.execute_batch("ROLLBACK");
}
}
let guard = EndSnapshot(&self.conn);
let out = f(self);
drop(guard);
out
}
pub fn total_rows(&self, source_id: i64, stream: &str) -> Result<u64> {
let total: i64 = self
.conn
.query_row(
"SELECT COALESCE(SUM(rows), 0) FROM segments WHERE source_id = ?1 AND stream = ?2",
rusqlite::params![source_id, stream],
|row| row.get(0),
)
.map_err(Error::sqlite(format!("failed to sum rows for {stream}")))?;
Ok(total as u64)
}
#[cfg(test)]
fn streams(&self, source_id: i64) -> Result<Vec<String>> {
let mut stmt = self
.conn
.prepare("SELECT DISTINCT stream FROM segments WHERE source_id = ?1 ORDER BY stream")
.map_err(Error::sqlite("failed to query streams"))?;
let rows = stmt
.query_map([source_id], |row| row.get::<_, String>(0))
.map_err(Error::sqlite("failed to query streams"))?;
let mut out = Vec::new();
for row in rows {
out.push(row.map_err(Error::sqlite("failed to read stream name"))?);
}
Ok(out)
}
pub fn all_streams(&self, source_id: i64) -> Result<Vec<String>> {
let mut stmt = self
.conn
.prepare(
"SELECT stream FROM segments WHERE source_id = ?1 \
UNION \
SELECT stream FROM wal WHERE source_id = ?1 \
ORDER BY stream",
)
.map_err(Error::sqlite("failed to query all_streams"))?;
let rows = stmt
.query_map([source_id], |row| row.get::<_, String>(0))
.map_err(Error::sqlite("failed to query all_streams"))?;
let mut out = Vec::new();
for row in rows {
out.push(row.map_err(Error::sqlite("failed to read stream name"))?);
}
Ok(out)
}
pub fn read_wal(&self, source_id: i64, stream: &str) -> Result<Vec<WalRow>> {
let mut stmt = self
.conn
.prepare(
"SELECT stream, ts, wall_offset, row FROM wal \
WHERE source_id = ?1 AND stream = ?2 ORDER BY ts",
)
.map_err(Error::sqlite(format!("failed to query WAL for {stream}")))?;
Self::collect_wal_rows(&mut stmt, source_id, stream)
}
pub fn live_wal(&self, source_id: i64, stream: &str) -> Result<Vec<WalRow>> {
let mut stmt = self
.conn
.prepare(&format!(
"SELECT stream, ts, wall_offset, row FROM wal \
WHERE {LIVE_WAL_PREDICATE} ORDER BY ts"
))
.map_err(Error::sqlite(format!(
"failed to query live WAL for {stream}"
)))?;
Self::collect_wal_rows(&mut stmt, source_id, stream)
}
pub fn live_wal_span(&self, source_id: i64, stream: &str) -> Result<Span> {
self.query_span(
&format!("SELECT COUNT(*), MIN(ts), MAX(ts) FROM wal WHERE {LIVE_WAL_PREDICATE}"),
source_id,
stream,
)
.map_err(Error::sqlite(format!(
"failed to measure the live WAL for {stream}"
)))
}
pub fn segment_span(&self, source_id: i64, stream: &str) -> Result<(u64, Span)> {
let segments: i64 = self
.conn
.query_row(
"SELECT COUNT(*) FROM segments WHERE source_id = ?1 AND stream = ?2",
rusqlite::params![source_id, stream],
|row| row.get(0),
)
.map_err(Error::sqlite(format!(
"failed to count segments for {stream}"
)))?;
let span = self
.query_span(
"SELECT COALESCE(SUM(rows), 0), MIN(first_ts), MAX(last_ts) FROM segments \
WHERE source_id = ?1 AND stream = ?2",
source_id,
stream,
)
.map_err(Error::sqlite(format!(
"failed to measure the segments of {stream}"
)))?;
Ok((segments as u64, span))
}
fn query_span(&self, sql: &str, source_id: i64, stream: &str) -> rusqlite::Result<Span> {
self.conn
.query_row(sql, rusqlite::params![source_id, stream], |row| {
Ok(Span {
rows: row.get::<_, i64>(0)? as u64,
first_ts: row.get::<_, Option<i64>>(1)?,
last_ts: row.get::<_, Option<i64>>(2)?,
})
})
}
fn collect_wal_rows(
stmt: &mut rusqlite::Statement<'_>,
source_id: i64,
stream: &str,
) -> Result<Vec<WalRow>> {
let rows = stmt
.query_map(rusqlite::params![source_id, stream], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, Vec<u8>>(3)?,
))
})
.map_err(Error::sqlite(format!(
"failed to query WAL rows for {stream}"
)))?;
let mut out = Vec::new();
for row in rows {
let (stream, ts, wall_offset, data) = row.map_err(Error::sqlite(format!(
"failed to read WAL row for {stream}"
)))?;
out.push(WalRow {
stream,
ts,
wall_offset,
row: data,
});
}
Ok(out)
}
pub fn segment_sizes(&self, source_id: i64) -> Result<Vec<(i64, u64)>> {
let mut stmt = self
.conn
.prepare(
"SELECT last_ts, length(bytes) FROM segments \
WHERE source_id = ?1 ORDER BY last_ts",
)
.map_err(Error::sqlite("failed to query segment sizes"))?;
let rows = stmt
.query_map([source_id], |row| {
Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)? as u64))
})
.map_err(Error::sqlite("failed to query segment sizes"))?;
let mut out = Vec::new();
for row in rows {
out.push(row.map_err(Error::sqlite("failed to read a segment size"))?);
}
Ok(out)
}
pub fn page_stats(&self) -> Result<PageStats> {
Ok(PageStats {
pages: self.pragma_u32("page_count")?,
free: self.pragma_u32("freelist_count")?,
page_size: self.pragma_u32("page_size")?,
})
}
pub fn archive_bytes(&self) -> Result<u64> {
let s = self.page_stats()?;
Ok(s.pages as u64 * s.page_size as u64)
}
pub fn vacuum_into(&self, dest: &Path) -> Result<()> {
let dest = dest
.to_str()
.ok_or_else(|| format!("dump destination {} is not valid UTF-8", dest.display()))?;
self.conn
.execute("VACUUM INTO ?1", [dest])
.map_err(Error::sqlite(format!("failed to write the dump to {dest}")))?;
Ok(())
}
#[cfg_attr(not(feature = "write"), allow(dead_code))]
pub(crate) fn source_time_span(&self, source_id: i64) -> Result<(Option<i64>, Option<i64>)> {
self.conn
.query_row(
&format!(
"SELECT MIN(first_ts), MAX(last_ts) FROM ( \
SELECT first_ts, last_ts FROM segments WHERE source_id = ?1 \
UNION ALL \
SELECT ts, ts FROM wal WHERE source_id = ?1 \
AND ({LIVE_WAL_PREDICATE_FOR_ROW}))"
),
[source_id],
|row| Ok((row.get::<_, Option<i64>>(0)?, row.get::<_, Option<i64>>(1)?)),
)
.map_err(Error::sqlite(format!(
"failed to measure source {source_id}"
)))
}
#[cfg(test)]
pub(crate) fn user_table_names(&self) -> Result<Vec<String>> {
let mut stmt = self
.conn
.prepare(
"SELECT name FROM sqlite_master WHERE type = 'table' AND name NOT LIKE 'sqlite_%'",
)
.map_err(Error::sqlite("failed to list tables"))?;
let rows = stmt
.query_map([], |r| r.get::<_, String>(0))
.map_err(Error::sqlite("failed to list tables"))?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::sqlite("failed to list tables"))
}
pub fn verify(&self, depth: Depth) -> Result<Report> {
let mut problems = Vec::new();
let check = match depth {
Depth::Quick => "PRAGMA quick_check",
Depth::Full => "PRAGMA integrity_check",
};
let mut stmt = self
.conn
.prepare(check)
.map_err(Error::sqlite("failed to check the archive"))?;
let lines = stmt
.query_map([], |row| row.get::<_, String>(0))
.map_err(Error::sqlite("failed to check the archive"))?;
for line in lines {
let line = line.map_err(Error::sqlite("failed to read a check result"))?;
if line != "ok" {
problems.push(Problem::Corrupt(line));
}
}
let corrupt = |e: Error, problems: &mut Vec<Problem>| -> Result<()> {
if e.sqlite_code() == Some(rusqlite::ErrorCode::DatabaseCorrupt) {
problems.push(Problem::Corrupt(e.to_string()));
Ok(())
} else {
Err(e)
}
};
let references = (|| -> Result<Vec<String>> {
let mut stmt = self
.conn
.prepare("PRAGMA foreign_key_check")
.map_err(Error::sqlite("failed to check references"))?;
let violations = stmt
.query_map([], |row| {
Ok(format!(
"{} row {:?} -> {}",
row.get::<_, String>(0)?,
row.get::<_, Option<i64>>(1)?,
row.get::<_, String>(2)?
))
})
.map_err(Error::sqlite("failed to check references"))?;
violations
.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::sqlite("failed to read a reference violation"))
})();
match references {
Ok(violations) => problems.extend(violations.into_iter().map(Problem::ForeignKey)),
Err(e) => corrupt(e, &mut problems)?,
}
let walked = self.read_snapshot(|db| {
let rows = db.read_sources()?;
let mut streams = 0usize;
let mut segments = 0usize;
let mut wal_rows = 0usize;
for src in &rows {
for stream in db.all_streams(src.id)? {
streams += 1;
for seg in db.read_segment_meta(src.id, &stream)? {
segments += 1;
let (seq, meta) = seg;
let mut bad = Vec::new();
if meta.first_ts > meta.last_ts {
bad.push(format!(
"spans [{}, {}], which runs backwards",
meta.first_ts, meta.last_ts
));
}
if meta.rows == 0 {
bad.push("claims no rows; such a segment should not exist".to_string());
}
for detail in bad {
problems.push(Problem::Segment {
source_id: src.id,
stream: stream.clone(),
seq,
detail,
});
}
}
let all = db.total_wal_rows(src.id, &stream)?;
let live = db.live_wal_span(src.id, &stream)?.rows as usize;
wal_rows += all;
if all > live {
problems.push(Problem::UnreadableWalRows {
source_id: src.id,
stream: stream.clone(),
rows: all - live,
});
}
}
}
Ok((rows.len(), streams, segments, wal_rows))
});
let (sources, streams, segments, wal_rows) = match walked {
Ok(counts) => counts,
Err(e) => {
corrupt(e, &mut problems)?;
(0, 0, 0, 0)
}
};
Ok(Report {
sources,
streams,
segments,
wal_rows,
problems,
})
}
pub(crate) fn total_wal_rows(&self, source_id: i64, stream: &str) -> Result<usize> {
self.conn
.query_row(
"SELECT COUNT(*) FROM wal WHERE source_id = ?1 AND stream = ?2",
rusqlite::params![source_id, stream],
|row| row.get::<_, i64>(0),
)
.map(|n| n as usize)
.map_err(Error::sqlite(format!(
"failed to count WAL rows for {stream}"
)))
}
pub fn stream_bytes(&self, source_id: i64, stream: &str) -> Result<u64> {
self.conn
.query_row(
"SELECT COALESCE(SUM(length(bytes)), 0) FROM segments \
WHERE source_id = ?1 AND stream = ?2",
rusqlite::params![source_id, stream],
|row| row.get::<_, i64>(0),
)
.map(|n| n as u64)
.map_err(Error::sqlite(format!("failed to size {stream}")))
}
pub fn sealed_watermarks(&self) -> Result<BTreeMap<i64, BTreeMap<String, i64>>> {
let mut stmt = self
.conn
.prepare(
"SELECT source_id, stream, MAX(last_ts) FROM segments \
GROUP BY source_id, stream",
)
.map_err(Error::sqlite("failed to query sealed watermarks"))?;
let rows = stmt
.query_map([], |row| {
Ok((
(row.get::<_, i64>(0)?, row.get::<_, String>(1)?),
row.get::<_, i64>(2)?,
))
})
.map_err(Error::sqlite("failed to query sealed watermarks"))?;
let mut out: BTreeMap<i64, BTreeMap<String, i64>> = BTreeMap::new();
for row in rows {
let ((source_id, stream), ts) =
row.map_err(Error::sqlite("failed to read a sealed watermark"))?;
out.entry(source_id).or_default().insert(stream, ts);
}
Ok(out)
}
#[cfg_attr(not(feature = "write"), allow(dead_code))]
pub(crate) fn next_seqs(&self) -> Result<BTreeMap<(i64, String), u64>> {
let mut stmt = self
.conn
.prepare(
"SELECT source_id, stream, MAX(seq) + 1 FROM segments \
GROUP BY source_id, stream",
)
.map_err(Error::sqlite("failed to query segment sequences"))?;
let rows = stmt
.query_map([], |row| {
Ok((
(row.get::<_, i64>(0)?, row.get::<_, String>(1)?),
row.get::<_, i64>(2)? as u64,
))
})
.map_err(Error::sqlite("failed to query segment sequences"))?;
let mut out = BTreeMap::new();
for row in rows {
let (key, next) = row.map_err(Error::sqlite("failed to read a segment sequence"))?;
out.insert(key, next);
}
Ok(out)
}
pub fn source_metadata(&self, source_id: i64) -> Result<BTreeMap<String, String>> {
let encoded: String = self
.conn
.query_row(
"SELECT metadata FROM sources WHERE id = ?1",
[source_id],
|row| row.get(0),
)
.map_err(Error::sqlite(format!(
"failed to read the metadata of source {source_id}"
)))?;
serde_json::from_str(&encoded)
.map_err(|e| Error::Message(format!("source {source_id} has invalid metadata: {e}")))
}
pub(crate) fn pragma_u32(&self, name: &str) -> Result<u32> {
let value = self.pragma_i64(name)?;
u32::try_from(value)
.map_err(|_| Error::Message(format!("pragma {name} is {value}, not a u32")))
}
pub(crate) fn pragma_i64(&self, name: &str) -> Result<i64> {
self.conn
.pragma_query_value(None, name, |row| row.get(0))
.map_err(Error::sqlite(format!("failed to read pragma {name}")))
}
pub fn read_clock_offsets(&self, source_id: i64) -> Result<Vec<(i64, i64)>> {
let mut stmt = self
.conn
.prepare("SELECT ts, offset_ns FROM clock_offsets WHERE source_id = ?1 ORDER BY ts")
.map_err(Error::sqlite("failed to query clock offsets"))?;
let rows = stmt
.query_map([source_id], |row| {
Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?))
})
.map_err(Error::sqlite("failed to query clock offsets"))?;
let mut out = Vec::new();
for row in rows {
let (ts, offset) = row.map_err(Error::sqlite("failed to read clock offset"))?;
out.push((ts, offset));
}
Ok(out)
}
pub fn read_caller_rows(
&self,
source_id: i64,
stream: &str,
start: i64,
end: i64,
) -> Result<Vec<CallerRow>> {
let mut stmt = self
.conn
.prepare(
"SELECT ts, blob FROM caller_rows \
WHERE source_id = ?1 AND stream = ?2 AND ts >= ?3 AND ts <= ?4 \
ORDER BY ts, rowid",
)
.map_err(Error::sqlite(format!(
"failed to query caller rows for {stream}"
)))?;
let rows = stmt
.query_map(rusqlite::params![source_id, stream, start, end], |row| {
Ok(CallerRow {
ts: row.get::<_, i64>(0)?,
blob: row.get::<_, Vec<u8>>(1)?,
})
})
.map_err(Error::sqlite(format!(
"failed to query caller rows for {stream}"
)))?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::sqlite(format!(
"failed to read caller rows for {stream}"
)))
}
pub fn caller_row_streams(&self, source_id: i64) -> Result<Vec<String>> {
let mut stmt = self
.conn
.prepare("SELECT DISTINCT stream FROM caller_rows WHERE source_id = ?1 ORDER BY stream")
.map_err(Error::sqlite("failed to query caller row streams"))?;
let rows = stmt
.query_map([source_id], |row| row.get::<_, String>(0))
.map_err(Error::sqlite("failed to query caller row streams"))?;
rows.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::sqlite("failed to read a caller row stream name"))
}
#[cfg(test)]
fn pragma_string(&self, name: &str) -> Result<String> {
self.conn
.pragma_query_value(None, name, |row| row.get(0))
.map_err(Error::sqlite(format!("failed to read pragma {name}")))
}
}
pub struct ArchiveMut {
db: Archive,
#[cfg(any(test, feature = "test-support"))]
commits: std::cell::Cell<u64>,
}
impl std::ops::Deref for ArchiveMut {
type Target = Archive;
fn deref(&self) -> &Archive {
&self.db
}
}
impl std::fmt::Debug for ArchiveMut {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ArchiveMut").field("db", &self.db).finish()
}
}
impl ArchiveMut {
fn wrap(db: Archive) -> Self {
ArchiveMut {
db,
#[cfg(any(test, feature = "test-support"))]
commits: std::cell::Cell::new(0),
}
}
pub fn open(path: &Path) -> Result<Self> {
let db = Archive::open_with_cache(path, READER_CACHE_SIZE_KIB, true)?;
Ok(ArchiveMut::wrap(db))
}
pub fn transaction<T>(&mut self, f: impl FnOnce(&Transaction<'_>) -> Result<T>) -> Result<T> {
let tx = Transaction {
tx: self
.db
.conn
.transaction()
.map_err(Error::sqlite("failed to begin transaction"))?,
};
let out = f(&tx)?;
tx.tx
.commit()
.map_err(Error::sqlite("failed to commit transaction"))?;
#[cfg(any(test, feature = "test-support"))]
self.commits.set(self.commits.get() + 1);
Ok(out)
}
pub fn insert_source(&mut self, meta: &SourceMeta) -> Result<i64> {
insert_source_sql(&self.db.conn, meta, None)
}
pub fn insert_source_with_uuid(
&mut self,
meta: &SourceMeta,
uuid: Option<&str>,
) -> Result<i64> {
insert_source_sql(&self.db.conn, meta, uuid)
}
pub fn insert_segment(
&mut self,
source_id: i64,
stream: &str,
seq: u64,
meta: &SegmentMeta,
bytes: &[u8],
) -> Result<()> {
insert_segment_sql(&self.db.conn, source_id, stream, seq, meta, bytes, None)
}
pub fn insert_segment_with_index(
&mut self,
source_id: i64,
stream: &str,
seq: u64,
meta: &SegmentMeta,
bytes: &[u8],
caller_index: Option<&[u8]>,
) -> Result<()> {
insert_segment_sql(
&self.db.conn,
source_id,
stream,
seq,
meta,
bytes,
caller_index,
)
}
pub fn insert_wal_rows(&mut self, source_id: i64, rows: &[WalRow]) -> Result<()> {
self.transaction(|tx| tx.insert_wal_rows(source_id, rows))
}
pub fn insert_wal_rows_batch(&mut self, ticks: &[(i64, Vec<WalRow>)]) -> Result<()> {
self.transaction(|tx| {
for (source_id, rows) in ticks {
tx.insert_wal_rows(*source_id, rows)?;
}
Ok(())
})
}
pub fn prune_wal(&mut self, source_id: i64, stream: &str, upto_ts: i64) -> Result<usize> {
self.db
.conn
.execute(
"DELETE FROM wal WHERE source_id = ?1 AND stream = ?2 AND ts <= ?3",
rusqlite::params![source_id, stream, upto_ts],
)
.map_err(Error::sqlite(format!("failed to prune WAL for {stream}")))
}
pub fn evict_before(&mut self, source_id: i64, cutoff_ts: i64) -> Result<Evicted> {
self.evict(
source_id,
"DELETE FROM segments WHERE source_id = ?1 AND last_ts < ?2",
"DELETE FROM wal WHERE source_id = ?1 AND ts < ?2",
cutoff_ts,
)
}
pub fn evict_streams_before(
&mut self,
source_id: i64,
cutoff_ts: i64,
evict: &dyn Fn(&str) -> bool,
) -> Result<Evicted> {
let mut streams: Vec<String> = self.db.all_streams(source_id)?;
for name in self.db.caller_row_streams(source_id)? {
if !streams.contains(&name) {
streams.push(name);
}
}
let streams: Vec<String> = streams.into_iter().filter(|s| evict(s)).collect();
self.transaction(|tx| {
let mut total = Evicted::default();
for stream in &streams {
let params = rusqlite::params![source_id, stream, cutoff_ts];
total.live_rows += tx
.tx
.query_row(
&format!("SELECT COUNT(*) FROM wal WHERE {LIVE_WAL_PREDICATE} AND ts < ?3"),
params,
|row| row.get::<_, i64>(0),
)
.map_err(Error::sqlite(format!(
"failed to count live {stream} rows before eviction"
)))? as usize;
total.segments += tx
.tx
.execute(
"DELETE FROM segments \
WHERE source_id = ?1 AND stream = ?2 AND last_ts < ?3",
params,
)
.map_err(Error::sqlite(format!("failed to evict {stream} segments")))?;
total.wal_rows += tx
.tx
.execute(
"DELETE FROM wal WHERE source_id = ?1 AND stream = ?2 AND ts < ?3",
params,
)
.map_err(Error::sqlite(format!("failed to evict {stream} WAL rows")))?;
total.caller_rows += tx
.tx
.execute(
"DELETE FROM caller_rows \
WHERE source_id = ?1 AND stream = ?2 AND ts < ?3",
params,
)
.map_err(Error::sqlite(format!(
"failed to evict {stream} caller rows"
)))?;
}
tx.tx
.execute(
"DELETE FROM clock_offsets WHERE source_id = ?1 AND ts < COALESCE( \
(SELECT MIN(oldest) FROM ( \
SELECT MIN(first_ts) AS oldest FROM segments WHERE source_id = ?1 \
UNION ALL \
SELECT MIN(ts) FROM wal WHERE source_id = ?1)), \
?2)",
rusqlite::params![source_id, cutoff_ts],
)
.map_err(Error::sqlite("failed to evict clock offsets"))?;
Ok(total)
})
}
fn evict(
&mut self,
source_id: i64,
segments_sql: &str,
wal_sql: &str,
cutoff_ts: i64,
) -> Result<Evicted> {
self.transaction(|tx| {
let params = rusqlite::params![source_id, cutoff_ts];
let live_rows: i64 = tx
.tx
.query_row(
&format!(
"SELECT COUNT(*) FROM wal WHERE source_id = ?1 AND ts < ?2 \
AND ({LIVE_WAL_PREDICATE_FOR_ROW})"
),
params,
|row| row.get(0),
)
.map_err(Error::sqlite("failed to count live rows before eviction"))?;
let segments = tx
.tx
.execute(segments_sql, params)
.map_err(Error::sqlite("failed to evict segments"))?;
let wal_rows = tx
.tx
.execute(wal_sql, params)
.map_err(Error::sqlite("failed to evict WAL rows"))?;
tx.tx
.execute(
"DELETE FROM clock_offsets WHERE source_id = ?1 AND ts < ?2",
rusqlite::params![source_id, cutoff_ts],
)
.map_err(Error::sqlite("failed to evict clock offsets"))?;
let caller_rows = tx
.tx
.execute(
"DELETE FROM caller_rows WHERE source_id = ?1 AND ts < ?2",
params,
)
.map_err(Error::sqlite("failed to evict caller rows"))?;
Ok(Evicted {
segments,
wal_rows,
live_rows: live_rows as usize,
caller_rows,
})
})
}
pub fn insert_caller_rows(
&mut self,
source_id: i64,
stream: &str,
rows: &[CallerRow],
) -> Result<()> {
self.transaction(|tx| tx.insert_caller_rows(source_id, stream, rows))
}
pub fn mark_complete(&mut self, source_id: i64) -> Result<()> {
self.transaction(|tx| tx.mark_complete(source_id))
}
pub fn patch_source_metadata(
&mut self,
source_id: i64,
patch: &BTreeMap<String, String>,
) -> Result<()> {
let mut metadata = self.db.source_metadata(source_id)?;
for (k, v) in patch {
metadata.insert(k.clone(), v.clone());
}
self.update_source_metadata(source_id, &metadata)
}
pub fn update_source_metadata(
&mut self,
source_id: i64,
metadata: &BTreeMap<String, String>,
) -> Result<()> {
let encoded = serde_json::to_string(metadata)
.map_err(|e| Error::Message(format!("failed to encode source metadata: {e}")))?;
let changed = self
.db
.conn
.execute(
"UPDATE sources SET metadata = ?1 WHERE id = ?2",
rusqlite::params![encoded, source_id],
)
.map_err(Error::sqlite("failed to update source metadata"))?;
if changed == 0 {
return Err(Error::Message(format!("no source with id {source_id}")));
}
Ok(())
}
pub fn checkpoint_passive(&mut self) -> Result<()> {
self.db
.conn
.execute_batch("PRAGMA wal_checkpoint(PASSIVE);")
.map_err(Error::sqlite("failed to checkpoint the WAL"))
}
#[cfg(any(test, feature = "test-support"))]
pub fn set_busy_timeout(&self, timeout: std::time::Duration) -> Result<()> {
self.db
.conn
.busy_timeout(timeout)
.map_err(Error::sqlite("failed to set busy_timeout"))
}
pub fn create(path: &Path) -> Result<Self> {
Archive::create_with_page_size(path, PAGE_SIZE)
}
pub fn create_in_memory() -> Result<Self> {
let conn = Connection::open_in_memory()
.map_err(Error::sqlite("failed to open an in-memory database"))?;
let db = Archive { conn };
db.set_pragma("page_size", PAGE_SIZE)?;
db.apply_connection_pragmas(WRITER_CACHE_SIZE_KIB)?;
db.conn
.execute_batch(SCHEMA_SQL)
.map_err(Error::sqlite("failed to create archive schema"))?;
db.stamp_header()?;
Ok(ArchiveMut::wrap(db))
}
#[cfg_attr(not(feature = "write"), allow(dead_code))]
pub(crate) fn open_for_write(path: &Path) -> Result<Self> {
let db = Archive::open_with_cache(path, WRITER_CACHE_SIZE_KIB, false)?;
Ok(ArchiveMut::wrap(db))
}
#[cfg(any(test, feature = "test-support"))]
pub fn commits(&self) -> u64 {
self.commits.get()
}
pub fn incremental_vacuum(&mut self, pages: u32) -> Result<()> {
let fail = |e| format!("failed to reclaim {pages} pages: {e}");
let mut stmt = self
.conn
.prepare(&format!("PRAGMA incremental_vacuum({pages})"))
.map_err(fail)?;
let mut rows = stmt.query([]).map_err(fail)?;
while rows.next().map_err(fail)?.is_some() {}
Ok(())
}
}
pub struct Transaction<'a> {
tx: rusqlite::Transaction<'a>,
}
impl Transaction<'_> {
pub fn insert_source(&self, meta: &SourceMeta) -> Result<i64> {
insert_source_sql(&self.tx, meta, None)
}
pub fn insert_source_with_uuid(&self, meta: &SourceMeta, uuid: Option<&str>) -> Result<i64> {
insert_source_sql(&self.tx, meta, uuid)
}
pub fn insert_segment(
&self,
source_id: i64,
stream: &str,
seq: u64,
meta: &SegmentMeta,
bytes: &[u8],
) -> Result<()> {
insert_segment_sql(&self.tx, source_id, stream, seq, meta, bytes, None)
}
pub fn delete_segment(&self, source_id: i64, stream: &str, seq: u64) -> Result<()> {
self.tx
.execute(
"DELETE FROM segments WHERE source_id = ?1 AND stream = ?2 AND seq = ?3",
rusqlite::params![source_id, stream, seq as i64],
)
.map_err(Error::sqlite(format!(
"failed to drop segment {stream}#{seq}"
)))?;
Ok(())
}
pub fn insert_segment_with_index(
&self,
source_id: i64,
stream: &str,
seq: u64,
meta: &SegmentMeta,
bytes: &[u8],
caller_index: Option<&[u8]>,
) -> Result<()> {
insert_segment_sql(&self.tx, source_id, stream, seq, meta, bytes, caller_index)
}
pub fn insert_wal_rows(&self, source_id: i64, rows: &[WalRow]) -> Result<()> {
let mut stmt = self
.tx
.prepare(
"INSERT INTO wal(source_id, stream, ts, wall_offset, row) \
VALUES (?1, ?2, ?3, ?4, ?5)",
)
.map_err(Error::sqlite("failed to prepare WAL insert"))?;
for r in rows {
stmt.execute(rusqlite::params![
source_id,
r.stream,
r.ts,
r.wall_offset,
r.row,
])
.map_err(Error::sqlite(format!(
"failed to insert WAL row for {}",
r.stream
)))?;
}
Ok(())
}
pub fn insert_clock_offset(&self, source_id: i64, ts: i64, offset_ns: i64) -> Result<()> {
self.tx
.execute(
"INSERT OR IGNORE INTO clock_offsets(source_id, ts, offset_ns) \
VALUES (?1, ?2, ?3)",
rusqlite::params![source_id, ts, offset_ns],
)
.map_err(Error::sqlite("failed to insert clock offset"))?;
Ok(())
}
pub fn insert_caller_rows(
&self,
source_id: i64,
stream: &str,
rows: &[CallerRow],
) -> Result<()> {
let mut stmt = self
.tx
.prepare("INSERT INTO caller_rows(source_id, stream, ts, blob) VALUES (?1, ?2, ?3, ?4)")
.map_err(Error::sqlite("failed to prepare caller row insert"))?;
for r in rows {
stmt.execute(rusqlite::params![source_id, stream, r.ts, r.blob])
.map_err(Error::sqlite(format!(
"failed to insert a caller row for {stream}"
)))?;
}
Ok(())
}
pub fn mark_incomplete(&self, source_id: i64) -> Result<()> {
let changed = self
.tx
.execute("UPDATE sources SET complete = 0 WHERE id = ?1", [source_id])
.map_err(Error::sqlite(format!(
"failed to reopen source {source_id}"
)))?;
if changed == 0 {
return Err(Error::Message(format!("no source with id {source_id}")));
}
Ok(())
}
pub fn mark_complete(&self, source_id: i64) -> Result<()> {
self.tx
.execute("UPDATE sources SET complete = 1 WHERE id = ?1", [source_id])
.map_err(Error::sqlite(format!(
"failed to mark source {source_id} complete"
)))?;
Ok(())
}
}
fn insert_source_sql(conn: &Connection, meta: &SourceMeta, uuid: Option<&str>) -> Result<i64> {
let labels = serde_json::to_string(&meta.labels)
.map_err(|e| Error::Message(format!("failed to encode source labels: {e}")))?;
let metadata = serde_json::to_string(&meta.metadata)
.map_err(|e| Error::Message(format!("failed to encode source metadata: {e}")))?;
let uuid = match uuid {
Some(u) => u.to_string(),
None => mint_uuid(conn)?,
};
conn.execute(
"INSERT INTO sources(labels, metadata, complete, clock_anchor_wall_ns, uuid) \
VALUES (?1, ?2, 0, ?3, ?4)",
rusqlite::params![labels, metadata, meta.clock_anchor_wall_ns, uuid],
)
.map_err(Error::sqlite("failed to insert source"))?;
Ok(conn.last_insert_rowid())
}
fn mint_uuid(conn: &Connection) -> Result<String> {
let mut b: Vec<u8> = conn
.query_row("SELECT randomblob(16)", [], |row| row.get(0))
.map_err(Error::sqlite("failed to mint a uuid"))?;
if b.len() != 16 {
return Err(Error::Message(format!(
"randomblob(16) returned {} bytes",
b.len()
)));
}
b[6] = (b[6] & 0x0f) | 0x40;
b[8] = (b[8] & 0x3f) | 0x80;
let hex: String = b.iter().map(|x| format!("{x:02x}")).collect();
Ok(format!(
"{}-{}-{}-{}-{}",
&hex[0..8],
&hex[8..12],
&hex[12..16],
&hex[16..20],
&hex[20..32]
))
}
fn has_column(conn: &Connection, table: &str, column: &str) -> Result<bool> {
let mut stmt = conn
.prepare(&format!("PRAGMA table_info({table})"))
.map_err(Error::sqlite(format!("failed to inspect {table}")))?;
let names = stmt
.query_map([], |row| row.get::<_, String>(1))
.map_err(Error::sqlite(format!("failed to inspect {table}")))?;
for name in names {
if name.map_err(Error::sqlite(format!("failed to inspect {table}")))? == column {
return Ok(true);
}
}
Ok(false)
}
fn read_segment_indexes_sql(
conn: &Connection,
source_id: i64,
stream: &str,
) -> Result<Vec<(u64, Option<Vec<u8>>)>> {
let column = if has_column(conn, "segments", "caller_index")? {
"caller_index"
} else {
"NULL"
};
let mut stmt = conn
.prepare(&format!(
"SELECT seq, {column} FROM segments \
WHERE source_id = ?1 AND stream = ?2 ORDER BY seq"
))
.map_err(Error::sqlite(format!("failed to query {stream} indexes")))?;
let rows = stmt
.query_map(rusqlite::params![source_id, stream], |row| {
Ok((
row.get::<_, i64>(0)? as u64,
row.get::<_, Option<Vec<u8>>>(1)?,
))
})
.map_err(Error::sqlite(format!("failed to query {stream} indexes")))?;
let mut out = Vec::new();
for row in rows {
out.push(row.map_err(Error::sqlite(format!("failed to read a {stream} index")))?);
}
Ok(out)
}
fn insert_segment_sql(
conn: &Connection,
source_id: i64,
stream: &str,
seq: u64,
meta: &SegmentMeta,
bytes: &[u8],
caller_index: Option<&[u8]>,
) -> Result<()> {
conn.execute(
"INSERT INTO segments(source_id, stream, seq, rows, first_ts, last_ts, bytes, \
caller_index) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
rusqlite::params![
source_id,
stream,
seq as i64,
meta.rows as i64,
meta.first_ts,
meta.last_ts,
bytes,
caller_index,
],
)
.map_err(Error::sqlite(format!(
"failed to insert segment {stream}#{seq}"
)))?;
Ok(())
}
const SCHEMA_SQL: &str = "
CREATE TABLE sources(
id INTEGER PRIMARY KEY,
labels TEXT NOT NULL, -- JSON
metadata TEXT NOT NULL, -- JSON
complete INTEGER NOT NULL DEFAULT 0,
clock_anchor_wall_ns INTEGER NOT NULL,
-- The source's identity across files: minted at insert, carried verbatim
-- by every copy, so whether two archives hold the same source is a
-- comparison rather than a guess from labels. NULL only in archives
-- written before the column existed; readers treat that as unknown.
uuid TEXT
);
CREATE TABLE segments(
source_id INTEGER NOT NULL REFERENCES sources(id),
stream TEXT NOT NULL,
seq INTEGER NOT NULL,
rows INTEGER NOT NULL,
first_ts INTEGER NOT NULL,
last_ts INTEGER NOT NULL,
bytes BLOB NOT NULL,
-- The caller's index over this segment, stored and never read. The
-- catalog knows a segment's stream and span and nothing about its
-- contents; this is where a caller that needs more puts it, so that an
-- archive with an index is still one file. Named for whose it is.
caller_index BLOB,
PRIMARY KEY (source_id, stream, seq)
);
-- The catalog half of the design: it makes retention
-- (`WHERE last_ts < cutoff`) and range reads indexed lookups rather than
-- scans. `live_wal`'s subquery (`SELECT MAX(last_ts) FROM segments WHERE
-- source_id = ? AND stream = ?`) already uses it — confirmed by
-- `EXPLAIN QUERY PLAN` during review — so this is not a speculative index
-- sitting unused; keep it maintained.
CREATE INDEX segments_by_time ON segments(source_id, stream, last_ts);
CREATE TABLE wal(
source_id INTEGER NOT NULL REFERENCES sources(id),
stream TEXT NOT NULL,
ts INTEGER NOT NULL,
wall_offset INTEGER NOT NULL,
row BLOB NOT NULL,
PRIMARY KEY (source_id, stream, ts)
);
-- `PRIMARY KEY (source_id, ts)`: at most one observation per source per
-- timestamp. Two seal batches landing on the same `last_ts` - streams sealed
-- one at a time, which is an ordinary thing to do - used to write two rows with
-- different offsets, and a consumer could not read the series uniformly. The
-- constraint is what makes that impossible rather than merely unlikely; see
-- `insert_clock_offset`, which is `INSERT OR IGNORE` so the first observation
-- at a timestamp wins.
CREATE TABLE clock_offsets(
source_id INTEGER NOT NULL REFERENCES sources(id),
ts INTEGER NOT NULL,
offset_ns INTEGER NOT NULL,
PRIMARY KEY (source_id, ts)
);
-- The caller's time-keyed store: never decoded, never merged, evicted by
-- timestamp. No primary key, so several rows may share a timestamp and
-- rowid keeps their insertion order.
CREATE TABLE caller_rows(
source_id INTEGER NOT NULL REFERENCES sources(id),
stream TEXT NOT NULL,
ts INTEGER NOT NULL,
blob BLOB NOT NULL
);
CREATE INDEX caller_rows_by_time ON caller_rows(source_id, stream, ts);
";
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn create_applies_the_one_way_pragmas() {
let dir = tempfile::tempdir().unwrap();
let db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
assert_eq!(db.pragma_u32("auto_vacuum").unwrap(), 2, "INCREMENTAL");
assert_eq!(db.pragma_string("journal_mode").unwrap(), "wal");
}
#[test]
fn create_honors_the_page_size_it_is_given() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.dendro");
let db = Archive::create_with_page_size(&path, 8192).unwrap();
assert_eq!(db.pragma_u32("page_size").unwrap(), 8192);
drop(db);
assert_eq!(
Archive::open(&path)
.unwrap()
.pragma_u32("page_size")
.unwrap(),
8192
);
}
#[test]
fn open_reapplies_the_per_connection_pragmas() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.dendro");
drop(ArchiveMut::create(&path).unwrap());
let db = ArchiveMut::open(&path).unwrap();
assert_eq!(
db.pragma_u32("wal_autocheckpoint").unwrap(),
WAL_AUTOCHECKPOINT_BYTES / PAGE_SIZE,
"byte-denominated cap, not SQLite's 1000-page default"
);
assert_eq!(
db.pragma_i64("cache_size").unwrap(),
READER_CACHE_SIZE_KIB as i64,
"256 MiB reader cache, not SQLite's -2000 default"
);
}
#[test]
fn a_writing_connection_does_not_get_the_reader_cache() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.dendro");
let created = ArchiveMut::create(&path).unwrap();
assert_eq!(
created.pragma_i64("cache_size").unwrap(),
WRITER_CACHE_SIZE_KIB as i64,
"a created (writing) connection takes the writer cache"
);
assert_eq!(
Archive::open(&path)
.unwrap()
.pragma_i64("cache_size")
.unwrap(),
READER_CACHE_SIZE_KIB as i64,
"an opened connection still takes the reader cache"
);
}
#[test]
fn effective_config_matches_what_was_measured() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.dendro");
let created = ArchiveMut::create(&path).unwrap();
let reopened = Archive::open(&path).unwrap();
assert_eq!(PAGE_SIZE, 4096, "the swept and measured page size");
for (which, db) in [("created", &*created), ("reopened", &reopened)] {
assert_eq!(db.pragma_u32("page_size").unwrap(), 4096, "{which}");
assert_eq!(
db.pragma_u32("synchronous").unwrap(),
2,
"{which}: FULL (3 is EXTRA)"
);
}
}
#[test]
fn schema_round_trips_a_source() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let id = db
.insert_source(&SourceMeta {
labels: [("host".to_string(), "h1".to_string())]
.into_iter()
.collect(),
metadata: [("source".to_string(), "weather".to_string())]
.into_iter()
.collect(),
clock_anchor_wall_ns: 1_700_000_000_000_000_000,
})
.unwrap();
let got = db.read_sources().unwrap();
assert_eq!(got.len(), 1);
assert_eq!(got[0].id, id);
assert_eq!(got[0].meta.labels["host"], "h1");
assert_eq!(got[0].meta.metadata["source"], "weather");
assert_eq!(got[0].meta.clock_anchor_wall_ns, 1_700_000_000_000_000_000);
assert!(!got[0].complete, "a fresh source is not complete");
}
#[test]
fn segments_read_back_in_seq_order_not_insertion_order() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let rid = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
db.insert_segment(
rid,
"cpu_usage",
1,
&SegmentMeta {
rows: 10,
first_ts: 90,
last_ts: 99,
},
b"seq-one-bytes",
)
.unwrap();
db.insert_segment(
rid,
"cpu_usage",
0,
&SegmentMeta {
rows: 5,
first_ts: 100,
last_ts: 200,
},
b"seq-zero-bytes",
)
.unwrap();
let got = db.read_segments(rid, "cpu_usage").unwrap();
assert_eq!(got.len(), 2);
assert_eq!(
got[0].seq, 0,
"seq 0 must come first despite being inserted second"
);
assert_eq!(got[0].bytes, b"seq-zero-bytes");
assert_eq!(got[1].seq, 1);
assert_eq!(got[1].bytes, b"seq-one-bytes");
}
#[test]
fn total_rows_sums_across_segments() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let rid = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
db.insert_segment(
rid,
"cpu_usage",
0,
&SegmentMeta {
rows: 3,
first_ts: 0,
last_ts: 29,
},
b"a",
)
.unwrap();
db.insert_segment(
rid,
"cpu_usage",
1,
&SegmentMeta {
rows: 2,
first_ts: 30,
last_ts: 49,
},
b"b",
)
.unwrap();
assert_eq!(db.total_rows(rid, "cpu_usage").unwrap(), 5);
}
#[test]
fn streams_lists_each_stream_once() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let rid = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
let meta = SegmentMeta {
rows: 1,
first_ts: 0,
last_ts: 9,
};
db.insert_segment(rid, "cpu_usage", 0, &meta, b"a").unwrap();
db.insert_segment(rid, "cpu_usage", 1, &meta, b"b").unwrap();
db.insert_segment(rid, "blockio", 0, &meta, b"c").unwrap();
assert_eq!(db.streams(rid).unwrap(), vec!["blockio", "cpu_usage"]);
}
#[test]
fn all_streams_includes_a_stream_that_has_never_sealed() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let rid = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
db.insert_segment(
rid,
"cpu_usage",
0,
&SegmentMeta {
rows: 1,
first_ts: 0,
last_ts: 9,
},
b"a",
)
.unwrap();
db.insert_wal_rows(rid, &[wal_row("drivehealth", 5)])
.unwrap();
assert_eq!(
db.streams(rid).unwrap(),
vec!["cpu_usage"],
"streams() legitimately does not see the WAL-only stream"
);
assert_eq!(
db.all_streams(rid).unwrap(),
vec!["cpu_usage", "drivehealth"],
"all_streams() must see it — this is the whole point of the accessor"
);
}
#[test]
fn an_unbounded_cutoff_evicts_everything_rather_than_nothing() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let id = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
db.insert_segment(
id,
"s",
0,
&SegmentMeta {
rows: 1,
first_ts: 10,
last_ts: 20,
},
b"x",
)
.unwrap();
db.insert_wal_rows(
id,
&[WalRow {
stream: "s".to_string(),
ts: 30,
wall_offset: 0,
row: vec![1],
}],
)
.unwrap();
let evicted = db.evict_before(id, i64::MAX).unwrap();
assert_eq!(evicted.segments, 1);
assert_eq!(evicted.wal_rows, 1);
assert!(db.all_streams(id).unwrap().is_empty());
}
#[test]
fn timestamps_span_the_whole_signed_range() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let id = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: i64::MIN,
})
.unwrap();
let extremes = [i64::MIN, -1_000_000_000, 0, 1, i64::MAX];
for ts in extremes {
db.insert_wal_rows(
id,
&[WalRow {
stream: "s".to_string(),
ts,
wall_offset: 0,
row: vec![1],
}],
)
.unwrap();
}
let back: Vec<i64> = db.live_wal(id, "s").unwrap().iter().map(|r| r.ts).collect();
assert_eq!(back, extremes, "every one round-trips, in order");
assert_eq!(
db.read_sources().unwrap()[0].meta.clock_anchor_wall_ns,
i64::MIN
);
db.insert_segment(
id,
"s",
0,
&SegmentMeta {
rows: 3,
first_ts: i64::MIN,
last_ts: 0,
},
b"x",
)
.unwrap();
let live: Vec<i64> = db.live_wal(id, "s").unwrap().iter().map(|r| r.ts).collect();
assert_eq!(live, vec![1, i64::MAX]);
}
#[test]
fn a_row_at_timestamp_zero_is_live() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let id = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
for ts in [0i64, 1, 2] {
db.insert_wal_rows(
id,
&[WalRow {
stream: "s".to_string(),
ts,
wall_offset: 0,
row: vec![1],
}],
)
.unwrap();
}
let live: Vec<i64> = db.live_wal(id, "s").unwrap().iter().map(|r| r.ts).collect();
assert_eq!(live, vec![0, 1, 2], "ts=0 is a timestamp like any other");
assert_eq!(db.live_wal_span(id, "s").unwrap().rows, 3);
db.insert_segment(
id,
"s",
0,
&SegmentMeta {
rows: 2,
first_ts: 0,
last_ts: 1,
},
b"x",
)
.unwrap();
let live: Vec<i64> = db.live_wal(id, "s").unwrap().iter().map(|r| r.ts).collect();
assert_eq!(live, vec![2]);
}
#[test]
fn per_stream_eviction_leaves_the_streams_it_was_not_given() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let id = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
let sm = |first_ts, last_ts| SegmentMeta {
rows: 1,
first_ts,
last_ts,
};
for stream in ["debug/a", "debug/b", "metric/c"] {
db.insert_segment(id, stream, 0, &sm(0, 9), b"old").unwrap();
db.insert_segment(id, stream, 1, &sm(100, 109), b"new")
.unwrap();
}
let evicted = db
.evict_streams_before(id, 50, &|s: &str| s.starts_with("debug/"))
.unwrap();
assert_eq!(
evicted.segments, 2,
"one old segment from each debug stream"
);
for stream in ["debug/a", "debug/b"] {
let got = db.read_segments(id, stream).unwrap();
assert_eq!(got.len(), 1, "{stream} keeps only its newer segment");
assert_eq!(got[0].bytes, b"new");
}
assert_eq!(
db.read_segments(id, "metric/c").unwrap().len(),
2,
"a stream the predicate rejected must be untouched"
);
}
#[test]
fn per_stream_eviction_takes_the_wal_rows_its_segments_shadowed() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let id = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
db.insert_segment(
id,
"s",
0,
&SegmentMeta {
rows: 2,
first_ts: 0,
last_ts: 20,
},
b"sealed",
)
.unwrap();
for ts in [10i64, 20, 30] {
db.insert_wal_rows(
id,
&[WalRow {
stream: "s".to_string(),
ts,
wall_offset: 0,
row: vec![1],
}],
)
.unwrap();
}
assert_eq!(db.live_wal(id, "s").unwrap().len(), 1, "only ts=30 is live");
db.evict_streams_before(id, 25, &|_| true).unwrap();
assert!(db.read_segments(id, "s").unwrap().is_empty());
let live = db.live_wal(id, "s").unwrap();
assert_eq!(
live.len(),
1,
"the shadowed rows went with the segment; only ts=30 remains, \
and it was live before"
);
assert_eq!(live[0].ts, 30);
}
#[test]
fn segment_sizes_are_oldest_first_and_measure_the_payload() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let id = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
db.insert_segment(
id,
"b",
0,
&SegmentMeta {
rows: 1,
first_ts: 90,
last_ts: 99,
},
&[0u8; 300],
)
.unwrap();
db.insert_segment(
id,
"a",
0,
&SegmentMeta {
rows: 1,
first_ts: 0,
last_ts: 9,
},
&[0u8; 100],
)
.unwrap();
assert_eq!(db.segment_sizes(id).unwrap(), vec![(9, 100), (99, 300)]);
}
#[test]
fn archive_bytes_counts_the_file_including_pages_not_yet_reclaimed() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let id = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
let empty = db.archive_bytes().unwrap();
for seq in 0..40i64 {
db.insert_segment(
id,
"s",
seq as u64,
&SegmentMeta {
rows: 1,
first_ts: seq * 10,
last_ts: seq * 10 + 9,
},
&[7u8; 4096],
)
.unwrap();
}
let full = db.archive_bytes().unwrap();
assert!(
full > empty,
"{empty} -> {full}: writing must grow the file"
);
db.evict_before(id, i64::MAX).unwrap();
assert_eq!(
db.archive_bytes().unwrap(),
full,
"eviction frees pages for reuse but does not return them - that is \
what `incremental_vacuum` is for, and a size cap has to know it"
);
}
#[test]
fn segments_are_scoped_per_source() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let meta = |labels: &str| SourceMeta {
labels: [("host".to_string(), labels.to_string())]
.into_iter()
.collect(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
};
let r1 = db.insert_source(&meta("h1")).unwrap();
let r2 = db.insert_source(&meta("h2")).unwrap();
let sm = SegmentMeta {
rows: 1,
first_ts: 0,
last_ts: 9,
};
db.insert_segment(r1, "cpu_usage", 0, &sm, b"r1-bytes")
.unwrap();
db.insert_segment(r2, "cpu_usage", 0, &sm, b"r2-bytes")
.unwrap();
let got1 = db.read_segments(r1, "cpu_usage").unwrap();
assert_eq!(got1.len(), 1);
assert_eq!(got1[0].bytes, b"r1-bytes");
let got2 = db.read_segments(r2, "cpu_usage").unwrap();
assert_eq!(got2.len(), 1);
assert_eq!(got2[0].bytes, b"r2-bytes");
assert_eq!(db.total_rows(r1, "cpu_usage").unwrap(), 1);
assert_eq!(db.streams(r1).unwrap(), vec!["cpu_usage"]);
}
#[test]
fn create_refuses_an_existing_file() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.dendro");
drop(ArchiveMut::create(&path).unwrap());
let err = match ArchiveMut::create(&path) {
Ok(_) => panic!("create clobbered an existing archive"),
Err(e) => e,
};
let already_exists = OpenOptions::new()
.write(true)
.create_new(true)
.open(&path)
.unwrap_err();
assert_eq!(already_exists.kind(), std::io::ErrorKind::AlreadyExists);
assert!(
err.to_string().contains(&already_exists.to_string()),
"{err:?} should carry the AlreadyExists error from create_new"
);
}
fn wal_test_db() -> (tempfile::TempDir, ArchiveMut, i64) {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let rid = db
.insert_source(&SourceMeta {
labels: BTreeMap::new(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
})
.unwrap();
(dir, db, rid)
}
fn wal_row(stream: &str, ts: i64) -> WalRow {
WalRow {
stream: stream.to_string(),
ts,
wall_offset: ts,
row: format!("row@{ts}").into_bytes(),
}
}
#[test]
fn live_wal_excludes_rows_already_covered_by_a_sealed_segment() {
let (_dir, mut db, rid) = wal_test_db();
db.insert_segment(
rid,
"cpu_usage",
0,
&SegmentMeta {
rows: 3,
first_ts: 10,
last_ts: 30,
},
b"sealed-bytes",
)
.unwrap();
db.insert_wal_rows(
rid,
&[
wal_row("cpu_usage", 10),
wal_row("cpu_usage", 20),
wal_row("cpu_usage", 30),
wal_row("cpu_usage", 40),
],
)
.unwrap();
let live = db.live_wal(rid, "cpu_usage").unwrap();
assert_eq!(live.len(), 1, "only ts=40 is past the sealed watermark");
assert_eq!(live[0].ts, 40);
let all = db.read_wal(rid, "cpu_usage").unwrap();
assert_eq!(all.len(), 4, "the raw WAL table is untouched by sealing");
}
#[test]
fn live_wal_returns_everything_when_nothing_has_sealed() {
let (_dir, mut db, rid) = wal_test_db();
db.insert_wal_rows(
rid,
&[wal_row("drivehealth", 5), wal_row("drivehealth", 15)],
)
.unwrap();
let live = db.live_wal(rid, "drivehealth").unwrap();
assert_eq!(
live.len(),
2,
"no segments sealed yet, so every row is live"
);
assert_eq!(live[0].ts, 5);
assert_eq!(live[1].ts, 15);
}
#[test]
fn live_wal_watermark_is_scoped_to_its_own_stream_and_source() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let meta = |host: &str| SourceMeta {
labels: [("host".to_string(), host.to_string())]
.into_iter()
.collect(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
};
let r1 = db.insert_source(&meta("h1")).unwrap();
let r2 = db.insert_source(&meta("h2")).unwrap();
db.insert_segment(
r1,
"cpu_usage",
0,
&SegmentMeta {
rows: 3,
first_ts: 10,
last_ts: 30,
},
b"r1-cpu_usage-sealed",
)
.unwrap();
db.insert_wal_rows(r1, &[wal_row("cpu_usage", 40), wal_row("blockio", 5)])
.unwrap();
db.insert_wal_rows(r2, &[wal_row("cpu_usage", 5)]).unwrap();
let r1_blockio = db.live_wal(r1, "blockio").unwrap();
assert_eq!(
r1_blockio.len(),
1,
"blockio in r1 must not inherit cpu_usage's sealed watermark"
);
assert_eq!(r1_blockio[0].ts, 5);
let r2_cpu_usage = db.live_wal(r2, "cpu_usage").unwrap();
assert_eq!(
r2_cpu_usage.len(),
1,
"cpu_usage in r2 must not inherit r1's sealed watermark"
);
assert_eq!(r2_cpu_usage[0].ts, 5);
let r1_cpu_usage = db.live_wal(r1, "cpu_usage").unwrap();
assert_eq!(r1_cpu_usage.len(), 1);
assert_eq!(r1_cpu_usage[0].ts, 40);
}
#[test]
fn segment_span_summarizes_the_catalog_without_reading_a_blob() {
let (_dir, mut db, rid) = wal_test_db();
db.insert_segment(
rid,
"cpu_usage",
0,
&SegmentMeta {
rows: 3,
first_ts: 10,
last_ts: 29,
},
b"not-parquet",
)
.unwrap();
db.insert_segment(
rid,
"cpu_usage",
1,
&SegmentMeta {
rows: 2,
first_ts: 30,
last_ts: 49,
},
b"not-parquet-either",
)
.unwrap();
db.insert_segment(
rid,
"blockio",
0,
&SegmentMeta {
rows: 99,
first_ts: 0,
last_ts: 99,
},
b"nor-this",
)
.unwrap();
let (segments, span) = db.segment_span(rid, "cpu_usage").unwrap();
assert_eq!(segments, 2);
assert_eq!(span.rows, 5, "the SUM of the catalog's row counts");
assert_eq!((span.first_ts, span.last_ts), (Some(10), Some(49)));
let (segments, span) = db.segment_span(rid, "drivehealth").unwrap();
assert_eq!(segments, 0);
assert_eq!(span.rows, 0);
assert_eq!((span.first_ts, span.last_ts), (None, None));
}
#[test]
fn live_wal_span_counts_the_same_rows_live_wal_returns() {
let (_dir, mut db, rid) = wal_test_db();
db.insert_segment(
rid,
"cpu_usage",
0,
&SegmentMeta {
rows: 3,
first_ts: 10,
last_ts: 30,
},
b"not-parquet",
)
.unwrap();
db.insert_wal_rows(
rid,
&[
wal_row("cpu_usage", 10),
wal_row("cpu_usage", 20),
wal_row("cpu_usage", 30),
wal_row("cpu_usage", 40),
wal_row("cpu_usage", 50),
],
)
.unwrap();
let span = db.live_wal_span(rid, "cpu_usage").unwrap();
assert_eq!(
span.rows,
db.live_wal(rid, "cpu_usage").unwrap().len() as u64,
"the depth must agree with the rows the reader will replay"
);
assert_eq!(span.rows, 2, "ts=40 and ts=50 are past the watermark");
assert_eq!((span.first_ts, span.last_ts), (Some(40), Some(50)));
db.insert_wal_rows(rid, &[wal_row("drivehealth", 5)])
.unwrap();
let span = db.live_wal_span(rid, "drivehealth").unwrap();
assert_eq!(span.rows, 1);
assert_eq!((span.first_ts, span.last_ts), (Some(5), Some(5)));
}
#[test]
fn prune_is_idempotent_and_bounded_to_one_stream() {
let (_dir, mut db, rid) = wal_test_db();
db.insert_wal_rows(
rid,
&[
wal_row("cpu_usage", 10),
wal_row("cpu_usage", 20),
wal_row("blockio", 10),
wal_row("blockio", 20),
],
)
.unwrap();
let deleted = db.prune_wal(rid, "cpu_usage", 10).unwrap();
assert_eq!(deleted, 1, "only cpu_usage's ts<=10 row");
let deleted_again = db.prune_wal(rid, "cpu_usage", 10).unwrap();
assert_eq!(deleted_again, 0);
let blockio = db.read_wal(rid, "blockio").unwrap();
assert_eq!(blockio.len(), 2, "blockio untouched by cpu_usage's prune");
let cpu_usage = db.read_wal(rid, "cpu_usage").unwrap();
assert_eq!(cpu_usage.len(), 1, "cpu_usage's ts=10 row is gone");
assert_eq!(cpu_usage[0].ts, 20);
}
#[test]
fn insert_wal_rows_inserts_every_row_in_the_batch() {
let (_dir, mut db, rid) = wal_test_db();
let streams: Vec<WalRow> = (0..26)
.map(|i| wal_row(&format!("stream_{i}"), 100))
.collect();
db.insert_wal_rows(rid, &streams).unwrap();
for i in 0..26 {
let stream = format!("stream_{i}");
let rows = db.read_wal(rid, &stream).unwrap();
assert_eq!(rows.len(), 1, "{stream} should have its tick's row");
}
}
#[test]
fn insert_wal_rows_is_one_transaction_for_the_whole_tick() {
let (_dir, mut db, rid) = wal_test_db();
let err = db
.insert_wal_rows(
rid,
&[
wal_row("cpu_usage", 10),
wal_row("blockio", 10),
wal_row("cpu_usage", 10), ],
)
.expect_err("a PRIMARY KEY collision must fail the whole call");
assert!(
{
let text = err.to_string().to_lowercase();
text.contains("unique") || text.contains("constraint")
},
"{err:?} should name the PK collision, not some other failure"
);
assert_eq!(
db.read_wal(rid, "cpu_usage").unwrap().len(),
0,
"a failed tick must leave NO rows, not the one that would have committed alone"
);
assert_eq!(
db.read_wal(rid, "blockio").unwrap().len(),
0,
"blockio's collision-free row must also be rolled back"
);
}
#[test]
fn a_transaction_commits_the_whole_batch_or_none_of_it() {
let (_dir, mut db, rid) = wal_test_db();
let meta = |first_ts, last_ts| SegmentMeta {
rows: 1,
first_ts,
last_ts,
};
db.transaction(|tx| {
tx.insert_segment(rid, "cpu_usage", 0, &meta(10, 19), b"cpu-0")?;
tx.insert_segment(rid, "blockio", 0, &meta(10, 19), b"blk-0")?;
tx.insert_wal_rows(rid, &[wal_row("cpu_usage", 20)])
})
.unwrap();
assert_eq!(db.read_segments(rid, "cpu_usage").unwrap().len(), 1);
assert_eq!(db.read_segments(rid, "blockio").unwrap().len(), 1);
assert_eq!(db.read_wal(rid, "cpu_usage").unwrap().len(), 1);
let err = db
.transaction(|tx| {
tx.insert_segment(rid, "cpu_usage", 1, &meta(20, 29), b"cpu-1")?;
tx.insert_segment(rid, "blockio", 1, &meta(20, 29), b"blk-1")?;
tx.insert_segment(rid, "cpu_usage", 1, &meta(20, 29), b"dup")
})
.expect_err("a PRIMARY KEY collision must fail the whole batch");
assert!(
{
let text = err.to_string().to_lowercase();
text.contains("unique") || text.contains("constraint")
},
"{err:?} should name the PK collision, not some other failure"
);
assert_eq!(
db.read_segments(rid, "cpu_usage").unwrap().len(),
1,
"seq 1 must be rolled back, leaving only the committed seq 0"
);
assert_eq!(
db.read_segments(rid, "blockio").unwrap().len(),
1,
"blockio's collision-free insert must be rolled back too"
);
}
#[test]
fn wal_rows_are_scoped_per_source() {
let dir = tempfile::tempdir().unwrap();
let mut db = ArchiveMut::create(&dir.path().join("t.dendro")).unwrap();
let meta = |host: &str| SourceMeta {
labels: [("host".to_string(), host.to_string())]
.into_iter()
.collect(),
metadata: BTreeMap::new(),
clock_anchor_wall_ns: 0,
};
let r1 = db.insert_source(&meta("h1")).unwrap();
let r2 = db.insert_source(&meta("h2")).unwrap();
db.insert_wal_rows(r1, &[wal_row("cpu_usage", 10)]).unwrap();
db.insert_wal_rows(r2, &[wal_row("cpu_usage", 20)]).unwrap();
let got1 = db.read_wal(r1, "cpu_usage").unwrap();
assert_eq!(got1.len(), 1);
assert_eq!(got1[0].ts, 10);
let got2 = db.read_wal(r2, "cpu_usage").unwrap();
assert_eq!(got2.len(), 1);
assert_eq!(got2[0].ts, 20);
let live1 = db.live_wal(r1, "cpu_usage").unwrap();
assert_eq!(live1.len(), 1);
assert_eq!(live1[0].ts, 10);
}
}