use alloc::{
string::{String, ToString},
vec::Vec,
};
use core::cmp::Ordering;
use std::{
collections::BTreeMap,
path::{Path, PathBuf},
};
use io_replica::{
collection::ReplicaCollectionId,
hub::{ReplicaSourceBinding, ReplicaSourceId},
object::ReplicaHash,
placement::{ReplicaLevel, ReplicaLinkId},
};
use rusqlite::{Connection, OpenFlags, OptionalExtension, named_params};
use crate::{
client::{
PimdirBlobs, PimdirCollection, PimdirConflict, PimdirError, PimdirItem, PimdirParkedAction,
PimdirPendingAction, PimdirPlacement, binding_from_row, check_rename_cascades,
check_version_agreement, collection_row, conflict_row, load_pending_actions,
read_hash_algo, read_item_from_row, rows,
},
codec::{self, PimdirAction},
hash::{PimdirHashAlgo, PimdirHasher},
sql,
};
pub struct PimdirReader {
pub(super) conn: Connection,
pub(super) dir: PathBuf,
pub(super) blobs: PathBuf,
pub(super) hash: PimdirHashAlgo,
overlay: bool,
}
impl PimdirReader {
pub fn open(dir: impl AsRef<Path>) -> Result<Self, PimdirError> {
let dir = dir.as_ref();
let flags = OpenFlags::SQLITE_OPEN_READ_ONLY
| OpenFlags::SQLITE_OPEN_URI
| OpenFlags::SQLITE_OPEN_NO_MUTEX;
let conn = Connection::open_with_flags(dir.join("pimdir.db"), flags)?;
conn.execute_batch("PRAGMA busy_timeout = 30000;")?;
let version: i64 = conn.pragma_query_value(None, "user_version", |r| r.get(0))?;
match version {
version if version == sql::VERSION => {}
0 => return Err(PimdirError::Uncreated),
found => return Err(PimdirError::Version { found }),
}
check_version_agreement(&conn, version)?;
check_rename_cascades(&conn)?;
let hash = read_hash_algo(&conn, None)?;
Ok(Self::over(
conn,
dir.to_path_buf(),
dir.join("objects"),
hash,
))
}
pub(super) fn over(
conn: Connection,
dir: PathBuf,
blobs: PathBuf,
hash: PimdirHashAlgo,
) -> Self {
Self {
conn,
dir,
blobs,
hash,
overlay: false,
}
}
pub fn with_pending(mut self) -> Self {
self.overlay = true;
self
}
pub fn overlays_pending(&self) -> bool {
self.overlay
}
}
impl PimdirReader {
pub fn hash_algo(&self) -> PimdirHashAlgo {
self.hash
}
pub fn blobs(&self) -> PimdirBlobs {
PimdirBlobs {
root: self.blobs.clone(),
hash: self.hash,
}
}
pub fn hash(&self, bytes: &[u8]) -> ReplicaHash {
self.hash.hash(bytes)
}
pub fn hasher(&self) -> PimdirHasher {
self.hash.hasher()
}
pub fn collection_account(
&self,
collection: &str,
) -> Result<Option<Option<String>>, PimdirError> {
Ok(self
.conn
.query_row(
sql::LOAD_ACCOUNT,
named_params! { ":collection": collection },
|r| r.get::<_, Option<String>>(0),
)
.optional()?)
}
pub fn collection_kind(&self, collection: &str) -> Result<Option<String>, PimdirError> {
Ok(self
.conn
.query_row(
sql::LOAD_KIND,
named_params! { ":collection": collection },
|r| r.get::<_, String>(0),
)
.optional()?)
}
pub fn list_collections(&self) -> Result<Vec<PimdirCollection>, PimdirError> {
Ok(rows(&self.conn, sql::LIST_COLLECTIONS, [], collection_row)?)
}
pub fn list_collections_by_account(
&self,
account: Option<&str>,
) -> Result<Vec<PimdirCollection>, PimdirError> {
Ok(rows(
&self.conn,
sql::LIST_COLLECTIONS_BY_ACCOUNT,
named_params! { ":account": account },
collection_row,
)?)
}
pub fn list_accounts(&self) -> Result<Vec<String>, PimdirError> {
Ok(rows(&self.conn, sql::LIST_ACCOUNTS, [], |r| r.get(0))?)
}
pub fn link_placements(&self, link_id: &str) -> Result<Vec<PimdirPlacement>, PimdirError> {
Ok(rows(
&self.conn,
sql::LIST_LINK_PLACEMENTS,
named_params! { ":link_id": link_id },
|r| {
Ok(PimdirPlacement {
collection: r.get(0)?,
account: r.get(1)?,
seq: r.get(2)?,
link_id: ReplicaLinkId(link_id.to_string()),
object: r.get::<_, Option<String>>(3)?.map(ReplicaHash),
flags: codec::flags_from_json(r.get::<_, Option<String>>(4)?.as_deref()),
level: codec::level_from_int(r.get(5)?),
})
},
)?)
}
pub fn object_placements(&self, hash: &str) -> Result<Vec<PimdirPlacement>, PimdirError> {
Ok(rows(
&self.conn,
sql::LIST_OBJECT_PLACEMENTS,
named_params! { ":hash": hash },
|r| {
Ok(PimdirPlacement {
collection: r.get(0)?,
account: r.get(1)?,
seq: r.get(2)?,
link_id: ReplicaLinkId(r.get(3)?),
object: Some(ReplicaHash(hash.to_string())),
flags: codec::flags_from_json(r.get::<_, Option<String>>(4)?.as_deref()),
level: codec::level_from_int(r.get(5)?),
})
},
)?)
}
pub fn list_items(
&self,
collection: &str,
after: Option<&str>,
limit: usize,
) -> Result<Vec<PimdirItem>, PimdirError> {
let after = after.unwrap_or("");
self.overlaid(
collection,
limit,
|limit| {
Ok(rows(
&self.conn,
sql::LIST_ITEMS_PAGE,
named_params! {
":collection": collection,
":after": after,
":limit": limit as i64,
},
read_item_from_row,
)?)
},
|item| item.link_id.0.as_str() > after,
|left, right| left.link_id.0.cmp(&right.link_id.0),
)
}
pub fn list_items_page_asc(
&self,
collection: &str,
after: Option<(&str, i64)>,
limit: usize,
) -> Result<Vec<PimdirItem>, PimdirError> {
let (key, seq) = after.unwrap_or(("", 0));
self.sorted_page(
sql::LIST_ITEMS_PAGE_ASC,
collection,
Some((key, seq)),
limit,
false,
)
}
pub fn list_items_page_desc(
&self,
collection: &str,
after: Option<(&str, i64)>,
limit: usize,
) -> Result<Vec<PimdirItem>, PimdirError> {
self.sorted_page(sql::LIST_ITEMS_PAGE_DESC, collection, after, limit, true)
}
fn sorted_page(
&self,
statement: &str,
collection: &str,
after: Option<(&str, i64)>,
limit: usize,
descending: bool,
) -> Result<Vec<PimdirItem>, PimdirError> {
let after = after.map(|(key, seq)| (key.to_string(), seq));
self.overlaid(
collection,
limit,
|limit| {
Ok(rows(
&self.conn,
statement,
named_params! {
":collection": collection,
":after_key": after.as_ref().map(|(key, _)| key.as_str()),
":after_seq": after.as_ref().map(|(_, seq)| *seq).unwrap_or_default(),
":limit": limit as i64,
},
read_item_from_row,
)?)
},
|item| {
let here = (item.sort_key.as_str(), item.seq);
match &after {
None => true,
Some((key, seq)) if descending => here < (key.as_str(), *seq),
Some((key, seq)) => here > (key.as_str(), *seq),
}
},
|left, right| {
let order =
(left.sort_key.as_str(), left.seq).cmp(&(right.sort_key.as_str(), right.seq));
if descending { order.reverse() } else { order }
},
)
}
pub fn get_item(&self, collection: &str, seq: i64) -> Result<Option<PimdirItem>, PimdirError> {
let item = self.committed_item(collection, seq)?;
if !self.overlay {
return Ok(item);
}
let pending = self.pending(collection)?;
let item = match item {
Some(item) => Some(item),
None => match pending.arrivals.get(&seq) {
Some(from) => self.committed_item(from, seq)?,
None => None,
},
};
Ok(item.and_then(|item| fold(item, pending.edits.get(&seq))))
}
pub fn seq_for_link(
&self,
collection: &str,
link_id: &str,
) -> Result<Option<i64>, PimdirError> {
Ok(self
.conn
.query_row(
sql::SEQ_BY_LINK,
named_params! { ":collection": collection, ":link_id": link_id },
|row| row.get(0),
)
.optional()?)
}
pub fn item_bindings(
&self,
collection: &str,
link_id: &str,
) -> Result<BTreeMap<ReplicaSourceId, ReplicaSourceBinding>, PimdirError> {
Ok(rows(
&self.conn,
sql::ITEM_BINDINGS,
named_params! { ":collection": collection, ":link_id": link_id },
binding_from_row,
)?
.into_iter()
.map(|(_, source, binding)| (source, binding))
.collect())
}
pub fn list_conflicts(
&self,
account: Option<&str>,
) -> Result<Vec<PimdirConflict>, PimdirError> {
Ok(rows(
&self.conn,
sql::LIST_CONFLICTED_BINDINGS,
named_params! { ":account": account },
conflict_row,
)?)
}
pub fn distinct_sources(&self) -> Result<Vec<String>, PimdirError> {
Ok(rows(&self.conn, sql::LIST_SOURCES, [], |r| r.get(0))?)
}
pub fn count_items(&self, collection: &str) -> Result<u64, PimdirError> {
let count: i64 = self.conn.query_row(
sql::COUNT_ITEMS,
named_params! { ":collection": collection },
|r| r.get(0),
)?;
let mut count = count.max(0) as u64;
if !self.overlay {
return Ok(count);
}
let pending = self.pending(collection)?;
for (seq, edits) in &pending.edits {
let Some(item) = self.committed_item(collection, *seq)? else {
continue;
};
if fold(item, Some(edits)).is_none() {
count -= 1;
}
}
Ok(count + self.arrived(&pending)?.len() as u64)
}
pub fn list_retained(
&self,
collection: &ReplicaCollectionId,
after: Option<i64>,
limit: usize,
) -> Result<Vec<PimdirItem>, PimdirError> {
Ok(rows(
&self.conn,
sql::LIST_RETAINED_PAGE,
named_params! {
":collection": collection.0,
":after": after.unwrap_or(0),
":limit": limit as i64,
},
read_item_from_row,
)?)
}
pub fn count_retained(&self, collection: &ReplicaCollectionId) -> Result<i64, PimdirError> {
Ok(self.conn.query_row(
sql::COUNT_RETAINED,
named_params! { ":collection": collection.0 },
|r| r.get(0),
)?)
}
pub fn retained_bytes(&self) -> Result<u64, PimdirError> {
let bytes: i64 = self.conn.query_row(sql::RETAINED_BYTES, [], |r| r.get(0))?;
Ok(bytes.max(0) as u64)
}
pub fn generation(&self, collection: &str) -> Result<Option<i64>, PimdirError> {
Ok(self
.conn
.query_row(
sql::LOAD_GENERATION,
named_params! { ":collection": collection },
|r| r.get(0),
)
.optional()?)
}
pub fn queued_collections(&self) -> Result<Vec<String>, PimdirError> {
Ok(rows(&self.conn, sql::LIST_QUEUED_COLLECTIONS, [], |r| {
r.get(0)
})?)
}
pub fn pending_actions(
&self,
collection: &str,
) -> Result<Vec<PimdirPendingAction>, PimdirError> {
load_pending_actions(&self.conn, collection)
}
pub fn parked_actions(&self) -> Result<Vec<PimdirParkedAction>, PimdirError> {
Ok(rows(&self.conn, sql::LOAD_PARKED_ACTIONS, [], |r| {
Ok(PimdirParkedAction {
id: r.get(0)?,
created_at: r.get(1)?,
producer: r.get(2)?,
collection: r.get(3)?,
action: r.get(4)?,
payload: r.get(5)?,
attempts: r.get(6)?,
error: r.get(7)?,
})
})?)
}
}
#[derive(Debug, Default)]
struct PimdirPending {
edits: BTreeMap<i64, Vec<PimdirAction>>,
arrivals: BTreeMap<i64, String>,
creates: usize,
}
impl PimdirPending {
fn removals(&self) -> usize {
self.edits
.values()
.filter(|actions| {
actions
.iter()
.any(|action| matches!(action, PimdirAction::Remove { .. }))
})
.count()
}
}
impl PimdirReader {
fn committed_item(
&self,
collection: &str,
seq: i64,
) -> Result<Option<PimdirItem>, PimdirError> {
Ok(self
.conn
.query_row(
sql::GET_ITEM,
named_params! { ":collection": collection, ":seq": seq },
read_item_from_row,
)
.optional()?)
}
fn pending(&self, collection: &str) -> Result<PimdirPending, PimdirError> {
let mut queued = Vec::new();
for from in self.queued_collections()? {
for action in load_pending_actions(&self.conn, &from)? {
queued.push((from.clone(), action));
}
}
queued.sort_by_key(|(_, action)| action.id);
let mut pending = PimdirPending::default();
for (from, action) in queued {
let here = from == collection;
match &action.action {
PimdirAction::Add { .. } if here => pending.creates += 1,
PimdirAction::SetFlags { seq, .. }
| PimdirAction::Update { seq, .. }
| PimdirAction::Remove { seq }
if here =>
{
pending.edits.entry(*seq).or_default().push(action.action);
}
PimdirAction::Move { seq, to } => {
if here && to.0 != collection {
pending
.edits
.entry(*seq)
.or_default()
.push(PimdirAction::Remove { seq: *seq });
}
if !here && to.0 == collection {
pending.arrivals.insert(*seq, from);
}
}
PimdirAction::Copy { seq, to } if !here && to.0 == collection => {
pending.arrivals.insert(*seq, from);
}
_ => {}
}
}
Ok(pending)
}
fn arrived(&self, pending: &PimdirPending) -> Result<Vec<PimdirItem>, PimdirError> {
let mut items = Vec::new();
for (seq, from) in &pending.arrivals {
let Some(item) = self.committed_item(from, *seq)? else {
continue;
};
if let Some(item) = fold(item, pending.edits.get(seq)) {
items.push(item);
}
}
Ok(items)
}
fn overlaid(
&self,
collection: &str,
limit: usize,
fetch: impl Fn(usize) -> Result<Vec<PimdirItem>, PimdirError>,
inside: impl Fn(&PimdirItem) -> bool,
order: impl Fn(&PimdirItem, &PimdirItem) -> Ordering,
) -> Result<Vec<PimdirItem>, PimdirError> {
if !self.overlay {
return fetch(limit);
}
let pending = self.pending(collection)?;
let page = fetch(limit + pending.removals())?;
let mut items: Vec<PimdirItem> = page
.into_iter()
.filter_map(|item| {
let edits = pending.edits.get(&item.seq);
fold(item, edits)
})
.collect();
for item in self.arrived(&pending)? {
if inside(&item) && !items.iter().any(|held| held.seq == item.seq) {
items.push(item);
}
}
items.sort_by(order);
items.truncate(limit);
Ok(items)
}
pub fn pending_creates(
&self,
collection: &str,
) -> Result<Vec<PimdirPendingAction>, PimdirError> {
Ok(self
.pending_actions(collection)?
.into_iter()
.filter(|queued| matches!(queued.action, PimdirAction::Add { .. }))
.collect())
}
pub fn count_pending_creates(&self, collection: &str) -> Result<usize, PimdirError> {
Ok(self.pending_creates(collection)?.len())
}
}
fn fold(mut item: PimdirItem, actions: Option<&Vec<PimdirAction>>) -> Option<PimdirItem> {
for action in actions.into_iter().flatten() {
match action {
PimdirAction::SetFlags { flags, .. } => item.flags = flags.clone(),
PimdirAction::Update { object, meta, .. } => {
item.object = Some(object.clone());
if meta.is_some() {
item.meta = meta.clone();
}
item.level = ReplicaLevel::Full;
}
PimdirAction::Remove { .. } => return None,
_ => {}
}
}
Some(item)
}