Skip to main content

faucet_core/
lib.rs

1#![cfg_attr(docsrs, feature(doc_cfg))]
2
3//! # faucet-core
4//!
5//! Shared types, traits, and utilities for the faucet-stream ecosystem.
6//!
7//! This crate provides the common foundation used by all faucet source and
8//! sink connectors:
9//!
10//! - [`FaucetError`] — unified error type
11//! - [`Source`] / [`Sink`] — async traits for data connectors
12//! - [`RecordTransform`] — record transformation pipeline
13//! - [`ReplicationMethod`] — incremental replication support
14//! - [`schema::infer_schema`] — JSON Schema inference from record samples
15
16pub mod adaptive;
17pub mod auth;
18pub mod check;
19pub mod cleanup;
20#[cfg(feature = "arrow")]
21pub mod columnar;
22pub mod config;
23#[cfg(feature = "contract")]
24pub mod contract;
25#[cfg(feature = "transform-cross-join")]
26pub mod cross_join;
27pub mod discover;
28pub mod dlq;
29pub mod drift;
30#[cfg(feature = "encryption")]
31pub mod encryption;
32pub mod error;
33pub mod idempotency;
34pub mod join;
35#[cfg(feature = "masking")]
36pub mod masking;
37pub mod metadata;
38pub mod observability;
39pub mod pipeline;
40#[cfg(feature = "quality")]
41pub mod quality;
42pub mod redact;
43pub mod replication;
44pub mod resilience;
45pub mod retry;
46pub mod schema;
47pub mod shard;
48pub mod stage;
49pub mod staging;
50pub mod state;
51pub mod tls;
52pub mod topology;
53pub mod traits;
54pub mod transform;
55pub mod transforming_source;
56#[cfg(feature = "transform-tree-flatten")]
57pub mod tree;
58pub mod util;
59pub mod verify;
60pub mod window;
61pub mod write_mode;
62#[cfg(feature = "transform-zip-columns")]
63pub mod zip_columns;
64
65#[cfg(feature = "compression")]
66pub mod compression;
67
68pub use adaptive::{
69    AdaptiveBatchConfig, AdjustDirection, AdjustReason, Adjustment, AimdController, Observation,
70};
71pub use auth::{
72    AuthProvider, AuthReference, AuthSpec, Credential, CredentialPlacement, RequestAuth,
73    SharedAuthProvider,
74};
75pub use check::{CheckContext, CheckReport, Probe, ProbeStatus};
76pub use cleanup::{CleanupMode, CleanupPolicy, DEFAULT_MAX_KEYS, SeenKeys};
77#[cfg(feature = "arrow")]
78pub use columnar::{
79    ColumnarPage, infer_arrow_schema, record_batch_to_values, values_to_record_batch,
80    values_to_record_batch_inferred,
81};
82#[cfg(feature = "transform-cross-join")]
83pub use cross_join::{CompiledCrossJoin, CrossJoinSpec, OnEmpty as CrossJoinOnEmpty};
84pub use discover::{DatasetDescriptor, columns_to_schema, nullable_type, sql_type_to_json_schema};
85pub use dlq::{
86    DlqConfig, DlqReason, DlqStats, EnvelopeError, OnBatchError, UnwrappedEnvelope, build_envelope,
87    unwrap_envelope,
88};
89pub use drift::{
90    ColumnChange, OnDrift, OnIncompatible, SchemaDiff, SchemaDriftPolicy, SchemaDriftSpec,
91    SchemaEvolution, SqlBaseType, adds_null, base_widened, json_schema_base_type,
92};
93#[cfg(feature = "encryption")]
94pub use encryption::{CompiledEncryption, EncryptionAlgorithm, EncryptionSpec};
95pub use error::FaucetError;
96pub use idempotency::{
97    DeliveryGuarantee, DeliveryMode, EffectivelyOnceMechanism, GuaranteeInputs, ReplayGuarantee,
98    SinkGuarantee, derive_delivery_guarantee, format_token, format_token_with_bookmark,
99    parse_token, parse_token_parts, unwrap_state, wrap_state,
100};
101pub use join::{
102    HashJoin, JoinConfig, JoinMode, JoinStats, KeyNormalize, OnCollision, OnDuplicate, Projection,
103};
104pub use metadata::{
105    CompiledMetadata, MetadataColumn, MetadataColumnsSpec, MetadataContext, MetadataSink,
106};
107#[cfg(feature = "contract")]
108pub use observability::instrumented_apply_contract;
109#[cfg(feature = "masking")]
110pub use observability::instrumented_apply_masking;
111#[cfg(feature = "quality")]
112pub use observability::instrumented_apply_quality;
113pub use observability::otel::{OtelConfig, OtelProtocol, OtelSignal, shutdown_otel};
114pub use observability::{
115    DurationGuard, InstallError, InstallReport, InstrumentedSink, InstrumentedSource,
116    InstrumentedStateStore, Labels, ObservabilityConfig, PrometheusConfig, RunStreamOptions,
117    TracingConfig, install_observability, instrumented_apply_stages, register_build_info,
118    update_bookmark_lag,
119};
120pub use pipeline::{
121    DEFAULT_BATCH_SIZE, MAX_BATCH_SIZE, Pipeline, PipelineResult, StreamPage, run_stream,
122    validate_batch_size,
123};
124pub use replication::{
125    BindFormat, BindTarget, ReplicationBind, ReplicationMethod, format_bookmark, format_instant,
126    json_gt, parse_instant,
127};
128pub use resilience::{
129    BackoffKind, CircuitBreaker, CircuitBreakerConfig, PoisonAction, PoisonPolicy,
130    ResiliencePolicy, RetryClass, RetryClassSet, RetryMetrics, RetryPolicy, classify,
131    execute_with_policy, execute_with_policy_metered,
132};
133pub use retry::execute_with_retry;
134pub use shard::ShardSpec;
135#[cfg(feature = "transform-cdc-unwrap")]
136pub use stage::CdcUnwrapSpec;
137#[cfg(feature = "transform-unpivot")]
138pub use stage::{CompiledUnpivot, UnpivotSpec};
139#[cfg(feature = "transform-explode")]
140pub use stage::{ExplodeSpec, OnMissing};
141#[cfg(feature = "transform-filter")]
142pub use stage::{FilterOp, FilterSpec};
143pub use stage::{TransformStage, apply_stages, compile_stage};
144pub use staging::{
145    StagedFile, StagingCleanup, StagingCompression, StagingFormat, StagingLocation, StagingScheme,
146    StagingSpec, serialize_records,
147};
148pub use state::{FileStateStore, MemoryStateStore, StateStore};
149pub use tls::TlsClientConfig;
150pub use topology::{
151    Edge, JoinNode, Node, NodeKind, Topology, TopologyBuilder, TopologyOnError, TopologyOptions,
152    TopologyResult,
153};
154pub use traits::{RowOutcome, Sink, Source};
155#[cfg(feature = "transform-json-parse")]
156pub use transform::JsonParseOnError;
157#[cfg(feature = "transform-keys-case")]
158pub use transform::KeyCaseMode;
159#[cfg(feature = "transform-lookup")]
160pub use transform::LookupOnMissing;
161pub use transform::RecordTransform;
162#[cfg(feature = "transform-value-case")]
163pub use transform::ValueCaseMode;
164#[cfg(feature = "transform-cast")]
165pub use transform::{CastOnError, CastType};
166#[cfg(feature = "transform-hash")]
167pub use transform::{HashAlgorithm, HashEncoding};
168pub use transforming_source::TransformingSource;
169#[cfg(feature = "transform-tree-flatten")]
170pub use tree::{AncestorsSpec, ColumnsSpec, CompiledTreeFlatten, TreeFlattenSpec};
171pub use util::redact_uri_credentials;
172pub use verify::{IntegrityCheck, LengthCheck, VerifyingReader};
173pub use window::{
174    WINDOW_PLACEHOLDER, Window, WindowBind, WindowSpec, enumerate_windows, parse_step,
175};
176pub use write_mode::{
177    DeleteMarker, KeyTuple, OverwriteScope, WriteMode, WritePlan, WriteSpec, key_to_doc_id,
178    key_to_filter, plan_writes,
179};
180#[cfg(feature = "transform-zip-columns")]
181pub use zip_columns::{CompiledZipColumns, ZipColumnsSpec};
182
183// Re-export dependencies that connector authors need, so they only depend on
184// `faucet-core` instead of adding `async-trait` and `serde_json` themselves.
185pub use async_stream;
186pub use async_trait::async_trait;
187pub use futures_core::{self, Stream};
188pub use schemars::{self, JsonSchema, schema_for};
189pub use serde_json::{self, Value, json};
190/// Re-exported so callers of [`Pipeline::with_cancel`](pipeline::Pipeline::with_cancel)
191/// / [`RunStreamOptions::with_cancel`] can name the token type without adding
192/// `tokio-util` themselves.
193pub use tokio_util::sync::CancellationToken;
194
195#[cfg(feature = "compression")]
196pub use compression::{Compression, CompressionConfig, compress_buf, warn_mismatch};
197
198#[cfg(feature = "quality")]
199pub use quality::{
200    BatchCheck, CheckTally, CompareOp, CompiledQuality, JsonType, OnFailure, QualityOutcome,
201    QualitySpec, QuarantinedRecord, RecordCheck, apply_quality,
202};
203
204#[cfg(feature = "contract")]
205pub use contract::{
206    CompiledContract, ContractFieldType, ContractOutcome, ContractSpec, ContractViolation,
207    FieldContract, OnBreach, ViolatingRecord, apply_contract, to_json_schema, to_openlineage_facet,
208};
209
210#[cfg(feature = "masking")]
211pub use masking::{
212    CompiledMasking, Detector, MaskAction, MaskHit, MaskRule, MaskingOutcome, MaskingSpec,
213    MatchSpec, apply_masking,
214};