#![cfg(target_os = "linux")]
#![cfg_attr(docsrs, feature(doc_cfg))]
#![deny(missing_docs)]
use std::ffi::CString;
use std::fmt::Debug;
use std::fmt::Formatter;
use bytes::Buf;
use opendal_core::raw::*;
use opendal_core::*;
use probe::probe_lazy;
#[derive(Clone, Default)]
#[non_exhaustive]
pub struct DtraceLayer {}
impl DtraceLayer {
pub fn new() -> Self {
Self::default()
}
}
impl<A: Access> Layer<A> for DtraceLayer {
type LayeredAccess = DTraceAccessor<A>;
fn layer(&self, inner: A) -> Self::LayeredAccess {
DTraceAccessor { inner }
}
}
#[doc(hidden)]
pub struct DTraceAccessor<A: Access> {
inner: A,
}
impl<A: Access> Debug for DTraceAccessor<A> {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
f.debug_struct("DTraceAccessor")
.field("inner", &self.inner)
.finish_non_exhaustive()
}
}
impl<A: Access> LayeredAccess for DTraceAccessor<A> {
type Inner = A;
type Reader = DtraceLayerWrapper<A::Reader>;
type Writer = DtraceLayerWrapper<A::Writer>;
type Lister = A::Lister;
type Deleter = A::Deleter;
type Copier = A::Copier;
fn inner(&self) -> &Self::Inner {
&self.inner
}
async fn create_dir(&self, path: &str, args: OpCreateDir) -> Result<RpCreateDir> {
let c_path = CString::new(path).unwrap();
probe_lazy!(opendal, create_dir_start, c_path.as_ptr());
let result = self.inner.create_dir(path, args).await;
probe_lazy!(opendal, create_dir_end, c_path.as_ptr());
result
}
async fn read(&self, path: &str, args: OpRead) -> Result<(RpRead, Self::Reader)> {
let c_path = CString::new(path).unwrap();
probe_lazy!(opendal, read_start, c_path.as_ptr());
let result = self
.inner
.read(path, args)
.await
.map(|(rp, r)| (rp, DtraceLayerWrapper::new(r, &path.to_string())));
probe_lazy!(opendal, read_end, c_path.as_ptr());
result
}
async fn write(&self, path: &str, args: OpWrite) -> Result<(RpWrite, Self::Writer)> {
let c_path = CString::new(path).unwrap();
probe_lazy!(opendal, write_start, c_path.as_ptr());
let result = self
.inner
.write(path, args)
.await
.map(|(rp, r)| (rp, DtraceLayerWrapper::new(r, &path.to_string())));
probe_lazy!(opendal, write_end, c_path.as_ptr());
result
}
async fn copy(
&self,
from: &str,
to: &str,
args: OpCopy,
opts: OpCopier,
) -> Result<(RpCopy, Self::Copier)> {
let c_from = CString::new(from).unwrap();
probe_lazy!(opendal, copy_start, c_from.as_ptr());
let result = self.inner.copy(from, to, args, opts).await;
probe_lazy!(opendal, copy_end, c_from.as_ptr());
result
}
async fn stat(&self, path: &str, args: OpStat) -> Result<RpStat> {
let c_path = CString::new(path).unwrap();
probe_lazy!(opendal, stat_start, c_path.as_ptr());
let result = self.inner.stat(path, args).await;
probe_lazy!(opendal, stat_end, c_path.as_ptr());
result
}
async fn delete(&self) -> Result<(RpDelete, Self::Deleter)> {
self.inner.delete().await
}
async fn list(&self, path: &str, args: OpList) -> Result<(RpList, Self::Lister)> {
let c_path = CString::new(path).unwrap();
probe_lazy!(opendal, list_start, c_path.as_ptr());
let result = self.inner.list(path, args).await;
probe_lazy!(opendal, list_end, c_path.as_ptr());
result
}
async fn presign(&self, path: &str, args: OpPresign) -> Result<RpPresign> {
let c_path = CString::new(path).unwrap();
probe_lazy!(opendal, presign_start, c_path.as_ptr());
let result = self.inner.presign(path, args).await;
probe_lazy!(opendal, presign_end, c_path.as_ptr());
result
}
}
#[doc(hidden)]
pub struct DtraceLayerWrapper<R> {
inner: R,
path: String,
}
impl<R> DtraceLayerWrapper<R> {
fn new(inner: R, path: &String) -> Self {
Self {
inner,
path: path.to_string(),
}
}
}
impl<R: oio::Read> oio::Read for DtraceLayerWrapper<R> {
async fn read(&mut self) -> Result<Buffer> {
let c_path = CString::new(self.path.clone()).unwrap();
probe_lazy!(opendal, reader_read_start, c_path.as_ptr());
match self.inner.read().await {
Ok(bs) => {
probe_lazy!(opendal, reader_read_ok, c_path.as_ptr(), bs.remaining());
Ok(bs)
}
Err(e) => {
probe_lazy!(opendal, reader_read_error, c_path.as_ptr());
Err(e)
}
}
}
}
impl<R: oio::Write> oio::Write for DtraceLayerWrapper<R> {
async fn write(&mut self, bs: Buffer) -> Result<()> {
let c_path = CString::new(self.path.clone()).unwrap();
probe_lazy!(opendal, writer_write_start, c_path.as_ptr());
self.inner
.write(bs)
.await
.map(|_| {
probe_lazy!(opendal, writer_write_ok, c_path.as_ptr());
})
.inspect_err(|_| {
probe_lazy!(opendal, writer_write_error, c_path.as_ptr());
})
}
async fn abort(&mut self) -> Result<()> {
let c_path = CString::new(self.path.clone()).unwrap();
probe_lazy!(opendal, writer_poll_abort_start, c_path.as_ptr());
self.inner
.abort()
.await
.map(|_| {
probe_lazy!(opendal, writer_poll_abort_ok, c_path.as_ptr());
})
.inspect_err(|_| {
probe_lazy!(opendal, writer_poll_abort_error, c_path.as_ptr());
})
}
async fn close(&mut self) -> Result<Metadata> {
let c_path = CString::new(self.path.clone()).unwrap();
probe_lazy!(opendal, writer_close_start, c_path.as_ptr());
self.inner
.close()
.await
.inspect(|_| {
probe_lazy!(opendal, writer_close_ok, c_path.as_ptr());
})
.inspect_err(|_| {
probe_lazy!(opendal, writer_close_error, c_path.as_ptr());
})
}
}