1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
//! High-performance, at-least-once ETL pipeline framework.
//!
//! `spate` is the facade crate for the Spate framework and the only crate
//! applications need to depend on. The engine ([`spate_core`]) is re-exported
//! at the root; connectors are enabled through cargo features:
//!
//! | Feature | Enables |
//! |---|---|
//! | `kafka` | [`kafka`] — Kafka source built on `rdkafka` (single consumer, partition-queue lanes) |
//! | `s3` | [`s3`] — coordinated bounded object-storage (S3) backfill source: the fleet's leader plans the prefix into splits, workers lease them through an injected coordinator (solo in-process by default; progress is ephemeral without a durable store), and the job self-terminates once the plan is final and every split completes (`refresh_listing` keeps it open) |
//! | `clickhouse` | [`clickhouse`] — ClickHouse sink (RowBinary, dedup tokens, replica rotation) |
//! | `clickhouse-uuid` | `uuid::Uuid` fields for `UUID` columns (`clickhouse::serde::uuid`) |
//! | `clickhouse-chrono` | `chrono` fields for `Date`/`DateTime`/`DateTime64`/`Time` columns |
//! | `clickhouse-time` | `time` crate fields for the same date/time columns |
//! | `clickhouse-rust-decimal` | `rust_decimal::Decimal` conversions for `Decimal` columns |
//! | `avro` | [`avro`] — Avro deserialization (Confluent wire format, schema registry) |
//! | `json` | [`json`] — JSON deserialization (single-document, NDJSON, top-level array) |
//! | `json-float-roundtrip`, `json-arbitrary-precision`, `json-raw-value` | Opt-in `serde_json` fidelity knobs for [`json`]. Not part of `full` — `arbitrary-precision` is crate-wide |
//! | `coordination` | [`coordination`] backend — multi-instance leader-assigned work distribution for broker-less sources: the protocol, `StoreCoordinator`, and the in-memory store (zero-infrastructure embedding). The seam types and the `CoordinationDriver` live in `spate::coordination` without any feature |
//! | `coordination-nats` | The production NATS JetStream KV store (server >= 2.11) on top of [`coordination`]; pulls the async-nats dependency tree |
//! | `datagen` | [`datagen`] — synthetic storefront-event source (orders, payments, refunds) for a pipeline that needs no broker, bucket or coordination store. A demo and test source: it keeps no durable progress, and is not part of `full` |
//! | `datagen-avro` | Avro payloads from [`datagen`] instead of JSON (implies `datagen`) |
//! | `full` | All connectors (`avro`, `json`, `kafka`, `clickhouse`, `coordination-nats`, `s3`) |
//!
//! # Anatomy of a pipeline
//!
//! One process runs one pipeline (the [user guide] carries the full
//! architecture and its rationale):
//!
//! [user guide]: https://spate.kainth.dev/docs/user-guide/
//!
//! ```text
//! ┌───────────────────────────────────────────────┐
//! pinned std thread │ lane.poll → deserialize (borrowed) → operator │──try_send──▶ per-shard
//! (× N) │ chain (map/filter/flat_map, monomorphized) │ bounded queues
//! └───────────────────────────────────────────────┘ │
//! ▲ full? pause lanes, keep polling ▼
//! │ ┌───────────────────────────┐
//! ┌──────────┐ acks (never block) │ sink workers: merge chunks,│
//! │ source │◀───────────────────────────────────── │ seal batches, rotate │
//! │ control │ watermarks → store/commit │ replicas, retry; admin │
//! └──────────┘ │ server (/metrics, probes) │
//! └───────────────────────────┘
//! ```
//!
//! Delivery is **at-least-once**: a batch's offsets commit only after every
//! record derived from it was durably written (or intentionally dropped).
//! Duplicates are possible after a crash, so design target tables to
//! tolerate replays.
//!
//! # A minimal pipeline
//!
//! Operators are stateful closures chained into a graph; the YAML
//! carries tuning and connector configuration; [`pipeline::Pipeline`]
//! assembles and runs the process. This compiles and runs against the
//! `spate-test` mocks; swap in `KafkaSource::from_component_config` and a
//! ClickHouse sink for the production version (see
//! `examples/kafka_avro_to_clickhouse.rs`):
//!
//! ```
//! use spate::prelude::*;
//! use spate_test::{BytesPassthrough, TestEncoder, capture_sink, memory_source};
//!
//! # fn main() -> Result<(), Box<dyn std::error::Error>> {
//! let config = PipelineConfig::from_str(
//! "pipeline: { name: demo, threads: 1, io_threads: 1 }\n\
//! admin: { listen: none }\n\
//! checkpoint: { interval: 100ms }\n\
//! metrics: { exporter: none }\n\
//! source: { memory: {} }\n\
//! sink: { capture: {} }",
//! )?;
//! let (source, handle) = memory_source();
//! let (sink, script) = capture_sink(1, 1);
//!
//! let runtime = Pipeline::from_config(config)?
//! .sink(sink)?
//! .chains(|ctx| {
//! // The sink's YAML `chunk:` block (or the default), bound before
//! // `with_metrics` moves `ctx.pipeline`.
//! let chunk_cfg = ctx.chunk();
//! chain_owned::<Vec<u8>, _>(BytesPassthrough)
//! .with_metrics(ctx.pipeline, "main")
//! .filter(|payload: &Vec<u8>| !payload.is_empty())
//! .sink(TestEncoder, KeyHashRouter, chunk_cfg,
//! ctx.queues, ctx.budget)
//! .build()
//! })
//! .runtime_options(RuntimeOptions { handle_signals: false, ..Default::default() })
//! .into_runtime(source)?;
//!
//! // Drive it: assign a lane, produce, wait for the durable commit.
//! let shutdown = runtime.shutdown_handle();
//! let join = std::thread::spawn(move || runtime.run());
//! handle.assign_lanes(&[(spate::source::LaneId(0), PartitionId(0))]);
//! let last = handle.push(PartitionId(0), None, b"hello");
//! while handle.last_committed(PartitionId(0)) != Some(last + 1) {
//! std::thread::sleep(std::time::Duration::from_millis(5));
//! }
//! shutdown.trigger();
//! let report = join.join().expect("join")?;
//! assert_eq!(report.exit_code(), 0);
//! assert!(!script.writes().is_empty());
//! # Ok(())
//! # }
//! ```
//!
//! # Where things live
//!
//! - [`ops`] — the chain builder and operator combinators.
//! - [`source`] / [`sink`] — the connector traits ([`source::Source`],
//! [`source::SourceLane`], [`sink::RowEncoder`], [`sink::ShardWriter`])
//! and the framework-owned sink pool.
//! - [`pipeline`] — the runtime: pinned threads, controller, shutdown.
//! - [`checkpoint`] / [`backpressure`] — acknowledgments and flow control.
//! - [`config`] — YAML with `${VAR:-default}` interpolation and opaque
//! per-connector sections.
//! - [`metrics`] / [`admin`] / [`telemetry`] — observability. The
//! [`metrics`](https://crates.io/crates/metrics) facade is the
//! instrumentation API, and [the metrics reference] carries the taxonomy.
//! - Testing your pipelines: the `spate-test` crate (in-memory sources and
//! sinks with scripting handles).
//!
//! [the metrics reference]: https://spate.kainth.dev/docs/METRICS
pub use *;
/// Curated imports for pipeline assemblies: one `use spate::prelude::*;`
/// brings in the builder, chain entry points, and the types every
/// assembly touches.
///
/// Connector constructors are excluded; import those from their
/// feature-gated modules ([`kafka`], [`clickhouse`], [`avro`]).
/// Additions to this module are semver-additive; nothing is ever removed.
/// Avro deserialization support (Confluent wire format, schema registry).
pub use spate_avro as avro;
/// JSON deserialization support (single-document, NDJSON, top-level array).
pub use spate_json as json;
/// Kafka source and producer-sink connector.
pub use spate_kafka as kafka;
/// ClickHouse sink connector.
pub use spate_clickhouse as clickhouse;
/// Bounded object-storage (S3) backfill source.
pub use spate_s3 as s3;
/// Synthetic storefront-event source for a pipeline with no prerequisites.
/// A demo and test source; it keeps no durable progress.
pub use spate_datagen as datagen;
// Without the feature, `spate::coordination` resolves to the seam module
// (traits, events, the driver) through the `spate_core::*` glob above. With
// it, this explicit re-export shadows the glob with the backend crate,
// which itself re-exports the same seam types, so the one path serves both
// worlds and the backend types appear alongside.
/// Multi-instance leader-assigned coordination backend (NATS JetStream KV
/// + in-memory store). The seam and driver need no feature.
pub use spate_coordination as coordination;
/// Compile-checks the README's example, so it cannot drift from the builder
/// API.
///
/// `cfg(doctest)` is stripped before macro expansion, so `include_str!` never
/// runs in an ordinary build. The README escaping this crate's directory
/// cannot break `cargo publish`, and the file is not appended to the rendered
/// documentation.
;