#[cfg(feature = "grid_fs")]
use crate::{core::find::FindManyCursor, error::Error, util::convert_bson_to_oid};
use futures_util::{AsyncRead, AsyncWrite, AsyncWriteExt};
#[cfg(feature = "grid_fs")]
#[cfg(feature = "grid_fs")]
use mongodb::{
bson::{doc, oid::ObjectId, Bson, Document},
gridfs::{FilesCollectionDocument, GridFsBucket, GridFsDownloadStream, GridFsUploadStream},
options::*,
Database,
};
#[cfg(feature = "grid_fs")]
pub struct GridFs {
pub bucket: GridFsBucket,
}
#[cfg(feature = "grid_fs")]
impl GridFs {
fn _build_download_options(
&self,
revision: Option<i32>,
) -> Option<GridFsDownloadByNameOptions> {
if let Some(revision) = revision {
Some(
GridFsDownloadByNameOptions::builder()
.revision(revision)
.build(),
)
} else {
None
}
}
pub fn new(db: &Database, options: Option<GridFsBucketOptions>) -> Self {
let bucket = db.gridfs_bucket(options);
Self { bucket }
}
pub async fn drop(&self) -> Result<(), Error> {
self.bucket.drop().await.map_err(Error::Mongo)
}
pub async fn size(&self) -> Result<(usize, usize), Error> {
let mut total_bytes = 0;
let mut total_files = 0;
let mut cursor = self.find_many(doc! {}, None).await?;
while let Some(FilesCollectionDocument { length, .. }) = cursor.next().await? {
total_files += 1;
total_bytes += length as usize;
}
Ok((total_bytes, total_files))
}
pub async fn exists(&self, filename: &str) -> Result<bool, Error> {
Ok(self
.find_many(
doc! { "filename": filename },
Some(GridFsFindOptions::builder().limit(1).build()),
)
.await?
.next()
.await?
.is_some())
}
pub async fn find_one(
&self,
filter: Document,
sort: Option<Document>,
) -> Result<Option<FilesCollectionDocument>, Error> {
Ok(self
.find_many(
filter,
Some(GridFsFindOptions::builder().limit(1).sort(sort).build()),
)
.await?
.next()
.await?)
}
async fn find_one_by_filename_revision(
&self,
filename: &str,
revision: i32,
) -> Result<Option<FilesCollectionDocument>, Error> {
let (sort, skip) = if revision >= 0 {
(1, revision)
} else {
(-1, -revision - 1)
};
let options = GridFsFindOptions::builder()
.sort(doc! { "uploadDate": sort })
.skip(skip as u64)
.limit(Some(1))
.build();
let mut cursor = self
.find_many(doc! { "filename": filename }, Some(options))
.await?;
cursor.next().await
}
pub async fn upload<S>(
&self,
filename: &str,
mut source: S,
options: Option<GridFsUploadOptions>,
) -> Result<ObjectId, Error>
where
S: AsyncRead + Unpin,
{
let mut upload_stream = self.open_upload_stream(filename, options).await?;
let upload_id = convert_bson_to_oid(upload_stream.id().clone())?;
futures_util::io::copy(&mut source, &mut upload_stream)
.await
.map_err(|err| Error::IO(err.kind()))?;
upload_stream
.close()
.await
.map_err(|err| Error::IO(err.kind()))?;
Ok(upload_id)
}
pub async fn open_upload_stream(
&self,
filename: &str,
options: Option<GridFsUploadOptions>,
) -> Result<GridFsUploadStream, Error> {
self.bucket
.open_upload_stream(filename)
.with_options(options)
.await
.map_err(Error::Mongo)
}
pub async fn download<D>(
&self,
filename: &str,
mut destination: &mut D,
revision: Option<i32>,
) -> Result<usize, Error>
where
D: AsyncWrite + Unpin,
{
let mut download_stream = self.open_download_stream(filename, revision).await?;
futures_util::io::copy(&mut download_stream, &mut destination)
.await
.map(|bytes| bytes as usize)
.map_err(|err| Error::IO(err.kind()))
}
pub async fn open_download_stream(
&self,
filename: &str,
revision: Option<i32>,
) -> Result<GridFsDownloadStream, Error> {
let options = self._build_download_options(revision);
self.bucket
.open_download_stream_by_name(filename)
.with_options(options)
.await
.map_err(Error::Mongo)
}
pub async fn find_many(
&self,
filter: Document,
options: Option<GridFsFindOptions>,
) -> Result<FindManyCursor<FilesCollectionDocument>, Error> {
let cursor = self
.bucket
.find(filter)
.with_options(options)
.await
.map_err(Error::Mongo)?;
Ok(FindManyCursor::from_cursor(cursor))
}
pub async fn delete_many(
&self,
filter: Document,
options: Option<GridFsFindOptions>,
) -> Result<usize, Error> {
let mut cursor = self.find_many(filter, options).await?;
let mut deletion_count = 0;
while let Some(FilesCollectionDocument { filename, .. }) = cursor.next().await? {
if let Some(ref filename) = filename {
deletion_count += self.delete_by_filename(filename, None).await?;
}
}
Ok(deletion_count)
}
pub async fn rename_by_filename(
&self,
filename: &str,
new_filename: &str,
revision: Option<i32>,
) -> Result<usize, Error> {
let mut rename_counter = 0;
if let Some(revision) = revision {
if let Some(FilesCollectionDocument { id, .. }) = self
.find_one_by_filename_revision(filename, revision)
.await?
{
self.rename_by_oid(convert_bson_to_oid(id)?, new_filename)
.await?;
rename_counter += 1;
}
return Ok(rename_counter);
}
let mut cursor = self.find_many(doc! { "filename": filename }, None).await?;
while let Some(FilesCollectionDocument { id, .. }) = cursor.next().await? {
self.rename_by_oid(convert_bson_to_oid(id)?, new_filename)
.await?;
rename_counter += 1;
}
Ok(rename_counter)
}
pub async fn rename_by_oid(&self, oid: ObjectId, new_filename: &str) -> Result<(), Error> {
self.bucket
.rename(Bson::ObjectId(oid), new_filename)
.await
.map_err(Error::Mongo)
}
pub async fn delete_by_filename(
&self,
filename: &str,
revision: Option<i32>,
) -> Result<usize, Error> {
let mut deletion_count = 0;
if let Some(revision) = revision {
if let Some(FilesCollectionDocument { id, .. }) = self
.find_one_by_filename_revision(filename, revision)
.await?
{
self.delete_by_oid(convert_bson_to_oid(id)?).await?;
deletion_count += 1;
}
return Ok(deletion_count);
}
let mut cursor = self.find_many(doc! { "filename": filename }, None).await?;
while let Some(FilesCollectionDocument { id, .. }) = cursor.next().await? {
self.delete_by_oid(convert_bson_to_oid(id)?).await?;
deletion_count += 1;
}
Ok(deletion_count)
}
pub async fn delete_by_oid(&self, oid: ObjectId) -> Result<(), Error> {
self.bucket
.delete(Bson::ObjectId(oid))
.await
.map_err(Error::Mongo)
}
}