use crate::driver::protocol::{DriverError, Response};
use crate::driver::DriverHandler;
use crate::storage::query_cache;
use crate::sync::protocol::Operation;
fn invalidate_query_cache(collection: &str) {
query_cache::get_query_cache().invalidate_collection(collection);
}
pub fn handle_get(
handler: &DriverHandler,
database: String,
collection: String,
key: String,
) -> Response {
match handler.get_collection(&database, &collection) {
Ok(coll) => match coll.get(&key) {
Ok(doc) => Response::ok(doc.to_value()),
Err(e) => Response::error(DriverError::DatabaseError(e.to_string())),
},
Err(e) => Response::error(e),
}
}
pub fn handle_insert(
handler: &DriverHandler,
database: String,
collection: String,
key: Option<String>,
document: serde_json::Value,
) -> Response {
match handler.get_collection(&database, &collection) {
Ok(coll) => {
let mut doc_data = document;
if let Some(k) = key {
if let Some(obj) = doc_data.as_object_mut() {
obj.insert("_key".to_string(), serde_json::json!(k));
}
}
match coll.insert(doc_data) {
Ok(doc) => {
let value = doc.to_value();
handler.log_replication(
&database,
&collection,
Operation::Insert,
&doc.key,
Some(&value),
);
invalidate_query_cache(&collection);
Response::ok(value)
}
Err(e) => Response::error(DriverError::DatabaseError(e.to_string())),
}
}
Err(e) => Response::error(e),
}
}
pub fn handle_update(
handler: &DriverHandler,
database: String,
collection: String,
key: String,
document: serde_json::Value,
merge: bool,
) -> Response {
match handler.get_collection(&database, &collection) {
Ok(coll) => {
let result = if merge {
match coll.get(&key) {
Ok(existing) => {
let mut merged = existing.data.clone();
if let (Some(base), Some(updates)) =
(merged.as_object_mut(), document.as_object())
{
for (k, v) in updates {
base.insert(k.clone(), v.clone());
}
}
coll.update(&key, merged)
}
Err(e) => Err(e),
}
} else {
coll.update(&key, document)
};
match result {
Ok(doc) => {
let value = doc.to_value();
handler.log_replication(
&database,
&collection,
Operation::Update,
&doc.key,
Some(&value),
);
invalidate_query_cache(&collection);
Response::ok(value)
}
Err(e) => Response::error(DriverError::DatabaseError(e.to_string())),
}
}
Err(e) => Response::error(e),
}
}
pub fn handle_delete(
handler: &DriverHandler,
database: String,
collection: String,
key: String,
) -> Response {
match handler.get_collection(&database, &collection) {
Ok(coll) => match coll.delete(&key) {
Ok(_) => {
handler.log_replication(&database, &collection, Operation::Delete, &key, None);
invalidate_query_cache(&collection);
Response::ok_empty()
}
Err(e) => Response::error(DriverError::DatabaseError(e.to_string())),
},
Err(e) => Response::error(e),
}
}
pub fn handle_list(
handler: &DriverHandler,
database: String,
collection: String,
limit: Option<usize>,
offset: Option<usize>,
) -> Response {
match handler.get_collection(&database, &collection) {
Ok(coll) => {
let all_docs = coll.scan(None);
let total = all_docs.len();
let offset = offset.unwrap_or(0);
let limit = limit.unwrap_or(100);
let docs: Vec<_> = all_docs
.into_iter()
.skip(offset)
.take(limit)
.map(|d| d.to_value())
.collect();
Response::Ok {
data: Some(serde_json::json!(docs)),
count: Some(total),
tx_id: None,
}
}
Err(e) => Response::error(e),
}
}
pub fn handle_bulk_insert(
handler: &DriverHandler,
database: String,
collection: String,
documents: Vec<serde_json::Value>,
) -> Response {
match handler.get_collection(&database, &collection) {
Ok(coll) => {
match coll.insert_batch(documents) {
Ok(docs) => {
handler.log_replication_batch(&database, &collection, Operation::Insert, &docs);
invalidate_query_cache(&collection);
Response::ok_count(docs.len())
}
Err(e) => Response::error(DriverError::DatabaseError(e.to_string())),
}
}
Err(e) => Response::error(e),
}
}