Skip to main content

TransformingSource

Struct TransformingSource 

Source
pub struct TransformingSource { /* private fields */ }
Expand description

Source decorator that applies a fixed list of compiled stages to every record. Emits faucet_transform_* metrics per page via instrumented_apply_stages.

§Example

use faucet_core::{RecordTransform, Source, TransformingSource};
use faucet_core::observability::Labels;
use faucet_core::stage::TransformStage;
use faucet_core::transform::KeyCaseMode;

let inner: Box<dyn Source> = build_inner();
let wrapped = TransformingSource::new(
    inner,
    vec![TransformStage::Map(RecordTransform::KeysCase { mode: KeyCaseMode::Snake, on_collision: Default::default() })],
    Labels::for_named("rest"),
).unwrap();

Implementations§

Source§

impl TransformingSource

Source

pub fn new( inner: Box<dyn Source>, stages: Vec<TransformStage>, labels: Labels, ) -> Result<Self, FaucetError>

Compile stages and wrap inner. Returns FaucetError::Transform if any stage’s compilation fails (e.g. invalid regex in RenameKeys). The chain stays on the Value path (no columnar batch forms).

Source

pub fn new_with_batches( inner: Box<dyn Source>, stages: Vec<TransformStage>, batch_fns: Vec<Option<PageFnBatchBox>>, labels: Labels, ) -> Result<Self, FaucetError>

Available on crate feature arrow only.

Like new, but each stage may carry an Arrow RecordBatch form (batch_fns[i] parallels stages[i]), so the chain can run on the columnar fast path (#375) when the inner source and sink are Arrow-native and every stage supplies one. Used by the CLI for sql transforms.

Trait Implementations§

Source§

impl Source for TransformingSource

Source§

fn supports_columnar(&self) -> bool

Available on crate feature arrow only.

Columnar only when the inner source is columnar and every stage has an Arrow batch form (today: the SQL transform). Any Value-only stage (Map / Filter / Explode / CdcUnwrap / Custom / plain PageFn) keeps the whole chain on the Value path (#375).

Source§

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

Available on crate feature arrow only.

Stream the inner source’s Arrow batches with every stage’s batch form applied in declared order — so parquet → sql → parquet runs Arrow end-to-end. Only reached when supports_columnar is true, i.e. every stage has a batch form.

Source§

fn fetch_with_context<'life0, 'life1, 'async_trait>( &'life0 self, ctx: &'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, ctx: &'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 stream_pages<'a>( &'a self, ctx: &'a HashMap<String, Value>, batch_size: usize, ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>>

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 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 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 dataset_uri(&self) -> String

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

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

Receive a pre-labelled handle for counting upstream round trips (#638) — calls this connector makes to its own backend. Read more
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 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 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<'a>( &'a self, context: &'a HashMap<String, Value>, format: NativeFormat, batch_size: usize, ) -> Pin<Box<dyn Stream<Item = Result<NativeBatch, FaucetError>> + Send + 'a>>

Stream the source natively as byte batches in format (#633). Read more
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 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§

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