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 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
143pub 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};
150pub 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};