Skip to main content

opendal_layer_logging/
lib.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18#![doc = include_str!("../README.md")]
19#![cfg_attr(docsrs, feature(doc_cfg))]
20#![cfg_attr(docsrs, doc(auto_cfg))]
21#![deny(missing_docs)]
22use std::fmt::Debug;
23use std::fmt::Display;
24use std::sync::Arc;
25
26use log::Level;
27use log::log;
28use opendal_core::raw::*;
29use opendal_core::*;
30
31static LOGGING_TARGET: &str = "opendal::services";
32
33/// `LoggingLayer` records every operation with
34/// [log](https://docs.rs/log/).
35///
36/// # Logging
37///
38/// - OpenDAL emits structured logs.
39/// - Every operation starts with a `started` log entry.
40/// - Every operation ends with one of the following statuses:
41///   - `succeeded`: The operation succeeded but might have more work to perform.
42///   - `finished`: The whole operation finished.
43///   - `failed`: The operation returned an error.
44/// - The default log level for expected errors is `Warn`.
45/// - The default log level for unexpected errors is `Error`.
46///
47/// # Examples
48///
49/// ```no_run
50/// # use opendal_core::services;
51/// # use opendal_core::Operator;
52/// # use opendal_core::Result;
53/// # use opendal_layer_logging::LoggingLayer;
54/// #
55/// # fn main() -> Result<()> {
56/// let _ = Operator::new(services::Memory::default())?
57///     .layer(LoggingLayer::default());
58/// # Ok(())
59/// # }
60/// ```
61///
62/// # Output
63///
64/// OpenDAL uses [`log`](https://docs.rs/log/latest/log/) internally.
65///
66/// To enable logging output, please set `RUST_LOG`:
67///
68/// ```shell
69/// RUST_LOG=debug ./app
70/// ```
71///
72/// To config logging output, please refer to [Configure Logging](https://rust-lang-nursery.github.io/rust-cookbook/development_tools/debugging/config_log.html):
73///
74/// ```shell
75/// RUST_LOG="info,opendal::services=debug" ./app
76/// ```
77///
78/// # Logging Interceptor
79///
80/// You can implement your own logging interceptor to customize the logging behavior.
81///
82/// ```no_run
83/// # use opendal_core::raw;
84/// # use opendal_core::services;
85/// # use opendal_core::Error;
86/// # use opendal_core::Operator;
87/// # use opendal_core::Result;
88/// # use opendal_layer_logging::LoggingInterceptor;
89/// # use opendal_layer_logging::LoggingLayer;
90/// #
91/// #[derive(Debug, Clone)]
92/// struct MyLoggingInterceptor;
93///
94/// impl LoggingInterceptor for MyLoggingInterceptor {
95///     fn log(
96///         &self,
97///         info: &raw::ServiceInfo,
98///         operation: raw::Operation,
99///         context: &[(&str, &str)],
100///         message: &str,
101///         err: Option<&Error>,
102///     ) {
103///         // log something
104///     }
105/// }
106///
107/// # fn main() -> Result<()> {
108/// let _ = Operator::new(services::Memory::default())?
109///     .layer(LoggingLayer::new(MyLoggingInterceptor));
110/// # Ok(())
111/// # }
112/// ```
113#[derive(Clone, Copy, Debug)]
114pub struct LoggingLayer<I = DefaultLoggingInterceptor> {
115    logger: I,
116}
117
118impl Default for LoggingLayer {
119    fn default() -> Self {
120        Self {
121            logger: DefaultLoggingInterceptor,
122        }
123    }
124}
125
126impl LoggingLayer {
127    /// Create the layer with specific logging interceptor.
128    pub fn new<I: LoggingInterceptor>(logger: I) -> LoggingLayer<I> {
129        LoggingLayer { logger }
130    }
131}
132
133impl<I: LoggingInterceptor> Layer for LoggingLayer<I> {
134    fn apply_service(&self, inner: Servicer) -> Servicer {
135        Arc::new(self.layer(inner))
136    }
137}
138
139impl<I: LoggingInterceptor> LoggingLayer<I> {
140    fn layer(&self, inner: Servicer) -> LoggingService<I> {
141        let info = inner.info();
142        LoggingService {
143            inner,
144            info,
145            logger: self.logger.clone(),
146        }
147    }
148}
149
150/// LoggingInterceptor customizes log emission.
151pub trait LoggingInterceptor: Debug + Clone + Send + Sync + Unpin + 'static {
152    /// Called for every log event.
153    ///
154    /// # Inputs
155    ///
156    /// - `info`: The service information used for this operation.
157    /// - `operation`: The operation being logged.
158    /// - `context`: Additional key-value context such as path, range, or counters.
159    /// - `message`: The event message, such as `started`, `finished`, or `failed`.
160    /// - `err`: The error associated with this event, if any.
161    ///
162    /// # Performance
163    ///
164    /// This method runs inline with the operation path. Avoid expensive I/O,
165    /// network calls, or long-running work here.
166    fn log(
167        &self,
168        info: &ServiceInfo,
169        operation: Operation,
170        context: &[(&str, &str)],
171        message: &str,
172        err: Option<&Error>,
173    );
174}
175
176/// The DefaultLoggingInterceptor will log the message by the standard logging macro.
177#[derive(Clone, Copy, Debug, Default)]
178pub struct DefaultLoggingInterceptor;
179
180impl LoggingInterceptor for DefaultLoggingInterceptor {
181    fn log(
182        &self,
183        info: &ServiceInfo,
184        operation: Operation,
185        context: &[(&str, &str)],
186        message: &str,
187        err: Option<&Error>,
188    ) {
189        if let Some(err) = err {
190            // Expected errors are logged as warnings; unexpected errors need stronger
191            // visibility and more diagnostic context.
192            let lvl = if err.kind() == ErrorKind::Unexpected {
193                Level::Error
194            } else {
195                Level::Warn
196            };
197
198            log!(
199                target: LOGGING_TARGET,
200                lvl,
201                "service={} name={}{}: {operation} {message} {}",
202                info.scheme(),
203                info.name(),
204                LoggingContext(context),
205                // Use Debug for unexpected errors to preserve more context.
206                // String avoids conditional format_args! temporaries in this log! argument.
207                if err.kind() != ErrorKind::Unexpected {
208                    format!("{err}")
209                } else {
210                    format!("{err:?}")
211                }
212            );
213        }
214
215        log!(
216            target: LOGGING_TARGET,
217            Level::Debug,
218            "service={} name={}{}: {operation} {message}",
219            info.scheme(),
220            info.name(),
221            LoggingContext(context),
222        );
223    }
224}
225
226struct LoggingContext<'a>(&'a [(&'a str, &'a str)]);
227
228impl Display for LoggingContext<'_> {
229    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
230        for (k, v) in self.0.iter() {
231            write!(f, " {k}={v}")?;
232        }
233        Ok(())
234    }
235}
236
237#[doc(hidden)]
238pub struct LoggingService<I: LoggingInterceptor> {
239    inner: Servicer,
240    info: ServiceInfo,
241    logger: I,
242}
243
244impl<I: LoggingInterceptor> Debug for LoggingService<I> {
245    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
246        f.debug_struct("LoggingService")
247            .field("inner", &self.inner)
248            .field("info", &self.info)
249            .finish_non_exhaustive()
250    }
251}
252
253impl<I: LoggingInterceptor> LoggingService<I> {
254    fn log_start(&self, op: Operation, context: &[(&str, &str)]) {
255        self.logger.log(&self.info, op, context, "started", None);
256    }
257
258    fn log_finish(&self, op: Operation, context: &[(&str, &str)], err: Option<&Error>) {
259        let message = if err.is_some() { "failed" } else { "finished" };
260        self.logger.log(&self.info, op, context, message, err);
261    }
262}
263
264impl<I: LoggingInterceptor> Service for LoggingService<I> {
265    type Reader = LoggingReader<oio::Reader, I>;
266    type Writer = LoggingWriter<oio::Writer, I>;
267    type Lister = LoggingLister<oio::Lister, I>;
268    type Deleter = LoggingDeleter<oio::Deleter, I>;
269    type Copier = LoggingCopier<oio::Copier, I>;
270
271    fn info(&self) -> ServiceInfo {
272        self.info.clone()
273    }
274
275    fn capability(&self) -> Capability {
276        self.inner.capability()
277    }
278
279    async fn create_dir(
280        &self,
281        ctx: &OperationContext,
282        path: &str,
283        args: OpCreateDir,
284    ) -> Result<RpCreateDir> {
285        self.log_start(Operation::CreateDir, &[("path", path)]);
286        let result = self.inner.create_dir(ctx, path, args).await;
287        self.log_finish(
288            Operation::CreateDir,
289            &[("path", path)],
290            result.as_ref().err(),
291        );
292        result
293    }
294
295    fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
296        self.log_start(Operation::Read, &[("path", path)]);
297        self.inner
298            .read(ctx, path, args)
299            .map(|r| {
300                self.logger.log(
301                    &self.info,
302                    Operation::Read,
303                    &[("path", path)],
304                    "created reader",
305                    None,
306                );
307                LoggingReader::new(self.info.clone(), self.logger.clone(), path, r)
308            })
309            .inspect_err(|err| {
310                self.logger.log(
311                    &self.info,
312                    Operation::Read,
313                    &[("path", path)],
314                    "failed",
315                    Some(err),
316                );
317            })
318    }
319
320    fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
321        self.log_start(Operation::Write, &[("path", path)]);
322        self.inner
323            .write(ctx, path, args)
324            .map(|w| {
325                self.logger.log(
326                    &self.info,
327                    Operation::Write,
328                    &[("path", path)],
329                    "created writer",
330                    None,
331                );
332                LoggingWriter::new(self.info.clone(), self.logger.clone(), path, w)
333            })
334            .inspect_err(|err| {
335                self.logger.log(
336                    &self.info,
337                    Operation::Write,
338                    &[("path", path)],
339                    "failed",
340                    Some(err),
341                );
342            })
343    }
344
345    async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
346        self.log_start(Operation::Stat, &[("path", path)]);
347        let result = self.inner.stat(ctx, path, args).await;
348        self.log_finish(Operation::Stat, &[("path", path)], result.as_ref().err());
349        result
350    }
351
352    fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
353        self.log_start(Operation::Delete, &[]);
354        self.inner
355            .delete(ctx)
356            .map(|d| {
357                self.logger
358                    .log(&self.info, Operation::Delete, &[], "finished", None);
359                LoggingDeleter::new(self.info.clone(), self.logger.clone(), d)
360            })
361            .inspect_err(|err| {
362                self.logger
363                    .log(&self.info, Operation::Delete, &[], "failed", Some(err));
364            })
365    }
366
367    fn copy(
368        &self,
369        ctx: &OperationContext,
370        from: &str,
371        to: &str,
372        args: OpCopy,
373    ) -> Result<Self::Copier> {
374        self.log_start(Operation::Copy, &[("from", from), ("to", to)]);
375        self.inner
376            .copy(ctx, from, to, args)
377            .map(|c| {
378                self.logger.log(
379                    &self.info,
380                    Operation::Copy,
381                    &[("from", from), ("to", to)],
382                    "created copier",
383                    None,
384                );
385                LoggingCopier::new(self.info.clone(), self.logger.clone(), from, to, c)
386            })
387            .inspect_err(|err| {
388                self.logger.log(
389                    &self.info,
390                    Operation::Copy,
391                    &[("from", from), ("to", to)],
392                    "failed",
393                    Some(err),
394                );
395            })
396    }
397
398    async fn rename(
399        &self,
400        ctx: &OperationContext,
401        from: &str,
402        to: &str,
403        args: OpRename,
404    ) -> Result<RpRename> {
405        self.log_start(Operation::Rename, &[("from", from), ("to", to)]);
406        let result = self.inner.rename(ctx, from, to, args).await;
407        self.log_finish(
408            Operation::Rename,
409            &[("from", from), ("to", to)],
410            result.as_ref().err(),
411        );
412        result
413    }
414
415    async fn restore(
416        &self,
417        ctx: &OperationContext,
418        path: &str,
419        args: OpRestore,
420    ) -> Result<RpRestore> {
421        self.log_start(Operation::Restore, &[("path", path)]);
422        let result = self.inner.restore(ctx, path, args).await;
423        self.log_finish(Operation::Restore, &[("path", path)], result.as_ref().err());
424        result
425    }
426
427    fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
428        self.log_start(Operation::List, &[("path", path)]);
429        self.inner
430            .list(ctx, path, args)
431            .map(|v| {
432                self.logger.log(
433                    &self.info,
434                    Operation::List,
435                    &[("path", path)],
436                    "created lister",
437                    None,
438                );
439                LoggingLister::new(self.info.clone(), self.logger.clone(), path, v)
440            })
441            .inspect_err(|err| {
442                self.logger.log(
443                    &self.info,
444                    Operation::List,
445                    &[("path", path)],
446                    "failed",
447                    Some(err),
448                );
449            })
450    }
451
452    async fn presign(
453        &self,
454        ctx: &OperationContext,
455        path: &str,
456        args: OpPresign,
457    ) -> Result<RpPresign> {
458        self.log_start(Operation::Presign, &[("path", path)]);
459        let result = self.inner.presign(ctx, path, args).await;
460        self.log_finish(Operation::Presign, &[("path", path)], result.as_ref().err());
461        result
462    }
463}
464
465#[doc(hidden)]
466pub struct LoggingReader<R, I: LoggingInterceptor> {
467    info: ServiceInfo,
468    logger: I,
469    path: String,
470    range: Option<BytesRange>,
471
472    read: u64,
473    inner: R,
474}
475
476impl<R, I: LoggingInterceptor> LoggingReader<R, I> {
477    fn new(info: ServiceInfo, logger: I, path: &str, reader: R) -> Self {
478        Self::with_range(info, logger, path, None, reader)
479    }
480
481    fn with_range(
482        info: ServiceInfo,
483        logger: I,
484        path: &str,
485        range: Option<BytesRange>,
486        reader: R,
487    ) -> Self {
488        Self {
489            info,
490            logger,
491            path: path.to_string(),
492            range,
493
494            read: 0,
495            inner: reader,
496        }
497    }
498
499    fn range_label(&self) -> String {
500        self.range
501            .map(|range| range.to_string())
502            .unwrap_or_default()
503    }
504}
505
506impl<R: oio::ReadStream, I: LoggingInterceptor> oio::ReadStream for LoggingReader<R, I> {
507    async fn read(&mut self) -> Result<Buffer> {
508        match self.inner.read().await {
509            Ok(bs) if bs.is_empty() => {
510                let range = self.range_label();
511                self.logger.log(
512                    &self.info,
513                    Operation::Read,
514                    &[
515                        ("path", &self.path),
516                        ("range", &range),
517                        ("read", &self.read.to_string()),
518                        ("size", &bs.len().to_string()),
519                    ],
520                    "finished",
521                    None,
522                );
523                Ok(bs)
524            }
525            Ok(bs) => {
526                self.read += bs.len() as u64;
527                Ok(bs)
528            }
529            Err(err) => {
530                let range = self.range_label();
531                self.logger.log(
532                    &self.info,
533                    Operation::Read,
534                    &[
535                        ("path", &self.path),
536                        ("range", &range),
537                        ("read", &self.read.to_string()),
538                    ],
539                    "failed",
540                    Some(&err),
541                );
542                Err(err)
543            }
544        }
545    }
546}
547
548impl<R: oio::Read, I: LoggingInterceptor> oio::Read for LoggingReader<R, I> {
549    async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
550        match self.inner.open(range).await {
551            Ok((rp, stream)) => Ok((
552                rp,
553                Box::new(LoggingReader::with_range(
554                    self.info.clone(),
555                    self.logger.clone(),
556                    &self.path,
557                    Some(range),
558                    stream,
559                )) as Box<dyn oio::ReadStreamDyn>,
560            )),
561            Err(err) => {
562                self.logger.log(
563                    &self.info,
564                    Operation::Read,
565                    &[("path", &self.path), ("range", &range.to_string())],
566                    "failed",
567                    Some(&err),
568                );
569                Err(err)
570            }
571        }
572    }
573
574    async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
575        match self.inner.read(range).await {
576            Ok((rp, buffer)) => {
577                self.logger.log(
578                    &self.info,
579                    Operation::Read,
580                    &[
581                        ("path", &self.path),
582                        ("range", &range.to_string()),
583                        ("size", &buffer.len().to_string()),
584                    ],
585                    "finished",
586                    None,
587                );
588                Ok((rp, buffer))
589            }
590            Err(err) => {
591                self.logger.log(
592                    &self.info,
593                    Operation::Read,
594                    &[("path", &self.path), ("range", &range.to_string())],
595                    "failed",
596                    Some(&err),
597                );
598                Err(err)
599            }
600        }
601    }
602}
603
604#[doc(hidden)]
605pub struct LoggingWriter<W, I> {
606    info: ServiceInfo,
607    logger: I,
608    path: String,
609
610    written: u64,
611    inner: W,
612}
613
614impl<W, I> LoggingWriter<W, I> {
615    fn new(info: ServiceInfo, logger: I, path: &str, writer: W) -> Self {
616        Self {
617            info,
618            logger,
619            path: path.to_string(),
620
621            written: 0,
622            inner: writer,
623        }
624    }
625}
626
627impl<W: oio::Write, I: LoggingInterceptor> oio::Write for LoggingWriter<W, I> {
628    async fn write(&mut self, bs: Buffer) -> Result<()> {
629        let size = bs.len();
630
631        match self.inner.write(bs).await {
632            Ok(_) => {
633                self.written += size as u64;
634                Ok(())
635            }
636            Err(err) => {
637                self.logger.log(
638                    &self.info,
639                    Operation::Write,
640                    &[
641                        ("path", &self.path),
642                        ("written", &self.written.to_string()),
643                        ("size", &size.to_string()),
644                    ],
645                    "failed",
646                    Some(&err),
647                );
648                Err(err)
649            }
650        }
651    }
652
653    async fn copy_from(&mut self, path: &str, args: OpRead, range: BytesRange) -> Result<()> {
654        let size = range
655            .size()
656            .expect("writer copy range must be absolute and bounded");
657
658        match self.inner.copy_from(path, args, range).await {
659            Ok(()) => {
660                self.written += size;
661                Ok(())
662            }
663            Err(err) => {
664                self.logger.log(
665                    &self.info,
666                    Operation::Write,
667                    &[
668                        ("path", &self.path),
669                        ("source", path),
670                        ("written", &self.written.to_string()),
671                        ("size", &size.to_string()),
672                    ],
673                    "failed",
674                    Some(&err),
675                );
676                Err(err)
677            }
678        }
679    }
680
681    async fn abort(&mut self) -> Result<()> {
682        match self.inner.abort().await {
683            Ok(_) => {
684                self.logger.log(
685                    &self.info,
686                    Operation::Write,
687                    &[("path", &self.path), ("written", &self.written.to_string())],
688                    "abort succeeded",
689                    None,
690                );
691                Ok(())
692            }
693            Err(err) => {
694                self.logger.log(
695                    &self.info,
696                    Operation::Write,
697                    &[("path", &self.path), ("written", &self.written.to_string())],
698                    "abort failed",
699                    Some(&err),
700                );
701                Err(err)
702            }
703        }
704    }
705
706    async fn close(&mut self) -> Result<Metadata> {
707        match self.inner.close().await {
708            Ok(meta) => {
709                self.logger.log(
710                    &self.info,
711                    Operation::Write,
712                    &[("path", &self.path), ("written", &self.written.to_string())],
713                    "close succeeded",
714                    None,
715                );
716                Ok(meta)
717            }
718            Err(err) => {
719                self.logger.log(
720                    &self.info,
721                    Operation::Write,
722                    &[("path", &self.path), ("written", &self.written.to_string())],
723                    "close failed",
724                    Some(&err),
725                );
726                Err(err)
727            }
728        }
729    }
730}
731
732#[doc(hidden)]
733pub struct LoggingLister<P, I: LoggingInterceptor> {
734    info: ServiceInfo,
735    logger: I,
736    path: String,
737
738    listed: usize,
739    inner: P,
740}
741
742impl<P, I: LoggingInterceptor> LoggingLister<P, I> {
743    fn new(info: ServiceInfo, logger: I, path: &str, inner: P) -> Self {
744        Self {
745            info,
746            logger,
747            path: path.to_string(),
748
749            listed: 0,
750            inner,
751        }
752    }
753}
754
755impl<P: oio::List, I: LoggingInterceptor> oio::List for LoggingLister<P, I> {
756    async fn next(&mut self) -> Result<Option<oio::Entry>> {
757        let res = self.inner.next().await;
758
759        match &res {
760            Ok(Some(_)) => {
761                self.listed += 1;
762            }
763            Ok(None) => {
764                self.logger.log(
765                    &self.info,
766                    Operation::List,
767                    &[("path", &self.path), ("listed", &self.listed.to_string())],
768                    "finished",
769                    None,
770                );
771            }
772            Err(err) => {
773                self.logger.log(
774                    &self.info,
775                    Operation::List,
776                    &[("path", &self.path), ("listed", &self.listed.to_string())],
777                    "failed",
778                    Some(err),
779                );
780            }
781        };
782
783        res
784    }
785}
786
787#[doc(hidden)]
788pub struct LoggingDeleter<D, I: LoggingInterceptor> {
789    info: ServiceInfo,
790    logger: I,
791
792    deleted: usize,
793    inner: D,
794}
795
796impl<D, I: LoggingInterceptor> LoggingDeleter<D, I> {
797    fn new(info: ServiceInfo, logger: I, inner: D) -> Self {
798        Self {
799            info,
800            logger,
801
802            deleted: 0,
803            inner,
804        }
805    }
806}
807
808impl<D: oio::Delete, I: LoggingInterceptor> oio::Delete for LoggingDeleter<D, I> {
809    async fn delete(&mut self, path: &str, args: OpDelete) -> Result<()> {
810        let version = args
811            .version()
812            .map(|v| v.to_string())
813            .unwrap_or_else(|| "<latest>".to_string());
814
815        let res = self.inner.delete(path, args).await;
816
817        match &res {
818            Ok(_) => {
819                self.deleted += 1;
820            }
821            Err(err) => {
822                self.logger.log(
823                    &self.info,
824                    Operation::Delete,
825                    &[
826                        ("path", path),
827                        ("version", &version),
828                        ("deleted", &self.deleted.to_string()),
829                    ],
830                    "failed",
831                    Some(err),
832                );
833            }
834        };
835
836        res
837    }
838
839    async fn close(&mut self) -> Result<()> {
840        let res = self.inner.close().await;
841
842        match &res {
843            Ok(_) => {
844                self.logger.log(
845                    &self.info,
846                    Operation::Delete,
847                    &[("deleted", &self.deleted.to_string())],
848                    "succeeded",
849                    None,
850                );
851            }
852            Err(err) => {
853                self.logger.log(
854                    &self.info,
855                    Operation::Delete,
856                    &[("deleted", &self.deleted.to_string())],
857                    "failed",
858                    Some(err),
859                );
860            }
861        };
862
863        res
864    }
865}
866
867#[doc(hidden)]
868pub struct LoggingCopier<C, I: LoggingInterceptor> {
869    info: ServiceInfo,
870    logger: I,
871    from: String,
872    to: String,
873
874    copied: u64,
875    inner: C,
876}
877
878impl<C, I: LoggingInterceptor> LoggingCopier<C, I> {
879    fn new(info: ServiceInfo, logger: I, from: &str, to: &str, inner: C) -> Self {
880        Self {
881            info,
882            logger,
883            from: from.to_string(),
884            to: to.to_string(),
885
886            copied: 0,
887            inner,
888        }
889    }
890}
891
892impl<C: oio::Copy, I: LoggingInterceptor> oio::Copy for LoggingCopier<C, I> {
893    async fn next(&mut self) -> Result<Option<usize>> {
894        match self.inner.next().await {
895            Ok(Some(n)) => {
896                self.copied += n as u64;
897                Ok(Some(n))
898            }
899            Ok(None) => {
900                self.logger.log(
901                    &self.info,
902                    Operation::Copy,
903                    &[
904                        ("from", &self.from),
905                        ("to", &self.to),
906                        ("copied", &self.copied.to_string()),
907                    ],
908                    "finished",
909                    None,
910                );
911                Ok(None)
912            }
913            Err(err) => {
914                self.logger.log(
915                    &self.info,
916                    Operation::Copy,
917                    &[
918                        ("from", &self.from),
919                        ("to", &self.to),
920                        ("copied", &self.copied.to_string()),
921                    ],
922                    "failed",
923                    Some(&err),
924                );
925                Err(err)
926            }
927        }
928    }
929
930    async fn close(&mut self) -> Result<Metadata> {
931        self.inner.close().await
932    }
933
934    async fn abort(&mut self) -> Result<()> {
935        match self.inner.abort().await {
936            Ok(_) => {
937                self.logger.log(
938                    &self.info,
939                    Operation::Copy,
940                    &[
941                        ("from", &self.from),
942                        ("to", &self.to),
943                        ("copied", &self.copied.to_string()),
944                    ],
945                    "abort succeeded",
946                    None,
947                );
948                Ok(())
949            }
950            Err(err) => {
951                self.logger.log(
952                    &self.info,
953                    Operation::Copy,
954                    &[
955                        ("from", &self.from),
956                        ("to", &self.to),
957                        ("copied", &self.copied.to_string()),
958                    ],
959                    "abort failed",
960                    Some(&err),
961                );
962                Err(err)
963            }
964        }
965    }
966}