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 opts: OpCopier,
374 ) -> Result<Self::Copier> {
375 self.log_start(Operation::Copy, &[("from", from), ("to", to)]);
376 self.inner
377 .copy(ctx, from, to, args, opts)
378 .map(|c| {
379 self.logger.log(
380 &self.info,
381 Operation::Copy,
382 &[("from", from), ("to", to)],
383 "created copier",
384 None,
385 );
386 LoggingCopier::new(self.info.clone(), self.logger.clone(), from, to, c)
387 })
388 .inspect_err(|err| {
389 self.logger.log(
390 &self.info,
391 Operation::Copy,
392 &[("from", from), ("to", to)],
393 "failed",
394 Some(err),
395 );
396 })
397 }
398
399 async fn rename(
400 &self,
401 ctx: &OperationContext,
402 from: &str,
403 to: &str,
404 args: OpRename,
405 ) -> Result<RpRename> {
406 self.log_start(Operation::Rename, &[("from", from), ("to", to)]);
407 let result = self.inner.rename(ctx, from, to, args).await;
408 self.log_finish(
409 Operation::Rename,
410 &[("from", from), ("to", to)],
411 result.as_ref().err(),
412 );
413 result
414 }
415
416 fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
417 self.log_start(Operation::List, &[("path", path)]);
418 self.inner
419 .list(ctx, path, args)
420 .map(|v| {
421 self.logger.log(
422 &self.info,
423 Operation::List,
424 &[("path", path)],
425 "created lister",
426 None,
427 );
428 LoggingLister::new(self.info.clone(), self.logger.clone(), path, v)
429 })
430 .inspect_err(|err| {
431 self.logger.log(
432 &self.info,
433 Operation::List,
434 &[("path", path)],
435 "failed",
436 Some(err),
437 );
438 })
439 }
440
441 async fn presign(
442 &self,
443 ctx: &OperationContext,
444 path: &str,
445 args: OpPresign,
446 ) -> Result<RpPresign> {
447 self.log_start(Operation::Presign, &[("path", path)]);
448 let result = self.inner.presign(ctx, path, args).await;
449 self.log_finish(Operation::Presign, &[("path", path)], result.as_ref().err());
450 result
451 }
452}
453
454#[doc(hidden)]
455pub struct LoggingReader<R, I: LoggingInterceptor> {
456 info: ServiceInfo,
457 logger: I,
458 path: String,
459 range: Option<BytesRange>,
460
461 read: u64,
462 inner: R,
463}
464
465impl<R, I: LoggingInterceptor> LoggingReader<R, I> {
466 fn new(info: ServiceInfo, logger: I, path: &str, reader: R) -> Self {
467 Self::with_range(info, logger, path, None, reader)
468 }
469
470 fn with_range(
471 info: ServiceInfo,
472 logger: I,
473 path: &str,
474 range: Option<BytesRange>,
475 reader: R,
476 ) -> Self {
477 Self {
478 info,
479 logger,
480 path: path.to_string(),
481 range,
482
483 read: 0,
484 inner: reader,
485 }
486 }
487
488 fn range_label(&self) -> String {
489 self.range
490 .map(|range| range.to_string())
491 .unwrap_or_default()
492 }
493}
494
495impl<R: oio::ReadStream, I: LoggingInterceptor> oio::ReadStream for LoggingReader<R, I> {
496 async fn read(&mut self) -> Result<Buffer> {
497 match self.inner.read().await {
498 Ok(bs) if bs.is_empty() => {
499 let range = self.range_label();
500 self.logger.log(
501 &self.info,
502 Operation::Read,
503 &[
504 ("path", &self.path),
505 ("range", &range),
506 ("read", &self.read.to_string()),
507 ("size", &bs.len().to_string()),
508 ],
509 "finished",
510 None,
511 );
512 Ok(bs)
513 }
514 Ok(bs) => {
515 self.read += bs.len() as u64;
516 Ok(bs)
517 }
518 Err(err) => {
519 let range = self.range_label();
520 self.logger.log(
521 &self.info,
522 Operation::Read,
523 &[
524 ("path", &self.path),
525 ("range", &range),
526 ("read", &self.read.to_string()),
527 ],
528 "failed",
529 Some(&err),
530 );
531 Err(err)
532 }
533 }
534 }
535}
536
537impl<R: oio::Read, I: LoggingInterceptor> oio::Read for LoggingReader<R, I> {
538 async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
539 match self.inner.open(range).await {
540 Ok((rp, stream)) => Ok((
541 rp,
542 Box::new(LoggingReader::with_range(
543 self.info.clone(),
544 self.logger.clone(),
545 &self.path,
546 Some(range),
547 stream,
548 )) as Box<dyn oio::ReadStreamDyn>,
549 )),
550 Err(err) => {
551 self.logger.log(
552 &self.info,
553 Operation::Read,
554 &[("path", &self.path), ("range", &range.to_string())],
555 "failed",
556 Some(&err),
557 );
558 Err(err)
559 }
560 }
561 }
562
563 async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
564 match self.inner.read(range).await {
565 Ok((rp, buffer)) => {
566 self.logger.log(
567 &self.info,
568 Operation::Read,
569 &[
570 ("path", &self.path),
571 ("range", &range.to_string()),
572 ("size", &buffer.len().to_string()),
573 ],
574 "finished",
575 None,
576 );
577 Ok((rp, buffer))
578 }
579 Err(err) => {
580 self.logger.log(
581 &self.info,
582 Operation::Read,
583 &[("path", &self.path), ("range", &range.to_string())],
584 "failed",
585 Some(&err),
586 );
587 Err(err)
588 }
589 }
590 }
591}
592
593#[doc(hidden)]
594pub struct LoggingWriter<W, I> {
595 info: ServiceInfo,
596 logger: I,
597 path: String,
598
599 written: u64,
600 inner: W,
601}
602
603impl<W, I> LoggingWriter<W, I> {
604 fn new(info: ServiceInfo, logger: I, path: &str, writer: W) -> Self {
605 Self {
606 info,
607 logger,
608 path: path.to_string(),
609
610 written: 0,
611 inner: writer,
612 }
613 }
614}
615
616impl<W: oio::Write, I: LoggingInterceptor> oio::Write for LoggingWriter<W, I> {
617 async fn write(&mut self, bs: Buffer) -> Result<()> {
618 let size = bs.len();
619
620 match self.inner.write(bs).await {
621 Ok(_) => {
622 self.written += size as u64;
623 Ok(())
624 }
625 Err(err) => {
626 self.logger.log(
627 &self.info,
628 Operation::Write,
629 &[
630 ("path", &self.path),
631 ("written", &self.written.to_string()),
632 ("size", &size.to_string()),
633 ],
634 "failed",
635 Some(&err),
636 );
637 Err(err)
638 }
639 }
640 }
641
642 async fn abort(&mut self) -> Result<()> {
643 match self.inner.abort().await {
644 Ok(_) => {
645 self.logger.log(
646 &self.info,
647 Operation::Write,
648 &[("path", &self.path), ("written", &self.written.to_string())],
649 "abort succeeded",
650 None,
651 );
652 Ok(())
653 }
654 Err(err) => {
655 self.logger.log(
656 &self.info,
657 Operation::Write,
658 &[("path", &self.path), ("written", &self.written.to_string())],
659 "abort failed",
660 Some(&err),
661 );
662 Err(err)
663 }
664 }
665 }
666
667 async fn close(&mut self) -> Result<Metadata> {
668 match self.inner.close().await {
669 Ok(meta) => {
670 self.logger.log(
671 &self.info,
672 Operation::Write,
673 &[("path", &self.path), ("written", &self.written.to_string())],
674 "close succeeded",
675 None,
676 );
677 Ok(meta)
678 }
679 Err(err) => {
680 self.logger.log(
681 &self.info,
682 Operation::Write,
683 &[("path", &self.path), ("written", &self.written.to_string())],
684 "close failed",
685 Some(&err),
686 );
687 Err(err)
688 }
689 }
690 }
691}
692
693#[doc(hidden)]
694pub struct LoggingLister<P, I: LoggingInterceptor> {
695 info: ServiceInfo,
696 logger: I,
697 path: String,
698
699 listed: usize,
700 inner: P,
701}
702
703impl<P, I: LoggingInterceptor> LoggingLister<P, I> {
704 fn new(info: ServiceInfo, logger: I, path: &str, inner: P) -> Self {
705 Self {
706 info,
707 logger,
708 path: path.to_string(),
709
710 listed: 0,
711 inner,
712 }
713 }
714}
715
716impl<P: oio::List, I: LoggingInterceptor> oio::List for LoggingLister<P, I> {
717 async fn next(&mut self) -> Result<Option<oio::Entry>> {
718 let res = self.inner.next().await;
719
720 match &res {
721 Ok(Some(_)) => {
722 self.listed += 1;
723 }
724 Ok(None) => {
725 self.logger.log(
726 &self.info,
727 Operation::List,
728 &[("path", &self.path), ("listed", &self.listed.to_string())],
729 "finished",
730 None,
731 );
732 }
733 Err(err) => {
734 self.logger.log(
735 &self.info,
736 Operation::List,
737 &[("path", &self.path), ("listed", &self.listed.to_string())],
738 "failed",
739 Some(err),
740 );
741 }
742 };
743
744 res
745 }
746}
747
748#[doc(hidden)]
749pub struct LoggingDeleter<D, I: LoggingInterceptor> {
750 info: ServiceInfo,
751 logger: I,
752
753 deleted: usize,
754 inner: D,
755}
756
757impl<D, I: LoggingInterceptor> LoggingDeleter<D, I> {
758 fn new(info: ServiceInfo, logger: I, inner: D) -> Self {
759 Self {
760 info,
761 logger,
762
763 deleted: 0,
764 inner,
765 }
766 }
767}
768
769impl<D: oio::Delete, I: LoggingInterceptor> oio::Delete for LoggingDeleter<D, I> {
770 async fn delete(&mut self, path: &str, args: OpDelete) -> Result<()> {
771 let version = args
772 .version()
773 .map(|v| v.to_string())
774 .unwrap_or_else(|| "<latest>".to_string());
775
776 let res = self.inner.delete(path, args).await;
777
778 match &res {
779 Ok(_) => {
780 self.deleted += 1;
781 }
782 Err(err) => {
783 self.logger.log(
784 &self.info,
785 Operation::Delete,
786 &[
787 ("path", path),
788 ("version", &version),
789 ("deleted", &self.deleted.to_string()),
790 ],
791 "failed",
792 Some(err),
793 );
794 }
795 };
796
797 res
798 }
799
800 async fn close(&mut self) -> Result<()> {
801 let res = self.inner.close().await;
802
803 match &res {
804 Ok(_) => {
805 self.logger.log(
806 &self.info,
807 Operation::Delete,
808 &[("deleted", &self.deleted.to_string())],
809 "succeeded",
810 None,
811 );
812 }
813 Err(err) => {
814 self.logger.log(
815 &self.info,
816 Operation::Delete,
817 &[("deleted", &self.deleted.to_string())],
818 "failed",
819 Some(err),
820 );
821 }
822 };
823
824 res
825 }
826}
827
828#[doc(hidden)]
829pub struct LoggingCopier<C, I: LoggingInterceptor> {
830 info: ServiceInfo,
831 logger: I,
832 from: String,
833 to: String,
834
835 copied: u64,
836 inner: C,
837}
838
839impl<C, I: LoggingInterceptor> LoggingCopier<C, I> {
840 fn new(info: ServiceInfo, logger: I, from: &str, to: &str, inner: C) -> Self {
841 Self {
842 info,
843 logger,
844 from: from.to_string(),
845 to: to.to_string(),
846
847 copied: 0,
848 inner,
849 }
850 }
851}
852
853impl<C: oio::Copy, I: LoggingInterceptor> oio::Copy for LoggingCopier<C, I> {
854 async fn next(&mut self) -> Result<Option<usize>> {
855 match self.inner.next().await {
856 Ok(Some(n)) => {
857 self.copied += n as u64;
858 Ok(Some(n))
859 }
860 Ok(None) => {
861 self.logger.log(
862 &self.info,
863 Operation::Copy,
864 &[
865 ("from", &self.from),
866 ("to", &self.to),
867 ("copied", &self.copied.to_string()),
868 ],
869 "finished",
870 None,
871 );
872 Ok(None)
873 }
874 Err(err) => {
875 self.logger.log(
876 &self.info,
877 Operation::Copy,
878 &[
879 ("from", &self.from),
880 ("to", &self.to),
881 ("copied", &self.copied.to_string()),
882 ],
883 "failed",
884 Some(&err),
885 );
886 Err(err)
887 }
888 }
889 }
890
891 async fn close(&mut self) -> Result<Metadata> {
892 self.inner.close().await
893 }
894
895 async fn abort(&mut self) -> Result<()> {
896 match self.inner.abort().await {
897 Ok(_) => {
898 self.logger.log(
899 &self.info,
900 Operation::Copy,
901 &[
902 ("from", &self.from),
903 ("to", &self.to),
904 ("copied", &self.copied.to_string()),
905 ],
906 "abort succeeded",
907 None,
908 );
909 Ok(())
910 }
911 Err(err) => {
912 self.logger.log(
913 &self.info,
914 Operation::Copy,
915 &[
916 ("from", &self.from),
917 ("to", &self.to),
918 ("copied", &self.copied.to_string()),
919 ],
920 "abort failed",
921 Some(&err),
922 );
923 Err(err)
924 }
925 }
926 }
927}