Struct PipelineBuilder
pub struct PipelineBuilder<Mode = Disk>where
Mode: BufferMode,{ /* private fields */ }dial9 and pipeline only.Expand description
Closure-scoped builder for assembling a custom processor pipeline.
Obtained via with_custom_pipeline(|p| ...) on the runtime builder. The
Mode type parameter binds the pipeline to the writer’s storage mode:
disk-only processors like write_back are not in
scope on PipelineBuilder<Memory>, so wiring write-back into an
in-memory pipeline is a compile error.
§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<Mode> PipelineBuilder<Mode>where
Mode: BufferMode,
impl<Mode> PipelineBuilder<Mode>where
Mode: BufferMode,
pub fn pipe<S>(self, processor: S) -> PipelineBuilder<Mode>where
S: SegmentProcessor + 'static,
pub fn pipe<S>(self, processor: S) -> PipelineBuilder<Mode>where
S: SegmentProcessor + 'static,
Append a user-supplied SegmentProcessor to the pipeline.
pub fn gzip(self) -> PipelineBuilder<Mode>
pub fn gzip(self) -> PipelineBuilder<Mode>
Gzip the segment payload in-memory.
§impl PipelineBuilder
Disk-only methods on the pipeline builder.
impl PipelineBuilder
Disk-only methods on the pipeline builder.
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.
pub fn write_back_to(self, dir: impl Into<PathBuf>) -> PipelineBuilder
pub fn write_back_to(self, dir: impl Into<PathBuf>) -> PipelineBuilder
Write the current payload bytes to a specific directory instead of
back alongside the original segment. The file name is preserved;
when the payload has been gzipped, a .gz suffix is appended.
The original sealed segment is removed after a successful write.
Trait Implementations§
Auto Trait Implementations§
impl<Mode = Disk> !RefUnwindSafe for PipelineBuilder<Mode>
impl<Mode = Disk> !Sync for PipelineBuilder<Mode>
impl<Mode = Disk> !UnwindSafe for PipelineBuilder<Mode>
impl<Mode> Freeze for PipelineBuilder<Mode>where
PhantomData<Mode>: Freeze,
impl<Mode> Send for PipelineBuilder<Mode>where
PhantomData<Mode>: Send,
impl<Mode> Unpin for PipelineBuilder<Mode>where
PhantomData<Mode>: Unpin,
impl<Mode> UnsafeUnpin for PipelineBuilder<Mode>where
PhantomData<Mode>: UnsafeUnpin,
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
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> ⓘ
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> ⓘ
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>
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>
Policy that returns Action::Follow only if self and other return
Action::Follow. Read more