Expand description
§faucet-core
Shared types, traits, and utilities for the faucet-stream ecosystem.
This crate provides the common foundation used by all faucet source and sink connectors:
FaucetError— unified error typeSource/Sink— async traits for data connectorsRecordTransform— record transformation pipelineReplicationMethod— incremental replication supportschema::infer_schema— JSON Schema inference from record samples
Re-exports§
pub use adaptive::AdaptiveBatchConfig;pub use adaptive::AdjustDirection;pub use adaptive::AdjustReason;pub use adaptive::Adjustment;pub use adaptive::AimdController;pub use adaptive::Observation;pub use anomaly::AnomalyMethod;pub use auth::AuthProvider;pub use auth::AuthReference;pub use auth::AuthSpec;pub use auth::Credential;pub use auth::CredentialPlacement;pub use auth::RequestAuth;pub use budget::BudgetKind;pub use budget::BudgetSink;pub use budget::BudgetSpec;pub use budget::BudgetState;pub use budget::BudgetTimer;pub use budget::BudgetVerdict;pub use check::CheckContext;pub use check::CheckReport;pub use check::Probe;pub use check::ProbeStatus;pub use cleanup::CleanupMode;pub use cleanup::CleanupPolicy;pub use cleanup::DEFAULT_MAX_KEYS;pub use cleanup::SeenKeys;pub use columnar::ColumnarPage;arrowpub use columnar::infer_arrow_schema;arrowpub use columnar::record_batch_to_values;arrowpub use columnar::values_to_record_batch;arrowpub use columnar::values_to_record_batch_inferred;arrowpub use create_table::PlannedColumn;pub use create_table::missing_target_error;pub use create_table::plan_columns;pub use create_table::plan_keyed_columns;pub use create_table::render_column_defs;pub use create_table::render_columns;pub use create_table::render_primary_key;pub use cross_join::CompiledCrossJoin;transform-cross-joinpub use cross_join::CrossJoinSpec;transform-cross-joinpub use cross_join::OnEmpty as CrossJoinOnEmpty;transform-cross-joinpub use diff::ContentDigest;pub use diff::Difference;pub use diff::DifferenceKind;pub use diff::DigestAccumulator;pub use diff::KeyRange;pub use diff::Normalizer;pub use diff::ServerDigest;pub use diff::VerifyReport;pub use diff::diff_rows;pub use diff::plan_ranges;pub use diff::row_hash;pub use discover::DatasetDescriptor;pub use discover::attach_primary_keys;pub use discover::columns_to_schema;pub use discover::nullable_type;pub use discover::sql_type_to_json_schema;pub use dlq::BatchAtomicity;pub use dlq::BatchOutcome;pub use dlq::BatchOutcomeCounters;pub use dlq::BatchOutcomes;pub use dlq::DlqConfig;pub use dlq::DlqReason;pub use dlq::DlqStats;pub use dlq::EnvelopeError;pub use dlq::OnBatchError;pub use dlq::UnwrappedEnvelope;pub use dlq::build_envelope;pub use dlq::check_dlq_all_policy;pub use dlq::dlq_all_is_safe;pub use dlq::dlq_all_refusal;pub use dlq::unwrap_envelope;pub use drift::ColumnChange;pub use drift::OnDrift;pub use drift::OnIncompatible;pub use drift::SchemaDiff;pub use drift::SchemaDriftPolicy;pub use drift::SchemaDriftSpec;pub use drift::SchemaEvolution;pub use drift::SqlBaseType;pub use drift::adds_null;pub use drift::base_widened;pub use drift::json_schema_base_type;pub use encryption::CompiledEncryption;encryptionpub use encryption::EncryptionAlgorithm;encryptionpub use encryption::EncryptionSpec;encryptionpub use error::FaucetError;pub use file_format::parquet_io::ParquetReadOptions;pub use file_format::AvroCodec;pub use file_format::AvroOptions;pub use file_format::ContainerDecoder;pub use file_format::CsvOptions;pub use file_format::CsvUnknownField;pub use file_format::ExcelOptions;pub use file_format::FileFormat;pub use file_format::FileInput;pub use file_format::FormatOptions;pub use file_format::OrcOptions;pub use file_format::XmlOptions;pub use idempotency::DeliveryGuarantee;pub use idempotency::DeliveryMode;pub use idempotency::EffectivelyOnceMechanism;pub use idempotency::GuaranteeInputs;pub use idempotency::ReplayGuarantee;pub use idempotency::SinkGuarantee;pub use idempotency::derive_delivery_guarantee;pub use idempotency::format_token;pub use idempotency::format_token_with_bookmark;pub use idempotency::parse_token;pub use idempotency::parse_token_parts;pub use idempotency::unwrap_state;pub use idempotency::wrap_state;pub use join::HashJoin;pub use join::JoinConfig;pub use join::JoinMode;pub use join::JoinStats;pub use join::KeyNormalize;pub use join::OnCollision;pub use join::OnDuplicate;pub use join::Projection;pub use lag::LagObserver;pub use lag::SourceLag;pub use local_outputs::LocalOutput;pub use local_outputs::LocalOutputLog;pub use local_outputs::probe_pre_existing;pub use metadata::CompiledMetadata;pub use metadata::MetadataColumn;pub use metadata::MetadataColumnsSpec;pub use metadata::MetadataContext;pub use metadata::MetadataSink;pub use native::CsvDialect;pub use native::NativeBatch;pub use native::NativeFormat;pub use native::NativeLoadCapability;pub use native::NativeLoadContext;pub use native::NativePayload;pub use native::NativePlan;pub use native::NativePlanInputs;pub use native::plan_native_transfer;pub use object_rollover::CompletedObject;pub use object_rollover::ObjectAccumulator;pub use object_rollover::PageAccumulator;pub use observability::instrumented_apply_contract;contractpub use observability::instrumented_apply_masking;maskingpub use observability::instrumented_apply_quality;qualitypub use observability::otel::OtelConfig;pub use observability::otel::OtelProtocol;pub use observability::otel::OtelSignal;pub use observability::otel::shutdown_otel;pub use observability::DurationGuard;pub use observability::InstallError;pub use observability::InstallReport;pub use observability::InstrumentedSink;pub use observability::InstrumentedSource;pub use observability::InstrumentedStateStore;pub use observability::Labels;pub use observability::ObservabilityConfig;pub use observability::PrometheusConfig;pub use observability::RunStreamOptions;pub use observability::TracingConfig;pub use observability::install_observability;pub use observability::instrumented_apply_stages;pub use observability::register_build_info;pub use observability::update_bookmark_lag;pub use pipeline::DEFAULT_BATCH_SIZE;pub use pipeline::MAX_BATCH_SIZE;pub use pipeline::Pipeline;pub use pipeline::PipelineResult;pub use pipeline::StreamPage;pub use pipeline::run_stream;pub use pipeline::validate_batch_size;pub use policy::ColumnFacts;policypub use policy::CompiledPolicy;policypub use policy::PolicyRule;policypub use policy::PolicyScope;policypub use policy::PolicySink;policypub use policy::PolicySpec;policypub use policy::RuntimeAction;policypub use policy::SinkFacts;policypub use policy::Violation;policypub use policy::ViolationKind;policypub use profiling::ColumnProfile;pub use profiling::DriftMetric;pub use profiling::OnProfileDrift;pub use profiling::ProfileDrift;pub use profiling::Profiler;pub use profiling::ProfilingSink;pub use profiling::ProfilingSpec;pub use profiling::RunProfile;pub use profiling::detect_drift;pub use replication::BindFormat;pub use replication::BindTarget;pub use replication::BindValueType;pub use replication::IncrementalFilter;pub use replication::OnMissingKey;pub use replication::ReplicationBind;pub use replication::ReplicationKey;pub use replication::ReplicationMethod;pub use replication::format_bookmark;pub use replication::format_instant;pub use replication::json_gt;pub use replication::parse_instant;pub use replication::set_body_pointer;pub use resilience::BackoffFrom;pub use resilience::BackoffKind;pub use resilience::CircuitBreaker;pub use resilience::CircuitBreakerConfig;pub use resilience::PoisonAction;pub use resilience::PoisonPolicy;pub use resilience::ResiliencePolicy;pub use resilience::RetryClass;pub use resilience::RetryClassSet;pub use resilience::RetryMatcher;pub use resilience::RetryMetrics;pub use resilience::RetryPolicy;pub use resilience::WaitUnit;pub use resilience::classify;pub use resilience::execute_with_policy;pub use resilience::execute_with_policy_metered;pub use resilience::execute_with_policy_recorded;pub use retry::execute_with_retry;pub use rollback::DEFAULT_RUN_ID_COLUMN;pub use rollback::PREVIOUS_TABLE_SUFFIX;pub use rollback::RUN_JOURNAL_TABLE;pub use rollback::RollbackMode;pub use rollback::RollbackOptions;pub use rollback::RollbackOutcome;pub use rollback::RollbackWriteSpec;pub use shard::ShardSpec;pub use stage::CdcUnwrapSpec;transform-cdc-unwrappub use stage::CompiledUnpivot;transform-unpivotpub use stage::UnpivotSpec;transform-unpivotpub use stage::ExplodeSpec;transform-explodepub use stage::OnMissing;transform-explodepub use stage::FilterOp;transform-filterpub use stage::FilterSpec;transform-filterpub use stage::TransformStage;pub use stage::apply_stages;pub use stage::compile_stage;pub use staging::StagedFile;pub use staging::StagingCleanup;pub use staging::StagingCompression;pub use staging::StagingFormat;pub use staging::StagingLocation;pub use staging::StagingScheme;pub use staging::StagingSpec;pub use staging::serialize_records;pub use state::FileStateStore;pub use state::MemoryStateStore;pub use state::StateExport;pub use state::StateStore;pub use state_version::ResolvedState;pub use state_version::STATE_FORMAT;pub use state_version::StateCompat;pub use state_version::StoredState;pub use state_version::check_compat;pub use state_version::peel_versioned;pub use state_version::resolve_for_source;pub use state_version::wrap_versioned;pub use tls::TlsClientConfig;pub use topology::Edge;pub use topology::JoinNode;pub use topology::Node;pub use topology::NodeKind;pub use topology::Topology;pub use topology::TopologyBuilder;pub use topology::TopologyOnError;pub use topology::TopologyOptions;pub use topology::TopologyResult;pub use traits::RowOutcome;pub use traits::Sink;pub use traits::Source;pub use transform::JsonParseOnError;transform-json-parsepub use transform::LookupOnMissing;transform-lookuppub use transform::RecordTransform;pub use transform::ValueCaseMode;transform-value-casepub use transform::CastOnError;transform-castpub use transform::CastType;transform-castpub use transform::HashAlgorithm;transform-hashpub use transform::HashEncoding;transform-hashpub use transform::KeyCaseMode;transform-keys-casepub use transform::KeyCollision;transform-keys-casepub use transforming_source::TransformingSource;pub use tree::AncestorsSpec;transform-tree-flattenpub use tree::ColumnsSpec;transform-tree-flattenpub use tree::CompiledTreeFlatten;transform-tree-flattenpub use tree::TreeFlattenSpec;transform-tree-flattenpub use usage::CostSignal;pub use usage::UsageMeter;pub use usage::UsageSide;pub use usage::UsageSnapshot;pub use usage::estimate_json_bytes;pub use util::redact_uri_credentials;pub use verify::IntegrityCheck;pub use verify::LengthCheck;pub use verify::VerifyingReader;pub use window::WINDOW_END_PLACEHOLDER;pub use window::WINDOW_PLACEHOLDER;pub use window::WINDOW_START_PLACEHOLDER;pub use window::Window;pub use window::WindowBind;pub use window::WindowSpec;pub use window::enumerate_windows;pub use window::parse_step;pub use write_mode::DeleteMarker;pub use write_mode::KeyTuple;pub use write_mode::OverwriteScope;pub use write_mode::WriteMode;pub use write_mode::WritePlan;pub use write_mode::WriteSpec;pub use write_mode::key_to_doc_id;pub use write_mode::key_to_filter;pub use write_mode::plan_writes;pub use write_mode::sql_literal;pub use zip_columns::ColumnGroupSpec;transform-zip-columnspub use zip_columns::CompiledZipColumns;transform-zip-columnspub use zip_columns::ZipColumnsSpec;transform-zip-columnspub use compression::Compression;compressionpub use compression::CompressionConfig;compressionpub use compression::compress_buf;compressionpub use compression::warn_mismatch;compressionpub use quality::BatchCheck;qualitypub use quality::CheckTally;qualitypub use quality::CompareOp;qualitypub use quality::CompiledQuality;qualitypub use quality::JsonType;qualitypub use quality::OnFailure;qualitypub use quality::QualityOutcome;qualitypub use quality::QualitySpec;qualitypub use quality::QuarantinedRecord;qualitypub use quality::RecordCheck;qualitypub use quality::apply_quality;qualitypub use contract::CompiledContract;contractpub use contract::ContractFieldType;contractpub use contract::ContractOutcome;contractpub use contract::ContractSpec;contractpub use contract::ContractViolation;contractpub use contract::FieldContract;contractpub use contract::OnBreach;contractpub use contract::ViolatingRecord;contractpub use contract::apply_contract;contractpub use contract::to_json_schema;contractpub use contract::to_openlineage_facet;contractpub use masking::CompiledMasking;maskingpub use masking::Detector;maskingpub use masking::MaskAction;maskingpub use masking::MaskHit;maskingpub use masking::MaskRule;maskingpub use masking::MaskingOutcome;maskingpub use masking::MaskingSpec;maskingpub use masking::MatchSpec;maskingpub use masking::apply_masking;maskingpub use async_stream;pub use self;pub use self;pub use self;
Modules§
- adaptive
- Adaptive batch sizing — an AIMD controller that auto-tunes the effective
write batch size per pipeline row from observed sink latency + error rate.
Pure logic (no I/O);
run_streamfeeds it observations and emits metrics. Seedocs/superpowers/specs/2026-05-31-adaptive-batch-sizing-design.md. - anomaly
- Learned-baseline anomaly detection shared by the CLI’s SLA volume check
(#202) and column profiling (#708): a new observation is flagged against a
rolling baseline of earlier observations by z-score or Tukey IQR fences.
Pure math over
f64— no I/O, no config loading. - auth
- Shared, connector-agnostic authentication abstraction.
- budget
- Run budgets (#703): hard ceilings on what one invocation may move, so an approved change (or any run) cannot exceed what was agreed.
- check
- Preflight check types for
faucet doctor(#126). - cleanup
- Scoped cleanup — delete destination rows a run did not write (#478).
- columnar
arrow - Opt-in Apache Arrow columnar record path (feature
arrow, RFC 0002 / #375). - columnar_
transform arrow - Arrow-native (vectorized) forms of the built-in record transforms (#636).
- compression
compression - Transparent gzip / zstd compression wrappers for file-shaped connectors.
- config
- Configuration loading utilities.
- contract
contract - Data contracts (issue #204): a declarative, versioned schema + semantic
contract for a pipeline’s output — required fields, types, nullability,
enum sets, regex patterns, numeric/length bounds — enforced per page after
transforms and quality checks and before the sink write. Breaches
fail,quarantine(via the DLQ), orwarnper the contract-level policy. - create_
table - Shared auto-create-table planning for table-based sinks (#580).
- cross_
join transform-cross-join - Inbuilt
cross_jointransform (#534): expand one record into the cartesian product of two or more of its sibling array fields, emitting one flat record per combination. - diff
- Content verification primitives for
faucet verify(#701). - discover
- Live source introspection (
faucet discover, #211). - dlq
- Dead-letter queue (DLQ) wiring shared by the pipeline runner.
- drift
- Schema-drift detection + policy types (issue #194).
- encryption
encryption - Encryption at rest for faucet-managed local files (#207).
- error
- Error types for faucet-stream.
- file_
format - One file-format vocabulary for every file and object-store connector (#604).
- idempotency
- Exactly-once / idempotent delivery primitives.
- join
- Pure hash-join logic for the topology
joinnode (issue #72). - lag
- Source lag (#733): how far a CDC / streaming pipeline is behind its source.
- local_
outputs - Local sink output tracking — the provenance record a retention GC deletes from (#587).
- masking
masking - PII detection + column-level masking policies (issue #206): a declarative,
destination-scoped policy that classifies sensitive fields — by field-name
pattern, by value detector (email / card-with-Luhn / SSN / phone / IPv4), or
by explicit field list — and rewrites them per action (
redact/hash/tokenize/partial). - metadata
- Optional
_faucet_*run/lineage metadata columns (#510). - native
- Native byte-passthrough fast-load (#633) — the third capability-negotiated transfer mechanism, alongside the Arrow columnar path (RFC 0002 / #375) and staged bulk load (#528).
- object_
rollover - Cross-page accumulation and rollover for object-store sinks (#618).
- observability
- Pipeline-internal observability: tracing spans and
metricscounters/ histograms wired automatically around every source, sink, transform, and state-store operation. Seedocs/superpowers/specs/2026-05-23-observability-otel-prometheus-design.md. - pipeline
- Source-to-sink pipeline orchestration.
- policy
policy - Data-flow policies (#702): organisation-wide rules about where a labelled column may land, declared once and enforced on every pipeline.
- profiling
- Learned column profiles with drift detection (#708).
- quality
quality - Data-quality checks: declarative per-record and per-batch assertions that
quarantine violating records to the DLQ or abort the run. Pure evaluation;
the pipeline wires the DLQ routing in
run_stream. - redact
- A process-global redaction hook (#456 H5).
- replication
- Incremental replication support.
- resilience
- Unified resilience policy: retry, backoff classification, circuit breaker,
and poison-pill row handling. See
docs/superpowers/specs/2026-06-17-resilience-policy-design.md. - retry
- Shared exponential-backoff retry executor for HTTP-style sources.
- rollback
- Undo a run (#706): the shared types behind
faucet rollback. - schema
- JSON Schema inference from record samples.
- shard
- Source sharding for clustered (Mode B) execution.
- stage
- Pipeline-level transform stages. A
TransformStagewraps one of six shapes: - staging
- Staged bulk load (#528): stage a page to object storage, then have the
warehouse pull it with its native bulk-load command (
COPY/COPY INTO/ load-from-URI /s3()table function). - state
- Pluggable state store for incremental replication bookmarks.
- state_
version - Versioned pipeline state (#736): upgrade-safe bookmarks and CDC positions.
- tls
- Shared mutual-TLS (client-certificate) configuration for HTTP connectors.
- topology
- Multi-edge pipeline topology — fan-out (tee), fan-in (merge), and hash-join over an explicit node graph (issues #71 and #72).
- traits
- Shared traits for faucet sources and sinks.
- transform
- Record transformation pipeline.
- transforming_
source - Wrap any
Sourcewith a fixed list ofTransformStages applied to every emitted record. The canonical way for library callers to attach stages (transforms wrapped viaTransformStage::Map, plusFilter/Explode/Custom); the CLI uses this same type internally. - tree
transform-tree-flatten - Recursive report-tree / matrix flatten transform (
tree_flatten, #530). - usage
- Per-run usage metering (#704): what a run moved and what it asked its backends to do, so the CLI can account for cost per pipeline / row / dataset and enforce run budgets (#703).
- util
- Shared utilities used across faucet source and sink crates.
- verify
- Read-time integrity verification for streaming object bodies.
- window
- In-run datetime window slicing for forward incremental (#527).
- write_
mode - Unified write-mode types + planner shared by every upsert-capable sink.
- zip_
columns transform-zip-columns - Inbuilt
zip_columnstransform (#551): turn a columnar payload —{ columns: [{name}, …], rows: [[v0, v1, …], …] }— into one object per row, keyed by column name.
Macros§
- json
- Construct a
serde_json::Valuefrom a JSON literal. - schema_
for - Generates a
Schemafor the given type using default settings. The default settings currently conform to JSON Schema 2020-12, but this is liable to change in a future version of Schemars if support for other JSON Schema versions is added.
Structs§
- Cancellation
Token - Re-exported so callers of
Pipeline::with_cancel/RunStreamOptions::with_cancelcan name the token type without addingtokio-utilthemselves. A token which can be used to signal a cancellation request to one or more tasks.
Enums§
- Value
- Represents any valid JSON value.
Traits§
- Json
Schema - A type which can be described as a JSON Schema document.
- Stream
- A stream of values produced asynchronously.
Attribute Macros§
Derive Macros§
- Json
Schema - Derive macro for
JsonSchematrait.