Struct PipelineBuilder
pub struct PipelineBuilder { /* private fields */ }Available on crate feature
dial9 only.Expand description
Closure-scoped builder for assembling a custom processor pipeline.
Obtained via with_custom_pipeline(|p| ...) on the runtime builder.
Built-in processors are reachable through dedicated methods
(gzip, write_back,
s3, symbolize); custom processors
are added with pipe.
§Example
ⓘ
struct Logger;
impl SegmentProcessor for Logger {
fn name(&self) -> &'static str { "logger" }
fn process(&mut self, data: SegmentData)
-> Pin<Box<dyn Future<Output = Result<SegmentData, ProcessError>> + Send + '_>>
{
Box::pin(async move {
println!("segment {} ({} bytes)", data.segment().index(), data.payload().len());
Ok(data)
})
}
}
builder.with_custom_pipeline(|p| p.pipe(Logger).gzip().write_back())Implementations§
§impl PipelineBuilder
impl PipelineBuilder
pub fn pipe<S>(self, processor: S) -> PipelineBuilderwhere
S: SegmentProcessor + 'static,
pub fn pipe<S>(self, processor: S) -> PipelineBuilderwhere
S: SegmentProcessor + 'static,
Append a user-supplied SegmentProcessor to the pipeline.
pub fn gzip(self) -> PipelineBuilder
pub fn gzip(self) -> PipelineBuilder
Gzip the segment payload in-memory.
pub fn write_back(self) -> PipelineBuilder
pub fn write_back(self) -> PipelineBuilder
Write the current payload bytes back to disk. When the payload has
been gzipped earlier in the pipeline, the file is written with a
.gz suffix and the original sealed segment is removed.
Trait Implementations§
Auto Trait Implementations§
impl !RefUnwindSafe for PipelineBuilder
impl !Sync for PipelineBuilder
impl !UnwindSafe for PipelineBuilder
impl Freeze for PipelineBuilder
impl Send for PipelineBuilder
impl Unpin for PipelineBuilder
impl UnsafeUnpin for PipelineBuilder
Blanket Implementations§
§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
§impl<T> FutureExt for T
impl<T> FutureExt for T
§fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
§fn with_current_context(self) -> WithContext<Self> ⓘ
fn with_current_context(self) -> WithContext<Self> ⓘ
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
Converts
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
Converts
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
Wrap the input message
T in a rama_grpc::Request§impl<T> Pointable for T
impl<T> Pointable for T
§impl<T> PolicyExt for Twhere
T: ?Sized,
impl<T> PolicyExt for Twhere
T: ?Sized,
§fn and<P, B, E>(self, other: P) -> And<T, P>
fn and<P, B, E>(self, other: P) -> And<T, P>
Create a new
Policy that returns Action::Follow only if self and other return
Action::Follow. Read more