Skip to main content

InstrumentedSource

Struct InstrumentedSource 

Source
pub struct InstrumentedSource<'a, S: Source + ?Sized> { /* private fields */ }
Expand description

Wraps a &dyn Source (or any &S: Source) and emits spans + metrics around every call. Constructed by Pipeline::run and never exposed to end users; the wrapped source remains the user-facing object.

Implementations§

Source§

impl<'a, S: Source + ?Sized> InstrumentedSource<'a, S>

Source

pub fn new(inner: &'a S, labels: Labels) -> Self

Source

pub fn with_meter(self, meter: Arc<UsageMeter>) -> Self

Tally records and estimated bytes into a run’s usage meter (#704).

Trait Implementations§

Source§

impl<'a, S: Source + ?Sized> Source for InstrumentedSource<'a, S>

Source§

fn set_roundtrip_recorder(&self, recorder: Arc<RoundtripRecorder>)

Forward the round-trip recorder to the wrapped connector (#638) — the decorator sits between the pipeline and the connector, so without this the hook would never reach the code that performs the I/O.

Source§

fn connector_name(&self) -> &'static str

Stable identifier used as the connector label on metrics and the connector attribute on spans. Defaults to the final segment of std::any::type_name::<Self>(), e.g. "RestSource". Built-in connectors override with a short, friendly snake_case name (e.g. "rest"). Must return a non-empty string; observability decorators fall back to "unknown" in release builds if it is empty (and debug_assert! in debug builds).
Source§

fn set_run_clock(&self, now: DateTime<Utc>)

Receive the run clock (#769): the instant ${now.*} renders from — faucet run --clock, a schedule tick, a backfill unit’s start — so a source that bounds reads by “now” (REST window slicing) reproduces the run as of that instant rather than the wall clock. Read more
Source§

fn state_key(&self) -> Option<String>

Stable key under which this source’s incremental-replication bookmark should be persisted in a StateStore. Read more
Source§

fn apply_start_bookmark<'life0, 'async_trait>( &'life0 self, bookmark: Value, ) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Apply a bookmark loaded from a StateStore as this run’s starting point. Read more
Source§

fn supports_exactly_once(&self) -> bool

Whether this source deterministically replays the same page sequence from a given bookmark — the requirement for the atomic-watermark effectively-once path (a non-deterministic replay could cause the pipeline to skip a page whose contents differ from the one already committed). Default: false. Read more
Source§

fn replay_guarantee(&self) -> ReplayGuarantee

The typed replay capability this source advertises — see ReplayGuarantee. Read more
Source§

fn capture_resume_position<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Option<Value>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Capture the source’s current replication position without consuming any changes, ensuring any server-side resource (e.g. a logical replication slot) needed to later resume from that position exists. Read more
Source§

fn lag<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Option<SourceLag>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

How far this source is behind its head (#733) — unread WAL bytes, binlog distance, unconsumed messages, or the age of the oldest unread change. Read more
Source§

fn state_schema(&self) -> u32

The version of this source’s bookmark shape (#736). Stored state carries it; bump it whenever the shape changes and teach migrate_state the step. Default 0. Decorators must forward this.
Source§

fn migrate_state(&self, from: u32, data: Value) -> Result<Value, FaucetError>

Bring a bookmark stored at shape version from up to state_schema (#736). Must be pure — the migrated value is only persisted by the next bookmark write, so a crash in between re-runs it. The default knows only its own version. Decorators must forward this.
Source§

fn record_table(&self, record: &Value) -> Option<String>

The dataset (table / collection) a change record belongs to, for a change stream that carries several tables (#731). A multi-table faucet mirror runs one stream and routes each record to its table’s pipeline by this name, which must match the name the paired bulk source’s discover reports (e.g. public.orders). None means the record belongs to no table (a DDL / control event) — and is the default, so a source that does not override this cannot be used for a multi-table mirror.
Source§

fn position_le(&self, a: &Value, b: &Value) -> Option<bool>

Order two of this source’s bookmarks (#731): Some(true) when stream position a is at or before b (every change up to a is also covered by b), Some(false) when it is not, and None when the source cannot tell (the default, or two positions it cannot relate). A multi-table mirror uses it to resume one shared stream from the earliest table and skip, per table, pages that table has already committed.
Source§

fn position_min(&self, positions: &[Value]) -> Option<Value>

The earliest stream position every one of positions can resume from (#731): the shared change stream restarts there and each table skips what it has already applied. The default picks the position that position_le orders at or before all the others, and None when there is none (an empty slice, or positions this source cannot order). Sources whose positions are only partially ordered (one cursor per capture instance, say) override it with a component-wise minimum.
Source§

fn fetch_with_context<'life0, 'life1, 'async_trait>( &'life0 self, context: &'life1 HashMap<String, Value>, ) -> Pin<Box<dyn Future<Output = Result<Vec<Value>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Primary fetch method. Receives context from a parent source’s records. Read more
Source§

fn fetch_with_context_incremental<'life0, 'life1, 'async_trait>( &'life0 self, context: &'life1 HashMap<String, Value>, ) -> Pin<Box<dyn Future<Output = Result<(Vec<Value>, Option<Value>), FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Incremental fetch with parent context support. Read more
Source§

fn supports_columnar(&self) -> bool

Available on crate feature arrow only.
Whether this source can emit columnar (ColumnarPage) pages via stream_batches. Default: false. Read more
Source§

fn stream_batches<'b>( &'b self, context: &'b HashMap<String, Value>, batch_size: usize, ) -> Pin<Box<dyn Stream<Item = Result<ColumnarPage, FaucetError>> + Send + 'b>>

Available on crate feature arrow only.
Stream the source natively as Arrow ColumnarPages. Read more
Source§

fn native_output_formats(&self) -> &'static [NativeFormat]

Wire formats this source can emit as raw bytes for the native byte-passthrough fast path (#633), in preference order (first = best). Default: &[] (no native fast path). Read more
Source§

fn stream_native<'b>( &'b self, context: &'b HashMap<String, Value>, format: NativeFormat, batch_size: usize, ) -> Pin<Box<dyn Stream<Item = Result<NativeBatch, FaucetError>> + Send + 'b>>

Stream the source natively as byte batches in format (#633). Read more
Source§

fn stream_pages<'b>( &'b self, context: &'b HashMap<String, Value>, batch_size: usize, ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'b>>

Stream records page-by-page so the pipeline can write to the sink as pages arrive instead of buffering the full result set. Read more
Source§

fn fetch_all<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Vec<Value>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Convenience: fetch with no parent context.
Source§

fn fetch_all_incremental<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<(Vec<Value>, Option<Value>), FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Convenience: incremental fetch with no parent context.
Source§

fn config_schema(&self) -> Value

Return a JSON Schema describing the configuration this source accepts.
Source§

fn is_shardable(&self) -> bool

Whether this source can split its work into independent shards for clustered (Mode B) execution. Default: false (single whole-dataset shard). Sources with a natural partition (object-store prefixes, table primary-key ranges) override this to true and implement enumerate_shards + apply_shard.
Source§

fn enumerate_shards<'life0, 'async_trait>( &'life0 self, _target: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<ShardSpec>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Enumerate the shards this source splits into, aiming for roughly target of them (a hint — the source may return fewer, e.g. when the dataset is small, or one per natural partition regardless of target). Read more
Source§

fn apply_shard<'life0, 'life1, 'async_trait>( &'life0 self, _shard: &'life1 ShardSpec, ) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Narrow this source instance to a single shard before streaming. Read more
Source§

fn range_digest<'life0, 'life1, 'life2, 'life3, 'async_trait>( &'life0 self, range: &'life1 KeyRange, key: &'life2 str, columns: &'life3 [String], ) -> Pin<Box<dyn Future<Output = Result<Option<ServerDigest>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait,

Compute a ServerDigest of the rows in range inside the backend, so a matching range ships no rows (#701). columns are the compared columns; key the integer key the range is over. Default Ok(None): not supported, the verifier streams the range instead. Two digests compare only when both sides report the same algorithm.
Source§

fn supports_discover(&self) -> bool

Whether this source can enumerate the datasets behind its connection via discover. Default: false. Sources backed by an introspectable catalog (database information_schema, MongoDB collections, Elasticsearch indices, object-store prefixes) override this to true.
Source§

fn discover<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Vec<DatasetDescriptor>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Enumerate the datasets living behind this source’s connection — one DatasetDescriptor per table / collection / index / prefix, each carrying a partial config override that selects it (used by faucet discover to scaffold one matrix row per dataset). Read more
Source§

fn dataset_uri(&self) -> String

Logical dataset identity for lineage emission, following OpenLineage naming conventions (https://openlineage.io/docs/spec/naming). Read more
Source§

fn check<'life0, 'life1, 'async_trait>( &'life0 self, ctx: &'life1 CheckContext, ) -> Pin<Box<dyn Future<Output = Result<CheckReport, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Run a fast, non-mutating preflight probe (used by faucet doctor). Read more

Auto Trait Implementations§

§

impl<'a, S> Freeze for InstrumentedSource<'a, S>
where &'a S: Freeze, S: ?Sized,

§

impl<'a, S> RefUnwindSafe for InstrumentedSource<'a, S>

§

impl<'a, S> Send for InstrumentedSource<'a, S>
where &'a S: Send, S: ?Sized,

§

impl<'a, S> Sync for InstrumentedSource<'a, S>
where &'a S: Sync, S: ?Sized,

§

impl<'a, S> Unpin for InstrumentedSource<'a, S>
where &'a S: Unpin, S: ?Sized,

§

impl<'a, S> UnsafeUnpin for InstrumentedSource<'a, S>
where &'a S: UnsafeUnpin, S: ?Sized,

§

impl<'a, S> UnwindSafe for InstrumentedSource<'a, S>
where &'a S: UnwindSafe, S: ?Sized,

Blanket Implementations§

Source§

impl<T> Allocation for T
where T: RefUnwindSafe + Send + Sync,

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

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self> ⓘ

Attaches the current Context to this type, returning a WithContext wrapper. Read more
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> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
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> 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 = !

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

fn try_from(value: U) -> Result<T, !>

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