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 5510b6cdc refactor(layers/logging): Don't trigger logigng in heavy IO
path (#6343)
5510b6cdc is described below
commit 5510b6cdc462afe79af4b36f7dfeff5ae18ace50
Author: Xuanwo <[email protected]>
AuthorDate: Mon Jun 30 14:51:00 2025 +0800
refactor(layers/logging): Don't trigger logigng in heavy IO path (#6343)
Signed-off-by: Xuanwo <[email protected]>
---
core/benches/vs_fs/src/main.rs | 2 +-
core/examples/basic/src/main.rs | 2 +-
core/src/layers/logging.rs | 121 ++++------------------------------------
3 files changed, 13 insertions(+), 112 deletions(-)
diff --git a/core/benches/vs_fs/src/main.rs b/core/benches/vs_fs/src/main.rs
index 78b8a7c80..1859899dd 100644
--- a/core/benches/vs_fs/src/main.rs
+++ b/core/benches/vs_fs/src/main.rs
@@ -76,7 +76,7 @@ fn prepare() -> String {
rng.fill_bytes(&mut content);
let name = uuid::Uuid::new_v4();
- let path = format!("/tmp/opendal/{}", name);
+ let path = format!("/tmp/opendal/{name}");
let _ = std::fs::write(path, content);
name.to_string()
diff --git a/core/examples/basic/src/main.rs b/core/examples/basic/src/main.rs
index b75a2aeed..c1b92dbdc 100644
--- a/core/examples/basic/src/main.rs
+++ b/core/examples/basic/src/main.rs
@@ -29,7 +29,7 @@ async fn example(op: Operator) -> Result<()> {
// Fetch metadata of s3.
let meta = op.stat("test.txt").await?;
- println!("stat: {:?}", meta);
+ println!("stat: {meta:?}");
// Delete data from s3.
op.delete("test.txt").await?;
diff --git a/core/src/layers/logging.rs b/core/src/layers/logging.rs
index da3dd814a..d4704c7e3 100644
--- a/core/src/layers/logging.rs
+++ b/core/src/layers/logging.rs
@@ -541,7 +541,6 @@ impl<A: Access, I: LoggingInterceptor> LayeredAccess for
LoggingAccessor<A, I> {
}
}
-/// `LoggingReader` is a wrapper of `BytesReader`, with logging functionality.
pub struct LoggingReader<R, I: LoggingInterceptor> {
info: Arc<AccessorInfo>,
logger: I,
@@ -566,17 +565,8 @@ impl<R, I: LoggingInterceptor> LoggingReader<R, I> {
impl<R: oio::Read, I: LoggingInterceptor> oio::Read for LoggingReader<R, I> {
async fn read(&mut self) -> Result<Buffer> {
- self.logger.log(
- &self.info,
- Operation::Read,
- &[("path", &self.path), ("read", &self.read.to_string())],
- "started",
- None,
- );
-
match self.inner.read().await {
- Ok(bs) => {
- self.read += bs.len() as u64;
+ Ok(bs) if bs.is_empty() => {
self.logger.log(
&self.info,
Operation::Read,
@@ -585,15 +575,15 @@ impl<R: oio::Read, I: LoggingInterceptor> oio::Read for
LoggingReader<R, I> {
("read", &self.read.to_string()),
("size", &bs.len().to_string()),
],
- if bs.is_empty() {
- "finished"
- } else {
- "succeeded"
- },
+ "finished",
None,
);
Ok(bs)
}
+ Ok(bs) => {
+ self.read += bs.len() as u64;
+ Ok(bs)
+ }
Err(err) => {
self.logger.log(
&self.info,
@@ -634,32 +624,9 @@ impl<W: oio::Write, I: LoggingInterceptor> oio::Write for
LoggingWriter<W, I> {
async fn write(&mut self, bs: Buffer) -> Result<()> {
let size = bs.len();
- self.logger.log(
- &self.info,
- Operation::Write,
- &[
- ("path", &self.path),
- ("written", &self.written.to_string()),
- ("size", &size.to_string()),
- ],
- "started",
- None,
- );
-
match self.inner.write(bs).await {
Ok(_) => {
self.written += size as u64;
- self.logger.log(
- &self.info,
- Operation::Write,
- &[
- ("path", &self.path),
- ("written", &self.written.to_string()),
- ("size", &size.to_string()),
- ],
- "succeeded",
- None,
- );
Ok(())
}
Err(err) => {
@@ -680,21 +647,13 @@ impl<W: oio::Write, I: LoggingInterceptor> oio::Write for
LoggingWriter<W, I> {
}
async fn abort(&mut self) -> Result<()> {
- self.logger.log(
- &self.info,
- Operation::Write,
- &[("path", &self.path), ("written", &self.written.to_string())],
- "started",
- None,
- );
-
match self.inner.abort().await {
Ok(_) => {
self.logger.log(
&self.info,
Operation::Write,
&[("path", &self.path), ("written",
&self.written.to_string())],
- "succeeded",
+ "abort succeeded",
None,
);
Ok(())
@@ -704,7 +663,7 @@ impl<W: oio::Write, I: LoggingInterceptor> oio::Write for
LoggingWriter<W, I> {
&self.info,
Operation::Write,
&[("path", &self.path), ("written",
&self.written.to_string())],
- "failed",
+ "abort failed",
Some(&err),
);
Err(err)
@@ -713,21 +672,13 @@ impl<W: oio::Write, I: LoggingInterceptor> oio::Write for
LoggingWriter<W, I> {
}
async fn close(&mut self) -> Result<Metadata> {
- self.logger.log(
- &self.info,
- Operation::Write,
- &[("path", &self.path), ("written", &self.written.to_string())],
- "started",
- None,
- );
-
match self.inner.close().await {
Ok(meta) => {
self.logger.log(
&self.info,
Operation::Write,
&[("path", &self.path), ("written",
&self.written.to_string())],
- "succeeded",
+ "close succeeded",
None,
);
Ok(meta)
@@ -737,7 +688,7 @@ impl<W: oio::Write, I: LoggingInterceptor> oio::Write for
LoggingWriter<W, I> {
&self.info,
Operation::Write,
&[("path", &self.path), ("written",
&self.written.to_string())],
- "failed",
+ "close failed",
Some(&err),
);
Err(err)
@@ -770,30 +721,11 @@ impl<P, I: LoggingInterceptor> LoggingLister<P, I> {
impl<P: oio::List, I: LoggingInterceptor> oio::List for LoggingLister<P, I> {
async fn next(&mut self) -> Result<Option<oio::Entry>> {
- self.logger.log(
- &self.info,
- Operation::List,
- &[("path", &self.path), ("listed", &self.listed.to_string())],
- "started",
- None,
- );
-
let res = self.inner.next().await;
match &res {
- Ok(Some(de)) => {
+ Ok(Some(_)) => {
self.listed += 1;
- self.logger.log(
- &self.info,
- Operation::List,
- &[
- ("path", &self.path),
- ("listed", &self.listed.to_string()),
- ("entry", de.path()),
- ],
- "succeeded",
- None,
- );
}
Ok(None) => {
self.logger.log(
@@ -848,31 +780,11 @@ impl<D: oio::Delete, I: LoggingInterceptor> oio::Delete
for LoggingDeleter<D, I>
.map(|v| v.to_string())
.unwrap_or_else(|| "<latest>".to_string());
- self.logger.log(
- &self.info,
- Operation::Delete,
- &[("path", path), ("version", &version)],
- "started",
- None,
- );
-
let res = self.inner.delete(path, args);
match &res {
Ok(_) => {
self.queued += 1;
- self.logger.log(
- &self.info,
- Operation::Delete,
- &[
- ("path", path),
- ("version", &version),
- ("queued", &self.queued.to_string()),
- ("deleted", &self.deleted.to_string()),
- ],
- "succeeded",
- None,
- );
}
Err(err) => {
self.logger.log(
@@ -894,17 +806,6 @@ impl<D: oio::Delete, I: LoggingInterceptor> oio::Delete
for LoggingDeleter<D, I>
}
async fn flush(&mut self) -> Result<usize> {
- self.logger.log(
- &self.info,
- Operation::Delete,
- &[
- ("queued", &self.queued.to_string()),
- ("deleted", &self.deleted.to_string()),
- ],
- "started",
- None,
- );
-
let res = self.inner.flush().await;
match &res {