Struct WriteBackProcessor
pub struct WriteBackProcessor { /* private fields */ }Available on crate features
dial9 and pipeline only.Expand description
Writes the current payload bytes back to disk. If a
write_back_extension metadata key is present, the bytes are written to
{original}{extension} and the original segment file is removed.
When dir is set, the file is written to that directory instead of
alongside the original.
Implementations§
§impl WriteBackProcessor
impl WriteBackProcessor
pub fn to_dir(dir: PathBuf) -> WriteBackProcessor
pub fn to_dir(dir: PathBuf) -> WriteBackProcessor
Write to dir instead of alongside the original segment.
Trait Implementations§
§impl Debug for WriteBackProcessor
impl Debug for WriteBackProcessor
§impl Default for WriteBackProcessor
impl Default for WriteBackProcessor
§fn default() -> WriteBackProcessor
fn default() -> WriteBackProcessor
Returns the “default value” for a type. Read more
§impl SegmentProcessor for WriteBackProcessor
impl SegmentProcessor for WriteBackProcessor
§fn process(
&mut self,
data: SegmentData,
) -> Pin<Box<dyn Future<Output = Result<SegmentData, ProcessError>> + Send + '_>>
fn process( &mut self, data: SegmentData, ) -> Pin<Box<dyn Future<Output = Result<SegmentData, ProcessError>> + Send + '_>>
Process a segment, transforming or consuming its data.
Returns the (possibly modified) data for the next processor,
or an error to skip this segment.
§fn initialize(
&mut self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + '_>>
fn initialize( &mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + '_>>
Initialize this processor on the worker’s Tokio runtime before the
pipeline begins processing segments. Read more
§fn finalize_dump(
&mut self,
completion: &DumpCompletion,
) -> Pin<Box<dyn Future<Output = Option<String>> + Send + '_>>
fn finalize_dump( &mut self, completion: &DumpCompletion, ) -> Pin<Box<dyn Future<Output = Option<String>> + Send + '_>>
Called once per finished dump in triggered mode (see
crate::dump), in pipeline order, so stages can flush any
per-dump state they accumulated. Return the S3 key of a manifest
written for this dump, or None; the last Some across the
pipeline lands on DumpReceipt::manifest_key. Read moreAuto Trait Implementations§
impl Freeze for WriteBackProcessor
impl RefUnwindSafe for WriteBackProcessor
impl Send for WriteBackProcessor
impl Sync for WriteBackProcessor
impl Unpin for WriteBackProcessor
impl UnsafeUnpin for WriteBackProcessor
impl UnwindSafe for WriteBackProcessor
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