use crate::{
Config, Error, MergePair, Node, NodeData, NodeHistory, NodeId, ObjectId, ObjectPayload, Owner,
Provenance, Result, TransactionId, TransactionPackage, TransactionSource, WriterId,
store::{self, Candidate, HistoryIndex, NodeFile, TxMeta},
wire::{self, NodeOperation, ObjectDeclaration, ParsedTransaction, UnsignedTransaction},
};
use chrono::Utc;
use ed25519_dalek::SigningKey;
use fs2::FileExt;
use std::{
cmp::Ordering,
collections::{BTreeMap, BTreeSet, VecDeque},
fmt,
fs::{self, File, OpenOptions},
path::{Path, PathBuf},
sync::{
Arc, Mutex, MutexGuard, RwLock,
atomic::{AtomicBool, Ordering as AtomicOrdering},
},
};
#[derive(Clone)]
pub struct KwebDb {
inner: Arc<Inner>,
}
struct Inner {
root: PathBuf,
_lock_file: File,
mutation: Mutex<()>,
visibility: RwLock<()>,
pending: Mutex<PendingPool>,
outbox_draining: AtomicBool,
signing_key: SigningKey,
local_writer: WriterId,
writers: Vec<WriterId>,
gossip: Arc<dyn crate::Gossip>,
poisoned: AtomicBool,
}
impl Drop for Inner {
fn drop(&mut self) {
let _ = store::clear_incoming(&self.root);
}
}
#[derive(Default)]
struct PendingPool {
transactions: BTreeMap<TransactionId, PendingTransaction>,
waiting_on: BTreeMap<TransactionId, BTreeSet<TransactionId>>,
ready: VecDeque<TransactionId>,
object_reservations: BTreeMap<ObjectId, TransactionId>,
}
struct PendingTransaction {
id: TransactionId,
transaction: Vec<u8>,
objects: Vec<store::SpooledObject>,
missing: BTreeSet<TransactionId>,
}
enum CommitObjects<'a> {
Memory(&'a [ObjectPayload]),
Spooled(&'a [store::SpooledObject]),
}
impl fmt::Debug for KwebDb {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("KwebDb")
.field("root", &self.inner.root)
.finish_non_exhaustive()
}
}
impl KwebDb {
pub fn open(path: impl AsRef<Path>, config: Config) -> Result<Self> {
validate_config(&config)?;
let root = path.as_ref().to_path_buf();
match fs::symlink_metadata(&root) {
Ok(metadata) => {
if !metadata.file_type().is_dir() || metadata.file_type().is_symlink() {
return Err(Error::invalid_config(
"database root must be a real directory",
));
}
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
fs::create_dir_all(&root)?;
let metadata = fs::symlink_metadata(&root)?;
if !metadata.file_type().is_dir() || metadata.file_type().is_symlink() {
return Err(Error::invalid_config(
"database root must be a real directory",
));
}
}
Err(error) => return Err(error.into()),
}
let lock_file = open_lock_file(&root.join("LOCK"))?;
FileExt::try_lock_exclusive(&lock_file).map_err(|error| {
Error::Busy(format!(
"cannot exclusively lock {}: {error}",
root.display()
))
})?;
store::open_root(&root, &config.writers_by_priority)?;
let signing_key = SigningKey::from_bytes(&config.signing_key);
let database = Self {
inner: Arc::new(Inner {
root,
_lock_file: lock_file,
mutation: Mutex::new(()),
visibility: RwLock::new(()),
pending: Mutex::new(PendingPool::default()),
outbox_draining: AtomicBool::new(false),
signing_key,
local_writer: WriterId::from_signing_key(&config.signing_key),
writers: config.writers_by_priority,
gossip: config.gossip,
poisoned: AtomicBool::new(false),
}),
};
database.drain_one_outbox()?;
Ok(database)
}
pub fn start_transaction(&self, provenance: Provenance) -> Result<Transaction<'_>> {
provenance.validate()?;
let guard = self
.inner
.mutation
.lock()
.map_err(|_| Error::corrupt("mutation mutex is poisoned"))?;
self.ensure_healthy()?;
let state = store::read_state(&self.inner.root)?;
Ok(Transaction {
db: self,
guard: Some(guard),
provenance,
heads: state.heads,
object_bytes: 0,
objects: BTreeMap::new(),
reserved_nodes: BTreeSet::new(),
creates: BTreeMap::new(),
updates: BTreeMap::new(),
merges: BTreeSet::new(),
})
}
pub fn accept_transaction(&self, package: TransactionPackage) -> Result<bool> {
self.accept_transaction_inner(package, None)
}
pub fn accept_gossip_transaction(
&self,
package: TransactionPackage,
source: &dyn TransactionSource,
) -> Result<bool> {
self.accept_transaction_inner(package, Some(source))
}
fn accept_transaction_inner(
&self,
package: TransactionPackage,
source: Option<&dyn TransactionSource>,
) -> Result<bool> {
let guard = self
.inner
.mutation
.lock()
.map_err(|_| Error::corrupt("mutation mutex is poisoned"))?;
self.ensure_healthy()?;
let parsed = wire::parse_signed_authorized(&package.transaction, &self.inner.writers)?;
if parsed.unsigned.heads.contains(&parsed.id) {
return Err(Error::invalid_transaction(
"transaction cannot name itself as a head",
));
}
if store::tx_exists(&self.inner.root, parsed.id) {
if store::read_signed_bytes(&self.inner.root, parsed.id)? == package.transaction {
drop(guard);
return Ok(false);
}
return Err(Error::corrupt(
"retained transaction ID has different signed bytes",
));
}
{
let pending = self
.inner
.pending
.lock()
.map_err(|_| Error::corrupt("pending mutex is poisoned"))?;
if let Some(existing) = pending.transactions.get(&parsed.id) {
if existing.transaction == package.transaction {
drop(pending);
drop(guard);
return Ok(false);
}
return Err(Error::corrupt(
"pending transaction ID has different signed bytes",
));
}
}
let missing = self.missing_heads(&parsed)?;
if source.is_none()
&& let Some(first) = missing.first()
{
return Err(Error::invalid_transaction(format!(
"transaction depends on {} uncommitted head(s), beginning with {first}; \
only gossip admission may wait for dependencies",
missing.len()
)));
}
verify_package(&parsed, &package)?;
self.ensure_object_ids_available(parsed.id, &parsed.unsigned.objects)?;
let requests = if missing.is_empty() {
let TransactionPackage {
transaction,
objects,
} = package;
let id = parsed.id;
self.commit_new(parsed, &transaction, CommitObjects::Memory(&objects))?;
self.process_committed_transaction(id)?;
Vec::new()
} else {
self.enqueue_pending(parsed, package, missing)?
};
drop(guard);
if !requests.is_empty() {
source
.ok_or_else(|| Error::corrupt("missing gossip transaction source"))?
.request_transactions(requests);
}
self.drain_one_outbox()?;
Ok(true)
}
pub fn get_node(&self, id: NodeId) -> Result<Node> {
let _guard = self
.inner
.visibility
.read()
.map_err(|_| Error::corrupt("visibility lock is poisoned"))?;
self.ensure_healthy()?;
match store::read_node(&self.inner.root, id) {
Ok(node) => Ok(node.node),
Err(Error::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => {
if store::read_history_index_optional(&self.inner.root, id)?
.is_some_and(|history| history.visible.is_some())
{
Err(Error::corrupt(format!(
"authoritative node file {id} is missing"
)))
} else {
Err(Error::not_found(format!("node {id}")))
}
}
Err(error) => Err(error),
}
}
pub fn get_node_history(&self, id: NodeId) -> Result<NodeHistory> {
let _guard = self
.inner
.visibility
.read()
.map_err(|_| Error::corrupt("visibility lock is poisoned"))?;
self.ensure_healthy()?;
let mut history =
store::read_history(&self.inner.root, id).map_err(|error| match error {
Error::Io(io) if io.kind() == std::io::ErrorKind::NotFound => {
Error::not_found(format!("node history {id}"))
}
other => other,
})?;
sort_history(&self.inner.root, &mut history)?;
Ok(history)
}
pub fn get_object(&self, id: ObjectId) -> Result<Vec<u8>> {
let _guard = self
.inner
.visibility
.read()
.map_err(|_| Error::corrupt("visibility lock is poisoned"))?;
self.ensure_healthy()?;
let (creator, bytes) =
store::read_object(&self.inner.root, id).map_err(|error| match error {
Error::Io(io) if io.kind() == std::io::ErrorKind::NotFound => {
Error::not_found(format!("object {id}"))
}
other => other,
})?;
if !store::transaction_committed(&self.inner.root, creator)? {
return Err(Error::not_found(format!("object {id} is not committed")));
}
Ok(bytes)
}
pub fn get_object_with_provenance(&self, id: ObjectId) -> Result<(Vec<u8>, Provenance)> {
let _guard = self
.inner
.visibility
.read()
.map_err(|_| Error::corrupt("visibility lock is poisoned"))?;
self.ensure_healthy()?;
let (creator, bytes) =
store::read_object(&self.inner.root, id).map_err(|error| match error {
Error::Io(io) if io.kind() == std::io::ErrorKind::NotFound => {
Error::not_found(format!("object {id}"))
}
other => other,
})?;
if !store::transaction_committed(&self.inner.root, creator)? {
return Err(Error::not_found(format!("object {id} is not committed")));
}
let signed = store::read_signed_bytes(&self.inner.root, creator)?;
let parsed = wire::parse_signed(&signed)?;
if parsed.id != creator {
return Err(Error::corrupt(format!(
"object {id} names a mismatched creating transaction"
)));
}
Ok((bytes, parsed.unsigned.provenance))
}
fn commit_new(
&self,
parsed: ParsedTransaction,
transaction: &[u8],
objects: CommitObjects<'_>,
) -> Result<()> {
let mut state = store::read_state(&self.inner.root)?;
for parent in &parsed.unsigned.heads {
if !store::transaction_committed(&self.inner.root, *parent)? {
return Err(Error::invalid_transaction(format!(
"transaction head {parent} has not been processed"
)));
}
}
let dag_generation = next_dag_generation(
parsed
.unsigned
.heads
.iter()
.map(|parent| store::dag_generation(&self.inner.root, *parent))
.collect::<Result<Vec<_>>>()?,
)?;
let meta = TxMeta {
format: store::record_version(),
id: parsed.id,
parents: parsed.unsigned.heads.clone(),
dag_generation,
};
let (nodes, histories) = self.project_transaction(&parsed)?;
let frame_length = store::log_frame_length(transaction.len())?;
let mut commit = store::Commit::new(&self.inner.root, parsed.id, &state)?;
commit.stage_log(parsed.id, transaction)?;
match objects {
CommitObjects::Memory(objects) => {
for payload in objects {
if store::object_exists(&self.inner.root, payload.id) {
return Err(Error::invalid_transaction(format!(
"object locator collision at {}",
payload.id
)));
}
commit.stage_object(store::object_rel(payload.id), parsed.id, payload)?;
}
}
CommitObjects::Spooled(objects) => {
if objects.len() != parsed.unsigned.objects.len() {
return Err(Error::corrupt(
"spooled object count differs from signed transaction",
));
}
for (object, declaration) in objects.iter().zip(&parsed.unsigned.objects) {
if object.id != declaration.id {
return Err(Error::corrupt(
"spooled object order differs from signed transaction",
));
}
if store::object_exists(&self.inner.root, object.id) {
return Err(Error::invalid_transaction(format!(
"object locator collision at {}",
object.id
)));
}
commit.stage_spooled_object(store::object_rel(object.id), object)?;
}
}
}
commit.stage_signed_transaction(store::tx_bytes_rel(parsed.id), parsed.id, transaction)?;
commit.stage_bytes(
store::tx_meta_rel(parsed.id),
&store::tx_meta_bytes(&meta)?,
true,
)?;
commit.stage_bytes(
store::outbox_rel(parsed.id),
&store::queue_bytes(parsed.id)?,
true,
)?;
stage_projection(&mut commit, nodes, histories)?;
for parent in &parsed.unsigned.heads {
state.heads.retain(|head| head != parent);
}
state.heads.push(parsed.id);
state.heads.sort();
state.heads.dedup();
state.generation = state
.generation
.checked_add(1)
.ok_or_else(|| Error::corrupt("database generation overflow"))?;
state.log_offset = state
.log_offset
.checked_add(frame_length)
.ok_or_else(|| Error::corrupt("transaction log offset overflow"))?;
commit.stage_bytes(
PathBuf::from("state.kws"),
&store::state_bytes(&state)?,
false,
)?;
let _visibility = self
.inner
.visibility
.write()
.map_err(|_| Error::corrupt("visibility lock is poisoned"))?;
self.finish_commit(commit)
}
fn project_transaction(
&self,
parsed: &ParsedTransaction,
) -> Result<(BTreeMap<NodeId, NodeFile>, BTreeMap<NodeId, HistoryUpdate>)> {
let overlay = store::DagOverlay {
id: parsed.id,
parents: &parsed.unsigned.heads,
};
let current_creates = parsed
.unsigned
.creates
.iter()
.map(|operation| operation.id)
.collect::<BTreeSet<_>>();
let current_objects = parsed
.unsigned
.objects
.iter()
.map(|object| object.id)
.collect::<BTreeSet<_>>();
let current_creates =
self.resolvable_current_creates(parsed, current_creates, ¤t_objects, &overlay)?;
let mut nodes = BTreeMap::new();
let mut histories = BTreeMap::new();
for (created, operation) in parsed
.unsigned
.creates
.iter()
.map(|operation| (true, operation))
.chain(
parsed
.unsigned
.updates
.iter()
.map(|operation| (false, operation)),
)
{
let mut history = store::read_history_index_optional(&self.inner.root, operation.id)?
.unwrap_or_else(|| store::empty_history(operation.id));
let existing_node = store::read_node_optional(&self.inner.root, operation.id)?;
let references = self.references_resolve(
parsed.id,
operation.id,
&operation.data,
¤t_creates,
¤t_objects,
&overlay,
)?;
let effective = if created {
references && !self.has_ancestor_create(&history, parsed.id, &overlay)?
} else {
references && self.has_ancestor_create(&history, parsed.id, &overlay)?
};
let node = if effective {
Some(self.apply_candidate(parsed, operation, existing_node, &overlay)?)
} else {
existing_node
};
let existing_entry =
store::read_history_entry_optional(&self.inner.root, operation.id, parsed.id)?;
if existing_entry.is_some() {
return Err(Error::corrupt(
"new transaction collides with an existing history entry",
));
}
let entry = store::history_entry(parsed, operation.data.clone(), created);
if effective && created {
history.creations.push(parsed.id);
history.creations.sort();
history.creations.dedup();
}
if let Some(node) = &node {
history.frontier = node
.frontier
.iter()
.map(|candidate| candidate.transaction)
.collect();
history.visible = Some(node.visible_transaction);
}
if let Some(node) = node {
nodes.insert(operation.id, node);
}
histories.insert(
operation.id,
HistoryUpdate {
index: history,
entry,
create_entry: true,
},
);
}
Ok((nodes, histories))
}
fn resolvable_current_creates(
&self,
parsed: &ParsedTransaction,
mut resolvable: BTreeSet<NodeId>,
current_objects: &BTreeSet<ObjectId>,
overlay: &store::DagOverlay<'_>,
) -> Result<BTreeSet<NodeId>> {
loop {
let mut invalid = Vec::new();
for operation in &parsed.unsigned.creates {
if resolvable.contains(&operation.id)
&& !self.references_resolve(
parsed.id,
operation.id,
&operation.data,
&resolvable,
current_objects,
overlay,
)?
{
invalid.push(operation.id);
}
}
if invalid.is_empty() {
return Ok(resolvable);
}
for id in invalid {
resolvable.remove(&id);
}
}
}
fn apply_candidate(
&self,
parsed: &ParsedTransaction,
operation: &NodeOperation,
existing: Option<NodeFile>,
overlay: &store::DagOverlay<'_>,
) -> Result<NodeFile> {
let mut frontier = existing.map_or_else(Vec::new, |node| node.frontier);
for candidate in &frontier {
if store::is_ancestor(
&self.inner.root,
parsed.id,
candidate.transaction,
Some(overlay),
)? {
return preferred_node(operation.id, frontier, &self.inner.writers);
}
}
let mut retained = Vec::with_capacity(frontier.len() + 1);
for candidate in frontier.drain(..) {
if !store::is_ancestor(
&self.inner.root,
candidate.transaction,
parsed.id,
Some(overlay),
)? {
retained.push(candidate);
}
}
retained.push(Candidate {
transaction: parsed.id,
writer: parsed.unsigned.writer,
committed_at: parsed.unsigned.committed_at,
provenance: parsed.unsigned.provenance.clone(),
data: operation.data.clone(),
});
retained.sort_by_key(|candidate| candidate.transaction);
preferred_node(operation.id, retained, &self.inner.writers)
}
fn has_ancestor_create(
&self,
history: &HistoryIndex,
transaction: TransactionId,
overlay: &store::DagOverlay<'_>,
) -> Result<bool> {
for creation in &history.creations {
if *creation != transaction
&& store::is_ancestor(&self.inner.root, *creation, transaction, Some(overlay))?
{
return Ok(true);
}
}
Ok(false)
}
fn references_resolve(
&self,
transaction: TransactionId,
self_id: NodeId,
data: &NodeData,
current_creates: &BTreeSet<NodeId>,
current_objects: &BTreeSet<ObjectId>,
overlay: &store::DagOverlay<'_>,
) -> Result<bool> {
let mut node_refs = data
.fixed_connections
.iter()
.chain(&data.recent_connections)
.copied()
.collect::<Vec<_>>();
if let Owner::Node(owner) = data.owner {
node_refs.push(owner);
}
for reference in node_refs {
if reference == self_id || current_creates.contains(&reference) {
continue;
}
let Some(history) = store::read_history_index_optional(&self.inner.root, reference)?
else {
return Ok(false);
};
if !self.has_ancestor_create(&history, transaction, overlay)? {
return Ok(false);
}
}
for object in &data.objects {
if current_objects.contains(object) {
continue;
}
let creator = match store::object_creator(&self.inner.root, *object) {
Ok(creator) => creator,
Err(Error::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => {
return Ok(false);
}
Err(error) => return Err(error),
};
if !store::transaction_committed(&self.inner.root, creator)?
|| !store::is_ancestor(&self.inner.root, creator, transaction, Some(overlay))?
{
return Ok(false);
}
}
Ok(true)
}
fn missing_heads(&self, parsed: &ParsedTransaction) -> Result<BTreeSet<TransactionId>> {
let mut missing = BTreeSet::new();
for head in &parsed.unsigned.heads {
if !store::transaction_committed(&self.inner.root, *head)? {
missing.insert(*head);
}
}
Ok(missing)
}
fn enqueue_pending(
&self,
parsed: ParsedTransaction,
package: TransactionPackage,
missing: BTreeSet<TransactionId>,
) -> Result<Vec<TransactionId>> {
let TransactionPackage {
transaction,
objects,
} = package;
let objects = store::spool_objects(&self.inner.root, parsed.id, objects)?;
let mut pending = self
.inner
.pending
.lock()
.map_err(|_| Error::corrupt("pending mutex is poisoned"))?;
let requests = missing
.iter()
.filter(|head| !pending.transactions.contains_key(head))
.copied()
.collect::<Vec<_>>();
for head in &missing {
pending
.waiting_on
.entry(*head)
.or_default()
.insert(parsed.id);
}
for object in &objects {
pending.object_reservations.insert(object.id, parsed.id);
}
pending.transactions.insert(
parsed.id,
PendingTransaction {
id: parsed.id,
transaction,
objects,
missing,
},
);
Ok(requests)
}
fn ensure_object_ids_available(
&self,
transaction: TransactionId,
declarations: &[ObjectDeclaration],
) -> Result<()> {
let pending = self
.inner
.pending
.lock()
.map_err(|_| Error::corrupt("pending mutex is poisoned"))?;
for declaration in declarations {
if store::object_exists(&self.inner.root, declaration.id)
|| pending
.object_reservations
.get(&declaration.id)
.is_some_and(|owner| *owner != transaction)
{
return Err(Error::invalid_transaction(format!(
"object locator collision at {}",
declaration.id
)));
}
}
Ok(())
}
fn process_committed_transaction(&self, committed: TransactionId) -> Result<()> {
self.release_pending_dependents(committed)?;
loop {
let pending_transaction = {
let mut pending = self
.inner
.pending
.lock()
.map_err(|_| Error::corrupt("pending mutex is poisoned"))?;
let Some(id) = pending.ready.pop_front() else {
return Ok(());
};
pending.transactions.remove(&id).ok_or_else(|| {
Error::corrupt("ready transaction is missing from the pending pool")
})?
};
let parsed = match wire::parse_signed_authorized(
&pending_transaction.transaction,
&self.inner.writers,
) {
Ok(parsed) => parsed,
Err(error) => {
self.requeue_pending(pending_transaction)?;
return Err(error);
}
};
let missing = match self.missing_heads(&parsed) {
Ok(missing) => missing,
Err(error) => {
self.requeue_pending(pending_transaction)?;
return Err(error);
}
};
if parsed.id != pending_transaction.id || !missing.is_empty() {
self.requeue_pending(pending_transaction)?;
return Err(Error::corrupt(
"ready transaction still has an unresolved dependency",
));
}
let id = pending_transaction.id;
let result = self.commit_new(
parsed,
&pending_transaction.transaction,
CommitObjects::Spooled(&pending_transaction.objects),
);
if let Err(error) = result {
self.requeue_pending(pending_transaction)?;
return Err(error);
}
self.complete_pending(id)?;
let _ = store::discard_spooled_objects(&self.inner.root, id);
self.release_pending_dependents(id)?;
}
}
fn release_pending_dependents(&self, committed: TransactionId) -> Result<()> {
let mut pending = self
.inner
.pending
.lock()
.map_err(|_| Error::corrupt("pending mutex is poisoned"))?;
let Some(dependents) = pending.waiting_on.remove(&committed) else {
return Ok(());
};
for dependent in dependents {
let transaction = pending.transactions.get_mut(&dependent).ok_or_else(|| {
Error::corrupt("dependency waiter is missing from the pending pool")
})?;
transaction.missing.remove(&committed);
if transaction.missing.is_empty() {
pending.ready.push_back(dependent);
}
}
Ok(())
}
fn complete_pending(&self, id: TransactionId) -> Result<()> {
let mut pending = self
.inner
.pending
.lock()
.map_err(|_| Error::corrupt("pending mutex is poisoned"))?;
pending
.object_reservations
.retain(|_, transaction| *transaction != id);
Ok(())
}
fn requeue_pending(&self, transaction: PendingTransaction) -> Result<()> {
let mut pending = self
.inner
.pending
.lock()
.map_err(|_| Error::corrupt("pending mutex is poisoned"))?;
pending.ready.push_front(transaction.id);
pending.transactions.insert(transaction.id, transaction);
Ok(())
}
fn drain_one_outbox(&self) -> Result<()> {
if self
.inner
.outbox_draining
.compare_exchange(false, true, AtomicOrdering::AcqRel, AtomicOrdering::Acquire)
.is_err()
{
return Ok(());
}
let _drain_guard = AtomicFlagGuard(&self.inner.outbox_draining);
store::recover_outbox_claim(&self.inner.root)?;
let Some(id) = store::next_outbox(&self.inner.root)? else {
return Ok(());
};
store::claim_outbox(&self.inner.root, id)?;
let package = match store::package_for(&self.inner.root, id) {
Ok(package) => package,
Err(error) => {
store::requeue_outbox_claim(&self.inner.root, id)?;
return Err(error);
}
};
if self.inner.gossip.announce(package) {
store::acknowledge_outbox_claim(&self.inner.root)?;
} else {
store::requeue_outbox_claim(&self.inner.root, id)?;
}
Ok(())
}
fn finish_commit(&self, commit: store::Commit) -> Result<()> {
match commit.finish() {
Ok(()) => Ok(()),
Err(failure) => {
if failure.prepared {
self.inner.poisoned.store(true, AtomicOrdering::Release);
}
Err(failure.error)
}
}
}
fn ensure_healthy(&self) -> Result<()> {
if self.inner.poisoned.load(AtomicOrdering::Acquire) {
Err(Error::corrupt(
"a prepared WAL could not be fully applied; close and reopen the database",
))
} else {
Ok(())
}
}
}
struct AtomicFlagGuard<'a>(&'a AtomicBool);
impl Drop for AtomicFlagGuard<'_> {
fn drop(&mut self) {
self.0.store(false, AtomicOrdering::Release);
}
}
pub struct Transaction<'a> {
db: &'a KwebDb,
guard: Option<MutexGuard<'a, ()>>,
provenance: Provenance,
heads: Vec<TransactionId>,
object_bytes: u64,
objects: BTreeMap<ObjectId, Vec<u8>>,
reserved_nodes: BTreeSet<NodeId>,
creates: BTreeMap<NodeId, NodeData>,
updates: BTreeMap<NodeId, NodeData>,
merges: BTreeSet<MergePair>,
}
impl fmt::Debug for Transaction<'_> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("Transaction")
.field("heads", &self.heads)
.field("object_count", &self.objects.len())
.field("object_bytes", &self.object_bytes)
.field("reserved_nodes", &self.reserved_nodes.len())
.field("creates", &self.creates.len())
.field("updates", &self.updates.len())
.field("merges", &self.merges.len())
.finish_non_exhaustive()
}
}
impl Transaction<'_> {
pub fn create_object(&mut self, bytes: Vec<u8>) -> Result<ObjectId> {
let length = bytes.len() as u64;
if length > crate::MAX_OBJECT_BYTES {
return Err(Error::invalid_input("object exceeds the 32 GiB limit"));
}
let object_bytes = self
.object_bytes
.checked_add(length)
.ok_or_else(|| Error::invalid_input("transaction object length overflow"))?;
if object_bytes > crate::MAX_TRANSACTION_OBJECT_BYTES {
return Err(Error::invalid_input(
"transaction object payload total exceeds 32 GiB",
));
}
let id = loop {
let candidate = ObjectId::random();
if !self.objects.contains_key(&candidate)
&& !store::object_exists(&self.db.inner.root, candidate)
&& !self
.db
.inner
.pending
.lock()
.map_err(|_| Error::corrupt("pending mutex is poisoned"))?
.object_reservations
.contains_key(&candidate)
{
break candidate;
}
};
self.object_bytes = object_bytes;
self.objects.insert(id, bytes);
Ok(id)
}
pub fn reserve_node_id(&mut self) -> Result<NodeId> {
let id = loop {
let candidate = NodeId::random();
if !self.reserved_nodes.contains(&candidate)
&& !self.creates.contains_key(&candidate)
&& !store::node_id_occupied(&self.db.inner.root, candidate)
{
break candidate;
}
};
self.reserved_nodes.insert(id);
Ok(id)
}
pub fn create_reserved_node(&mut self, id: NodeId, data: NodeData) -> Result<()> {
if self.creates.contains_key(&id) {
return Err(Error::invalid_input(format!(
"reserved node {id} has already been materialized"
)));
}
if !self.reserved_nodes.contains(&id) {
return Err(Error::invalid_input(format!(
"node ID {id} was not reserved by this transaction"
)));
}
data.validate()?;
self.reserved_nodes.remove(&id);
self.creates.insert(id, data);
Ok(())
}
pub fn create_node(&mut self, data: NodeData) -> Result<NodeId> {
data.validate()?;
let id = self.reserve_node_id()?;
self.create_reserved_node(id, data)?;
Ok(id)
}
pub fn update_node(&mut self, id: NodeId, data: NodeData) -> Result<()> {
if self.reserved_nodes.contains(&id) || self.creates.contains_key(&id) {
return Err(Error::invalid_input(
"one transaction cannot create and update the same node",
));
}
data.validate()?;
if self.updates.insert(id, data).is_some() {
return Err(Error::invalid_input(
"one transaction cannot update the same node twice",
));
}
Ok(())
}
pub fn merge(&mut self, first: TransactionId, second: TransactionId) -> Result<()> {
let pair = MergePair::new(first, second)?;
if !self.merges.insert(pair) {
return Err(Error::invalid_input("duplicate merge pair"));
}
Ok(())
}
pub fn finalize(mut self) -> Result<TransactionId> {
self.validate_local()?;
let unsigned = UnsignedTransaction {
writer: self.db.inner.local_writer,
committed_at: Utc::now(),
heads: self.heads.clone(),
provenance: self.provenance.clone(),
merge_pairs: self.merges.iter().copied().collect(),
objects: self
.objects
.iter()
.map(|(id, bytes)| ObjectDeclaration {
id: *id,
length: bytes.len() as u64,
sha256: wire::object_hash(bytes),
})
.collect(),
creates: self
.creates
.iter()
.map(|(id, data)| NodeOperation {
id: *id,
data: data.clone(),
})
.collect(),
updates: self
.updates
.iter()
.map(|(id, data)| NodeOperation {
id: *id,
data: data.clone(),
})
.collect(),
};
let transaction = wire::build_signed(&unsigned, &self.db.inner.signing_key)?;
let parsed = wire::parse_signed(&transaction)?;
let id = parsed.id;
let objects = self
.objects
.into_iter()
.map(|(id, bytes)| ObjectPayload { id, bytes })
.collect::<Vec<_>>();
self.db
.commit_new(parsed, &transaction, CommitObjects::Memory(&objects))?;
drop(objects);
self.db.process_committed_transaction(id)?;
drop(self.guard.take());
self.db.drain_one_outbox()?;
Ok(id)
}
fn validate_local(&self) -> Result<()> {
if !self.reserved_nodes.is_empty() {
return Err(Error::invalid_input(format!(
"transaction has {} unmaterialized reserved node ID(s)",
self.reserved_nodes.len()
)));
}
let created = self.creates.keys().copied().collect::<BTreeSet<_>>();
let objects = self.objects.keys().copied().collect::<BTreeSet<_>>();
for (id, data) in self.creates.iter().chain(&self.updates) {
if self.updates.contains_key(id) && !store::node_exists(&self.db.inner.root, *id) {
return Err(Error::invalid_input(format!(
"cannot update nonvisible node {id}"
)));
}
validate_local_references(&self.db.inner.root, *id, data, &created, &objects)?;
}
for pair in &self.merges {
let valid = self.updates.keys().any(|id| {
store::read_node(&self.db.inner.root, *id)
.map(|node| {
let frontier = node
.frontier
.iter()
.map(|candidate| candidate.transaction)
.collect::<BTreeSet<_>>();
frontier.contains(&pair.first) && frontier.contains(&pair.second)
})
.unwrap_or(false)
});
if !valid {
return Err(Error::invalid_input(
"merge pair is not an exact current frontier for an updated node",
));
}
}
Ok(())
}
}
fn validate_local_references(
root: &Path,
self_id: NodeId,
data: &NodeData,
created: &BTreeSet<NodeId>,
objects: &BTreeSet<ObjectId>,
) -> Result<()> {
let mut node_refs = data
.fixed_connections
.iter()
.chain(&data.recent_connections)
.copied()
.collect::<Vec<_>>();
if let Owner::Node(owner) = data.owner {
node_refs.push(owner);
}
for reference in node_refs {
if reference != self_id
&& !created.contains(&reference)
&& !store::node_exists(root, reference)
{
return Err(Error::invalid_input(format!(
"node reference {reference} is not locally resolvable"
)));
}
}
for object in &data.objects {
if !objects.contains(object) && !store::object_exists(root, *object) {
return Err(Error::invalid_input(format!(
"object {object} is not locally resolvable"
)));
}
}
Ok(())
}
fn preferred_node(id: NodeId, frontier: Vec<Candidate>, writers: &[WriterId]) -> Result<NodeFile> {
let visible = frontier
.iter()
.max_by(|left, right| compare_preference(left, right, writers))
.ok_or_else(|| Error::corrupt("node frontier is empty"))?;
Ok(NodeFile {
format: store::record_version(),
id,
node: Node {
id,
data: visible.data.clone(),
last_author: visible.provenance.author.clone(),
committed_at: visible.committed_at,
},
visible_transaction: visible.transaction,
frontier,
})
}
fn compare_preference(left: &Candidate, right: &Candidate, writers: &[WriterId]) -> Ordering {
let left_rank = writers
.iter()
.position(|writer| *writer == left.writer)
.unwrap_or(usize::MAX);
let right_rank = writers
.iter()
.position(|writer| *writer == right.writer)
.unwrap_or(usize::MAX);
right_rank
.cmp(&left_rank)
.then_with(|| left.transaction.cmp(&right.transaction))
}
fn next_dag_generation(generations: Vec<u64>) -> Result<u64> {
match generations.into_iter().max() {
Some(generation) => generation
.checked_add(1)
.ok_or_else(|| Error::corrupt("transaction DAG generation overflow")),
None => Ok(0),
}
}
fn sort_history(root: &Path, history: &mut NodeHistory) -> Result<()> {
let mut generations = BTreeMap::new();
for entry in &history.entries {
let generation = store::read_tx_meta(root, entry.transaction_id)?.dag_generation;
generations.insert(entry.transaction_id, generation);
}
history.entries.sort_by(|left, right| {
(generations[&right.transaction_id], right.transaction_id)
.cmp(&(generations[&left.transaction_id], left.transaction_id))
});
Ok(())
}
fn stage_projection(
commit: &mut store::Commit,
nodes: BTreeMap<NodeId, NodeFile>,
histories: BTreeMap<NodeId, HistoryUpdate>,
) -> Result<()> {
for (id, node) in nodes {
commit.stage_bytes(store::node_rel(id), &store::node_bytes(&node)?, false)?;
}
for (id, history) in histories {
commit.stage_bytes(
store::history_index_rel(id),
&store::history_index_bytes(&history.index)?,
false,
)?;
commit.stage_bytes(
store::history_entry_rel(id, history.entry.transaction_id),
&store::history_entry_bytes(&history.entry)?,
history.create_entry,
)?;
}
Ok(())
}
struct HistoryUpdate {
index: HistoryIndex,
entry: crate::HistoryEntry,
create_entry: bool,
}
fn verify_package(parsed: &ParsedTransaction, package: &TransactionPackage) -> Result<()> {
if parsed.unsigned.objects.len() != package.objects.len() {
return Err(Error::invalid_transaction(
"package does not contain exactly every declared object",
));
}
let mut total = 0_u64;
for (declaration, payload) in parsed.unsigned.objects.iter().zip(&package.objects) {
if declaration.id != payload.id {
return Err(Error::invalid_transaction(
"package object order or ID differs",
));
}
let length = payload.bytes.len() as u64;
total = total
.checked_add(length)
.ok_or_else(|| Error::invalid_transaction("package object length overflow"))?;
if length != declaration.length
|| length > crate::MAX_OBJECT_BYTES
|| total > crate::MAX_TRANSACTION_OBJECT_BYTES
{
return Err(Error::invalid_transaction(
"package object length is invalid",
));
}
if wire::object_hash(&payload.bytes) != declaration.sha256 {
return Err(Error::invalid_transaction(
"package object SHA-256 mismatch",
));
}
}
Ok(())
}
fn validate_config(config: &Config) -> Result<()> {
if config.writers_by_priority.is_empty() {
return Err(Error::invalid_config("writers_by_priority cannot be empty"));
}
let writers = config
.writers_by_priority
.iter()
.copied()
.collect::<BTreeSet<_>>();
if writers.len() != config.writers_by_priority.len() {
return Err(Error::invalid_config("writers_by_priority must be unique"));
}
if config
.writers_by_priority
.iter()
.any(|writer| ed25519_dalek::VerifyingKey::from_bytes(&writer.0).is_err())
{
return Err(Error::invalid_config(
"writers_by_priority contains an invalid Ed25519 public key",
));
}
if !writers.contains(&WriterId::from_signing_key(&config.signing_key)) {
return Err(Error::invalid_config(
"writers_by_priority must contain the local writer",
));
}
Ok(())
}
fn open_lock_file(path: &Path) -> Result<File> {
let file = match fs::symlink_metadata(path) {
Ok(metadata) => {
if metadata.file_type().is_symlink() || !metadata.file_type().is_file() {
return Err(Error::invalid_config(
"database LOCK must be a regular file",
));
}
OpenOptions::new().read(true).write(true).open(path)?
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => OpenOptions::new()
.read(true)
.write(true)
.create_new(true)
.open(path)?,
Err(error) => return Err(error.into()),
};
if !file.metadata()?.file_type().is_file() {
return Err(Error::invalid_config(
"database LOCK must be a regular file",
));
}
Ok(file)
}