This is an automated email from the ASF dual-hosted git repository.
xuanwo pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/opendal.git
The following commit(s) were added to refs/heads/main by this push:
new 275224dc70 feat(raw/oio): Use `Buffer` as cache in `OneshotWrite`
(#4477)
275224dc70 is described below
commit 275224dc70c87851f8da28aaf5a4517942ee612b
Author: Weijie Guo <[email protected]>
AuthorDate: Sat Apr 13 11:12:41 2024 +0800
feat(raw/oio): Use `Buffer` as cache in `OneshotWrite` (#4477)
* feat(raw/oio): Use `Buffer` as cache in `OneshotWrite`
* fix impl
* clippy
* fmt
---
core/src/raw/oio/write/one_shot_write.rs | 10 ++++------
core/src/services/azdls/writer.rs | 12 ++++--------
core/src/services/azfile/writer.rs | 5 ++---
core/src/services/chainsafe/core.rs | 2 +-
core/src/services/chainsafe/writer.rs | 3 +--
core/src/services/dbfs/writer.rs | 7 ++++---
core/src/services/dropbox/writer.rs | 5 ++---
core/src/services/gdrive/core.rs | 6 +++---
core/src/services/gdrive/writer.rs | 4 ++--
core/src/services/github/backend.rs | 2 +-
core/src/services/github/core.rs | 6 +++---
core/src/services/github/writer.rs | 3 +--
core/src/services/ipmfs/backend.rs | 3 +--
core/src/services/ipmfs/writer.rs | 3 +--
core/src/services/koofr/core.rs | 2 +-
core/src/services/koofr/writer.rs | 3 +--
core/src/services/onedrive/writer.rs | 8 ++++----
core/src/services/pcloud/core.rs | 7 ++-----
core/src/services/pcloud/writer.rs | 3 +--
core/src/services/seafile/writer.rs | 3 +--
core/src/services/supabase/writer.rs | 5 ++---
core/src/services/swift/writer.rs | 5 ++---
core/src/services/vercel_artifacts/writer.rs | 5 ++---
core/src/services/webdav/writer.rs | 10 ++--------
core/src/services/yandex_disk/writer.rs | 5 ++---
25 files changed, 50 insertions(+), 77 deletions(-)
diff --git a/core/src/raw/oio/write/one_shot_write.rs
b/core/src/raw/oio/write/one_shot_write.rs
index 0cd6b2a97a..09e9b4f8f0 100644
--- a/core/src/raw/oio/write/one_shot_write.rs
+++ b/core/src/raw/oio/write/one_shot_write.rs
@@ -17,8 +17,6 @@
use std::future::Future;
-use bytes::Bytes;
-
use crate::raw::*;
use crate::*;
@@ -32,13 +30,13 @@ pub trait OneShotWrite: Send + Sync + Unpin + 'static {
/// write_once write all data at once.
///
/// Implementations should make sure that the data is written correctly at
once.
- fn write_once(&self, bs: Bytes) -> impl Future<Output = Result<()>> +
MaybeSend;
+ fn write_once(&self, bs: Buffer) -> impl Future<Output = Result<()>> +
MaybeSend;
}
/// OneShotWrite is used to implement [`Write`] based on one shot.
pub struct OneShotWriter<W: OneShotWrite> {
inner: W,
- buffer: Option<Bytes>,
+ buffer: Option<Buffer>,
}
impl<W: OneShotWrite> OneShotWriter<W> {
@@ -60,7 +58,7 @@ impl<W: OneShotWrite> oio::Write for OneShotWriter<W> {
)),
None => {
let size = bs.len();
- self.buffer = Some(bs.to_bytes());
+ self.buffer = Some(bs);
Ok(size)
}
}
@@ -69,7 +67,7 @@ impl<W: OneShotWrite> oio::Write for OneShotWriter<W> {
async fn close(&mut self) -> Result<()> {
match self.buffer.clone() {
Some(bs) => self.inner.write_once(bs).await,
- None => self.inner.write_once(Bytes::new()).await,
+ None => self.inner.write_once(Buffer::new()).await,
}
}
diff --git a/core/src/services/azdls/writer.rs
b/core/src/services/azdls/writer.rs
index 2f1bb012e8..2738cb9e88 100644
--- a/core/src/services/azdls/writer.rs
+++ b/core/src/services/azdls/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::AzdlsCore;
@@ -41,7 +40,7 @@ impl AzdlsWriter {
}
impl oio::OneShotWrite for AzdlsWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let mut req =
self.core
.azdls_create_request(&self.path, "file", &self.op,
Buffer::new())?;
@@ -60,12 +59,9 @@ impl oio::OneShotWrite for AzdlsWriter {
}
}
- let mut req = self.core.azdls_update_request(
- &self.path,
- Some(bs.len() as u64),
- 0,
- Buffer::from(bs),
- )?;
+ let mut req = self
+ .core
+ .azdls_update_request(&self.path, Some(bs.len() as u64), 0, bs)?;
self.core.sign(&mut req).await?;
diff --git a/core/src/services/azfile/writer.rs
b/core/src/services/azfile/writer.rs
index 52cf6d5ac6..b111aae312 100644
--- a/core/src/services/azfile/writer.rs
+++ b/core/src/services/azfile/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::AzfileCore;
@@ -40,7 +39,7 @@ impl AzfileWriter {
}
impl oio::OneShotWrite for AzfileWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let resp = self
.core
.azfile_create_file(&self.path, bs.len(), &self.op)
@@ -58,7 +57,7 @@ impl oio::OneShotWrite for AzfileWriter {
let resp = self
.core
- .azfile_update(&self.path, bs.len() as u64, 0, Buffer::from(bs))
+ .azfile_update(&self.path, bs.len() as u64, 0, bs)
.await?;
let status = resp.status();
match status {
diff --git a/core/src/services/chainsafe/core.rs
b/core/src/services/chainsafe/core.rs
index 11554a6c67..e864c41cb2 100644
--- a/core/src/services/chainsafe/core.rs
+++ b/core/src/services/chainsafe/core.rs
@@ -163,7 +163,7 @@ impl ChainsafeCore {
self.send(req).await
}
- pub async fn upload_object(&self, path: &str, bs: Bytes) ->
Result<Response<Buffer>> {
+ pub async fn upload_object(&self, path: &str, bs: Buffer) ->
Result<Response<Buffer>> {
let path = build_abs_path(&self.root, path);
let url = format!(
diff --git a/core/src/services/chainsafe/writer.rs
b/core/src/services/chainsafe/writer.rs
index bf9547d5b1..976f0c29de 100644
--- a/core/src/services/chainsafe/writer.rs
+++ b/core/src/services/chainsafe/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::ChainsafeCore;
@@ -44,7 +43,7 @@ impl ChainsafeWriter {
}
impl oio::OneShotWrite for ChainsafeWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let resp = self.core.upload_object(&self.path, bs).await?;
let status = resp.status();
diff --git a/core/src/services/dbfs/writer.rs b/core/src/services/dbfs/writer.rs
index 390f973269..e1e261c75e 100644
--- a/core/src/services/dbfs/writer.rs
+++ b/core/src/services/dbfs/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::error::parse_error;
@@ -39,7 +38,7 @@ impl DbfsWriter {
}
impl oio::OneShotWrite for DbfsWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let size = bs.len();
// MAX_BLOCK_SIZE_EXCEEDED will be thrown if this limit(1MB) is
exceeded.
@@ -50,7 +49,9 @@ impl oio::OneShotWrite for DbfsWriter {
));
}
- let req = self.core.dbfs_create_file_request(&self.path, bs)?;
+ let req = self
+ .core
+ .dbfs_create_file_request(&self.path, bs.to_bytes())?;
let resp = self.core.client.send(req).await?;
diff --git a/core/src/services/dropbox/writer.rs
b/core/src/services/dropbox/writer.rs
index 7c9aaeef50..3adbf063f2 100644
--- a/core/src/services/dropbox/writer.rs
+++ b/core/src/services/dropbox/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::DropboxCore;
@@ -38,10 +37,10 @@ impl DropboxWriter {
}
impl oio::OneShotWrite for DropboxWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let resp = self
.core
- .dropbox_update(&self.path, Some(bs.len()), &self.op,
Buffer::from(bs))
+ .dropbox_update(&self.path, Some(bs.len()), &self.op, bs)
.await?;
let status = resp.status();
match status {
diff --git a/core/src/services/gdrive/core.rs b/core/src/services/gdrive/core.rs
index 9130f8cabc..6e3840872e 100644
--- a/core/src/services/gdrive/core.rs
+++ b/core/src/services/gdrive/core.rs
@@ -183,7 +183,7 @@ impl GdriveCore {
&self,
path: &str,
size: u64,
- body: Bytes,
+ body: Buffer,
) -> Result<Response<Buffer>> {
let parent = self.path_cache.ensure_dir(get_parent(path)).await?;
@@ -233,7 +233,7 @@ impl GdriveCore {
&self,
file_id: &str,
size: u64,
- body: Bytes,
+ body: Buffer,
) -> Result<Response<Buffer>> {
let url = format!(
"https://www.googleapis.com/upload/drive/v3/files/{}?uploadType=media",
@@ -244,7 +244,7 @@ impl GdriveCore {
.header(header::CONTENT_TYPE, "application/octet-stream")
.header(header::CONTENT_LENGTH, size)
.header("X-Upload-Content-Length", size)
- .body(Buffer::from(body))
+ .body(body)
.map_err(new_request_build_error)?;
self.sign(&mut req).await?;
diff --git a/core/src/services/gdrive/writer.rs
b/core/src/services/gdrive/writer.rs
index ebd91afbdd..46b6f85e3b 100644
--- a/core/src/services/gdrive/writer.rs
+++ b/core/src/services/gdrive/writer.rs
@@ -17,7 +17,7 @@
use std::sync::Arc;
-use bytes::{Buf, Bytes};
+use bytes::Buf;
use http::StatusCode;
use super::core::GdriveCore;
@@ -46,7 +46,7 @@ impl GdriveWriter {
}
impl oio::OneShotWrite for GdriveWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let size = bs.len();
let resp = if let Some(file_id) = &self.file_id {
diff --git a/core/src/services/github/backend.rs
b/core/src/services/github/backend.rs
index 71e4815f84..29182ea419 100644
--- a/core/src/services/github/backend.rs
+++ b/core/src/services/github/backend.rs
@@ -254,7 +254,7 @@ impl Accessor for GithubBackend {
}
async fn create_dir(&self, path: &str, _: OpCreateDir) ->
Result<RpCreateDir> {
- let empty_bytes = bytes::Bytes::new();
+ let empty_bytes = Buffer::new();
let resp = self
.core
diff --git a/core/src/services/github/core.rs b/core/src/services/github/core.rs
index 05b1213b9b..35baad60bb 100644
--- a/core/src/services/github/core.rs
+++ b/core/src/services/github/core.rs
@@ -158,7 +158,7 @@ impl GithubCore {
self.send(req).await
}
- pub async fn upload(&self, path: &str, bs: Bytes) ->
Result<Response<Buffer>> {
+ pub async fn upload(&self, path: &str, bs: Buffer) ->
Result<Response<Buffer>> {
let sha = self.get_file_sha(path).await?;
let path = build_abs_path(&self.root, path);
@@ -176,7 +176,7 @@ impl GithubCore {
let mut req_body = CreateOrUpdateContentsRequest {
message: format!("Write {} at {} via opendal", path,
chrono::Local::now()),
- content: base64::engine::general_purpose::STANDARD.encode(&bs),
+ content:
base64::engine::general_purpose::STANDARD.encode(&bs.to_bytes()),
sha: None,
};
@@ -188,7 +188,7 @@ impl GithubCore {
let req = req
.header("Accept", "application/vnd.github+json")
- .body(Buffer::from(Bytes::from(req_body)))
+ .body(Buffer::from(req_body))
.map_err(new_request_build_error)?;
self.send(req).await
diff --git a/core/src/services/github/writer.rs
b/core/src/services/github/writer.rs
index 5b6d3f6c1f..095da01572 100644
--- a/core/src/services/github/writer.rs
+++ b/core/src/services/github/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::GithubCore;
@@ -39,7 +38,7 @@ impl GithubWriter {
}
impl oio::OneShotWrite for GithubWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let resp = self.core.upload(&self.path, bs).await?;
let status = resp.status();
diff --git a/core/src/services/ipmfs/backend.rs
b/core/src/services/ipmfs/backend.rs
index 52062a2db4..4862c41585 100644
--- a/core/src/services/ipmfs/backend.rs
+++ b/core/src/services/ipmfs/backend.rs
@@ -22,7 +22,6 @@ use std::sync::Arc;
use async_trait::async_trait;
use bytes::Buf;
-use bytes::Bytes;
use http::Request;
use http::Response;
use http::StatusCode;
@@ -248,7 +247,7 @@ impl IpmfsBackend {
}
/// Support write from reader.
- pub async fn ipmfs_write(&self, path: &str, body: Bytes) ->
Result<Response<Buffer>> {
+ pub async fn ipmfs_write(&self, path: &str, body: Buffer) ->
Result<Response<Buffer>> {
let p = build_rooted_abs_path(&self.root, path);
let url = format!(
diff --git a/core/src/services/ipmfs/writer.rs
b/core/src/services/ipmfs/writer.rs
index f27fe5c77b..ec0811bf83 100644
--- a/core/src/services/ipmfs/writer.rs
+++ b/core/src/services/ipmfs/writer.rs
@@ -15,7 +15,6 @@
// specific language governing permissions and limitations
// under the License.
-use bytes::Bytes;
use http::StatusCode;
use super::backend::IpmfsBackend;
@@ -36,7 +35,7 @@ impl IpmfsWriter {
}
impl oio::OneShotWrite for IpmfsWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let resp = self.backend.ipmfs_write(&self.path, bs).await?;
let status = resp.status();
diff --git a/core/src/services/koofr/core.rs b/core/src/services/koofr/core.rs
index 0b92f60bef..892baf0d68 100644
--- a/core/src/services/koofr/core.rs
+++ b/core/src/services/koofr/core.rs
@@ -259,7 +259,7 @@ impl KoofrCore {
self.send(req).await
}
- pub async fn put(&self, path: &str, bs: Bytes) -> Result<Response<Buffer>>
{
+ pub async fn put(&self, path: &str, bs: Buffer) ->
Result<Response<Buffer>> {
let path = build_rooted_abs_path(&self.root, path);
let filename = get_basename(&path);
diff --git a/core/src/services/koofr/writer.rs
b/core/src/services/koofr/writer.rs
index ff8039e8af..ca1cccbb39 100644
--- a/core/src/services/koofr/writer.rs
+++ b/core/src/services/koofr/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::KoofrCore;
@@ -39,7 +38,7 @@ impl KoofrWriter {
}
impl oio::OneShotWrite for KoofrWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
self.core.ensure_dir_exists(&self.path).await?;
let resp = self.core.put(&self.path, bs).await?;
diff --git a/core/src/services/onedrive/writer.rs
b/core/src/services/onedrive/writer.rs
index f62bd3bd6b..3db9121ae6 100644
--- a/core/src/services/onedrive/writer.rs
+++ b/core/src/services/onedrive/writer.rs
@@ -44,13 +44,13 @@ impl OneDriveWriter {
}
impl oio::OneShotWrite for OneDriveWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let size = bs.len();
if size <= Self::MAX_SIMPLE_SIZE {
self.write_simple(bs).await?;
} else {
- self.write_chunked(bs).await?;
+ self.write_chunked(bs.to_bytes()).await?;
}
Ok(())
@@ -58,10 +58,10 @@ impl oio::OneShotWrite for OneDriveWriter {
}
impl OneDriveWriter {
- async fn write_simple(&self, bs: Bytes) -> Result<()> {
+ async fn write_simple(&self, bs: Buffer) -> Result<()> {
let resp = self
.backend
- .onedrive_upload_simple(&self.path, Some(bs.len()), &self.op,
Buffer::from(bs))
+ .onedrive_upload_simple(&self.path, Some(bs.len()), &self.op, bs)
.await?;
let status = resp.status();
diff --git a/core/src/services/pcloud/core.rs b/core/src/services/pcloud/core.rs
index 3a27ca108b..4a146d3697 100644
--- a/core/src/services/pcloud/core.rs
+++ b/core/src/services/pcloud/core.rs
@@ -19,7 +19,6 @@ use std::fmt::Debug;
use std::fmt::Formatter;
use bytes::Buf;
-use bytes::Bytes;
use http::header;
use http::Request;
use http::Response;
@@ -313,7 +312,7 @@ impl PcloudCore {
self.send(req).await
}
- pub async fn upload_file(&self, path: &str, bs: Bytes) ->
Result<Response<Buffer>> {
+ pub async fn upload_file(&self, path: &str, bs: Buffer) ->
Result<Response<Buffer>> {
let path = build_abs_path(&self.root, path);
let (name, path) = (get_basename(&path),
get_parent(&path).trim_end_matches('/'));
@@ -330,9 +329,7 @@ impl PcloudCore {
let req = Request::put(url);
// set body
- let req = req
- .body(Buffer::from(bs))
- .map_err(new_request_build_error)?;
+ let req = req.body(bs).map_err(new_request_build_error)?;
self.send(req).await
}
diff --git a/core/src/services/pcloud/writer.rs
b/core/src/services/pcloud/writer.rs
index d4d3bd53ca..1c185dec1c 100644
--- a/core/src/services/pcloud/writer.rs
+++ b/core/src/services/pcloud/writer.rs
@@ -18,7 +18,6 @@
use std::sync::Arc;
use bytes::Buf;
-use bytes::Bytes;
use http::StatusCode;
use super::core::PcloudCore;
@@ -41,7 +40,7 @@ impl PcloudWriter {
}
impl oio::OneShotWrite for PcloudWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
self.core.ensure_dir_exists(&self.path).await?;
let resp = self.core.upload_file(&self.path, bs).await?;
diff --git a/core/src/services/seafile/writer.rs
b/core/src/services/seafile/writer.rs
index ad7cc98b2e..718d06cf5c 100644
--- a/core/src/services/seafile/writer.rs
+++ b/core/src/services/seafile/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::header;
use http::Request;
use http::StatusCode;
@@ -46,7 +45,7 @@ impl SeafileWriter {
}
impl oio::OneShotWrite for SeafileWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let upload_url = self.core.get_upload_url().await?;
let req = Request::post(upload_url);
diff --git a/core/src/services/supabase/writer.rs
b/core/src/services/supabase/writer.rs
index 90a9b14336..4c7feae7c0 100644
--- a/core/src/services/supabase/writer.rs
+++ b/core/src/services/supabase/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::*;
@@ -43,12 +42,12 @@ impl SupabaseWriter {
}
impl oio::OneShotWrite for SupabaseWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let mut req = self.core.supabase_upload_object_request(
&self.path,
Some(bs.len()),
self.op.content_type(),
- Buffer::from(bs),
+ bs,
)?;
self.core.sign(&mut req)?;
diff --git a/core/src/services/swift/writer.rs
b/core/src/services/swift/writer.rs
index 3a097e4376..e2b51dab18 100644
--- a/core/src/services/swift/writer.rs
+++ b/core/src/services/swift/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::SwiftCore;
@@ -37,10 +36,10 @@ impl SwiftWriter {
}
impl oio::OneShotWrite for SwiftWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let resp = self
.core
- .swift_create_object(&self.path, bs.len() as u64, Buffer::from(bs))
+ .swift_create_object(&self.path, bs.len() as u64, bs)
.await?;
let status = resp.status();
diff --git a/core/src/services/vercel_artifacts/writer.rs
b/core/src/services/vercel_artifacts/writer.rs
index 0e52788766..9ef4a08e1c 100644
--- a/core/src/services/vercel_artifacts/writer.rs
+++ b/core/src/services/vercel_artifacts/writer.rs
@@ -15,7 +15,6 @@
// specific language governing permissions and limitations
// under the License.
-use bytes::Bytes;
use http::StatusCode;
use super::backend::VercelArtifactsBackend;
@@ -41,10 +40,10 @@ impl VercelArtifactsWriter {
}
impl oio::OneShotWrite for VercelArtifactsWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let resp = self
.backend
- .vercel_artifacts_put(self.path.as_str(), bs.len() as u64,
Buffer::from(bs))
+ .vercel_artifacts_put(self.path.as_str(), bs.len() as u64, bs)
.await?;
let status = resp.status();
diff --git a/core/src/services/webdav/writer.rs
b/core/src/services/webdav/writer.rs
index 28f744530c..7e1951c311 100644
--- a/core/src/services/webdav/writer.rs
+++ b/core/src/services/webdav/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::StatusCode;
use super::core::*;
@@ -39,15 +38,10 @@ impl WebdavWriter {
}
impl oio::OneShotWrite for WebdavWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
let resp = self
.core
- .webdav_put(
- &self.path,
- Some(bs.len() as u64),
- &self.op,
- Buffer::from(bs),
- )
+ .webdav_put(&self.path, Some(bs.len() as u64), &self.op, bs)
.await?;
let status = resp.status();
diff --git a/core/src/services/yandex_disk/writer.rs
b/core/src/services/yandex_disk/writer.rs
index 31b2371d2f..58cefe5ae1 100644
--- a/core/src/services/yandex_disk/writer.rs
+++ b/core/src/services/yandex_disk/writer.rs
@@ -17,7 +17,6 @@
use std::sync::Arc;
-use bytes::Bytes;
use http::Request;
use http::StatusCode;
@@ -40,13 +39,13 @@ impl YandexDiskWriter {
}
impl oio::OneShotWrite for YandexDiskWriter {
- async fn write_once(&self, bs: Bytes) -> Result<()> {
+ async fn write_once(&self, bs: Buffer) -> Result<()> {
self.core.ensure_dir_exists(&self.path).await?;
let upload_url = self.core.get_upload_url(&self.path).await?;
let req = Request::put(upload_url)
- .body(Buffer::from(bs))
+ .body(bs)
.map_err(new_request_build_error)?;
let resp = self.core.send(req).await?;