Skip to main content

ProcessorContext

Struct ProcessorContext 

Source
pub struct ProcessorContext<'ctx, 'd, KOut, VOut> { /* private fields */ }
Expand description

Handed to Processor::process. forward boxes the record and queues it for each child node (the driver drains the queue).

Two lifetimes: 'ctx is the borrow of the Dispatch reference itself; 'd is the lifetime of the data inside Dispatch (buffers, slices, etc.). Keeping them separate avoids lifetime-invariance issues when constructing a ProcessorContext from a &mut Dispatch<'d> with an independently-scoped outer borrow 'ctx.

Implementations§

Source§

impl<'ctx, 'd, KOut, VOut> ProcessorContext<'ctx, 'd, KOut, VOut>
where KOut: Any + Send + Clone, VOut: Any + Send + Clone,

Source

pub fn forward(&mut self, record: Record<KOut, VOut>)

Forward a record to all child nodes. The record is cloned per child for fan-out; the last child receives the original by move (so the common single-child case performs zero clones). Mirrors the JVM ProcessorContext.forward(Record), which takes the record by value.

Source

pub fn get_state_store<K2: Send + Sync + 'static, V2: Send + 'static>( &mut self, name: &str, ) -> Option<&mut dyn KeyValueStore<K2, V2>>

Access a connected state store, typed. None if absent or the K/V types don’t match. Fetch it per-record (do not hold across process calls).

Source

pub async fn global_get<GK: Send + Sync + 'static, VG: Send + 'static>( &mut self, store: &str, key: &GK, ) -> Option<VG>

Look up a value in a connected GLOBAL store (fully-replicated, shared across tasks). Returns an owned value — no borrow escapes the shared manager’s lock, so the lookup future need not be held across forward. None on miss / type mismatch. Fetch it per-record (do not hold across process calls).

Source

pub fn get_window_store<K2: Send + Sync + 'static, V2: Send + 'static>( &mut self, name: &str, ) -> Option<&mut dyn WindowStore<K2, V2>>

Access a connected window store, typed. None if absent or the K/V types don’t match. Fetch it per-record (do not hold across process calls).

Source

pub fn get_join_window_store<K2: Send + Sync + 'static, V2: Send + 'static>( &mut self, name: &str, ) -> Option<&mut dyn JoinWindowStore<K2, V2>>

Access a connected join-window store (retainDuplicates), typed. None if absent or the K/V types don’t match. Fetch it per-record (do not hold across process calls).

Source

pub fn get_session_store<K2: Send + Sync + 'static, V2: Send + 'static>( &mut self, name: &str, ) -> Option<&mut dyn SessionStore<K2, V2>>

Access a connected session store, typed. None if absent or the K/V types don’t match. Fetch it per-record (do not hold across process calls).

Source

pub fn get_versioned_store<K2: Send + Sync + 'static, V2: Send + 'static>( &mut self, name: &str, ) -> Option<&mut dyn VersionedKeyValueStore<K2, V2>>

Access a connected versioned store (KIP-889), typed. None if absent or the K/V types don’t match. Fetch it per-record (do not hold across process calls).

Source

pub fn record_context(&self) -> &RecordContext

Metadata of the source record currently being processed.

Source

pub fn store_is_cached(&self, name: &str) -> bool

Whether the named KV state store is record-cached (so this processor should suppress its immediate forward and let the cache flush forward the deduped change). False for absent/non-KV/uncached stores.

Source

pub fn schedule<P>( &mut self, interval: Duration, ty: PunctuationType, punctuator: P, ) -> Cancellable
where P: Punctuator<KOut, VOut>,

Schedule a periodic Punctuator. Callable from init or process. interval must be positive. Returns a Cancellable to stop it.

Auto Trait Implementations§

§

impl<'ctx, 'd, KOut, VOut> !RefUnwindSafe for ProcessorContext<'ctx, 'd, KOut, VOut>

§

impl<'ctx, 'd, KOut, VOut> !Sync for ProcessorContext<'ctx, 'd, KOut, VOut>

§

impl<'ctx, 'd, KOut, VOut> !UnwindSafe for ProcessorContext<'ctx, 'd, KOut, VOut>

§

impl<'ctx, 'd, KOut, VOut> Freeze for ProcessorContext<'ctx, 'd, KOut, VOut>

§

impl<'ctx, 'd, KOut, VOut> Send for ProcessorContext<'ctx, 'd, KOut, VOut>

§

impl<'ctx, 'd, KOut, VOut> Unpin for ProcessorContext<'ctx, 'd, KOut, VOut>

§

impl<'ctx, 'd, KOut, VOut> UnsafeUnpin for ProcessorContext<'ctx, 'd, KOut, VOut>

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> Downcast for T
where T: Any,

Source§

fn into_any(self: Box<T>) -> Box<dyn Any>

Converts Box<dyn Trait> (where Trait: Downcast) to Box<dyn Any>, which can then be downcast into Box<dyn ConcreteType> where ConcreteType implements Trait.
Source§

fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>

Converts Rc<Trait> (where Trait: Downcast) to Rc<Any>, which can then be further downcast into Rc<ConcreteType> where ConcreteType implements Trait.
Source§

fn as_any(&self) -> &(dyn Any + 'static)

Converts &Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot generate &Any’s vtable from &Trait’s.
Source§

fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)

Converts &mut Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot generate &mut Any’s vtable from &mut Trait’s.
Source§

impl<T> DowncastSend for T
where T: Any + Send,

Source§

fn into_any_send(self: Box<T>) -> Box<dyn Any + Send>

Converts Box<Trait> (where Trait: DowncastSend) to Box<dyn Any + Send>, which can then be downcast into Box<ConcreteType> where ConcreteType implements Trait.
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Fruit for T
where T: Send + Downcast,

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

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
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more