1#![cfg_attr(docsrs, feature(doc_cfg))]
2
3pub mod adaptive;
17pub mod anomaly;
18pub mod auth;
19pub mod budget;
20pub mod check;
21pub mod cleanup;
22#[cfg(feature = "arrow")]
23pub mod columnar;
24#[cfg(feature = "arrow")]
25pub mod columnar_transform;
26pub mod config;
27#[cfg(feature = "contract")]
28pub mod contract;
29pub mod create_table;
30#[cfg(feature = "transform-cross-join")]
31pub mod cross_join;
32pub mod diff;
33pub mod discover;
34pub mod dlq;
35pub mod drift;
36#[cfg(feature = "encryption")]
37pub mod encryption;
38pub mod error;
39pub mod file_format;
40pub mod idempotency;
41pub mod join;
42pub mod lag;
43pub mod local_outputs;
44#[cfg(feature = "masking")]
45pub mod masking;
46pub mod metadata;
47pub mod native;
48pub mod object_rollover;
49pub mod observability;
50pub mod pipeline;
51#[cfg(feature = "policy")]
52pub mod policy;
53pub mod profiling;
54#[cfg(feature = "quality")]
55pub mod quality;
56pub mod redact;
57pub mod replication;
58pub mod resilience;
59pub mod retry;
60pub mod rollback;
61pub mod schema;
62pub mod shard;
63pub mod stage;
64pub mod staging;
65pub mod state;
66pub mod state_version;
67pub mod tls;
68pub mod topology;
69pub mod traits;
70pub mod transform;
71pub mod transforming_source;
72#[cfg(feature = "transform-tree-flatten")]
73pub mod tree;
74pub mod usage;
75pub mod util;
76pub mod verify;
77pub mod window;
78pub mod write_mode;
79#[cfg(feature = "transform-zip-columns")]
80pub mod zip_columns;
81
82#[cfg(feature = "compression")]
83pub mod compression;
84
85pub use adaptive::{
86 AdaptiveBatchConfig, AdjustDirection, AdjustReason, Adjustment, AimdController, Observation,
87};
88pub use anomaly::AnomalyMethod;
89pub use auth::{
90 AuthProvider, AuthReference, AuthSpec, Credential, CredentialPlacement, RequestAuth,
91 SharedAuthProvider,
92};
93pub use budget::{BudgetKind, BudgetSink, BudgetSpec, BudgetState, BudgetTimer, BudgetVerdict};
94pub use check::{CheckContext, CheckReport, Probe, ProbeStatus};
95pub use cleanup::{CleanupMode, CleanupPolicy, DEFAULT_MAX_KEYS, SeenKeys};
96#[cfg(feature = "arrow")]
97pub use columnar::{
98 ColumnarPage, infer_arrow_schema, record_batch_to_values, values_to_record_batch,
99 values_to_record_batch_inferred,
100};
101pub use create_table::{
102 PlannedColumn, missing_target_error, plan_columns, plan_keyed_columns, render_column_defs,
103 render_columns, render_primary_key,
104};
105#[cfg(feature = "transform-cross-join")]
106pub use cross_join::{CompiledCrossJoin, CrossJoinSpec, OnEmpty as CrossJoinOnEmpty};
107pub use diff::{
108 ContentDigest, Difference, DifferenceKind, DigestAccumulator, KeyRange, Normalizer,
109 ServerDigest, VerifyReport, diff_rows, plan_ranges, row_hash,
110};
111pub use discover::{
112 DatasetDescriptor, attach_primary_keys, columns_to_schema, nullable_type,
113 sql_type_to_json_schema,
114};
115pub use dlq::{
116 BatchAtomicity, BatchOutcome, BatchOutcomeCounters, BatchOutcomes, DlqConfig, DlqReason,
117 DlqStats, EnvelopeError, OnBatchError, UnwrappedEnvelope, build_envelope, check_dlq_all_policy,
118 dlq_all_is_safe, dlq_all_refusal, unwrap_envelope,
119};
120pub use drift::{
121 ColumnChange, OnDrift, OnIncompatible, SchemaDiff, SchemaDriftPolicy, SchemaDriftSpec,
122 SchemaEvolution, SqlBaseType, adds_null, base_widened, json_schema_base_type,
123};
124#[cfg(feature = "encryption")]
125pub use encryption::{CompiledEncryption, EncryptionAlgorithm, EncryptionSpec};
126pub use error::FaucetError;
127pub use file_format::parquet_io::ParquetReadOptions;
128pub use file_format::{
129 AvroCodec, AvroOptions, ContainerDecoder, CsvOptions, CsvUnknownField, ExcelOptions,
130 FileFormat, FileInput, FormatOptions, OrcOptions, XmlOptions,
131};
132pub use idempotency::{
133 DeliveryGuarantee, DeliveryMode, EffectivelyOnceMechanism, GuaranteeInputs, ReplayGuarantee,
134 SinkGuarantee, derive_delivery_guarantee, format_token, format_token_with_bookmark,
135 parse_token, parse_token_parts, unwrap_state, wrap_state,
136};
137pub use join::{
138 HashJoin, JoinConfig, JoinMode, JoinStats, KeyNormalize, OnCollision, OnDuplicate, Projection,
139};
140pub use lag::{LagObserver, SourceLag};
141pub use local_outputs::{LocalOutput, LocalOutputLog, probe_pre_existing};
142pub use metadata::{
143 CompiledMetadata, MetadataColumn, MetadataColumnsSpec, MetadataContext, MetadataSink,
144};
145pub use native::{
146 CsvDialect, NativeBatch, NativeFormat, NativeLoadCapability, NativeLoadContext, NativePayload,
147 NativePlan, NativePlanInputs, plan_native_transfer,
148};
149pub use object_rollover::{CompletedObject, ObjectAccumulator, PageAccumulator};
150#[cfg(feature = "contract")]
151pub use observability::instrumented_apply_contract;
152#[cfg(feature = "masking")]
153pub use observability::instrumented_apply_masking;
154#[cfg(feature = "quality")]
155pub use observability::instrumented_apply_quality;
156pub use observability::otel::{OtelConfig, OtelProtocol, OtelSignal, shutdown_otel};
157pub use observability::{
158 DurationGuard, InstallError, InstallReport, InstrumentedSink, InstrumentedSource,
159 InstrumentedStateStore, Labels, ObservabilityConfig, PrometheusConfig, RunStreamOptions,
160 TracingConfig, install_observability, instrumented_apply_stages, register_build_info,
161 update_bookmark_lag,
162};
163pub use pipeline::{
164 DEFAULT_BATCH_SIZE, MAX_BATCH_SIZE, Pipeline, PipelineResult, StreamPage, run_stream,
165 validate_batch_size,
166};
167#[cfg(feature = "policy")]
168pub use policy::{
169 ColumnFacts, CompiledPolicy, PolicyRule, PolicyScope, PolicySink, PolicySpec, RuntimeAction,
170 SinkFacts, Violation, ViolationKind,
171};
172pub use profiling::{
173 ColumnProfile, DriftMetric, OnProfileDrift, ProfileDrift, Profiler, ProfilingSink,
174 ProfilingSpec, RunProfile, detect_drift,
175};
176pub use replication::{
177 BindFormat, BindTarget, BindValueType, IncrementalFilter, OnMissingKey, ReplicationBind,
178 ReplicationKey, ReplicationMethod, format_bookmark, format_instant, json_gt, parse_instant,
179 set_body_pointer,
180};
181pub use resilience::{
182 BackoffFrom, BackoffKind, CircuitBreaker, CircuitBreakerConfig, PoisonAction, PoisonPolicy,
183 ResiliencePolicy, RetryClass, RetryClassSet, RetryMatcher, RetryMetrics, RetryPolicy, WaitUnit,
184 classify, execute_with_policy, execute_with_policy_metered, execute_with_policy_recorded,
185};
186pub use retry::execute_with_retry;
187pub use rollback::{
188 DEFAULT_RUN_ID_COLUMN, PREVIOUS_TABLE_SUFFIX, RUN_JOURNAL_TABLE, RollbackMode, RollbackOptions,
189 RollbackOutcome, RollbackWriteSpec,
190};
191pub use shard::ShardSpec;
192#[cfg(feature = "transform-cdc-unwrap")]
193pub use stage::CdcUnwrapSpec;
194#[cfg(feature = "transform-unpivot")]
195pub use stage::{CompiledUnpivot, UnpivotSpec};
196#[cfg(feature = "transform-explode")]
197pub use stage::{ExplodeSpec, OnMissing};
198#[cfg(feature = "transform-filter")]
199pub use stage::{FilterOp, FilterSpec};
200pub use stage::{TransformStage, apply_stages, compile_stage};
201pub use staging::{
202 StagedFile, StagingCleanup, StagingCompression, StagingFormat, StagingLocation, StagingScheme,
203 StagingSpec, serialize_records,
204};
205pub use state::{FileStateStore, MemoryStateStore, StateExport, StateStore};
206pub use state_version::{
207 ResolvedState, STATE_FORMAT, StateCompat, StoredState, check_compat, peel_versioned,
208 resolve_for_source, wrap_versioned,
209};
210pub use tls::TlsClientConfig;
211pub use topology::{
212 Edge, JoinNode, Node, NodeKind, Topology, TopologyBuilder, TopologyOnError, TopologyOptions,
213 TopologyResult,
214};
215pub use traits::{RowOutcome, Sink, Source};
216#[cfg(feature = "transform-json-parse")]
217pub use transform::JsonParseOnError;
218#[cfg(feature = "transform-lookup")]
219pub use transform::LookupOnMissing;
220pub use transform::RecordTransform;
221#[cfg(feature = "transform-value-case")]
222pub use transform::ValueCaseMode;
223#[cfg(feature = "transform-cast")]
224pub use transform::{CastOnError, CastType};
225#[cfg(feature = "transform-hash")]
226pub use transform::{HashAlgorithm, HashEncoding};
227#[cfg(feature = "transform-keys-case")]
228pub use transform::{KeyCaseMode, KeyCollision};
229pub use transforming_source::TransformingSource;
230#[cfg(feature = "transform-tree-flatten")]
231pub use tree::{AncestorsSpec, ColumnsSpec, CompiledTreeFlatten, TreeFlattenSpec};
232pub use usage::{CostSignal, UsageMeter, UsageSide, UsageSnapshot, estimate_json_bytes};
233pub use util::redact_uri_credentials;
234pub use verify::{IntegrityCheck, LengthCheck, VerifyingReader};
235pub use window::{
236 WINDOW_END_PLACEHOLDER, WINDOW_PLACEHOLDER, WINDOW_START_PLACEHOLDER, Window, WindowBind,
237 WindowSpec, enumerate_windows, parse_step,
238};
239pub use write_mode::{
240 DeleteMarker, KeyTuple, OverwriteScope, WriteMode, WritePlan, WriteSpec, key_to_doc_id,
241 key_to_filter, plan_writes, sql_literal,
242};
243#[cfg(feature = "transform-zip-columns")]
244pub use zip_columns::{ColumnGroupSpec, CompiledZipColumns, ZipColumnsSpec};
245
246pub use async_stream;
249pub use async_trait::async_trait;
250pub use futures_core::{self, Stream};
251pub use schemars::{self, JsonSchema, schema_for};
252pub use serde_json::{self, Value, json};
253pub use tokio_util::sync::CancellationToken;
257
258#[cfg(feature = "compression")]
259pub use compression::{Compression, CompressionConfig, compress_buf, warn_mismatch};
260
261#[cfg(feature = "quality")]
262pub use quality::{
263 BatchCheck, CheckTally, CompareOp, CompiledQuality, JsonType, OnFailure, QualityOutcome,
264 QualitySpec, QuarantinedRecord, RecordCheck, apply_quality,
265};
266
267#[cfg(feature = "contract")]
268pub use contract::{
269 CompiledContract, ContractFieldType, ContractOutcome, ContractSpec, ContractViolation,
270 FieldContract, OnBreach, ViolatingRecord, apply_contract, to_json_schema, to_openlineage_facet,
271};
272
273#[cfg(feature = "masking")]
274pub use masking::{
275 CompiledMasking, Detector, MaskAction, MaskHit, MaskRule, MaskingOutcome, MaskingSpec,
276 MatchSpec, apply_masking,
277};