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;
19#[cfg(feature = "arrow")]
20pub mod columnar;
21pub mod config;
22#[cfg(feature = "contract")]
23pub mod contract;
24pub mod discover;
25pub mod dlq;
26pub mod drift;
27#[cfg(feature = "encryption")]
28pub mod encryption;
29pub mod error;
30pub mod idempotency;
31pub mod join;
32#[cfg(feature = "masking")]
33pub mod masking;
34pub mod observability;
35pub mod pipeline;
36#[cfg(feature = "quality")]
37pub mod quality;
38pub mod replication;
39pub mod resilience;
40pub mod retry;
41pub mod schema;
42pub mod shard;
43pub mod stage;
44pub mod state;
45pub mod topology;
46pub mod traits;
47pub mod transform;
48pub mod transforming_source;
49pub mod util;
50pub mod verify;
51pub mod write_mode;
52
53#[cfg(feature = "compression")]
54pub mod compression;
55
56pub use adaptive::{
57    AdaptiveBatchConfig, AdjustDirection, AdjustReason, Adjustment, AimdController, Observation,
58};
59pub use auth::{AuthProvider, AuthReference, AuthSpec, Credential, SharedAuthProvider};
60pub use check::{CheckContext, CheckReport, Probe, ProbeStatus};
61#[cfg(feature = "arrow")]
62pub use columnar::{
63    ColumnarPage, infer_arrow_schema, record_batch_to_values, values_to_record_batch,
64    values_to_record_batch_inferred,
65};
66pub use discover::{DatasetDescriptor, columns_to_schema, nullable_type, sql_type_to_json_schema};
67pub use dlq::{
68    DlqConfig, DlqReason, DlqStats, EnvelopeError, OnBatchError, UnwrappedEnvelope, build_envelope,
69    unwrap_envelope,
70};
71pub use drift::{
72    ColumnChange, OnDrift, OnIncompatible, SchemaDiff, SchemaDriftPolicy, SchemaDriftSpec,
73    SchemaEvolution, SqlBaseType, adds_null, base_widened, json_schema_base_type,
74};
75#[cfg(feature = "encryption")]
76pub use encryption::{CompiledEncryption, EncryptionAlgorithm, EncryptionSpec};
77pub use error::FaucetError;
78pub use idempotency::{
79    DeliveryGuarantee, DeliveryMode, EffectivelyOnceMechanism, GuaranteeInputs, ReplayGuarantee,
80    SinkGuarantee, derive_delivery_guarantee, format_token, format_token_with_bookmark,
81    parse_token, parse_token_parts, unwrap_state, wrap_state,
82};
83pub use join::{
84    HashJoin, JoinConfig, JoinMode, JoinStats, KeyNormalize, OnCollision, OnDuplicate, Projection,
85};
86#[cfg(feature = "contract")]
87pub use observability::instrumented_apply_contract;
88#[cfg(feature = "masking")]
89pub use observability::instrumented_apply_masking;
90#[cfg(feature = "quality")]
91pub use observability::instrumented_apply_quality;
92pub use observability::otel::{OtelConfig, OtelProtocol, OtelSignal, shutdown_otel};
93pub use observability::{
94    DurationGuard, InstallError, InstallReport, InstrumentedSink, InstrumentedSource,
95    InstrumentedStateStore, Labels, ObservabilityConfig, PrometheusConfig, RunStreamOptions,
96    TracingConfig, install_observability, instrumented_apply_stages, register_build_info,
97    update_bookmark_lag,
98};
99pub use pipeline::{
100    DEFAULT_BATCH_SIZE, MAX_BATCH_SIZE, Pipeline, PipelineResult, StreamPage, run_stream,
101    validate_batch_size,
102};
103pub use replication::{ReplicationMethod, json_gt};
104pub use resilience::{
105    BackoffKind, CircuitBreaker, CircuitBreakerConfig, PoisonAction, PoisonPolicy,
106    ResiliencePolicy, RetryClass, RetryClassSet, RetryMetrics, RetryPolicy, classify,
107    execute_with_policy, execute_with_policy_metered,
108};
109pub use retry::execute_with_retry;
110pub use shard::ShardSpec;
111#[cfg(feature = "transform-cdc-unwrap")]
112pub use stage::CdcUnwrapSpec;
113#[cfg(feature = "transform-explode")]
114pub use stage::{ExplodeSpec, OnMissing};
115#[cfg(feature = "transform-filter")]
116pub use stage::{FilterOp, FilterSpec};
117pub use stage::{TransformStage, compile_stage};
118pub use state::{FileStateStore, MemoryStateStore, StateStore};
119pub use topology::{
120    Edge, JoinNode, Node, NodeKind, Topology, TopologyBuilder, TopologyOnError, TopologyOptions,
121    TopologyResult,
122};
123pub use traits::{RowOutcome, Sink, Source};
124#[cfg(feature = "transform-json-parse")]
125pub use transform::JsonParseOnError;
126#[cfg(feature = "transform-keys-case")]
127pub use transform::KeyCaseMode;
128pub use transform::RecordTransform;
129#[cfg(feature = "transform-value-case")]
130pub use transform::ValueCaseMode;
131#[cfg(feature = "transform-cast")]
132pub use transform::{CastOnError, CastType};
133#[cfg(feature = "transform-hash")]
134pub use transform::{HashAlgorithm, HashEncoding};
135pub use transforming_source::TransformingSource;
136pub use util::redact_uri_credentials;
137pub use verify::{IntegrityCheck, LengthCheck, VerifyingReader};
138pub use write_mode::{
139    DeleteMarker, KeyTuple, WriteMode, WritePlan, WriteSpec, key_to_doc_id, key_to_filter,
140    plan_writes,
141};
142
143// Re-export dependencies that connector authors need, so they only depend on
144// `faucet-core` instead of adding `async-trait` and `serde_json` themselves.
145pub use async_stream;
146pub use async_trait::async_trait;
147pub use futures_core::{self, Stream};
148pub use schemars::{self, JsonSchema, schema_for};
149pub use serde_json::{self, Value, json};
150/// Re-exported so callers of [`Pipeline::with_cancel`](pipeline::Pipeline::with_cancel)
151/// / [`RunStreamOptions::with_cancel`] can name the token type without adding
152/// `tokio-util` themselves.
153pub use tokio_util::sync::CancellationToken;
154
155#[cfg(feature = "compression")]
156pub use compression::{Compression, CompressionConfig, compress_buf, warn_mismatch};
157
158#[cfg(feature = "quality")]
159pub use quality::{
160    BatchCheck, CheckTally, CompareOp, CompiledQuality, JsonType, OnFailure, QualityOutcome,
161    QualitySpec, QuarantinedRecord, RecordCheck, apply_quality,
162};
163
164#[cfg(feature = "contract")]
165pub use contract::{
166    CompiledContract, ContractFieldType, ContractOutcome, ContractSpec, ContractViolation,
167    FieldContract, OnBreach, ViolatingRecord, apply_contract, to_json_schema, to_openlineage_facet,
168};
169
170#[cfg(feature = "masking")]
171pub use masking::{
172    CompiledMasking, Detector, MaskAction, MaskHit, MaskRule, MaskingOutcome, MaskingSpec,
173    MatchSpec, apply_masking,
174};