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
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
//! # Krafka
//!
//! A pure Rust, async-native Apache Kafka client.
//!
//! Krafka provides high-performance, safe, and idiomatic Rust APIs for
//! producing and consuming messages from Apache Kafka clusters.
//!
//! ## Features
//!
//! - **Pure Rust by default**: No librdkafka or C bindings; the optional `zstd`
//! compression feature links against `zstd-sys` and requires a C toolchain
//! - **Async-native**: Built on Tokio for non-blocking I/O
//! - **High-performance**: Zero-copy buffers, minimal allocations
//! - **Safe**: No unsafe code by default
//! - **Cloud-native**: First-class AWS MSK support including IAM auth
//!
//! ## Thread Safety
//!
//! All main types in Krafka implement `Send + Sync`:
//!
//! - [`Producer`](producer::Producer) - can be shared across tasks with `Arc`
//! - [`Consumer`](consumer::Consumer) - can be shared across tasks with `Arc`
//! - [`AdminClient`](admin::AdminClient) - can be shared across tasks with `Arc`
//!
//! This allows safe concurrent access from multiple Tokio tasks:
//!
//! ```rust,no_run
//! use std::sync::Arc;
//! use krafka::producer::Producer;
//!
//! # async fn example() -> Result<(), krafka::error::KrafkaError> {
//! let producer = Arc::new(Producer::builder()
//! .bootstrap_servers("localhost:9092")
//! .build()
//! .await?);
//!
//! // Spawn multiple tasks sharing the producer
//! for i in 0..10 {
//! let producer = producer.clone();
//! tokio::spawn(async move {
//! let _ = producer.send("topic", None, Some(b"message")).await;
//! });
//! }
//! # Ok(())
//! # }
//! ```
//!
//! ## Quick Start
//!
//! ### Producer
//!
//! ```rust,no_run
//! use krafka::producer::Producer;
//!
//! # async fn example() -> Result<(), krafka::error::KrafkaError> {
//! let producer = Producer::builder()
//! .bootstrap_servers("localhost:9092")
//! .build()
//! .await?;
//!
//! producer.send("my-topic", Some(b"key"), Some(b"value")).await?;
//!
//! // A `None` value is a tombstone: on a compacted topic it deletes the key.
//! producer.send("my-topic", Some(b"key"), None).await?;
//! # Ok(())
//! # }
//! ```
//!
//! ### Consumer
//!
//! ```rust,no_run
//! use krafka::consumer::Consumer;
//!
//! # async fn example() -> Result<(), krafka::error::KrafkaError> {
//! let consumer = Consumer::builder()
//! .bootstrap_servers("localhost:9092")
//! .group_id("my-group")
//! .build()
//! .await?;
//!
//! consumer.subscribe(&["my-topic"]).await?;
//!
//! loop {
//! match consumer.recv().await {
//! Ok(msg) => println!("{:?}", msg),
//! Err(krafka::RecvError::Closed) => break,
//! Err(krafka::RecvError::Error(e)) => return Err(e),
//! Err(_) => break,
//! }
//! }
//! # Ok(())
//! # }
//! ```
//!
//! ## Cargo Features
//!
//! | Feature | Default | Description |
//! |---------|---------|-------------|
//! | `compression` | **yes** | Enables pure-Rust compression codecs (`gzip` + `snappy` + `lz4`). |
//! | `compression-all` | no | Enables all compression codecs, including `zstd`. |
//! | `gzip` | via `compression` | Gzip record batch compression via `flate2`. |
//! | `snappy` | via `compression` | Snappy compression via `snap`. |
//! | `lz4` | via `compression` | LZ4 compression via `lz4_flex`. |
//! | `zstd` | no | Zstd compression via `zstd` (requires C toolchain). |
//! | `aws-msk` | no | AWS MSK IAM authentication with SDK credential chain. |
//! | `oauth-oidc` | no | Built-in OIDC token provider for SASL/OAUTHBEARER: the `client_credentials` grant (KIP-768) and RFC 7523 client assertions (KIP-1258). Adds no cryptography dependency — assertions are supplied pre-signed. |
//! | `socks5` | no | SOCKS5 proxy support via `tokio-socks`. |
//! | `telemetry` | no | OpenTelemetry exporter for producer/consumer metrics. |
//! | `unstable-protocol` | no | Enables protocol versions Kafka marks `latestVersionUnstable` — a released broker does not advertise them without `unstable.api.versions.enable=true`. Covers `ApiVersions` v5 (KIP-1242), `InitProducerId` v6 (KIP-939) and the Share Consumer (KIP-932). APIs under this feature may change without semver notice. |
//! | `ring` | **yes** | rustls crypto backend using `ring` (pure Rust). |
//! | `rustls-aws-lc-rs` | no | rustls crypto backend using `aws-lc-rs`. Preferred on AWS Graviton and for FIPS deployments. |
//! | `native-tls-roots` | no | Load platform-native root certificates via `rustls-native-certs`. |
//! | `test-broker` | no | In-process fake Kafka broker for testing your own code against a real client. Not for production builds. |
//!
//! ## TLS crypto backend
//!
//! Exactly one rustls crypto backend is used at runtime. `ring` is the default;
//! `rustls-aws-lc-rs` selects `aws-lc-rs` instead.
//!
//! These features are **additive**, as Cargo requires: if two crates in your
//! dependency graph each select a different backend, the build still succeeds
//! and `rustls-aws-lc-rs` wins. To pin a specific backend regardless of what
//! your dependency graph enabled, install it as the process default before
//! constructing any krafka client:
//!
//! ```rust,ignore
//! rustls::crypto::ring::default_provider().install_default().ok();
//! ```
//!
//! Enabling **neither** backend is a compile error — rustls cannot build a
//! `ClientConfig` without a crypto provider.
//!
//! To disable the default features and pick only what you need, remember that
//! `default-features = false` also drops `ring`:
//!
//! ```toml
//! [dependencies]
//! # `ring` (or `rustls-aws-lc-rs`) is required — without it the build fails.
//! cargo add krafka --no-default-features --features lz4,ring
//! ```
// Cargo features must be *additive*: if crate A depends on krafka with `ring`
// and crate B depends on krafka with `rustls-aws-lc-rs`, Cargo unifies the two
// feature sets and both are enabled. Rejecting that combination at compile
// time would make the two crates impossible to use together, with no recourse
// for the application author. Instead the backends are additive and
// `rustls-aws-lc-rs` deterministically wins when both are active — see
// `auth::tls::resolve_crypto_provider`.
//
// At least one backend is required, because rustls cannot construct a
// `ClientConfig` without a crypto provider.
compile_error!;
// krafka's metrics layer relies on 64-bit atomic operations (AtomicU64).
// 32-bit targets without hardware AtomicU64 support (e.g. Cortex-M3) are not
// supported. Fail fast with a clear diagnostic rather than a confusing link
// error or silent correctness bug.
compile_error!;
/// Tracks in-flight operations so a flush or a close can wait for work that has
/// started but not yet reached the structure it will land in.
///
/// Shared by the producer (draining a transaction before `EndTxn`) and the
/// share consumer (draining acknowledgements before `close`), which is why it
/// sits at the crate root rather than inside either.
/// Minimal async HTTP/1.1 client used by the OIDC token provider.
///
/// Compiled only when `oauth-oidc` is enabled. This is an implementation
/// detail and not part of the stable public API.
/// Cluster metadata cache and refresh logic.
///
/// This is an implementation detail of the consumer and producer. Types are
/// accessible for advanced use but are **not** part of the stable public API.
/// Network connection pool and transport layer.
///
/// This is an implementation detail. Types are accessible for advanced use
/// (e.g. custom authentication) but are **not** part of the stable public API.
/// Kafka wire-protocol encode/decode layer.
///
/// This is an implementation detail. Types are accessible for advanced use
/// (e.g. benchmarks, raw record batch construction) but are **not** part of
/// the stable public API.
/// In-process fake Kafka broker for deterministic client tests.
///
/// Enabled by the `test-broker` feature. Not compiled into production builds.
/// The types you need to write a producer, a consumer or an admin client.
///
/// ```rust
/// use krafka::prelude::*;
/// ```
///
/// A glob import, so an explicit `use` of the same type shadows it rather than
/// colliding. Everything here is re-exported from its own module; the prelude
/// exists so that the common case is one line instead of eight, and so that the
/// documentation snippets have a single import that keeps them honest.
///
/// Deliberately excluded: anything feature-gated (`ShareConsumer`, the
/// telemetry reporter, the fake broker), and anything a caller is unlikely to
/// name more than once per program (`ConnectionConfig`, the metrics exporters).
/// A prelude that pulls in names you did not ask for is worse than no prelude.
pub use ;
pub use MetadataRecoveryStrategy;
// Re-export user-facing protocol types at a stable path so callers do not
// need to reach into the hidden `protocol` module.
pub use ;
// Re-export user-facing network/auth types at a stable path so callers do
// not need to reach into the hidden `network` module.
pub use ;
/// Kafka protocol API version.
pub type ApiVersion = i16;
/// Kafka correlation ID for request/response matching.
pub type CorrelationId = i32;
/// Kafka partition ID.
pub type PartitionId = i32;
/// Kafka broker ID.
pub type BrokerId = i32;
/// Kafka offset.
pub type Offset = i64;
/// Kafka timestamp (milliseconds since epoch).
pub type Timestamp = i64;