use super::*;
use crate::error::{DbError, DbResult};
use crate::storage::serializer::serialize_doc;
use crate::transaction::lock_manager::LockManager;
use crate::transaction::wal::WalWriter;
use crate::transaction::{Operation, Transaction, TransactionId};
use rust_rocksdb::WriteBatch;
use serde_json::Value;
use uuid;
impl Collection {
fn parse_db_coll(&self) -> (String, String) {
let (a, b) = self.name.split_once(':').unwrap_or(("", &self.name));
(a.to_string(), b.to_string())
}
pub fn get_tx(
&self,
tx_id: TransactionId,
lock_manager: &Arc<LockManager>,
key: &str,
) -> DbResult<Option<Document>> {
let (db_name, coll_name) = self.parse_db_coll();
lock_manager.acquire_shared(tx_id, &db_name, &coll_name, key)?;
match self.get(key) {
Ok(doc) => Ok(Some(doc)),
Err(DbError::DocumentNotFound(_)) => Ok(None),
Err(e) => Err(e),
}
}
pub fn insert_tx(
&self,
tx: &mut Transaction,
_wal: &Arc<WalWriter>,
lock_manager: &Arc<LockManager>,
mut data: Value,
) -> DbResult<Document> {
if self.collection_type.read().as_str() == "edge" {
self.validate_edge_document(&data)?;
}
let key = if let Some(obj) = data.as_object_mut() {
if let Some(key_value) = obj.remove("_key") {
if let Some(key_str) = key_value.as_str() {
key_str.to_string()
} else {
return Err(DbError::InvalidDocument(
"_key must be a string".to_string(),
));
}
} else {
uuid::Uuid::new_v7(uuid::Timestamp::now(uuid::NoContext)).to_string()
}
} else {
uuid::Uuid::new_v7(uuid::Timestamp::now(uuid::NoContext)).to_string()
};
let (db_name, coll_name) = self.parse_db_coll();
lock_manager.acquire_exclusive(tx.id, &db_name, &coll_name, &key)?;
let doc = Document::with_key(&self.name, key.clone(), data);
self.check_unique_constraints(&key, &doc.to_value())?;
tx.add_operation(Operation::Insert {
database: db_name,
collection: coll_name,
key: key.clone(),
data: doc.to_value(),
});
Ok(doc)
}
pub fn update_tx(
&self,
tx: &mut Transaction,
_wal: &Arc<WalWriter>,
lock_manager: &Arc<LockManager>,
key: &str,
data: Value,
) -> DbResult<Document> {
if self.collection_type.read().as_str() == "timeseries" {
return Err(DbError::OperationNotSupported(
"Update operations are not allowed on timeseries collections".to_string(),
));
}
let (db_name, coll_name) = self.parse_db_coll();
lock_manager.acquire_exclusive(tx.id, &db_name, &coll_name, key)?;
let mut doc = self.get(key)?;
let old_data = doc.to_value();
doc.update(data);
if self.collection_type.read().as_str() == "edge" {
self.validate_edge_document(&doc.to_value())?;
}
tx.add_operation(Operation::Update {
database: db_name,
collection: coll_name,
key: key.to_string(),
old_data,
new_data: doc.to_value(),
});
Ok(doc)
}
pub fn delete_tx(
&self,
tx: &mut Transaction,
_wal: &Arc<WalWriter>,
lock_manager: &Arc<LockManager>,
key: &str,
) -> DbResult<()> {
let (db_name, coll_name) = self.parse_db_coll();
lock_manager.acquire_exclusive(tx.id, &db_name, &coll_name, key)?;
let doc = self.get(key)?;
let old_data = doc.to_value();
tx.add_operation(Operation::Delete {
database: db_name,
collection: coll_name,
key: key.to_string(),
old_data,
});
Ok(())
}
pub fn apply_transaction_operations(&self, operations: Vec<Operation>) -> DbResult<()> {
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
let mut batch = WriteBatch::default();
let mut change_events = Vec::new();
for op in &operations {
match op {
Operation::Insert { key, data, .. } => {
let document = Document::with_key(&self.name, key.clone(), data.clone());
let doc_bytes = serialize_doc(&document)?;
batch.put_cf(&cf, Self::doc_key(key), &doc_bytes);
if self.is_versioned() {
self.append_version_to_batch(&mut batch, &cf, key, Some(data));
}
let (regular_entries, geo_entries) =
self.compute_index_entries_for_insert(key, data)?;
for (entry_key, entry_value) in regular_entries {
batch.put_cf(&cf, entry_key, entry_value);
}
for (entry_key, entry_value) in geo_entries {
batch.put_cf(&cf, entry_key, entry_value);
}
let fulltext_entries = self.compute_fulltext_entries_for_insert(key, data);
for (entry_key, entry_value) in fulltext_entries {
batch.put_cf(&cf, entry_key, entry_value);
}
let ttl_expiry_entries = self.compute_ttl_expiry_entries_for_insert(key, data);
for (entry_key, _entry_value) in ttl_expiry_entries {
batch.put_cf(&cf, entry_key, Vec::new());
}
change_events.push(ChangeEvent {
type_: ChangeType::Insert,
key: key.clone(),
data: Some(data.clone()),
old_data: None,
});
}
Operation::Update {
key,
old_data,
new_data,
..
} => {
let document = Document::with_key(&self.name, key.clone(), new_data.clone());
let doc_bytes = serialize_doc(&document)?;
batch.put_cf(&cf, Self::doc_key(key), &doc_bytes);
if self.is_versioned() {
self.append_version_to_batch(&mut batch, &cf, key, Some(new_data));
}
let (entries_to_add, keys_to_remove, geo_entries_to_add, geo_keys_to_remove) =
self.compute_index_entries_for_update(key, old_data, new_data)?;
for key_to_remove in keys_to_remove {
batch.delete_cf(&cf, key_to_remove);
}
for geo_key in geo_keys_to_remove {
batch.delete_cf(&cf, geo_key);
}
for (entry_key, entry_value) in entries_to_add {
batch.put_cf(&cf, entry_key, entry_value);
}
for (entry_key, entry_value) in geo_entries_to_add {
batch.put_cf(&cf, entry_key, entry_value);
}
let fulltext_keys_to_remove =
self.compute_fulltext_entries_for_delete(key, old_data);
for key_to_remove in fulltext_keys_to_remove {
batch.delete_cf(&cf, key_to_remove);
}
let fulltext_entries_to_add =
self.compute_fulltext_entries_for_insert(key, new_data);
for (entry_key, entry_value) in fulltext_entries_to_add {
batch.put_cf(&cf, entry_key, entry_value);
}
let (ttl_entries_to_add, ttl_keys_to_remove) =
self.compute_ttl_expiry_entries_for_update(key, old_data, new_data);
for key_to_remove in ttl_keys_to_remove {
batch.delete_cf(&cf, key_to_remove);
}
for (entry_key, _entry_value) in ttl_entries_to_add {
batch.put_cf(&cf, entry_key, Vec::new());
}
change_events.push(ChangeEvent {
type_: ChangeType::Update,
key: key.clone(),
data: Some(new_data.clone()),
old_data: Some(old_data.clone()),
});
}
Operation::Delete { key, old_data, .. } => {
batch.delete_cf(&cf, Self::doc_key(key));
if self.is_versioned() {
self.append_version_to_batch(&mut batch, &cf, key, None);
}
let (regular_keys, geo_keys) =
self.compute_index_entries_for_delete(key, old_data)?;
for key_to_remove in regular_keys {
batch.delete_cf(&cf, key_to_remove);
}
for key_to_remove in geo_keys {
batch.delete_cf(&cf, key_to_remove);
}
let fulltext_keys = self.compute_fulltext_entries_for_delete(key, old_data);
for key_to_remove in fulltext_keys {
batch.delete_cf(&cf, key_to_remove);
}
let ttl_keys = self.compute_ttl_expiry_entries_for_delete(key, old_data);
for key_to_remove in ttl_keys {
batch.delete_cf(&cf, key_to_remove);
}
change_events.push(ChangeEvent {
type_: ChangeType::Delete,
key: key.clone(),
data: None,
old_data: Some(old_data.clone()),
});
}
_ => {} }
}
db.write(&batch).map_err(|e| {
DbError::InternalError(format!("Failed to commit transaction batch: {}", e))
})?;
for op in &operations {
match op {
Operation::Insert { key, data, .. } => {
self.update_vector_indexes_on_upsert(key, data);
self.increment_count();
}
Operation::Update { key, new_data, .. } => {
self.update_vector_indexes_on_delete(key);
self.update_vector_indexes_on_upsert(key, new_data);
}
Operation::Delete { .. } => {
self.decrement_count();
}
_ => {}
}
}
for event in change_events {
let _ = self.change_sender.send(event);
}
Ok(())
}
}