Skip to main content

SourceCtx

Struct SourceCtx 

Source
#[non_exhaustive]
pub struct SourceCtx { pub issuer: AckIssuer, pub meter: Option<Meter>, pub stage_metrics: Option<Arc<SourceMetrics>>, pub per_partition_detail: bool, }
Expand description

Everything a source receives at Source::open.

Fields (Non-exhaustive)§

This struct is marked as non-exhaustive
Non-exhaustive structs could have additional fields added in future. Therefore, non-exhaustive structs cannot be constructed in external crates using the traditional Struct { .. } syntax; cannot be matched against without a wildcard ..; and struct update syntax will not work.
§issuer: AckIssuer

Issuer for batch acknowledgement handles. Sources clone it into every lane they construct; each lane issues one AckRef per poll batch (issue(partition, last_offset)).

§meter: Option<Meter>

A Meter scoped spate_<component_type>_source_* for the source’s own metric families (e.g. consumer lag, broker statistics), pre-labelled with the standard pipeline/component/component_type. None unless the source declared a Source::component_type that is a usable, non-reserved namespace — a reserved default ("source") opts out silently, a malformed value is logged and also yields None. Resolve handles from it once here in open; never on the poll path.

§stage_metrics: Option<Arc<SourceMetrics>>

The framework’s own source-stage handles (spate_source_*), shared with the controller. A source that can observe its own consumer lag publishes it here — SourceMetrics::set_partition_lag and SourceMetrics::retain_partitions have no other caller, because the framework cannot compute lag without the client’s view of the log end, nor tell which partitions the client still owns. None only when the source is driven outside a pipeline (tests, or a direct open call), in which case lag simply goes unpublished.

Everything else on these handles — records, bytes, poll duration, rebalances, active lanes — is recorded by the runtime itself; a source must not touch those.

§per_partition_detail: bool

Whether cardinality-sensitive per-partition series are enabled (metrics.per_partition_detail). Gates a connector’s own per-partition families: when false, register and emit only aggregate (per-component or per-broker) series. It does not gate spate_source_lag_records — consumer lag has no aggregate series to fall back to, so it always publishes per partition.

Implementations§

Source§

impl SourceCtx

Source

pub fn new(issuer: AckIssuer) -> Self

Context wrapping the checkpointer’s issuer. The custom-metrics meter is None; the runtime attaches one via with_meter.

Source

pub fn with_meter(self, meter: Option<Meter>) -> Self

Attach the source’s custom-metrics scope. Called by the runtime, which builds it from the source’s component_type.

Source

pub fn with_stage_metrics(self, metrics: Option<Arc<SourceMetrics>>) -> Self

Share the framework’s source-stage handles so the source can publish the one series only it can measure: consumer lag. Called by the runtime with the same instance the controller records against.

Source

pub fn with_partition_detail(self, enabled: bool) -> Self

Enable cardinality-sensitive per-partition series. Called by the runtime from metrics.per_partition_detail.

Trait Implementations§

Source§

impl Debug for SourceCtx

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

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<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> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<A, B, T> HttpServerConnExec<A, B> for T
where B: Body,

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> 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> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

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