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 auth::AuthProvider;pub use auth::AuthReference;pub use auth::AuthSpec;pub use auth::Credential;pub use auth::CredentialPlacement;pub use auth::RequestAuth;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 cross_join::CompiledCrossJoin;transform-cross-joinpub use cross_join::CrossJoinSpec;transform-cross-joinpub use cross_join::OnEmpty as CrossJoinOnEmpty;transform-cross-joinpub use discover::DatasetDescriptor;pub use discover::columns_to_schema;pub use discover::nullable_type;pub use discover::sql_type_to_json_schema;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::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 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 metadata::CompiledMetadata;pub use metadata::MetadataColumn;pub use metadata::MetadataColumnsSpec;pub use metadata::MetadataContext;pub use metadata::MetadataSink;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 replication::BindFormat;pub use replication::BindTarget;pub use replication::ReplicationBind;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 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::RetryMetrics;pub use resilience::RetryPolicy;pub use resilience::classify;pub use resilience::execute_with_policy;pub use resilience::execute_with_policy_metered;pub use retry::execute_with_retry;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::StateStore;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::KeyCaseMode;transform-keys-casepub 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 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 util::redact_uri_credentials;pub use verify::IntegrityCheck;pub use verify::LengthCheck;pub use verify::VerifyingReader;pub use window::WINDOW_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 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 futures_core;pub use schemars;pub use serde_json;
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. - auth
- Shared, connector-agnostic authentication abstraction.
- 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). - 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. - 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. - 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.
- idempotency
- Exactly-once / idempotent delivery primitives.
- join
- Pure hash-join logic for the topology
joinnode (issue #72). - 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). - 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.
- 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.
- 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.
- 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). - 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.