pub mod backend;
pub mod backup;
pub mod lease;
pub mod metrics;
pub mod node;
pub mod server;
pub use backend::{LogBackend, NodeStore, StoreSpec};
use corium_core::{EntityId, IndexOrder, KeywordInterner, Partition, Schema};
use corium_db::{Db, FIRST_USER_ID, Idents};
use corium_index::Segment;
use corium_log::{LogError, TransactionLog, TxRecord};
use corium_store::{BlobStore, RootStore, StoreError};
use corium_tx::{PreparedTx, TxError, TxItem, prepare};
use std::{
sync::{Arc, Mutex, mpsc},
time::{SystemTime, UNIX_EPOCH},
};
use thiserror::Error;
#[derive(Clone, Debug)]
pub struct TxReport {
pub db_before: Db,
pub db_after: Db,
pub tx: PreparedTx,
pub tx_instant: i64,
}
#[derive(Debug, Error)]
pub enum TransactError {
#[error(transparent)]
Tx(#[from] TxError),
#[error(transparent)]
Log(#[from] LogError),
#[error(transparent)]
Store(#[from] StoreError),
#[error("system clock is before Unix epoch")]
Clock,
#[error("index task failed: {0}")]
IndexTask(String),
#[error("deposed: database root is owned by lease version {published}")]
Deposed {
published: u64,
},
}
struct State {
db: Db,
next_user: u64,
last_instant: i64,
subscribers: Vec<mpsc::Sender<TxReport>>,
}
fn next_user_id<'a>(datoms: impl Iterator<Item = &'a corium_core::Datom>, floor: u64) -> u64 {
datoms
.filter(|d| d.e.partition() == Partition::User as u32)
.map(|d| d.e.sequence() + 1)
.fold(floor, u64::max)
}
pub struct EmbeddedTransactor {
log: Arc<dyn TransactionLog>,
state: Mutex<State>,
}
impl EmbeddedTransactor {
pub fn recover(schema: Schema, log: Arc<dyn TransactionLog>) -> Result<Self, TransactError> {
Self::recover_from(Db::new(schema), log)
}
pub fn recover_from(base: Db, log: Arc<dyn TransactionLog>) -> Result<Self, TransactError> {
let mut db = base;
let mut last_instant = i64::MIN;
for record in log.replay()? {
db = db.with_transaction(record.t, &record.datoms);
last_instant = last_instant.max(record.tx_instant);
}
let next_user = next_user_id(db.recorded_datoms().iter(), FIRST_USER_ID);
Ok(Self {
log,
state: Mutex::new(State {
db,
next_user,
last_instant,
subscribers: Vec::new(),
}),
})
}
pub fn recover_from_snapshot(
snapshot: Db,
next_entity_id: u64,
last_tx_instant: i64,
log: Arc<dyn TransactionLog>,
) -> Result<Self, TransactError> {
let mut db = snapshot;
let index_basis = db.basis_t();
let mut last_instant = last_tx_instant;
let mut next_user = next_entity_id.max(FIRST_USER_ID);
for record in log.tx_range(index_basis + 1, None)? {
db = db.with_transaction(record.t, &record.datoms);
last_instant = last_instant.max(record.tx_instant);
next_user = next_user.max(next_user_id(record.datoms.iter(), next_user));
}
Ok(Self {
log,
state: Mutex::new(State {
db,
next_user,
last_instant,
subscribers: Vec::new(),
}),
})
}
fn recovery_snapshot(&self) -> (Db, u64, i64) {
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
(state.db.clone(), state.next_user, state.last_instant)
}
#[must_use]
pub fn db(&self) -> Db {
self.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.db
.clone()
}
pub fn subscribe(&self) -> mpsc::Receiver<TxReport> {
let (tx, rx) = mpsc::channel();
self.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.subscribers
.push(tx);
rx
}
pub fn transact(
&self,
items: impl IntoIterator<Item = TxItem>,
) -> Result<TxReport, TransactError> {
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let before = state.db.clone();
let t = before.basis_t() + 1;
let tx_id = EntityId::new(Partition::Tx as u32, t);
let prepared = prepare(&before, items, tx_id, state.next_user)?;
let millis = i64::try_from(
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_err(|_| TransactError::Clock)?
.as_millis(),
)
.unwrap_or(i64::MAX);
let tx_instant = millis.max(state.last_instant.saturating_add(1));
self.log.append(&TxRecord {
t,
tx_instant,
datoms: prepared.datoms.clone(),
})?;
state.db = before.with_transaction(t, &prepared.datoms);
state.last_instant = tx_instant;
state.next_user = prepared
.tempids
.values()
.filter(|e| e.partition() == Partition::User as u32)
.map(|e| e.sequence() + 1)
.max()
.unwrap_or(state.next_user)
.max(state.next_user);
let report = TxReport {
db_before: before,
db_after: state.db.clone(),
tx: prepared,
tx_instant,
};
state
.subscribers
.retain(|subscriber| subscriber.send(report.clone()).is_ok());
Ok(report)
}
pub fn update_naming(&self, idents: Idents, interner: KeywordInterner) {
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.db = state.db.clone().with_naming(idents, interner);
}
pub async fn publish_indexes(
&self,
store: &(impl BlobStore + RootStore),
root_name: &str,
lease_version: u64,
) -> Result<DbRoot, TransactError> {
let (snapshot, next_entity_id, last_tx_instant) = self.recovery_snapshot();
let datoms = snapshot.datoms();
let chunked = tokio::task::spawn_blocking(move || {
[
IndexOrder::Eavt,
IndexOrder::Aevt,
IndexOrder::Avet,
IndexOrder::Vaet,
]
.into_iter()
.map(|order| {
let segment = Segment::build(order, datoms.clone());
corium_store::chunk_segment_keys(segment.entries().map(|(key, _)| key.as_slice()))
})
.collect::<Vec<_>>()
})
.await
.map_err(|error| TransactError::IndexTask(error.to_string()))?;
let mut ids = Vec::new();
for chunks in chunked {
let mut children = Vec::new();
for chunk in chunks {
children.push(store.put_if_absent(&chunk).await?);
}
let manifest = corium_store::encode_index_manifest(&children);
ids.push(store.put_if_absent(&manifest).await?);
}
let root = DbRoot {
format_version: corium_store::FORMAT_VERSION,
lease_version,
owner: String::new(),
lease_expires_unix_ms: 0,
owner_endpoint: String::new(),
index_basis_t: snapshot.basis_t(),
roots: Some([
ids[0].clone(),
ids[1].clone(),
ids[2].clone(),
ids[3].clone(),
]),
next_entity_id,
last_tx_instant,
};
publish_root(store, root_name, &root).await?;
Ok(root)
}
}
pub async fn publish_root(
store: &dyn RootStore,
root_name: &str,
root: &DbRoot,
) -> Result<(), TransactError> {
loop {
let previous = store.get_root(root_name).await?;
let stored = previous.as_deref().and_then(DbRoot::decode);
let mut next = root.clone();
if let Some(stored) = stored {
if stored.lease_version > root.lease_version {
return Err(TransactError::Deposed {
published: stored.lease_version,
});
}
if stored.lease_version == root.lease_version
&& stored.index_basis_t >= root.index_basis_t
{
return Ok(());
}
if stored.lease_version == root.lease_version {
next.owner = stored.owner;
next.lease_expires_unix_ms = stored.lease_expires_unix_ms;
next.owner_endpoint = stored.owner_endpoint;
}
}
match store
.cas_root(root_name, previous.as_deref(), &next.encode())
.await
{
Ok(()) => return Ok(()),
Err(StoreError::CasFailed { .. }) => {}
Err(error) => return Err(error.into()),
}
}
}
pub use corium_store::{DbRoot, db_root_name};