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