Skip to main content

sasl/
sasl.rs

1//! Produce one record over SASL (PLAIN, SCRAM-SHA-256, or SCRAM-SHA-512).
2//!
3//! Broker on `KAFKA_BOOTSTRAP`. Topic `KAFKA_TOPIC` (default `partitionline`).
4//! `SASL_MECHANISM` defaults to `SCRAM-SHA-256`. `KAFKA_USERNAME` /
5//! `KAFKA_PASSWORD` default to `alice` / `secret`.
6//!
7//! Set `TLS_CA_PEM` (and optionally `TLS_SERVER_NAME`) to use SASL over TLS
8//! (`SASL_SSL`) — the common production path. See `scripts/ci-auth-smoke.sh`.
9//!
10//! For OAUTHBEARER / OIDC, see `examples/oauth.rs`.
11
12use partitionline::{ProduceRecord, Producer, ProducerConfig, Sasl, TlsConfig};
13
14#[tokio::main]
15async fn main() -> partitionline::Result<()> {
16    let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
17    let topic = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
18    let username = std::env::var("KAFKA_USERNAME").unwrap_or_else(|_| "alice".into());
19    let password = std::env::var("KAFKA_PASSWORD").unwrap_or_else(|_| "secret".into());
20    let mechanism = std::env::var("SASL_MECHANISM").unwrap_or_else(|_| "SCRAM-SHA-256".into());
21    let sasl = match mechanism.as_str() {
22        "PLAIN" => Sasl::plain(username, password),
23        "SCRAM-SHA-256" => Sasl::scram_sha256(username, password),
24        "SCRAM-SHA-512" => Sasl::scram_sha512(username, password),
25        other => {
26            return Err(partitionline::Error::protocol(format!(
27                "unsupported SASL_MECHANISM {other}; use PLAIN, SCRAM-SHA-256, or SCRAM-SHA-512"
28            )));
29        }
30    };
31
32    let mut cfg = ProducerConfig::bootstrap([bootstrap]).sasl(sasl);
33    if let Ok(ca_path) = std::env::var("TLS_CA_PEM") {
34        let mut tls =
35            TlsConfig::default().ca_pem(tokio::fs::read(&ca_path).await.map_err(|e| {
36                partitionline::Error::protocol(format!("read TLS_CA_PEM {ca_path}: {e}"))
37            })?);
38        if let Ok(name) = std::env::var("TLS_SERVER_NAME") {
39            if !name.is_empty() {
40                tls = tls.server_name(name);
41            }
42        }
43        cfg = cfg.tls(tls);
44    }
45
46    let producer = Producer::new(cfg).await?;
47    let md = producer
48        .send(ProduceRecord::to(topic).value(&b"hello over sasl"[..]))
49        .await?;
50    println!("{}-{}@{}", md.topic, md.partition, md.offset);
51    producer.close().await?;
52    Ok(())
53}