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