This is an automated email from the ASF dual-hosted git repository. xuanwo pushed a commit to branch reduce-log-rate in repository https://gitbox.apache.org/repos/asf/opendal.git
commit 8b62f20bb1e200989b523dbb50bcab4eed530476 Author: Xuanwo <[email protected]> AuthorDate: Sun Jun 29 13:13:42 2025 +0800 refactor(layers/logging): Don't trigger logigng in heavy IO path 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 {
