#![cfg_attr(docsrs, feature(doc_cfg))]
pub mod adaptive;
pub mod anomaly;
pub mod auth;
pub mod budget;
pub mod check;
pub mod cleanup;
#[cfg(feature = "arrow")]
pub mod columnar;
#[cfg(feature = "arrow")]
pub mod columnar_transform;
pub mod config;
#[cfg(feature = "contract")]
pub mod contract;
pub mod create_table;
#[cfg(feature = "transform-cross-join")]
pub mod cross_join;
pub mod diff;
pub mod discover;
pub mod dlq;
pub mod drift;
#[cfg(feature = "encryption")]
pub mod encryption;
pub mod error;
pub mod file_format;
pub mod idempotency;
pub mod join;
pub mod lag;
pub mod local_outputs;
#[cfg(feature = "masking")]
pub mod masking;
pub mod metadata;
pub mod native;
pub mod object_rollover;
pub mod observability;
pub mod pipeline;
#[cfg(feature = "policy")]
pub mod policy;
pub mod profiling;
#[cfg(feature = "quality")]
pub mod quality;
pub mod redact;
pub mod replication;
pub mod resilience;
pub mod retry;
pub mod rollback;
pub mod schema;
pub mod shard;
pub mod stage;
pub mod staging;
pub mod state;
pub mod state_version;
pub mod tls;
pub mod topology;
pub mod traits;
pub mod transform;
pub mod transforming_source;
#[cfg(feature = "transform-tree-flatten")]
pub mod tree;
pub mod usage;
pub mod util;
pub mod verify;
pub mod window;
pub mod write_mode;
#[cfg(feature = "transform-zip-columns")]
pub mod zip_columns;
#[cfg(feature = "compression")]
pub mod compression;
pub use adaptive::{
AdaptiveBatchConfig, AdjustDirection, AdjustReason, Adjustment, AimdController, Observation,
};
pub use anomaly::AnomalyMethod;
pub use auth::{
AuthProvider, AuthReference, AuthSpec, Credential, CredentialPlacement, RequestAuth,
SharedAuthProvider,
};
pub use budget::{BudgetKind, BudgetSink, BudgetSpec, BudgetState, BudgetTimer, BudgetVerdict};
pub use check::{CheckContext, CheckReport, Probe, ProbeStatus};
pub use cleanup::{CleanupMode, CleanupPolicy, DEFAULT_MAX_KEYS, SeenKeys};
#[cfg(feature = "arrow")]
pub use columnar::{
ColumnarPage, infer_arrow_schema, record_batch_to_values, values_to_record_batch,
values_to_record_batch_inferred,
};
pub use create_table::{
PlannedColumn, missing_target_error, plan_columns, plan_keyed_columns, render_column_defs,
render_columns, render_primary_key,
};
#[cfg(feature = "transform-cross-join")]
pub use cross_join::{CompiledCrossJoin, CrossJoinSpec, OnEmpty as CrossJoinOnEmpty};
pub use diff::{
ContentDigest, Difference, DifferenceKind, DigestAccumulator, KeyRange, Normalizer,
ServerDigest, VerifyReport, diff_rows, plan_ranges, row_hash,
};
pub use discover::{
DatasetDescriptor, attach_primary_keys, columns_to_schema, nullable_type,
sql_type_to_json_schema,
};
pub use dlq::{
BatchAtomicity, BatchOutcome, BatchOutcomeCounters, BatchOutcomes, DlqConfig, DlqReason,
DlqStats, EnvelopeError, OnBatchError, UnwrappedEnvelope, build_envelope, check_dlq_all_policy,
dlq_all_is_safe, dlq_all_refusal, unwrap_envelope,
};
pub use drift::{
ColumnChange, OnDrift, OnIncompatible, SchemaDiff, SchemaDriftPolicy, SchemaDriftSpec,
SchemaEvolution, SqlBaseType, adds_null, base_widened, json_schema_base_type,
};
#[cfg(feature = "encryption")]
pub use encryption::{CompiledEncryption, EncryptionAlgorithm, EncryptionSpec};
pub use error::FaucetError;
pub use file_format::parquet_io::ParquetReadOptions;
pub use file_format::{
AvroCodec, AvroOptions, ContainerDecoder, CsvOptions, CsvUnknownField, ExcelOptions,
FileFormat, FileInput, FormatOptions, OrcOptions, XmlOptions,
};
pub use idempotency::{
DeliveryGuarantee, DeliveryMode, EffectivelyOnceMechanism, GuaranteeInputs, ReplayGuarantee,
SinkGuarantee, derive_delivery_guarantee, format_token, format_token_with_bookmark,
parse_token, parse_token_parts, unwrap_state, wrap_state,
};
pub use join::{
HashJoin, JoinConfig, JoinMode, JoinStats, KeyNormalize, OnCollision, OnDuplicate, Projection,
};
pub use lag::{LagObserver, SourceLag};
pub use local_outputs::{LocalOutput, LocalOutputLog, probe_pre_existing};
pub use metadata::{
CompiledMetadata, MetadataColumn, MetadataColumnsSpec, MetadataContext, MetadataSink,
};
pub use native::{
CsvDialect, NativeBatch, NativeFormat, NativeLoadCapability, NativeLoadContext, NativePayload,
NativePlan, NativePlanInputs, plan_native_transfer,
};
pub use object_rollover::{CompletedObject, ObjectAccumulator, PageAccumulator};
#[cfg(feature = "contract")]
pub use observability::instrumented_apply_contract;
#[cfg(feature = "masking")]
pub use observability::instrumented_apply_masking;
#[cfg(feature = "quality")]
pub use observability::instrumented_apply_quality;
pub use observability::otel::{OtelConfig, OtelProtocol, OtelSignal, shutdown_otel};
pub use observability::{
DurationGuard, InstallError, InstallReport, InstrumentedSink, InstrumentedSource,
InstrumentedStateStore, Labels, ObservabilityConfig, PrometheusConfig, RunStreamOptions,
TracingConfig, install_observability, instrumented_apply_stages, register_build_info,
update_bookmark_lag,
};
pub use pipeline::{
DEFAULT_BATCH_SIZE, MAX_BATCH_SIZE, Pipeline, PipelineResult, StreamPage, run_stream,
validate_batch_size,
};
#[cfg(feature = "policy")]
pub use policy::{
ColumnFacts, CompiledPolicy, PolicyRule, PolicyScope, PolicySink, PolicySpec, RuntimeAction,
SinkFacts, Violation, ViolationKind,
};
pub use profiling::{
ColumnProfile, DriftMetric, OnProfileDrift, ProfileDrift, Profiler, ProfilingSink,
ProfilingSpec, RunProfile, detect_drift,
};
pub use replication::{
BindFormat, BindTarget, BindValueType, IncrementalFilter, OnMissingKey, ReplicationBind,
ReplicationKey, ReplicationMethod, format_bookmark, format_instant, json_gt, parse_instant,
set_body_pointer,
};
pub use resilience::{
BackoffFrom, BackoffKind, CircuitBreaker, CircuitBreakerConfig, PoisonAction, PoisonPolicy,
ResiliencePolicy, RetryClass, RetryClassSet, RetryMatcher, RetryMetrics, RetryPolicy, WaitUnit,
classify, execute_with_policy, execute_with_policy_metered, execute_with_policy_recorded,
};
pub use retry::execute_with_retry;
pub use rollback::{
DEFAULT_RUN_ID_COLUMN, PREVIOUS_TABLE_SUFFIX, RUN_JOURNAL_TABLE, RollbackMode, RollbackOptions,
RollbackOutcome, RollbackWriteSpec,
};
pub use shard::ShardSpec;
#[cfg(feature = "transform-cdc-unwrap")]
pub use stage::CdcUnwrapSpec;
#[cfg(feature = "transform-unpivot")]
pub use stage::{CompiledUnpivot, UnpivotSpec};
#[cfg(feature = "transform-explode")]
pub use stage::{ExplodeSpec, OnMissing};
#[cfg(feature = "transform-filter")]
pub use stage::{FilterOp, FilterSpec};
pub use stage::{TransformStage, apply_stages, compile_stage};
pub use staging::{
StagedFile, StagingCleanup, StagingCompression, StagingFormat, StagingLocation, StagingScheme,
StagingSpec, serialize_records,
};
pub use state::{FileStateStore, MemoryStateStore, StateExport, StateStore};
pub use state_version::{
ResolvedState, STATE_FORMAT, StateCompat, StoredState, check_compat, peel_versioned,
resolve_for_source, wrap_versioned,
};
pub use tls::TlsClientConfig;
pub use topology::{
Edge, JoinNode, Node, NodeKind, Topology, TopologyBuilder, TopologyOnError, TopologyOptions,
TopologyResult,
};
pub use traits::{RowOutcome, Sink, Source};
#[cfg(feature = "transform-json-parse")]
pub use transform::JsonParseOnError;
#[cfg(feature = "transform-lookup")]
pub use transform::LookupOnMissing;
pub use transform::RecordTransform;
#[cfg(feature = "transform-value-case")]
pub use transform::ValueCaseMode;
#[cfg(feature = "transform-cast")]
pub use transform::{CastOnError, CastType};
#[cfg(feature = "transform-hash")]
pub use transform::{HashAlgorithm, HashEncoding};
#[cfg(feature = "transform-keys-case")]
pub use transform::{KeyCaseMode, KeyCollision};
pub use transforming_source::TransformingSource;
#[cfg(feature = "transform-tree-flatten")]
pub use tree::{AncestorsSpec, ColumnsSpec, CompiledTreeFlatten, TreeFlattenSpec};
pub use usage::{CostSignal, UsageMeter, UsageSide, UsageSnapshot, estimate_json_bytes};
pub use util::redact_uri_credentials;
pub use verify::{IntegrityCheck, LengthCheck, VerifyingReader};
pub use window::{
WINDOW_END_PLACEHOLDER, WINDOW_PLACEHOLDER, WINDOW_START_PLACEHOLDER, Window, WindowBind,
WindowSpec, enumerate_windows, parse_step,
};
pub use write_mode::{
DeleteMarker, KeyTuple, OverwriteScope, WriteMode, WritePlan, WriteSpec, key_to_doc_id,
key_to_filter, plan_writes, sql_literal,
};
#[cfg(feature = "transform-zip-columns")]
pub use zip_columns::{ColumnGroupSpec, CompiledZipColumns, ZipColumnsSpec};
pub use async_stream;
pub use async_trait::async_trait;
pub use futures_core::{self, Stream};
pub use schemars::{self, JsonSchema, schema_for};
pub use serde_json::{self, Value, json};
pub use tokio_util::sync::CancellationToken;
#[cfg(feature = "compression")]
pub use compression::{Compression, CompressionConfig, compress_buf, warn_mismatch};
#[cfg(feature = "quality")]
pub use quality::{
BatchCheck, CheckTally, CompareOp, CompiledQuality, JsonType, OnFailure, QualityOutcome,
QualitySpec, QuarantinedRecord, RecordCheck, apply_quality,
};
#[cfg(feature = "contract")]
pub use contract::{
CompiledContract, ContractFieldType, ContractOutcome, ContractSpec, ContractViolation,
FieldContract, OnBreach, ViolatingRecord, apply_contract, to_json_schema, to_openlineage_facet,
};
#[cfg(feature = "masking")]
pub use masking::{
CompiledMasking, Detector, MaskAction, MaskHit, MaskRule, MaskingOutcome, MaskingSpec,
MatchSpec, apply_masking,
};