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
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}