1use 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}