use alloc::{
format,
string::{String, ToString},
vec,
vec::Vec,
};
use core::sync::atomic::{AtomicU64, Ordering};
use std::{
collections::{BTreeMap, BTreeSet, HashMap},
fmt, fs,
io::{self, ErrorKind, Write},
ops::{Deref, DerefMut},
path::{Path, PathBuf},
sync::Arc,
};
use io_replica::{
change::{ReplicaDropReason, ReplicaWriteOp},
client::ReplicaStorage,
collection::{ReplicaCheckpoint, ReplicaCollectionId},
coroutine::{ReplicaArg, ReplicaCoroutine, ReplicaCoroutineState, ReplicaYield},
hub::{ReplicaHub, ReplicaHubConflict, ReplicaHubItem, ReplicaSourceBinding, ReplicaSourceId},
mutate::{ReplicaMutate, ReplicaMutation},
object::{ReplicaHash, ReplicaObject},
placement::{
ReplicaBase, ReplicaFlags, ReplicaHandle, ReplicaLevel, ReplicaLinkId, ReplicaMeta,
ReplicaPlacement, ReplicaSortKey, ReplicaStatus,
},
storage::{ReplicaLoadScope, ReplicaLoaded},
};
use rusqlite::{
Connection, ErrorCode, OpenFlags, OptionalExtension, Params, Row, TransactionBehavior,
named_params, params, types::ToSql,
};
use crate::{
client::{lock::PimdirLock, reader::PimdirReader},
codec::{self, PimdirAction, PimdirActionError},
hash::{PimdirHashAlgo, PimdirHasher},
sql,
};
pub mod diagnostics;
pub mod reader;
mod lock;
pub struct PimdirStore {
reader: PimdirReader,
_lock: Option<Arc<PimdirLock>>,
account: Option<String>,
}
impl Deref for PimdirStore {
type Target = PimdirReader;
fn deref(&self) -> &Self::Target {
&self.reader
}
}
impl DerefMut for PimdirStore {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.reader
}
}
pub struct PimdirSourceStore {
store: PimdirStore,
source: ReplicaSourceId,
residual: HashMap<(ReplicaCollectionId, ReplicaHandle), ReplicaPlacement>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PimdirCollection {
pub id: String,
pub account: Option<String>,
pub kind: String,
pub name: String,
pub parent: Option<String>,
pub color: Option<String>,
pub description: Option<String>,
pub sort_order: Option<i64>,
pub generation: i64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PimdirPlacement {
pub collection: String,
pub account: Option<String>,
pub seq: i64,
pub link_id: ReplicaLinkId,
pub object: Option<ReplicaHash>,
pub flags: ReplicaFlags,
pub level: ReplicaLevel,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PimdirConflict {
pub collection: String,
pub link_id: ReplicaLinkId,
pub source: ReplicaSourceId,
pub handle: ReplicaHandle,
pub conflict_revision: Option<String>,
pub base_object: Option<ReplicaHash>,
pub object: Option<ReplicaHash>,
pub conflict_object: Option<ReplicaHash>,
}
fn conflict_row(r: &rusqlite::Row<'_>) -> rusqlite::Result<PimdirConflict> {
Ok(PimdirConflict {
collection: r.get(0)?,
link_id: ReplicaLinkId(r.get(1)?),
source: ReplicaSourceId(r.get(2)?),
handle: ReplicaHandle(r.get(3)?),
conflict_revision: r.get(4)?,
base_object: r.get::<_, Option<String>>(5)?.map(ReplicaHash),
object: r.get::<_, Option<String>>(6)?.map(ReplicaHash),
conflict_object: r.get::<_, Option<String>>(7)?.map(ReplicaHash),
})
}
fn collection_row(r: &rusqlite::Row<'_>) -> rusqlite::Result<PimdirCollection> {
Ok(PimdirCollection {
id: r.get(0)?,
account: r.get(1)?,
kind: r.get(2)?,
name: r.get(3)?,
parent: r.get(4)?,
color: r.get(5)?,
description: r.get(6)?,
sort_order: r.get(7)?,
generation: r.get(8)?,
})
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PimdirItem {
pub seq: i64,
pub link_id: ReplicaLinkId,
pub flags: ReplicaFlags,
pub meta: Option<ReplicaMeta>,
pub sort_key: String,
pub object: Option<ReplicaHash>,
pub level: ReplicaLevel,
pub retention: Option<PimdirRetention>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PimdirRetention {
pub at: String,
pub by: Option<String>,
pub size: Option<u64>,
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct PimdirPurgeReport {
pub items: usize,
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct PimdirGcReport {
pub objects: usize,
pub blobs: usize,
pub bytes: u64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PimdirPendingAction {
pub id: i64,
pub created_at: String,
pub producer: String,
pub action: PimdirAction,
pub attempts: i64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PimdirParkedAction {
pub id: i64,
pub created_at: String,
pub producer: String,
pub collection: String,
pub action: String,
pub payload: String,
pub attempts: i64,
pub error: String,
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct PimdirDrainReport {
pub applied: usize,
pub parked: usize,
pub skipped: usize,
}
impl PimdirStore {
pub fn open(dir: impl AsRef<Path>) -> Result<Self, PimdirError> {
Self::open_with_hash(dir, None)
}
pub fn open_with_hash(
dir: impl AsRef<Path>,
hash: Option<PimdirHashAlgo>,
) -> Result<Self, PimdirError> {
let dir = dir.as_ref();
fs::create_dir_all(dir)?;
let blobs = dir.join("objects");
fs::create_dir_all(&blobs)?;
let lock = PimdirLock::own(dir)?;
let mut conn = Connection::open(dir.join("pimdir.db"))?;
conn.execute_batch(
"PRAGMA journal_mode = WAL; PRAGMA foreign_keys = ON; PRAGMA busy_timeout = 30000;",
)?;
init_schema(&mut conn, hash.unwrap_or_default())?;
let hash = read_hash_algo(&conn, hash)?;
Ok(Self {
reader: PimdirReader::over(conn, dir.to_path_buf(), blobs, hash),
_lock: Some(lock),
account: None,
})
}
#[deprecated(
since = "0.3.0",
note = "use `PimdirReader::open`, which carries the reads and no write at all"
)]
pub fn open_read_only(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 {
reader: PimdirReader::over(conn, dir.to_path_buf(), dir.join("objects"), hash),
_lock: None,
account: None,
})
}
pub fn for_account(mut self, account: impl Into<String>) -> Self {
self.account = Some(account.into());
self
}
pub fn account(&self) -> Option<&str> {
self.account.as_deref()
}
pub fn for_source(self, source: impl Into<String>) -> PimdirSourceStore {
PimdirSourceStore {
store: self,
source: ReplicaSourceId(source.into()),
residual: HashMap::new(),
}
}
pub fn load_hub(&self, collection: &str) -> Result<ReplicaHub, PimdirError> {
Ok(load_hub(&self.conn, collection)?)
}
pub fn ensure_collection(&self, collection: &str, kind: &str) -> Result<(), PimdirError> {
self.conn.execute(
sql::SET_COLLECTION_KIND,
named_params! {
":collection": collection,
":account": self.account.as_deref(),
":kind": kind,
},
)?;
Ok(())
}
pub fn set_collection_account(
&self,
collection: &str,
account: Option<&str>,
) -> Result<(), PimdirError> {
self.conn.execute(
sql::SET_COLLECTION_ACCOUNT,
named_params! { ":collection": collection, ":account": account },
)?;
Ok(())
}
pub fn set_sort_key(
&self,
collection: &str,
link_id: &str,
sort_key: &str,
) -> Result<(), PimdirError> {
self.conn.execute(
sql::SET_SORT_KEY,
named_params! {
":collection": collection,
":link_id": link_id,
":sort_key": sort_key,
},
)?;
Ok(())
}
pub fn rename_collection(&self, collection: &str, new_id: &str) -> Result<(), PimdirError> {
self.conn.execute(
sql::RENAME_COLLECTION,
named_params! { ":collection": collection, ":new_id": new_id },
)?;
Ok(())
}
}
impl PimdirStore {
pub fn purge(
&mut self,
collection: &ReplicaCollectionId,
seq: i64,
) -> Result<bool, PimdirError> {
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
let pinned: Option<(Option<String>, Option<String>)> = tx
.prepare(sql::PURGE_ITEM)?
.query_row(
named_params! { ":collection": collection.0, ":seq": seq },
|row| Ok((row.get(0)?, row.get(1)?)),
)
.optional()?;
let Some((object, conflict_object)) = pinned else {
return Ok(false);
};
release_pins(&tx, [object, conflict_object].into_iter().flatten())?;
tx.commit().map_err(busy_or_sql)?;
Ok(true)
}
pub fn purge_retained_before(
&mut self,
cutoff: &str,
) -> Result<PimdirPurgeReport, PimdirError> {
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
let pinned: Vec<(Option<String>, Option<String>)> = rows(
&tx,
sql::PURGE_RETAINED_BEFORE,
named_params! { ":cutoff": cutoff },
|row| Ok((row.get(0)?, row.get(1)?)),
)?;
let items = pinned.len();
release_pins(
&tx,
pinned
.into_iter()
.flat_map(|(object, conflict)| [object, conflict])
.flatten(),
)?;
tx.commit().map_err(busy_or_sql)?;
Ok(PimdirPurgeReport { items })
}
}
impl PimdirStore {
pub fn collect_garbage(&mut self) -> Result<PimdirGcReport, PimdirError> {
let _staging = PimdirLock::collect(&self.dir)?;
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
let objects = tx.execute(sql::DELETE_GARBAGE_OBJECTS, [])?;
tx.commit().map_err(busy_or_sql)?;
let mut report = PimdirGcReport {
objects,
..Default::default()
};
let mut exists = self.conn.prepare(sql::OBJECT_EXISTS)?;
for blob in self.blobs().files()? {
if exists.exists(named_params! { ":hash": blob.hash })? {
continue;
}
fs::remove_file(&blob.path)?;
report.blobs += 1;
report.bytes += blob.size;
}
drop(exists);
Ok(report)
}
pub fn recompute_refcounts(&self) -> Result<usize, PimdirError> {
Ok(self.conn.execute(sql::RECOMPUTE_REFCOUNTS, [])?)
}
pub fn clear_dangling_bindings(&self) -> Result<usize, PimdirError> {
Ok(self.conn.execute(sql::DELETE_DANGLING_BINDINGS, [])?)
}
}
fn rows<T>(
conn: &Connection,
sql: &str,
params: impl Params,
map: impl FnMut(&Row) -> rusqlite::Result<T>,
) -> rusqlite::Result<Vec<T>> {
conn.prepare(sql)?.query_map(params, map)?.collect()
}
fn release_pins(
conn: &Connection,
hashes: impl Iterator<Item = String>,
) -> Result<(), PimdirError> {
let hashes: Vec<String> = hashes.collect();
if hashes.is_empty() {
return Ok(());
}
conn.execute(
sql::RELEASE_PINS,
named_params! { ":hashes": serde_json::to_string(&hashes)? },
)?;
Ok(())
}
impl PimdirStore {
pub fn cancel_action(dir: impl AsRef<Path>, id: i64) -> Result<bool, PimdirError> {
let dir = dir.as_ref();
if !dir.join("pimdir.db").is_file() {
return Err(PimdirError::Uncreated);
}
Self::open(dir)?.drop_action(id)
}
pub fn drop_action(&mut self, id: i64) -> Result<bool, PimdirError> {
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
let hash: Option<Option<String>> = tx
.query_row(sql::LOAD_ACTION_ROW, named_params! { ":id": id }, |r| {
r.get(1)
})
.optional()?;
let Some(hash) = hash else {
return Ok(false);
};
tx.execute(sql::CANCEL_ACTION, named_params! { ":id": id })?;
release_pins(&tx, hash.into_iter())?;
tx.commit().map_err(busy_or_sql)?;
Ok(true)
}
pub fn fail_action(&self, id: i64, error: Option<&str>) -> Result<(), PimdirError> {
let Some(error) = error else {
self.conn
.execute(sql::BUMP_ATTEMPTS, named_params! { ":id": id })?;
return Ok(());
};
let attempts: Option<i64> = self
.conn
.query_row(sql::LOAD_ACTION_ROW, named_params! { ":id": id }, |r| {
r.get(0)
})
.optional()?;
if let Some(attempts) = attempts {
self.conn.execute(
sql::PARK_ACTION,
named_params! { ":id": id, ":attempts": attempts + 1, ":error": error },
)?;
}
Ok(())
}
}
impl PimdirSourceStore {
pub fn source(&self) -> &str {
&self.source.0
}
pub fn for_account(mut self, account: impl Into<String>) -> Self {
self.store = self.store.for_account(account);
self
}
pub fn write_rekeyed(
&mut self,
collection: &str,
ops: Vec<ReplicaWriteOp>,
) -> Result<i64, PimdirError> {
stage_blobs(&self.store.reader.blobs, &ops)?;
let tx = self
.store
.reader
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
apply_ops(
&tx,
&self.store.reader.blobs,
&self.source,
self.store.account.as_deref(),
&mut self.residual,
ops,
)?;
tx.execute(
sql::ENSURE_COLLECTION,
named_params! { ":collection": collection, ":account": self.store.account.as_deref() },
)?;
let generation: i64 = tx.query_row(
sql::BUMP_GENERATION,
named_params! { ":collection": collection },
|r| r.get(0),
)?;
tx.commit().map_err(busy_or_sql)?;
Ok(generation)
}
pub fn drain_collection(&mut self, collection: &str) -> Result<PimdirDrainReport, PimdirError> {
let pending: Vec<QueueRow> = rows(
&self.store.reader.conn,
sql::LOAD_PENDING_ACTIONS,
named_params! { ":collection": collection },
|r| {
Ok(QueueRow {
id: r.get(0)?,
action: r.get(3)?,
payload: r.get(4)?,
object_hash: r.get(5)?,
})
},
)?;
let mut report = PimdirDrainReport::default();
for row in pending {
let action = match codec::action_from_payload(&row.action, &row.payload) {
Ok(action) => action,
Err(err) => {
self.fail_action(row.id, Some(&err.to_string()))?;
report.parked += 1;
continue;
}
};
if matches!(action, PimdirAction::Unknown { .. }) {
report.skipped += 1;
continue;
}
match self.apply_queued(collection, &row, &action) {
Ok(None) => report.applied += 1,
Ok(Some(PimdirRefusal::Skip)) => report.skipped += 1,
Ok(Some(PimdirRefusal::Park(reason))) => {
self.fail_action(row.id, Some(&reason))?;
report.parked += 1;
}
Err(err) => {
self.store
.conn
.execute(sql::BUMP_ATTEMPTS, named_params! { ":id": row.id })?;
return Err(err);
}
}
}
Ok(report)
}
fn apply_queued(
&mut self,
collection: &str,
row: &QueueRow,
action: &PimdirAction,
) -> Result<Option<PimdirRefusal>, PimdirError> {
let tx = self
.store
.reader
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
let claimed = tx
.prepare(sql::CLAIM_ACTION)?
.query_row(named_params! { ":id": row.id }, |r| r.get::<_, i64>(0))
.optional()?;
if claimed.is_none() {
return Ok(None);
}
let ops = match stage_action(&tx, &self.source, collection, row.id, action)? {
Ok(ops) => ops,
Err(refusal) => return Ok(Some(refusal)),
};
apply_ops(
&tx,
&self.store.reader.blobs,
&self.source,
self.store.account.as_deref(),
&mut self.residual,
ops,
)?;
if let Some(hash) = &row.object_hash {
tx.execute(
sql::ADJUST_REFCOUNT,
named_params! { ":delta": -1, ":hash": hash },
)?;
}
tx.commit().map_err(busy_or_sql)?;
Ok(None)
}
}
impl ReplicaStorage for PimdirSourceStore {
type Error = PimdirError;
fn load(
&self,
collection: &ReplicaCollectionId,
scope: &ReplicaLoadScope,
) -> Result<ReplicaLoaded, Self::Error> {
let hub = match scope {
ReplicaLoadScope::All => load_hub(&self.store.reader.conn, &collection.0)?,
ReplicaLoadScope::Links(links) => {
let links: Vec<String> = links.iter().map(|l| l.0.clone()).collect();
load_hub_by_link(&self.store.reader.conn, &collection.0, &links)?
}
ReplicaLoadScope::Handles(handles) => {
let mut links = Vec::new();
for handle in handles {
let link = self
.store
.reader
.conn
.query_row(
sql::LINK_FOR_HANDLE,
named_params! {
":collection": collection.0,
":source": self.source.0,
":handle": handle.0,
},
|r| r.get::<_, String>(0),
)
.optional()?;
links.extend(link);
}
load_hub_by_link(&self.store.reader.conn, &collection.0, &links)?
}
};
let mut placements = hub.project(collection, &self.source);
placements.extend(
self.residual
.values()
.filter(|p| &p.collection == collection)
.cloned(),
);
let checkpoint = self
.store
.reader
.conn
.query_row(
sql::LOAD_CHECKPOINT,
named_params! { ":collection": collection.0, ":source": self.source.0 },
|r| r.get::<_, Option<Vec<u8>>>(0),
)
.optional()?
.flatten()
.map(ReplicaCheckpoint);
Ok(ReplicaLoaded {
placements,
checkpoint,
})
}
fn lookup_objects(
&self,
links: &[ReplicaLinkId],
) -> Result<BTreeMap<ReplicaLinkId, ReplicaHash>, Self::Error> {
let ids: Vec<&str> = links.iter().map(|l| l.0.as_str()).collect();
let json = serde_json::to_string(&ids)?;
let found = rows(
&self.store.reader.conn,
sql::LOOKUP_OBJECTS,
named_params! { ":links": json, ":account": self.store.account.as_deref() },
|r| {
Ok((
ReplicaLinkId(r.get::<_, String>(0)?),
ReplicaHash(r.get::<_, String>(1)?),
))
},
)?;
let mut map: BTreeMap<ReplicaLinkId, ReplicaHash> = found.into_iter().collect();
for placement in self.residual.values() {
if let (Some(link), Some(object)) = (&placement.link_id, &placement.object) {
if links.contains(link) {
map.entry(link.clone()).or_insert_with(|| object.clone());
}
}
}
Ok(map)
}
fn write(&mut self, ops: Vec<ReplicaWriteOp>) -> Result<(), Self::Error> {
stage_blobs(&self.store.reader.blobs, &ops)?;
let tx = self
.store
.reader
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
apply_ops(
&tx,
&self.store.reader.blobs,
&self.source,
self.store.account.as_deref(),
&mut self.residual,
ops,
)?;
tx.commit().map_err(busy_or_sql)?;
Ok(())
}
}
impl Deref for PimdirSourceStore {
type Target = PimdirStore;
fn deref(&self) -> &Self::Target {
&self.store
}
}
impl DerefMut for PimdirSourceStore {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.store
}
}
struct QueueRow {
id: i64,
action: String,
payload: String,
object_hash: Option<String>,
}
fn load_pending_actions(
conn: &Connection,
collection: &str,
) -> Result<Vec<PimdirPendingAction>, PimdirError> {
let pending = rows(
conn,
sql::LOAD_PENDING_ACTIONS,
named_params! { ":collection": collection },
|r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
r.get::<_, String>(3)?,
r.get::<_, String>(4)?,
r.get::<_, i64>(6)?,
))
},
)?;
let mut actions = Vec::new();
for (id, created_at, producer, kind, payload, attempts) in pending {
actions.push(PimdirPendingAction {
id,
created_at,
producer,
action: codec::action_from_payload(&kind, &payload)?,
attempts,
});
}
Ok(actions)
}
enum PimdirRefusal {
Park(String),
Skip,
}
fn stage_action(
tx: &Connection,
source: &ReplicaSourceId,
collection: &str,
row_id: i64,
action: &PimdirAction,
) -> Result<Result<Vec<ReplicaWriteOp>, PimdirRefusal>, PimdirError> {
let collection_id = ReplicaCollectionId(collection.to_string());
if let PimdirAction::Add {
link_id,
flags,
object,
meta,
handle,
} = action
{
let link = link_id
.clone()
.or_else(|| object.as_ref().map(|hash| ReplicaLinkId(hash.0.clone())));
let Some(link) = link else {
return Ok(Err(PimdirRefusal::Park(
"add carries neither link_id nor object".to_string(),
)));
};
let live = tx
.query_row(
sql::LIVE_ITEM_FOR_LINK,
named_params! { ":collection": collection, ":link_id": link.0 },
|r| r.get::<_, i64>(0),
)
.optional()?;
if live.is_some() {
return Ok(Err(PimdirRefusal::Park(format!(
"link id already present: {}",
link.0
))));
}
let level = match (object, meta) {
(Some(_), _) => ReplicaLevel::Full,
(None, Some(_)) => ReplicaLevel::Meta,
(None, None) => ReplicaLevel::Probed,
};
let create = ReplicaPlacement {
collection: collection_id,
handle: handle
.clone()
.unwrap_or_else(|| ReplicaHandle(format!("queue-{row_id}"))),
link_id: Some(link),
object: object.clone(),
level,
meta: meta.clone(),
sort_key: ReplicaSortKey::default(),
flags: flags.clone(),
status: ReplicaStatus::Created,
conflict_revision: None,
conflict_object: None,
base: None,
origin: None,
};
return Ok(Ok(vec![ReplicaWriteOp::UpsertPlacement(create)]));
}
let (seq, removes) = match action {
PimdirAction::SetFlags { seq, .. }
| PimdirAction::Move { seq, .. }
| PimdirAction::Copy { seq, .. }
| PimdirAction::Update { seq, .. } => (*seq, false),
PimdirAction::Remove { seq } => (*seq, true),
PimdirAction::Add { .. } => unreachable!("add staged above"),
PimdirAction::Unknown { .. } => unreachable!("unknown kinds are skipped, never staged"),
};
let item = tx
.query_row(
sql::GET_ITEM,
named_params! { ":collection": collection, ":seq": seq },
read_item_from_row,
)
.optional()?;
let Some(item) = item else {
return if removes {
Ok(Ok(Vec::new()))
} else {
Ok(Err(PimdirRefusal::Park(format!("unknown seq: {seq}"))))
};
};
let handle = tx
.query_row(
sql::HANDLE_FOR_LINK,
named_params! {
":collection": collection,
":link_id": item.link_id.0,
":source": source.0,
},
|r| r.get::<_, String>(0),
)
.optional()?;
let handle = match handle {
Some(handle) => ReplicaHandle(handle),
None if removes => return Ok(Ok(Vec::new())),
None => return Ok(Err(PimdirRefusal::Skip)),
};
let mutation = match action {
PimdirAction::SetFlags { flags, .. } => ReplicaMutation::SetFlags {
handle,
flags: flags.clone(),
},
PimdirAction::Remove { .. } => ReplicaMutation::Remove(handle),
PimdirAction::Move { to, .. } => ReplicaMutation::Move {
handle,
target: to.clone(),
placeholder: ReplicaHandle(format!("queue-{row_id}")),
},
PimdirAction::Copy { to, .. } => ReplicaMutation::Copy {
handle,
target: to.clone(),
placeholder: ReplicaHandle(format!("queue-{row_id}")),
},
PimdirAction::Update { object, meta, .. } => ReplicaMutation::Edit {
handle,
object: ReplicaObject {
hash: object.clone(),
size: 0,
},
body: Vec::new(),
meta: meta.clone(),
sort_key: None,
},
PimdirAction::Add { .. } => unreachable!("add staged above"),
PimdirAction::Unknown { .. } => unreachable!("unknown kinds are skipped, never staged"),
};
let mut mutate = ReplicaMutate::new(collection_id.clone(), mutation);
let _ = mutate.resume(None);
let placements = load_hub_by_link(tx, collection, core::slice::from_ref(&item.link_id.0))?
.project(&collection_id, source);
let loaded = ReplicaLoaded {
placements,
checkpoint: None,
};
match mutate.resume(Some(ReplicaArg::Load(loaded))) {
ReplicaCoroutineState::Yielded(ReplicaYield::WantsWrite(ops)) => {
let ops = ops
.into_iter()
.filter(|op| !matches!(op, ReplicaWriteOp::StoreObject { .. }))
.collect();
Ok(Ok(ops))
}
ReplicaCoroutineState::Complete(Err(err)) => Ok(Err(PimdirRefusal::Park(err.to_string()))),
state => Ok(Err(PimdirRefusal::Park(format!(
"unexpected mutate state: {state:?}"
)))),
}
}
pub struct PimdirProducer {
conn: Connection,
_lock: PimdirLock,
producer: String,
hash: PimdirHashAlgo,
account: Option<String>,
}
impl PimdirProducer {
pub fn open(dir: impl AsRef<Path>, producer: impl Into<String>) -> Result<Self, PimdirError> {
let dir = dir.as_ref();
let flags = OpenFlags::SQLITE_OPEN_READ_WRITE
| OpenFlags::SQLITE_OPEN_URI
| OpenFlags::SQLITE_OPEN_NO_MUTEX;
let conn = Connection::open_with_flags(dir.join("pimdir.db"), flags)?;
conn.execute_batch(
"PRAGMA journal_mode = WAL; PRAGMA foreign_keys = ON; 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 {
conn,
_lock: PimdirLock::stage(dir)?,
producer: producer.into(),
hash,
account: None,
})
}
pub fn hash_algo(&self) -> PimdirHashAlgo {
self.hash
}
pub fn hash(&self, bytes: &[u8]) -> ReplicaHash {
self.hash.hash(bytes)
}
pub fn hasher(&self) -> PimdirHasher {
self.hash.hasher()
}
pub fn for_account(mut self, account: impl Into<String>) -> Self {
self.account = Some(account.into());
self
}
pub fn enqueue(
&mut self,
collection: &str,
action: &PimdirAction,
object_size: Option<u64>,
created_at: &str,
) -> Result<i64, PimdirError> {
let hash = action.object_hash().cloned();
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
tx.execute(
sql::ENSURE_COLLECTION,
named_params! { ":collection": collection, ":account": self.account.as_deref() },
)?;
if let (Some(hash), Some(size)) = (&hash, object_size) {
tx.execute(
sql::STORE_OBJECT,
named_params! { ":hash": hash.0, ":size": size as i64 },
)?;
}
tx.execute(
sql::ENQUEUE_ACTION,
named_params! {
":created_at": created_at,
":producer": self.producer,
":collection": collection,
":action": action.kind(),
":payload": codec::action_to_payload(action),
":object_hash": hash.as_ref().map(|h| h.0.as_str()),
},
)?;
if let Some(hash) = &hash {
tx.execute(
sql::ADJUST_REFCOUNT,
named_params! { ":delta": 1, ":hash": hash.0 },
)?;
}
let id = tx.last_insert_rowid();
tx.commit().map_err(busy_or_sql)?;
Ok(id)
}
pub fn pending_actions(
&self,
collection: &str,
) -> Result<Vec<PimdirPendingAction>, PimdirError> {
load_pending_actions(&self.conn, collection)
}
}
#[derive(Clone, Debug)]
pub struct PimdirBlobs {
root: PathBuf,
hash: PimdirHashAlgo,
}
impl PimdirBlobs {
pub fn open(dir: impl AsRef<Path>, hash: PimdirHashAlgo) -> Self {
Self {
root: dir.as_ref().join("objects"),
hash,
}
}
pub fn hash_algo(&self) -> PimdirHashAlgo {
self.hash
}
pub fn hash(&self, bytes: &[u8]) -> ReplicaHash {
self.hash.hash(bytes)
}
pub fn hasher(&self) -> PimdirHasher {
self.hash.hasher()
}
pub fn path(&self, hash: &ReplicaHash) -> PathBuf {
blob_path(&self.root, &hash.0)
}
pub fn get(&self, hash: &ReplicaHash) -> io::Result<Option<Vec<u8>>> {
match fs::read(blob_path(&self.root, &hash.0)) {
Ok(bytes) => Ok(Some(bytes)),
Err(err) if err.kind() == ErrorKind::NotFound => Ok(None),
Err(err) => Err(err),
}
}
pub fn reader(&self, hash: &ReplicaHash) -> io::Result<Option<fs::File>> {
match fs::File::open(blob_path(&self.root, &hash.0)) {
Ok(file) => Ok(Some(file)),
Err(err) if err.kind() == ErrorKind::NotFound => Ok(None),
Err(err) => Err(err),
}
}
pub fn writer(&self) -> io::Result<PimdirBlobWriter> {
fs::create_dir_all(&self.root)?;
let seq = TMP_SEQ.fetch_add(1, Ordering::Relaxed);
let tmp = self.root.join(format!(".tmp-{}-{seq}", std::process::id()));
let file = fs::File::create(&tmp)?;
Ok(PimdirBlobWriter {
root: self.root.clone(),
tmp,
file: Some(file),
written: 0,
})
}
pub fn files(&self) -> io::Result<Vec<PimdirBlobFile>> {
let mut files = Vec::new();
if self.root.is_dir() {
walk_blobs(&self.root, &mut files)?;
}
Ok(files)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PimdirBlobFile {
pub hash: String,
pub path: PathBuf,
pub size: u64,
}
fn walk_blobs(dir: &Path, files: &mut Vec<PimdirBlobFile>) -> io::Result<()> {
for entry in fs::read_dir(dir)? {
let entry = entry?;
let name = entry.file_name().to_string_lossy().to_string();
if name.starts_with('.') {
continue;
}
let metadata = entry.metadata()?;
if metadata.is_dir() {
walk_blobs(&entry.path(), files)?;
} else if metadata.is_file() {
files.push(PimdirBlobFile {
hash: name,
path: entry.path(),
size: metadata.len(),
});
}
}
Ok(())
}
static TMP_SEQ: AtomicU64 = AtomicU64::new(0);
pub struct PimdirBlobWriter {
root: PathBuf,
tmp: PathBuf,
file: Option<fs::File>,
written: u64,
}
impl PimdirBlobWriter {
pub fn commit(mut self, hash: &ReplicaHash) -> io::Result<u64> {
let file = self.file.take().expect("writer open");
file.sync_all()?;
drop(file);
let path = blob_path(&self.root, &hash.0);
if path.exists() {
let _ = fs::remove_file(&self.tmp);
return Ok(self.written);
}
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
fs::rename(&self.tmp, &path)?;
if let Some(parent) = path.parent() {
sync_dir(parent)?;
}
Ok(self.written)
}
}
impl Write for PimdirBlobWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
let file = self.file.as_mut().expect("writer open");
let n = file.write(buf)?;
self.written += n as u64;
Ok(n)
}
fn flush(&mut self) -> io::Result<()> {
self.file.as_mut().expect("writer open").flush()
}
}
impl Drop for PimdirBlobWriter {
fn drop(&mut self) {
if self.file.is_some() {
let _ = fs::remove_file(&self.tmp);
}
}
}
fn apply_ops(
tx: &Connection,
blobs: &Path,
source: &ReplicaSourceId,
account: Option<&str>,
residual: &mut HashMap<(ReplicaCollectionId, ReplicaHandle), ReplicaPlacement>,
ops: Vec<ReplicaWriteOp>,
) -> Result<(), PimdirError> {
let mut hub_ops: BTreeMap<String, Vec<ReplicaWriteOp>> = BTreeMap::new();
let mut superseded: BTreeMap<String, BTreeSet<ReplicaHandle>> = BTreeMap::new();
for op in ops {
match op {
ReplicaWriteOp::StoreObject { object, body } => {
if let Some(body) = body {
write_blob(blobs, &object.hash.0, &body)?;
}
tx.execute(
sql::STORE_OBJECT,
named_params! { ":hash": object.hash.0, ":size": object.size as i64 },
)?;
}
ReplicaWriteOp::SetCheckpoint {
collection,
checkpoint,
} => {
tx.execute(
sql::ENSURE_COLLECTION,
named_params! { ":collection": collection.0, ":account": account },
)?;
tx.execute(
sql::UPSERT_CHECKPOINT,
named_params! {
":collection": collection.0,
":source": source.0,
":checkpoint": checkpoint.0,
},
)?;
}
ReplicaWriteOp::UpsertPlacement(placement) => {
if placement.link_id.is_some() {
drop_residual(residual, &placement.collection, &placement.handle);
hub_ops
.entry(placement.collection.0.clone())
.or_default()
.push(ReplicaWriteOp::UpsertPlacement(placement));
} else {
let key = (placement.collection.clone(), placement.handle.clone());
residual.insert(key, placement);
}
}
ReplicaWriteOp::DropPlacement {
collection,
handle,
reason,
} => {
drop_residual(residual, &collection, &handle);
if reason == ReplicaDropReason::Superseded {
superseded
.entry(collection.0.clone())
.or_default()
.insert(handle.clone());
}
hub_ops.entry(collection.0.clone()).or_default().push(
ReplicaWriteOp::DropPlacement {
collection,
handle,
reason,
},
);
}
}
}
for (collection, ops) in hub_ops {
refuse_colliding_upserts(&collection, source, &ops)?;
let links = batch_links(tx, &collection, source, &ops)?;
let old_hub = load_hub_by_link(tx, &collection, &links)?;
let mut new_hub = old_hub.clone();
new_hub.absorb(source, &ops);
let superseded = superseded.remove(&collection).unwrap_or_default();
save_hub_diff(
tx,
&collection,
source,
account,
&old_hub,
&new_hub,
&superseded,
)?;
adjust_refcounts(tx, &object_refs(&old_hub), &object_refs(&new_hub))?;
}
Ok(())
}
fn refuse_colliding_upserts(
collection: &str,
source: &ReplicaSourceId,
ops: &[ReplicaWriteOp],
) -> Result<(), PimdirError> {
let mut claimed: BTreeMap<&ReplicaLinkId, &ReplicaHandle> = BTreeMap::new();
for op in ops {
let ReplicaWriteOp::UpsertPlacement(placement) = op else {
continue;
};
let Some(link) = placement.link_id.as_ref() else {
continue;
};
match claimed.insert(link, &placement.handle) {
Some(bound) if *bound != placement.handle => {
return Err(PimdirError::Rebind {
collection: collection.into(),
link_id: link.0.clone(),
source: source.0.clone(),
bound: bound.0.clone(),
incoming: placement.handle.0.clone(),
});
}
_ => {}
}
}
Ok(())
}
fn init_schema(conn: &mut Connection, hash: PimdirHashAlgo) -> Result<(), PimdirError> {
let version: i64 = conn.pragma_query_value(None, "user_version", |r| r.get(0))?;
if version > sql::VERSION {
return Err(PimdirError::Version { found: version });
}
if version == sql::VERSION {
check_version_agreement(conn, version)?;
check_rename_cascades(conn)?;
return reconcile_draft_shape(conn);
}
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
tx.execute_batch(sql::MIGRATION_0001)?;
tx.execute(
"INSERT OR IGNORE INTO store_meta(id, version, hash_algo, created_at) \
VALUES(1, ?1, ?2, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
params![sql::VERSION, hash.as_str()],
)?;
tx.pragma_update(None, "user_version", sql::VERSION)?;
tx.commit().map_err(busy_or_sql)?;
Ok(())
}
fn check_rename_cascades(conn: &Connection) -> Result<(), PimdirError> {
const CASCADING: [(&str, &str); 5] = [
("collections", "collections"),
("sources", "collections"),
("items", "collections"),
("bindings", "items"),
("queue", "collections"),
];
for (table, parent) in CASCADING {
let mut stmt = conn.prepare(&format!(
"SELECT on_update FROM pragma_foreign_key_list('{table}') WHERE \"table\" = '{parent}'"
))?;
let mut rows = stmt.query([])?;
while let Some(row) = rows.next()? {
let on_update: String = row.get(0)?;
if on_update != "CASCADE" {
return Err(PimdirError::Unreconcilable { table });
}
}
}
Ok(())
}
fn reconcile_draft_shape(conn: &mut Connection) -> Result<(), PimdirError> {
const FOLDED_IN: [(&str, &str, &str); 9] = [
("bindings", "conflicted", "INTEGER NOT NULL DEFAULT 0"),
("bindings", "conflict_revision", "TEXT"),
(
"bindings",
"conflict_object",
"TEXT REFERENCES objects(hash)",
),
("bindings", "shared_object", "TEXT"),
("items", "retained_at", "TEXT"),
("items", "retained_by", "TEXT"),
("collections", "account", "TEXT"),
("items", "sort_key", "TEXT NOT NULL DEFAULT ''"),
("bindings", "base_present", "INTEGER NOT NULL DEFAULT 0"),
];
const FOLDED_OUT: [(&str, &str); 1] = [("bindings", "ambiguous_handles")];
let mut missing = Vec::new();
for (table, column, decl) in FOLDED_IN {
if !has_column(conn, table, column)? {
missing.push((table, column, decl));
}
}
let mut stale = Vec::new();
for (table, column) in FOLDED_OUT {
if has_column(conn, table, column)? {
stale.push((table, column));
}
}
let mut reshaped = Vec::new();
for (index, columns) in sql::RESHAPED_INDEXES {
if index_columns(conn, index)?.is_some_and(|held| held != *columns) {
reshaped.push(*index);
}
}
let backfill_shared = missing
.iter()
.any(|(table, column, _)| (*table, *column) == ("bindings", "shared_object"));
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(busy_or_sql)?;
for (table, column, decl) in missing {
tx.execute_batch(&format!("ALTER TABLE {table} ADD COLUMN {column} {decl}"))?;
}
if backfill_shared {
tx.execute_batch(sql::BACKFILL_SHARED_OBJECT)?;
}
for (table, column) in stale {
tx.execute_batch(&format!("ALTER TABLE {table} DROP COLUMN {column}"))?;
}
for index in reshaped {
tx.execute_batch(&format!("DROP INDEX IF EXISTS {index}"))?;
}
tx.execute_batch(sql::ENSURE_INDEXES)?;
tx.commit().map_err(busy_or_sql)?;
Ok(())
}
fn read_hash_algo(
conn: &Connection,
declared: Option<PimdirHashAlgo>,
) -> Result<PimdirHashAlgo, PimdirError> {
let stored: Option<String> = conn
.query_row("SELECT hash_algo FROM store_meta WHERE id = 1", [], |row| {
row.get(0)
})
.optional()?;
let Some(stored) = stored else {
return Ok(declared.unwrap_or_default());
};
let Some(algo) = PimdirHashAlgo::parse(&stored) else {
return Err(PimdirError::HashAlgo {
found: stored,
declared: declared.map(|a| a.as_str()),
});
};
match declared {
Some(declared) if declared != algo => Err(PimdirError::HashAlgo {
found: stored,
declared: Some(declared.as_str()),
}),
_ => Ok(algo),
}
}
fn check_version_agreement(conn: &Connection, user_version: i64) -> Result<(), PimdirError> {
let stamped: Option<i64> = conn
.query_row("SELECT version FROM store_meta WHERE id = 1", [], |row| {
row.get(0)
})
.optional()?;
match stamped {
Some(store_meta) if store_meta != user_version => Err(PimdirError::VersionMismatch {
user_version,
store_meta,
}),
_ => Ok(()),
}
}
fn has_column(conn: &Connection, table: &str, column: &str) -> rusqlite::Result<bool> {
let mut stmt = conn.prepare(&format!("PRAGMA table_info({table})"))?;
let mut rows = stmt.query([])?;
while let Some(row) = rows.next()? {
let name: String = row.get(1)?;
if name == column {
return Ok(true);
}
}
Ok(false)
}
fn index_columns(conn: &Connection, index: &str) -> rusqlite::Result<Option<Vec<String>>> {
let columns = rows(conn, &format!("PRAGMA index_info({index})"), [], |row| {
row.get::<_, String>(2)
})?;
Ok((!columns.is_empty()).then_some(columns))
}
fn drop_residual(
residual: &mut HashMap<(ReplicaCollectionId, ReplicaHandle), ReplicaPlacement>,
collection: &ReplicaCollectionId,
handle: &ReplicaHandle,
) {
residual.remove(&(collection.clone(), handle.clone()));
}
fn load_hub(conn: &Connection, collection: &str) -> rusqlite::Result<ReplicaHub> {
read_hub(conn, collection, None)
}
fn batch_links(
conn: &Connection,
collection: &str,
source: &ReplicaSourceId,
ops: &[ReplicaWriteOp],
) -> rusqlite::Result<Vec<String>> {
let mut links: BTreeSet<String> = BTreeSet::new();
for op in ops {
match op {
ReplicaWriteOp::UpsertPlacement(placement) => {
if let Some(link) = &placement.link_id {
links.insert(link.0.clone());
}
}
ReplicaWriteOp::DropPlacement { handle, .. } => {
let link = conn
.query_row(
sql::LINK_FOR_HANDLE,
named_params! {
":collection": collection,
":source": source.0,
":handle": handle.0,
},
|r| r.get::<_, String>(0),
)
.optional()?;
links.extend(link);
}
_ => {}
}
}
Ok(links.into_iter().collect())
}
fn load_hub_by_link(
conn: &Connection,
collection: &str,
links: &[String],
) -> rusqlite::Result<ReplicaHub> {
read_hub(conn, collection, Some(links))
}
fn read_hub(
conn: &Connection,
collection: &str,
links: Option<&[String]>,
) -> rusqlite::Result<ReplicaHub> {
let mut hub = ReplicaHub::default();
if let Some(policy) = conn
.query_row(
sql::LOAD_CONFLICT,
named_params! { ":collection": collection },
|r| r.get::<_, String>(0),
)
.optional()?
{
hub.conflict = conflict_from_str(&policy);
}
let scope = links.map(|links| serde_json::to_string(links).unwrap_or_else(|_| "[]".into()));
let (items_sql, bindings_sql) = match scope {
Some(_) => (sql::LOAD_ITEMS_BY_LINK, sql::LOAD_BINDINGS_BY_LINK),
None => (sql::LOAD_ITEMS, sql::LOAD_BINDINGS),
};
let mut params: Vec<(&str, &dyn ToSql)> = vec![(":collection", &collection)];
if let Some(scope) = &scope {
params.push((":links", scope));
}
for (link, item) in rows(conn, items_sql, params.as_slice(), item_from_row)? {
hub.items.insert(link, item);
}
for (link, source, binding) in rows(conn, bindings_sql, params.as_slice(), binding_from_row)? {
if let Some(item) = hub.items.get_mut(&link) {
item.sources.insert(source, binding);
}
}
Ok(hub)
}
fn save_hub_diff(
conn: &Connection,
collection: &str,
source: &ReplicaSourceId,
account: Option<&str>,
old: &ReplicaHub,
new: &ReplicaHub,
superseded: &BTreeSet<ReplicaHandle>,
) -> Result<(), PimdirError> {
conn.execute(
sql::ENSURE_COLLECTION,
named_params! { ":collection": collection, ":account": account },
)?;
if old.conflict != new.conflict {
conn.execute(
sql::SET_CONFLICT,
named_params! { ":collection": collection, ":conflict": conflict_to_str(new.conflict) },
)?;
}
for (link, item) in &old.items {
if new.items.contains_key(link) {
continue;
}
conn.execute(
sql::RETAIN_ITEM,
named_params! { ":collection": collection, ":link_id": link.0, ":source": source.0 },
)?;
conn.execute(
sql::DELETE_ITEM_BINDINGS,
named_params! { ":collection": collection, ":link_id": link.0 },
)?;
for hash in [item.object.as_ref(), item.conflict_object.as_ref()]
.into_iter()
.flatten()
{
conn.execute(
sql::ADJUST_REFCOUNT,
named_params! { ":delta": 1, ":hash": hash.0 },
)?;
}
}
for (link, item) in &new.items {
match old.items.get(link) {
None => insert_item(conn, collection, link, item)?,
Some(prev) => {
if !item_columns_eq(prev, item) {
update_item(conn, collection, link, item)?;
}
save_bindings_diff(conn, collection, link, prev, item, superseded)?;
}
}
}
Ok(())
}
fn item_columns_eq(a: &ReplicaHubItem, b: &ReplicaHubItem) -> bool {
a.flags == b.flags
&& a.object == b.object
&& a.meta == b.meta
&& a.sort_key == b.sort_key
&& a.level == b.level
&& a.deleted == b.deleted
&& a.conflicted == b.conflicted
&& a.conflict_object == b.conflict_object
}
fn insert_item(
conn: &Connection,
collection: &str,
link: &ReplicaLinkId,
item: &ReplicaHubItem,
) -> rusqlite::Result<()> {
if revive_item(conn, collection, link, item)? {
return Ok(());
}
let seq: i64 = match conn
.query_row(
sql::SEQ_FOR_LINK_ANY,
named_params! { ":link_id": link.0 },
|row| row.get(0),
)
.optional()?
{
Some(existing) => existing,
None => conn.query_row(sql::BUMP_NEXT_SEQ, [], |row| row.get(0))?,
};
conn.execute(
sql::INSERT_ITEM,
named_params! {
":collection": collection,
":link_id": link.0,
":seq": seq,
":flags": codec::flags_to_json(&item.flags),
":object_hash": item.object.as_ref().map(|o| o.0.as_str()),
":meta": item.meta.as_ref().map(|m| m.0.as_str()),
":sort_key": item.sort_key.0.as_str(),
":level": codec::level_to_int(item.level),
":deleted": item.deleted as i64,
":conflicted": item.conflicted as i64,
":conflict_object": item.conflict_object.as_ref().map(|o| o.0.as_str()),
},
)?;
for (source, binding) in &item.sources {
insert_binding(conn, collection, link, source, binding)?;
}
Ok(())
}
fn revive_item(
conn: &Connection,
collection: &str,
link: &ReplicaLinkId,
item: &ReplicaHubItem,
) -> rusqlite::Result<bool> {
let pinned: Option<(Option<String>, Option<String>)> = conn
.query_row(
sql::RETAINED_ITEM,
named_params! { ":collection": collection, ":link_id": link.0 },
|row| Ok((row.get(1)?, row.get(2)?)),
)
.optional()?;
let Some((object, conflict_object)) = pinned else {
return Ok(false);
};
conn.execute(
sql::REVIVE_ITEM,
named_params! { ":collection": collection, ":link_id": link.0 },
)?;
update_item(conn, collection, link, item)?;
for hash in [object, conflict_object].into_iter().flatten() {
conn.execute(
sql::ADJUST_REFCOUNT,
named_params! { ":delta": -1, ":hash": hash },
)?;
}
for (source, binding) in &item.sources {
insert_binding(conn, collection, link, source, binding)?;
}
Ok(true)
}
fn update_item(
conn: &Connection,
collection: &str,
link: &ReplicaLinkId,
item: &ReplicaHubItem,
) -> rusqlite::Result<()> {
conn.execute(
sql::UPDATE_ITEM,
named_params! {
":collection": collection,
":link_id": link.0,
":flags": codec::flags_to_json(&item.flags),
":object_hash": item.object.as_ref().map(|o| o.0.as_str()),
":meta": item.meta.as_ref().map(|m| m.0.as_str()),
":sort_key": item.sort_key.0.as_str(),
":level": codec::level_to_int(item.level),
":deleted": item.deleted as i64,
":conflicted": item.conflicted as i64,
":conflict_object": item.conflict_object.as_ref().map(|o| o.0.as_str()),
},
)?;
Ok(())
}
fn save_bindings_diff(
conn: &Connection,
collection: &str,
link: &ReplicaLinkId,
old: &ReplicaHubItem,
new: &ReplicaHubItem,
superseded: &BTreeSet<ReplicaHandle>,
) -> Result<(), PimdirError> {
for source in old.sources.keys() {
if !new.sources.contains_key(source) {
conn.execute(
sql::DELETE_BINDING,
named_params! { ":collection": collection, ":link_id": link.0, ":source": source.0 },
)?;
}
}
for (source, binding) in &new.sources {
match old.sources.get(source) {
None => insert_binding(conn, collection, link, source, binding)?,
Some(prev) if binding.handle != prev.handle && superseded.contains(&prev.handle) => {
conn.execute(
sql::DELETE_BINDING,
named_params! { ":collection": collection, ":link_id": link.0, ":source": source.0 },
)?;
insert_binding(conn, collection, link, source, binding)?
}
Some(prev) if binding.handle != prev.handle => {
return Err(PimdirError::Rebind {
collection: collection.into(),
link_id: link.0.clone(),
source: source.0.clone(),
bound: prev.handle.0.clone(),
incoming: binding.handle.0.clone(),
});
}
Some(prev) if prev != binding => {
update_binding(conn, collection, link, source, binding)?
}
Some(_) => {}
}
}
Ok(())
}
fn insert_binding(
conn: &Connection,
collection: &str,
link: &ReplicaLinkId,
source: &ReplicaSourceId,
binding: &ReplicaSourceBinding,
) -> rusqlite::Result<()> {
conn.execute(
sql::INSERT_BINDING,
named_params! {
":collection": collection,
":link_id": link.0,
":source": source.0,
":handle": binding.handle.0,
":base_flags": binding.base.as_ref().map(|b| codec::flags_to_json(&b.flags)),
":base_object": binding.base.as_ref().and_then(|b| b.object.as_ref()).map(|o| o.0.as_str()),
":base_revision": binding.base.as_ref().and_then(|b| b.revision.as_deref()),
":base_present": binding.base.is_some() as i64,
":conflicted": binding.conflicted as i64,
":conflict_revision": binding.conflicted.then_some(binding.conflict_revision.as_deref()).flatten(),
":conflict_object": conflict_object(binding).map(|hash| hash.0.as_str()),
":shared_object": binding.shared_object.as_ref().map(|hash| hash.0.as_str()),
},
)?;
Ok(())
}
fn update_binding(
conn: &Connection,
collection: &str,
link: &ReplicaLinkId,
source: &ReplicaSourceId,
binding: &ReplicaSourceBinding,
) -> rusqlite::Result<()> {
conn.execute(
sql::UPDATE_BINDING,
named_params! {
":collection": collection,
":link_id": link.0,
":source": source.0,
":base_flags": binding.base.as_ref().map(|b| codec::flags_to_json(&b.flags)),
":base_object": binding.base.as_ref().and_then(|b| b.object.as_ref()).map(|o| o.0.as_str()),
":base_revision": binding.base.as_ref().and_then(|b| b.revision.as_deref()),
":base_present": binding.base.is_some() as i64,
":conflicted": binding.conflicted as i64,
":conflict_revision": binding.conflicted.then_some(binding.conflict_revision.as_deref()).flatten(),
":conflict_object": conflict_object(binding).map(|hash| hash.0.as_str()),
":shared_object": binding.shared_object.as_ref().map(|hash| hash.0.as_str()),
},
)?;
Ok(())
}
fn conflict_object(binding: &ReplicaSourceBinding) -> Option<&ReplicaHash> {
binding
.conflicted
.then_some(binding.conflict_object.as_ref())
.flatten()
}
fn object_refs(hub: &ReplicaHub) -> HashMap<String, i64> {
let mut refs: HashMap<String, i64> = HashMap::new();
let mut bump = |hash: &ReplicaHash| *refs.entry(hash.0.clone()).or_insert(0) += 1;
for item in hub.items.values() {
if let Some(object) = &item.object {
bump(object);
}
if let Some(conflict) = &item.conflict_object {
bump(conflict);
}
for binding in item.sources.values() {
if let Some(object) = binding.base.as_ref().and_then(|b| b.object.as_ref()) {
bump(object);
}
if let Some(conflict) = conflict_object(binding) {
bump(conflict);
}
}
}
refs
}
fn adjust_refcounts(
conn: &Connection,
old: &HashMap<String, i64>,
new: &HashMap<String, i64>,
) -> rusqlite::Result<()> {
for (hash, new_count) in new {
let delta = new_count - old.get(hash).copied().unwrap_or(0);
if delta != 0 {
conn.execute(
sql::ADJUST_REFCOUNT,
named_params! { ":delta": delta, ":hash": hash },
)?;
}
}
for (hash, old_count) in old {
if !new.contains_key(hash) {
conn.execute(
sql::ADJUST_REFCOUNT,
named_params! { ":delta": -old_count, ":hash": hash },
)?;
}
}
Ok(())
}
fn read_item_from_row(row: &Row) -> rusqlite::Result<PimdirItem> {
let seq: i64 = row.get(0)?;
let link: String = row.get(1)?;
let flags: Option<String> = row.get(2)?;
let object: Option<String> = row.get(3)?;
let meta: Option<String> = row.get(4)?;
let sort_key: String = row.get(5)?;
let level: i64 = row.get(6)?;
let retention = match row.as_ref().column_count() > 7 {
true => Some(PimdirRetention {
at: row.get(7)?,
by: row.get(8)?,
size: row.get::<_, Option<i64>>(9)?.map(|size| size.max(0) as u64),
}),
false => None,
};
Ok(PimdirItem {
seq,
link_id: ReplicaLinkId(link),
flags: codec::flags_from_json(flags.as_deref()),
meta: meta.map(ReplicaMeta),
sort_key,
object: object.map(ReplicaHash),
level: codec::level_from_int(level),
retention,
})
}
fn item_from_row(row: &Row) -> rusqlite::Result<(ReplicaLinkId, ReplicaHubItem)> {
let link: String = row.get(0)?;
let flags: Option<String> = row.get(1)?;
let object: Option<String> = row.get(2)?;
let meta: Option<String> = row.get(3)?;
let sort_key: String = row.get(4)?;
let level: i64 = row.get(5)?;
let deleted: i64 = row.get(6)?;
let conflicted: i64 = row.get(7)?;
let conflict_object: Option<String> = row.get(8)?;
Ok((
ReplicaLinkId(link),
ReplicaHubItem {
flags: codec::flags_from_json(flags.as_deref()),
object: object.map(ReplicaHash),
meta: meta.map(ReplicaMeta),
sort_key: ReplicaSortKey(sort_key),
level: codec::level_from_int(level),
deleted: deleted != 0,
conflicted: conflicted != 0,
conflict_object: conflict_object.map(ReplicaHash),
sources: BTreeMap::new(),
},
))
}
fn binding_from_row(
row: &Row,
) -> rusqlite::Result<(ReplicaLinkId, ReplicaSourceId, ReplicaSourceBinding)> {
let link: String = row.get(0)?;
let source: String = row.get(1)?;
let handle: String = row.get(2)?;
let base_flags: Option<String> = row.get(3)?;
let base_object: Option<String> = row.get(4)?;
let base_revision: Option<String> = row.get(5)?;
let base_present: i64 = row.get(6)?;
let conflicted: i64 = row.get(7)?;
let conflict_revision: Option<String> = row.get(8)?;
let conflict_object: Option<String> = row.get(9)?;
let shared_object: Option<String> = row.get(10)?;
let base = if base_present != 0
|| base_flags.is_some()
|| base_object.is_some()
|| base_revision.is_some()
{
Some(ReplicaBase {
flags: codec::flags_from_json(base_flags.as_deref()),
revision: base_revision,
object: base_object.map(ReplicaHash),
})
} else {
None
};
let conflicted = conflicted != 0;
Ok((
ReplicaLinkId(link),
ReplicaSourceId(source),
ReplicaSourceBinding {
handle: ReplicaHandle(handle),
base,
conflicted,
conflict_revision: conflicted.then_some(conflict_revision).flatten(),
conflict_object: conflicted
.then_some(conflict_object)
.flatten()
.map(ReplicaHash),
shared_object: shared_object.map(ReplicaHash),
},
))
}
fn conflict_from_str(value: &str) -> ReplicaHubConflict {
match value {
"prefer-incoming" => ReplicaHubConflict::PreferIncoming,
"prefer-existing" => ReplicaHubConflict::PreferExisting,
_ => ReplicaHubConflict::Manual,
}
}
fn conflict_to_str(policy: ReplicaHubConflict) -> &'static str {
match policy {
ReplicaHubConflict::Manual => "manual",
ReplicaHubConflict::PreferIncoming => "prefer-incoming",
ReplicaHubConflict::PreferExisting => "prefer-existing",
}
}
fn blob_path(blobs: &Path, hash: &str) -> PathBuf {
if hash.len() >= 4 {
blobs.join(&hash[0..2]).join(&hash[2..4]).join(hash)
} else {
blobs.join(hash)
}
}
fn stage_blobs(blobs: &Path, ops: &[ReplicaWriteOp]) -> io::Result<()> {
for op in ops {
if let ReplicaWriteOp::StoreObject {
object,
body: Some(body),
} = op
{
write_blob(blobs, &object.hash.0, body)?;
}
}
Ok(())
}
fn write_blob(blobs: &Path, hash: &str, body: &[u8]) -> io::Result<()> {
let path = blob_path(blobs, hash);
if path.exists() {
return Ok(());
}
let parent = path.parent().unwrap_or(blobs);
fs::create_dir_all(parent)?;
let tmp = parent.join(format!(".{hash}.tmp"));
{
let mut file = fs::File::create(&tmp)?;
file.write_all(body)?;
file.sync_all()?;
}
fs::rename(&tmp, &path)?;
sync_dir(parent)
}
fn sync_dir(dir: &Path) -> io::Result<()> {
fs::File::open(dir)?.sync_all()
}
#[derive(Debug)]
pub enum PimdirError {
Sql(rusqlite::Error),
Io(io::Error),
Json(serde_json::Error),
Action(PimdirActionError),
Rebind {
collection: String,
link_id: String,
source: String,
bound: String,
incoming: String,
},
Version {
found: i64,
},
Uncreated,
VersionMismatch {
user_version: i64,
store_meta: i64,
},
Unreconcilable {
table: &'static str,
},
HashAlgo {
found: String,
declared: Option<&'static str>,
},
Busy,
Owned(PathBuf),
Staging(PathBuf),
}
impl fmt::Display for PimdirError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
PimdirError::Sql(err) => write!(f, "pimdir SQL error: {err}"),
PimdirError::Io(err) => write!(f, "pimdir I/O error: {err}"),
PimdirError::Json(err) => write!(f, "pimdir JSON error: {err}"),
PimdirError::Action(err) => write!(f, "pimdir action error: {err}"),
PimdirError::Rebind {
collection,
link_id,
source,
bound,
incoming,
} => write!(
f,
"pimdir binding {collection}/{link_id} on source {source} holds handle {bound}, and this write carries {incoming}: a binding pins one handle, and a second copy of an identity is stored under a key of its own"
),
PimdirError::Version { found } => write!(
f,
"pimdir store schema version {found} is unsupported (this crate services version {})",
sql::VERSION
),
PimdirError::Uncreated => write!(
f,
"pimdir store has no schema yet: its owner has to create it first"
),
PimdirError::VersionMismatch {
user_version,
store_meta,
} => write!(
f,
"pimdir store is corrupt: PRAGMA user_version is {user_version} but store_meta records {store_meta}"
),
PimdirError::Unreconcilable { table } => write!(
f,
"pimdir store predates the ON UPDATE CASCADE on `{table}`, which cannot be added in place: delete the store and let it resync"
),
PimdirError::HashAlgo {
found,
declared: Some(declared),
} => write!(
f,
"pimdir store names its objects with `{found}`, not the `{declared}` this handle declared"
),
PimdirError::HashAlgo {
found,
declared: None,
} => write!(
f,
"pimdir store names its objects with `{found}`, which this crate does not compute"
),
PimdirError::Owned(store) => write!(
f,
"pimdir store at {} is owned by another process",
store.display()
),
PimdirError::Staging(store) => write!(
f,
"pimdir store at {} has a producer staging a body",
store.display()
),
PimdirError::Busy => write!(
f,
"pimdir store is busy: another writer holds the write lock; retry once it releases"
),
}
}
}
fn busy_or_sql(err: rusqlite::Error) -> PimdirError {
match &err {
rusqlite::Error::SqliteFailure(e, _)
if matches!(e.code, ErrorCode::DatabaseBusy | ErrorCode::DatabaseLocked) =>
{
PimdirError::Busy
}
_ => PimdirError::Sql(err),
}
}
impl std::error::Error for PimdirError {}
impl From<rusqlite::Error> for PimdirError {
fn from(err: rusqlite::Error) -> Self {
PimdirError::Sql(err)
}
}
impl From<io::Error> for PimdirError {
fn from(err: io::Error) -> Self {
PimdirError::Io(err)
}
}
impl From<serde_json::Error> for PimdirError {
fn from(err: serde_json::Error) -> Self {
PimdirError::Json(err)
}
}
impl From<PimdirActionError> for PimdirError {
fn from(err: PimdirActionError) -> Self {
PimdirError::Action(err)
}
}