use serde::Serialize;
use bson::Document;
use std::borrow::Borrow;
use std::sync::{Mutex, Weak};
use serde::de::DeserializeOwned;
use crate::{ClientSession, DbErr, DbResult};
use crate::db::db_inner::DatabaseInner;
use crate::results::{DeleteResult, InsertManyResult, InsertOneResult, UpdateResult};
pub struct Collection<T> {
db: Weak<Mutex<DatabaseInner>>,
name: String,
_phantom: std::marker::PhantomData<T>,
}
impl<T> Collection<T>
{
pub(super) fn new(db: Weak<Mutex<DatabaseInner>>, name: &str) -> Collection<T> {
Collection {
db,
name: name.into(),
_phantom: std::default::Default::default(),
}
}
pub fn name(&self) -> &str {
&self.name
}
pub fn count_documents(&self) -> DbResult<u64> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.count_documents(&self.name, &mut session)
}
pub fn count_documents_with_session(&self, session: &mut ClientSession) -> DbResult<u64> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.count_documents(&self.name, &mut session.inner)
}
pub fn update_one(&self, query: Document, update: Document) -> DbResult<UpdateResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.update_one(
&self.name,
Some(&query),
&update,
&mut session,
)
}
pub fn update_one_with_session(&self, query: Document, update: Document, session: &mut ClientSession) -> DbResult<UpdateResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.update_one(
&self.name,
Some(&query),
&update,
&mut session.inner,
)
}
pub fn update_many(&self, query: Document, update: Document) -> DbResult<UpdateResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.update_many(&self.name, query, update, &mut session)
}
pub fn update_many_with_session(&self, query: Document, update: Document, session: &mut ClientSession) -> DbResult<UpdateResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.update_many(&self.name, query, update, &mut session.inner)
}
pub fn delete_one(&self, query: Document) -> DbResult<DeleteResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.delete_one(&self.name, query, &mut session)
}
pub fn delete_one_with_session(&self, query: Document, session: &mut ClientSession) -> DbResult<DeleteResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.delete_one(&self.name, query, &mut session.inner)
}
pub fn delete_many(&self, query: Document) -> DbResult<DeleteResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.delete_many(&self.name, query, &mut session)
}
pub fn delete_many_with_session(&self, query: Document, session: &mut ClientSession) -> DbResult<DeleteResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.delete_many(&self.name, query, &mut session.inner)
}
#[allow(dead_code)]
fn create_index(&self, _keys: &Document, _options: Option<&Document>) -> DbResult<()> {
unimplemented!()
}
pub fn drop(&self) -> DbResult<()> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.drop_collection(&self.name, &mut session)
}
pub fn drop_with_session(&self, session: &mut ClientSession) -> DbResult<()> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.drop_collection(&self.name, &mut session.inner)
}
}
impl<T> Collection<T>
where
T: Serialize,
{
pub fn insert_one(&self, doc: impl Borrow<T>) -> DbResult<InsertOneResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.insert_one(
&self.name,
bson::to_document(doc.borrow())?,
&mut session,
)
}
pub fn insert_one_with_session(&self, doc: impl Borrow<T>, session: &mut ClientSession) -> DbResult<InsertOneResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.insert_one(
&self.name,
bson::to_document(doc.borrow())?,
&mut session.inner,
)
}
pub fn insert_many(&self, docs: impl IntoIterator<Item = impl Borrow<T>>) -> DbResult<InsertManyResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.insert_many(&self.name, docs, &mut session)
}
pub fn insert_many_with_session(
&self,
docs: impl IntoIterator<Item = impl Borrow<T>>,
session: &mut ClientSession
) -> DbResult<InsertManyResult> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.insert_many(&self.name, docs, &mut session.inner)
}
}
impl<T> Collection<T>
where
T: DeserializeOwned,
{
pub fn find_many(&self, filter: impl Into<Option<Document>>) -> DbResult<Vec<T>> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.find_many(&self.name, filter, &mut session)
}
pub fn find_many_with_session(&self, filter: impl Into<Option<Document>>, session: &mut ClientSession) -> DbResult<Vec<T>> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.find_many(&self.name, filter, &mut session.inner)
}
pub fn find_one(&self, filter: impl Into<Option<Document>>) -> DbResult<Option<T>> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
let mut session = db.start_session()?;
db.find_one(&self.name, filter, &mut session)
}
pub fn find_one_with_session(&self, filter: impl Into<Option<Document>>, session: &mut ClientSession) -> DbResult<Option<T>> {
let db_ref = self.db.upgrade().ok_or(DbErr::DbIsClosed)?;
let mut db = db_ref.lock()?;
db.find_one(&self.name, filter, &mut session.inner)
}
}