Skip to main content

oauth/
oauth.rs

1//! Produce one record with SASL OAUTHBEARER (unsecured JWT or OIDC token).
2//!
3//! Broker on `KAFKA_BOOTSTRAP`. Topic `KAFKA_TOPIC` (default `partitionline`).
4//!
5//! Unsecured JWT (local / test brokers that accept librdkafka-style tokens):
6//!
7//! ```text
8//! SASL_OAUTH_PRINCIPAL=alice cargo run --example oauth
9//! ```
10//!
11//! OIDC client-credentials (token URL must be reachable from this process):
12//!
13//! ```text
14//! OIDC_TOKEN_URL=https://issuer.example/oauth/token \
15//! OIDC_CLIENT_ID=… OIDC_CLIENT_SECRET=… \
16//! cargo run --example oauth
17//! ```
18//!
19//! Set `TLS_CA_PEM` (and optionally `TLS_SERVER_NAME`) for SASL_SSL — see
20//! `scripts/ci-auth-smoke.sh`.
21
22use partitionline::{OidcConfig, ProduceRecord, Producer, ProducerConfig, Sasl, TlsConfig};
23
24#[tokio::main]
25async fn main() -> partitionline::Result<()> {
26    let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
27    let topic = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
28
29    let sasl = if let (Ok(url), Ok(id), Ok(secret)) = (
30        std::env::var("OIDC_TOKEN_URL"),
31        std::env::var("OIDC_CLIENT_ID"),
32        std::env::var("OIDC_CLIENT_SECRET"),
33    ) {
34        Sasl::oidc(OidcConfig::new(url, id, secret))
35    } else {
36        let principal = std::env::var("SASL_OAUTH_PRINCIPAL").unwrap_or_else(|_| "alice".into());
37        Sasl::oauthbearer(principal)
38    };
39
40    let mut cfg = ProducerConfig::bootstrap([bootstrap]).sasl(sasl);
41    if let Ok(ca_path) = std::env::var("TLS_CA_PEM") {
42        let mut tls =
43            TlsConfig::default().ca_pem(tokio::fs::read(&ca_path).await.map_err(|e| {
44                partitionline::Error::protocol(format!("read TLS_CA_PEM {ca_path}: {e}"))
45            })?);
46        if let Ok(name) = std::env::var("TLS_SERVER_NAME") {
47            if !name.is_empty() {
48                tls = tls.server_name(name);
49            }
50        }
51        cfg = cfg.tls(tls);
52    }
53
54    let producer = Producer::new(cfg).await?;
55    let md = producer
56        .send(ProduceRecord::to(topic).value(&b"hello over oauthbearer"[..]))
57        .await?;
58    println!("{}-{}@{}", md.topic, md.partition, md.offset);
59    producer.close().await?;
60    Ok(())
61}