1#![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#[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 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
150pub trait LoggingInterceptor: Debug + Clone + Send + Sync + Unpin + 'static {
152 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#[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 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 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 type Composer = oio::Composer;
271
272 fn info(&self) -> ServiceInfo {
273 self.info.clone()
274 }
275
276 fn capability(&self) -> Capability {
277 self.inner.capability()
278 }
279
280 fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result<Self::Composer> {
281 self.inner.compose(ctx, to, args)
282 }
283
284 async fn create_dir(
285 &self,
286 ctx: &OperationContext,
287 path: &str,
288 args: OpCreateDir,
289 ) -> Result<RpCreateDir> {
290 self.log_start(Operation::CreateDir, &[("path", path)]);
291 let result = self.inner.create_dir(ctx, path, args).await;
292 self.log_finish(
293 Operation::CreateDir,
294 &[("path", path)],
295 result.as_ref().err(),
296 );
297 result
298 }
299
300 fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
301 self.log_start(Operation::Read, &[("path", path)]);
302 self.inner
303 .read(ctx, path, args)
304 .map(|r| {
305 self.logger.log(
306 &self.info,
307 Operation::Read,
308 &[("path", path)],
309 "created reader",
310 None,
311 );
312 LoggingReader::new(self.info.clone(), self.logger.clone(), path, r)
313 })
314 .inspect_err(|err| {
315 self.logger.log(
316 &self.info,
317 Operation::Read,
318 &[("path", path)],
319 "failed",
320 Some(err),
321 );
322 })
323 }
324
325 fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
326 self.log_start(Operation::Write, &[("path", path)]);
327 self.inner
328 .write(ctx, path, args)
329 .map(|w| {
330 self.logger.log(
331 &self.info,
332 Operation::Write,
333 &[("path", path)],
334 "created writer",
335 None,
336 );
337 LoggingWriter::new(self.info.clone(), self.logger.clone(), path, w)
338 })
339 .inspect_err(|err| {
340 self.logger.log(
341 &self.info,
342 Operation::Write,
343 &[("path", path)],
344 "failed",
345 Some(err),
346 );
347 })
348 }
349
350 async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
351 self.log_start(Operation::Stat, &[("path", path)]);
352 let result = self.inner.stat(ctx, path, args).await;
353 self.log_finish(Operation::Stat, &[("path", path)], result.as_ref().err());
354 result
355 }
356
357 fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
358 self.log_start(Operation::Delete, &[]);
359 self.inner
360 .delete(ctx)
361 .map(|d| {
362 self.logger
363 .log(&self.info, Operation::Delete, &[], "finished", None);
364 LoggingDeleter::new(self.info.clone(), self.logger.clone(), d)
365 })
366 .inspect_err(|err| {
367 self.logger
368 .log(&self.info, Operation::Delete, &[], "failed", Some(err));
369 })
370 }
371
372 fn copy(
373 &self,
374 ctx: &OperationContext,
375 from: &str,
376 to: &str,
377 args: OpCopy,
378 ) -> Result<Self::Copier> {
379 self.log_start(Operation::Copy, &[("from", from), ("to", to)]);
380 self.inner
381 .copy(ctx, from, to, args)
382 .map(|c| {
383 self.logger.log(
384 &self.info,
385 Operation::Copy,
386 &[("from", from), ("to", to)],
387 "created copier",
388 None,
389 );
390 LoggingCopier::new(self.info.clone(), self.logger.clone(), from, to, c)
391 })
392 .inspect_err(|err| {
393 self.logger.log(
394 &self.info,
395 Operation::Copy,
396 &[("from", from), ("to", to)],
397 "failed",
398 Some(err),
399 );
400 })
401 }
402
403 async fn rename(
404 &self,
405 ctx: &OperationContext,
406 from: &str,
407 to: &str,
408 args: OpRename,
409 ) -> Result<RpRename> {
410 self.log_start(Operation::Rename, &[("from", from), ("to", to)]);
411 let result = self.inner.rename(ctx, from, to, args).await;
412 self.log_finish(
413 Operation::Rename,
414 &[("from", from), ("to", to)],
415 result.as_ref().err(),
416 );
417 result
418 }
419
420 async fn restore(
421 &self,
422 ctx: &OperationContext,
423 path: &str,
424 args: OpRestore,
425 ) -> Result<RpRestore> {
426 self.log_start(Operation::Restore, &[("path", path)]);
427 let result = self.inner.restore(ctx, path, args).await;
428 self.log_finish(Operation::Restore, &[("path", path)], result.as_ref().err());
429 result
430 }
431
432 fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
433 self.log_start(Operation::List, &[("path", path)]);
434 self.inner
435 .list(ctx, path, args)
436 .map(|v| {
437 self.logger.log(
438 &self.info,
439 Operation::List,
440 &[("path", path)],
441 "created lister",
442 None,
443 );
444 LoggingLister::new(self.info.clone(), self.logger.clone(), path, v)
445 })
446 .inspect_err(|err| {
447 self.logger.log(
448 &self.info,
449 Operation::List,
450 &[("path", path)],
451 "failed",
452 Some(err),
453 );
454 })
455 }
456
457 async fn presign(
458 &self,
459 ctx: &OperationContext,
460 path: &str,
461 args: OpPresign,
462 ) -> Result<RpPresign> {
463 self.log_start(Operation::Presign, &[("path", path)]);
464 let result = self.inner.presign(ctx, path, args).await;
465 self.log_finish(Operation::Presign, &[("path", path)], result.as_ref().err());
466 result
467 }
468}
469
470#[doc(hidden)]
471pub struct LoggingReader<R, I: LoggingInterceptor> {
472 info: ServiceInfo,
473 logger: I,
474 path: String,
475 range: Option<BytesRange>,
476
477 read: u64,
478 inner: R,
479}
480
481impl<R, I: LoggingInterceptor> LoggingReader<R, I> {
482 fn new(info: ServiceInfo, logger: I, path: &str, reader: R) -> Self {
483 Self::with_range(info, logger, path, None, reader)
484 }
485
486 fn with_range(
487 info: ServiceInfo,
488 logger: I,
489 path: &str,
490 range: Option<BytesRange>,
491 reader: R,
492 ) -> Self {
493 Self {
494 info,
495 logger,
496 path: path.to_string(),
497 range,
498
499 read: 0,
500 inner: reader,
501 }
502 }
503
504 fn range_label(&self) -> String {
505 self.range
506 .map(|range| range.to_string())
507 .unwrap_or_default()
508 }
509}
510
511impl<R: oio::ReadStream, I: LoggingInterceptor> oio::ReadStream for LoggingReader<R, I> {
512 async fn read(&mut self) -> Result<Buffer> {
513 match self.inner.read().await {
514 Ok(bs) if bs.is_empty() => {
515 let range = self.range_label();
516 self.logger.log(
517 &self.info,
518 Operation::Read,
519 &[
520 ("path", &self.path),
521 ("range", &range),
522 ("read", &self.read.to_string()),
523 ("size", &bs.len().to_string()),
524 ],
525 "finished",
526 None,
527 );
528 Ok(bs)
529 }
530 Ok(bs) => {
531 self.read += bs.len() as u64;
532 Ok(bs)
533 }
534 Err(err) => {
535 let range = self.range_label();
536 self.logger.log(
537 &self.info,
538 Operation::Read,
539 &[
540 ("path", &self.path),
541 ("range", &range),
542 ("read", &self.read.to_string()),
543 ],
544 "failed",
545 Some(&err),
546 );
547 Err(err)
548 }
549 }
550 }
551}
552
553impl<R: oio::Read, I: LoggingInterceptor> oio::Read for LoggingReader<R, I> {
554 async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
555 match self.inner.open(range).await {
556 Ok((rp, stream)) => Ok((
557 rp,
558 Box::new(LoggingReader::with_range(
559 self.info.clone(),
560 self.logger.clone(),
561 &self.path,
562 Some(range),
563 stream,
564 )) as Box<dyn oio::ReadStreamDyn>,
565 )),
566 Err(err) => {
567 self.logger.log(
568 &self.info,
569 Operation::Read,
570 &[("path", &self.path), ("range", &range.to_string())],
571 "failed",
572 Some(&err),
573 );
574 Err(err)
575 }
576 }
577 }
578
579 async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
580 match self.inner.read(range).await {
581 Ok((rp, buffer)) => {
582 self.logger.log(
583 &self.info,
584 Operation::Read,
585 &[
586 ("path", &self.path),
587 ("range", &range.to_string()),
588 ("size", &buffer.len().to_string()),
589 ],
590 "finished",
591 None,
592 );
593 Ok((rp, buffer))
594 }
595 Err(err) => {
596 self.logger.log(
597 &self.info,
598 Operation::Read,
599 &[("path", &self.path), ("range", &range.to_string())],
600 "failed",
601 Some(&err),
602 );
603 Err(err)
604 }
605 }
606 }
607}
608
609#[doc(hidden)]
610pub struct LoggingWriter<W, I> {
611 info: ServiceInfo,
612 logger: I,
613 path: String,
614
615 written: u64,
616 inner: W,
617}
618
619impl<W, I> LoggingWriter<W, I> {
620 fn new(info: ServiceInfo, logger: I, path: &str, writer: W) -> Self {
621 Self {
622 info,
623 logger,
624 path: path.to_string(),
625
626 written: 0,
627 inner: writer,
628 }
629 }
630}
631
632impl<W: oio::Write, I: LoggingInterceptor> oio::Write for LoggingWriter<W, I> {
633 async fn write(&mut self, bs: Buffer) -> Result<()> {
634 let size = bs.len();
635
636 match self.inner.write(bs).await {
637 Ok(_) => {
638 self.written += size as u64;
639 Ok(())
640 }
641 Err(err) => {
642 self.logger.log(
643 &self.info,
644 Operation::Write,
645 &[
646 ("path", &self.path),
647 ("written", &self.written.to_string()),
648 ("size", &size.to_string()),
649 ],
650 "failed",
651 Some(&err),
652 );
653 Err(err)
654 }
655 }
656 }
657
658 async fn copy_from(&mut self, path: &str, args: OpRead, range: BytesRange) -> Result<()> {
659 let size = range
660 .size()
661 .expect("writer copy range must be absolute and bounded");
662
663 match self.inner.copy_from(path, args, range).await {
664 Ok(()) => {
665 self.written += size;
666 Ok(())
667 }
668 Err(err) => {
669 self.logger.log(
670 &self.info,
671 Operation::Write,
672 &[
673 ("path", &self.path),
674 ("source", path),
675 ("written", &self.written.to_string()),
676 ("size", &size.to_string()),
677 ],
678 "failed",
679 Some(&err),
680 );
681 Err(err)
682 }
683 }
684 }
685
686 async fn abort(&mut self) -> Result<()> {
687 match self.inner.abort().await {
688 Ok(_) => {
689 self.logger.log(
690 &self.info,
691 Operation::Write,
692 &[("path", &self.path), ("written", &self.written.to_string())],
693 "abort succeeded",
694 None,
695 );
696 Ok(())
697 }
698 Err(err) => {
699 self.logger.log(
700 &self.info,
701 Operation::Write,
702 &[("path", &self.path), ("written", &self.written.to_string())],
703 "abort failed",
704 Some(&err),
705 );
706 Err(err)
707 }
708 }
709 }
710
711 async fn close(&mut self) -> Result<Metadata> {
712 match self.inner.close().await {
713 Ok(meta) => {
714 self.logger.log(
715 &self.info,
716 Operation::Write,
717 &[("path", &self.path), ("written", &self.written.to_string())],
718 "close succeeded",
719 None,
720 );
721 Ok(meta)
722 }
723 Err(err) => {
724 self.logger.log(
725 &self.info,
726 Operation::Write,
727 &[("path", &self.path), ("written", &self.written.to_string())],
728 "close failed",
729 Some(&err),
730 );
731 Err(err)
732 }
733 }
734 }
735}
736
737#[doc(hidden)]
738pub struct LoggingLister<P, I: LoggingInterceptor> {
739 info: ServiceInfo,
740 logger: I,
741 path: String,
742
743 listed: usize,
744 inner: P,
745}
746
747impl<P, I: LoggingInterceptor> LoggingLister<P, I> {
748 fn new(info: ServiceInfo, logger: I, path: &str, inner: P) -> Self {
749 Self {
750 info,
751 logger,
752 path: path.to_string(),
753
754 listed: 0,
755 inner,
756 }
757 }
758}
759
760impl<P: oio::List, I: LoggingInterceptor> oio::List for LoggingLister<P, I> {
761 async fn next(&mut self) -> Result<Option<oio::Entry>> {
762 let res = self.inner.next().await;
763
764 match &res {
765 Ok(Some(_)) => {
766 self.listed += 1;
767 }
768 Ok(None) => {
769 self.logger.log(
770 &self.info,
771 Operation::List,
772 &[("path", &self.path), ("listed", &self.listed.to_string())],
773 "finished",
774 None,
775 );
776 }
777 Err(err) => {
778 self.logger.log(
779 &self.info,
780 Operation::List,
781 &[("path", &self.path), ("listed", &self.listed.to_string())],
782 "failed",
783 Some(err),
784 );
785 }
786 };
787
788 res
789 }
790}
791
792#[doc(hidden)]
793pub struct LoggingDeleter<D, I: LoggingInterceptor> {
794 info: ServiceInfo,
795 logger: I,
796
797 deleted: usize,
798 inner: D,
799}
800
801impl<D, I: LoggingInterceptor> LoggingDeleter<D, I> {
802 fn new(info: ServiceInfo, logger: I, inner: D) -> Self {
803 Self {
804 info,
805 logger,
806
807 deleted: 0,
808 inner,
809 }
810 }
811}
812
813impl<D: oio::Delete, I: LoggingInterceptor> oio::Delete for LoggingDeleter<D, I> {
814 async fn delete(&mut self, path: &str, args: OpDelete) -> Result<()> {
815 let version = args
816 .version()
817 .map(|v| v.to_string())
818 .unwrap_or_else(|| "<latest>".to_string());
819
820 let res = self.inner.delete(path, args).await;
821
822 match &res {
823 Ok(_) => {
824 self.deleted += 1;
825 }
826 Err(err) => {
827 self.logger.log(
828 &self.info,
829 Operation::Delete,
830 &[
831 ("path", path),
832 ("version", &version),
833 ("deleted", &self.deleted.to_string()),
834 ],
835 "failed",
836 Some(err),
837 );
838 }
839 };
840
841 res
842 }
843
844 async fn close(&mut self) -> Result<()> {
845 let res = self.inner.close().await;
846
847 match &res {
848 Ok(_) => {
849 self.logger.log(
850 &self.info,
851 Operation::Delete,
852 &[("deleted", &self.deleted.to_string())],
853 "succeeded",
854 None,
855 );
856 }
857 Err(err) => {
858 self.logger.log(
859 &self.info,
860 Operation::Delete,
861 &[("deleted", &self.deleted.to_string())],
862 "failed",
863 Some(err),
864 );
865 }
866 };
867
868 res
869 }
870}
871
872#[doc(hidden)]
873pub struct LoggingCopier<C, I: LoggingInterceptor> {
874 info: ServiceInfo,
875 logger: I,
876 from: String,
877 to: String,
878
879 copied: u64,
880 inner: C,
881}
882
883impl<C, I: LoggingInterceptor> LoggingCopier<C, I> {
884 fn new(info: ServiceInfo, logger: I, from: &str, to: &str, inner: C) -> Self {
885 Self {
886 info,
887 logger,
888 from: from.to_string(),
889 to: to.to_string(),
890
891 copied: 0,
892 inner,
893 }
894 }
895}
896
897impl<C: oio::Copy, I: LoggingInterceptor> oio::Copy for LoggingCopier<C, I> {
898 async fn next(&mut self) -> Result<Option<usize>> {
899 match self.inner.next().await {
900 Ok(Some(n)) => {
901 self.copied += n as u64;
902 Ok(Some(n))
903 }
904 Ok(None) => {
905 self.logger.log(
906 &self.info,
907 Operation::Copy,
908 &[
909 ("from", &self.from),
910 ("to", &self.to),
911 ("copied", &self.copied.to_string()),
912 ],
913 "finished",
914 None,
915 );
916 Ok(None)
917 }
918 Err(err) => {
919 self.logger.log(
920 &self.info,
921 Operation::Copy,
922 &[
923 ("from", &self.from),
924 ("to", &self.to),
925 ("copied", &self.copied.to_string()),
926 ],
927 "failed",
928 Some(&err),
929 );
930 Err(err)
931 }
932 }
933 }
934
935 async fn close(&mut self) -> Result<Metadata> {
936 self.inner.close().await
937 }
938
939 async fn abort(&mut self) -> Result<()> {
940 match self.inner.abort().await {
941 Ok(_) => {
942 self.logger.log(
943 &self.info,
944 Operation::Copy,
945 &[
946 ("from", &self.from),
947 ("to", &self.to),
948 ("copied", &self.copied.to_string()),
949 ],
950 "abort succeeded",
951 None,
952 );
953 Ok(())
954 }
955 Err(err) => {
956 self.logger.log(
957 &self.info,
958 Operation::Copy,
959 &[
960 ("from", &self.from),
961 ("to", &self.to),
962 ("copied", &self.copied.to_string()),
963 ],
964 "abort failed",
965 Some(&err),
966 );
967 Err(err)
968 }
969 }
970 }
971}