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