-
Notifications
You must be signed in to change notification settings - Fork 1
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
11 changed files
with
369 additions
and
66 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
20 changes: 20 additions & 0 deletions
20
src/db/server/source_sink/effects_sink/apply_effects/apply_effects_batch.rs
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,20 @@ | ||
use crate::db::{ | ||
db_error::DbError, reference::Effect, server::lockable_db::transaction_or_db::TransactionOrDb, | ||
}; | ||
|
||
use super::{ | ||
apply_effects_batch_db::apply_effects_batch_db, | ||
apply_effects_batch_transaction_or_db::apply_effects_batch_transaction_or_db, | ||
}; | ||
|
||
pub(super) fn apply_effects_batch<'a>( | ||
transaction_or_db: &TransactionOrDb<'a>, | ||
effects: &[Effect], | ||
) -> Result<(), DbError> { | ||
match transaction_or_db { | ||
TransactionOrDb::Transaction(_, _) => { | ||
apply_effects_batch_transaction_or_db(transaction_or_db, effects) | ||
} | ||
TransactionOrDb::Db(db) => apply_effects_batch_db(db, effects), | ||
} | ||
} |
26 changes: 26 additions & 0 deletions
26
src/db/server/source_sink/effects_sink/apply_effects/apply_effects_batch_db.rs
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,26 @@ | ||
use crate::db::{db_error::DbError, reference::Effect}; | ||
use rocksdb::{TransactionDB, WriteBatchWithTransaction}; | ||
|
||
use super::make_access_effect_batch::make_access_effect_batch; | ||
|
||
pub(super) fn apply_effects_batch_db( | ||
db: &TransactionDB, | ||
effects: &[Effect], | ||
) -> Result<(), DbError> { | ||
let mut batch = WriteBatchWithTransaction::default(); | ||
|
||
for effect in effects { | ||
match effect { | ||
Effect::Access(access) => { | ||
make_access_effect_batch(db, access, &mut batch)?; | ||
} | ||
_ => unreachable!(), | ||
} | ||
} | ||
|
||
db | ||
.write(batch) | ||
.map_err(|e| DbError::TantivyError(e.to_string()))?; | ||
|
||
Ok(()) | ||
} |
50 changes: 50 additions & 0 deletions
50
...db/server/source_sink/effects_sink/apply_effects/apply_effects_batch_transaction_or_db.rs
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,50 @@ | ||
use tonic::Status; | ||
|
||
use crate::db::{ | ||
db_error::DbError, reference::{Effect, AccessEffect}, server::{lockable_db::transaction_or_db::TransactionOrDb, db_error_to_status::DbErrorToStatus}, | ||
}; | ||
|
||
pub(super) fn apply_effects_batch_transaction_or_db<'a>( | ||
transaction_or_db: &TransactionOrDb<'a>, | ||
effects: &[Effect], | ||
) -> Result<(), DbError> { | ||
for effect in effects { | ||
println!("Effect: {:?}", effect); | ||
match effect { | ||
Effect::Access(access) => { | ||
apply_access_effect(&transaction_or_db, &access) | ||
.map_err(|e| DbError::TantivyError(e.to_string()))?; | ||
} | ||
_ => unreachable!(), | ||
} | ||
} | ||
|
||
Ok(()) | ||
} | ||
|
||
pub(crate) fn apply_access_effect<'a>( | ||
db: &TransactionOrDb<'a>, | ||
access_effect: &AccessEffect, | ||
) -> Result<(), Status> { | ||
match access_effect { | ||
AccessEffect::DatabaseServerStoredEffect(effect) => { | ||
super::super::super::database_server_sink::apply_effect(&db, effect).map_db_err_to_status()?; | ||
} | ||
AccessEffect::DomainStoredEffect(effect) => { | ||
super::super::super::domain_sink::apply_effect(&db, effect).map_db_err_to_status()?; | ||
} | ||
AccessEffect::TableStoredEffect(effect) => { | ||
super::super::super::table_sink::apply_effect(&db, effect).map_db_err_to_status()?; | ||
} | ||
AccessEffect::TableValueEffect(effect) => { | ||
super::super::super::table_value_sink::apply_effect(&db, effect).map_db_err_to_status()?; | ||
} | ||
AccessEffect::IndexValueEffect(effect) => { | ||
super::super::super::index_value_sink::apply_effect(&db, effect).map_db_err_to_status()?; | ||
} | ||
AccessEffect::ColumnValueEffect(effect) => { | ||
super::super::super::column_value_sink::apply_effect(&db, effect).map_db_err_to_status()?; | ||
} | ||
} | ||
Ok(()) | ||
} |
39 changes: 39 additions & 0 deletions
39
src/db/server/source_sink/effects_sink/apply_effects/make_access_effect_batch.rs
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,39 @@ | ||
use super::{ | ||
make_column_value_effect_batch::make_column_value_effect_batch, | ||
make_database_server_stored_effect_batch::make_database_server_stored_effect_batch, | ||
make_domain_stored_effect_batch::make_domain_stored_effect_batch, | ||
make_index_value_effect_batch::make_index_value_effect_batch, | ||
make_table_stored_effect_batch::make_table_stored_effect_batch, | ||
make_table_value_effect_batch::make_table_value_effect_batch, | ||
}; | ||
use crate::db::{db_error::DbError, reference::AccessEffect}; | ||
use rocksdb::{TransactionDB, WriteBatchWithTransaction}; | ||
|
||
pub(super) fn make_access_effect_batch<'a>( | ||
db: &TransactionDB, | ||
access_effect: &AccessEffect, | ||
batch: &mut WriteBatchWithTransaction<true>, | ||
) -> Result<(), DbError> { | ||
match access_effect { | ||
AccessEffect::DatabaseServerStoredEffect(effect) => { | ||
make_database_server_stored_effect_batch(db, effect, batch)?; | ||
} | ||
AccessEffect::DomainStoredEffect(effect) => { | ||
make_domain_stored_effect_batch(db, effect, batch)?; | ||
} | ||
AccessEffect::TableStoredEffect(effect) => { | ||
make_table_stored_effect_batch(db, effect, batch)?; | ||
} | ||
AccessEffect::TableValueEffect(effect) => { | ||
make_table_value_effect_batch(db, effect, batch)?; | ||
} | ||
AccessEffect::IndexValueEffect(effect) => { | ||
make_index_value_effect_batch(db, effect, batch)?; | ||
} | ||
AccessEffect::ColumnValueEffect(effect) => { | ||
make_column_value_effect_batch(db, effect, batch)?; | ||
} | ||
} | ||
|
||
Ok(()) | ||
} |
30 changes: 30 additions & 0 deletions
30
src/db/server/source_sink/effects_sink/apply_effects/make_column_value_effect_batch.rs
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,30 @@ | ||
use crate::db::{ | ||
db_error::DbError, entity::OndoKey, reference::ColumnValueEffect, | ||
server::source_sink::ondo_serializer::OndoSerializer, | ||
}; | ||
use rocksdb::WriteBatchWithTransaction; | ||
use serde_json::Value; | ||
|
||
pub(super) fn make_column_value_effect_batch( | ||
db: &rocksdb::TransactionDB, | ||
effect: &ColumnValueEffect, | ||
batch: &mut WriteBatchWithTransaction<true>, | ||
) -> Result<(), DbError> { | ||
match effect { | ||
ColumnValueEffect::Put(cf_name, key, value) => { | ||
let ondo_key = OndoKey::ondo_serialize(key)?; | ||
let ondo_value = Value::ondo_serialize(value)?; | ||
let cf = db.cf_handle(cf_name).ok_or(DbError::CfNotFound)?; | ||
|
||
batch.put_cf(&cf, ondo_key, ondo_value); | ||
} | ||
ColumnValueEffect::Delete(cf_name, key) => { | ||
let ondo_key = OndoKey::ondo_serialize(key)?; | ||
let cf = db.cf_handle(cf_name).ok_or(DbError::CfNotFound)?; | ||
|
||
batch.delete_cf(&cf, ondo_key); | ||
} | ||
} | ||
|
||
Ok(()) | ||
} |
35 changes: 35 additions & 0 deletions
35
...server/source_sink/effects_sink/apply_effects/make_database_server_stored_effect_batch.rs
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,35 @@ | ||
use crate::db::{ | ||
db_error::DbError, | ||
entity::DatabaseServerStored, | ||
reference::{DatabaseServerName, DatabaseServerStoredEffect}, | ||
server::source_sink::ondo_serializer::OndoSerializer, | ||
}; | ||
use rocksdb::{TransactionDB, WriteBatchWithTransaction}; | ||
|
||
pub(super) fn make_database_server_stored_effect_batch( | ||
db: &TransactionDB, | ||
effect: &DatabaseServerStoredEffect, | ||
batch: &mut WriteBatchWithTransaction<true>, | ||
) -> Result<(), DbError> { | ||
match effect { | ||
DatabaseServerStoredEffect::Put(cf_name, key, database_server_stored) => { | ||
let ondo_key = DatabaseServerName::ondo_serialize(key)?; | ||
let ondo_value = DatabaseServerStored::ondo_serialize(database_server_stored)?; | ||
let cf = db | ||
.cf_handle(cf_name) | ||
.ok_or(DbError::CfNotFound)?; | ||
|
||
batch.put_cf(&cf, ondo_key, ondo_value); | ||
} | ||
DatabaseServerStoredEffect::Delete(cf_name, key) => { | ||
let ondo_key = DatabaseServerName::ondo_serialize(key)?; | ||
let cf = db | ||
.cf_handle(cf_name) | ||
.ok_or(DbError::CfNotFound)?; | ||
|
||
batch.delete_cf(&cf, ondo_key); | ||
} | ||
} | ||
|
||
Ok(()) | ||
} |
38 changes: 38 additions & 0 deletions
38
src/db/server/source_sink/effects_sink/apply_effects/make_domain_stored_effect_batch.rs
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,38 @@ | ||
use crate::db::{ | ||
db_error::DbError, | ||
entity::DomainStored, | ||
reference::{DomainName, DomainStoredEffect}, | ||
server::{ | ||
|
||
source_sink::ondo_serializer::OndoSerializer, | ||
}, | ||
}; | ||
use rocksdb::{WriteBatchWithTransaction, TransactionDB}; | ||
|
||
pub(super) fn make_domain_stored_effect_batch<'a>( | ||
db: &TransactionDB, | ||
effect: &DomainStoredEffect, | ||
batch: &mut WriteBatchWithTransaction<true>, | ||
) -> Result<(), DbError> { | ||
match effect { | ||
DomainStoredEffect::Put(cf_name, key, domain_stored) => { | ||
let ondo_key = DomainName::ondo_serialize(key)?; | ||
let ondo_value = DomainStored::ondo_serialize(domain_stored)?; | ||
let cf = db | ||
.cf_handle(cf_name) | ||
.ok_or(DbError::CfNotFound)?; | ||
|
||
batch.put_cf(&cf, ondo_key, ondo_value); | ||
} | ||
DomainStoredEffect::Delete(cf_name, key) => { | ||
let ondo_key = DomainName::ondo_serialize(key)?; | ||
let cf = db | ||
.cf_handle(cf_name) | ||
.ok_or(DbError::CfNotFound)?; | ||
|
||
batch.delete_cf(&cf, ondo_key); | ||
} | ||
} | ||
|
||
Ok(()) | ||
} |
Oops, something went wrong.