Skip to main content

faucet_stream/
lib.rs

1#![cfg_attr(docsrs, feature(doc_cfg))]
2
3//! # faucet-stream
4//!
5//! A declarative, config-driven data pipeline with pluggable source and sink
6//! connectors.
7//!
8//! 📖 **Guide, tutorials & cookbook:** <https://pawansikawat.github.io/faucet-stream/>
9//!
10//! ## Feature flags
11//!
12//! | Feature | Description |
13//! |---------|-------------|
14//! | `source-rest` *(default)* | REST API source with pagination, auth, transforms |
15//! | `source-graphql` | GraphQL API source with cursor pagination |
16//! | `source-xml` | XML/SOAP API source with XML-to-JSON conversion |
17//! | `source-grpc` | gRPC source with dynamic protobuf messages |
18//! | `source-postgres` | PostgreSQL query source |
19//! | `source-postgres-cdc` | PostgreSQL CDC source (logical replication) |
20//! | `source-mysql` | MySQL query source |
21//! | `source-mssql` | Microsoft SQL Server query source |
22//! | `source-sqlite` | SQLite query source |
23
24//! | `source-s3` | AWS S3 file source |
25//! | `source-mongodb` | MongoDB query source |
26//! | `source-mongodb-cdc` | MongoDB CDC source (Change Streams) |
27//! | `source-mysql-cdc` | MySQL CDC source (binlog replication) |
28//! | `source-redis` | Redis source (streams, lists, keys) |
29//! | `source-webhook` | Webhook HTTP receiver source |
30//! | `source-websocket` | WebSocket streaming source |
31//! | `source-csv` | CSV file source |
32//! | `source-elasticsearch` | Elasticsearch search/scroll source |
33//! | `source-kafka` | Apache Kafka consumer source |
34//! | `source-kinesis` | AWS Kinesis Data Streams source |
35//! | `source-spanner` | Google Cloud Spanner query source |
36//! | `source-parquet` | Apache Parquet file source (local, glob, S3) |
37//! | `source-delta` | Apache Delta Lake source (local FS or S3/Azure/GCS, time travel) |
38//! | `source-databricks` | Databricks SQL query source (Statement Execution API) |
39//! | `sink-bigquery` | Google BigQuery streaming insert sink |
40//! | `sink-iceberg` | Apache Iceberg sink (append-only, REST/Glue/SQL/HMS catalogs) |
41//! | `sink-postgres` | PostgreSQL sink (jsonb or auto-mapped columns) |
42//! | `sink-jsonl` | JSON Lines file sink |
43//! | `sink-snowflake` | Snowflake SQL REST API sink |
44//! | `sink-mysql` | MySQL sink |
45//! | `sink-mssql` | Microsoft SQL Server sink |
46//! | `sink-sqlite` | SQLite sink |
47
48//! | `sink-s3` | AWS S3 file sink |
49//! | `sink-mongodb` | MongoDB insert sink |
50//! | `sink-redis` | Redis sink (streams, lists, key-value) |
51//! | `sink-csv` | CSV file sink |
52//! | `sink-elasticsearch` | Elasticsearch bulk index sink |
53//! | `sink-http` | HTTP POST sink |
54//! | `sink-kafka` | Apache Kafka producer sink |
55//! | `sink-kinesis` | AWS Kinesis Data Streams sink |
56//! | `sink-spanner` | Google Cloud Spanner mutation sink |
57//! | `sink-parquet` | Apache Parquet file sink (local, S3) |
58//! | `encryption` | AES-256-GCM at-rest sealing for file state-store bookmarks and per-line JSONL/DLQ output |
59//! | `sink-delta` | Apache Delta Lake sink (append-only; local FS or S3/Azure/GCS) |
60//! | `kafka-schema-registry` | Schema Registry support for Kafka connectors |
61//! | `source` | All source connectors |
62//! | `sink` | All sink connectors |
63//! | `full` | Every connector |
64
65// Always re-export core types and traits.
66pub use faucet_core::*;
67
68// Explicit re-exports for the library-side transforms wrapper and observability
69// labels so users can import them via the umbrella path.
70pub use faucet_core::TransformingSource;
71pub use faucet_core::observability::Labels;
72
73// ── Shared auth providers ────────────────────────────────────────────────────
74/// Single-flight OAuth2 / token-endpoint auth providers (enable the `auth`
75/// feature). Share one across connectors via `with_auth_provider` or the CLI
76/// `auth: { ref }` catalog.
77#[cfg(feature = "auth")]
78pub mod auth {
79    pub use faucet_auth::*;
80}
81
82// ── Source connectors ────────────────────────────────────────────────────────
83
84#[cfg(feature = "source-rest")]
85pub mod source {
86    pub mod rest {
87        pub use faucet_source_rest::*;
88    }
89
90    #[cfg(feature = "source-graphql")]
91    pub mod graphql {
92        pub use faucet_source_graphql::*;
93    }
94
95    #[cfg(feature = "source-xml")]
96    pub mod xml {
97        pub use faucet_source_xml::*;
98    }
99
100    #[cfg(feature = "source-grpc")]
101    pub mod grpc {
102        pub use faucet_source_grpc::*;
103    }
104
105    #[cfg(feature = "source-postgres")]
106    pub mod postgres {
107        pub use faucet_source_postgres::*;
108    }
109
110    #[cfg(feature = "source-postgres-cdc")]
111    pub mod postgres_cdc {
112        pub use faucet_source_postgres_cdc::*;
113    }
114
115    #[cfg(feature = "source-mysql")]
116    pub mod mysql {
117        pub use faucet_source_mysql::*;
118    }
119
120    #[cfg(feature = "source-mssql")]
121    pub mod mssql {
122        pub use faucet_source_mssql::*;
123    }
124
125    #[cfg(feature = "source-sqlite")]
126    pub mod sqlite {
127        pub use faucet_source_sqlite::*;
128    }
129
130    #[cfg(feature = "source-s3")]
131    pub mod s3 {
132        pub use faucet_source_s3::*;
133    }
134
135    #[cfg(feature = "source-mongodb")]
136    pub mod mongodb {
137        pub use faucet_source_mongodb::*;
138    }
139
140    #[cfg(feature = "source-mongodb-cdc")]
141    pub mod mongodb_cdc {
142        pub use faucet_source_mongodb_cdc::*;
143    }
144
145    #[cfg(feature = "source-mysql-cdc")]
146    pub mod mysql_cdc {
147        pub use faucet_source_mysql_cdc::*;
148    }
149
150    #[cfg(feature = "source-redis")]
151    pub mod redis {
152        pub use faucet_source_redis::*;
153    }
154
155    #[cfg(feature = "source-webhook")]
156    pub mod webhook {
157        pub use faucet_source_webhook::*;
158    }
159
160    #[cfg(feature = "source-websocket")]
161    pub mod websocket {
162        pub use faucet_source_websocket::*;
163    }
164
165    #[cfg(feature = "source-csv")]
166    pub mod csv {
167        pub use faucet_source_csv::*;
168    }
169
170    #[cfg(feature = "source-elasticsearch")]
171    pub mod elasticsearch {
172        pub use faucet_source_elasticsearch::*;
173    }
174
175    #[cfg(feature = "source-kafka")]
176    pub mod kafka {
177        pub use faucet_source_kafka::*;
178    }
179
180    #[cfg(feature = "source-kinesis")]
181    pub mod kinesis {
182        pub use faucet_source_kinesis::*;
183    }
184
185    #[cfg(feature = "source-bigquery")]
186    pub mod bigquery {
187        pub use faucet_source_bigquery::*;
188    }
189
190    #[cfg(feature = "source-snowflake")]
191    pub mod snowflake {
192        pub use faucet_source_snowflake::*;
193    }
194
195    #[cfg(feature = "source-spanner")]
196    pub mod spanner {
197        pub use faucet_source_spanner::*;
198    }
199
200    #[cfg(feature = "source-parquet")]
201    pub mod parquet {
202        pub use faucet_source_parquet::*;
203    }
204
205    #[cfg(feature = "source-delta")]
206    pub mod delta {
207        pub use faucet_source_delta::*;
208    }
209
210    #[cfg(feature = "source-databricks")]
211    pub mod databricks {
212        pub use faucet_source_databricks::*;
213    }
214
215    #[cfg(feature = "source-gcs")]
216    pub mod gcs {
217        pub use faucet_source_gcs::*;
218    }
219
220    #[cfg(feature = "source-mssql-cdc")]
221    pub mod mssql_cdc {
222        pub use faucet_source_mssql_cdc::*;
223    }
224
225    #[cfg(feature = "source-redshift")]
226    pub mod redshift {
227        pub use faucet_source_redshift::*;
228    }
229
230    #[cfg(feature = "source-pubsub")]
231    pub mod pubsub {
232        pub use faucet_source_pubsub::*;
233    }
234
235    #[cfg(feature = "source-clickhouse")]
236    pub mod clickhouse {
237        pub use faucet_source_clickhouse::*;
238    }
239
240    #[cfg(feature = "source-azure-blob")]
241    pub mod azure_blob {
242        pub use faucet_source_azure_blob::*;
243    }
244
245    #[cfg(feature = "source-singer")]
246    pub mod singer {
247        pub use faucet_source_singer::*;
248    }
249}
250
251// Source modules available without source-rest (when only other sources are enabled).
252#[cfg(not(feature = "source-rest"))]
253pub mod source {
254    #[cfg(feature = "source-graphql")]
255    pub mod graphql {
256        pub use faucet_source_graphql::*;
257    }
258
259    #[cfg(feature = "source-xml")]
260    pub mod xml {
261        pub use faucet_source_xml::*;
262    }
263
264    #[cfg(feature = "source-grpc")]
265    pub mod grpc {
266        pub use faucet_source_grpc::*;
267    }
268
269    #[cfg(feature = "source-postgres")]
270    pub mod postgres {
271        pub use faucet_source_postgres::*;
272    }
273
274    #[cfg(feature = "source-postgres-cdc")]
275    pub mod postgres_cdc {
276        pub use faucet_source_postgres_cdc::*;
277    }
278
279    #[cfg(feature = "source-mysql")]
280    pub mod mysql {
281        pub use faucet_source_mysql::*;
282    }
283
284    #[cfg(feature = "source-mssql")]
285    pub mod mssql {
286        pub use faucet_source_mssql::*;
287    }
288
289    #[cfg(feature = "source-sqlite")]
290    pub mod sqlite {
291        pub use faucet_source_sqlite::*;
292    }
293
294    #[cfg(feature = "source-s3")]
295    pub mod s3 {
296        pub use faucet_source_s3::*;
297    }
298
299    #[cfg(feature = "source-mongodb")]
300    pub mod mongodb {
301        pub use faucet_source_mongodb::*;
302    }
303
304    #[cfg(feature = "source-mongodb-cdc")]
305    pub mod mongodb_cdc {
306        pub use faucet_source_mongodb_cdc::*;
307    }
308
309    #[cfg(feature = "source-mysql-cdc")]
310    pub mod mysql_cdc {
311        pub use faucet_source_mysql_cdc::*;
312    }
313
314    #[cfg(feature = "source-redis")]
315    pub mod redis {
316        pub use faucet_source_redis::*;
317    }
318
319    #[cfg(feature = "source-webhook")]
320    pub mod webhook {
321        pub use faucet_source_webhook::*;
322    }
323
324    #[cfg(feature = "source-websocket")]
325    pub mod websocket {
326        pub use faucet_source_websocket::*;
327    }
328
329    #[cfg(feature = "source-csv")]
330    pub mod csv {
331        pub use faucet_source_csv::*;
332    }
333
334    #[cfg(feature = "source-elasticsearch")]
335    pub mod elasticsearch {
336        pub use faucet_source_elasticsearch::*;
337    }
338
339    #[cfg(feature = "source-kafka")]
340    pub mod kafka {
341        pub use faucet_source_kafka::*;
342    }
343
344    #[cfg(feature = "source-kinesis")]
345    pub mod kinesis {
346        pub use faucet_source_kinesis::*;
347    }
348
349    #[cfg(feature = "source-bigquery")]
350    pub mod bigquery {
351        pub use faucet_source_bigquery::*;
352    }
353
354    #[cfg(feature = "source-snowflake")]
355    pub mod snowflake {
356        pub use faucet_source_snowflake::*;
357    }
358
359    #[cfg(feature = "source-spanner")]
360    pub mod spanner {
361        pub use faucet_source_spanner::*;
362    }
363
364    #[cfg(feature = "source-parquet")]
365    pub mod parquet {
366        pub use faucet_source_parquet::*;
367    }
368
369    #[cfg(feature = "source-delta")]
370    pub mod delta {
371        pub use faucet_source_delta::*;
372    }
373
374    #[cfg(feature = "source-databricks")]
375    pub mod databricks {
376        pub use faucet_source_databricks::*;
377    }
378
379    #[cfg(feature = "source-gcs")]
380    pub mod gcs {
381        pub use faucet_source_gcs::*;
382    }
383
384    #[cfg(feature = "source-mssql-cdc")]
385    pub mod mssql_cdc {
386        pub use faucet_source_mssql_cdc::*;
387    }
388
389    #[cfg(feature = "source-redshift")]
390    pub mod redshift {
391        pub use faucet_source_redshift::*;
392    }
393
394    #[cfg(feature = "source-pubsub")]
395    pub mod pubsub {
396        pub use faucet_source_pubsub::*;
397    }
398
399    #[cfg(feature = "source-clickhouse")]
400    pub mod clickhouse {
401        pub use faucet_source_clickhouse::*;
402    }
403
404    #[cfg(feature = "source-azure-blob")]
405    pub mod azure_blob {
406        pub use faucet_source_azure_blob::*;
407    }
408
409    #[cfg(feature = "source-singer")]
410    pub mod singer {
411        pub use faucet_source_singer::*;
412    }
413}
414
415// Backwards-compatible flat re-exports for existing users who depend on
416// `faucet-stream::{RestStream, Auth, ...}` without the `source::rest::` path.
417#[cfg(feature = "source-rest")]
418pub use faucet_source_rest::{
419    Auth, DEFAULT_EXPIRY_RATIO, DEFAULT_TOKEN_ENDPOINT_EXPIRY_RATIO, PaginationStyle,
420    ResponseValidator, RestStream, RestStreamConfig, fetch_oauth2_token, fetch_token_from_endpoint,
421};
422
423// ── Sink connectors ──────────────────────────────────────────────────────────
424
425pub mod sink {
426    #[cfg(feature = "sink-bigquery")]
427    pub mod bigquery {
428        pub use faucet_sink_bigquery::*;
429    }
430
431    #[cfg(feature = "sink-iceberg")]
432    pub mod iceberg {
433        pub use faucet_sink_iceberg::*;
434    }
435
436    #[cfg(feature = "sink-postgres")]
437    pub mod postgres {
438        pub use faucet_sink_postgres::*;
439    }
440
441    #[cfg(feature = "sink-jsonl")]
442    pub mod jsonl {
443        pub use faucet_sink_jsonl::*;
444    }
445
446    #[cfg(feature = "sink-snowflake")]
447    pub mod snowflake {
448        pub use faucet_sink_snowflake::*;
449    }
450
451    #[cfg(feature = "sink-mysql")]
452    pub mod mysql {
453        pub use faucet_sink_mysql::*;
454    }
455
456    #[cfg(feature = "sink-mssql")]
457    pub mod mssql {
458        pub use faucet_sink_mssql::*;
459    }
460
461    #[cfg(feature = "sink-sqlite")]
462    pub mod sqlite {
463        pub use faucet_sink_sqlite::*;
464    }
465
466    #[cfg(feature = "sink-s3")]
467    pub mod s3 {
468        pub use faucet_sink_s3::*;
469    }
470
471    #[cfg(feature = "sink-mongodb")]
472    pub mod mongodb {
473        pub use faucet_sink_mongodb::*;
474    }
475
476    #[cfg(feature = "sink-redis")]
477    pub mod redis {
478        pub use faucet_sink_redis::*;
479    }
480
481    #[cfg(feature = "sink-csv")]
482    pub mod csv {
483        pub use faucet_sink_csv::*;
484    }
485
486    #[cfg(feature = "sink-elasticsearch")]
487    pub mod elasticsearch {
488        pub use faucet_sink_elasticsearch::*;
489    }
490
491    #[cfg(feature = "sink-http")]
492    pub mod http {
493        pub use faucet_sink_http::*;
494    }
495
496    #[cfg(feature = "sink-stdout")]
497    pub mod stdout {
498        pub use faucet_sink_stdout::*;
499    }
500
501    #[cfg(feature = "sink-kafka")]
502    pub mod kafka {
503        pub use faucet_sink_kafka::*;
504    }
505
506    #[cfg(feature = "sink-kinesis")]
507    pub mod kinesis {
508        pub use faucet_sink_kinesis::*;
509    }
510
511    #[cfg(feature = "sink-spanner")]
512    pub mod spanner {
513        pub use faucet_sink_spanner::*;
514    }
515
516    #[cfg(feature = "sink-parquet")]
517    pub mod parquet {
518        pub use faucet_sink_parquet::*;
519    }
520
521    #[cfg(feature = "sink-delta")]
522    pub mod delta {
523        pub use faucet_sink_delta::*;
524    }
525
526    #[cfg(feature = "sink-gcs")]
527    pub mod gcs {
528        pub use faucet_sink_gcs::*;
529    }
530
531    #[cfg(feature = "sink-redshift")]
532    pub mod redshift {
533        pub use faucet_sink_redshift::*;
534    }
535
536    #[cfg(feature = "sink-pubsub")]
537    pub mod pubsub {
538        pub use faucet_sink_pubsub::*;
539    }
540
541    #[cfg(feature = "sink-clickhouse")]
542    pub mod clickhouse {
543        pub use faucet_sink_clickhouse::*;
544    }
545
546    #[cfg(feature = "sink-azure-blob")]
547    pub mod azure_blob {
548        pub use faucet_sink_azure_blob::*;
549    }
550}
551
552// ── GCS common types ─────────────────────────────────────────────────────────
553
554#[cfg(any(feature = "source-gcs", feature = "sink-gcs"))]
555pub mod common_gcs {
556    pub use faucet_common_gcs::*;
557}
558
559// ── Kafka common types ───────────────────────────────────────────────────────
560
561#[cfg(any(feature = "source-kafka", feature = "sink-kafka"))]
562pub mod common_kafka {
563    pub use faucet_common_kafka::*;
564}
565
566/// Shared AWS Kinesis types (credentials enum, client builder), re-exported
567/// for library callers when either Kinesis connector is enabled.
568#[cfg(any(feature = "source-kinesis", feature = "sink-kinesis"))]
569pub mod common_kinesis {
570    pub use faucet_common_kinesis::*;
571}
572
573/// Shared Cloud Spanner types (credentials enum, connection block, value
574/// conversion), re-exported for library callers when either Spanner
575/// connector is enabled.
576#[cfg(any(feature = "source-spanner", feature = "sink-spanner"))]
577pub mod common_spanner {
578    pub use faucet_common_spanner::*;
579}
580
581/// Shared Amazon Redshift types (credentials enum, connection block, pool
582/// builder), re-exported when either Redshift connector is enabled.
583#[cfg(any(feature = "source-redshift", feature = "sink-redshift"))]
584pub mod common_redshift {
585    pub use faucet_common_redshift::*;
586}
587
588/// Shared Google Cloud Pub/Sub types (credentials enum, connection block,
589/// client builder), re-exported when either Pub/Sub connector is enabled.
590#[cfg(any(feature = "source-pubsub", feature = "sink-pubsub"))]
591pub mod common_pubsub {
592    pub use faucet_common_pubsub::*;
593}
594
595/// Shared ClickHouse types (connection block, HTTP client builder), re-exported
596/// when either ClickHouse connector is enabled.
597#[cfg(any(feature = "source-clickhouse", feature = "sink-clickhouse"))]
598pub mod common_clickhouse {
599    pub use faucet_common_clickhouse::*;
600}
601
602/// Shared Azure Blob / ADLS Gen2 types (credentials enum, object-store builder),
603/// re-exported when either Azure Blob connector is enabled.
604#[cfg(any(feature = "source-azure-blob", feature = "sink-azure-blob"))]
605pub mod common_azure {
606    pub use faucet_common_azure::*;
607}
608
609// ── State-store backends ─────────────────────────────────────────────────────
610
611pub mod state {
612    #[cfg(feature = "state-redis")]
613    pub mod redis {
614        pub use faucet_state_redis::*;
615    }
616
617    #[cfg(feature = "state-postgres")]
618    pub mod postgres {
619        pub use faucet_state_postgres::*;
620    }
621}
622
623// ── Lineage (OpenLineage emission) ───────────────────────────────────────────
624/// OpenLineage event emission for pipeline runs (enable the `lineage` feature;
625/// `lineage-kafka` adds the Kafka transport). The CLI wires this automatically
626/// from a `lineage:` config block; library callers can build a
627/// [`lineage::LineageEmitter`] directly.
628#[cfg(feature = "lineage")]
629pub use faucet_lineage as lineage;
630
631// ── SQL transform (embedded DuckDB) ──────────────────────────────────────────
632/// SQL-as-transform: run DuckDB SQL over each pipeline page (the `batch`
633/// relation). Enable the `transform-sql` feature. The CLI wires this via the
634/// `sql` transform; library callers build [`transform_sql::SqlTransform`] and
635/// attach it with [`TransformingSource`].
636#[cfg(feature = "transform-sql")]
637pub use faucet_transform_sql as transform_sql;