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